effectmq
Reference

TaskQueue reference

Queue descriptors, offer and completion operations, durable handles, waiting, and events.

The TaskQueue namespace is the high-level queue API. It binds a stable queue name to one TaskDefinition and performs typed operations through TaskEngine.

TaskQueue.make

make(name, taskDefinition): TaskQueue

Creates a pure queue descriptor. Redis state is created by the first operation that needs it.

TaskQueue.offer

offer(queue, payload, options?): Effect<OfferOutcome, OfferError, OfferRequirements>

Encodes and enqueues a payload. Returns:

{
  _tag: "TaskCreated" | "TaskExisting"
  task: Task
  handle: TaskHandle
}

Before encoding or Redis access, offer checks the task definition's retry, storage, and retention invariants. Invalid definition configuration is a defect, not an OfferError.

Offer options

OptionTypeDefaultDescription
taskIdstringdefinition keyExplicit identity override.
delaynumber0Milliseconds before acquisition eligibility.
maxRetriesnumberdefinition capPer-generation handler retry override.
maxStalledCountnumber1Expired lease recoveries allowed before terminal stalled failure.
onSuccessPolicyCompletionPolicy"delete"Record/index policy after success.
onFailurePolicyCompletionPolicy"delete"Record/index policy after terminal failure.
retainResultUntil"current-task-settles"unsetAdds a result-retention hold from the current managed handler.
onDuplicate"return-existing" | "new-generation""return-existing"Duplicate identity behavior.

Completion policies are delete, keep, mark-as-success, and mark-as-failure. delete removes the task record when no retention hold remains; the terminal result follows its separate result retention. keep retains the record without terminal-index membership. The two mark policies retain the record in the named terminal index.

TaskHandle

A handle identifies exactly one durable generation:

FieldType
_tag"TaskHandle"
queuestring
taskIdstring
generationnumber
cursorstring
taskNamestring
schemaIdstring
protocolVersion1

The phantom success and error members retain the handle's result types.

TaskQueue.complete

complete(queue, handler): Effect<string, CompleteError, CompleteRequirements>
complete(handler)(queue): Effect<string, CompleteError, CompleteRequirements>

Polls until one task is available, supervises its lease, invokes the handler, persists the handler outcome, and returns the task id. A typed handler failure is stored and routed through retry/failure policy; it does not fail the completion Effect.

Definition invariants are checked before acquisition. An invalid definition is a defect, not a CompleteError.

The handler type is:

(task: Task, context: TaskHistory.Context<Progress>) => Effect<Success, TypedFailure | ProgressWriteError, HandlerRequirements>

TaskQueue.completeOne

completeOne(queue, handler, processing?): Effect<boolean, CompleteOneError, CompleteRequirements>

Attempts one non-polling acquisition. Returns false when no task is currently available and true after an attempt is processed. Managed workers use this operation internally.

CompleteOneError adds ProcessingConfigurationError to CompleteError for invalid lease or heartbeat options.

Processing options

OptionDefault
lockTimeout30 seconds
lockRefresh10 seconds
heartbeatRetryDelay250 milliseconds
heartbeatRetryCount3

The heartbeat retry count is reduced when necessary to fit inside the lease safety window.

TaskQueue.wait

wait(queue, handle, options?): Effect<Success, WaitError, WaitRequirements>

Waits for the exact generation in the handle. It reads durable state, subscribes from the handle cursor, and rechecks state after subscription.

Retriable failures do not end the wait. A terminal failure is wrapped in TaskFailed.failure. A handle for a different queue or task name fails with TaskHandleMismatch; schema identity is checked separately.

options.timeout accepts Duration.Input. A timeout fails with CallerTimeout and does not cancel queue work.

TaskQueue.execute

execute(queue, payload, offerOptions?): Effect<Success, ExecuteError, ExecuteRequirements>

Equivalent to offer followed by wait on the returned handle. It has no wait timeout option.

TaskQueue.stream

stream(queue, { cursor?, pollInterval? }): Stream<Event, ...>

Defaults to the current stream position, so an omitted cursor observes future events only. pollInterval accepts Duration.Duration and defaults to a two-second blocking Redis read; it is not a delay between non-blocking polls.

Event tagTyped payload
task.createdNew decoded task and execution state.
task.updatedExisting and replacement decoded tasks plus state.
task.failedTyped failure, policy, retry time, failure kind, attempt, and terminal flag.
task.completedTyped success value and policy.
task.movedPrevious/new list and state plus attempt counters.

All events contain id, taskId, generation, protocolVersion, and schemaId.

Requirements and failures

OfferRequirements, CompleteRequirements, WaitRequirements, and ExecuteRequirements expose the exact Effect service unions for generic programs. OfferError, CompleteError, WaitError, and ExecuteError expose their exact failure unions. See the error reference.

TaskQueue.readEvents

readEvents(queue, handle, { after?, limit? }): Effect<TaskHistory.Page<Progress>, ...>

Reads the exact generation's typed progress and compact lifecycle entries in append order. Requires TaskEngine and the progress schema's decoding services. The page contains entries, cursor, hasMore, and truncated. Each entry contains id, taskId, generation, attempt, timestamp: Date, and event: { _tag: "Progress", data } | { _tag: "Lifecycle", data }. Lifecycle data has a task.created, task.updated, task.moved, task.failed, or task.completed tag and transition metadata, without payloads or results.

limit defaults to 100 and is bounded to 1–1,000. Pass the opaque returned cursor as after; an empty page can be polled again. Cursors are bound to queue, task, and generation. Reads are non-destructive and independent for each reader.

History is unlimited unless maxHistoryEntries is configured. Trimmed history sets truncated; missing entries after an explicit cursor cause HistoryCursorExpired with an earliestCursor recovery position. Disabled, unavailable, malformed-cursor, and incompatible/corrupt-data cases fail with typed errors. Task deletion also deletes history, even if its result survives.

context.progress(value) on managed handlers encodes the declared progress schema and returns Effect<string, ProgressWriteError, Progress.EncodingServices>. Unhandled progress failures escape through the operational channel without business-error encoding or success settlement. See the complete example.

On this page