Skip to content

Pipeline Basics — Inventory-Management

End-to-end showcase of the M.1 step-vocabulary: every Tier-1 step exercised in a single inventory-management feature, with a real Postgres + Redis integration test that proves the pipeline form behaves end-to-end against the worktree’s framework source.

StepWhere it shows up
r.step.aggregate.createproduct:create
r.step.aggregate.updateproduct:rename, product:adjust-stock, product:bulk-adjust
r.step.aggregate.appendEventproduct:adjust-stock, product:archive, report:archive-low-stock-products
r.step.read.findOneproduct:rename, product:adjust-stock, product:bulk-adjust
r.step.read.findManyreport:archive-low-stock-products
r.step.computeproduct:adjust-stock, product:bulk-adjust
r.step.branchproduct:rename, product:adjust-stock
r.step.forEachproduct:bulk-adjust, report:archive-low-stock-products
r.step.returnevery handler
r.step.unsafeProjectionUpsertproduct:adjust-stock
r.step.unsafeProjectionDeleteproduct:adjust-stock, product:archive, report:archive-low-stock-products

A product has an SKU, a name, and a current-stock counter. Operators adjust stock individually (product:adjust-stock) or in bulk (product:bulk-adjust). Whenever stock crosses below a threshold (LOW_STOCK_THRESHOLD = 10), an alert row is upserted into a custom projection (read_inventory_low_stock_alerts); when it returns above, the alert is deleted — both inline in the same TX as the aggregate write.

Open src/feature.ts and read top-to-bottom. Each handler is preceded by a short comment block that names the steps it exercises and why the combination is the production-realistic shape (not just an API tour).

If you’ve only ever seen the free-form handler signature (r.writeHandler(name, schema, async fn, opts)), compare:

  • samples/recipes/custom-handlers/src/feature.ts — a similar domain, written in the free-form style.
  • this file — the same kind of operations expressed as defineWriteHandler({ perform: pipeline(...) }).

Two patterns the pipeline form catches that the free-form does not:

  1. Skip-if-noop: product:rename runs read.findOne followed by branch — the aggregate.update step only fires when the name actually changed. In free-form code this kind of check tends to be inlined inconsistently per handler; the pipeline form makes it a visible step.
  2. Read-then-loop: report:archive-low-stock-products calls read.findMany and pipes the result into forEach whose body appends an event + deletes the projection-row per item. The sub-pipeline boundary is explicit, the boot-validator walks into it, and the Designer (M.5) can render the loop body as a nested step group.
Terminal window
# From THIS sample directory
cd samples/recipes/pipeline-basics
bun test

The test relies on Postgres + Redis from docker compose up (framework dev stack — not the published image).

This sample may need tsconfig paths so @cosmicdrift/kumiko-framework resolves to packages/framework/src when testing unreleased engine APIs. The repo-wide integration test config excludes this sample until the engine APIs land in a published release.

  • Pure-ES vs CRUD-update: product:adjust-stock is written in the pure event-sourcing shape — it appends a domain event (product-stock-adjusted) and lets the inline product-stock-counter projection update currentStock. The alternative (calling aggregate.update and appendEvent in one handler) inflates the stream-version mid-handler and trips optimistic-locking on the next call. The other handlers (rename, bulk-adjust) use aggregate.update directly because they don’t combine it with appendEvent.
  • unsafeProjection.* allowlist: every projection-table written via unsafeProjectionUpsert/Delete must be declared via r.requires.projection("table_name") in the same feature. The boot-validator enforces this and rejects writes against aggregate-tables (registered via r.entity) — those have to go through r.step.aggregate.*.
  • PipelineCtx import: one resolver in product:adjust-stock annotates its argument as PipelineCtx because payload: StepResolver<unknown> can’t infer the destructure. This is a known M.1 DX gap (Followup #4 — TData-inference); the rest of this sample’s as-cast lastiness will get cleaner once that followup is closed.
  • Worktree-local test setup: tsconfig paths may alias @cosmicdrift/kumiko-framework to the worktree source when testing unreleased engine APIs. Once published, the alias can go away.

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):

// Pipeline Basics — Inventory-Management Showcase
//
// Production-pattern sample that exercises every Tier-1 step from
// the M.1 step-vocabulary in a single feature. The handlers below
// are the canonical reference for "what the pipeline form looks
// like in real code"; the M.3 codemod (later) will translate the
// free-form `custom-handlers` sample into roughly this shape.
//
// Steps covered:
// - r.step.aggregate.create / aggregate.update / aggregate.appendEvent
// - r.step.read.findOne / read.findMany
// - r.step.compute / branch / forEach / return
// - r.step.unsafeProjectionUpsert / unsafeProjectionDelete
//
// Domain: products with a current-stock counter. Stock-adjustments
// are tracked as domain events on the product stream (in addition
// to the auto-CRUD updated-event). A custom non-aggregate
// projection (`low_stock_alerts`) gets upserted/deleted inline
// whenever stock crosses a threshold — exactly the use-case the
// `unsafe`-prefix is designed for: framework-author opting into
// raw projection-writes after seeing the prefix at every call site.
import type { EntityTableMeta } from "@cosmicdrift/kumiko-framework/db";
import {
buildEntityTable,
createEventStoreExecutor,
integer,
selectMany,
table,
text,
updateMany,
uuid,
} from "@cosmicdrift/kumiko-framework/db";
import type { PipelineCtx } from "@cosmicdrift/kumiko-framework/engine";
import {
createEntity,
createNumberField,
createTextField,
defineFeature,
defineWriteHandler,
pipeline,
} from "@cosmicdrift/kumiko-framework/engine";
import { z } from "zod";
const LOW_STOCK_THRESHOLD = 10;
export const productEntity = createEntity({
table: "read_inventory_products",
fields: {
sku: createTextField({ required: true }),
name: createTextField({ required: true }),
currentStock: createNumberField({ default: 0 }),
},
});
// Custom non-aggregate projection. The boot-validator allows direct
// upsert/delete only because `r.requires.projection(...)` is declared
// inside the feature below — without that line, boot fails fast.
export const lowStockAlertsTable: EntityTableMeta = table("read_inventory_low_stock_alerts", {
productId: uuid("product_id").primaryKey(),
sku: text("sku").notNull(),
currentStock: integer("current_stock").notNull(),
threshold: integer("threshold").notNull(),
});
// Module-level drizzle-table + executor — the steps below capture
// `productExecutor` from this scope, and the integration test
// imports `productTable` for raw selects.
export const productTable = buildEntityTable("product", productEntity);
const productExecutor = createEventStoreExecutor(productTable, productEntity, {
entityName: "product",
});
export const inventoryFeature = defineFeature("inventory", (r) => {
r.entity("product", productEntity);
r.requires.projection("read_inventory_low_stock_alerts");
// Domain-events on the product stream (alongside the auto-CRUD
// events). r.defineEvent registers them globally — the `type`
// field of aggregate.appendEvent is a plain string, but using
// `<eventDef>.name` here gives type-safe aliasing if the event
// is later renamed.
const stockAdjusted = r.defineEvent(
"product-stock-adjusted",
z.object({
delta: z.number().int(),
reason: z.string(),
newStock: z.number().int(),
}),
);
const archived = r.defineEvent("product-archived", z.object({ reason: z.string() }));
// Inline projection that maintains `currentStock` from the
// stock-adjusted domain event. Pure event-sourcing pattern: the
// adjust-stock handler does NOT call aggregate.update — it only
// appends the domain event, and this projection (running in the
// same TX) replays its effect onto the product row.
//
// Why pure-ES (rather than aggregate.update + appendEvent in one
// handler): combining the two inflates the stream-version mid-
// handler, which would trip optimistic-locking on a back-to-back
// adjust-stock call. The split-here is the canonical Marten/CQRS
// shape — domain events are the source of truth, projections are
// the queryable view.
r.projection({
name: "product-stock-counter",
source: "product",
table: productTable,
apply: {
[stockAdjusted.name]: async (event, tx, table) => {
const p = event.payload as { newStock: number };
await updateMany(tx, table, { currentStock: p.newStock }, { id: event.aggregateId });
},
},
});
// -------------------------------------------------------------
// 1. inventory:product:create — straight aggregate.create.
// Simplest pipeline form: one mutating step + return.
// -------------------------------------------------------------
r.writeHandler(
defineWriteHandler({
name: "product:create",
schema: z.object({
sku: z.string().min(1),
name: z.string().min(1),
initialStock: z.number().int().min(0).default(0),
}),
access: { roles: ["Admin"] },
perform: pipeline<{ sku: string; name: string; initialStock: number }, { id: string }>(
({ event, r }) => [
r.step.aggregate.create("product", {
executor: productExecutor,
data: () => ({
sku: event.payload.sku,
name: event.payload.name,
currentStock: event.payload.initialStock,
}),
}),
r.step.return(({ steps }) => ({
isSuccess: true as const,
data: { id: (steps["product"] as { id: string }).id },
})),
],
),
}),
);
// -------------------------------------------------------------
// 2. inventory:product:rename — read + branch (skip-if-noop).
// Demonstrates the read-then-conditional-write pattern.
// The aggregate.update inside `onTrue` only runs when the
// name actually changed — saves an event + a projection
// write per redundant rename.
// -------------------------------------------------------------
r.writeHandler(
defineWriteHandler({
name: "product:rename",
schema: z.object({ id: z.uuid(), name: z.string().min(1) }),
access: { roles: ["Admin"] },
perform: pipeline<{ id: string; name: string }, { id: string; renamed: boolean }>(
({ event, r }) => [
r.step.read.findOne("current", {
table: productTable,
where: () => ({ id: event.payload.id }),
}),
r.step.branch({
if: ({ steps }) => {
const cur = steps["current"] as { name: string } | null;
return cur !== null && cur.name !== event.payload.name;
},
onTrue: [
r.step.aggregate.update("product", {
executor: productExecutor,
id: () => event.payload.id,
version: ({ steps }) => (steps["current"] as { version?: number } | null)?.version,
changes: () => ({ name: event.payload.name }),
}),
],
}),
r.step.return(({ steps }) => ({
isSuccess: true as const,
data: {
id: event.payload.id,
// Truthy iff the aggregate.update step actually ran —
// the executor lands its SaveContext under steps.product.
renamed: steps["product"] !== undefined,
},
})),
],
),
}),
);
// -------------------------------------------------------------
// 3. inventory:product:adjust-stock — every Tier-1 step in one
// handler. Production-realistic: read aggregate, compute new
// stock, persist update, append a domain-event for the
// audit-trail, then upsert OR delete the low-stock-alert
// based on whether the new stock crosses the threshold.
// -------------------------------------------------------------
r.writeHandler(
defineWriteHandler({
name: "product:adjust-stock",
schema: z.object({
id: z.uuid(),
delta: z.number().int(),
reason: z.string().min(1),
}),
access: { roles: ["Admin", "User"] },
perform: pipeline<
{ id: string; delta: number; reason: string },
{ id: string; newStock: number }
>(({ event, r }) => [
r.step.read.findOne("current", {
table: productTable,
where: () => ({ id: event.payload.id }),
}),
r.step.compute("newStock", ({ steps }) => {
const cur = steps["current"] as { currentStock: number } | null;
if (!cur) {
throw new Error(`product not found: ${event.payload.id}`);
}
return cur.currentStock + event.payload.delta;
}),
// Pure-ES: emit the domain event; the inline `product-stock-counter`
// projection above replays it onto the product row in the same TX.
// No aggregate.update here — see the projection's comment for why
// the split is intentional.
r.step.aggregate.appendEvent({
aggregateId: () => event.payload.id,
aggregateType: "product",
type: stockAdjusted.name,
// Explicit ctx-type: payload's resolver is StepResolver<unknown>,
// which prevents TS from inferring `ctx.steps` shape via the
// destructure. Followup #4 (TData-Inference) tracks the DX-fix.
payload: ({ steps }: PipelineCtx) => ({
delta: event.payload.delta,
reason: event.payload.reason,
newStock: steps["newStock"] as number,
}),
}),
// Inline projection-maintenance: cross the threshold in
// either direction in the same TX as the aggregate write.
// Async multi-stream projections would lag here — this is
// the use-case for inline unsafeProjection.*.
r.step.branch({
if: ({ steps }) => (steps["newStock"] as number) < LOW_STOCK_THRESHOLD,
onTrue: [
r.step.unsafeProjectionUpsert({
table: lowStockAlertsTable,
on: ["productId"],
row: ({ steps }) => {
const cur = steps["current"] as { sku: string };
return {
productId: event.payload.id,
sku: cur.sku,
currentStock: steps["newStock"] as number,
threshold: LOW_STOCK_THRESHOLD,
};
},
}),
],
onFalse: [
r.step.unsafeProjectionDelete({
table: lowStockAlertsTable,
where: () => ({ productId: event.payload.id }),
}),
],
}),
r.step.return(({ steps }) => ({
isSuccess: true as const,
data: {
id: event.payload.id,
newStock: steps["newStock"] as number,
},
})),
]),
}),
);
// -------------------------------------------------------------
// 4. inventory:product:bulk-adjust — forEach over a list of
// adjustments. Each iteration is its own read+update mini-
// pipeline; `scope.adj` is the per-iteration item.
// Sequential by design (concurrency=1 in M.1.6).
// -------------------------------------------------------------
r.writeHandler(
defineWriteHandler({
name: "product:bulk-adjust",
schema: z.object({
adjustments: z.array(z.object({ id: z.uuid(), delta: z.number().int() })).min(1),
}),
access: { roles: ["Admin"] },
perform: pipeline<
{ adjustments: ReadonlyArray<{ id: string; delta: number }> },
{ processed: number }
>(({ event, r }) => [
r.step.forEach({
over: () => event.payload.adjustments,
as: "adj",
do: [
r.step.read.findOne("current", {
table: productTable,
where: ({ scope }) => ({ id: (scope["adj"] as { id: string }).id }),
}),
r.step.compute("newStock", ({ steps, scope }) => {
const cur = steps["current"] as { currentStock: number } | null;
const adj = scope["adj"] as { delta: number };
if (!cur) throw new Error("product not found in bulk-adjust");
return cur.currentStock + adj.delta;
}),
r.step.aggregate.update("product", {
executor: productExecutor,
id: ({ scope }) => (scope["adj"] as { id: string }).id,
version: ({ steps }) => (steps["current"] as { version?: number } | null)?.version,
changes: ({ steps }) => ({ currentStock: steps["newStock"] as number }),
// Pure update path inside forEach — no appendEvent
// alongside, so optimistic-locking stays meaningful.
}),
],
}),
r.step.return(() => ({
isSuccess: true as const,
data: { processed: event.payload.adjustments.length },
})),
]),
}),
);
// -------------------------------------------------------------
// 5. inventory:product:archive — domain-event + projection
// cleanup. The aggregate is preserved (event-sourced!), the
// side-projection row gets purged in the same TX. Reads
// against archived products still work via loadAggregate;
// list-views skip them by filtering on the archived event.
// -------------------------------------------------------------
r.writeHandler(
defineWriteHandler({
name: "product:archive",
schema: z.object({ id: z.uuid(), reason: z.string().min(1) }),
access: { roles: ["Admin"] },
perform: pipeline<{ id: string; reason: string }, { id: string }>(({ event, r }) => [
r.step.aggregate.appendEvent({
aggregateId: () => event.payload.id,
aggregateType: "product",
type: archived.name,
payload: () => ({ reason: event.payload.reason }),
}),
r.step.unsafeProjectionDelete({
table: lowStockAlertsTable,
where: () => ({ productId: event.payload.id }),
}),
r.step.return(() => ({
isSuccess: true as const,
data: { id: event.payload.id },
})),
]),
}),
);
// -------------------------------------------------------------
// 6. inventory:report:bulk-archive-stale — forEach + read.findMany.
// Demonstrates the read-list-then-loop pattern: pull a slice
// of low-stock alerts via read.findMany, archive each. Used
// by the integration test to prove findMany lands rows under
// steps.<name> as an array.
// -------------------------------------------------------------
r.writeHandler(
defineWriteHandler({
name: "report:archive-low-stock-products",
schema: z.object({ reason: z.string().min(1) }),
access: { roles: ["Admin"] },
perform: pipeline<{ reason: string }, { archivedCount: number }>(({ event, r }) => [
r.step.read.findMany("alerts", {
table: lowStockAlertsTable,
}),
r.step.forEach<{ productId: string }>({
over: ({ steps }) => (steps["alerts"] as ReadonlyArray<{ productId: string }>) ?? [],
as: "alert",
do: [
r.step.aggregate.appendEvent({
aggregateId: ({ scope }) => (scope["alert"] as { productId: string }).productId,
aggregateType: "product",
type: archived.name,
payload: () => ({ reason: event.payload.reason }),
}),
r.step.unsafeProjectionDelete({
table: lowStockAlertsTable,
where: ({ scope }) => ({
productId: (scope["alert"] as { productId: string }).productId,
}),
}),
],
}),
r.step.return(({ steps }) => ({
isSuccess: true as const,
data: {
archivedCount: (steps["alerts"] as ReadonlyArray<unknown>).length,
},
})),
]),
}),
);
// -------------------------------------------------------------
// Query handlers — free-form. M.1's pipeline-engine is write-
// only; queries get the existing handler signature. The list
// below is what the integration test reads against.
// -------------------------------------------------------------
r.queryHandler(
"low-stock-alerts:list",
z.object({}),
async (_query, ctx) => {
const rows = await selectMany(ctx.db, lowStockAlertsTable);
return { rows };
},
{ access: { roles: ["Admin"] } },
);
});

📄 On GitHub: samples/recipes/pipeline-basics/src/feature.ts