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: stepsPipeline(...) }).

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,
stepsPipeline,
} 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(),
}),
{ piiFields: "none" },
);
const archived = r.defineEvent("product-archived", z.object({ reason: z.string() }), {
piiFields: "none",
});
// 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: stepsPipeline<{ 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: stepsPipeline<{ 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: stepsPipeline<
{ 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: stepsPipeline<
{ 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: stepsPipeline<{ 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: stepsPipeline<{ 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