Webhook Step
Workflow / pipeline step that calls an HTTP webhook.
What it shows
Section titled “What it shows”- Webhook step registration in a workflow
- Failure surface for outbound HTTP (one delivery attempt per request)
- Authenticated webhooks via a tenant-owned secret (
incident:open-authenticated)
Auth secret
Section titled “Auth secret”r.step.webhook.send’s auth.secret is a name inside the tenant-owned
step-dispatcher:webhook-auth.<name> namespace of the secrets feature —
never a process-wide env var. The secrets feature must be mounted
(createSecretsFeature() + a MasterKeyProvider) alongside
step-dispatcher. Set the secret per tenant before dispatching:
await stack.http.writeOk( "secrets:write:set", { key: "step-dispatcher:webhook-auth.incident-hook", value: "<token>" }, tenantAdmin,);Each tenant’s secret is isolated — a tenant with no matching secret (or one
stored without the step-dispatcher:webhook-auth. prefix) gets a generic
webhook auth secret is not available step.dispatch-failed event, never
another tenant’s credential.
Never combine auth.secret with a caller-controlled url: the caller would
receive the secret as the Authorization header. incident:open-authenticated
therefore posts to a fixed URL.
Idempotency-Key
Section titled “Idempotency-Key”Every webhook request carries Idempotency-Key: <dispatch stream id>. The
value is the same when the same dispatch request is delivered again, so a
receiver can deduplicate redelivered calls. An Idempotency-Key set explicitly in
headers (any casing) wins.
Source
Section titled “Source”Feature entry point: src/feature.ts.
Needs a running Postgres and TEST_DATABASE_URL set (e.g. postgres://postgres:[email protected]:5432/postgres, see demo/.env.example).
cd samples/recipes/webhook-stepbun testSource 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):
// webhook-step Sample — Tier-2 r.step.webhook.send showcase.//// A minimal "incident-open" handler that opens an incident aggregate and// dispatches a Zapier-style webhook in deferred mode. The webhook fires// only after the TX commits — if the aggregate.create rolls back, the// webhook does NOT go out (the step.dispatch-requested event vanishes// with the rollback).
import { buildEntityTable, createEventStoreExecutor } from "@cosmicdrift/kumiko-framework/db";import { createEntity, createTextField, defineFeature, defineWriteHandler, type PipelineCtx, stepsPipeline,} from "@cosmicdrift/kumiko-framework/engine";import { z } from "zod";
export const incidentEntity = createEntity({ table: "read_webhook_demo_incidents", fields: { title: createTextField({ personal: false, reason: "is_business_data", required: true }), severity: createTextField({ personal: false, reason: "is_business_data", required: true }), },});
export const incidentTable = buildEntityTable("incident", incidentEntity);const incidentExecutor = createEventStoreExecutor(incidentTable, incidentEntity, { entityName: "incident",});
const INCIDENT_HOOK_URL = "https://hooks.example/incident";
export const webhookDemoFeature = defineFeature("webhook-demo", (r) => { r.entity("incident", incidentEntity); r.requires.step("webhook.send"); r.requires.step("mail.send"); r.requires.step("callFeature");
r.writeHandler( defineWriteHandler({ name: "incident:open", schema: z.object({ title: z.string().min(1), severity: z.enum(["low", "medium", "high"]), webhookUrl: z.string(), }), access: { roles: ["Admin"] }, perform: stepsPipeline< { title: string; severity: "low" | "medium" | "high"; webhookUrl: string }, { id: string } >(({ event, r }) => [ r.step.aggregate.create("incident", { executor: incidentExecutor, data: () => ({ title: event.payload.title, severity: event.payload.severity }), }), r.step.webhook.send({ url: () => event.payload.webhookUrl, mode: "deferred", body: ({ steps }: PipelineCtx) => ({ event: "incident-opened", id: (steps["incident"] as { id: string }).id, title: event.payload.title, severity: event.payload.severity, }), }), r.step.return(({ steps }) => ({ isSuccess: true as const, data: { id: (steps["incident"] as { id: string }).id }, })), ]), }), );
// auth variant — the URL is fixed server-side: a caller-supplied URL would // receive the tenant's bearer secret. The webhook secret resolves per-tenant // via the secrets feature under `step-dispatcher:webhook-auth.incident-hook`, never a // process-wide env var. See README for how to set it via // `secrets:write:set`. r.writeHandler( defineWriteHandler({ name: "incident:open-authenticated", schema: z.object({ title: z.string().min(1), severity: z.enum(["low", "medium", "high"]), }), access: { roles: ["Admin"] }, perform: stepsPipeline< { title: string; severity: "low" | "medium" | "high" }, { id: string } >(({ event, r }) => [ r.step.aggregate.create("incident", { executor: incidentExecutor, data: () => ({ title: event.payload.title, severity: event.payload.severity }), }), r.step.webhook.send({ url: () => INCIDENT_HOOK_URL, mode: "deferred", auth: { kind: "bearer", secret: "incident-hook" }, body: ({ steps }: PipelineCtx) => ({ event: "incident-opened", id: (steps["incident"] as { id: string }).id, title: event.payload.title, severity: event.payload.severity, }), }), r.step.return(({ steps }) => ({ isSuccess: true as const, data: { id: (steps["incident"] as { id: string }).id }, })), ]), }), );
// mail.send variant — same deferred shape, different stepKind. // Step-dispatcher MSP routes to performMailDispatch. r.writeHandler( defineWriteHandler({ name: "incident:notify-via-mail", schema: z.object({ to: z.string(), title: z.string(), severity: z.enum(["low", "medium", "high"]), }), access: { roles: ["Admin"] }, perform: stepsPipeline< { to: string; title: string; severity: "low" | "medium" | "high" }, { id: string } >(({ event, r }) => [ r.step.aggregate.create("incident", { executor: incidentExecutor, data: () => ({ title: event.payload.title, severity: event.payload.severity }), }), r.step.mail.send({ to: () => event.payload.to, subject: () => `Incident: ${event.payload.title}`, body: () => `Severity ${event.payload.severity}`, mode: "deferred", }), r.step.return(({ steps }) => ({ isSuccess: true as const, data: { id: (steps["incident"] as { id: string }).id }, })), ]), }), );
// callFeature variant — sync sub-command on the same feature's // incident:open. Threads the result of the callFeature step into // the wrapper's response. r.writeHandler( defineWriteHandler({ name: "incident:open-via-call", schema: z.object({ title: z.string(), severity: z.enum(["low", "medium", "high"]) }), access: { roles: ["Admin"] }, perform: stepsPipeline< { title: string; severity: "low" | "medium" | "high" }, { id: string } >(({ event, r }) => [ r.step.callFeature("inner", { handler: "webhook-demo:write:incident:open", payload: () => ({ title: event.payload.title, severity: event.payload.severity, webhookUrl: "https://hooks.example/from-callFeature", }), }), r.step.return(({ steps }) => ({ isSuccess: true as const, data: { id: (steps["inner"] as { id: string }).id }, })), ]), }), );
// Negative test handler: aggregate.create succeeds, then a compute step // throws. Proves the webhook event is rolled back too — no fetch should // ever fire because the step.dispatch-requested event vanishes with // the TX rollback. r.writeHandler( defineWriteHandler({ name: "incident:open-then-fail", schema: z.object({ title: z.string().min(1), webhookUrl: z.string() }), access: { roles: ["Admin"] }, perform: stepsPipeline<{ title: string; webhookUrl: string }, never>(({ event, r }) => [ r.step.aggregate.create("incident", { executor: incidentExecutor, data: () => ({ title: event.payload.title, severity: "high" }), }), r.step.webhook.send({ url: () => event.payload.webhookUrl, mode: "deferred", body: () => ({ event: "incident-opened-but-rolled-back" }), }), r.step.compute("explode", () => { throw new Error("rollback-test: throwing AFTER webhook.send"); }), r.step.return({ isSuccess: true as const, data: undefined as never }), ]), }), );});📄 On GitHub: samples/recipes/webhook-step/src/feature.ts