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

Concurrency

This specification owns concurrent work in Flow: the Task<T> value, the four operations task.spawn, task.await, task.stop and task.first, who owns a task and when it ends, what stopping does, how faults in tasks surface, the states a task passes through, what the runtime may choose and must record, the library functions built over the four operations, and timeouts.

It does not own the rest. semantics.md owns function calls, anonymous functions, evaluation order and the tail-call guarantee this file relies on. effects.md owns why timing is not an effect. errors.md owns faults and how a run ends. history.md owns how spawns, choices and stops are recorded and identified, and which tasks a checkpoint may hold. tools.md owns the Tool call a task makes. stdlib.md owns the task and clock modules as modules.

Concurrency adds no grammar. Every operation here is an ordinary function in std's task module, and a program branches on concurrent results with match.

The operations

use "std@v1/task"

@lang(task) pub type Task<T>
pub type Stop<T> { Finished(T), Stopped }

pub flow spawn<T>(work: () -> T) -> Task<T>                 // start work now
pub flow await<T>(child: Task<T>) -> T                      // wait for it
pub flow stop<T>(child: Task<T>) -> Stop<T>                 // Finished(value) if it had returned, else Stopped
pub flow first<T>(children: List<Task<T>>) -> T             // whichever finishes first; the others keep running
Open in playground →
  • spawn starts work as a new task and returns at once. The work is any function value with no parameters, usually an anonymous one: task.spawn(flow() = investigate(input.a)). The effect of work passes through to the spawn call (effects.md), so spawning work that makes Tool calls is !tool.
  • await waits until the task finishes and returns its value. Awaiting a finished task again returns the same value.
  • stop stops the task and says whether it had already finished: Finished(value) if it had returned, Stopped otherwise. A value already produced is handed back rather than discarded. Stopping below says what stopping does.
  • first waits until at least one of the given tasks finishes and returns that task's value. It stops nothing: the others keep running, and the program can await, pass on or stop them.

await, stop and first carry no effect of their own. Their answers depend on timing, and every timing choice is recorded, so replay is exact (effects.md).

Top-level constants can't use task operations.

Tasks are shared values

Task<T> is an ordinary immutable value. A program can bind it, capture it in an anonymous function, store it in a record, list, Dict value or another task's result, pass it and return it. Nothing requires a task to be awaited, or to be awaited only once.

A task is not data (types.md): it has no equality, ordering or text, can't be a Dict key or Set element, and never crosses a Tool boundary or appears in a run's input or output (tools.md). Two tasks that return the same value are still different work.

A task receives nothing once it has started. Everything it uses was captured when its work was created. It is not an addressable recipient, and tasks do not send each other anything; one task observes another only by awaiting, racing or stopping it.

Ownership

A task belongs to the call of the function that spawned it, its owning call. When the owning call returns, every task it still owns that is still running is stopped, exactly as Stopping describes. A function's body therefore says what can still be running after it returns.

  • A tail call continues the same call. A loop written as tail recursion keeps its tasks alive from one round to the next (semantics.md guarantees the tail call).
  • Returning a value that contains a task hands the task to the caller. The caller's call owns it from then on. This covers a task anywhere inside the returned value, and a task captured by a returned anonymous function. Values are immutable and the returned value is known exactly, so the hand-off is deterministic.
  • A task spawned inside a callback belongs to that callback's call. An anonymous function is a function, and its call owns what it spawns until it returns it.
  • A task's work that returns a value containing a task hands that task to the owning call of the task that returned it.

Results nobody read when the owning call returns are dropped. History keeps their replies.

A task spawned in a call lives until that call returns, so a task kept alive on purpose for the whole of a long-running loop is allowed. A loop that never returns never stops the losers of its races by itself: it stops them by hand with task.stop.

Waiting on different kinds of work

first takes tasks of one result type. To wait on different kinds of work, tag each task's result with a constructor when it is spawned, then match on what first returns:

type Event {
    Generated(Generated),
    Received(Result<Text, ToolProblem<ConversationError>>),
}

let generation = task.spawn(flow() = Generated(generate(history)))
let incoming = task.spawn(flow() = Received(conversation.next_message()))
match task.first([generation, incoming]) {
    Generated(generated) => ...,
    Received(Ok(message)) => ...,
    Received(Err(problem)) => ...,
}
Open in playground →

The event type doubles as a list of everything the code reacts to. task.map(t, f) re-tags a task that is already running (Library functions below).

Stopping

Stopping a task means it takes no further steps. It is not rollback, and it does not undo anything a Tool call did outside.

task.stop(child):

  1. records that the program stopped the child, and whether it had already finished;
  2. if the child had returned, returns Finished(value) and stops nothing else;
  3. otherwise tears down the child and every task it owns, requesting cancellation of each open call from the host, and returns Stopped.

Teardown follows a fixed order, never the order replies happened to arrive in; what it records, and in what order, is history.md's.

A stopped child may leave its own work unfinished, including replies it received but never looked at; none of them reaches a program value. A program that must not lose a reply keeps the receiving task alive rather than racing it and stopping the loser; the second example under Examples does that.

A host implementation may ignore a cancellation request or be unable to honor it. Flow claims nothing more about the world than history records.

Halting a whole run is not a Flow outcome (execution.md).

Serves foundations: honest uncertainty.

Faults in tasks

A fault in a task is not a value (errors.md), and it is never lost.

  • A fault reaches whoever awaits the task. await, stop or first applied to a faulted task passes the fault on to the call that applied it.
  • If nobody is awaiting it, the fault reaches the owning call at once. The owning call stops every other task it owns and faults, even when it is a loop that never returns.
  • A fault in a task whose result nobody read still faults its owner.

A fault that reaches a call ends the run as faulted, as every fault does (errors.md), after the stops this section requires.

Task states

A task is running, finished with a value, faulted, or stopped.

State await stop in first
Running waits stops it; Stopped waits
Finished with v v Finished(v) may be chosen; gives v
Faulted passes the fault on passes the fault on may be chosen; passes the fault on
Stopped faults Stopped again faults

Waiting on a stopped task, with await or first, is a bug, so it faults. stop may be called again on a task it stopped and returns the same answer.

task.first([]) faults. A program that keeps a changing set of running tasks keeps them in a Dict keyed by name and tags each result with its key, so the finished one can be removed.

What the runtime chooses

  • Every ready task eventually runs. Nothing more is promised about the order in which ready tasks run.
  • When more than one task given to first has finished, the runtime picks one and records the choice. Replay follows the record. Completion time and arrival order are not language rules; a live runtime may use them to choose, and the recorded choice is what later runs consume.
  • Whether stop found the task finished is recorded the same way, unless the stopping task already knows. A task stopped before has nothing to find. A task knows another has finished once it has seen that task's result through await, first or stop, or seen the result of a task, or the value of a map view, computed by a task that knew, or when the task that spawned it knew at the spawn.
  • Spawning, hand-offs and task identity are computed, not recorded. Identity comes from position (history.md).
  • A losing task's calls and replies are recorded like any other. Its pure work is recomputed on replay.

Ordinary computation inside tasks stays deterministic. Values are immutable and tasks share nothing that can change, so there is no data race inside a run. Flow infers no conflict between Tool calls that touch the same outside resource; the Tool or the application owns that.

Serves foundations: concurrency can select an outcome and deterministic reconstruction.

Library functions

std's task module also has these functions. detach, all, map_all and race are written in Flow over the four operations. map_limit and map are declared with @lang: map cannot be written over the four operations without starting a second task, which would have its own id and owner, and map_limit follows it.

pub flow detach<T>(_: Task<T>) -> Unit
pub flow all<T>(work: List<() -> T>) -> List<T>
pub flow map_all<A, B>(items: List<A>, work: (A) -> B) -> List<B>
@lang(task_map_limit) pub flow map_limit<A, B>(items: List<A>, limit: Int, work: (A) -> B) -> List<B>
pub flow race<T>(work: List<() -> T>) -> T
@lang(task_map) pub flow map<A, B>(child: Task<A>, f: (A) -> B) -> Task<B>
Open in playground →
  • detach lets a task run without waiting for it, until its owner returns and stops it. It is how a task that is neither awaited nor stopped is kept, since let _ cannot drop one (checking.md).
  • all starts every piece of work, waits for all of them, and returns the results in input order.
  • map_all starts one task per item and returns the results in input order.
  • map_limit does the same with at most limit tasks running at once: items start in input order, and each time one finishes the next starts, so limit run while items remain. A limit below 1 is treated as 1.
  • race starts every piece of work; the first to finish wins, and the rest are stopped before it returns.
  • map re-tags a running task without starting new work: awaiting the result awaits child and applies f. f is applied once, by the first task to wait on the result once child has finished, and every wait on the result gets that value; if that task is stopped before f returns, the next task to wait applies it.

Their effects pass through from the functions they are given, so task.race([flow() = web.search(q), flow() = docs.search(q)]) is !tool.

Timeouts

Time is outside the run, so Flow has no timeout clock of its own. Std's clock module declares the clock as Tools the host implements, and its timeout, whose signature and Timed<T> result stdlib.md declare, races work against clock.sleep. clock.timeout(limit, work) returns Done(value) when the work finishes first, TimedOut when the sleep does, and ClockFailed(problem) when the clock call itself fails, so a failed clock is never mistaken for an elapsed deadline. The loser is stopped: a clock that fails before the work finishes stops the work too, and its result is not kept. A program that would rather go on without a deadline when the clock fails races the work and clock.sleep with task.first itself and keeps waiting for the work. On replay the clock's reply comes from history, so the race is reproduced rather than measured again.

A host may impose its own deadline on a whole run. That halts the run; it is not a branch in the program.

Examples

Selective cancellation. A, B and C run; one later observation stops A and C, and B continues:

type Report {
    a: task.Stop<Finding>,
    b: Finding,
    c: task.Stop<Finding>,
}

pub flow main(input: Input) -> Report !tool = {
    let a = task.spawn(flow() = investigate(input.a))
    let b = task.spawn(flow() = investigate(input.b))
    let c = task.spawn(flow() = investigate(input.c))
    match triage.next() {
        Ok(MootAandC) => Report { a: task.stop(a), c: task.stop(c), b: task.await(b) },
        _ => Report {
            a: Finished(task.await(a)),
            b: task.await(b),
            c: Finished(task.await(c)),
        },
    }
}
Open in playground →

New input stops generation, and no message is lost. The pending read is a task passed from round to round. Tail calls keep it alive, and only generation is ever stopped:

flow converse(history: List<Turn>, incoming: task.Task<Event>) -> Ending !tool = {
    let generation = task.spawn(flow() = Generated(generate(history)))
    match task.first([generation, incoming]) {
        Generated(generated) => respond(generated, history, incoming),
        Received(Ok(message)) => {
            let _ = task.stop(generation)
            let next = task.spawn(flow() = Received(conversation.next_message()))
            converse(history ++ [User(message)], next)
        },
        Received(Err(problem)) => Lost { problem: describe(problem), history },
    }
}
Open in playground →

A person's answer against a deadline. Whichever task loses is stopped when the call returns:

flow timed_confirm(question: Text, wait: duration.Duration) -> ConfirmOutcome !tool = {
    let answer = task.spawn(flow() = Replied(human.confirm(question)))
    let timer = task.spawn(flow() = Expired(clock.sleep(wait)))
    match task.first([answer, timer]) {
        Replied(Ok(Yes)) => Confirmed,
        Replied(Ok(No)) => Declined,
        Replied(Err(problem)) => AskFailed(describe(problem)),
        Expired(Ok(())) => TimedOut,
        Expired(Err(problem)) => ClockFailed(describe(problem)),
    }
}
Open in playground →

Children started over many turns. The state holds running children as ordinary values, each tagged with its name. Each turn is a tail call, so they stay alive:

Spawn(name, brief) => {
    let child = task.spawn(flow() = Settled(name, delegate(brief, depth + 1)))
    step(State { running: state.running ++ [child], ..state })
},
Wait => match task.first(state.running) {
    Settled(name, outcome) => step(settle(state, name, outcome)),
},
Open in playground →

Tasks that wait on other tasks. Each task captures the tasks it depends on and awaits them inside its own work. A task can capture only tasks started before it, so waiting can never form a cycle:

let started = earlier ++ [Started { index: node.index, child: task.spawn(flow() = execute(node, earlier)) }]
Open in playground →

A parallel map.

let views = task.map_all(refs, flow(r) = describe(catalog.read(r)))
Open in playground →

Refusals

The compiler rejects:

  • a task operation in a top-level constant;
  • comparing, ordering or interpolating a task, or using one as a Dict key or Set element (types.md);
  • a task in a Tool argument, a Tool reply, or a run's input or output (tools.md);
  • tasks of different result types in one first call.

At run time, task.first([]) and waiting on a stopped task, with await or first, fault.

Not in the language

  • A cleanup form. The owner acquires a resource, stops the child, then releases it. A cleanup form can be added later without changing any program.
  • Streams as a mechanism. A stream is repeated ordinary calls, such as stream.next(handle).
  • Channels between tasks, and any addressable recipient inside a run.
  • A select expression, in any spelling. Tagging results and match cover it.
  • A group value passed to every helper that starts work, and tasks owned by the whole run, which would hide what is still running after a function returns.
  • Joins only. They can't stop some children and keep others, start children over many turns, or keep a read alive across a race.
  • Implicit parallelism, which would reorder Tool calls with no data dependency between them.
  • Stopping unreachable tasks automatically at a tail call, which would stop tasks kept alive on purpose.