Skip to content

Run steps in parallel

A fanout step runs several tools at once, on the same input, and collects their outputs into one. acme.assess-order scores an order's risk and quotes its shipping at the same time, with two new read-only tools:

tools/score-risk/index.ts
import { defineTool } from '@kindgi/sdk/define';
import type { ToolId } from '@kindgi/sdk/types';
import { z } from 'zod';
const defined = defineTool({
id: 'acme.score-risk' as ToolId,
description: 'Scores the fraud risk of an order.',
version: '0.1.0',
input: z.object({ total: z.number() }),
output: z.object({ level: z.enum(['low', 'high']) }),
effects: [],
mutating: false,
handler: async ({ total }) => ({ level: total > 1000 ? ('high' as const) : ('low' as const) }),
});
if (defined.kind === 'err') throw new Error(defined.error.message);
export default defined.value;
tools/quote-shipping/index.ts
import { defineTool } from '@kindgi/sdk/define';
import type { ToolId } from '@kindgi/sdk/types';
import { z } from 'zod';
const defined = defineTool({
id: 'acme.quote-shipping' as ToolId,
description: 'Quotes how many days an order takes to ship.',
version: '0.1.0',
input: z.object({ items: z.array(z.object({ sku: z.string(), quantity: z.number() })) }),
output: z.object({ days: z.number().int() }),
effects: [],
mutating: false,
handler: async ({ items }) => ({ days: 2 + items.length }),
});
if (defined.kind === 'err') throw new Error(defined.error.message);
export default defined.value;
flows/assess-order/index.ts
import { defineFlow } from '@kindgi/sdk/define';
const defined = defineFlow({
id: 'acme.assess-order',
version: '0.1.0',
name: 'Assess an order',
description: 'Scores the risk and quotes shipping for an order, at the same time.',
nodes: [
{
id: 'order',
kind: 'tool',
ref: 'acme.get-order',
inputMapping: { orderId: { path: 'runInput.orderId' } },
},
{
id: 'assess',
kind: 'fanout',
convergence: 'all-succeed',
branches: [
{
branchId: 'risk',
handler: 'acme.score-risk',
outputSchema: { type: 'object', properties: { level: { type: 'string' } }, required: ['level'] },
},
{
branchId: 'shipping',
handler: 'acme.quote-shipping',
outputSchema: { type: 'object', properties: { days: { type: 'integer' } }, required: ['days'] },
},
],
},
],
edges: [
{ id: 'e1', from: '$start', to: 'order' },
{ id: 'e2', from: 'order', to: 'assess' },
{ id: 'e3', from: 'assess', to: '$end' },
],
output: {
mapping: {
risk: { path: 'nodeOutputs.assess.risk.level' },
shippingDays: { path: 'nodeOutputs.assess.shipping.days' },
},
},
});
if (defined.kind === 'err') throw new Error(defined.error.message);
export default defined.value;
  • Each branch has a branchId, the tool it runs (handler), and the outputSchema its output must match.
  • Every branch gets the fanout step's input: the output of the step before it, here the whole order. Each tool takes the fields its input declares (total, items) and drops the rest.
  • convergence decides when the step is done and what it returns (below).
Terminal window
kindgi runs start --flow=acme.assess-order --input='{"orderId":"A-200"}'
{
…
"status": "completed",
…
"output": {
"risk": "high",
"shippingDays": 4
},
…
}

In the journal, both branches are dispatched before either completes:

{"sequence": 6, "kind": "fanout.dispatched", "nodeId": "assess", "payload": {…, "handler": "acme.score-risk", "branchId": "risk", "fanoutNodeId": "assess"}, …}
{"sequence": 7, "kind": "fanout.dispatched", "nodeId": "assess", "payload": {…, "handler": "acme.quote-shipping", "branchId": "shipping", "fanoutNodeId": "assess"}, …}
{"sequence": 8, "kind": "fanout.branch-completed", "nodeId": "assess", "payload": {"output": {"days": 4}, "branchId": "shipping", "fanoutNodeId": "assess"}, …}
{"sequence": 9, "kind": "fanout.branch-completed", "nodeId": "assess", "payload": {"output": {"level": "high"}, "branchId": "risk", "fanoutNodeId": "assess"}, …}
{"sequence": 10, "kind": "fanout.converged", "nodeId": "assess", "payload": {"output": {"risk": {"level": "high"}, "shipping": {"days": 4}}, "convergence": "all-succeed", "fanoutNodeId": "assess"}, …}
convergence Done when The step's output
all-succeed every branch succeeded; the first failure fails the step and cancels the others { "<branchId>": <output>, … }
any-succeed the first branch succeeds (the others are cancelled); fails only if every branch fails { "winnerBranchId": "<branchId>", "output": <output> }
settle-all every branch finished, succeeded or failed; the step itself succeeds { "<branchId>": { "status": "succeeded", "output": … } | { "status": "failed", "error": "…" }, … }

With settle-all, a failed branch is data for the next step, not a failed run. Adding a branch that runs acme.check-stock (whose input the order doesn't fit) to the fanout above gives:

{
"risk": { "output": { "level": "high" }, "status": "succeeded" },
"stock": {
"error": "input-validation-failed: Input for tool \"acme.check-stock\" failed validation",
"status": "failed"
}
}

With all-succeed, the same branch fails the run:

[branch-failure] Fanout "assess" all-succeed mode: branch "stock" failed — input-validation-failed: Input for tool "acme.check-stock" failed validation

A fanout runs different tools on the same input. To run the same steps on each element of a list, several at a time, use a foreach loop with concurrency: Repeat steps in a loop.