/** * @beignet/core/outbox * * Durable outbox primitives for transactionally recording events and jobs that * should be delivered after the owning database transaction commits. */ import { type EventPayloadDef, type EventPublishOptions, type EventTransportValue, type InferEventPayload, prepareEventPayloadForTransport, } from "../events/index.js"; import { getJobRetryDelayMs, getJobRetryMaxAttempts, type InferJobPayload, type JobDef, parseJobPayload, SINGLE_ATTEMPT_DISPATCH, type SingleAttemptJobDispatch, shouldRetryJob, } from "../jobs/index.js"; import type { JobDispatcherPort } from "../ports/events.js"; import type { DomainEventRecorderPort } from "../ports/unit-of-work.js"; import { type BaseProviderInstrumentationEvent, createProviderInstrumentation, type ProviderInstrumentationTarget, } from "../providers/index.js"; import { captureTraceCarrier, parseTraceCarrier, resolveTracingPort, runWithTracing, type TraceCarrier, type TracingPort, } from "../tracing/index.js"; /** * Value or promise of that value. */ export type MaybePromise = T | Promise; /** * Default lease duration for claimed outbox messages. */ export const DEFAULT_OUTBOX_LEASE_MS = 30_000; /** * Default maximum messages handled by one bounded drain pass. */ export const DEFAULT_OUTBOX_BATCH_SIZE = 100; /** * Default number of outbox messages delivered concurrently. */ export const DEFAULT_OUTBOX_CONCURRENCY = 1; /** * Default maximum time Beignet renews a claim for one delivery. */ export const DEFAULT_OUTBOX_MAX_ACTIVE_MS = 300_000; /** * Default maximum delivery attempts before a message is dead-lettered. */ export const DEFAULT_OUTBOX_MAX_ATTEMPTS = 3; /** * Message kinds supported by the Beignet outbox. */ export type OutboxMessageKind = "event" | "job"; /** * Delivery status for an outbox message. */ export type OutboxMessageStatus = | "pending" | "claimed" | "delivered" | "deadLettered"; /** * JSON-serializable value accepted by the outbox. */ export type OutboxJsonValue = | null | string | number | boolean | readonly OutboxJsonValue[] | { readonly [key: string]: OutboxJsonValue }; /** * JSON object accepted by outbox payload helpers. */ export type OutboxJsonObject = { readonly [key: string]: OutboxJsonValue }; /** * Serialized delivery error stored on failed messages. */ export interface OutboxErrorInfo { /** * Error name when available. */ name?: string; /** * Error message. */ message: string; /** * Error stack when available. */ stack?: string; } /** * Input for enqueueing a raw outbox message. */ export interface OutboxEnqueueInput { /** * Optional caller-provided message ID. */ id?: string; /** * Message kind. */ kind: OutboxMessageKind; /** * Event or job name. */ name: string; /** * JSON-serializable payload. */ payload: OutboxJsonValue; /** Versioned trace context captured when the message was recorded. */ trace?: TraceCarrier; /** * Earliest time the message may be claimed. */ availableAt?: Date; /** * Maximum delivery attempts before dead-lettering. */ maxAttempts?: number; } /** * Durable outbox message record. */ export interface OutboxMessage { /** * Stable message ID. */ id: string; /** * Message kind. */ kind: OutboxMessageKind; /** * Event or job name. */ name: string; /** * JSON-serializable payload. */ payload: OutboxJsonValue; /** Versioned trace context captured when the message was recorded. */ trace?: TraceCarrier; /** * Current delivery status. */ status: OutboxMessageStatus; /** * Number of claim attempts. */ attempts: number; /** * Maximum delivery attempts before dead-lettering. */ maxAttempts: number; /** * Earliest time the message may be claimed. */ availableAt: Date; /** * Last claim timestamp. */ claimedAt: Date | null; /** * Lease expiration timestamp for claimed messages. */ lockedUntil: Date | null; /** * Token required to mark a claimed message delivered or failed. */ claimToken: string | null; /** * Delivery timestamp. */ deliveredAt: Date | null; /** * Last delivery error. */ lastError: OutboxErrorInfo | null; /** * Creation timestamp. */ createdAt: Date; /** * Last update timestamp. */ updatedAt: Date; } /** * Message returned from a successful outbox claim. */ export interface ClaimedOutboxMessage extends Omit< OutboxMessage, "claimToken" | "claimedAt" | "lockedUntil" | "status" > { status: "claimed"; claimToken: string; claimedAt: Date; lockedUntil: Date; } /** * Options for leasing a batch of pending outbox messages for delivery. */ export interface OutboxClaimBatchOptions { /** * Maximum eligible messages to claim or reconcile in one batch. */ limit: number; /** * Claim timestamp. */ now?: Date; /** * Lease duration in milliseconds. */ leaseMs?: number; } /** * Result of atomically selecting one bounded set of eligible messages. */ export interface OutboxClaimBatchResult { /** Messages claimed for delivery by the current worker. */ claimed: readonly ClaimedOutboxMessage[]; /** * Eligible messages moved directly to dead letter because their claim * attempt budget was already exhausted. */ deadLettered: readonly OutboxMessage[]; } /** * Input for extending one active outbox claim. */ export interface OutboxRenewClaimInput { /** Claimed message ID. */ id: string; /** Claim token returned by `claimBatch(...)`. */ claimToken: string; /** Renewal timestamp. */ now?: Date; /** New lease duration measured from `now`. */ leaseMs?: number; } /** * Confirmed result of extending one active outbox claim. */ export interface OutboxRenewClaimResult { /** New confirmed lease expiration timestamp. */ lockedUntil: Date; } /** * Input for marking a claimed message delivered. */ export interface OutboxMarkDeliveredInput { /** * Claimed message ID. */ id: string; /** * Claim token returned by `claimBatch(...)`. */ claimToken: string; /** * Delivery timestamp. */ now?: Date; } /** * Input for marking a claimed message failed. */ export interface OutboxMarkFailedInput { /** * Claimed message ID. */ id: string; /** * Claim token returned by `claimBatch(...)`. */ claimToken: string; /** * Delivery error. */ error?: unknown; /** * Next time the message may be claimed. */ retryAt?: Date; /** * Whether this failure should dead-letter the message. */ deadLetter?: boolean; /** * Failure timestamp. */ now?: Date; } /** * Default maximum messages returned from outbox admin list calls. */ export const DEFAULT_OUTBOX_ADMIN_LIST_LIMIT = 50; /** * Shared filters for outbox admin read operations. */ export interface OutboxMessageQuery { /** * Status or statuses to include. */ status?: OutboxMessageStatus | readonly OutboxMessageStatus[]; /** * Message kind to include. */ kind?: OutboxMessageKind; /** * Event or job name to include. */ name?: string; /** * Include messages last updated before this timestamp. */ updatedBefore?: Date; /** * Include delivered messages delivered before this timestamp. */ deliveredBefore?: Date; } /** * Options for listing outbox messages through the admin port. */ export interface OutboxListMessagesOptions extends OutboxMessageQuery { /** * Maximum messages to return. Defaults to * `DEFAULT_OUTBOX_ADMIN_LIST_LIMIT`. */ limit?: number; } /** * Options for counting outbox messages through the admin port. */ export interface OutboxCountMessagesOptions extends OutboxMessageQuery {} /** * Input for returning a dead-lettered message to the pending queue. */ export interface OutboxRequeueMessageInput { /** * Dead-lettered message ID. */ id: string; /** * Earliest time the message may be claimed again. Defaults to `now`. */ availableAt?: Date; /** * Reset attempts to zero before requeueing. Defaults to preserving the * attempt count so operators can decide whether to grant a fresh retry * budget. */ resetAttempts?: boolean; /** * Requeue timestamp. */ now?: Date; } /** * Input for purging dead-lettered messages. */ export interface OutboxPurgeDeadLetteredInput { /** * Only purge messages last updated before this timestamp. Omit only when the * caller intentionally wants to purge all dead-lettered messages. */ before?: Date; /** * Maximum messages to purge. */ limit?: number; } /** * Input for pruning delivered messages. */ export interface OutboxPruneDeliveredInput { /** * Prune delivered messages delivered before this timestamp. */ before: Date; /** * Maximum messages to prune. */ limit?: number; } /** * Result for outbox admin delete operations. */ export interface OutboxDeleteResult { /** * Number of rows deleted. */ deleted: number; } /** * App-facing outbox storage port. * * Durable adapters should claim messages atomically and require `claimToken` * for delivery/failure updates. */ export interface OutboxPort { /** * Enqueue a new pending message. */ enqueue(input: OutboxEnqueueInput): Promise; /** * Atomically claim eligible messages for one worker. */ claimBatch(options: OutboxClaimBatchOptions): Promise; /** * Extend an unexpired claim owned by the supplied claim token. */ renewClaim(input: OutboxRenewClaimInput): Promise; /** * Mark a claimed message delivered. */ markDelivered(input: OutboxMarkDeliveredInput): Promise; /** * Mark a claimed message failed, retryable, or dead-lettered. */ markFailed(input: OutboxMarkFailedInput): Promise; } /** * Operational outbox admin port. * * Keep this separate from `OutboxPort` so request and drain contexts can expose * only the hot-path delivery operations. Wire this port into maintenance * contexts for CLI commands, runbooks, and devtools. */ export interface OutboxAdminPort { /** * List messages ordered by newest update first. */ listMessages( options?: OutboxListMessagesOptions, ): Promise; /** * Count messages matching admin filters. */ countMessages(options?: OutboxCountMessagesOptions): Promise; /** * Fetch one message by ID. */ getMessage(id: string): Promise; /** * Return a dead-lettered message to pending state. */ requeueMessage(input: OutboxRequeueMessageInput): Promise; /** * Delete dead-lettered messages. */ purgeDeadLettered( input?: OutboxPurgeDeadLetteredInput, ): Promise; /** * Delete delivered messages older than a retention cutoff. */ pruneDelivered(input: OutboxPruneDeliveredInput): Promise; } /** * In-memory outbox for tests and local examples. */ export interface MemoryOutboxPort extends OutboxPort, OutboxAdminPort { /** * Current message snapshots. */ readonly messages: readonly OutboxMessage[]; /** * Remove all messages. */ clear(): void; } /** * Options for `createOutboxMessage(...)`. */ export interface CreateOutboxMessageOptions { /** * Generated message ID override. Wins over `input.id`. */ id?: string; /** * Fallback ID factory used when neither `id` nor `input.id` is provided. */ createId?: () => string; /** * Timestamp used for created/updated/available dates. */ now?: Date; } /** * Options for typed event/job enqueue helpers. */ export interface EnqueueTypedOutboxOptions { /** * Optional caller-provided message ID. */ id?: string; /** * Earliest time the message may be claimed. */ availableAt?: Date; /** * Maximum delivery attempts before dead-lettering. */ maxAttempts?: number; /** Explicit trace context to persist with this message. */ trace?: TraceCarrier; /** Tracing port used to capture the active context at enqueue time. */ tracing?: TracingPort; } /** * Registry of definitions that `drainOutbox(...)` can deliver. */ export interface OutboxRegistry { /** * Event definitions keyed by event name. */ readonly events: ReadonlyMap; /** * Job definitions keyed by job name. */ readonly jobs: ReadonlyMap; } /** * Input for defining an outbox registry. */ export interface DefineOutboxRegistryInput { /** * Events that may be delivered from the outbox. */ events?: readonly EventPayloadDef[]; /** * Jobs that may be delivered from the outbox. */ jobs?: readonly JobDef[]; } /** * Correlation fields attached to outbox instrumentation events. */ export type OutboxInstrumentationContext = Pick< BaseProviderInstrumentationEvent, "requestId" | "traceId" | "spanId" | "parentSpanId" | "traceparent" >; /** Wait primitive used by outbox heartbeats and bounded settlement retries. */ export type OutboxDrainWait = ( delayMs: number, signal: AbortSignal, ) => Promise; /** Structured failure from a claim heartbeat or active-delivery boundary. */ export interface OutboxLeaseFailure { /** Underlying renewal, ownership, or duration error. */ error: unknown; /** Message whose claim could not be kept active. */ message: ClaimedOutboxMessage; /** Lease phase that surfaced the failure. */ operation: "renewClaim" | "maxActiveDuration"; /** Current lease state after handling the failure. */ state: "recovered" | "degraded" | "lost"; /** Whether the worker no longer has confirmed ownership. */ confirmedLost: boolean; } /** Structured failure from a token-guarded outbox settlement. */ export interface OutboxSettlementFailure { /** Storage failure returned by the settlement operation. */ error: unknown; /** Message whose final storage state is unknown. */ message: ClaimedOutboxMessage; /** Settlement operation that failed. */ operation: "markDelivered" | "markFailed"; /** Whether the external delivery completed successfully. */ deliverySucceeded: boolean; /** Original delivery error when `deliverySucceeded` is false. */ deliveryError?: unknown; } /** * Options for draining one outbox batch. */ export interface DrainOutboxOptions { /** * Outbox storage port. */ outbox: OutboxPort; /** * Registry used to resolve message names to event/job definitions. */ registry: OutboxRegistry; /** * Event bus used for event messages. Required when the registry contains * events. */ eventBus?: { publish( event: E, payload: InferEventPayload, options?: EventPublishOptions, ): MaybePromise; }; /** * Job dispatcher used for job messages. Required when the registry contains * jobs. */ jobs?: JobDispatcherPort; /** * Maximum eligible messages to handle in one drain pass. */ batchSize?: number; /** * Maximum messages delivered concurrently. Defaults to serial delivery. * Values greater than one do not preserve delivery order. */ concurrency?: number; /** * Clock used independently for claiming, renewal, settlement, and retry * scheduling. Defaults to the system clock. */ now?: () => Date; /** * Claim lease duration in milliseconds. */ leaseMs?: number; /** * Interval between serialized claim renewals. Defaults to one third of the * lease duration and must remain shorter than the lease. */ heartbeatMs?: number; /** * Maximum time Beignet renews a claim for one delivery. When exceeded, the * drain stops renewing and leaves final recovery to the last lease expiry. */ maxActiveMs?: number; /** * Abort-aware wait implementation. Inject a deterministic implementation in * tests; production callers normally use the default timer. */ wait?: OutboxDrainWait; /** * Retry delay in milliseconds or function for per-message delay. */ retryDelayMs?: | number | ((args: { message: ClaimedOutboxMessage; error: unknown; now: Date; }) => number); /** * Optional instrumentation target for delivery, retry, and dead-letter * visibility. */ instrumentation?: ProviderInstrumentationTarget; /** * Optional correlation fields attached to outbox instrumentation events. */ instrumentationContext?: OutboxInstrumentationContext; /** * Observer called when delivery fails. Observer failures are ignored so the * original delivery failure still controls retry/dead-letter behavior. */ onError?: ( error: unknown, message: ClaimedOutboxMessage, ) => MaybePromise; /** * Observer called after a failed delivery is successfully moved to the dead * letter state. Observer failures are ignored. */ onDeadLetter?: (error: unknown, message: OutboxMessage) => MaybePromise; /** * Observer called when claim renewal degrades or ownership is lost. * Observer failures are ignored. */ onLeaseError?: (failure: OutboxLeaseFailure) => MaybePromise; /** * Observer called when a delivery outcome cannot be settled durably. * Observer failures are ignored. */ onSettlementError?: (failure: OutboxSettlementFailure) => MaybePromise; } /** * Summary returned from one `drainOutbox(...)` pass. */ export interface DrainOutboxResult { /** * Messages claimed in this batch. */ claimed: number; /** * Messages delivered successfully. */ delivered: number; /** * Messages scheduled for retry. */ retried: number; /** * Messages moved to dead letter state. */ deadLettered: number; /** * Dead-lettered messages whose claim attempt budget was already exhausted. * This is a subset of `deadLettered`. */ abandonedDeadLettered: number; /** Messages delivered or failed whose final storage state is unknown. */ settlementFailed: number; /** Messages whose active claim could no longer be confirmed. */ leaseLost: number; } /** * Error thrown when an outbox payload is not JSON serializable. */ export class OutboxSerializationError extends Error { constructor(message: string) { super(message); this.name = "OutboxSerializationError"; } } /** * Error thrown when an outbox message cannot be resolved through the registry. */ export class OutboxRegistryError extends Error { constructor(message: string) { super(message); this.name = "OutboxRegistryError"; } } /** * Error thrown when a claimed message cannot be updated with the supplied token. */ export class OutboxClaimError extends Error { /** * Message ID involved in the claim error. */ readonly id: string; constructor(args: { id: string; message: string }) { super(args.message); this.name = "OutboxClaimError"; this.id = args.id; } } /** * Error persisted when an eligible message has exhausted its claim attempts * without reaching a terminal settlement. */ export class OutboxAbandonedClaimError extends Error { /** Message ID whose attempt budget was exhausted. */ readonly id: string; /** Number of claims already made. */ readonly attempts: number; /** Maximum permitted claim attempts. */ readonly maxAttempts: number; constructor(args: { id: string; attempts: number; maxAttempts: number }) { super( `Outbox message "${args.id}" exhausted ${args.maxAttempts} claim attempts without a terminal settlement.`, ); this.name = "OutboxAbandonedClaimError"; this.id = args.id; this.attempts = args.attempts; this.maxAttempts = args.maxAttempts; } } /** Error surfaced when an outbox worker no longer has a confirmed claim. */ export class OutboxLeaseLostError extends Error { /** Message ID whose ownership was lost. */ readonly id: string; constructor(args: { id: string; message?: string; cause?: unknown }) { super( args.message ?? `Outbox claim for message "${args.id}" is no longer active.`, args.cause === undefined ? undefined : { cause: args.cause }, ); this.name = "OutboxLeaseLostError"; this.id = args.id; } } /** Error surfaced when delivery exceeds the configured active-claim window. */ export class OutboxMaxActiveDurationError extends Error { /** Message ID whose active window elapsed. */ readonly id: string; /** Configured maximum active duration. */ readonly maxActiveMs: number; constructor(args: { id: string; maxActiveMs: number }) { super( `Outbox delivery for message "${args.id}" exceeded the ${args.maxActiveMs}ms active-claim limit.`, ); this.name = "OutboxMaxActiveDurationError"; this.id = args.id; this.maxActiveMs = args.maxActiveMs; } } /** * Error thrown when an outbox admin operation cannot be completed safely. */ export class OutboxAdminError extends Error { /** * Message ID involved in the admin error, when applicable. */ readonly id?: string; constructor(args: { message: string; id?: string }) { super(args.message); this.name = "OutboxAdminError"; this.id = args.id; } } function assertNonEmptyString(name: string, value: string): void { if (typeof value !== "string" || value.trim().length === 0) { throw new Error(`${name} must be a non-empty string`); } } function assertPositiveInteger(name: string, value: number): void { if (!Number.isInteger(value) || value <= 0) { throw new Error(`${name} must be a positive integer`); } } function assertValidDate(name: string, value: Date): void { if (!(value instanceof Date) || Number.isNaN(value.getTime())) { throw new Error(`${name} must be a valid Date`); } } function cloneDate(value: Date): Date { return new Date(value.getTime()); } function createId(): string { if (!globalThis.crypto?.randomUUID) { throw new Error("crypto.randomUUID is required to create outbox IDs."); } return globalThis.crypto.randomUUID(); } function assertJsonValue( value: unknown, path: readonly string[] = [], seen: WeakSet = new WeakSet(), ): OutboxJsonValue { const label = path.length > 0 ? path.join(".") : "payload"; if (value === null) return null; if ( typeof value === "string" || typeof value === "boolean" || typeof value === "number" ) { if (typeof value === "number" && !Number.isFinite(value)) { throw new OutboxSerializationError( `Outbox ${label} must be a finite number.`, ); } return value; } if (Array.isArray(value)) { if (seen.has(value)) { throw new OutboxSerializationError( `Outbox ${label} must be JSON serializable. Circular references are not supported.`, ); } seen.add(value); try { return value.map((item, index) => assertJsonValue(item, [...path, String(index)], seen), ); } finally { seen.delete(value); } } if (typeof value === "object") { if (value instanceof Date) { throw new OutboxSerializationError( `Outbox ${label} must be JSON serializable. Convert Date values to strings before enqueueing.`, ); } if (seen.has(value)) { throw new OutboxSerializationError( `Outbox ${label} must be JSON serializable. Circular references are not supported.`, ); } seen.add(value); try { const record = value as Record; const output: Record = {}; for (const key of Object.keys(record)) { const child = record[key]; if (child === undefined) { throw new OutboxSerializationError( `Outbox ${[...path, key].join(".")} cannot be undefined.`, ); } output[key] = assertJsonValue(child, [...path, key], seen); } return output; } finally { seen.delete(value); } } throw new OutboxSerializationError( `Outbox ${label} must be JSON serializable. Received ${typeof value}.`, ); } /** * Convert an unknown value to an outbox-safe JSON value. * * Dates, undefined values, functions, non-finite numbers, symbols, and circular * references are rejected so durable adapters can store the payload safely. */ export function toOutboxJsonValue(value: unknown): OutboxJsonValue { const jsonValue = assertJsonValue(value); return JSON.parse(JSON.stringify(jsonValue)) as OutboxJsonValue; } /** * Serialize an unknown delivery error into outbox error metadata. */ export function serializeOutboxError(error: unknown): OutboxErrorInfo { if (error instanceof Error) { return { name: error.name, message: error.message, stack: error.stack, }; } if (typeof error === "string") { return { message: error }; } return { message: "Unknown outbox delivery error" }; } function copyMessage(message: OutboxMessage): OutboxMessage { return { ...message, availableAt: cloneDate(message.availableAt), claimedAt: message.claimedAt ? cloneDate(message.claimedAt) : null, lockedUntil: message.lockedUntil ? cloneDate(message.lockedUntil) : null, deliveredAt: message.deliveredAt ? cloneDate(message.deliveredAt) : null, createdAt: cloneDate(message.createdAt), updatedAt: cloneDate(message.updatedAt), lastError: message.lastError ? { ...message.lastError } : null, }; } function normalizeStatusFilter( status: OutboxMessageQuery["status"], ): ReadonlySet | undefined { if (status === undefined) return undefined; return new Set(Array.isArray(status) ? status : [status]); } function messageMatchesQuery( message: OutboxMessage, query: OutboxMessageQuery = {}, ): boolean { const statuses = normalizeStatusFilter(query.status); if (statuses && !statuses.has(message.status)) return false; if (query.kind !== undefined && message.kind !== query.kind) return false; if (query.name !== undefined && message.name !== query.name) return false; if ( query.updatedBefore !== undefined && message.updatedAt.getTime() >= query.updatedBefore.getTime() ) { return false; } if ( query.deliveredBefore !== undefined && (message.deliveredAt === null || message.deliveredAt.getTime() >= query.deliveredBefore.getTime()) ) { return false; } return true; } function compareMessagesForAdminList( left: OutboxMessage, right: OutboxMessage, ): number { const updated = right.updatedAt.getTime() - left.updatedAt.getTime(); if (updated !== 0) return updated; const created = right.createdAt.getTime() - left.createdAt.getTime(); if (created !== 0) return created; return right.id.localeCompare(left.id); } function compareMessagesForAdminCleanup( left: OutboxMessage, right: OutboxMessage, ): number { const updated = left.updatedAt.getTime() - right.updatedAt.getTime(); if (updated !== 0) return updated; const created = left.createdAt.getTime() - right.createdAt.getTime(); if (created !== 0) return created; return left.id.localeCompare(right.id); } function compareDeliveredMessagesForPrune( left: OutboxMessage, right: OutboxMessage, ): number { const delivered = (left.deliveredAt?.getTime() ?? 0) - (right.deliveredAt?.getTime() ?? 0); if (delivered !== 0) return delivered; return compareMessagesForAdminCleanup(left, right); } function resolveAdminListLimit(limit: number | undefined): number { const resolved = limit ?? DEFAULT_OUTBOX_ADMIN_LIST_LIMIT; assertPositiveInteger("limit", resolved); return resolved; } function toClaimedMessage(message: OutboxMessage): ClaimedOutboxMessage { if ( message.status !== "claimed" || !message.claimToken || !message.claimedAt || !message.lockedUntil ) { throw new OutboxClaimError({ id: message.id, message: `Outbox message "${message.id}" is not claimed.`, }); } return { ...copyMessage(message), status: "claimed", claimToken: message.claimToken, claimedAt: cloneDate(message.claimedAt), lockedUntil: cloneDate(message.lockedUntil), }; } /** * Create a validated pending outbox message. */ export function createOutboxMessage( input: OutboxEnqueueInput, options: CreateOutboxMessageOptions = {}, ): OutboxMessage { assertNonEmptyString("kind", input.kind); assertNonEmptyString("name", input.name); if (input.id !== undefined) assertNonEmptyString("id", input.id); if (input.maxAttempts !== undefined) { assertPositiveInteger("maxAttempts", input.maxAttempts); } const now = options.now ?? new Date(); const trace = parseTraceCarrier(input.trace); return { id: options.id ?? input.id ?? options.createId?.() ?? createId(), kind: input.kind, name: input.name, payload: toOutboxJsonValue(input.payload), ...(trace ? { trace } : {}), status: "pending", attempts: 0, maxAttempts: input.maxAttempts ?? DEFAULT_OUTBOX_MAX_ATTEMPTS, availableAt: input.availableAt ? cloneDate(input.availableAt) : cloneDate(now), claimedAt: null, lockedUntil: null, claimToken: null, deliveredAt: null, lastError: null, createdAt: cloneDate(now), updatedAt: cloneDate(now), }; } function isEligible(message: OutboxMessage, now: Date): boolean { if (message.status === "pending") { return message.availableAt.getTime() <= now.getTime(); } return ( message.status === "claimed" && message.lockedUntil !== null && message.lockedUntil.getTime() <= now.getTime() ); } /** * Options for `createMemoryOutbox(...)`. */ export interface MemoryOutboxOptions { /** * Message and claim-token ID factory. Defaults to `crypto.randomUUID()`. */ id?: () => string; /** * Clock used for enqueue, claim, and completion timestamps when a call does * not supply its own `now`. Defaults to the system clock. */ now?: () => Date; } /** * Create an in-memory outbox for tests and local examples. * * The memory outbox is process-local and not durable. */ export function createMemoryOutbox( storeOptions: MemoryOutboxOptions = {}, ): MemoryOutboxPort { const createStoreId = storeOptions.id ?? createId; const storeNow = storeOptions.now ?? (() => new Date()); const messages = new Map(); function getClaimedOrThrow(id: string, claimToken: string): OutboxMessage { const message = messages.get(id); if (!message) { throw new OutboxClaimError({ id, message: `Outbox message "${id}" does not exist.`, }); } if (message.status !== "claimed" || message.claimToken !== claimToken) { throw new OutboxClaimError({ id, message: `Outbox message "${id}" is not claimed by this worker.`, }); } return message; } return { get messages() { return [...messages.values()].map(copyMessage); }, async listMessages(options = {}) { const limit = resolveAdminListLimit(options.limit); if (options.updatedBefore) { assertValidDate("updatedBefore", options.updatedBefore); } if (options.deliveredBefore) { assertValidDate("deliveredBefore", options.deliveredBefore); } return [...messages.values()] .filter((message) => messageMatchesQuery(message, options)) .sort(compareMessagesForAdminList) .slice(0, limit) .map(copyMessage); }, async countMessages(options = {}) { if (options.updatedBefore) { assertValidDate("updatedBefore", options.updatedBefore); } if (options.deliveredBefore) { assertValidDate("deliveredBefore", options.deliveredBefore); } return [...messages.values()].filter((message) => messageMatchesQuery(message, options), ).length; }, async getMessage(id) { assertNonEmptyString("id", id); const message = messages.get(id); return message ? copyMessage(message) : null; }, async requeueMessage(input) { assertNonEmptyString("id", input.id); const message = messages.get(input.id); if (!message) { throw new OutboxAdminError({ id: input.id, message: `Outbox message "${input.id}" does not exist.`, }); } if (message.status !== "deadLettered") { throw new OutboxAdminError({ id: input.id, message: `Outbox message "${input.id}" is not dead-lettered.`, }); } const now = input.now ?? storeNow(); const availableAt = input.availableAt ?? now; assertValidDate("now", now); assertValidDate("availableAt", availableAt); message.status = "pending"; message.availableAt = cloneDate(availableAt); message.claimToken = null; message.claimedAt = null; message.lockedUntil = null; if (input.resetAttempts) message.attempts = 0; message.updatedAt = cloneDate(now); return copyMessage(message); }, async purgeDeadLettered(input = {}) { const limit = input.limit === undefined ? undefined : resolveAdminListLimit(input.limit); if (input.before) assertValidDate("before", input.before); const candidates = [...messages.values()] .filter( (message) => message.status === "deadLettered" && (input.before === undefined || message.updatedAt.getTime() < input.before.getTime()), ) .sort(compareMessagesForAdminCleanup) .slice(0, limit); for (const message of candidates) { messages.delete(message.id); } return { deleted: candidates.length }; }, async pruneDelivered(input) { assertValidDate("before", input.before); const limit = input.limit === undefined ? undefined : resolveAdminListLimit(input.limit); const candidates = [...messages.values()] .filter( (message) => message.status === "delivered" && message.deliveredAt !== null && message.deliveredAt.getTime() < input.before.getTime(), ) .sort(compareDeliveredMessagesForPrune) .slice(0, limit); for (const message of candidates) { messages.delete(message.id); } return { deleted: candidates.length }; }, async enqueue(input) { const message = createOutboxMessage(input, { createId: createStoreId, now: storeNow(), }); if (messages.has(message.id)) { throw new Error(`Outbox message "${message.id}" already exists.`); } messages.set(message.id, message); return copyMessage(message); }, async claimBatch(options) { assertPositiveInteger("limit", options.limit); const now = options.now ?? storeNow(); const leaseMs = options.leaseMs ?? DEFAULT_OUTBOX_LEASE_MS; assertValidDate("now", now); assertPositiveInteger("leaseMs", leaseMs); const lockedUntil = new Date(now.getTime() + leaseMs); const claimed: ClaimedOutboxMessage[] = []; const deadLettered: OutboxMessage[] = []; const eligible = [...messages.values()] .filter((message) => isEligible(message, now)) .sort((a, b) => { const available = a.availableAt.getTime() - b.availableAt.getTime(); if (available !== 0) return available; return a.createdAt.getTime() - b.createdAt.getTime(); }) .slice(0, options.limit); for (const message of eligible) { if (message.attempts >= message.maxAttempts) { message.status = "deadLettered"; message.lastError = serializeOutboxError( new OutboxAbandonedClaimError({ id: message.id, attempts: message.attempts, maxAttempts: message.maxAttempts, }), ); message.claimToken = null; message.claimedAt = null; message.lockedUntil = null; message.updatedAt = cloneDate(now); deadLettered.push(copyMessage(message)); continue; } message.status = "claimed"; message.attempts += 1; message.claimToken = createStoreId(); message.claimedAt = cloneDate(now); message.lockedUntil = cloneDate(lockedUntil); message.updatedAt = cloneDate(now); claimed.push(toClaimedMessage(message)); } return { claimed, deadLettered }; }, async renewClaim(input) { assertNonEmptyString("id", input.id); assertNonEmptyString("claimToken", input.claimToken); const message = getClaimedOrThrow(input.id, input.claimToken); const now = input.now ?? storeNow(); const leaseMs = input.leaseMs ?? DEFAULT_OUTBOX_LEASE_MS; assertValidDate("now", now); assertPositiveInteger("leaseMs", leaseMs); if ( message.lockedUntil === null || message.lockedUntil.getTime() <= now.getTime() ) { throw new OutboxClaimError({ id: input.id, message: `Outbox message "${input.id}" no longer has an active claim to renew.`, }); } const lockedUntil = new Date(now.getTime() + leaseMs); message.lockedUntil = cloneDate(lockedUntil); message.updatedAt = cloneDate(now); return { lockedUntil: cloneDate(lockedUntil) }; }, async markDelivered(input) { assertNonEmptyString("id", input.id); assertNonEmptyString("claimToken", input.claimToken); const message = getClaimedOrThrow(input.id, input.claimToken); const now = input.now ?? storeNow(); assertValidDate("now", now); if ( message.lockedUntil === null || message.lockedUntil.getTime() <= now.getTime() ) { throw new OutboxClaimError({ id: input.id, message: `Outbox message "${input.id}" no longer has an active claim to settle.`, }); } message.status = "delivered"; message.deliveredAt = cloneDate(now); message.claimToken = null; message.claimedAt = null; message.lockedUntil = null; message.updatedAt = cloneDate(now); }, async markFailed(input) { assertNonEmptyString("id", input.id); assertNonEmptyString("claimToken", input.claimToken); const message = getClaimedOrThrow(input.id, input.claimToken); const now = input.now ?? storeNow(); assertValidDate("now", now); if (input.retryAt) assertValidDate("retryAt", input.retryAt); if ( message.lockedUntil === null || message.lockedUntil.getTime() <= now.getTime() ) { throw new OutboxClaimError({ id: input.id, message: `Outbox message "${input.id}" no longer has an active claim to settle.`, }); } message.status = input.deadLetter ? "deadLettered" : "pending"; message.lastError = serializeOutboxError(input.error); message.availableAt = input.retryAt ? cloneDate(input.retryAt) : cloneDate(now); message.claimToken = null; message.claimedAt = null; message.lockedUntil = null; message.updatedAt = cloneDate(now); }, clear() { messages.clear(); }, }; } function mapDefinitions( kind: string, defs: readonly T[], ): ReadonlyMap { const map = new Map(); for (const def of defs) { if (map.has(def.name)) { throw new OutboxRegistryError( `Duplicate ${kind} definition "${def.name}" in outbox registry.`, ); } map.set(def.name, def); } return map; } /** * Define the events and jobs that an outbox drain worker can deliver. * * Duplicate names throw because message delivery resolves by name. */ export function defineOutboxRegistry( input: DefineOutboxRegistryInput, ): OutboxRegistry { return { events: mapDefinitions("event", input.events ?? []), jobs: mapDefinitions("job", input.jobs ?? []), }; } /** * Validate an event payload and enqueue it as an outbox message. */ export async function enqueueEvent( outbox: OutboxPort, event: E, payload: InferEventPayload, options: EnqueueTypedOutboxOptions = {}, ): Promise { const prepared = await prepareEventPayloadForTransport(event, payload); return await enqueueTransportEvent( outbox, event, prepared.transportValue, options, ); } async function enqueueTransportEvent( outbox: OutboxPort, event: E, payload: EventTransportValue, options: EnqueueTypedOutboxOptions, ): Promise { const trace = parseTraceCarrier(options.trace) ?? captureTraceCarrier(options.tracing); return outbox.enqueue({ id: options.id, kind: "event", name: event.name, payload, trace, availableAt: options.availableAt, maxAttempts: options.maxAttempts, }); } /** * Validate a job payload and enqueue it as an outbox message. */ export async function enqueueJob( outbox: OutboxPort, job: J, payload: InferJobPayload, options: EnqueueTypedOutboxOptions = {}, ): Promise { await parseJobPayload(job, payload); const trace = parseTraceCarrier(options.trace) ?? captureTraceCarrier(options.tracing); return outbox.enqueue({ id: options.id, kind: "job", name: job.name, payload: toOutboxJsonValue(payload), trace, availableAt: options.availableAt, maxAttempts: options.maxAttempts ?? getJobRetryMaxAttempts(job.retry), }); } /** * Create a domain event recorder that writes events to the outbox. */ export function createOutboxEventRecorder( outbox: OutboxPort, options: EnqueueTypedOutboxOptions = {}, ): DomainEventRecorderPort { return { async record(event, payload, publishOptions) { const prepared = await prepareEventPayloadForTransport( event, payload, publishOptions, ); await enqueueTransportEvent(outbox, event, prepared.transportValue, { ...options, trace: prepared.publishOptions.trace ?? options.trace, }); }, }; } /** * Create a job dispatcher that writes jobs to the outbox. */ export function createOutboxJobDispatcher( outbox: OutboxPort, options: EnqueueTypedOutboxOptions = {}, ): JobDispatcherPort { return { async dispatch(job, payload, dispatchOptions) { await enqueueJob(outbox, job, payload, { ...options, trace: dispatchOptions?.trace ?? options.trace, }); }, }; } function resolveRetryDelayMs( options: DrainOutboxOptions, message: ClaimedOutboxMessage, error: unknown, now: Date, ): number { if (typeof options.retryDelayMs === "function") { const delay = options.retryDelayMs({ message, error, now }); assertPositiveInteger("retryDelayMs", delay); return delay; } if (options.retryDelayMs !== undefined) { assertPositiveInteger("retryDelayMs", options.retryDelayMs); return options.retryDelayMs; } if (message.kind === "job") { const job = options.registry.jobs.get(message.name); if (job?.retry) { return getJobRetryDelayMs(job.retry, { attempt: message.attempts, error, jobName: message.name, }); } } return Math.min(60_000, 1000 * 2 ** Math.max(0, message.attempts - 1)); } function shouldRetryOutboxMessage( options: DrainOutboxOptions, message: ClaimedOutboxMessage, error: unknown, ): boolean { if (message.kind !== "job") { return message.attempts < message.maxAttempts; } const job = options.registry.jobs.get(message.name); return shouldRetryJob(job?.retry, { attempt: message.attempts, error, jobName: message.name, maxAttempts: message.maxAttempts, }); } function outboxInstrumentationDetails( message: OutboxMessage, details?: Record, ): Record { return { attempt: message.attempts, maxAttempts: message.maxAttempts, messageId: message.id, messageKind: message.kind, messageName: message.name, ...details, }; } async function deliverOutboxMessage( options: DrainOutboxOptions, message: ClaimedOutboxMessage, trace?: TraceCarrier, ): Promise { if (message.kind === "event") { if (!options.eventBus) { throw new OutboxRegistryError( `Cannot deliver event "${message.name}" without an event bus.`, ); } const event = options.registry.events.get(message.name); if (!event) { throw new OutboxRegistryError( `Outbox registry does not include event "${message.name}".`, ); } const prepared = await prepareEventPayloadForTransport( event, message.payload, trace ? { trace } : undefined, ); await options.eventBus.publish( event, prepared.payload, prepared.publishOptions, ); return; } if (!options.jobs) { throw new OutboxRegistryError( `Cannot deliver job "${message.name}" without a job dispatcher.`, ); } const job = options.registry.jobs.get(message.name); if (!job) { throw new OutboxRegistryError( `Outbox registry does not include job "${message.name}".`, ); } await parseJobPayload(job, message.payload); // The drain owns execution retries: failed deliveries are rescheduled with // the job's own policy via markFailed/retryAt. When the dispatcher exposes // a single-attempt dispatch (the inline dispatcher does), use it so the // retry policy runs in exactly one layer. Durable providers do not expose // it — for them dispatch is an enqueue and the queue owns execution. const singleAttempt = ( options.jobs as { [SINGLE_ATTEMPT_DISPATCH]?: SingleAttemptJobDispatch; } )[SINGLE_ATTEMPT_DISPATCH]; if (singleAttempt) { await singleAttempt(job, message.payload as never, { attempt: message.attempts, maxAttempts: message.maxAttempts, trace, }); return; } await options.jobs.dispatch( job, message.payload as never, trace ? { trace } : undefined, ); } const MAX_TIMER_DELAY_MS = 2_147_483_647; const OUTBOX_SETTLEMENT_ATTEMPTS = 3; const OUTBOX_SETTLEMENT_RETRY_DELAY_MS = 100; type ResolvedDrainRuntime = { batchSize: number; concurrency: number; leaseMs: number; heartbeatMs: number; maxActiveMs: number; now: () => Date; wait: OutboxDrainWait; }; type MessageDrainOutcome = { delivered: number; retried: number; deadLettered: number; settlementFailed: number; leaseLost: number; }; type ClaimHeartbeat = { lost: Promise; stopAndExtend(): Promise<{ lockedUntil: Date; renewalError?: unknown; renewalRecovered?: boolean; failure?: OutboxLeaseFailure; }>; stop(): Promise; }; type BoundedRenewalResult = | { kind: "succeeded"; lockedUntil: Date } | { kind: "failed"; error: unknown } | { kind: "deadline" } | { kind: "stopped" } | { kind: "waitFailed"; error: unknown }; type BoundedSettlementResult = | { kind: "succeeded" } | { kind: "failed"; error: unknown } | { kind: "deadline" } | { kind: "waitFailed"; error: unknown }; function defaultOutboxWait( delayMs: number, signal: AbortSignal, ): Promise { return new Promise((resolve) => { if (signal.aborted) { resolve(); return; } let timer: ReturnType | undefined; const finish = () => { if (timer !== undefined) clearTimeout(timer); signal.removeEventListener("abort", finish); resolve(); }; timer = setTimeout(finish, delayMs); signal.addEventListener("abort", finish, { once: true }); }); } function readOutboxNow(now: () => Date): Date { const value = now(); assertValidDate("now", value); return value; } function assertTimerDuration(name: string, value: number): void { assertPositiveInteger(name, value); if (value > MAX_TIMER_DELAY_MS) { throw new Error(`${name} must be at most ${MAX_TIMER_DELAY_MS}`); } } function resolveDrainRuntime( options: DrainOutboxOptions, ): ResolvedDrainRuntime { const batchSize = options.batchSize ?? DEFAULT_OUTBOX_BATCH_SIZE; const concurrency = options.concurrency ?? DEFAULT_OUTBOX_CONCURRENCY; const leaseMs = options.leaseMs ?? DEFAULT_OUTBOX_LEASE_MS; const heartbeatMs = options.heartbeatMs ?? Math.max(1, Math.floor(leaseMs / 3)); const maxActiveMs = options.maxActiveMs ?? DEFAULT_OUTBOX_MAX_ACTIVE_MS; assertPositiveInteger("batchSize", batchSize); assertPositiveInteger("concurrency", concurrency); if (concurrency > batchSize) { throw new Error("concurrency must be less than or equal to batchSize"); } assertTimerDuration("leaseMs", leaseMs); if (leaseMs < 2) throw new Error("leaseMs must be at least 2"); assertTimerDuration("heartbeatMs", heartbeatMs); if (heartbeatMs >= leaseMs) { throw new Error("heartbeatMs must be shorter than leaseMs"); } assertTimerDuration("maxActiveMs", maxActiveMs); return { batchSize, concurrency, leaseMs, heartbeatMs, maxActiveMs, now: options.now ?? (() => new Date()), wait: options.wait ?? defaultOutboxWait, }; } async function notifyLeaseFailure( options: DrainOutboxOptions, failure: OutboxLeaseFailure, ): Promise { try { await options.onLeaseError?.(failure); } catch { // Lease observers cannot change delivery ownership or recovery behavior. } } async function notifySettlementFailure( options: DrainOutboxOptions, failure: OutboxSettlementFailure, ): Promise { try { await options.onSettlementError?.(failure); } catch { // Settlement observers cannot replace the unknown storage outcome. } } async function runBoundedClaimRenewal(options: { outbox: OutboxPort; message: ClaimedOutboxMessage; runtime: ResolvedDrainRuntime; renewalAt: Date; deadline: Date; signal: AbortSignal; }): Promise { const remainingMs = options.deadline.getTime() - options.renewalAt.getTime(); if (remainingMs <= 0) return { kind: "deadline" }; if (options.signal.aborted) return { kind: "stopped" }; const waitController = new AbortController(); const stopWaiting = () => waitController.abort(); options.signal.addEventListener("abort", stopWaiting, { once: true }); const renewal = Promise.resolve() .then(() => options.outbox.renewClaim({ id: options.message.id, claimToken: options.message.claimToken, now: options.renewalAt, leaseMs: options.runtime.leaseMs, }), ) .then( (result) => ({ kind: "succeeded" as const, result }), (error: unknown) => ({ kind: "failed" as const, error }), ); const deadline = Promise.resolve() .then(() => options.runtime.wait(remainingMs, waitController.signal)) .then( () => options.signal.aborted ? { kind: "stopped" as const } : { kind: "deadline" as const }, (error: unknown) => options.signal.aborted ? { kind: "stopped" as const } : { kind: "waitFailed" as const, error }, ); const result = await Promise.race([renewal, deadline]); options.signal.removeEventListener("abort", stopWaiting); waitController.abort(); if (result.kind !== "succeeded") return result; try { const completedAt = readOutboxNow(options.runtime.now); if (completedAt.getTime() >= options.deadline.getTime()) { return { kind: "deadline" }; } assertValidDate("lockedUntil", result.result.lockedUntil); if (result.result.lockedUntil.getTime() <= completedAt.getTime()) { return { kind: "failed", error: new Error( `Outbox claim renewal for message "${options.message.id}" returned an expired lease.`, ), }; } return { kind: "succeeded", lockedUntil: cloneDate(result.result.lockedUntil), }; } catch (error) { return { kind: "failed", error }; } } function createClaimHeartbeat( options: DrainOutboxOptions, runtime: ResolvedDrainRuntime, message: ClaimedOutboxMessage, activeUntil: Date, ): ClaimHeartbeat { const controller = new AbortController(); let stopped = false; let lockedUntil = cloneDate(message.lockedUntil); let renewing = false; let firstRenewalError: unknown; let terminalFailure: OutboxLeaseFailure | undefined; let resolveLost: (failure: OutboxLeaseFailure) => void = () => {}; const lost = new Promise((resolve) => { resolveLost = resolve; }); const fail = (failure: OutboxLeaseFailure) => { if (terminalFailure) return; terminalFailure = failure; stopped = true; controller.abort(); resolveLost(failure); }; const run = async () => { try { let delayMs = runtime.heartbeatMs; while (!stopped) { try { await runtime.wait(delayMs, controller.signal); } catch (error) { if (stopped || controller.signal.aborted) return; fail({ error: new OutboxLeaseLostError({ id: message.id, message: `Outbox heartbeat scheduling failed for message "${message.id}".`, cause: error, }), message, operation: "renewClaim", state: "lost", confirmedLost: false, }); return; } if (stopped || controller.signal.aborted) return; const renewalAt = readOutboxNow(runtime.now); const renewalDeadline = new Date( Math.min(lockedUntil.getTime(), activeUntil.getTime()), ); const activeDeadlineEndsFirst = activeUntil.getTime() <= lockedUntil.getTime(); if (renewalAt.getTime() >= renewalDeadline.getTime()) { fail({ error: activeDeadlineEndsFirst ? new OutboxMaxActiveDurationError({ id: message.id, maxActiveMs: runtime.maxActiveMs, }) : new OutboxLeaseLostError({ id: message.id }), message, operation: activeDeadlineEndsFirst ? "maxActiveDuration" : "renewClaim", state: "lost", confirmedLost: !activeDeadlineEndsFirst, }); return; } renewing = true; let renewal: BoundedRenewalResult; try { renewal = await runBoundedClaimRenewal({ outbox: options.outbox, message, runtime, renewalAt, deadline: renewalDeadline, signal: controller.signal, }); } catch (error) { renewal = { kind: "failed", error }; } finally { renewing = false; } if (renewal.kind === "stopped") return; if (renewal.kind === "waitFailed") { fail({ error: new OutboxLeaseLostError({ id: message.id, message: `Could not enforce the renewal deadline for outbox message "${message.id}".`, cause: renewal.error, }), message, operation: "renewClaim", state: "lost", confirmedLost: false, }); return; } if (renewal.kind === "deadline") { fail({ error: activeDeadlineEndsFirst ? new OutboxMaxActiveDurationError({ id: message.id, maxActiveMs: runtime.maxActiveMs, }) : new OutboxLeaseLostError({ id: message.id, message: `Outbox claim renewal for message "${message.id}" did not complete before the lease deadline.`, }), message, operation: activeDeadlineEndsFirst ? "maxActiveDuration" : "renewClaim", state: "lost", confirmedLost: false, }); return; } if (renewal.kind === "succeeded") { lockedUntil = renewal.lockedUntil; delayMs = runtime.heartbeatMs; continue; } const error = renewal.error; if (error instanceof OutboxClaimError) { fail({ error, message, operation: "renewClaim", state: "lost", confirmedLost: true, }); return; } firstRenewalError ??= error; if (stopped) return; const retryAt = readOutboxNow(runtime.now); const remainingMs = renewalDeadline.getTime() - retryAt.getTime(); if (remainingMs <= 0) { fail({ error: activeDeadlineEndsFirst ? new OutboxMaxActiveDurationError({ id: message.id, maxActiveMs: runtime.maxActiveMs, }) : new OutboxLeaseLostError({ id: message.id, message: `Outbox claim renewal for message "${message.id}" did not recover before the lease expired.`, cause: error, }), message, operation: activeDeadlineEndsFirst ? "maxActiveDuration" : "renewClaim", state: "lost", confirmedLost: false, }); return; } delayMs = Math.max( 1, Math.min(runtime.heartbeatMs, Math.floor(remainingMs / 3)), ); } } catch (error) { if (stopped || controller.signal.aborted) return; fail({ error: new OutboxLeaseLostError({ id: message.id, message: `Outbox claim renewal failed unexpectedly for message "${message.id}".`, cause: error, }), message, operation: "renewClaim", state: "lost", confirmedLost: false, }); } }; const running = run(); const stop = async (abortRenewal = true) => { stopped = true; if (abortRenewal || !renewing) controller.abort(); await running; controller.abort(); }; return { lost, stop: () => stop(true), async stopAndExtend() { await stop(false); if (terminalFailure) { return { lockedUntil, failure: terminalFailure }; } const renewalAt = readOutboxNow(runtime.now); const activeDeadlineEndsFirst = activeUntil.getTime() <= lockedUntil.getTime(); const finalRenewalDeadline = new Date( Math.min(lockedUntil.getTime(), activeUntil.getTime()), ); if (renewalAt.getTime() >= finalRenewalDeadline.getTime()) { return { lockedUntil, failure: { error: activeDeadlineEndsFirst ? new OutboxMaxActiveDurationError({ id: message.id, maxActiveMs: runtime.maxActiveMs, }) : new OutboxLeaseLostError({ id: message.id }), message, operation: activeDeadlineEndsFirst ? "maxActiveDuration" : "renewClaim", state: "lost", confirmedLost: !activeDeadlineEndsFirst, }, }; } const finalController = new AbortController(); const renewal = await runBoundedClaimRenewal({ outbox: options.outbox, message, runtime, renewalAt, deadline: finalRenewalDeadline, signal: finalController.signal, }); finalController.abort(); if (renewal.kind === "succeeded") { return { lockedUntil: renewal.lockedUntil, renewalError: firstRenewalError, renewalRecovered: firstRenewalError === undefined ? undefined : true, }; } if (renewal.kind === "failed") { if (!(renewal.error instanceof OutboxClaimError)) { const failedAt = readOutboxNow(runtime.now); if ( failedAt.getTime() < lockedUntil.getTime() && failedAt.getTime() < activeUntil.getTime() ) { return { lockedUntil, renewalError: renewal.error, renewalRecovered: false, }; } if ( activeDeadlineEndsFirst && failedAt.getTime() >= activeUntil.getTime() ) { return { lockedUntil, failure: { error: new OutboxMaxActiveDurationError({ id: message.id, maxActiveMs: runtime.maxActiveMs, }), message, operation: "maxActiveDuration", state: "lost", confirmedLost: false, }, }; } } return { lockedUntil, failure: { error: renewal.error, message, operation: "renewClaim", state: "lost", confirmedLost: renewal.error instanceof OutboxClaimError, }, }; } return { lockedUntil, failure: { error: renewal.kind === "waitFailed" ? new OutboxLeaseLostError({ id: message.id, message: activeDeadlineEndsFirst ? `Could not enforce the maximum active duration for outbox message "${message.id}".` : `Could not enforce the final renewal deadline for outbox message "${message.id}".`, cause: renewal.error, }) : activeDeadlineEndsFirst ? new OutboxMaxActiveDurationError({ id: message.id, maxActiveMs: runtime.maxActiveMs, }) : new OutboxLeaseLostError({ id: message.id, message: `Outbox claim renewal for message "${message.id}" did not complete before settlement.`, }), message, operation: activeDeadlineEndsFirst ? "maxActiveDuration" : "renewClaim", state: "lost", confirmedLost: false, }, }; }, }; } async function settleClaim(options: { operation: OutboxSettlementFailure["operation"]; message: ClaimedOutboxMessage; lockedUntil: Date; activeUntil: Date; runtime: ResolvedDrainRuntime; settle(now: Date): Promise; }): Promise<{ ok: true } | { ok: false; error: unknown; claimLost: boolean }> { let lastError: unknown; for (let attempt = 1; attempt <= OUTBOX_SETTLEMENT_ATTEMPTS; attempt += 1) { const settlementAt = readOutboxNow(options.runtime.now); const settlementDeadline = new Date( Math.min(options.lockedUntil.getTime(), options.activeUntil.getTime()), ); const settlement = await runBoundedSettlement({ runtime: options.runtime, settlementAt, deadline: settlementDeadline, settle: options.settle, }); if (settlement.kind === "succeeded") { return { ok: true }; } if (settlement.kind === "deadline") { return { ok: false, error: new OutboxLeaseLostError({ id: options.message.id, message: `Outbox ${options.operation} for message "${options.message.id}" did not complete before its settlement deadline.`, }), claimLost: true, }; } if (settlement.kind === "waitFailed") { return { ok: false, error: settlement.error, claimLost: false }; } const error = settlement.error; lastError = error; if (error instanceof OutboxClaimError) { return { ok: false, error, claimLost: true }; } const remainingMs = settlementDeadline.getTime() - readOutboxNow(options.runtime.now).getTime(); if (attempt >= OUTBOX_SETTLEMENT_ATTEMPTS || remainingMs <= 0) break; const waitController = new AbortController(); try { await options.runtime.wait( Math.max( 1, Math.min( OUTBOX_SETTLEMENT_RETRY_DELAY_MS * 2 ** (attempt - 1), remainingMs, ), ), waitController.signal, ); } catch (error) { return { ok: false, error, claimLost: false }; } } return { ok: false, error: lastError, claimLost: false }; } async function runBoundedSettlement(options: { runtime: ResolvedDrainRuntime; settlementAt: Date; deadline: Date; settle(now: Date): Promise; }): Promise { const remainingMs = options.deadline.getTime() - options.settlementAt.getTime(); if (remainingMs <= 0) return { kind: "deadline" }; const waitController = new AbortController(); const settlement = Promise.resolve() .then(() => options.settle(options.settlementAt)) .then( () => ({ kind: "succeeded" as const }), (error: unknown) => ({ kind: "failed" as const, error }), ); const deadline = Promise.resolve() .then(() => options.runtime.wait(remainingMs, waitController.signal)) .then( () => ({ kind: "deadline" as const }), (error: unknown) => ({ kind: "waitFailed" as const, error }), ); const result = await Promise.race([settlement, deadline]); waitController.abort(); return result; } async function recordDeadLetter( options: DrainOutboxOptions, instrumentation: ReturnType, jobInstrumentation: ReturnType, message: OutboxMessage, error: unknown, details: Record = {}, ): Promise { try { await options.onDeadLetter?.(error, message); } catch { // Dead-letter observers must not change the settled message state. } instrumentation.record({ type: "outbox", ...options.instrumentationContext, messageId: message.id, messageKind: message.kind, messageName: message.name, status: "deadLettered", details: outboxInstrumentationDetails(message, { ...details, error: serializeOutboxError(error), }), }); if (message.kind === "job") { jobInstrumentation.record({ type: "job", ...options.instrumentationContext, jobName: message.name, status: "deadLettered", details: outboxInstrumentationDetails(message, { ...details, error: serializeOutboxError(error), }), }); } } async function processClaimedMessage( options: DrainOutboxOptions, runtime: ResolvedDrainRuntime, instrumentation: ReturnType, jobInstrumentation: ReturnType, tracing: TracingPort | undefined, message: ClaimedOutboxMessage, ): Promise { const emptyOutcome = (): MessageDrainOutcome => ({ delivered: 0, retried: 0, deadLettered: 0, settlementFailed: 0, leaseLost: 0, }); const outcome = emptyOutcome(); const startedAt = readOutboxNow(runtime.now); const activeUntil = new Date(startedAt.getTime() + runtime.maxActiveMs); const heartbeat = createClaimHeartbeat( options, runtime, message, activeUntil, ); const lifetimeController = new AbortController(); const parentTrace = parseTraceCarrier(message.trace); const traceAttributes = { "beignet.outbox.message_kind": message.kind, "beignet.outbox.message_name": message.name, } as const; const delivery = Promise.resolve() .then(() => runWithTracing( tracing, { name: `beignet.outbox deliver ${message.name}`, type: "outbox", kind: "consumer", parent: parentTrace, attributes: traceAttributes, metricAttributes: traceAttributes, }, (span) => deliverOutboxMessage( options, message, captureTraceCarrier(span?.context ?? parentTrace), ), ), ) .then( () => ({ kind: "succeeded" as const }), (error: unknown) => ({ kind: "failed" as const, error }), ); const maximumActive = Promise.resolve() .then(() => runtime.wait(runtime.maxActiveMs, lifetimeController.signal)) .then( () => ({ kind: "maxActive" as const }), (error: unknown) => ({ kind: "waitFailed" as const, error }), ); const leaseLost = heartbeat.lost.then((failure) => ({ kind: "leaseLost" as const, failure, })); const deliveryOutcome = await Promise.race([ delivery, maximumActive, leaseLost, ]); lifetimeController.abort(); if (deliveryOutcome.kind === "maxActive") { await heartbeat.stop(); const failure: OutboxLeaseFailure = { error: new OutboxMaxActiveDurationError({ id: message.id, maxActiveMs: runtime.maxActiveMs, }), message, operation: "maxActiveDuration", state: "lost", confirmedLost: false, }; await notifyLeaseFailure(options, failure); instrumentation.custom({ name: "outbox.lease.lost", label: "Outbox claim no longer confirmed", summary: `Stopped renewing ${message.kind} "${message.name}" after its active-delivery limit`, details: outboxInstrumentationDetails(message, { operation: failure.operation, state: failure.state, error: serializeOutboxError(failure.error), }), }); outcome.leaseLost = 1; return outcome; } if (deliveryOutcome.kind === "waitFailed") { await heartbeat.stop(); const failure: OutboxLeaseFailure = { error: new OutboxLeaseLostError({ id: message.id, message: `Could not enforce the active-delivery limit for outbox message "${message.id}".`, cause: deliveryOutcome.error, }), message, operation: "maxActiveDuration", state: "lost", confirmedLost: false, }; await notifyLeaseFailure(options, failure); instrumentation.custom({ name: "outbox.lease.lost", label: "Outbox claim no longer confirmed", summary: `Could not enforce the active-delivery limit for ${message.kind} "${message.name}"`, details: outboxInstrumentationDetails(message, { operation: failure.operation, state: failure.state, error: serializeOutboxError(failure.error), }), }); outcome.leaseLost = 1; return outcome; } if (deliveryOutcome.kind === "leaseLost") { await heartbeat.stop(); await notifyLeaseFailure(options, deliveryOutcome.failure); instrumentation.custom({ name: "outbox.lease.lost", label: "Outbox claim lost", summary: `Could not keep the claim for ${message.kind} "${message.name}" active`, details: outboxInstrumentationDetails(message, { operation: deliveryOutcome.failure.operation, state: deliveryOutcome.failure.state, confirmedLost: deliveryOutcome.failure.confirmedLost, error: serializeOutboxError(deliveryOutcome.failure.error), }), }); outcome.leaseLost = 1; return outcome; } const lease = await heartbeat.stopAndExtend(); if (lease.failure) { await notifyLeaseFailure(options, lease.failure); instrumentation.custom({ name: "outbox.lease.lost", label: "Outbox claim lost", summary: `Could not confirm the claim for ${message.kind} "${message.name}" before settlement`, details: outboxInstrumentationDetails(message, { operation: lease.failure.operation, state: lease.failure.state, confirmedLost: lease.failure.confirmedLost, error: serializeOutboxError(lease.failure.error), }), }); outcome.leaseLost = 1; return outcome; } const reportRenewalOutcome = async () => { if (lease.renewalError === undefined) return; const failure: OutboxLeaseFailure = { error: lease.renewalError, message, operation: "renewClaim", state: lease.renewalRecovered ? "recovered" : "degraded", confirmedLost: false, }; await notifyLeaseFailure(options, failure); instrumentation.custom({ name: lease.renewalRecovered ? "outbox.lease.renewal.recovered" : "outbox.lease.renewal.degraded", label: lease.renewalRecovered ? "Outbox claim renewal recovered" : "Outbox claim renewal degraded", summary: lease.renewalRecovered ? `Recovered claim renewal for ${message.kind} "${message.name}"` : `Continued ${message.kind} "${message.name}" settlement under its last confirmed lease`, details: outboxInstrumentationDetails(message, { state: failure.state, error: serializeOutboxError(lease.renewalError), }), }); }; if (deliveryOutcome.kind === "succeeded") { const settlement = await settleClaim({ operation: "markDelivered", message, lockedUntil: lease.lockedUntil, activeUntil, runtime, settle: (now) => options.outbox.markDelivered({ id: message.id, claimToken: message.claimToken, now, }), }); await reportRenewalOutcome(); if (!settlement.ok) { const failure: OutboxSettlementFailure = { error: settlement.error, message, operation: "markDelivered", deliverySucceeded: true, }; await notifySettlementFailure(options, failure); instrumentation.custom({ name: "outbox.settlement.failed", label: "Outbox settlement failed", summary: `Delivered ${message.kind} "${message.name}", but could not confirm its durable acknowledgement`, details: outboxInstrumentationDetails(message, { operation: failure.operation, deliverySucceeded: true, settlementError: serializeOutboxError(settlement.error), }), }); outcome.settlementFailed = 1; if (settlement.claimLost) outcome.leaseLost = 1; return outcome; } instrumentation.record({ type: "outbox", ...options.instrumentationContext, messageId: message.id, messageKind: message.kind, messageName: message.name, status: "delivered", details: outboxInstrumentationDetails(message), }); outcome.delivered = 1; return outcome; } const deliveryError = deliveryOutcome.error; const failedAt = readOutboxNow(runtime.now); const shouldRetry = shouldRetryOutboxMessage(options, message, deliveryError); const deadLetter = !shouldRetry; const retryDelayMs = deadLetter ? 0 : resolveRetryDelayMs(options, message, deliveryError, failedAt); const retryAt = deadLetter ? undefined : new Date(failedAt.getTime() + retryDelayMs); const settlement = await settleClaim({ operation: "markFailed", message, lockedUntil: lease.lockedUntil, activeUntil, runtime, settle: (now) => options.outbox.markFailed({ id: message.id, claimToken: message.claimToken, error: deliveryError, deadLetter, now, retryAt, }), }); await reportRenewalOutcome(); try { await options.onError?.(deliveryError, message); } catch { // Delivery observers cannot change retry, dead-letter, or recovery state. } if (!settlement.ok) { const failure: OutboxSettlementFailure = { error: settlement.error, message, operation: "markFailed", deliverySucceeded: false, deliveryError, }; await notifySettlementFailure(options, failure); instrumentation.custom({ name: "outbox.settlement.failed", label: "Outbox settlement failed", summary: `Could not settle failed ${message.kind} "${message.name}"`, details: outboxInstrumentationDetails(message, { operation: failure.operation, deliverySucceeded: false, deliveryError: serializeOutboxError(deliveryError), settlementError: serializeOutboxError(settlement.error), }), }); outcome.settlementFailed = 1; if (settlement.claimLost) outcome.leaseLost = 1; return outcome; } if (deadLetter) { await recordDeadLetter( options, instrumentation, jobInstrumentation, message, deliveryError, ); outcome.deadLettered = 1; return outcome; } instrumentation.record({ type: "outbox", ...options.instrumentationContext, messageId: message.id, messageKind: message.kind, messageName: message.name, status: "retryScheduled", details: outboxInstrumentationDetails(message, { retryDelayMs, retryAt: retryAt?.toISOString(), error: serializeOutboxError(deliveryError), }), }); if (message.kind === "job") { jobInstrumentation.record({ type: "job", ...options.instrumentationContext, jobName: message.name, status: "retryScheduled", details: outboxInstrumentationDetails(message, { retryDelayMs, retryAt: retryAt?.toISOString(), error: serializeOutboxError(deliveryError), }), }); } outcome.retried = 1; return outcome; } /** * Claim and deliver one bounded set of outbox messages. * * The drain claims only enough messages to fill active delivery slots and * renews each active claim until delivery settles or reaches its configured * maximum duration. It remains an at-least-once transport: a process can still * terminate after the external effect succeeds but before acknowledgement. */ export async function drainOutbox( options: DrainOutboxOptions, ): Promise { const runtime = resolveDrainRuntime(options); assertOutboxDrainPort(options.outbox); assertOutboxDeliveryCapabilities(options); const instrumentation = createProviderInstrumentation( options.instrumentation, { providerName: "outbox", watcher: "outbox", }, ); const jobInstrumentation = createProviderInstrumentation( options.instrumentation, { providerName: "outbox", watcher: "jobs", }, ); const tracing = resolveTracingPort(options.instrumentation); const result: DrainOutboxResult = { claimed: 0, delivered: 0, retried: 0, deadLettered: 0, abandonedDeadLettered: 0, settlementFailed: 0, leaseLost: 0, }; let remaining = runtime.batchSize; let sourceExhausted = false; while (remaining > 0 && !sourceExhausted) { const active: ClaimedOutboxMessage[] = []; const abandoned: OutboxMessage[] = []; while (active.length < runtime.concurrency && remaining > 0) { const limit = Math.min(runtime.concurrency - active.length, remaining); const selected = await options.outbox.claimBatch({ limit, now: readOutboxNow(runtime.now), leaseMs: runtime.leaseMs, }); const selectedCount = selected.claimed.length + selected.deadLettered.length; if (selectedCount === 0) { sourceExhausted = true; break; } if (selectedCount > limit) { throw new Error( `Outbox claimBatch returned ${selectedCount} messages for a limit of ${limit}.`, ); } remaining -= selectedCount; result.claimed += selected.claimed.length; active.push(...selected.claimed); abandoned.push(...selected.deadLettered); } // Construct active delivery promises first so their heartbeats protect // freshly claimed rows while abandoned-message observers run. const outcomePromises = active.map((message) => processClaimedMessage( options, runtime, instrumentation, jobInstrumentation, tracing, message, ), ); const abandonedPromises = abandoned.map(async (message) => { const error = new OutboxAbandonedClaimError({ id: message.id, attempts: message.attempts, maxAttempts: message.maxAttempts, }); await recordDeadLetter( options, instrumentation, jobInstrumentation, message, error, { abandoned: true }, ); }); const [outcomes] = await Promise.all([ Promise.all(outcomePromises), Promise.all(abandonedPromises), ]); result.deadLettered += abandoned.length; result.abandonedDeadLettered += abandoned.length; for (const outcome of outcomes) { result.delivered += outcome.delivered; result.retried += outcome.retried; result.deadLettered += outcome.deadLettered; result.settlementFailed += outcome.settlementFailed; result.leaseLost += outcome.leaseLost; } } return result; } function assertOutboxDrainPort(outbox: OutboxPort): void { const candidate = outbox as unknown as Record; const missing = [ "claimBatch", "renewClaim", "markDelivered", "markFailed", ].filter((method) => typeof candidate[method] !== "function"); if (missing.length > 0) { throw new Error( `Cannot drain this outbox: the outbox port is missing ${missing .map((method) => `${method}()`) .join(", ")}.`, ); } } function assertOutboxDeliveryCapabilities(options: DrainOutboxOptions): void { const missing: string[] = []; if (options.registry.events.size > 0 && !options.eventBus) { missing.push("events require an event bus"); } if (options.registry.jobs.size > 0 && !options.jobs) { missing.push("jobs require a job dispatcher"); } if (missing.length > 0) { throw new OutboxRegistryError( `Cannot drain this outbox registry: ${missing.join("; ")}.`, ); } } /** * Domain event recorder port re-exported for outbox integrations. */ export type { DomainEventRecorderPort };