@xemahq/durable-outbox
v0.3.0
Published
Storage-neutral durable outbox for Xema services: a Temporal relay per committed row, deduplicated by row id, with a scheduled sweep for crashes between commit and enqueue.
Readme
@xemahq/durable-outbox
Reliable post-transaction delivery coordination
Overview
This package coordinates durable, retrying delivery from service-owned outbox storage. It keeps storage and transport behind small public interfaces while standardizing claims, leases, bounded backoff, and delivery results.
When to use it
- Use this when an external side effect, or a Temporal start, must follow a committed state change without being lost or repeated.
Installation
pnpm add @xemahq/durable-outboxUsage
import { drainDurableOutboxBatch } from '@xemahq/durable-outbox';
await drainDurableOutboxBatch({
adapter,
deliver: async ({ payload }) => transport.send(payload),
options,
});The Temporal relay (@xemahq/durable-outbox/relay)
Write the outbox row in the domain transaction. After the commit, start its
relay: one Standalone Activity per row, with the id outbox:<table>:<rowId>
and idConflict: use-existing, so every enqueuer of a row gets the same run.
import {
createOutboxRelayActivities,
enqueueOutboxRelay,
outboxSweepSchedule,
startOutboxWorkflow,
} from '@xemahq/durable-outbox/relay';
import { upsertPlatformSchedule } from '@xemahq/temporal-runtime';
const definition = {
target: { table: 'audit-mirror', taskQueue, startToCloseTimeout: '30s' },
timing: { claimLeaseMs: 60_000, retryFloorMs: 5_000, retryCeilingMs: 900_000, sweepBatchSize: 100 },
adapter, // claim / markDelivered / markRetry: row CAS, fenced on attempt
deliver: (claim) => startOutboxWorkflow(client, {
workflowType, workflowId: `audit:${claim.id}`, taskQueue: targetQueue, orgId, args: [claim.payload],
}),
};
// The worker on `taskQueue` registers the relay and sweep activities, and
// `outboxSweepWorkflow` from `@xemahq/durable-outbox/workflows`.
const activities = createOutboxRelayActivities(client, definition);
// After the commit:
await enqueueOutboxRelay(client, definition.target, definition.timing, { id: rowId, orgId });
// Once, at boot: the sweep that covers a crash between commit and enqueue.
await upsertPlatformSchedule(client, outboxSweepSchedule(definition.target, {
every: '1m', sweepTimeout: '30s', scheduleId: `outbox-sweep:${service}:audit-mirror`,
}));deliver must be idempotent on the row id: a relay that dies after delivering
and before marking is retried. startOutboxWorkflow is idempotent by
construction (a deterministic workflow id, USE_EXISTING on a running run and
REJECT_DUPLICATE on a closed one). The claim lease must exceed one attempt's
startToCloseTimeout. No session advisory lock is involved; a pooled client is
safe.
The retry schedule, on its own
outboxRetryDelayMs is the schedule the drain uses, exported so an outbox that
is not drained by this engine still computes its next attempt the same way:
exponential from a floor, saturating at a ceiling, and deliberately not
jittered, so the stored next-attempt time is one an operator can read and reason
about.
createOutboxRetryPolicy adds a hard attempt budget to it, for an outbox whose
work is bounded rather than delivered-forever.
import { createOutboxRetryPolicy } from '@xemahq/durable-outbox';
const policy = createOutboxRetryPolicy({
drainIntervalMs: 5_000,
maxAttempts: 64,
maxBackoffMs: 60_000,
});
// Consult the budget BEFORE spending the attempt, not only after a caught
// failure: every other way an attempt is spent — a rollout, an OOM, a throwing
// failure-write — never reaches a `catch`.
if (policy.isBudgetExhausted(row.attempt)) {
await park(row);
} else {
await claim(row, policy.backoffMs(row.attempt + 1));
}A budget is safe only where a terminal row bounds ONE unit of work and never the
aggregate — that is, where new work enqueues a new row with a fresh budget.
drainDurableOutboxBatch deliberately takes no budget: an undelivered event or
audit record has no attempt count at which giving up is correct.
License
LicenseRef-Xema-BSL-1.1 © Xema — xema.dev
