RFC 016 Chunk 4: observer example with ps-style + tree dump
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.
This commit is contained in:
@@ -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"]
|
||||
|
||||
@@ -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<u64> = 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 {
|
||||
"<root>".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::<u64>();
|
||||
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();
|
||||
}
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user