Chapter 7
Concurrency
Rill uses Go's model: many lightweight threads over rendezvous channels. Rill calls one of them a strand β a strand of a rope, of which a program is many, twisted together. (Go calls the same thing a goroutine, Erlang a process, Java a virtual thread. The concept is the same; the word here is Rill's own. Note for anyone coming from C++ networking: Boost.Asio's strand is a different idea β a serialised execution context β and the two are unrelated.)
fn worker(jobs, results) = n = recv(jobs) send(results, n * n) worker(jobs, results) fn main() = jobs = channel() results = channel() spawn worker(jobs, results) # a new strand send(jobs, 7) println(recv(results))
spawn exprruns the expression on a new strand β all of it, sospawn handle(accept(s))accepts on the new strand, and a loop written that way never waits: it spawns strand after strand, each parked on the accept, and spins. Bind first βconn = accept(s)β so the loop is the one that waits, and hand the connection to the strand.channel()creates aChan(a);send(ch, v)andrecv(ch)rendezvous β a sender blocks until a receiver arrives and vice versa.channel(n)buffersnvalues, so a sender only blocks once the buffer is full and a receiver only blocks once it is empty.mailbox()buffers without end: its buffer grows rather than making a sender wait. It is what amachinereads from, since two machines sending each other signals at once must never hold each other up β SDL's input port refuses nothing β and the price is that nothing tells a sender to slow down.selectwaits on several channels at once and runs whichever arm becomes ready first; adefaultarm turns it into a non-blocking poll, and atimeout(ms)arm puts a deadline on the wait.close(ch)says nothing more will be sent, and wakes everything parked on the channel at once;closed(ch)asks whether it has been.- The program exits when
mainreturns; strands still running are abandoned, as in Go. - A program where every strand is blocked reports a deadlock and aborts.
Strands are cheap: 100,000 of them cost about 45 MB and start in microseconds.
fn drain(a, b) = got = select x = recv(a) -> x # bind what arrives y = recv(b) -> y send(done, 1) -> 0 # a send arm is ready when it can proceed got fn poll(ch) = v = select x = recv(ch) -> x default -> -1 # nothing was ready v
Arms are tried in order, so an earlier one wins a tie. Like match, select needs an indented block, so it goes in a body, a statement or a binding rather than inside parentheses.
A deadline on a wait
A timeout(ms) arm becomes ready when nothing else has, for that long:
fn ask(ch, ms) = r = select v = recv(ch) -> "answered " + int_to_str(v) timeout(ms) -> "gave up" r
Before this a select that had to give up needed a strand and a channel to be its clock, which is the right price for a rate and the wrong one for a deadline: a server putting a limit on every request in flight would pay a strand and a channel per request, for a clock each of them looks at once. The arm parks in the poller instead, beside the channels, and costs nothing while it waits.
The shortest deadline is the one that answers, so a wait may carry both a warning and a limit. A deadline already past is ready at once, which makes timeout(0) a poll β default by another name, and the better spelling where the number is the point. Writing default and timeout in the same select is refused, because default means "do not wait" and the deadline could then never be reached. With no channel arm at all it is a sleep with a body, which falls out of the rule rather than being a special case.
examples/timeouts.rill is all five of those.
A machine: a state and a mailbox
A strand that waits for a signal, does what the signal asks in the state it is in, and waits again is what the telephone exchanges called a process and wrote in SDL: a state machine with a mailbox. In Rill it is a tail-recursive function over a channel and a state, and machine is that function written as the machine it is:
type Signal OffHook Digit(n: Int) OnHook Ring Answer Report Status(text: Str) type Line Idle Dialling(digits: List(Int)) Ringing(caller: Chan(Mail(Signal))) Talking(peer: Chan(Mail(Signal))) machine line(me: Chan(Mail(Signal)), log: Chan(Str)) state Idle OffHook -> Dialling(Nil) Ring -> Ringing(sender) Report -> reply(Status("idle")) stay state Dialling(digits) Digit(n) if n >= 0 -> Dialling(Cons(n, digits)) OnHook -> Idle after 10000 -> send(log, "gave up") Idle state Ringing(caller) Answer -> output(caller, Answer) Talking(caller) OnHook -> Idle state Talking(peer) OnHook -> output(peer, OnHook) stop fn main() = env = mailbox() a = mailbox() spawn line(a, channel(16), Idle) post(a, env, OffHook)
The states are the constructors of a type declared as any other, so their fields say their types once. The machine is a function of that name β a machine name is lowercase, as a function's is β whose first parameter is the mailbox its signals arrive on, and whose last parameter, unwritten, is the state: it is started with its initial state as one more argument, spawn line(a, log, Idle), or called in the strand at hand.
A mailbox is made with mailbox(), which grows rather than making a sender wait, so two machines sending each other signals at once never hold each other up; a channel(n) serves too, with that risk. What arrives on the mailbox is a Mail: the signal, and the mailbox it came from β so the mailbox is a Chan(Mail(Signal)), and in every arm sender is who sent the signal being answered, as SDL had it. reply(sig) sends to it, and output(to, sig) to any mailbox, both with this machine as the sender; a signal that wants an answer no longer has to carry a channel. From outside a machine, post(to, from, sig) says which mailbox the answer comes to β main above makes env for the purpose, and reads recv(env).signal β because in SDL the environment is a process too, and nothing sends to one without being somewhere a reply can go. Two machines talk when they share a signal type; a timer's signal has the machine itself as its sender. Each state names its constructor with its fields, and under it the signals it answers, as the arms of a match on the signal: patterns, guards, and a body whose value is the next state. A body may be a block, its last line the next state; stop in that place ends the strand instead, and stay is the state the machine is in. after ms -> is the arm taken when no signal has arrived for that long, and a state has at most one. A signal no arm names is consumed and the state kept β what SDL does with an unexpected signal β unless an arm _ -> says what to do with it. Every state of the type has to appear, or the machine is refused naming the one missing.
state * is the arms every state answers β a Report that says where the machine is, a Kill that stops it, an after for a machine that has heard nothing β put after each state's own, so a state that answers the same signal itself wins, and a state with an after of its own keeps it. stay at the tail of an arm is the state the machine is in, SDL's nextstate -; in a shared arm it is the only way to say it.
Two more things a state may do, as SDL's processes did. A timer is a signal the machine sends itself later: set(Retry, 200) in an arm's body starts one that delivers Retry in 200 ms, set(Tick(n), 10) one that carries a value, reset(Retry) stops it, and active(Retry) asks whether it is running. A timer is named by its signal, setting it again moves it, and it keeps running across states β that is what after cannot do β until it runs out and its signal is answered by whatever state the machine is in then, like any other signal. save Digit(_) in a state keeps a signal that arrives there for a later state rather than answering or dropping it: the saved signals wait in the order they came, and each state takes the first of them it does not save before it reads its mailbox β so a Data that arrives while Connecting is answered by On the moment the link is up. A save takes a pattern and a guard, like an arm, and no body.
A system is the block diagram. It names the routes between the kinds of machine and the environment, with the signals each carries:
system exchange env -> line: OffHook, Digit, OnHook, Answer, Report line -> env: Status line <-> line: Ring, Answer, OnHook
The compiler checks it against the machines, as SDL's tools checked a block diagram against its processes: a signal a machine sends with output or reply that no route out of it carries is refused, naming the machine and the signal; so is a signal a machine answers that no route into it brings, a route to something that is not a machine of the file or env, and a signal that does not exist. rill diagram draws it first, before the machines, as a flowchart or with --svg as boxes and arrows. Instances are not declared β they are made with spawn, as many as the program wants β so the diagram is of kinds, which is what a block diagram is.
A run can say what its machines did. With RILL_TRACE=1 in the environment, every machine writes one line an event on standard error β a signal received in a state and from whom, the next state, a signal sent and to whom, a timer set or reset, a signal saved or dropped, an after that ran out, a stop β numbered by the runtime in the order the events happened across every strand, with each machine named by the number of its mailbox: recv line#2 Idle OffHook 1. rill msc draws those lines as the message sequence chart SDL's tools drew from a simulation, in Mermaid or, with --svg, as a picture: a lifeline per machine and per mailbox of the environment, an arrow per signal, a note where a state changed. A run without the variable pays one load per event and says nothing.
A model runs on a clock of its own. With RILL_VIRTUAL_TIME=1 in the environment, time does not pass while anything runs; when every strand is waiting and the soonest thing waited for is a timer, the clock jumps to it β the clock SDL's simulators kept. An after 50 fires the moment nothing else can happen, an hour's sleep takes none, now_millis and now_nanos tell the virtual time, and the trace of a run is the same every run, which is what makes a model's behaviour something to test. Sockets and file syncs still wait in real time, with the virtual clock standing still meanwhile; one worker, since a jump made by two would be a race.
rill check says what a machine drops. For each state, the signals some other state answers and this one neither answers nor saves β the ones SDL's tools list as a state's implicit consumption β as a note with the line; and once per machine, the signals no state answers at all, which are either not inputs of this machine (a reply it sends, like Status) or a mistake in every state. A state with _ -> stay drops nothing. A build says nothing of this, since dropping is what an SDL process does by default; the notes are for reading a model, and the language server shows them as information.
priority Sig in a state answers that signal before any that came earlier, whatever its place in the mailbox β SDL's priority input, for the abort that must not wait behind the work. It takes a pattern and a guard like save, and the arm that answers the signal is written as any other.
A machine can be a procedure. Called inside another machine's transition, on the same mailbox β if confirm(me, log, Asking(name)) then β¦ β it runs in that strand with states of its own, reading the signals as they come, and done(v) at an arm's tail is its return, as stop ends a machine with nothing; the caller's transition goes on with the value. That is SDL's procedure with states. The caller's timers do not fire while it runs, and what the caller had saved stays saved.
Two of SDL's conditions. save Sig if !ready before an arm for Sig is input Sig provided ready: a signal that arrives before its condition holds waits, and is answered in the order it came once the condition does. when cond -> body is a transition taken with no signal at all β SDL's continuous signal: when the machine is in the state and nothing pending can be answered, the first when whose condition holds is taken, before the machine waits for anything. A signal already in the mailbox is answered first, as in SDL, so a crate that ships at three items ships four if the fourth had arrived. A when that leads back to the same state with its condition still true is a loop, in SDL as here.
What the parser makes of it is the function you would have written: a match on the state, and in each state a recv β or a select with a timeout where there is an after β followed by a match on the signal whose arms call the machine again with the next state. A machine with timers or saves carries them as two more parameters no program can name, a table of what is due and the list of what was kept, and its wait's deadline is the soonest of its timers and its after. It costs what that costs, which is nothing beyond the wait; rill explain shows the calls as the tail calls they are, and rill fmt writes the machine back as it was declared. examples/exchange.rill is a small exchange of them.
A machine is a diagram before it is code, and rill diagram file.rill draws the file's machines from the declaration: as Mermaid (stateDiagram-v2, one arrow per transition labelled with the signal, its guard and what goes out on the way) to paste into a page, or with --svg as the process diagram SDL drew β the state at the top, under it a column per signal with the input symbol, the outputs and tasks in order, a decision where the arm branches on an if or a match, and at the bottom the next state or the stop mark β a timer set or reset as a task with the hourglass, a save as SDL's notched symbol. -o writes it to a file. The diagram is read off the code, never kept beside it, so it cannot fall behind.
Stopping
A strand parked on recv is parked until a value arrives. So tearing down whatever was going to send one leaves the strand exactly where it is, holding its stack, for as long as the program runs β and a server that has a strand per connection has as many of those as it has had connections. close is what there was no way to say.
fn produce(ch, n) = if n <= 0 then close(ch) else send(ch, n) produce(ch, n - 1) fn drain(ch, acc) = match recv_opt(ch) None -> acc Some(v) -> drain(ch, acc + v)
Closing is a broadcast, which is the thing sending a value cannot be: it reaches everyone parked on the channel at once, however many that is, and a program tearing down does not have to know the number. Closing twice is not an error, so a teardown that runs from two directions does not have to arrange which of them gets there first. What was already buffered survives, so a producer may close the moment it has nothing more to say and the consumer still sees everything in flight.
recv_opt(ch) is recv for a channel that ends: Some(v) while there is something, None once it is closed and empty. recv itself still answers the element type, and that is a decision rather than an oversight β most channels never end, and making every one of them answer an Option would put a match in front of every receive in the language to pay for the few that do. A recv on a closed channel ends the program and says to use recv_opt, which is the same answer an index outside a buffer gets and for the same reason: there is no value to hand back that would not be a lie.
Sending into a closed channel ends the program too. The rule that avoids it is the one Go arrived at β whoever sends is whoever closes β and where a strand really is parked on a send that nobody will ever take, select is how it is told:
fn push(out, stop, n) = r = select send(out, n) -> push(out, stop, n + 1) recv(stop) -> "gave up" r
A closed channel makes its receive arm ready, so close(stop) reaches that strand whatever the other arm is doing. The arm binds nothing, and that is what makes it work: there is no value for a closed channel to produce. An arm written x = recv(stop) -> would have to produce one, and ends the program saying so.
chan_wait and chan_taken are the two builtins recv_opt is written from, in the way map_get is written from map_find and map_val_at. They are there for a receive of a shape the prelude did not think of; a program that wants recv_opt should write recv_opt.
Being told to stop
wait_signal() parks until the program is asked something, and answers which: sig_hup(), sig_int() or sig_term(), whose numbers are 1, 2 and 15 on every Unix. It is the other end of the same idea as close, and the two are written together:
fn shutdown(done) = wait_signal() close(done)
One strand waits and tells the others, so the waiter does not have to know how many others there are. Waiting costs a worker nothing β the strand is parked in the same poller a socket would have parked it in β and a program whose only strand is waiting for a signal is not a deadlock, because it is waiting on something that answers.
A signal arrives on whichever thread the operating system picked, between any two instructions, possibly while a strand holds a lock. Nothing a program has written may run there, and nothing does: the handler writes one byte to a pipe and everything else happens on the strand.
Two things follow from how signals work rather than from anything Rill decided. The handlers are installed by the first wait_signal, so a signal sent before that does what it always did β and after it, a program that catches an interrupt and ignores it cannot be stopped with one. And SIGUSR1 and SIGUSR2 are not here: they are 30 and 31 on a Mac and 10 and 12 on Linux, and a number meaning one thing where it was written and another where it runs is worse than not having it. SIGHUP is the traditional "read your configuration again" and covers what a server would have wanted them for.
Sixteen strands may wait at once, each on a pipe of its own; the seventeenth is told to have one wait and close a channel the others watch, which is what it should have been doing.
Waiting on the clock
sleep_ms(ms) parks the strand for that long and hands its worker to whoever else is runnable β the same manoeuvre a strand waiting on a socket makes, in the same poller. A thousand strands can be asleep at once and the program is using no processor at all; a program whose every strand is asleep is not a deadlock, because it is waiting on something that always answers.
There is no Ticker type, because it is three lines and a channel:
fn ticker(ch, ms) = sleep_ms(ms) send(ch, 1) ticker(ch, ms)
Give it a channel(1) and it composes with everything else through select, which is what a loop that has to answer a client and keep a rate needs.
Ask for the time with now_nanos(), which is monotonic β only differences mean anything. A fixed-rate loop should sleep the distance to its next deadline rather than a constant, because the operating system's timers have slack in them and a constant accumulates it:
fn tick(next, period) = now = now_nanos() sleep_ms(if next > now then (next - now) / 1000000 else 0) # ... a tick's work here ... tick(next + period, period)
Running on several cores
By default every strand runs on one OS thread, which is what lets reference counting be non-atomic and therefore fast. Programs that want real parallelism opt in twice β once at build time, once at run time:
rill build main.rill -o main --parallel # refcounts become MT-safe RILL_THREADS=8 ./main # run on eight workers
--parallel makes reference counting check a flag and use atomic updates when more than one worker is active; without it the runtime refuses RILL_THREADS and says so. Workers steal strands that have not started yet from each other; one that has already run stays on its worker, because its stack frames live at that worker's addresses.
Parallelism helps when strands compute independently. Work that is a chain of rendezvous β one strand handing a value to the next β has nothing to overlap and runs fastest on a single worker. A library that wants to split a big piece of work asks workers() how many there are and spawns that many strands over ranges of it; lib/num/grid.rill does this for anything over a few million elements, and never spawns on a single-worker run.
A program that spawns anything is compiled with a check at every recursive tail call, so a strand that loops forever hands the worker over now and then (Β§7 above). A for loop whose body calls nothing is bounded by its range and is left alone, which is what keeps a loop over ten million floats in vector instructions; a while is bounded by nothing but its condition and is checked whatever it calls, so a strand waiting in one for something another strand sets does not wait forever. A long for that should share its worker says so itself, between chunks of the work: yield() gives up the turn if another strand is waiting on the worker, and returns at once if none is.
What strands may share
Storage β a Buf, a Map, a StrBuf, or anything holding one β is the one kind of value two strands can disagree about: a write from one under a read or a write from another is a race, and on several workers a race is what the machine makes of it. So a --parallel build checks what strands share, before the program runs, and rill check --parallel asks the same question of a file:
- A value handed to a strand β captured by a
spawn, or sent down a channel β and not used afterwards by the side that gave it is the strand's alone. Anything goes. - One that is used afterwards is shared. Sharing to read is fine: three strands summing thirds of one buffer while the fourth reads it too. Sharing and writing β the strand writes it, or the strand reads it while the side that kept it writes β is refused, naming the value and the line.
- A function value is opaque: what it holds and writes is not in its type, so a closure shared between strands is refused unless the program says what it knows about it.
What the program knows is said with apart(x), which is x and a claim: that this value is shared on purpose and the writes to it are kept apart β each strand its own range of a buffer, one writer and readers of a word, a closure that reads what it holds. The compiler cannot check the claim; it can make sure there is one to find. spawn work(apart(f), lo, hi, done) is how lib/num/grid.rill splits an operation across the workers, and how lib/http hands every connection's strand the one handler.
fn fill(b: Buf(Int), v) = for i in 0..#b b[i] = v fn main() = b = buf_int(8) spawn fill(b, 1) # refused: `b` is written by this strand and still used here spawn fill(b, 2) println(int_to_str(b[0]))
The check follows a value through bindings, fields, matches and calls that answer with what they were given, and knows which of a function's parameters it writes, through every call. Where it cannot see β an extern, a method call, a closure, a runtime call handed storage β it assumes a write, so what it refuses may be safe, and what it accepts without a claim is: no strand writes storage another can still reach, unless apart says so somewhere in the program. A one-worker build interleaves strands only where they wait, so the same program is not checked there.