effectmq
Reference

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

OptionDefaultBehavior
onCompletion"delete"Delete terminal events or retain them with "archive".
ttlMsnullOptional event deadline, measured from emission in milliseconds.
archiveRetentionMsnullOptional 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

OperationResult 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.

On this page