API/Worker-Split
Run Kumiko as two separate processes — an HTTP API that only enqueues jobs,
and a worker that executes them and applies the read models the API skips.
This is the deploy topology behind runSingleInstance: false +
createWorkerEntrypoint.
What it shows
Section titled “What it shows”- Two entrypoints, one feature —
bin/api.tsstarts the API process (runProdApp({ runSingleInstance: false })),bin/worker.tsthe worker (runWorkerApp). Both build the sameordersfeature; they differ only in what they run. - API = enqueuer-only — the API writes events and pushes
runIn: "worker"jobs onto the BullMQ queue. It consumes nothing and applies no multiStreamProjections. - Worker = the read side — the worker’s event-dispatcher applies the
order-activityprojection and its BullMQ runner executesprocess-order. - The sharp edge, live — with
runSingleInstance: falsethe API process applies no multiStreamProjections. No worker running means theorder-activityread side stays empty, silently (the 2026-06-11 incident class). The integration test asserts exactly this. - Result written back through the worker dispatcher —
JobContexthas nowrite/query(fw#1717).bin/worker.tswires a write-bridge fromdispatchSystemWriteand the job uses it to create afulfillmentrow (same pattern asinbound-mail-foundation/watch-supervisor.ts).
Requires Postgres + Redis (env: DATABASE_URL, REDIS_URL, JWT_SECRET).
bun installbun run schema:apply # creates framework + entity tables (kumiko/migrations)bun run api # terminal 1 — HTTP + enqueuebun run worker # terminal 2 — consumes jobs, applies projectionsProve the topology. orders:write:order:create declares openToAll, which
still requires a valid JWT (any authenticated user, no specific role) — grab
one the same way src/__tests__/feature.integration.test.ts does, via
api.jwt.sign(adminUser) on the entrypoint returned by
createApiEntrypoint:
curl -X POST http://localhost:3000/api/write \ -H 'Content-Type: application/json' \ -H "Authorization: Bearer $TOKEN" \ -d '{"type":"orders:write:order:create","payload":{"customerName":"Acme GmbH","amount":499}}'With only the API running, read_orders gains the row but
read_order_activity stays empty. Start the worker and both the activity
row and the read_fulfillments row appear.
Needs a running Postgres and TEST_DATABASE_URL set (e.g. postgres://postgres:[email protected]:5432/postgres, see demo/.env.example).
bun test src/__tests__/feature.integration.test.tsRuns both entrypoints in-process against real Postgres + Redis, at the
createApiEntrypoint/createWorkerEntrypoint level (not bin/api.ts/
bin/worker.ts themselves), and asserts the API-only read side stays empty
before the worker picks up the job and writes the fulfillment back.
Source code
Section titled “Source code”The feature entry point — embedded straight from the source file, so the code here is exactly what runs. Multi-file samples keep their remaining files next to it on GitHub (link below):
// API/Worker-Split Sample//// Proves the split deployment topology end-to-end://// - the API process runs `runSingleInstance: false` — it serves HTTP,// writes events, and ENQUEUES worker-lane jobs. It applies no// multiStreamProjections and consumes no jobs.// - the worker process consumes the worker-lane queue and applies the// multiStreamProjections the API skipped (the 2026-06-11 sharp edge:// without a worker the read-side stays empty, silently).// - a worker-lane job writes its RESULT back through the WORKER's// dispatcher — JobContext has no write/query (fw#1717), so the app// wires a write-bridge at boot (see setOrderFulfillWrite below).
import { insertOne } from "@cosmicdrift/kumiko-framework/bun-db";import { table, text, uuid } from "@cosmicdrift/kumiko-framework/db";import { createEntity, createNumberField, createTextField, defineFeature, type WriteResult,} from "@cosmicdrift/kumiko-framework/engine";import { z } from "zod";
const openAccess = { access: { openToAll: { reason: "demo recipe: any signed-in user manages every order and fulfillment; there " + "is no per-user ownership in this sample", }, },} as const;
// --- Entities ------------------------------------------------------------// Both are plain event-sourced CRUD entities. The write handler rows are// written synchronously by the event-store executor — the API process// always sees them, split or not.
export const orderEntity = createEntity({ table: "read_orders", fields: { customerName: createTextField({ personal: { of: "id" }, find: "none", required: true }), amount: createNumberField({ required: true }), status: createTextField({ personal: false, reason: "is_business_data", default: "pending" }), },});
export const fulfillmentEntity = createEntity({ table: "read_fulfillments", fields: { orderKey: createTextField({ personal: false, reason: "is_system_identifier", required: true }), carrier: createTextField({ personal: false, reason: "is_business_data", required: true }), label: createTextField({ personal: false, reason: "is_business_data", required: true }), },});
// --- Worker-applied read side (the sharp edge) ----------------------------// This projection is applied by the event-dispatcher of the WORKER process// only. `runSingleInstance: false` turns the API's local dispatcher off, so// this table stays empty until a worker runs — exactly the 2026-06-11// incident class (#512).export const orderActivityTable = table("read_order_activity", { id: uuid("id").primaryKey().defaultRandom(), tenantId: uuid("tenant_id").notNull(), orderKey: text("order_key").notNull(),});
// --- The write-bridge (fw#1717) -------------------------------------------// JobContext deliberately has no write/query — a background job must not// bypass the write path (idempotency, PII, transition guards). The app// therefore wires a component with the worker's dispatcher at boot and the// job calls it. Same pattern as inbound-mail-foundation's watch-supervisor.type FulfillWrite = (args: { readonly handlerQn: string; readonly payload: Record<string, unknown>; readonly tenantId: string;}) => Promise<WriteResult>;
let fulfillWrite: FulfillWrite | undefined;
/** Called by bin/worker.ts (wireComponents) with dispatchSystemWrite. */export function setOrderFulfillWrite(fn: FulfillWrite): void { fulfillWrite = fn;}
// Job payloads arrive as `Record<string, unknown>` (JobHandlerFn) — an// unchecked `as string` cast here would let a missing/non-string field// through silently, writing `orderKey: undefined` / `label: "label-// undefined"` instead of failing loudly. This recipe is a copy-paste// template, so the cast would get copied along with it.const processOrderPayloadSchema = z.object({ customerName: z.string().min(1),});
export function createApiWorkerSplitFeature() { return defineFeature("orders", (r) => { r.crud("order", orderEntity, { write: openAccess, read: openAccess }); r.crud("fulfillment", fulfillmentEntity, { write: openAccess, read: openAccess });
// Heavy follow-up work, pinned to the worker lane: the API process only // enqueues (its BullMQ client, no consumer), the worker executes. r.job( "process-order", { trigger: { on: "orders:write:order:create" }, runIn: "worker", }, async (payload, context) => { if (fulfillWrite === undefined) { throw new Error( "process-order: no fulfill-write bridge wired — bin/worker.ts must call setOrderFulfillWrite at boot", ); } // Every event-triggered job carries the tenant of the write that // fired it (job-runner sets `_tenantId` from the triggering user) // — reuse it instead of a fixed tenant, or a job triggered by one // tenant's write silently fulfills into another tenant. const tenantId = context.triggeredBy?.tenantId; if (tenantId === undefined) { throw new Error( "process-order: job context has no triggeredBy.tenantId — expected an order.created-triggered job to always carry the triggering write's tenant", ); } const { customerName } = processOrderPayloadSchema.parse(payload); const result = await fulfillWrite({ handlerQn: "orders:write:fulfillment:create", payload: { orderKey: customerName, carrier: "DHL", label: `label-${customerName}`, }, tenantId, }); if (!result.isSuccess) { // Throwing lets BullMQ retry the job — a silently dropped // fulfillment would leave the order stuck "pending" forever. throw new Error( `process-order: fulfillment create failed for order "${customerName}": ${result.error?.code ?? "unknown"}`, ); } }, );
// Async read model, applied by the WORKER's event-dispatcher. Entity ids // are generated at write time and are not part of the job payload — the // projection keys off the natural key (customerName), the same way a // background write-back would. Note the event type: the write-handler QN // ("orders:write:order:create") triggers jobs synchronously, but the // STORED event is the entity event "order.created" — that's what the // async dispatcher routes on. r.multiStreamProjection({ name: "order-activity", table: orderActivityTable, apply: { "order.created": async (event, tx) => { const payload = event.payload as { customerName: string }; await insertOne(tx, orderActivityTable, { tenantId: event.tenantId, orderKey: payload.customerName, }); }, }, }); });}📄 On GitHub: samples/recipes/api-worker-split/src/feature.ts