Skip to content

Retry a failing step

A step that calls something flaky can be retried before its failure fails the run. The retry policy goes on the edge into the step. acme.reserve-stock calls a warehouse that is busy for the first two calls per order:

tools/reserve-stock/index.ts
import { defineTool } from '@kindgi/sdk/define';
import type { ToolId } from '@kindgi/sdk/types';
import { z } from 'zod';
// The warehouse API is busy for the first two calls per order.
const calls = new Map<string, number>();
const defined = defineTool({
id: 'acme.reserve-stock' as ToolId,
description: 'Reserves the stock for an order at the warehouse.',
version: '0.1.0',
input: z.object({ orderId: z.string() }),
output: z.object({ orderId: z.string(), reserved: z.boolean(), attempt: z.number().int() }),
effects: [],
handler: async ({ orderId }) => {
const attempt = (calls.get(orderId) ?? 0) + 1;
calls.set(orderId, attempt);
if (attempt < 3) throw new Error('Warehouse busy (503)');
return { orderId, reserved: true, attempt };
},
});
if (defined.kind === 'err') throw new Error(defined.error.message);
export default defined.value;
flows/reserve-order/index.ts
import { defineFlow } from '@kindgi/sdk/define';
const defined = defineFlow({
id: 'acme.reserve-order',
version: '0.1.0',
name: 'Reserve an order',
description: 'Reserves the stock for an order, retrying while the warehouse is busy.',
nodes: [{ id: 'reserve', kind: 'tool', ref: 'acme.reserve-stock' }],
edges: [
{
id: 'e1',
from: '$start',
to: 'reserve',
policy: {
retry: { maxAttempts: 4, delayMs: 500, backoff: 'exponential', maxDelayMs: 5_000 },
timeoutMs: 10_000,
},
},
{ id: 'e2', from: 'reserve', to: '$end' },
],
});
if (defined.kind === 'err') throw new Error(defined.error.message);
export default defined.value;
Terminal window
kindgi runs start --flow=acme.reserve-order --input='{"orderId":"A-200"}'
{
…
"status": "completed",
…
"output": {
"attempt": 3,
"orderId": "A-200",
"reserved": true
},
…
}

The journal has each retry, with the error that caused it and the wait before the next attempt:

{"sequence": 2, "kind": "step.started", "nodeId": "reserve", "payload": {"input": {"orderId": "A-200"}}, …}
{"sequence": 3, "kind": "step.retry-scheduled", "nodeId": "reserve", "payload": {"nodeId": "reserve", "attempt": 1, "nextDelayMs": 500, "previousError": "handler-error: Tool \"acme.reserve-stock\" handler threw: … Warehouse busy (503)"}, …}
{"sequence": 4, "kind": "step.retry-scheduled", "nodeId": "reserve", "payload": {"nodeId": "reserve", "attempt": 2, "nextDelayMs": 1000, "previousError": "handler-error: Tool \"acme.reserve-stock\" handler threw: … Warehouse busy (503)"}, …}
{"sequence": 5, "kind": "step.completed", "nodeId": "reserve", "payload": {"output": {"attempt": 3, "orderId": "A-200", "reserved": true}}, …}

policy.retry:

  • maxAttempts: attempts in all, the first one included (at most 10).
  • delayMs: the wait before the first retry (0 by default).
  • backoff: fixed (the default: the same wait each time), linear (delayMs × the retry's number) or exponential (doubling each time: 500, 1000, 2000 ms here). exponential needs maxDelayMs, the longest wait.

When the last attempt fails, the step fails and so does the run, with the last error.

policy.timeoutMs limits each attempt. An attempt that takes longer is stopped and counts as a failed attempt; without retries, the step fails with reason: "timeout":

{"sequence": 3, "kind": "step.failed", "nodeId": "reserve", "payload": {"reason": "timeout", "limitMs": 5, "message": "Handler for node \"reserve\" exceeded timeout of 5ms", "elapsedMs": 68}, …}

(That was timeoutMs: 5, to show it.)

A policy applies to the step its edge leads to, when that step has exactly one incoming edge. A step that joins branches (several incoming edges) ignores the policies on them. The other policy fields, concurrencyKey and priority, are in the flow schema.