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:
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;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;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;from typing import Literal
from pydantic import BaseModel
from kindgi import tool
class RiskInput(BaseModel): total: float
class Risk(BaseModel): level: Literal["low", "high"]
@tool(id="acme.score-risk", mutating=False)def score_risk(input: RiskInput) -> Risk: """Scores the fraud risk of an order.""" return Risk(level="high" if input.total > 1000 else "low")from pydantic import BaseModel
from kindgi import tool
class Item(BaseModel): sku: str quantity: int
class ShippingInput(BaseModel): items: list[Item]
class Shipping(BaseModel): days: int
@tool(id="acme.quote-shipping", mutating=False)def quote_shipping(input: ShippingInput) -> Shipping: """Quotes how many days an order takes to ship.""" return Shipping(days=2 + len(input.items))from kindgi import Flow
assess_order = Flow( 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"}, } },)- Each branch has a
branchId, the tool it runs (handler), and theoutputSchemaits 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. convergencedecides when the step is done and what it returns (below).
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
Section titled “Convergence”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 validationFanout or loop
Section titled “Fanout or loop”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.