Skip to content

Webhook Step

Workflow / pipeline step that calls an HTTP webhook.

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

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.

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.

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

Terminal window
cd samples/recipes/webhook-step
bun test

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