Outbox

The outbox pattern records side-effect intent in the same database transaction as the business write. A separate worker drains those records after commit and delivers events or jobs with retries.

Use an outbox when an event or job must not be lost after the database commit: notifications, integrations, billing syncs, audit streams, search indexing, or other workflow-critical side effects.

bun add @beignet/core

Why it exists

After-commit publishing avoids one class of bug: listeners, jobs, and mail do not see data that later rolls back. It still leaves a different failure window:

  1. The database transaction commits.
  2. The process starts publishing the event or dispatching the job.
  3. The process crashes or the provider call fails.
  4. The business record is durable, but the side effect is lost.

The outbox closes that gap by writing the event or job to an outbox_messages table inside the transaction. Delivery becomes a retryable background workflow.

Outbox delivery is at least once, not exactly once. If the worker delivers a message and crashes before marking it delivered, the message may be delivered again. Use Idempotency inside listeners or job handlers when duplicate delivery would be harmful.

Use Workflow primitives to decide whether the side effect should be modeled as an event, job, notification, idempotent command, schedule, or outbox record, and see side effects after commit for the rule the outbox makes durable. Payment workflows often pair verified webhooks with outbox-backed entitlement, notification, or integration side effects; see Payments and billing.

Add an outbox to the starter

Start with the working SQLite Todos app from Quickstart. This example records TodoCompleted when an owner completes a todo and delivers it to a logging listener. Logging makes delivery visible without another service; required integrations use the same transaction and drain path.

Run these commands from the app root:

bun beignet make event todos.completed
bun beignet make listener todos.log-completed --event todos.completed
bun install
bun beignet db generate
bun beignet db migrate
bun beignet db status

The event generator also adds the outbox table schema, app ports, server/outbox.ts, and app/api/cron/outbox/drain/route.ts. The listener generator registers the listener and its lifecycle provider. Keep the generated event and listener: both use a payload containing the todo's id. Copy the generated CRON_SECRET entry from .env.example into .env.local for the drain route. Review and commit the generated migration alongside the code changes.

The generated listener error callback logs and reports errors. For this durable workflow, also rethrow the error so a failed listener makes the drain retry its message. Apply this change inside listenersProvider in server/providers.ts:

--- a/server/providers.ts
+++ b/server/providers.ts
@@ -85,4 +85,5 @@
               },
             });
+            throw error;
           },
         });

The default memory event bus waits for handlers. Keep that delivery mode for this example; detached delivery cannot tell the drain whether a listener failed. For a remote event bus, an outbox acknowledgement means the transport accepted the event, and consumer retries belong to that transport's delivery model.

Transaction-scoped recording

Keep the event API in the use case. Complete the business write and record TodoCompleted through tx.events inside the same transaction.

Record the event

Update updateTodoUseCase in features/todos/use-cases.ts and import the event:

--- a/features/todos/use-cases.ts
+++ b/features/todos/use-cases.ts
@@ -1,3 +1,4 @@
+import { TodoCompleted } from "./domain/events/completed";
 import "@beignet/core/server-only";
 import { requireUser } from "@/lib/auth";
 import { useCase } from "@/lib/use-case";
@@ -55,7 +56,8 @@
   .command("todos.update")
   .input(UpdateTodoInputSchema)
   .output(TodoSchema)
-  .run(async ({ ctx, input }) => {
+  .emits([TodoCompleted])
+  .run(async ({ ctx, input, events }) => {
     requireUser(ctx);

     return ctx.ports.uow.transaction(async (tx) => {
@@ -66,6 +68,10 @@

       await ctx.gate.authorize("todos.update", todo);

-      return tx.todos.update(input.id, { completed: input.completed });
+      const updated = await tx.todos.update(input.id, { completed: input.completed });
+      if (!todo.completed && updated.completed) {
+        await events.record(tx.events, TodoCompleted, { id: updated.id });
+      }
+      return updated;
     });
   });

Authorization still runs before the write. The event is recorded only when a todo changes from incomplete to complete. Retrying a completed update does not record another event, but delivery of an existing event can still happen more than once.

Connect transaction ports

Expose an event recorder on the transaction ports, then wire it to the outbox in infrastructure. Apply these changes to the generated port and repository factory:

--- a/ports/index.ts
+++ b/ports/index.ts
@@ -1,2 +1,3 @@
+import type { DomainEventRecorderPort } from "@beignet/core/ports";
 import type { ErrorReporterPort } from "@beignet/core/error-reporting";
 import type { IdempotencyPort } from "@beignet/core/idempotency";
@@ -14,4 +15,5 @@

 export type AppTransactionPorts = {
+  events: DomainEventRecorderPort;
   idempotency: IdempotencyPort;
   todos: TodoRepository;
--- a/infra/db/repositories.ts
+++ b/infra/db/repositories.ts
@@ -6,5 +6,5 @@
 export function createRepositories(
   db: DrizzleSqliteDatabase<typeof schema>,
-): Omit<AppTransactionPorts, "idempotency" | "outbox"> {
+): Omit<AppTransactionPorts, "idempotency" | "outbox" | "events"> {
   return {
     todos: createDrizzleTodoRepository(db),

Keep the transaction factory in infra/db/provider.ts so you can continue using beignet make event and beignet make job as your app grows. Export it to use the same transaction ports in your database tests. Every repository and outbox adapter below uses the same transaction handle:

--- a/infra/db/provider.ts
+++ b/infra/db/provider.ts
@@ -1,2 +1,4 @@
 import "@beignet/core/server-only";
+import { createOutboxEventRecorder } from "@beignet/core/outbox";
+import type { UnitOfWorkPort } from "@beignet/core/ports";
 import { createProvider } from "@beignet/core/providers";
@@ -1,1 +1,1 @@
-import type { AppPorts } from "@/ports";
+import type { AppPorts, AppTransactionPorts } from "@/ports";
@@ -36,11 +36,4 @@
       outbox,
       outboxAdmin,
-      uow: createDrizzleSqliteUnitOfWork({
-        db: dbPort.drizzle,
-        createTransactionPorts: (tx) => ({
-          ...createRepositories(tx),
-          idempotency: createDrizzleSqliteIdempotencyPort(tx),
-          outbox: createDrizzleSqliteOutboxPort(tx),
-        }),
-      }),
+      uow: createAppUnitOfWork(dbPort.drizzle),
     };
@@ -1,1 +1,15 @@
 });
+
+export function createAppUnitOfWork(
+  db: DbPort<typeof schema>["drizzle"],
+): UnitOfWorkPort<AppTransactionPorts> {
+  return createDrizzleSqliteUnitOfWork({
+    db,
+    createTransactionPorts: (tx) => ({
+      ...createRepositories(tx),
+      idempotency: createDrizzleSqliteIdempotencyPort(tx),
+      outbox: createDrizzleSqliteOutboxPort(tx),
+      events: createOutboxEventRecorder(createDrizzleSqliteOutboxPort(tx)),
+    }),
+  });
+}

Keep tests aligned

Give the existing in-memory use-case and route tests a transaction-local recorder too:

--- a/features/todos/tests/todos.test.ts
+++ b/features/todos/tests/todos.test.ts
@@ -1,2 +1,3 @@
+import { createDomainEventRecorder } from "@beignet/core/ports";
 import { describe, expect, it } from "@/lib/beignet-test";
 import { createUseCaseTester } from "@beignet/core/application";
@@ -38,4 +39,5 @@
       ports: (ports) => ({
         ...ports,
+        events: createDomainEventRecorder(),
         todos,
       }),
--- a/features/todos/tests/routes.test.ts
+++ b/features/todos/tests/routes.test.ts
@@ -1,2 +1,3 @@
+import { createDomainEventRecorder } from "@beignet/core/ports";
 import { describe, expect, it } from "@/lib/beignet-test";
 import { createStaticAuth } from "@beignet/core/ports";
@@ -41,4 +42,5 @@
       ports: (ports) => ({
         ...ports,
+        events: createDomainEventRecorder(),
         todos,
       }),

Verify delivery and rollback

bun run typecheck
bun run test
bun beignet doctor --strict
bun run dev

Sign in, create a todo, and complete it. Expect the checkbox to remain completed after a reload and a Listener handled log naming todos.log-completed. The generated Next.js provider requests a bounded drain after commit, so the message may already be delivered by the time you inspect it. From another terminal in the app directory, run:

bun beignet outbox list --status pending
bun beignet outbox drain --batch-size 100
bun beignet outbox list --status delivered

Expect a delivered todos.completed row. A drain reporting zero claimed rows is normal if the after-commit trigger already handled it. In production, schedule the generated authenticated drain route as a recovery sweep; the after-commit trigger alone cannot guarantee delivery. See Push-assisted polling in Next.js.

To check the transaction guarantee, add this complete test. It uses the starter's isolated SQLite database and the same Unit of Work factory as the app. It also checks that a failed delivery is retried before the row becomes delivered:

// infra/db/outbox.test.ts
import { expect, test } from "@/lib/beignet-test";
import { createMemoryEventBus } from "@beignet/provider-event-bus-memory";
import { defineOutboxRegistry, drainOutbox } from "@beignet/core/outbox";
import {
  createDrizzleSqliteOutboxAdminPort,
  createDrizzleSqliteOutboxPort,
} from "@beignet/provider-db-drizzle/sqlite";
import { TodoCompleted } from "@/features/todos/domain/events/completed";
import { user } from "./schema";
import { createTestDatabase } from "./test-database";
import { createAppUnitOfWork } from "./provider";

test("todo and event commit together; failed delivery retries; rollback records neither", async () => {
  const database = await createTestDatabase();
  const eventBus = createMemoryEventBus();
  const received: string[] = [];
  let failDelivery = true;
  const subscription = eventBus.subscribe(TodoCompleted, ({ id }) => {
    if (failDelivery) throw new Error("Temporary delivery failure");
    received.push(id);
  });
  await subscription.ready;
  try {
    await database.db.insert(user).values({
      id: "outbox-owner", name: "Owner", email: "outbox@example.com",
      createdAt: new Date(), updatedAt: new Date(),
    });
    const uow = createAppUnitOfWork(database.db);
    const outbox = createDrizzleSqliteOutboxPort(database.db);
    const admin = createDrizzleSqliteOutboxAdminPort(database.db);
    const registry = defineOutboxRegistry({ events: [TodoCompleted] });
    const saved = await uow.transaction(async (tx) => {
      const todo = await tx.todos.create({ userId: "outbox-owner", title: "Durable delivery" });
      await tx.todos.update(todo.id, { completed: true });
      await tx.events.record(TodoCompleted, { id: todo.id });
      return todo;
    });
    expect((await database.repositories.todos.findById(saved.id))?.completed).toBe(true);
    expect(await admin.countMessages({ status: "pending" })).toBe(1);
    expect(received).toEqual([]);

    let now = new Date();
    const failed = await drainOutbox({ outbox, registry, eventBus, now: () => now, retryDelayMs: 1000 });
    expect(failed.retried).toBe(1);
    expect(await admin.countMessages({ status: "pending" })).toBe(1);
    failDelivery = false;
    now = new Date(now.getTime() + 1000);
    const delivered = await drainOutbox({ outbox, registry, eventBus, now: () => now });
    expect(delivered.delivered).toBe(1);
    expect(received).toEqual([saved.id]);
    expect(await admin.countMessages({ status: "delivered" })).toBe(1);

    let rolledBackId = "";
    await expect(uow.transaction(async (tx) => {
      const todo = await tx.todos.create({ userId: "outbox-owner", title: "Roll back" });
      rolledBackId = todo.id;
      await tx.events.record(TodoCompleted, { id: todo.id });
      throw new Error("Cancel transaction");
    })).rejects.toThrow("Cancel transaction");
    expect(await database.repositories.todos.findById(rolledBackId)).toBeNull();
    expect(await admin.countMessages()).toBe(1);
  } finally {
    await subscription.unsubscribe();
    await database.close();
  }
});

Run bun run test infra/db/outbox.test.ts. Expect the commit, retry, delivery, and rollback assertions to pass. Memory-only Unit of Work tests cannot establish the database rollback guarantee.

Transactional job enqueueing

Use createOutboxJobDispatcher(...) when a use case should enqueue a job inside the same transaction as the business write. Wire it as the transaction-scoped jobs port the same way createOutboxEventRecorder(...) backs tx.events above. Add jobs: JobDispatcherPort to AppTransactionPorts, exclude it from the repository factory's return type, and return jobs: createOutboxJobDispatcher(createDrizzleSqliteOutboxPort(tx)) from createAppUnitOfWork's transaction factory in infra/db/provider.ts. After the authorized todo update, enqueue the job from the Jobs guide:

// features/todos/use-cases.ts (excerpt)
await tx.jobs.dispatch(LogCompletionJob, { id: updated.id });

Use this for direct transactional job enqueueing. For event-driven workflows, prefer recording an event and letting a listener enqueue the job during outbox drain.

Registry and draining

Use bun beignet make outbox to add delivery infrastructure without a new event or job. In the starter walkthrough above, the event generator already did this. The generated server/outbox.ts registers feature events and jobs and creates a service context for drains; keep central registration in that file.

After server/outbox.ts exists, beignet make event and beignet make job append new feature registries to it, and beignet doctor warns about feature events and jobs the registry cannot deliver; beignet doctor --fix registers them.

Registration and delivery wiring are separate. Events in the registry require ctx.ports.eventBus; jobs require ctx.ports.jobs. beignet make job declares and defers jobs and installs an app-owned inline dispatcher when no dispatcher is already wired, while preserving an existing direct binding or provider such as Inngest. A later Inngest or BullMQ preset replaces only that marked generated fallback and rejects unmarked custom inline wiring. beignet doctor reports BEIGNET_OUTBOX_JOB_DISPATCHER_MISSING when the registry contains jobs but the configured ports and port-wiring files do not declare and bind or defer that port. Keep the jobs key explicit in definePorts(...); when custom spreads or helper calls prevent doctor from verifying required port wiring, it reports the uncertainty, and beignet make job stops instead of replacing a custom dispatcher it cannot verify.

Drain messages from a cron route, worker process, queue consumer, or scheduled task. In Next.js apps, prefer createOutboxDrainRoute(...) so the outbox runs as a bounded serverless invocation instead of a provider startup loop:

// app/api/cron/outbox/drain/route.ts
import { createOutboxDrainRoute } from "@beignet/next";
import { env } from "@/lib/env";
import { getServer } from "@/server";
import { outboxRegistry } from "@/server/outbox";

export const runtime = "nodejs";

export const { GET, POST } = createOutboxDrainRoute({
  server: getServer,
  registry: outboxRegistry,
  secret: env.CRON_SECRET,
});

The route verifies CRON_SECRET before resolving the server or assembling app context, then runs authenticated requests through the normal raw-route pipeline.

For non-Next runtimes, call drainOutbox(...) from the host's bounded background entrypoint. Do not start setInterval polling from provider lifecycle hooks in serverless apps.

For a local, CI, or worker-hosted drain, use the same server/outbox.ts module with the CLI:

bun beignet outbox drain --batch-size 100 --concurrency 4

The CLI loads outboxRegistry, creates the app context through createOutboxDrainContext(...), drains one batch, records instrumentation, then calls stopOutboxDrainContext(...) when present.

The default batch size is 100 and the default concurrency is 1. Increase --concurrency only when handlers are safe to run in parallel; a value above one does not preserve delivery order. The CLI exposes throughput controls but uses the core lease, heartbeat, and maximum-active defaults. When those limits need tuning, use an app-owned worker that calls drainOutbox(...) or configure createOutboxDrainRoute(...) or createNextOutboxDrainTrigger(...).

See Runtime recipes for the difference between cron routes, worker-hosted drains, and command-based drains.

Push-assisted polling in Next.js

On Next.js 15.1 or newer, use after() to request one bounded drain after a successful Unit of Work transaction. This removes the one-minute polling floor for newly committed messages without running an idle poller. Keep the Next-specific wrapper in server/providers.ts, after the database provider that installs uow:

import { createObservedUnitOfWork } from "@beignet/core/ports";
import { createProvider } from "@beignet/core/providers";
import { createNextOutboxDrainTrigger } from "@beignet/next";
import { after } from "next/server";
import type { AppContext } from "@/app-context";
import type { AppPorts } from "@/ports";
import type { AppServiceContextInput } from "./context";

const outboxDrainProvider = createProvider<
  Pick<AppPorts, "uow">,
  AppContext,
  AppServiceContextInput
>()({
  name: "outbox-drain-trigger",
  setup({ ports, createServiceContext }): { ports: Pick<AppPorts, "uow"> } {
    const trigger: () => void = createNextOutboxDrainTrigger({
      defer: after,
      createContext: () => createServiceContext(undefined),
      registry: async () => (await import("./outbox")).outboxRegistry,
    });

    return {
      ports: {
        uow: createObservedUnitOfWork({
          unitOfWork: ports.uow,
          afterCommit: trigger,
        }),
      },
    };
  },
});

Register outboxDrainProvider after the database provider. This keeps Next-specific server composition out of infra/db/provider.ts; the lazy registry import avoids the cycle between server/providers.ts and server/outbox.ts. Pass the service-context input your app requires; the generated starter accepts undefined, while tenant-aware apps may provide an explicit service actor or tenant.

This model is push-assisted polling, not durable execution. after() may be missed during a crash or deployment, performs only one batch, and does not wait for a message's future availableAt. Keep createOutboxDrainRoute(...) and schedule it as a recovery sweep. About every 15 minutes is a useful default when immediate first delivery comes from after(); shorten it when the app's retry-latency requirement demands it.

Repeated triggers coalesce while a deferred drain is scheduled or running. This prevents transactions started by an outbox handler from recursively scheduling more drain callbacks; the recovery sweep handles messages left beyond the bounded batch.

Delayed messages and scheduled retries still depend on the recovery cron or a durable queue scheduler. Use a durable jobs provider when retries require a guaranteed low-latency wake-up. Apps on older Next.js versions remain valid with cron-only draining.

Delivery guarantees and tuning

drainOutbox(...) first validates that every transport required by the registry is available, then claims messages, validates payloads, publishes events through eventBus, dispatches jobs through jobs, and marks each message delivered. Missing transport wiring therefore fails before a claim and does not consume delivery attempts. Failed deliveries are retried with backoff until maxAttempts, then dead-lettered.

Each active delivery has a renewable lease. By default, Beignet claims for 30 seconds, renews every 10 seconds, and stops renewing a single delivery after 5 minutes. drainOutbox(...) claims only enough work to fill its active concurrency slots, so queued work does not sit behind a serial batch while its lease expires. Set leaseMs, heartbeatMs, and maxActiveMs together when a provider or handler needs different limits. heartbeatMs must be shorter than leaseMs. Reaching maxActiveMs stops further renewals; the last confirmed lease can remain active until its own lockedUntil timestamp.

Outbox adapters compare application-supplied timestamps when claiming and renewing leases. Synchronize every drain host with a reliable clock source. Clock skew can make another worker treat a live claim as expired; leave enough lease margin for expected scheduling delay and host-clock drift. Event outbox rows store parsed, canonical Standard Schema output. Before writing the row, Beignet verifies that the output is JSON-safe and unchanged when its decoded JSON is validated again. An invalid event rejects with EventTransportError instead of creating a row that would fail or change during delivery. The drain repeats the check for existing rows before publishing. Job rows preserve the original JSON-safe payload until the worker's handler-facing parse.

Delivery, lease, and settlement failures are isolated per claimed message. The drain retries a transient settlement write within the confirmed lease. If a successful external delivery cannot be acknowledged, Beignet never calls markFailed(...): doing so could label a delivered side effect as failed. The result instead increments settlementFailed. A claim that cannot be renewed or confirmed increments leaseLost. In either case, the durable final state is unknown and another worker may later deliver the message again.

The drain result contains claimed, delivered, retried, deadLettered, abandonedDeadLettered, settlementFailed, and leaseLost counters. abandonedDeadLettered is included in deadLettered. A Next drain route returns HTTP 500 when either uncertainty counter is nonzero, the CLI prints the complete report and exits nonzero, and the MCP outbox_run tool returns the same JSON report with an error result. Treat those outcomes as an operational incident rather than retrying the command blindly. A settlement that reaches its lease deadline can increment both uncertainty counters for the same message.

When you pass ctx.ports.devtools or another instrumentation port to drainOutbox(...), Beignet records first-class outbox events for delivered, retried, and dead-lettered messages, including attempt counts, retry timing, and a redacted error summary. createOutboxDrainRoute(...) passes the drain request's requestId and traceId into those rows so the devtools request view can expand into the messages delivered by that cron invocation.

When the drain context includes ports.errorReporter, the Next route and CLI runner report dead-lettered messages, lease failures, and settlement failures. Drain-level claim/infrastructure failures are also reported. Scheduled retries remain instrumentation events and do not create incidents. Direct drainOutbox(...) callers can use onDeadLetter, onLeaseError, and onSettlementError to apply the same ownership policy. Lease callbacks classify the current state as recovered, degraded, or lost; only lost increments the uncertainty counter.

Outbox delivery uses the same retry vocabulary as jobs. A retried message is marked pending again with a future availableAt computed from the backoff, maxAttempts caps total delivery attempts, and a dead-lettered message is in the terminal outbox state and is no longer retried automatically.

For job messages, enqueueJob(...) and createOutboxJobDispatcher(...) use the job definition's retry policy by default, and the drain owns execution retries: when delivery goes through the inline dispatcher, the drain detects it and runs the handler exactly once per pass, so the job's policy is applied by outbox rescheduling rather than stacked in-process retries. A configured inline-dispatcher onError observer still runs, but it cannot swallow that single-attempt failure: the drain marks the message failed and retries or dead-letters it. Durable providers are unaffected — for them dispatch is an enqueue and the queue owns execution. Customize the retry delay for outbox delivery when the worker needs an override:

await drainOutbox({
  outbox,
  registry,
  eventBus,
  jobs,
  instrumentation: ports,
  retryDelayMs: ({ message }) => Math.min(60_000, 1000 * message.attempts),
});

Dead-letter recovery and cleanup

server/outbox.ts can use the same service context for drains and admin commands when that context exposes ports.outboxAdmin. The Drizzle providers export createDrizzleSqliteOutboxAdminPort(...), createDrizzlePostgresOutboxAdminPort(...), and createDrizzleMysqlOutboxAdminPort(...) for this root maintenance port.

Inspect dead-lettered rows:

bun beignet outbox list --status deadLettered
bun beignet outbox show <message-id>

Requeue a reviewed dead-lettered message:

bun beignet outbox requeue <message-id> --reset-attempts

Clean up terminal rows only after review:

bun beignet outbox purge --before 2026-01-01T00:00:00.000Z --dry-run
bun beignet outbox purge --before 2026-01-01T00:00:00.000Z

Prune delivered rows by retention cutoff:

bun beignet outbox prune --before 2026-01-01T00:00:00.000Z --dry-run
bun beignet outbox prune --before 2026-01-01T00:00:00.000Z

purge targets only deadLettered rows using updatedAt as the cutoff. prune targets only delivered rows using deliveredAt as the cutoff. Both commands support --limit for bounded maintenance passes, delete the oldest eligible rows first, and support --json for runbooks or dashboards.

Drizzle SQLite

@beignet/provider-db-drizzle includes a durable outbox adapter on its /sqlite subpath:

bun add @beignet/provider-db-drizzle
import {
  createDrizzleSqliteOutboxAdminPort,
  createDrizzleSqliteOutboxPort,
  createDrizzleSqliteOutboxSetupStatements,
} from "@beignet/provider-db-drizzle/sqlite";

const outbox = createDrizzleSqliteOutboxPort(db);
const outboxAdmin = createDrizzleSqliteOutboxAdminPort(db);

Beignet does not hide migrations. Add the setup statements to your app-owned migration flow. The starter walkthrough uses generated Drizzle migrations instead of these setup statements. Do not execute setup DDL at server startup:

for (const statement of createDrizzleSqliteOutboxSetupStatements()) {
  await client.execute(statement);
}

The default table is outbox_messages; pass { tableName: "app_outbox_messages" } to both the setup statements and the port to override it.

Drizzle-backed outboxes store the carrier in nullable trace_context_json. If upgrading a table created before trace propagation was supported, add the missing column through a migration before deploying the updated adapter. Current generators already include it:

-- SQLite and Postgres
ALTER TABLE outbox_messages ADD COLUMN trace_context_json text;

-- MySQL
ALTER TABLE outbox_messages ADD COLUMN trace_context_json longtext;

Testing

Use the memory adapter in a use-case test fixture with an existing todos repository. This excerpt records intent but does not provide SQL rollback:

import {
  createMemoryOutbox,
  createOutboxEventRecorder,
} from "@beignet/core/outbox";

import { createNoopUnitOfWork } from "@beignet/core/ports";

const outbox = createMemoryOutbox();

const uow = createNoopUnitOfWork(() => ({
  todos,
  events: createOutboxEventRecorder(outbox),
  outbox,
}));

Then assert pending messages or drain them:

expect(outbox.messages).toMatchObject([
  {
    kind: "event",
    name: "todos.completed",
    status: "pending",
  },
]);

Use direct in-memory event recorders and inline jobs when durability is not the behavior under test.

Core API

createOutboxEventRecorder(...) exposes only record(...), which writes through the transaction-scoped outbox immediately. Use createDomainEventRecorder() when an in-memory test needs entries(), clear(), or flush().

For trace propagation, pass the installed tracing port as the second argument: createOutboxEventRecorder(outbox, { tracing: ports.tracing }). The job dispatcher accepts the same tracing option. Keep this wiring in infrastructure.

The generated outbox reference lists signatures and options. Adapter implementers can use the operations below.

Use @beignet/core/outbox for typed messages, registries, memory test storage, and the drain worker:

import {
  createMemoryOutbox,
  defineOutboxRegistry,
  drainOutbox,
} from "@beignet/core/outbox";

The app-facing delivery port is intentionally small:

import type { OutboxAdminPort, OutboxPort } from "@beignet/core/outbox";

export type AppPorts = {
  outbox: OutboxPort;
  outboxAdmin: OutboxAdminPort;
};

Production adapters implement:

OperationPurpose
enqueue(...)Store a pending event or job
claimBatch(...)Atomically claim available messages and reconcile eligible rows whose attempt budget is exhausted
renewClaim(...)Extend an active, unexpired claim with the current claim token
markDelivered(...)Ack a claimed message with its claim token
markFailed(...)Retry or dead-letter a claimed message

Renew, ack, and fail operations require both the current claimToken and an unexpired lease. This prevents an old worker from changing a message after its claim expires. claimBatch(...) returns claimed messages separately from rows it moved directly to deadLettered because a crashed worker had already used the final attempt.

Keep OutboxAdminPort on operational contexts rather than transaction-scoped use-case ports. It supports:

OperationPurpose
listMessages(...)Inspect pending, claimed, delivered, or dead-lettered rows
countMessages(...)Count rows before cleanup or alerting
getMessage(...)Inspect one payload, attempts, and last error
requeueMessage(...)Return a dead-lettered row to pending
purgeDeadLettered(...)Delete terminal dead-letter rows after review
pruneDelivered(...)Delete delivered rows older than a retention cutoff

When not to use it

Do not force every side effect through the outbox.

Use direct after-commit event publishing or inline jobs for low-stakes local workflows, tests, and single-process development. Use the outbox when losing the side effect would create user-visible, financial, compliance, or workflow correctness issues.