EventQueue reference
Durable application events with named subscriptions, independent acknowledgements, optional expiration, and archival.
EventQueue broadcasts each emitted event to the subscriptions registered at
emission time. Each subscription owes one acknowledgement. Multiple workers
sharing a name share its deliveries; registering that name again is idempotent.
Define and use a queue
import { Effect, Schema } from "effect";
import { EventEngine, EventQueue } from "@effectmq/core";
const orders = EventQueue.make(
"orders",
Schema.Struct({ orderId: Schema.String }),
{ onCompletion: "archive" },
);
const program = Effect.gen(function* () {
const billing = yield* EventQueue.subscribe(orders, "billing");
const email = yield* EventQueue.subscribe(orders, "email");
const event = yield* EventQueue.emit(orders, { orderId: "order-123" });
yield* EventQueue.processOne(orders, billing, (event) =>
Effect.log(`Bill ${event.payload.orderId}`),
);
yield* EventQueue.processOne(orders, email, (event) =>
Effect.log(`Email ${event.payload.orderId}`),
);
return yield* EventQueue.get(orders, event.id);
}).pipe(Effect.provide(EventEngine.layer()));Register recipients before emitting. Later registrations receive future events only. An event with zero recipients completes immediately.
Queue policy
| Option | Default | Behavior |
|---|---|---|
onCompletion | "delete" | Delete terminal events or retain them with "archive". |
ttlMs | null | Optional event deadline, measured from emission in milliseconds. |
archiveRetentionMs | null | Optional archive deadline, measured from settlement. |
Policy is persisted at first use. All clients must use the same policy;
conflicting definitions fail with ConfigurationConflict.
emit(queue, payload, { ttlMs }) overrides the delivery deadline. Passing null
explicitly allows indefinite waiting even when the queue has a default deadline.
An event completes when all obligations are acknowledged or waived by removal. An unfinished event past its deadline expires. Archives preserve that distinction, the payload, and each recipient's status. Indefinite events and archives are never evicted by application retention simply because time passed.
Operations
| Operation | Result and behavior |
|---|---|
subscribe(queue, name) | Durable { queue, name, generation } identity. |
unsubscribe(queue, subscription) | true when this active generation was removed; otherwise false. |
emit(queue, payload, options?) | Typed event with id, timestamps, status, and recipients. |
get(queue, id) | Typed active/archived event, or null after deletion. |
take(queue, subscription, options?) | One fenced delivery, or null; default lease is 30 seconds. |
acknowledge(queue, delivery) | acknowledged, already-acknowledged, or gone. |
renew(queue, delivery, leaseMs?) | Extend current ownership; defaults to 30 seconds. |
release(queue, delivery, delayMs?) | Retry without resolving the obligation; default delay is zero. |
processOne(queue, subscription, handler, options?) | Supervise one attempt and acknowledge on success; false if nothing is available. |
listArchived(queue, { offset?, limit? }) | Archive ids in settlement order, default limit 100, maximum 1,000. |
maintain(queue) | Bounded cleanup; returns { processed, pending }. |
runMaintenance(queue, intervalMs?) | Interruptible maintenance loop; default interval is one second. |
Timestamps are Unix milliseconds. Deliveries contain { event, subscription, leaseToken }; preserve this identity when acknowledging. Repeated acknowledgements
cannot count twice, and stale attempts cannot mutate newer leases. A gone result
means the event is no longer retained, not proof that this attempt completed it.
Managed processing
processOne options are leaseMs (30,000), renewEveryMs (one third of the
lease), and retryDelayMs (1,000). Renewal must be positive and less than the
lease. Renewal failure interrupts the handler; typed handler failure releases
for retry and propagates. Crashes and interruption recover after lease expiry.
Repeat processOne with a polling interval for continuous consumption.
Run runMaintenance alongside consumers, including on queues with no active
workers. Removal and cleanup are bounded; repeated maintenance drains the backlog.
Each pass processes at most maintenanceBatchSize records per cleanup category
(default 100, configurable from 1 to 1,000 in EventEngine.layer). Reads and
transitions also enforce deadlines and removal waivers. Acquisition can return
null after a bounded pass over stale entries; callers should continue polling.
Subscription removal
Disconnects preserve obligations. Explicit removal waives them and prevents future delivery to that generation. Re-registering a removed name creates a fresh generation; stale handles cannot affect it. A removal before the event deadline can complete it even when maintenance runs after the deadline.
Guarantees and limits
Delivery is at least once, without a processing-order guarantee. Make side effects
idempotent. Registration idempotency does not deduplicate emissions: each emit
call creates a new event. An IndeterminateWrite error includes operation and,
when applicable, event identity for inspecting uncertain outcomes before retrying.
Corrupt records, schema mismatches, and invalid replies fail explicitly.
Names are 1–256 UTF-8 bytes. Each queue permits 1,000 active subscriptions. Payloads use the existing 1 MiB storage limit. Durations are integer milliseconds up to 100 years. Archive offsets are not snapshot cursors and can shift during concurrent settlement or retention cleanup.
EventEngine.layer() supplies the Node Redis and Crypto services.
EventEngine.layerNoDeps() requires a custom RedisPool; callers must also provide
Crypto for identity-generating operations. Supported deployments are standalone
Redis and Sentinel; Cluster remains unsupported. EventEngine and EventRecord
also have root namespace exports and matching package subpaths.