From ecb0835aa7fc8d701692003dd193f2f8b7138e6d Mon Sep 17 00:00:00 2001 From: smarm-agent Date: Thu, 18 Jun 2026 11:02:35 +0000 Subject: [PATCH] examples: typed_actor, named_genserver, worker_pool for the new surface MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit typed_actor and named_genserver land as written (the target-ergonomics spec): identity-bound Pid vs durable Name, and a gen_server addressed by a durable ServerName. worker_pool is reworked from the spec into a self-draining program: workers retire on a sentinel and report their tally, and the dispatcher waits the pool out — so it terminates on its own rather than leaning on root-exit teardown, while still exercising the full typed surface (spawn_addr / join / pick_as / members_as / dispatch) alongside the untyped escape hatch (members / pick + send_dyn). All three build unchanged-against-spec (typed_actor, named_genserver) and run to completion. --- examples/named_genserver.rs | 68 ++++++++++++++++++++++++++ examples/typed_actor.rs | 78 ++++++++++++++++++++++++++++++ examples/worker_pool.rs | 95 +++++++++++++++++++++++++++++++++++++ 3 files changed, 241 insertions(+) create mode 100644 examples/named_genserver.rs create mode 100644 examples/typed_actor.rs create mode 100644 examples/worker_pool.rs diff --git a/examples/named_genserver.rs b/examples/named_genserver.rs new file mode 100644 index 0000000..379987e --- /dev/null +++ b/examples/named_genserver.rs @@ -0,0 +1,68 @@ +//! Addressing a gen_server by a durable name. +//! +//! A gen_server is multi-message (call / cast over one inbox), so it is named +//! by the *server* type rather than by a single message type. Registering it +//! under a [`ServerName`] lets clients `call` and `cast` by name, resolving on +//! every use — so the address keeps working across a supervised restart, with +//! no stale [`ServerRef`] to refresh. + +use smarm::{call, cast, run, whereis_server, GenServer, ServerBuilder, ServerName, ServerRef}; + +/// A counter server: synchronous `Get`, asynchronous `Inc` / `Add`. +struct Counter { + n: u64, +} +enum Query { + Get, +} +enum Update { + Inc, + Add(u64), +} +impl GenServer for Counter { + type Call = Query; + type Reply = u64; + type Cast = Update; + type Info = (); + + fn handle_call(&mut self, q: Query) -> u64 { + match q { + Query::Get => self.n, + } + } + fn handle_cast(&mut self, u: Update) { + match u { + Update::Inc => self.n += 1, + Update::Add(k) => self.n += k, + } + } +} + +/// A durable name typed by the server, so by-name `call` / `cast` check against +/// `Counter`'s `Call` / `Cast` / `Reply`. +const COUNTER: ServerName = ServerName::new("counter"); + +fn main() { + run(|| { + // Start the server and bind its name in one step. A named start is + // fallible: the name may already be held by another live server. + ServerBuilder::new(Counter { n: 0 }) + .named(COUNTER) + .start() + .unwrap(); + + // Address it purely by name. Each call/cast resolves through the + // registry, so a server restarted under the same name is reached + // transparently. + cast(COUNTER, Update::Inc).unwrap(); + cast(COUNTER, Update::Add(41)).unwrap(); + assert_eq!(call(COUNTER, Query::Get).unwrap(), 42); + + // When you want a handle to hold or pass on rather than resolve per + // call, recover a typed `ServerRef` from the name. + let svc: Option> = whereis_server(COUNTER); + if let Some(svc) = svc { + let _ = svc.call(Query::Get); + } + }); +} diff --git a/examples/typed_actor.rs b/examples/typed_actor.rs new file mode 100644 index 0000000..7fc655d --- /dev/null +++ b/examples/typed_actor.rs @@ -0,0 +1,78 @@ +//! Addressing a single-message actor two ways: by identity and by name. +//! +//! A single-message actor (one implementing [`Addressable`]) is reachable +//! through either of smarm's two address kinds: +//! +//! - [`Pid`] — identity-bound. Names the exact incarnation, never +//! redirects, and stops resolving once that actor dies. +//! - [`Name`] — durable. Re-resolves through the registry on every send, +//! so it always reaches whoever currently holds the name. +//! +//! A name is typed by the *message* it carries; a pid by the *actor* type +//! (whose [`Addressable::Msg`] is that message). + +use smarm::{ + channel, lookup_as, register, run, send, send_to, spawn, spawn_addr, Addressable, Name, Pid, + Receiver, +}; + +/// A single-message actor. Its one message type is what a `Pid` delivers. +struct Echo; +impl Addressable for Echo { + type Msg = EchoMsg; +} +enum EchoMsg { + Say(String), + Stop, +} + +/// A durable, message-typed name. Declared as a constant and shared freely. +const ECHO: Name = Name::new("echo"); + +fn main() { + run(|| { + // --- Identity-bound: Pid --------------------------------------- + // + // `spawn_addr` makes the actor's inbox, hands the body its receiver, + // installs the sender, and returns a typed `Pid`. The inbox is + // published before the pid is returned, so the address is live the + // instant we hold it — an immediate `send_to` always resolves. + let echo: Pid = spawn_addr::(|rx: Receiver| { + while let Ok(msg) = rx.recv() { + match msg { + EchoMsg::Say(s) => println!("echo: {s}"), + EchoMsg::Stop => break, + } + } + }); + send_to(echo, EchoMsg::Say("hello".into())).unwrap(); + send_to(echo, EchoMsg::Stop).unwrap(); + + // --- Durable: Name --------------------------------------------- + // + // An actor claims a name for its own inbox; senders resolve it on every + // send, so the binding outlives any single holder. + let (ready_tx, ready_rx) = channel::<()>(); + spawn(move || { + let (tx, rx) = channel::(); + register(ECHO, tx).unwrap(); + ready_tx.send(()).unwrap(); + while let Ok(msg) = rx.recv() { + match msg { + EchoMsg::Say(s) => println!("named echo: {s}"), + EchoMsg::Stop => break, + } + } + }); + ready_rx.recv().unwrap(); // the actor has claimed the name + + send(ECHO, EchoMsg::Say("by name".into())).unwrap(); + + // Recover a typed pid from the name when you want identity-bound sends: + // `lookup_as` re-types the registry's erased pid as `Pid`, so you + // drop back onto the compile-checked `send_to`. + if let Some(p) = lookup_as::("echo") { + send_to(p, EchoMsg::Stop).unwrap(); + } + }); +} diff --git a/examples/worker_pool.rs b/examples/worker_pool.rs new file mode 100644 index 0000000..26ab15d --- /dev/null +++ b/examples/worker_pool.rs @@ -0,0 +1,95 @@ +//! A typed worker pool over process groups (`pg`). +//! +//! Workers enroll in a named group; the dispatcher reaches them through it. +//! Because the pool is homogeneous — every member is a `Worker` — the group's +//! typed reads (`pick_as` / `members_as`) and the `dispatch` combinator hand +//! members back as `Pid`, so every send is an ordinary compile-checked +//! `send_to` rather than the untyped escape hatch. The untyped reads +//! (`members` / `pick` + `send_dyn`) stay available for identity-only use — +//! counting, logging, monitoring — where the message type isn't known. +//! +//! The pool drains itself: each worker retires on a sentinel and reports its +//! tally back, and the dispatcher waits the pool out before returning. (A pool +//! left running would be stopped anyway when the root actor exits, but draining +//! explicitly keeps the example deterministic.) + +use smarm::{ + channel, dispatch, join, members, members_as, pick, pick_as, run, send_dyn, send_to, + spawn_addr, Addressable, Pid, +}; + +/// The pool's worker actor. One message type, carried by `Pid`. +struct Worker; +impl Addressable for Worker { + type Msg = Job; +} + +/// A unit of work, or the sentinel that retires a worker. +enum Job { + Task { id: u64 }, + Retire, +} + +const POOL: &str = "pool"; +const WORKERS: u64 = 4; + +fn main() { + run(|| { + // Mint four typed workers and enroll them. `spawn_addr` yields a + // `Pid`; `join` takes a typed pid directly, erasing internally. + // Each worker reports how many tasks it handled over a shared channel, + // so the dispatcher can wait the pool out at the end. + let (done_tx, done_rx) = channel::<(u64, u64)>(); + for id in 0..WORKERS { + let done_tx = done_tx.clone(); + let w: Pid = spawn_addr::(move |rx| { + let mut handled = 0; + while let Ok(job) = rx.recv() { + match job { + Job::Task { id: job } => { + println!("worker {id} handling job {job}"); + handled += 1; + } + Job::Retire => break, + } + } + done_tx.send((id, handled)).unwrap(); + }); + join(POOL, w); + } + drop(done_tx); // from here only the workers hold senders + + // Pick one live member and hand it a job — typed end to end. + if let Some(w) = pick_as::(POOL) { + send_to(w, Job::Task { id: 1 }).unwrap(); + } + + // `dispatch` rolls pick-a-live-member-and-send into one call, returning + // the member it reached (or handing the job back if the pool is empty). + let _ = dispatch::(POOL, Job::Task { id: 2 }); + + // Fan a job out to the whole pool — typed, so each send is `send_to`. + for w in members_as::(POOL) { + let _ = send_to(w, Job::Task { id: 3 }); + } + + // The untyped reads stay available for identity-only use. To message a + // member reached this way, the explicit `send_dyn` escape hatch names + // the message type. + println!("pool size: {}", members(POOL).len()); + if let Some(any) = pick(POOL) { + let _ = send_dyn::(any, Job::Task { id: 4 }); + } + + // Retire every worker, then drain their tallies so the program winds + // down on its own. + for w in members_as::(POOL) { + let _ = send_to(w, Job::Retire); + } + for _ in 0..WORKERS { + if let Ok((id, handled)) = done_rx.recv() { + println!("worker {id} retired after {handled} jobs"); + } + } + }); +}