Flow examples
This is the pattern catalog that goes with getting-started.md. Each example is a small program for one pattern that orchestration programs keep needing. Each is written in the decided syntax and links the specifications that own its rules; an example illustrates a rule and never overrides one.
Every example imports its Tools from a module. A comment at the top restates the
declarations it relies on, so the example can be read without opening that
module. The repositories under github.com/example/ are placeholders.
| # | Pattern | Mainly |
|---|---|---|
| 1 | Fan out, a few at a time | task.map_limit |
| 2 | Take results as they arrive | task.first |
| 3 | Race a call against a deadline | clock.timeout |
| 4 | Stop some work and keep the rest | task.stop |
| 5 | A model chooses the next action | an Action enum |
| 6 | Structured output from one generic Tool | extract<T> |
| 7 | Wait for people | pauses |
| 8 | A bounded loop that checkpoints | guard, tail calls |
| 9 | An uncertain submit, settled by lookup | Unknown |
| 10 | Run code a model wrote | nested runs |
1. Fan out, a few at a time
Runnable, with tests: examples/fan-out.
Summarize many documents, at most four at once. A failure in one is a value in its slot; it stops nothing else.
// catalog.flow declares:
// pub type Document { pub id: Text, pub text: Text }
// pub type ReadError { Missing, Unavailable(Text) }
// @tool pub flow read(reference: Text) -> Result<Document, ReadError>
// model.flow declares:
// pub type ModelError { pub detail: Text }
// @tool pub flow summarize(text: Text) -> Result<Text, ModelError>
use "std@v1/task"
use "github.com/example/docs@v1/catalog"
use "github.com/example/models@v1/model"
pub type Summary {
Summarized { reference: Text, text: Text },
Unread { reference: Text, problem: ToolProblem<catalog.ReadError> },
Unsummarized { reference: Text, problem: ToolProblem<model.ModelError> },
}
pub flow main(references: List<Text>) -> List<Summary> !tool =
task.map_limit(references, 4, summarize_one)
flow summarize_one(reference: Text) -> Summary !tool = {
let document = match catalog.read(reference) {
Ok(document) => document,
Err(problem) => return Unread { reference, problem },
}
match model.summarize(document.text) {
Ok(text) => Summarized { reference, text },
Err(problem) => Unsummarized { reference, problem },
}
}Open in playground →task.map_limitreturns results in input order, whatever order the calls finish in.task.map_allis the same with no limit, andlist.mapruns the calls one after another (concurrency.md).- Passing
summarize_one, a!toolfunction, makes themap_limitcall!tool: an effect passes through the functions a call is given (effects.md). - Each task's calls are recorded under its own id,
root/1#1,root/2#1, and so on. The order the tasks ran in is not recorded, because it changes no value the program sees.
2. Take results as they arrive
Runnable, with tests: examples/first-to-arrive.
Plan several searches, run them all at once, and keep the findings in the
order they finish. The running tasks sit in a Dict keyed by index, and each
result carries its key so the finished task can be removed.
// research.flow declares:
// pub type Item { pub query: Text, pub reason: Text }
// pub type ModelError { pub detail: Text }
// @tool pub flow plan(question: Text) -> Result<List<Item>, ModelError>
// @tool pub flow search(item: Item) -> Result<Text, ModelError>
// @tool pub flow write(question: Text, found: List<Text>) -> Result<Text, ModelError>
use "std@v1/task"
use "github.com/example/research@v1/research"
type Arrival { Arrived { index: Int, reply: Result<Text, ToolProblem<research.ModelError>> } }
pub type Outcome {
Written { report: Text, missed: Int },
PlanFailed(ToolProblem<research.ModelError>),
WriteFailed(ToolProblem<research.ModelError>),
}
pub flow main(question: Text) -> Outcome !tool = {
let items = match research.plan(question) {
Ok(items) => items,
Err(problem) => return PlanFailed(problem),
}
let running = items.enumerate().fold([:], flow(started, (index, item)) =
started.insert(index, task.spawn(flow() = Arrived { index, reply: research.search(item) })))
collect(question, running, [], 0)
}
flow collect(question: Text, running: Dict<Int, task.Task<Arrival>>, found: List<Text>, missed: Int) -> Outcome !tool =
match running.values() {
[] => match research.write(question, found) {
Ok(report) => Written { report, missed },
Err(problem) => WriteFailed(problem),
},
pending => {
let Arrived { index, reply } = task.first(pending)
let rest = running.remove(index)
match reply {
Ok(text) => collect(question, rest, found ++ [text], missed),
Err(_) => collect(question, rest, found, missed + 1),
}
},
}Open in playground →task.firstreturns the value of a task that has finished and stops nothing; the others keep running (concurrency.md).- Which task finished first is not derivable, so every pick is recorded as a Choice. Replay follows the record and rebuilds the same order without calling anything (history.md).
- The tasks are spawned inside
fold's callback and handed back with the accumulator, somainowns them, and its tail call tocollectkeeps them alive (concurrency.md). - Each task is parked on its only Tool call, wrapped in a constructor, so
collectcan checkpoint while searches are still out (history.md).
3. Race a call against a deadline
Runnable, with tests: examples/deadline.
// web.flow declares:
// pub type Hit { pub url: Text, pub content: Text }
// pub type SearchError { pub detail: Text }
// @tool pub flow search(query: Text) -> Result<List<Hit>, SearchError>
use "std@v1/clock"
use "std@v1/duration"
use "github.com/example/web@v1/web"
pub type Lookup {
Found(List<web.Hit>),
SearchFailed(ToolProblem<web.SearchError>),
TooSlow,
NoClock(ToolProblem<Never>),
}
pub flow main(query: Text, limit: duration.Duration) -> Lookup !tool =
match clock.timeout(limit, flow() = web.search(query)) {
Done(Ok(hits)) => Found(hits),
Done(Err(problem)) => SearchFailed(problem),
TimedOut => TooSlow,
ClockFailed(problem) => NoClock(problem),
}Open in playground →- Flow has no clock of its own.
clock.sleepis a Tool the host implements, andclock.timeoutis ordinary Flow that races the work against it and stops the loser (concurrency.md). ClockFailedkeeps a failed clock call from being read as an elapsed deadline.- On replay the clock's reply and the race's Choice come from history, so the outcome is reproduced, not measured again.
- To race a person's answer against a timer and tell the two apart, spawn both and tag each result, as the person-and-deadline example in concurrency.md does.
4. Stop some work and keep the rest
Runnable, with tests: examples/selective-stop.
Three investigations start at once. A triage reply that arrives while they run may make A and C moot; B always finishes.
// probe.flow declares:
// pub type Finding { pub summary: Text }
// pub type ProbeError { pub detail: Text }
// pub type Verdict { DropAandC, KeepAll }
// @tool pub flow investigate(topic: Text) -> Result<Finding, ProbeError>
// @tool pub flow triage() -> Result<Verdict, ProbeError>
use "std@v1/task"
use "github.com/example/probes@v1/probe"
pub type Probe = Result<probe.Finding, ToolProblem<probe.ProbeError>>
pub type Report { pub a: task.Stop<Probe>, pub b: Probe, pub c: task.Stop<Probe> }
pub flow main(a: Text, b: Text, c: Text) -> Report !tool = {
let task_a = task.spawn(flow() = probe.investigate(a))
let task_b = task.spawn(flow() = probe.investigate(b))
let task_c = task.spawn(flow() = probe.investigate(c))
match probe.triage() {
Ok(DropAandC) => Report { a: task.stop(task_a), b: task.await(task_b), c: task.stop(task_c) },
_ => Report {
a: Finished(task.await(task_a)),
b: task.await(task_b),
c: Finished(task.await(task_c)),
},
}
}Open in playground →task.stopreturnsFinished(value)when the task had already returned, so a finished result is handed back rather than thrown away, andStoppedotherwise. Whether it had finished is recorded (concurrency.md).- Stopping requests cancellation of the task's open calls from the host and rolls nothing back. History records the stop, then how each open call was closed.
- A group join can't express this: it stops all children or none.
5. A model chooses the next action
Runnable, with tests: examples/model-chooses.
A model that decides what to do next replies like any other Tool, so its reply has a type: an enum of the actions the program allows.
// planner.flow declares:
// pub type Action { Read(Text), Search(Text), Finish(Text) }
// pub type Turn { Asked(Text), Saw(Text), Missed(Text) }
// pub type ModelError { pub detail: Text }
// @tool pub flow next(context: List<Turn>) -> Result<Action, ModelError>
// files.flow and web.flow each declare a Tool returning Result<Text, ...Error>.
use "github.com/example/models@v1/planner"
use "github.com/example/workspace@v1/files"
use "github.com/example/web@v1/web"
pub type Step {
Continue(planner.Turn),
Done(Text),
Stuck(ToolProblem<planner.ModelError>),
}
pub flow step(context: List<planner.Turn>) -> Step !tool =
match planner.next(context) {
Ok(Read(path)) => match files.read(path) {
Ok(text) => Continue(Saw(text)),
Err(problem) => Continue(Missed("read ${path}: ${problem}")),
},
Ok(Search(query)) => match web.search(query) {
Ok(text) => Continue(Saw(text)),
Err(problem) => Continue(Missed("search ${query}: ${problem}")),
},
Ok(Finish(answer)) => Done(answer),
Err(problem) => Stuck(problem),
}Open in playground →- Every action the model can request is visible in one enum, and the
matchmust cover all of them. The implementation shows the model theActionschema; an answer outside it isBadReply, not a crash (tools.md). - Flow has no call that names a Tool by text. When the host, not the program,
owns the set of tools, the program passes
json.Valuearguments to one generic Tool the host provides; tool-registries.md compares the two.
6. Structured output from one generic Tool
Runnable, with tests: examples/structured-output.
// model.flow declares:
// pub type ModelError { pub detail: Text }
// @tool pub flow extract<T>(prompt: Text) -> Result<T, ModelError>
use "github.com/example/models@v1/model"
pub type Severity { Low, Medium, High }
pub type Triage {
pub severity: Severity,
pub owner: Text,
pub summary: Text,
pub duplicate_of: Option<Text>,
}
pub type Outcome { Triaged(Triage), Unreadable(ToolProblem<model.ModelError>) }
pub flow main(report: Text) -> Outcome !tool = {
let prompt = """
Triage this bug report. Leave duplicate_of out unless you are sure.
${report}
"""
let reply: Result<Triage, ToolProblem<model.ModelError>> = model.extract(prompt)
match reply {
Ok(triage) => Triaged(triage),
Err(problem) => Unreadable(problem),
}
}Open in playground →- The expected type picks
T. The call sends the schema ofTriagewith the request, the reply is decoded into aTriageor becomesBadReply, and history records the concrete type with the call (tools.md). - A missing
duplicate_ofdecodes asNone; only anOptionfield may be missing. Decoding never guesses, fills a default or repairs a value (abi.md). Tcan't be anopaquetype: a value a module hands out only after its own check can't be forged by a decoder (tools.md).
7. Wait for people
Runnable, with tests: examples/wait-for-people.
Two reviewers read a draft at the same time. Each review is a Tool call that may take hours; while both are open, the run is paused on two calls.
// person.flow declares:
// pub type Verdict { Approve, Object(Text) }
// pub type AskError { pub detail: Text }
// @tool pub flow review(reviewer: Text, draft: Text) -> Result<Verdict, AskError>
use "std@v1/task"
use "github.com/example/people@v1/person"
pub type Signoff { Cleared, Objections(List<Text>) }
pub flow main(draft: Text) -> Result<Signoff, ToolProblem<person.AskError>> !tool = {
let legal = task.spawn(flow() = person.review("legal", draft))
let editor = task.spawn(flow() = person.review("editor", draft))
let notes = objection(task.await(legal)?) ++ objection(task.await(editor)?)
if notes == [] { Ok(Cleared) } else { Ok(Objections(notes)) }
}
flow objection(verdict: person.Verdict) -> List<Text> = match verdict {
Approve => [],
Object(note) => [note],
}Open in playground →A pause is an open call and nothing else (history.md). This run pauses with two pending calls,
root/1#1androot/2#1, and the host may resume them in either order, one reply per resume:flow resume run-41c2 --call 'root/2#1' --reply '{"Ok": "Approve"}' flow resume run-41c2 --call 'root/1#1' --reply '{"Ok": {"Object": "Cite the policy"}}'Resuming never repeats a call that already has a reply: replay serves it from history.
If the first answer is an
Err,?returns frommain, and the other review is stopped because its owning call returned. History records the stop and closes the open call with no reply.Who is asked, how, and what a reviewer may approve belong to the application that implements
person.review.
8. A bounded loop that checkpoints
Runnable, with tests: examples/bounded-loop.
A model reads catalog entries in batches until it can answer, for at most
max_rounds rounds.
// catalog.flow declares:
// pub type Document { pub id: Text, pub text: Text }
// pub type ReadError { Missing, Unavailable(Text) }
// @tool pub flow read(reference: Text) -> Result<Document, ReadError>
// model.flow declares:
// pub type ModelError { pub detail: Text }
// @tool pub flow extract<T>(prompt: Text) -> Result<T, ModelError>
use "github.com/example/docs@v1/catalog"
use "github.com/example/models@v1/model"
pub type Action { Finish(Text), Read(List<Text>) }
pub type Note { Found(catalog.Document), Unread(Text) }
pub type Completion {
Finished { answer: Text, rounds: Int },
OutOfRounds(List<Note>),
ModelFailed(ToolProblem<model.ModelError>),
}
pub flow main(goal: Text, max_rounds: Int) -> Completion !tool = step(goal, max_rounds, [], 0)
flow step(goal: Text, max_rounds: Int, notes: List<Note>, rounds: Int) -> Completion !tool = {
guard rounds < max_rounds else { return OutOfRounds(notes) }
let prompt = """
Goal: ${goal}
Notes so far: ${notes}
Answer Read with the catalog references you need next, or Finish with the answer.
"""
let decision: Result<Action, ToolProblem<model.ModelError>> = model.extract(prompt)
match decision {
Ok(Finish(answer)) => Finished { answer, rounds },
Ok(Read(references)) => step(goal, max_rounds, notes ++ references.map(read), rounds + 1),
Err(problem) => ModelFailed(problem),
}
}
flow read(reference: Text) -> Note !tool = match catalog.read(reference) {
Ok(document) => Found(document),
Err(problem) => Unread("${reference}: ${problem}"),
}Open in playground →- There is no loop statement. A loop is a tail call, and a tail call never grows the stack (semantics.md).
guardbounds the loop in the program itself, where a reader can see it. Spending limits on top of that are the host's.- Each round's call to
stepis a tail call to a top-level function from the root task, with no other live task, so the runtime may checkpoint there. A checkpoint holdsstepand its four arguments, and resume or replay can start from it instead of from the first round (history.md). - A loop that never returns, parked on a Tool call between rounds, takes the same shape; long-running-loops.md shows it.
9. An uncertain submit, settled by lookup
Runnable, with tests: examples/uncertain-submit.
Submit a job exactly once. If the submit may have happened, never submit again: look it up.
// jobs.flow declares:
// pub type Ack { pub request_id: Text, pub receipt: Text }
// pub type SubmitError { pub reason: Text }
// pub type LookupError { pub detail: Text }
// pub type Lookup { FoundAck(Ack), Pending, Absent }
// @tool pub flow submit(id: Text, job: Text) -> Result<Ack, SubmitError>
// @tool pub flow lookup(id: Text) -> Result<Lookup, LookupError>
use "github.com/example/jobs@v1/jobs"
pub type Output {
Accepted(jobs.Ack),
Rejected(jobs.SubmitError),
NotSubmitted(Text),
Inconclusive { id: Text, uncertainty: Text },
}
pub flow main(id: Text, job: Text) -> Output !tool =
match jobs.submit(id, job) {
Ok(ack) => Accepted(ack),
Err(Failed(error)) => Rejected(error),
Err(NotRun(reason)) => NotSubmitted(reason),
Err(Unknown(reason)) | Err(BadReply(reason)) => resolve(id, reason, 2),
}
flow resolve(id: Text, uncertainty: Text, tries: Int) -> Output !tool = {
guard tries > 0 else { return Inconclusive { id, uncertainty } }
match jobs.lookup(id) {
Ok(FoundAck(ack)) => Accepted(ack),
Ok(Pending) | Err(_) => resolve(id, uncertainty, tries - 1),
Ok(Absent) => Inconclusive { id, uncertainty },
}
}Open in playground →NotRunclaims the call did not happen, so it is safe to report as not submitted.Unknownmakes no such claim, andBadReplysays a reply arrived, so the job may exist either way (errors.md).- Flow never repeats a call on its own, on resume or on replay. After a crash
with the submit still out, the host records
Unknownby default, and this program handles it the same way (history.md). - The language assumes no lookup, idempotency key or compensation. This program
works because the
jobsservice offers a lookup by the caller's id; a real one might also wait between lookups withclock.sleep.
10. Run code a model wrote
A model writes a Flow program, repairs it from the checker's diagnostics, runs
it, and answers each call the inner run pauses on by asking a person. Checking,
running and resuming other Flow programs is a Tool the host provides, declared
in std's runner module
(execution.md).
// author.flow declares:
// pub type AuthorError { pub detail: Text }
// @tool pub flow write(goal: Text) -> Result<Text, AuthorError>
// @tool pub flow repair(goal: Text, source: Text, diagnostics: Text) -> Result<Text, AuthorError>
// person.flow declares:
// pub type AskError { pub detail: Text }
// @tool pub flow answer(call: runner.Pending) -> Result<Text, AskError>
use "std@v1/json"
use "std@v1/runner"
use "github.com/example/models@v1/author"
use "github.com/example/people@v1/person"
pub type Input { pub goal: Text }
pub type Ending {
Finished(Text),
InnerFaulted(Text),
GaveUp(List<runner.Diagnostic>),
NothingPending,
AuthorFailed(ToolProblem<author.AuthorError>),
CheckFailed(ToolProblem<List<runner.Diagnostic>>),
RunFailed(ToolProblem<runner.RunError>),
PersonFailed(ToolProblem<person.AskError>),
}
pub flow main(goal: Text) -> Ending !tool =
match author.write(goal) {
Ok(source) => attempt(goal, source, 0),
Err(problem) => AuthorFailed(problem),
}
flow attempt(goal: Text, source: Text, repairs: Int) -> Ending !tool =
match runner.check(source) {
Ok(_) => match runner.run(source, json.encode(Input { goal })) {
Ok(state) => drive(state),
Err(problem) => RunFailed(problem),
},
Err(Failed(diagnostics)) if repairs >= 3 => GaveUp(diagnostics),
Err(Failed(diagnostics)) => match author.repair(goal, source, "${diagnostics}") {
Ok(repaired) => attempt(goal, repaired, repairs + 1),
Err(problem) => AuthorFailed(problem),
},
Err(problem) => CheckFailed(problem),
}
flow drive(state: runner.RunState) -> Ending !tool =
match state {
Completed(output) => Finished(output),
Faulted(message) => InnerFaulted(message),
Paused { pending: [], .. } => NothingPending,
Paused { run, pending: [call, ..] } => match person.answer(call) {
Ok(reply) => match runner.resume(run, call.id, reply) {
Ok(next) => drive(next),
Err(problem) => RunFailed(problem),
},
Err(problem) => PersonFailed(problem),
},
}Open in playground →- A nested run is a run. It has its own id and its own history, linked to the call that started it. The outer history holds only the outer run's calls and replies, and replaying the outer run serves them without running the inner program again (history.md).
- The outer program can't know the inner program's types, so input, output
and replies cross as canonical JSON text, and
json.encodebuilds the input. - An inner fault is a value. It comes back as
Faultedand never faults the outer run (errors.md). - A paused inner run is a host value,
runner.Run, which the program holds and hands back but cannot build or inspect (tools.md). While the person answers, the outer run is itself paused onperson.answer. - The model's code reaches only the Tools its own imports declare, bound by the host like any others. What it may do there is the host's decision. recursive-delegation.md builds deeper delegation from the same pieces.