std/task/lib
std/task/src/lib.trb
Concurrent work: Task, a computation that finishes later, and Channel, a queue between two tasks.
Asynchrony lives in the type: a function that returns Task<Value> may call await(), and its body produces the
Value - the same way a function that returns Result may use ?. Task and Channel are the two shared types
that connect concurrent work; everything else is copied when it crosses into one.
Every task is cancellable (Task.cancel), and a cancellation is passed on rather than answered: await() answers
the Value, and a task that awaits a task that ended cancelled is cancelled itself, right there - its scopes end and
whoever awaits it is cancelled in turn. So an IO line is file.addLine(text).await()?, one ? for the failure of
the work and nothing for a cancellation. Task.result is the one way to observe a cancellation as a value
(Cancelled), for a supervisor or for code that cancelled a task and wants to confirm it.
type Task
native shared type Task<Value>
A computation that finishes later, running from the moment it is created (spawn { ... }, or calling a function
that returns one).
Reach for it wherever work should happen while the caller keeps going. map and flatMap mean the same as
everywhere else and work on the finished value.
Pitfalls
- Dropping the last handle does not stop the task: a task runs until it finishes or somebody calls
Task.cancel.
Related
fn await
fn await(): Value
Waits for the value. Allowed in functions that return a Task, in closures passed to spawn and at the top level
of entry files and scripts.
A cancellation is passed on, not answered. Where the task ends cancelled, the waiter is cancelled too, at this
very point: it stops here, its scopes end in reverse order (every using is closed) exactly as when it is cancelled
at a loop turn, and whoever awaits the waiter is cancelled in turn. At the top level of an entry file that ends the
program with cancelled: the program waited for a task that was cancelled and exit code 130. A waiter that is
cancelled itself never sees an answer either: it stops where it waits.
Panics
Never where the task ended cancelled: the waiter stops before the answer, as above. The panic in the body guards
a state the wait rules out, with a task read the value of a task that was cancelled.
Related
Task.result- waits the same way and answers a cancellation asFail(Cancelled)instead of passing it on.
fn result
fn result(): Result<Value, Cancelled>
Waits like Task.await, and answers a cancellation instead of passing it on: Ok(value) where the task finished,
Fail(Cancelled) where it ended cancelled. The waiter goes on either way. Allowed where await() is.
It is for the code that has to know: a supervisor that restarts what was cancelled, a test, the task that called
Task.cancel and wants to confirm it. Everything else writes await(), and a cancellation stops it.
fn cancel
var fn cancel()
Asks the task to stop at its next suspension point or loop turn. A request and not a kill; asking twice, or asking a task that has finished, changes nothing. The tasks this one started are cancelled with it.
fn within
var fn within(limit: Duration): Task<Result<Value, TimedOut>>
Cancels the task once the limit has passed: Ok(value) where it finished in time, Fail(TimedOut) where the
limit passed first. A task that finished first is unaffected. Cancelling the answer does not cancel this task, and
a task somebody else cancelled cancels the answer, as every await() of it would.
fn map
fn map<Output>(transform: Transform<Value, Output>): Task<Output>
Changes the value once the task finishes, without waiting for it here. A cancelled task stays cancelled.
fn flatMap
fn flatMap<Output>(transform: Transform<Value, Task<Output>>): Task<Output>
Continues with another task once this one finishes, and keeps one level: the async form of Result.flatMap.
fn all
static fn all(tasks: Iterate<Task<Value>>): Task<List<Value>>
A task for every item, running at the same time. The values arrive in the order of the tasks. Where one of them is cancelled the answer cannot be built, so the others are cancelled too and this task ends as cancelled - and whoever awaits it with it.
fn spawn
native fn spawn<Value>(body: () => Value): Task<Value>
Runs the closure as a new task, in parallel. It gets copies of what it captures and cannot capture a var.
fn both
fn both<First, Second>(first: Task<First>, second: Task<Second>): Task<(First, Second)>
Waits for two tasks of different types at once: const (user, posts) = both(fetchUser(1), fetchPosts(1)).await().
Where one of them is cancelled, the other is cancelled too and this task ends as cancelled - and whoever awaits it
with it.
It is both and not all: it takes exactly two tasks, each of its own type, while Task.all waits for a list of
tasks of one type and Iterate.all asks a predicate of every item.
type Workers
native type Workers
The worker pool of this process: how many operating-system threads run its tasks. A reading namespace - the count is decided before the first line of the program runs and never changes, because every worker owns a heap.
Related
fn count
native static fn count(): Int
How many workers this process runs: TORB_WORKERS where it is set, and otherwise one per logical processor, at
most 1024. What parallel(workers:) defaults to.
fn blocking
native static fn blocking(): Int
How many threads the blocking pool has, which runs the bodies of offload: TORB_BLOCKING where it is set, and
otherwise 4, at most 1024. The threads start the first time a body moves there.
fn pause
native fn pause(): Task<Void>
Puts this task at the back of the queue and lets every other ready task run first. It is what makes a long loop fair; a loop of a task is cancellable at every turn without it.
fn offload
fn offload<Value>(body: () => Value): Task<Value>
Runs body on a thread of the blocking pool instead of on a worker, so a call that blocks - reading a file, running a
child process, a C library that waits - does not stall the other tasks of the worker that asked.
const text = offload({ File.readText("notes.txt") }).await()
Reach for it around one blocking call, not around a computation: a computation that is merely long belongs on the
workers, with spawn or parallel(), because the pool has only Workers.blocking threads.
Pitfalls
- The body moves to the pool only where the closure may cross to another thread: where it captures nothing counted,
or only literals and other values nothing else holds a count of. A closure that captures a
Stringbuilt at run time, a list or ashared typeobject runs on the worker that asked - the same answer, but it blocks that worker as a direct call would (docs/design/CONCURRENCY.md section 16, "The blocking pool, as built"). - The body cannot
await. Cancelling the task does not interrupt a body that already runs: it runs to its end on its thread, and what it answers is thrown away.
type Channel
native shared type Channel<Item>
A stream in memory, of which one holder has both ends: a queue between two tasks.
const channel = Channel<Report>(capacity: 8)
const writer = spawn { channel.sink().fill(reports).await() }
const total = channel.source().count().await()
channel.source() and channel.sink() are a Source and a Sink and can be handed out separately, so a producer
never sees the reading end and a consumer never sees the writing one. They are methods and not fields because they are
computed, and because a field would have to be passed to the constructor. They answer the concrete ChannelSource
and ChannelSink, so a call on one is a direct call; where a Source or a Sink is asked for, either of them is one.
Reading cannot fail: a closed channel is the end of the stream, not a failure, and a cancellation is passed on
rather than answered. The writing end fails with ChannelClosed once the reader is gone, which is how a producer
learns that nobody wants its items any more - the cancellation signal of Source.produce.
field capacity
capacity: Int = 0
0 hands every item over directly: add waits until somebody pulls.
fn source
fn source(): ChannelSource<Item>
Pulls items in the order they were added, and ends when the writing end is finished or closed.
fn sink
fn sink(): ChannelSink<Item>
add waits while the channel is full, end() tells the reading end that the stream ended.
type ChannelSource
shared type ChannelSource<Item> with Source<Item, Never>
The reading end of a Channel, a Source<Item, Never>: every next() waits for the next item, and close() tells
the writers to stop.
fn of
static fn of(channel: Channel<Item>): ChannelSource<Item>
The reading end of channel: what Channel.source answers.
fn next
var fn next(): Task<Result<Item?, Never>>
The task that waits for the item ends as cancelled where the stream ended, which is None here.
fn close
var fn close()
Not documented.
type ChannelSink
shared type ChannelSink<Item> with Sink<Item, ChannelClosed>
The writing end of a Channel, a Sink<Item, ChannelClosed>: add finishes once a reader took the item, end() and
close() end the stream.
fn of
static fn of(channel: Channel<Item>): ChannelSink<Item>
The writing end of channel: what Channel.sink answers.
fn add
var fn add(item: Item): Task<Result<Void, ChannelClosed>>
The task that offers the item ends as cancelled where nobody reads any more, which is ChannelClosed here.
fn end
var fn end(): Task<Result<Void, ChannelClosed>>
Not documented.
fn close
var fn close()
Not documented.
type ChannelClosed
type ChannelClosed with Show, Error
Nobody is reading any more: the reading end of the channel was closed. The only way a Channel's writing end fails,
and a value rather than a panic, because a producer with nobody left to feed should stop, not crash.
fn show
fn show(): String
"the channel is closed".
extend ChannelClosed with From<ChannelClosed>
extend ChannelClosed with From<ChannelClosed>
A producer whose own failure is ChannelClosed needs no conversion to learn that nobody reads any more.
fn from
static fn from(value: ChannelClosed): ChannelClosed
type Cancelled
type Cancelled with Show, Error
A task that stopped because somebody asked it to, as Task.result answers it. It carries no reason: who asked and why
is the canceller's knowledge, and a deadline has a failure of its own, TimedOut. await() never answers it - a
cancellation it meets is passed on to the waiter instead.
fn show
fn show(): String
"the task was cancelled".
type TimedOut
type TimedOut with Show, Error
A deadline that passed: what Task.within answers when the limit came first.