/** * sqlfu/outbox — a small transactional-outbox / job-queue built on any sqlfu * `Client` (sync or async). Sync clients get plain return values from * `setup()`, `emit()`, and `claim()`; async clients get Promises. `tick()` * is always async because consumer handlers are. * * SQLite serialises writers for us, so claim-and-lease works as a plain * `BEGIN; select pending; update to running; commit` — no row-locking dance * needed. * * Zero Node-only dependencies: causation propagation is explicit (the handler * receives an `emit` helper already bound to its own job context) rather than * via `AsyncLocalStorage`, which keeps this module runnable in browsers, edge * workers, and anywhere else sqlfu already runs. * * The body of each method is written once as a generator. `client.sync` picks * the right driver (`driveSync` / `driveAsync`) so the runtime shape matches * the input client without two parallel implementations to keep in step. */ import type { Client, SyncClient } from '../types.js'; export type TimeUnit = 's' | 'm' | 'h' | 'd'; export type TimePeriod = `${number}${TimeUnit}`; export type JobContext = { jobId: number; eventId: number; attempt: number; consumerName: string; }; export type Causation = { eventId: number; consumerName: string; jobId: number; }; export type WhenFn = (input: { payload: TPayload; }) => boolean | null | undefined | ''; export type DelayFn = (input: { payload: TPayload; }) => TimePeriod; export type RetryOutcome = { retry: false; reason: string; delay?: never; } | { retry: true; reason: string; delay: TimePeriod; }; export type RetryFn = (job: JobContext, error: unknown) => RetryOutcome; export type EmitInput = { name: K; payload: TEvents[K]; }; export type EmitOptions = { /** Pass a transaction client to make emit atomic with the surrounding domain write. */ client?: TClient; }; export type EmitResult = { eventId: number; }; /** * Pick `TSync` for `SyncClient`s, `TAsync` for `AsyncClient`s. When `TClient` * is the open `Client` union, `TSync | TAsync` falls out — that's what * external callers typing `Outbox` (without parameterising over the * client) get. */ type MaybeAsync = TClient extends SyncClient ? TSync : TAsync; export type EmitFn = (event: EmitInput, options?: EmitOptions) => MaybeAsync>; export type ConsumerHandlerInput = { payload: TPayload; eventId: number; eventName: string; job: { id: number; attempt: number; }; /** * Emit a follow-up event from inside this handler. This `emit` is pre-bound * to the running job's causation, so the downstream event's * `context.causedBy` points back to this job/consumer/event automatically. * * Always async-shaped, regardless of the underlying client: handlers are * always async, so awaiting the bound emit costs nothing and keeps consumer * code uniform. */ emit: EmitFn; }; export type ConsumerDefinition = { name: string; when?: WhenFn; delay?: DelayFn; retry?: RetryFn; visibilityTimeout?: TimePeriod; handler: (input: ConsumerHandlerInput) => Promise; }; export type OutboxDefaults = { visibilityTimeout?: TimePeriod; retry?: RetryFn; /** Hard cap on attempts: once reached after a failed run, the job transitions to `failed` regardless of `retry`. */ maxAttempts?: number; /** Application environment tag written onto every event — useful when the same DB sees events from dev + CI. */ environment?: string; /** How many jobs to claim per tick(). */ batchSize?: number; /** * Hook for bookkeeping errors — updates that fail after a handler has * already run its side effects. Default is `console.warn`; pass a custom * logger or Sentry hook to surface these in production. Distinct from * handler errors, which are routed through the retry policy. */ onBookkeepingError?: (error: unknown, job: JobContext) => void; }; export type EventMap = Record; export type OutboxConsumers = { [K in keyof TEvents]?: ConsumerDefinition[]; }; export type OutboxConfig = { client: TClient; consumers: OutboxConsumers; now?: () => Date; defaults?: OutboxDefaults; }; export type ClaimedJob = { id: number; event_id: number; consumer_name: string; event_name: string; event_payload: string; event_context: string; attempt: number; vt_until: number; }; export type TickResult = { claimed: number; succeeded: number; failed: number; retried: number; }; export interface Outbox { setup(): MaybeAsync>; emit: EmitFn; tick(): Promise; claim(input?: { limit?: number; }): MaybeAsync>; } export declare function defineConsumer(definition: ConsumerDefinition): ConsumerDefinition; export declare function createOutbox(config: OutboxConfig): Outbox; export {};