Monitors could be installed but never taken down. That gap was about to
bite: the gen_server call timeout we want next is the Erlang dance —
monitor the server, wait for the reply or a Down or a deadline, then
demonitor — and without a way to remove a registration, every timed-out
call would leak a monitor on the server's slot and risk a stale Down
arriving later. So this lands the cleanup primitive before anything
depends on the old shape.
The decision flagged in the roadmap ("decide the monitor API NOW") is
resolved by giving each registration a process-unique MonitorId and
returning it to the caller. monitor() now hands back a
Monitor { id, target, rx } rather than a bare Receiver<Down>: read the
notice from rx as before, and pass &Monitor to demonitor to tear exactly
that registration down. The id comes from a monotonic counter in shared
state, bumped under the same lock that does the registration, so it's
deterministic and never reused — which is what lets demonitor name one of
several monitors on the same target unambiguously. target rode along on
the struct (over the sketched {id, rx}) purely so demonitor can go
straight to the slot instead of scanning every slot for the id.
demonitor returns Option<MonitorId> rather than a bool: Some(id) when a
live registration was found and removed, None when there was nothing left
to remove — it already fired (the registration is drained on finalize),
it was a NoProc, or the slot has been reclaimed. The generation half of
the pid quietly protects that last case: a recycled slot index fails
slot_mut's generation check, so a late demonitor is a clean no-op and can
never strip a different actor's monitor that happens to share the index.
The one piece of real care is reentrancy. Removing a registration drops
the slot's Sender, and Sender::drop can unpark a parked receiver, which
re-enters the shared mutex — which is not reentrant. So demonitor moves
the sender out of the Vec under the lock and lets it drop only after the
lock is released, the same discipline finalize_actor already follows for
its monitor and supervisor sends.
Flushing a Down the target already queued isn't a separate flag; it falls
out of dropping the Monitor. demonitor(&m); drop(m) stops future notices
and discards any queued one — the gen_server-call cleanup in one move.
Storage is now Vec<(MonitorId, Sender<Down>)>. The three slot-reset sites
were left alone on purpose: they clear/rebuild the Vec, which doesn't care
about the element type, so there's no fourth reset obligation. finalize
just destructures (_, m). Because chained_spawn and yield_many register no
monitors, that Vec stays empty on the hot path — taking an empty Vec costs
the same and the notify loop runs zero times — and a before/after general
probe confirmed both medians sit within noise.
Tests cover the three behaviours that matter: a demonitored watcher gets
no Down (its channel closes), demonitoring one of several leaves the
siblings firing normally, and demonitoring after the Down has already
fired reports None.
103 lines
4.3 KiB
Markdown
103 lines
4.3 KiB
Markdown
# smarm
|
|
|
|
> SMARM — Smarm, Marks Actor Runtime Machinery. A proof-of-concept green-thread actor runtime for Rust.
|
|
|
|
Implements the core ideas in [`Achitecture.md`](.docs/Architecture.md): green-thread actors on a
|
|
shared heap, scheduled cooperatively, communicating only by `Send` messages.
|
|
Erlang's isolation model without Erlang's copying GC, Rust's zero-copy
|
|
ownership transfers without async's function colouring.
|
|
|
|
The scheduler is multi-threaded — one OS thread per available CPU, all drawing
|
|
from a shared run queue. The single-threaded `run()` entry point is kept as a
|
|
convenience wrapper around `runtime::init(Config::exact(1)).run(f)`.
|
|
|
|
## What's here
|
|
|
|
| Module | What it does |
|
|
|--------------|------------------------------------------------------------------------|
|
|
| `stack` | `mmap`'d growable stack with guard page; SIGSEGV on overflow |
|
|
| `context` | `#[naked]` x86-64 context-switch shims, callee-saved regs only |
|
|
| `preempt` | Allocator-driven preemption; `check!()` macro for no-alloc loops |
|
|
| `pid` | `(index, generation)` PIDs; stale handles are detectable, not silent |
|
|
| `actor` | Trampoline + `catch_unwind` boundary at the actor entry point |
|
|
| `scheduler` | Run queue, slot table, spawn/join, parking, idle path |
|
|
| `channel` | Unbounded MPSC channel; `recv` parks the actor |
|
|
| `mutex` | `Mutex<T>` with mandatory timeout; FIFO waiters; parks the green thread |
|
|
| `timer` | Min-heap of `(deadline, reason)`; `Sleep` and `WaitTimeout` reasons |
|
|
| `io` | `block_on_io` for blocking work; `wait_readable`/`wait_writable` + `read`/`write` via epoll |
|
|
| `supervisor` | `Signal::Exit`/`Panic`/`Stopped` funnelled to a parent; `OneForOne`/`OneForAll`/`RestForOne` strategies + restart-intensity cap |
|
|
| `monitor` | `monitor(pid)` → `Monitor { id, target, rx }`; one-shot `Down` via `rx`; `demonitor(&m)` tears one registration down; unidirectional death notice |
|
|
| `link` | bidirectional `link`/`unlink`; abnormal death propagates (cooperative stop, or an `ExitSignal` message under `trap_exit`) |
|
|
| `gen_server` | `call` (sync request-reply) / `cast` (async) over one inbox; `ServerRef` + `init`/`terminate` hooks; server-down via channel closure |
|
|
|
|
## Quick taste
|
|
|
|
```rust
|
|
use smarm::{run, spawn, channel};
|
|
|
|
run(|| {
|
|
let (tx, rx) = channel::<i64>();
|
|
let h = spawn(move || {
|
|
for _ in 0..3 {
|
|
let v = rx.recv().unwrap();
|
|
println!("got {v}");
|
|
}
|
|
});
|
|
for v in 1..=3i64 {
|
|
tx.send(v).unwrap();
|
|
}
|
|
h.join().unwrap();
|
|
});
|
|
```
|
|
|
|
## Layout
|
|
|
|
```
|
|
src/
|
|
stack.rs context.rs preempt.rs pid.rs actor.rs
|
|
scheduler.rs channel.rs mutex.rs timer.rs io.rs
|
|
supervisor.rs monitor.rs link.rs runtime.rs
|
|
gen_server.rs lib.rs
|
|
tests/
|
|
per-module integration tests
|
|
benches/
|
|
primes.rs fan-out/fan-in compute, vs tokio current_thread
|
|
```
|
|
|
|
## Building and running
|
|
|
|
Standard Cargo. Requires Rust 1.95 or newer (the `#[naked]` attribute went stable
|
|
in 1.88; we use a few unrelated post-1.88 features). x86-64 Linux only —
|
|
ARM64 and macOS are on the deferred list because of the assembly shim and the
|
|
epoll dependency.
|
|
|
|
```sh
|
|
cargo test # all tests
|
|
cargo test --test mutex # one module
|
|
cargo bench # primes benchmark vs tokio
|
|
```
|
|
|
|
## What's not here
|
|
|
|
See the **Defer** section of `Architecture.md`.
|
|
`join!` for handle groups, stack growth via remap,
|
|
hierarchical timer wheel, fd-wait timeouts, `Signal::Timeout`. Each is
|
|
mechanism we know how to add; none belongs in this iteration.
|
|
|
|
## Docs
|
|
|
|
| Document | What it covers |
|
|
|---|---|
|
|
| [`Architecture.md`](./docs/Architecture.md) | Design intent, runtime model, and deferred work |
|
|
| [`smarm - Deep Dive.html`](./docs/smarm%20-%20Deep%20Dive.html) | Generated walkthrough of the system; good starting point |
|
|
| [`BENCHMARKS_AND_TUNING.md`](./docs/BENCHMARKS_AND_TUNING.md) | Where smarm wins and loses vs tokio, preemption knob recommendations |
|
|
| [`benchmarks.md`](./docs/benchmarks.md) | Raw benchmark results, methodology, and tuning experiment log |
|
|
|
|
## Contributing
|
|
|
|
This is a personal proof-of-concept. There's no PR workflow. If you fork it and do something interesting, just send me an email. If it's nice, I'll upstream the changes.
|
|
|
|
---
|
|
|
|
|