Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .agents/rules/handlers.md
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@ export const processOrder = declareWorkflow({
// context.handleSignal/handleQuery/handleUpdate — handler binding
// context.executeChildWorkflow / context.startChildWorkflow
// context.cancellableScope / context.nonCancellableScope — see below
// context.saga — steps with compensating undos, unwound LIFO

const inventory = await context.activities.validateInventory({ orderId: args.orderId });
if (inventory.isDefect()) {
Expand Down
54 changes: 54 additions & 0 deletions .changeset/workflow-saga.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
---
"@temporal-contract/worker": minor
---

`context.saga()`: steps with compensating undos, unwound LIFO — with the
machinery failures exempt by default.

`declareWorkflow` handed a workflow `context.activities` and `context.errors`
and nothing for the walk-back, so every saga wrote its own. The LIFO machinery
is now `@unthrown/saga`'s. What this adds is the decision that belongs to
Temporal rather than to a `Result` combinator: **which failures compensate.**

```ts
const fulfilled = await context
.saga()
.step(
() => context.activities.reserveStock(order),
(reservation) => context.activities.releaseStock({ id: reservation.id }),
)
.step(
() => context.activities.chargeCard(order),
(charge) => context.activities.refund({ id: charge.id }),
)
.step(() => context.activities.ship(order))
.run();
```

The undos run on a **declared contract error** — a permanent domain answer,
where what the step did before saying no is knowable. They do **not** run on an
`ActivityError`, a `ChildWorkflowError` or a defect: a step that failed
unmodelled left state nobody can see, and un-deciding what you cannot see is a
second bug. That failure propagates untouched, so `propagateActivityFailure`
still re-raises Temporal's original failure — which deletes the per-step

```ts
.with(P.tag(ACTIVITY_ERROR_TAG), P.tag(ACTIVITY_CANCELLED_ERROR_TAG), (error) => ErrAsync(error))
```

arm that had to be repeated, was easy to omit, and was invisible when omitted.

Cancellation is the one case a caller may opt back in to, with
`saga({ compensateOnCancellation: true })`. Every undo runs inside a
non-cancellable scope: a cancelled scope schedules no activity at all, so
without one that opt-in could never compensate for the failure it exists for.
A compensation that itself fails
becomes a defect carrying its own failure, which outranks the failure that
triggered the unwind — a refund that never happened is worse news than the order
that could not ship — and the remaining undos still run first.

`workflowSaga` is the same function, exported from
`@temporal-contract/worker/workflow` for a workflow that composes its steps in a
helper.

Closes #413.
55 changes: 55 additions & 0 deletions docs/reference/worker-surface.md
Original file line number Diff line number Diff line change
Expand Up @@ -318,6 +318,61 @@ to run cleanup that must not be interrupted.
In both, a **non-cancellation** throw is an unmodeled failure and rides the
defect channel, so the modeled error channel stays exactly one type.

#### `saga(options?)`

```typescript
(options?: { compensateOnCancellation?: boolean }) => WorkflowSagaBuilder<undefined, never>;
```

A sequence of steps whose compensations are unwound **LIFO** when a later step
fails. `step(run, undo?)` takes thunks — `run` receives nothing, `undo` receives
the value its own step produced — and both may answer a plain `Result` as well
as an `AsyncResult`, so an undo is written as the ordinary activity call it is.

```typescript
const fulfilled = await context
.saga()
.step(
() => context.activities.reserveStock(order),
(reservation) => context.activities.releaseStock({ id: reservation.id }),
)
.step(
() => context.activities.chargeCard(order),
(charge) => context.activities.refund({ id: charge.id }),
)
.step(() => context.activities.ship(order))
.run();
```

**Which failures compensate is the decision this makes for you.** The undos run
on a **declared contract error** — a permanent domain answer, where what the
step did before saying no is knowable. They do **not** run on an
`ActivityError`, a `ChildWorkflowError` or a defect: a step that failed
unmodelled left state nobody can see, and un-deciding what you cannot see is a
second bug. That failure propagates untouched, so
[`propagateActivityFailure`](#propagateactivityfailure-result) still re-raises
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Temporal's original failure.

Cancellation is the one case a caller may opt back in to, with
`saga({ compensateOnCancellation: true })` — for steps holding something a
cancellation has to release anyway: a seat, a reservation, a lock. Every undo
runs inside a **non-cancellable scope**, so a cancellation cannot interrupt the
walk-back it triggered — and, more to the point, a cancelled scope schedules no
activity at all, so the opt-in would otherwise be unable to compensate for the
very failure it exists for.

`run()` answers the last step's value, and the failure comes back **unchanged**,
so a caller triages exactly what it would have without the saga. A compensation
that itself **fails** becomes a defect carrying its own failure, which outranks
the failure that triggered the unwind — a refund that never happened is worse
news than the order that could not ship — and the remaining undos still run
first.

It is pure control flow — no timers, no clock, no randomness — so it replays
deterministically inside the sandbox. `workflowSaga` is the same function,
exported from `@temporal-contract/worker/workflow` for a workflow that composes
its steps in a helper.

#### `continueAsNew(...)`

```typescript
Expand Down
3 changes: 2 additions & 1 deletion packages/worker/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,8 @@
},
"dependencies": {
"@standard-schema/spec": "catalog:",
"@temporal-contract/contract": "workspace:*"
"@temporal-contract/contract": "workspace:*",
"@unthrown/saga": "catalog:"
},
"devDependencies": {
"@arethetypeswrong/cli": "catalog:",
Expand Down
74 changes: 74 additions & 0 deletions packages/worker/src/__tests__/saga.contract.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
import { defineActivity, defineContract, defineWorkflow } from "@temporal-contract/contract";
import { z } from "zod";

/**
* A three-step fulfilment whose last step fails, plus the two activities that
* take the first two back. `mode` chooses how the last step fails, which is
* the only variable the saga's policy reads.
*/
const reserve = defineActivity({
input: z.object({}),
output: z.object({ reservationId: z.string() }),
activityOptions: { retry: { maximumAttempts: 1 } },
});

const charge = defineActivity({
input: z.object({ sleepMs: z.number() }),
output: z.object({ chargeId: z.string() }),
// Temporal delivers a cancellation notification only in the response to a
// heartbeat RPC, so `cancelled` is unobservable without this.
activityOptions: { heartbeatTimeout: "2 seconds", retry: { maximumAttempts: 1 } },
});

/**
* `declared` answers the contract error the walk-back exists for; `unmodelled`
* fails as an `ActivityError`, which must leave the earlier steps standing.
*/
const ship = defineActivity({
input: z.object({ mode: z.enum(["declared", "unmodelled"]) }),
output: z.object({ shipmentId: z.string() }),
errors: {
OutOfStock: { data: z.object({ sku: z.string() }), nonRetryable: true },
},
activityOptions: { retry: { maximumAttempts: 1 } },
});

const release = defineActivity({
input: z.object({}),
output: z.object({}),
activityOptions: { retry: { maximumAttempts: 1 } },
});

const refund = defineActivity({
input: z.object({}),
output: z.object({}),
activityOptions: { retry: { maximumAttempts: 1 } },
});

const fulfil = defineWorkflow({
input: z.object({ mode: z.enum(["declared", "unmodelled"]) }),
output: z.object({ failedWith: z.string() }),
idempotency: "allow-duplicate",
activities: { ship, refund },
});

/**
* The `compensateOnCancellation` branch: step two blocks until the workflow
* is cancelled, and the undo of step one must still run — which it can only
* do from a non-cancellable scope, since a cancelled scope schedules nothing.
*/
const fulfilUntilCancelled = defineWorkflow({
input: z.object({}),
output: z.object({ failedWith: z.string() }),
idempotency: "allow-duplicate",
activities: {},
});

export const sagaContract = defineContract({
taskQueue: "saga-tests",
// `reserve`, `charge` and `release` are global: both workflows use them, and
// activities share one flat namespace at runtime, so a per-workflow copy
// would be two implementations of one name.
activities: { reserve, charge, release },
workflows: { fulfil, fulfilUntilCancelled },
});
Loading
Loading