From 2d15834b24e6183af8cc01f522e0b3f29ac7cbca Mon Sep 17 00:00:00 2001 From: smarm-agent Date: Fri, 19 Jun 2026 10:20:09 +0000 Subject: [PATCH] RFC 016 Chunk 4: observer example with ps-style + tree dump MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A runnable examples/observer.rs (required-features = ["observer"]) that stands up a named service + two parked workers, starts the observer, and renders a snapshot as a ps-style table and the parentage forest indented. The observer appears in its own dump, caught running while it serves the snapshot call — transport over the same read every consumer sees. Complements the runnable doctest already on observer::start. --- Cargo.toml | 5 ++ examples/observer.rs | 150 +++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 155 insertions(+) create mode 100644 examples/observer.rs diff --git a/Cargo.toml b/Cargo.toml index 15403b9..b9171c7 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -75,3 +75,8 @@ harness = false [[bench]] name = "switch_cost" harness = false + +# RFC 016 Chunk 4 — the live observer dump. Needs the optional gen_server. +[[example]] +name = "observer" +required-features = ["observer"] diff --git a/examples/observer.rs b/examples/observer.rs new file mode 100644 index 0000000..b65c669 --- /dev/null +++ b/examples/observer.rs @@ -0,0 +1,150 @@ +//! The live observer (RFC 016 Chunk 4) producing an OTP `observer`-flavoured +//! dump of a running system. +//! +//! Run it with the feature on: +//! +//! ```text +//! cargo run --example observer --features observer +//! ``` +//! +//! It stands up a tiny tree — a named service plus two workers parked on a gate +//! — starts the [`observer`](smarm::observer) gen_server, then asks it for a +//! snapshot and a tree over the call channel and renders both. The observer is +//! pure transport: every line below is the Chunk-1 read +//! ([`snapshot`](smarm::snapshot) / [`tree`](smarm::tree)) marshalled across a +//! `call`, nothing more. + +use smarm::observer::{self, ObserverReply, ObserverRequest}; +use smarm::{channel, register, run, spawn, ActorState, Name, RuntimeSnapshot, RuntimeTree, TreeNode}; + +const ECHO: Name = Name::new("echo"); + +fn state_glyph(s: ActorState) -> &'static str { + match s { + ActorState::Queued => "queued", + ActorState::Running => "running", + ActorState::Notified => "notified", + ActorState::Parked => "parked", + ActorState::Done => "done", + } +} + +/// A `ps`-style table over the flat snapshot. +fn print_snapshot(snap: &RuntimeSnapshot) { + println!("snapshot (format v{}, {} actors)", snap.format_version, snap.actors.len()); + println!( + " {:<10} {:<9} {:<10} {:>4} {:>4} {:>4} {:>4} {:>5} {}", + "pid", "state", "parent", "mon", "lnk", "joi", "mbox", "msgs", "names" + ); + for a in &snap.actors { + let parent = if a.supervisor.index() == u32::MAX { + "".to_string() + } else { + format!("{}.{}", a.supervisor.index(), a.supervisor.generation()) + }; + println!( + " {:<10} {:<9} {:<10} {:>4} {:>4} {:>4} {:>4} {:>5} {}", + format!("{}.{}", a.pid.index(), a.pid.generation()), + state_glyph(a.state), + parent, + a.monitors, + a.links, + a.joiners, + a.mailbox_depth, + a.messages_received, + if a.names.is_empty() { "-".to_string() } else { a.names.join(",") }, + ); + } +} + +/// The parentage forest, indented. +fn print_tree(t: &RuntimeTree) { + println!("tree (format v{})", t.format_version); + fn walk(node: &TreeNode, depth: usize) { + let indent = " ".repeat(depth + 1); + let flag = if node.orphaned { " [orphaned]" } else { "" }; + let names = if node.info.names.is_empty() { + String::new() + } else { + format!(" ({})", node.info.names.join(",")) + }; + println!( + "{indent}{}.{} {}{names}{flag}", + node.info.pid.index(), + node.info.pid.generation(), + state_glyph(node.info.state), + ); + for child in &node.children { + walk(child, depth + 1); + } + } + for root in &t.roots { + walk(root, 0); + } +} + +fn main() { + run(|| { + // A named echo service and two anonymous workers, all parked on a gate + // so the system holds still while we observe it. Each gets its own gate + // receiver (a Receiver is single-consumer); we keep the senders to + // release them at the end. + let (ready_tx, ready_rx) = channel::<()>(); + let mut gates = Vec::new(); + + let svc = { + let (gate_tx, gate_rx) = channel::<()>(); + gates.push(gate_tx); + let ready_tx = ready_tx.clone(); + spawn(move || { + let (cmd_tx, cmd_rx) = channel::(); + register(ECHO, cmd_tx).unwrap(); + ready_tx.send(()).unwrap(); + gate_rx.recv().unwrap(); + drop(cmd_rx); + }) + }; + let workers: Vec<_> = (0..2) + .map(|_| { + let (gate_tx, gate_rx) = channel::<()>(); + gates.push(gate_tx); + let ready_tx = ready_tx.clone(); + spawn(move || { + ready_tx.send(()).unwrap(); + gate_rx.recv().unwrap(); + }) + }) + .collect(); + + // Wait until all three have announced and parked. + for _ in 0..3 { + ready_rx.recv().unwrap(); + } + // Queue two commands at the echo service so its mailbox depth is visible. + smarm::send(ECHO, 1).unwrap(); + smarm::send(ECHO, 2).unwrap(); + + // Start the observer and dump the system through it. + let obs = observer::start(); + + let ObserverReply::Snapshot(snap) = obs.call(ObserverRequest::Snapshot).unwrap() else { + unreachable!() + }; + let ObserverReply::Tree(t) = obs.call(ObserverRequest::Tree).unwrap() else { + unreachable!() + }; + + print_snapshot(&snap); + println!(); + print_tree(&t); + + // Release everyone and drain. + for gate_tx in gates { + gate_tx.send(()).unwrap(); + } + svc.join().unwrap(); + for w in workers { + w.join().unwrap(); + } + }); +}