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 runningOpen in playground →spawnstartsworkas 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 ofworkpasses through to thespawncall (effects.md), so spawning work that makes Tool calls is!tool.awaitwaits until the task finishes and returns its value. Awaiting a finished task again returns the same value.stopstops the task and says whether it had already finished:Finished(value)if it had returned,Stoppedotherwise. A value already produced is handed back rather than discarded. Stopping below says what stopping does.firstwaits 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):
- records that the program stopped the child, and whether it had already finished;
- if the child had returned, returns
Finished(value)and stops nothing else; - 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,stoporfirstapplied 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
firsthas 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
stopfound 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 throughawait,firstorstop, or seen the result of a task, or the value of amapview, 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 →detachlets 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, sincelet _cannot drop one (checking.md).allstarts every piece of work, waits for all of them, and returns the results in input order.map_allstarts one task per item and returns the results in input order.map_limitdoes the same with at mostlimittasks running at once: items start in input order, and each time one finishes the next starts, solimitrun while items remain. Alimitbelow 1 is treated as 1.racestarts every piece of work; the first to finish wins, and the rest are stopped before it returns.mapre-tags a running task without starting new work: awaiting the result awaitschildand appliesf.fis applied once, by the first task to wait on the result oncechildhas finished, and every wait on the result gets that value; if that task is stopped beforefreturns, 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
Dictkey orSetelement (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
firstcall.
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
selectexpression, in any spelling. Tagging results andmatchcover 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.