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.