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:
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;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;from pydantic import BaseModel, Field
from kindgi import tool
# The warehouse API is busy for the first two calls per order.calls: dict[str, int] = {}
class ReserveInput(BaseModel): order_id: str = Field(alias="orderId")
class Reserved(BaseModel): order_id: str = Field(alias="orderId") reserved: bool attempt: int
@tool(id="acme.reserve-stock")def reserve_stock(input: ReserveInput) -> Reserved: """Reserves the stock for an order at the warehouse.""" attempt = calls.get(input.order_id, 0) + 1 calls[input.order_id] = attempt if attempt < 3: raise RuntimeError("Warehouse busy (503)") return Reserved(orderId=input.order_id, reserved=True, attempt=attempt)from kindgi import Flow
reserve_order = Flow( 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"}, ],)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}}, …}The policy
Section titled “The policy”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) orexponential(doubling each time: 500, 1000, 2000 ms here).exponentialneedsmaxDelayMs, 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.)
Where a policy applies
Section titled “Where a policy applies”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.