Browse docs

Start here

examplesGetting started with Flowdocumentation

Design

FoundationsLanguage architecturePhilosophy

Language specification

Program checkingConcurrencyDataEffectsResults, Tool problems, and faultsGrammarHistoryModules and importsLanguage specificationEvaluationStandard libraryToolsTypes

Runtime

Runtime architectureThe host boundaryDiagnosticsRunning a programThe history format

Guides

Writing programs that reach checkpointsImplementing Tools with a toolkitLoops that never returnRecursive delegationSharing types between Tool modulesHarnesses over tool registries

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_limit returns results in input order, whatever order the calls finish in. task.map_all is the same with no limit, and list.map runs the calls one after another (concurrency.md).
  • Passing summarize_one, a !tool function, makes the map_limit call !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.first returns 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, so main owns them, and its tail call to collect keeps them alive (concurrency.md).
  • Each task is parked on its only Tool call, wrapped in a constructor, so collect can 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.sleep is a Tool the host implements, and clock.timeout is ordinary Flow that races the work against it and stops the loser (concurrency.md).
  • ClockFailed keeps 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.stop returns Finished(value) when the task had already returned, so a finished result is handed back rather than thrown away, and Stopped otherwise. 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 match must cover all of them. The implementation shows the model the Action schema; an answer outside it is BadReply, 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.Value arguments 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 of Triage with the request, the reply is decoded into a Triage or becomes BadReply, and history records the concrete type with the call (tools.md).
  • A missing duplicate_of decodes as None; only an Option field may be missing. Decoding never guesses, fills a default or repairs a value (abi.md).
  • T can't be an opaque type: 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#1 and root/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 from main, 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).
  • guard bounds 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 step is 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 holds step and 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 →
  • NotRun claims the call did not happen, so it is safe to report as not submitted. Unknown makes no such claim, and BadReply says 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 Unknown by 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 jobs service offers a lookup by the caller's id; a real one might also wait between lookups with clock.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.encode builds the input.
  • An inner fault is a value. It comes back as Faulted and 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 on person.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.