effectmq
How-to guides

Report and read task progress

Emit typed progress during execution and page task-owned history.

Declare progress on a task to enable its Redis history. Enabled generations record both your progress values and compact lifecycle events. Tasks without that declaration do not create a history stream. Use progress: Schema.Never when you want lifecycle history without custom progress.

Run a worker and poll its history

Save this example as progress.ts in a project with @effectmq/core, effect, @effect/platform-node, and tsx installed. Start Redis on port 6379 and run pnpm exec tsx progress.ts. The worker emits updates while the reader polls. Use a fresh requestId for a new run.

import { NodeRuntime } from "@effect/platform-node"
import { Task, TaskEngine, TaskQueue, Worker } from "@effectmq/core"
import { Console, Effect, Fiber, Schema } from "effect"

const Greet = Task.make({
  name: "greet",
  schemaId: "greet/v1",
  payload: { requestId: Schema.String, name: Schema.String },
  success: Schema.String,
  error: Schema.Never,
  progress: Schema.Union([
    Schema.Struct({ _tag: Schema.Literal("Message"), text: Schema.String }),
    Schema.Struct({ _tag: Schema.Literal("Percent"), value: Schema.Number })
  ]),
  idempotencyKey: ({ requestId }) => requestId
})
const greetings = TaskQueue.make("greetings", Greet)
const worker = Worker.make(greetings, ({ payload }, context) =>
  Effect.gen(function* () {
    yield* context.progress({ _tag: "Message", text: "Starting greeting" })
    for (const value of [25, 50, 75, 100]) {
      yield* Effect.sleep("250 millis")
      yield* context.progress({ _tag: "Percent", value })
    }
    return `Hello, ${payload.name}!`
  })
)

const program = Effect.gen(function* () {
  const { handle } = yield* TaskQueue.offer(
    greetings,
    { requestId: "greeting-1", name: "Julia" },
    { onSuccessPolicy: "keep", onFailurePolicy: "keep" }
  )
  const running = yield* Effect.forkChild(Worker.run(worker))
  let after: string | undefined
  let terminal = false
  while (true) {
    const page = yield* TaskQueue.readEvents(greetings, handle, { after, limit: 20 })
    for (const entry of page.entries) {
      yield* Console.log(entry.attempt, entry.event)
      if (entry.event._tag === "Lifecycle" && entry.event.data.terminal) {
        terminal = true
      }
    }
    after = page.cursor
    if (terminal && !page.hasMore) break
    if (!page.hasMore) yield* Effect.sleep("100 millis")
  }
  yield* Console.log(yield* TaskQueue.wait(greetings, handle))
  yield* Fiber.interrupt(running)
})

program.pipe(
  Effect.provide(TaskEngine.layer({ redis: { url: "redis://127.0.0.1:6379" } })),
  NodeRuntime.runMain
)

The handler's second argument provides progress(value), returning the stored event ID. One-argument handlers remain valid. The same context is available in TaskQueue.complete and TaskQueue.completeOne. It belongs to that attempt: lease loss, deadline expiry, settlement, or removal prevents further writes.

Retain history after completion

History has exactly the task record's lifetime. The default completion policy is delete, so it removes history too, even while TaskQueue.wait can still read a separately retained result. Choose keep or a mark policy to retain the record and history until task-record retention expires. Retention holds delay record disposal and preserve history until the final hold is released.

Retries share one history. Entries carry taskId, generation, attempt, timestamp, and id; events before the first acquisition use attempt zero. A replacement generation gets a fresh history. Existing generations retain the enablement and count limit with which they were originally offered.

Bound history and recover from gaps

History has no count limit by default. Omitted or null storageLimits.maxHistoryEntries keeps every entry until task disposal. Set a positive safe integer to retain only the newest entries; lifecycle and custom progress share this exact cap. This is independent of maxEventEntries, which bounds the queue-wide event stream.

Each reader owns its cursor. Pass the returned page.cursor as after to read strictly later entries. An empty page preserves that cursor and can be polled again; it does not mean the task is complete. Pages are live views, not a snapshot across requests. Page size defaults to 100 and must be 1–1,000.

A new reader sees truncated: true if older entries were removed. If entries after an explicit cursor were trimmed, the read fails with TaskHistory.HistoryCursorExpired. Its earliestCursor resumes immediately before the oldest retained entry. Show the gap in your UI before choosing to resume there. A cursor whose own entry was trimmed still works if all later entries remain. With a small cap, terminal lifecycle entries may themselves be trimmed; use TaskQueue.wait to determine the outcome rather than waiting for a particular event forever.

Handle progress failures

Schema, value-size, ownership, and Redis failures are wrapped in TaskHistory.ProgressWriteError, separate from the task's business-error schema. An uncaught progress failure leaves the attempt unsettled for existing lease recovery. You can catch that error explicitly if progress is optional. Interruption remains interruption.

reason: "IndeterminateWrite" means Redis may already have stored the event. The library does not replay that append. Application retries and repeated task attempts can duplicate application progress; include an application event ID in your schema if your UI needs deduplication.

The read API distinguishes disabled history, unavailable generations, expired or invalid cursors, and schema/storage decoding failures. Progress schema encoding services are required only when emitting; decoding services are required when reading.

Deploy the feature

Upgrade every producer, worker, scheduler, and maintenance process before opting any task into history. Old records without history metadata remain disabled. For rollback, first stop creating enabled generations and dispose of existing enabled records through the upgraded engine; old engines cannot maintain or clean up their history. See the repository's upgrade and rollback runbook for the deployment sequence.

On this page