Reference

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

  • spawn - starts a task from a closure.
  • both - waits for two tasks of different types at once; Task.all does the same for any number of the same type.

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 as Fail(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

  • spawn - a task an idle worker may take before it first runs.
  • offload - a body that blocks, run on the blocking pool.

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 String built at run time, a list or a shared type object 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.

field limit

limit: Duration

How long the waiter was willing to wait.

fn show

fn show(): String

"the task did not finish within {limit}".