Large-scale Orchestration Of Models
Loom
MapReduce rethought for pipelines whose operators are model calls.
Declarative dataflow over records, where every task runs under an explicit least-privilege envelope, the scheduler speaks the language of rate limits and dollar budgets, invalid model output escalates to stronger models automatically, and completed AI work is never paid for twice.
Classic data frameworks assume the wrong things
They assume operators are cheap, deterministic and trusted, and that throughput is bounded by cores. AI workloads violate all four. Loom makes each difference a first-class concept rather than something you paper over in user code.
Calls spend real money
A run-level dollar and token budget governor stops a run with graceful partial
results instead of a surprise invoice — and loom.Explain prices the
whole run before a single call is made.
Output can be wrong while the call succeeds
A returned 200 is not a returned answer. Validate is a semantic gate,
and a record that fails it climbs the model escalation ladder rather than landing
in your results.
The ceiling is rate limits, not cores
Per-model token-bucket admission control on both requests per minute and tokens per minute, so the scheduler plans against the limit that actually binds.
Operators run model-derived behaviour
Every task carries a serializable envelope — model binding, capability grants, secret references, egress allowlist, budget, sandbox profile. The planner assembles the minimal one; executors enforce it at the moment of use, and audit it.
A pipeline, end to end
Classify a queue of support tickets on a cheap model, escalate only the records whose output doesn't validate, keep the urgent ones, and roll them up into one briefing.
p := pipeline.New("ticket-triage")
src := p.FromRecords("tickets", tickets)
classified := src.Infer("classify", pipeline.InferSpec{
Binding: model.Binding{Tier: model.TierFast, // run cheap,
Escalation: []string{"claude-sonnet-5"}}, // escalate when output is invalid
System: "You classify support tickets.",
Prompt: "Classify this ticket: {{.subject}}",
ParseJSON: true, // parse output into the record
Validate: func(r core.Record) error { ... }, // semantic gate
})
classified.
Filter("urgent-only", func(r core.Record) (bool, error) {
b, _ := r.Data["urgent"].(bool); return b, nil
}).
ReduceAI("briefing", pipeline.ReduceAISpec{ // hierarchical tree reduce
Binding: model.Binding{Model: "claude-opus-4-8"},
Prompt: "Summarize {{.Count}} items:\n{{range .Items}}- {{.}}\n{{end}}",
FanIn: 8,
})
res, err := loom.Run(ctx, p,
loom.WithRegistry(reg),
loom.WithRunBudget(core.Budget{MaxCostUSD: 5.00}), // hard dollar cap
loom.WithStateDir("./state"), // cache = checkpoint = resume
)
fmt.Print(res.Report) // cost, tokens, retries, p95
Branching builds DAGs, and the planner fuses adjacent pure stages. Go-function stages run
in-process; give one pipeline.WithVersion("v1") to make it cacheable, and bump
the version when its behaviour changes.
What you get
The pieces that only exist because the operators are model calls — each of them documented, tested, and runnable offline.
Declarative pipelines
Map, Filter, FlatMap and Combine,
plus AI-native Infer — templated per-record inference with JSON parsing
and validation — and ReduceAI for parallel tree aggregation.
Cost before you spend it
loom.Explain projects a run without making one model call: per-stage
calls, rendered prompt sizes, prompt-cache economics, priced cost, and a ceiling that
rests on no assumption — the number to hand your budget.
Caching is checkpointing
Task results are keyed by op fingerprint plus input content, so reruns and crash recovery replay finished AI work at zero cost — across process restarts, and across the processes of a fleet.
Fleets
Any number of pipelines as one engine: one rate limiter, one budget, one cache, one pool of slots. A contended slot goes to the agent whose program has been served least, so a three-call summary overtakes a 10,000-record sweep.
Iteration, with a pluggable algorithm
BSP message passing, self-critique refinement and beam search all ship. A vertex is keyed by its state and its inbox rather than its round number, so cost per round falls as the loop converges.
Stream mode
An input that never ends: windows cut it into finite sets, watermarks say when a set is complete, and one checkpoint ties window state, source positions and sink commits together. Delivery at-least-once, spend exactly-once.
MCP tools, under the envelope
A stage declares what it may call; the planner turns that into a grant per tool, the server's host on the egress allowlist, and the digest of the descriptors it compiled against. One connection per host, shared by the whole fleet.
Worker processes
A durable queue with leases, heartbeats and fencing tokens. A worker killed mid-call loses its claim rather than the task, and at-least-once delivery still produces exactly-once work.
Watch it happen
The constellation view draws every task and executor as a star — prompts, responses, lineage, forecast against actual. Loom Studio makes the pipeline a canvas that prices itself between keystrokes.
Run it
The full test suite needs no network and no API keys. Fifteen of the twenty examples run offline against a deterministic mock provider — the same seam the Anthropic, OpenAI and local llama.cpp adapters sit behind.
go test ./... # full suite, no network or keys needed
go run ./examples/triage # a complete pipeline on a mock model, offline
# watch a run as a sky of stars: http://localhost:8077
go run ./examples/constellation
# cache = checkpoint: the second run makes zero model calls
LOOM_STATE=/tmp/loom go run ./examples/triage
LOOM_STATE=/tmp/loom go run ./examples/triage
# real models
ANTHROPIC_API_KEY=sk-... go run ./examples/anthropic-review
OPENAI_API_KEY=sk-... go run ./examples/openai-review
Documentation
The design, component by component — and the honest account of what is left out.