Concurrency
Mesh uses the actor model for concurrency, inspired by Erlang/Elixir. Actors are lightweight processes that communicate by message passing. Jobs, scheduler-aware timers, and bounded channels cover common work that does not need a long-lived actor.
Autonomous clusters: This guide teaches the concurrency model. Continue with Autonomous Clusters for runtime-owned elasticity and Distributed Proof for release verification.
The Actor Model
In Mesh, actors are independent units of computation that do not share memory. Each actor:
- Has its own mailbox for receiving messages
- Communicates exclusively via message passing
- Runs concurrently with other actors
- Is isolated by default -- an unlinked actor crashing does not bring down other actors
This model avoids shared-memory data races and shared-state corruption. Application-level waiting cycles are still possible—for example, two services that synchronously call each other—so keep synchronous dependencies acyclic.
Spawning Actors
Define an actor with the actor keyword and start it with spawn:
actor greeter() do
receive do
_ -> println("actor received")
end
end
fn main() do
let pid = spawn(greeter)
send(pid, 1)
println("main done")
endThe spawn function returns a PID (process identifier) that you use to communicate with the actor. Actors run concurrently with the function that spawned them. Actor parameters become arguments to spawn:
actor counter(total :: Int) do
receive do
amount -> counter(total + amount)
end
end
fn main() do
let pid = spawn(counter, 0)
send(pid, 5)
endMesh infers Pid<T> from the actor's received message type. Sending a value of the wrong type is a compile-time error. Inside an actor, self() returns its own PID.
PIDs compare with == and !=, and show as <0.12> in an interpolation (<node.id.creation> for a process on another node). No process has PID 0: it is what a lookup that finds nothing and a remote spawn that fails return, shown as <0.0>, and a send to it goes nowhere.
Message Passing
Actors communicate by sending and receiving messages. Use send to deliver a message to an actor's mailbox, and receive to wait for the next message:
actor worker() do
receive do
_ -> println("worker done")
end
end
fn main() do
let w1 = spawn(worker)
let w2 = spawn(worker)
let w3 = spawn(worker)
send(w1, 1)
send(w2, 2)
send(w3, 3)
println("main sent all")
endKey points about message passing:
- Messages are processed one at a time from the actor's mailbox
sendreturns0when enqueued or written,1when the target is missing,2when the mailbox is full,3when the message exceeds the mailbox byte limit,4when a remote node is unavailable,5when a remote write fails, and6when a message for another node holds a function or runtime handle, which cannot leave its nodereceiveblocks until the next message arrivesreceivearms match the message likecasearms, with patterns andwhenguards, and together must cover the actor's message type- You can spawn multiple actors and send messages to each independently
send can take either argument from a pipe: pid |> send(message) passes the PID, and the slot pipe message |2> send(pid) passes the message.
Receive arms can match on the message's shape. An arm that calls the actor again continues with new arguments; an arm that does not ends the actor:
type Command do
Add(amount :: Int)
Report
Stop
end
actor tally(total :: Int) do
receive do
Add(amount) when amount > 0 -> tally(total + amount)
Add(_) -> tally(total)
Report ->
println("total: #{total}")
tally(total)
Stop -> println("stopped at #{total}")
end
end
fn main() do
let pid = spawn(tally, 0)
send(pid, Add(5))
send(pid, Add(-2))
send(pid, Report)
send(pid, Stop)
Timer.sleep(50)
endA receive can provide a timeout in milliseconds. The receive expression returns either the message arm's value or the timeout arm's value:
actor worker() do
let result = receive do
value -> value
after 1_000 -> 0
end
println("#{result}")
endActors can also perform computation before responding. Here is an actor that runs a function when it receives a message:
fn count_loop(n :: Int, target :: Int) -> Int do
if n >= target do
n
else
count_loop(n + 1, target)
end
end
actor worker() do
receive do
_ -> println("${count_loop(0, 100)}")
end
end
fn main() do
let pid = spawn(worker)
send(pid, 1)
endLinking and Monitoring
Actors can be linked so that failures propagate between them. If one linked actor crashes, the other is notified:
actor linked_worker() do
receive do
_ -> println("linked worker done")
end
end
actor linker() do
let worker = spawn(linked_worker)
link(worker)
receive do
message -> send(worker, message)
end
end
fn main() do
let linker_pid = spawn(linker)
send(linker_pid, 1)
endlink(pid)-- bidirectionally links two actors. If one ends abnormally (a panic or runtime error), the other ends too; a normal exit only removes the link. Exit signals never arrive as messages in an actor's ownreceive; only supervisors observe them.Process.monitor(pid, message)-- a one-way monitor: whenpidends, for whatever reason, the calling actor is sentmessage, one of its own messages. It is sent at once when there is no such process. Returns the monitor's reference.Process.demonitor(reference)-- removes a monitor, so its message is never sent; returns0on success or1when the reference is unknown.
actor worker() do
receive do
n -> println("working on #{n}")
end
end
actor watcher() do
let pid :: Pid<Int> = spawn(worker)
Process.monitor(pid, "worker ended")
send(pid, 1)
receive do
text -> println(text)
end
endLinking is the foundation for building fault-tolerant systems: supervisors use links to detect and restart failed actors.
An actor can also declare a termination callback for cleanup:
actor worker() do
receive do
_ -> println("work complete")
end
terminate do
println("cleaning up")
end
endSupervision
Supervisors are special actors that monitor and restart child actors when they fail. Define a supervisor with the supervisor keyword:
actor worker() do
receive do
_ -> println("worker got message")
end
end
supervisor WorkerSup do
strategy: one_for_one
max_restarts: 3
max_seconds: 5
child w1 do
start: fn -> spawn(worker) end
restart: permanent
shutdown: 5000
end
end
fn main() do
let sup = spawn(WorkerSup)
println("supervisor started")
endSupervision Strategies
| Strategy | Behavior |
|---|---|
one_for_one | Only the failed child is restarted |
one_for_all | All children are restarted when one fails |
rest_for_one | The failed child and all children started after it are restarted |
simple_one_for_one | Accepted template strategy; dynamic child start/terminate operations are not exposed by the current Mesh source API |
Child Specifications
Each child block configures how the supervisor manages that actor:
| Option | Purpose |
|---|---|
start | Function that spawns the child actor |
restart | Restart policy: permanent (always), transient (only on abnormal exit), temporary (never) |
shutdown | Positive milliseconds to wait for graceful shutdown, or brutal_kill |
Restart Limits
max_restarts-- maximum number of restarts allowed within the time windowmax_seconds-- the time window in seconds
If a child exceeds the restart limit, the supervisor itself shuts down, escalating the failure to its parent supervisor.
When omitted, the strategy defaults to one_for_one, max_restarts to 3, and max_seconds to 5.
Runtime Errors and Panics
A runtime error, such as List.get past the end of a list or head of an empty one, and a call to panic(message) end only the actor they happen in. The runtime prints Mesh panic: message to standard error, other actors keep running, and a supervised actor is restarted according to its restart policy:
actor worker() do
receive do
0 -> panic("cannot handle 0")
n -> println("handled #{n}")
end
worker()
end
fn main() do
let pid = spawn(worker)
send(pid, 1)
send(pid, 0)
Timer.sleep(100)
println("main still running")
endThis prints handled 1, then Mesh panic: cannot handle 0 on standard error, then main still running. The same failure in main itself ends the program with exit status 101.
Services (GenServer)
Services are stateful actors that follow the GenServer pattern. They provide a structured way to manage state with synchronous calls and asynchronous casts:
service Counter do
fn init(start_val :: Int) -> Int do
start_val
end
call GetCount() :: Int do |count|
(count, count)
end
call Increment(amount :: Int) :: Int do |count|
(count + amount, count + amount)
end
cast Reset() do |_count|
0
end
end
fn main() do
let pid = Counter.start(10)
let c1 = Counter.get_count(pid)
println("${c1}")
let c2 = Counter.increment(pid, 5)
println("${c2}")
Counter.reset(pid)
let c3 = Counter.get_count(pid)
println("${c3}")
endService Anatomy
init-- called when the service starts, returns the initial statecall-- synchronous request/response. The handler receives the current state and returns a tuple(new_state, reply)cast-- asynchronous fire-and-forget. The handler receives the current state and returns the new state
Starting and Calling Services
The compiler auto-generates snake_case methods from your PascalCase definitions:
service Store do
fn init(start_val :: Int) -> Int do
start_val
end
call Get() :: Int do |state|
(state, state)
end
call Set(value :: Int) :: Int do |_state|
(value, value)
end
cast Clear() do |_state|
0
end
end
fn main() do
let pid = Store.start(100)
let v1 = Store.get(pid)
println("${v1}")
let v2 = Store.set(pid, 200)
println("${v2}")
Store.clear(pid)
let v3 = Store.get(pid)
println("${v3}")
end| Definition | Generated method |
|---|---|
Store.start(100) | Starts the service with initial value |
Store.get(pid) | Calls the Get handler |
Store.set(pid, 200) | Calls the Set handler |
Store.clear(pid) | Casts the Clear handler |
Services with no init arguments use start() with no parameters:
service Accumulator do
fn init() -> Int do
0
end
call Add(n :: Int) :: Int do |state|
(state + n, state + n)
end
end
fn main() do
let pid = Accumulator.start()
Accumulator.add(pid, 1)
Accumulator.add(pid, 2)
let result = Accumulator.add(pid, 3)
println("${result}")
endArguments and replies are copied between the caller and the service, as actor messages are, so a reply stays valid after the service's state changes or the service stops.
Service calls are synchronous and wait for a reply. Casts are asynchronous. Prefer a cast, a direct actor message, or a job when the caller must remain independent of the service's response time.
A call waits for its own reply only: messages already in the caller's mailbox stay there for its next receive. A call to a service that has stopped, or that stops before it replies (a handler that panics, say), panics in the caller, as does a call to a service whose mailbox is full.
Jobs
Job runs finite work on lightweight actors and returns failures through Result.
fn main() do
let job = Job.async(fn() -> 21 * 2 end)
case Job.await_timeout(job, 1_000) do
Ok(value) -> println("#{value}")
Err(error) -> println("job failed: #{error}")
end
end| Function | Returns | Description |
|---|---|---|
Job.async(fn) | Pid<T> | Run a zero-argument function in an actor of its own |
Job.await(job) | Result<T, String> | Wait without a timeout |
Job.await_timeout(job, timeout_ms) | Result<T, String> | Wait up to the given number of milliseconds |
Job.map(values, fn) | List<Result<U, String>> | Run one job per list element and return results in input order |
Job replies go to the actor that started them, and each await selects the requested job without consuming unrelated mailbox messages; an actor's own receive never sees them. A job that fails (a panic, say) is an Err naming the failure to whoever awaits it, which goes on; a job whose caller fails is stopped with it. Await a job from its original caller. A timed-out job is not cancelled; its eventual reply remains available to a later await. Job.map exposes each completion's success or failure instead of failing the whole batch at the first error.
Timers
Timers use monotonic deadlines. Sleeping an actor yields its scheduler worker so other actors can run.
| Function | Returns | Description |
|---|---|---|
Timer.sleep(milliseconds) | Unit | Suspend the current actor until the delay expires |
Timer.send_after(pid, milliseconds, message) | Unit | Deliver a typed message after a delay |
Timer.apply_after(milliseconds, function) | Unit | Call a zero-argument function after a delay, in an actor of its own |
actor reminder() do
receive do
message -> println(message)
after 5_000 -> println("no reminder received")
end
end
fn main() do
let pid = spawn(reminder)
Timer.send_after(pid, 100, "time to stretch")
Timer.sleep(200)
endTimer.send_after delivers a plain message, which an actor's receive sees but a service's cast handler does not: a service dispatches on a tag that only its generated functions know. To reach a service after a delay, schedule a function that casts:
service Writer do
fn init(start :: Int) -> Int do
start
end
cast Flush(reason :: String) do |flushed|
println("flush: #{reason}")
flushed + 1
end
end
fn main() do
let writer = Writer.start(0)
Timer.apply_after(100, fn () -> Writer.flush(writer, "timer") end)
Timer.sleep(200)
endThe function runs in its own actor, which is not linked to the caller: a failing callback does not take the caller down. Values it captures stay alive until it has run. As with send_after, the program does not wait for a pending timer once main has returned and every remaining actor is idle.
Bounded Channels
Channel provides bounded, in-process queues. A Channel<T> carries values of one type T, of any type; creation and queue operations return Result, and producers use try_send and never wait for space.
fn round_trip(value :: Int) -> Int!String do
let channel = Channel.bounded_bytes(128, 1_024, :reject_newest)?
Channel.try_send(channel, value)?
let timeout = Duration.millis(10)?
Channel.recv(channel, timeout)
end
fn main() do
case round_trip(42) do
Ok(value) -> println("#{value}")
Err(error) -> println(error)
end
end| Function | Returns | Description |
|---|---|---|
Channel.bounded(item_capacity, policy) | Result<Channel<T>, String> | Create a queue bounded by item count |
Channel.bounded_bytes(item_capacity, byte_capacity, policy) | Result<Channel<T>, String> | Apply both item and byte bounds |
Channel.try_send(channel, value) | Result<Int, String> | Enqueue a value immediately or report backpressure |
Channel.recv(channel, timeout_nanos) | Result<T, String> | Dequeue, waiting for at most that many nanoseconds; 0 polls |
Channel.depth(channel) | Int | Current item count, or -1 for an unknown handle |
Channel.byte_depth(channel) | Int | Current queued bytes, or -1 for an unknown handle |
Channel.dropped(channel) | Int | Values rejected or replaced, or -1 for an unknown handle |
Overflow policies are:
:reject_newest— keep queued values and returnErr("channel full")for the new value.:drop_oldest— discard the oldest value and enqueue the new value.:latest_only— discard every queued value and retain only the newest.
Channel.recv takes nanoseconds. Use Duration.millis or Duration.seconds to make the unit explicit and to detect overflow.
A queued value takes eight bytes plus the size of whatever it references: a String, a list, or a struct is copied out of the sender's heap as it is sent, and into the receiver's as it is received, so actors can share a channel. bounded_bytes bounds that total too; a value larger than the whole byte capacity is rejected, whatever the policy, with Err("value exceeds the channel byte capacity"). try_send returns Ok(0) on acceptance and may return "channel busy" instead of waiting for the shared registry lock. An actor waiting in recv is suspended until a value arrives or the timeout passes, so its scheduler worker runs other actors meanwhile.
Process Names and Shutdown
Process manages actors within the current runtime and coordinates graceful shutdown.
| Function | Returns | Description |
|---|---|---|
Process.register(name, pid) | Int | Register a local name; 0 is success and 1 is failure |
Process.whereis(name) | Pid | Resolve a local name; PID 0 (<0.0>) means not found |
Process.monitor(pid, message) | Int | Send the calling actor message when pid ends; returns the monitor's reference |
Process.demonitor(reference) | Int | Remove a monitor; 0 is success and 1 is failure |
Process.install_shutdown_signals() | Unit | Treat native SIGINT and SIGTERM as shutdown requests |
Process.shutdown_requested() | Bool | Read the process-wide shutdown flag |
Process.request_shutdown() | Unit | Set the shutdown flag programmatically |
Process.exit(status) | Unit | Exit the native process with a status code, 0 to 255 (any other status exits with 1) |
Process.whereis returns an untyped Pid, which accepts any message, so a send through it does not tell the compiler what the actor receives. When an actor is reached only that way, give its pid a message type where it is spawned; otherwise the actor is error E0079:
actor counter() do
receive do
n -> println("got #{n}")
end
end
fn main() do
let pid :: Pid<Int> = spawn(counter)
Process.register("counter", pid)
send(Process.whereis("counter"), 5)
Timer.sleep(50)
endHTTP servers observe the same shutdown flag: after shutdown is requested they stop accepting new connections and drain connections already accepted.
Create and remove monitors from the actor that owns them: a monitor's message is one of that actor's messages, so Process.monitor outside an actor is a compile-time error (E0086). Process.demonitor returns 1 for a reference the calling actor does not hold.
Deterministic Random Values
Random is a deterministic, explicitly state-threaded pseudo-random generator. It is useful for simulations, scheduling choices, and reproducible tests; it is not a cryptographic random source.
fn main() do
let state = Random.seed(42)
let (next_state, value) = Random.next_int(state, 1, 100)
println("state=#{next_state}, value=#{value}")
end| Function | Returns | Description |
|---|---|---|
Random.seed(seed) | Int | Normalize a deterministic generator state |
Random.next_int(state, min, max) | (Int, Int) | Return (next_state, value) with min <= value <= max; panics when min > max |
Random.next_unit_ppm(state) | (Int, Int) | Return (next_state, value) where 0 <= value < 1_000_000 |
Next Steps
- Type System -- structs, generics, traits, and deriving
- Distributed Mesh -- remote actors and global process names
- Syntax Cheatsheet -- quick reference for all Mesh syntax