From 160967939b51bbd742f77d05b1cb63cb9bf46cfa Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 15 Aug 2026 07:07:14 +0000 Subject: [PATCH] =?UTF-8?q?feat(cluster):=20RFC=20010=20c7a=20=E2=80=94=20?= =?UTF-8?q?membership=20events=20+=20view=20at=20the=20manager?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit node_up/node_down are derived facts of the manager's own register/remove events, so the membership state lives in the manager (no cross-actor race between 'connection exists' and 'node is up'); src/cluster/membership.rs is the consumer surface: NodeEvent/NodeInfo, subscribe(), view(). The conn table stays private — no consumer touches it (roadmap-binding). Ratified semantics: subscribe() is snapshot-then-stream — one NodeUp per live peer is queued before the subscription joins the list, exact because gen_server handlers are serialized. Dropped subscribers are pruned on the next emit (closed channel), no monitor needed. One-viable call, flagged: NodeId is memoized per (name, incarnation) — a compact local alias for the wire identity, per pg.rs's framing. A reconnect blip at the same incarnation keeps its id; a restart (new incarnation) gets a fresh one, so a ghost and its successor are always distinguishable. Allocation starts at 1; NodeId(0) stays pg::DEFAULT_NODE_ID (self). Call::Register now carries the whole handshake Peer (the path already has it; node_up needs incarnation + meta). tests/cluster_membership.rs 4/0 stable over 5 runs (live up/down over localhost TCP with commanded and EOF teardown; late-subscriber snapshot + view agreement; restart-vs-blip id identity; dead-subscriber pruning). All cluster suites regression-clean; clippy --lib green both configs; fmt clean; default build compiles. --- src/cluster.rs | 2 + src/cluster/conn.rs | 4 +- src/cluster/manager.rs | 119 +++++++++++++--- src/cluster/membership.rs | 97 ++++++++++++++ tests/cluster_membership.rs | 261 ++++++++++++++++++++++++++++++++++++ 5 files changed, 465 insertions(+), 18 deletions(-) create mode 100644 src/cluster/membership.rs create mode 100644 tests/cluster_membership.rs diff --git a/src/cluster.rs b/src/cluster.rs index aea18d8..286ed0c 100644 --- a/src/cluster.rs +++ b/src/cluster.rs @@ -13,6 +13,7 @@ pub mod connect; pub mod envelope; pub mod handshake; pub mod manager; +pub mod membership; pub mod transport; use std::time::Duration; @@ -25,6 +26,7 @@ use crate::supervisor::{ChildSpec, OneForOne, Restart}; pub use conn::{spawn_established, ConnHandle}; pub use connect::{dial, spawn_acceptor, AcceptorHandle}; pub use manager::{Manager, MANAGER}; +pub use membership::{subscribe, view, MembershipEvents, NodeEvent, NodeInfo}; /// c6d — the derived build hash for [`handshake::LocalNode::build_hash`]: /// two builds may mesh only when this matches, and it is a pure function of diff --git a/src/cluster/conn.rs b/src/cluster/conn.rs index 3eee6c5..b2953b5 100644 --- a/src/cluster/conn.rs +++ b/src/cluster/conn.rs @@ -76,12 +76,12 @@ pub struct RegisterRefused; /// actor. A refusal has already stopped the actor and closed the socket. pub fn spawn_established(framed: FramedConn, peer: Peer) -> Result { let (cmd_tx, cmd_rx) = channel(); - let name = peer.node_name.clone(); + let reg_peer = peer.clone(); let pid = spawn(move || run(framed, peer, cmd_rx)).pid(); match gen_server::call( MANAGER, Call::Register { - name, + peer: reg_peer, pid, handle: ConnHandle { cmd_tx }, }, diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs index b5153b1..4a40152 100644 --- a/src/cluster/manager.rs +++ b/src/cluster/manager.rs @@ -10,17 +10,25 @@ //! that panics, is cancelled, or closes cleanly is removed without //! cooperation from the dying actor. //! -//! Membership as consumers will see it (the `node_up`/`node_down` interface and -//! the view) and the connector dial loop are c7, built on top of this table. -//! What lives here is only the table itself and the uniqueness rule the -//! handshake's `NameTaken` verdict depends on. +//! The manager also holds the **membership state** (c7a): `node_up` fires on +//! a successful registration and `node_down` on removal — they are derived +//! facts of the exact events this table already owns, so holding the view +//! here means no cross-actor race between "connection exists" and "node is +//! up". The consumer surface (event types, [`subscribe`], [`view`], +//! semantics) is [`membership`](crate::cluster::membership); no consumer +//! ever touches the table itself. +//! +//! The connector dial loop is c7b, built on top of both. use std::collections::HashMap; +use crate::channel::Sender; use crate::cluster::conn::ConnHandle; -use crate::cluster::handshake::HelloCtx; +use crate::cluster::handshake::{HelloCtx, Peer}; +use crate::cluster::membership::{NodeEvent, NodeInfo}; use crate::gen_server::{GenServer, GenServerCtx, GenServerName, Watcher}; use crate::monitor::{monitor, Down}; +use crate::pg::NodeId; use crate::pid::Pid; /// Well-known name of the singleton manager within a runtime. Connection @@ -28,15 +36,18 @@ use crate::pid::Pid; /// manager is always found at the same key. pub const MANAGER: GenServerName = GenServerName::new("smarm.cluster.manager"); -/// One live connection's entry: the actor running it, and the handle whose -/// lifetime *is* the connection's (see the module docs). +/// One live connection's entry: the actor running it, the handle whose +/// lifetime *is* the connection's (see the module docs), and the peer's +/// membership identity (what `node_up` announced and `node_down` will name). struct ConnEntry { pid: Pid, + info: NodeInfo, _handle: ConnHandle, } /// The connection registry: peer name → the connection actor that owns that -/// peer's control connection. +/// peer's control connection. Plus the membership state layered on it (c7a): +/// subscribers, and the `(name, incarnation)` → [`NodeId`] memo. pub struct Manager { conns: HashMap, /// In-flight dial intents: peer name -> the actor performing the dial. @@ -45,6 +56,16 @@ pub struct Manager { /// the dial resolves, and — because the dialer is monitored — on the /// dialer's death, so a panicking dial can never wedge the tie-break. dials: HashMap, + /// Membership subscribers; a closed channel is pruned on the next emit. + subscribers: Vec>, + /// The [`NodeId`] memo: a reconnect at the same incarnation keeps its id, + /// a restart (new incarnation) allocates a fresh one. Grows one entry per + /// distinct `(name, incarnation)` ever seen — unbounded in principle, + /// bounded in practice by restarts actually happening. + ids: HashMap<(String, u32), NodeId>, + /// Next id to allocate. Starts at 1: id 0 is + /// [`DEFAULT_NODE_ID`](crate::pg::DEFAULT_NODE_ID), the local node. + next_id: u32, watcher: Option>, } @@ -53,9 +74,31 @@ impl Manager { Manager { conns: HashMap::new(), dials: HashMap::new(), + subscribers: Vec::new(), + ids: HashMap::new(), + next_id: 1, watcher: None, } } + + /// The memoized id for `(name, incarnation)` — see the field docs. + fn node_id(&mut self, name: &str, incarnation: u32) -> NodeId { + *self + .ids + .entry((name.to_string(), incarnation)) + .or_insert_with(|| { + let id = NodeId::new(self.next_id); + self.next_id += 1; + id + }) + } + + /// Deliver `event` to every live subscriber, pruning the dead: a closed + /// channel means the subscriber dropped its [`MembershipEvents`] + /// (crate::cluster::membership::MembershipEvents). + fn emit(&mut self, event: &NodeEvent) { + self.subscribers.retain(|tx| tx.send(event.clone()).is_ok()); + } } impl Default for Manager { @@ -77,12 +120,13 @@ pub enum Registered { /// Requests to the manager. pub enum Call { /// The path claims its peer's name for a freshly-established connection, - /// handing the manager the actor's [`ConnHandle`]. The manager monitors - /// `pid` and holds the handle for as long as the entry lives; a + /// handing the manager the actor's [`ConnHandle`] and the handshake's + /// [`Peer`] (the membership identity `node_up` announces). The manager + /// monitors `pid` and holds the handle for as long as the entry lives; a /// [`Registered::Duplicate`] verdict drops the handle here, which stops /// the refused actor. Register { - name: String, + peer: Peer, pid: Pid, handle: ConnHandle, }, @@ -101,6 +145,15 @@ pub enum Call { /// The [`HelloCtx`] for an inbound `Hello` offering `peer_name` — the /// accept path asks this between reading the frame and judging it. HelloCtx { peer_name: String }, + /// Subscribe `tx` to membership events, snapshot-then-stream: one + /// [`NodeEvent::NodeUp`] per live peer is queued into `tx` before this + /// call answers, so the stream is exact from its first event (handlers + /// are serialized — nothing interleaves with the snapshot). Use + /// [`subscribe`](crate::cluster::membership::subscribe). + Subscribe { tx: Sender }, + /// The current view: every live peer's [`NodeInfo`], unordered. Use + /// [`view`](crate::cluster::membership::view). + View, } /// Replies from the manager. @@ -113,6 +166,8 @@ pub enum Reply { DialBegan(bool), DialEnded, HelloCtx(HelloCtx), + Subscribed, + View(Vec), } impl GenServer for Manager { @@ -128,26 +183,38 @@ impl GenServer for Manager { fn handle_call(&mut self, request: Call) -> Reply { match request { - Call::Register { name, pid, handle } => { - if self.conns.contains_key(&name) { + Call::Register { peer, pid, handle } => { + if self.conns.contains_key(&peer.node_name) { // `handle` drops here: the refused actor stops itself. return Reply::Registered(Registered::Duplicate); } if let Some(w) = &self.watcher { w.watch(monitor(pid)); } + let info = NodeInfo { + node: self.node_id(&peer.node_name, peer.incarnation.get()), + name: peer.node_name.clone(), + incarnation: peer.incarnation, + meta: peer.meta, + }; self.conns.insert( - name, + peer.node_name, ConnEntry { pid, + info: info.clone(), _handle: handle, }, ); + self.emit(&NodeEvent::NodeUp(info)); Reply::Registered(Registered::Ok) } Call::Disconnect { name } => { // Dropping the entry drops the handle, which stops the actor. - self.conns.remove(&name); + if let Some(entry) = self.conns.remove(&name) { + self.emit(&NodeEvent::NodeDown { + node: entry.info.node, + }); + } Reply::Disconnected } Call::Peers => { @@ -173,13 +240,33 @@ impl GenServer for Manager { name_claimed: self.conns.contains_key(&peer_name), dialing_this_peer: self.dials.contains_key(&peer_name), }), + Call::Subscribe { tx } => { + // The snapshot: queued before `tx` joins the list, and — the + // handlers being serialized — before any later event. + for entry in self.conns.values() { + let _ = tx.send(NodeEvent::NodeUp(entry.info.clone())); + } + self.subscribers.push(tx); + Reply::Subscribed + } + Call::View => Reply::View(self.conns.values().map(|e| e.info.clone()).collect()), } } fn handle_cast(&mut self, _request: ()) {} fn handle_down(&mut self, down: Down) { - self.conns.retain(|_, entry| entry.pid != down.pid); + let mut downs = Vec::new(); + self.conns.retain(|_, entry| { + let dead = entry.pid == down.pid; + if dead { + downs.push(entry.info.node); + } + !dead + }); + for node in downs { + self.emit(&NodeEvent::NodeDown { node }); + } self.dials.retain(|_, pid| *pid != down.pid); } } diff --git a/src/cluster/membership.rs b/src/cluster/membership.rs new file mode 100644 index 0000000..4a97390 --- /dev/null +++ b/src/cluster/membership.rs @@ -0,0 +1,97 @@ +//! RFC 010 c7a — membership: `node_up`/`node_down` events and the view. +//! +//! The membership *state* lives inside the [`manager`](crate::cluster::manager) +//! — `node_up` and `node_down` are derived facts of the exact events the +//! manager already owns (a successful registration; a reap or `Disconnect`), +//! so holding the view anywhere else would only add a cross-actor ordering +//! seam. This module is the consumer surface: the event and view types, and +//! the [`subscribe`]/[`view`] entry points. No consumer ever touches the +//! connection table (roadmap-binding, enforced by module privacy: the table +//! is a private field, and nothing here exposes names→pids). +//! +//! ## Subscription semantics (ratified 2026-08-15) +//! +//! [`subscribe`] is **snapshot-then-stream**: the returned receiver first +//! yields one [`NodeEvent::NodeUp`] per currently-live peer, then live events +//! as they happen. Because the manager is a `gen_server` (handlers are +//! serialized), the snapshot is exact — no event can interleave with it, and +//! per-subscriber ordering matches the manager's processing order. There is +//! no join-race for late subscribers and no separate "get, then diff" dance; +//! [`view`] exists for observation, not for synchronization. +//! +//! A dropped subscriber is pruned on the next emission (its channel reports +//! closed) — no monitor needed, the sender itself tells us. +//! +//! ## NodeId identity +//! +//! A [`NodeId`] is a compact **local alias for the wire identity** +//! `(node_name, incarnation)`, memoized by the manager: a reconnect blip at +//! the same incarnation keeps its id (down, then up, same id), while a +//! restart — a new incarnation — gets a fresh one, so a node's ghost and its +//! successor are always distinguishable. Ids are allocated from 1; +//! [`DEFAULT_NODE_ID`](crate::pg::DEFAULT_NODE_ID) (0) remains the local +//! node, per [`pg`](crate::pg)'s framing. + +use crate::channel::{channel, Receiver}; +use crate::cluster::envelope::NodeMeta; +use crate::cluster::manager::{Call, Reply, MANAGER}; +use crate::gen_server; +use crate::pg::{Incarnation, NodeId}; + +/// One live remote node, as the view and [`NodeEvent::NodeUp`] describe it. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct NodeInfo { + /// The local alias for `(name, incarnation)` — see the module docs. + pub node: NodeId, + /// The peer's claimed node name (handshake-verified). + pub name: String, + /// The peer's incarnation epoch, as offered in its `Hello`. + pub incarnation: Incarnation, + /// The peer's `Hello` metadata. + pub meta: NodeMeta, +} + +/// A membership change, as delivered to subscribers. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum NodeEvent { + /// A peer's control connection established and registered. + NodeUp(NodeInfo), + /// That peer's connection ended — reaped, commanded down, or the manager + /// itself shut down. Which [`NodeInfo`] this id named was delivered in + /// the corresponding `NodeUp`. + NodeDown { node: NodeId }, +} + +/// A live membership subscription: the receiving end of the event stream +/// (the [`Monitor`](crate::monitor::Monitor) shape — read from [`rx`], drop +/// to unsubscribe). +/// +/// [rx]: MembershipEvents::rx +pub struct MembershipEvents { + /// The event stream: the snapshot's `NodeUp`s first, then live events. + /// Fold it into a `select` from a plain actor, or pipe it into a + /// `gen_server` via `with_info`. + pub rx: Receiver, +} + +/// Subscribe to membership events (snapshot-then-stream — see the module +/// docs). `None`: the manager is not running. Must be called from inside an +/// actor. +pub fn subscribe() -> Option { + let (tx, rx) = channel(); + match gen_server::call(MANAGER, Call::Subscribe { tx }) { + Ok(Reply::Subscribed) => Some(MembershipEvents { rx }), + _ => None, + } +} + +/// The current view: every live peer's [`NodeInfo`], unordered. For +/// observation and tests; consumers that need to *track* the view should +/// [`subscribe`] instead (the snapshot makes the stream self-sufficient). +/// `None`: the manager is not running. Must be called from inside an actor. +pub fn view() -> Option> { + match gen_server::call(MANAGER, Call::View) { + Ok(Reply::View(v)) => Some(v), + _ => None, + } +} diff --git a/tests/cluster_membership.rs b/tests/cluster_membership.rs new file mode 100644 index 0000000..cf2eb67 --- /dev/null +++ b/tests/cluster_membership.rs @@ -0,0 +1,261 @@ +//! RFC 010 c7a — membership events and the view, at the manager. +//! +//! Same construction as the c6a lifecycle suite: the handshake is bypassed, +//! connections are built already-established over localhost TCP pairs with +//! fabricated `Peer`s, and the manager is started plainly so the test can +//! terminate. What is under test is the membership layer that c7 adds to the +//! manager: `node_up`/`node_down` events to subscribers (snapshot-then-stream), +//! the view, and NodeId identity — memoized per `(name, incarnation)`, so a +//! reconnect blip keeps its id and a restart (new incarnation) gets a fresh one. +#![cfg(feature = "cluster")] + +use std::time::Duration; + +use smarm::cluster::envelope::NodeMeta; +use smarm::cluster::handshake::Peer; +use smarm::cluster::manager::{Call, Manager, Reply, MANAGER}; +use smarm::cluster::membership::{subscribe, view, MembershipEvents, NodeEvent}; +use smarm::cluster::spawn_established; +use smarm::cluster::transport::tcp::TcpTransport; +use smarm::cluster::transport::{Conn, FramedConn, Transport}; +use smarm::gen_server::{self, GenServerBuilder}; +use smarm::pg::{Incarnation, NodeId}; +use smarm::run; + +/// A fabricated post-handshake peer identity, with the incarnation under the +/// test's control (it is identity-bearing here, unlike in the c6a suite). +fn peer(name: &str, inc: u32) -> Peer { + Peer { + node_name: name.to_string(), + incarnation: Incarnation::new(inc), + meta: NodeMeta { + role: "test".to_string(), + region: "test".to_string(), + }, + } +} + +/// One established transport pair over localhost (TCP backlog covers the +/// sequential dial-then-accept, as in the c3 conformance suite). +fn pair(t: &dyn Transport) -> (Box, Box) { + let mut l = t.listen("127.0.0.1:0").unwrap(); + let a = t.dial(&l.local_addr()).unwrap(); + let b = l.accept().unwrap(); + (a, b) +} + +/// The next event, or a panic naming the wait. The bound is generous against +/// a sub-millisecond real cost. +fn next_event(ev: &MembershipEvents, waiting_for: &str) -> NodeEvent { + ev.rx + .recv_timeout(Duration::from_secs(5)) + .unwrap_or_else(|e| panic!("timed out waiting for {waiting_for}: {e:?}")) +} + +/// Assert the subscription is drained: no event is pending. +fn assert_quiet(ev: &MembershipEvents) { + assert!(matches!(ev.rx.try_recv(), Ok(None))); +} + +fn disconnect(name: &str) { + assert!(matches!( + gen_server::call( + MANAGER, + Call::Disconnect { + name: name.to_string() + } + ), + Ok(Reply::Disconnected) + )); +} + +/// Live subscription: an empty snapshot, then `NodeUp` on registration and +/// `NodeDown` (same id) on commanded disconnect and on peer EOF alike. +#[test] +fn subscriber_sees_up_and_down() { + run(|| { + let mgr = GenServerBuilder::new(Manager::new()) + .named(MANAGER) + .start() + .expect("manager name is free"); + + let ev = subscribe().expect("manager is up"); + assert_quiet(&ev); // nothing live: the snapshot is empty + + let t = TcpTransport; + let (a1, b1) = pair(&t); + let (a2, b2) = pair(&t); + spawn_established(FramedConn::new(a1), peer("node-b", 1)).expect("node-b registers"); + spawn_established(FramedConn::new(a2), peer("node-c", 1)).expect("node-c registers"); + + let up_b = match next_event(&ev, "node_up(node-b)") { + NodeEvent::NodeUp(info) => { + assert_eq!(info.name, "node-b"); + assert_eq!(info.incarnation, Incarnation::new(1)); + assert_eq!(info.meta.role, "test"); + info + } + other => panic!("expected node_up(node-b), got {other:?}"), + }; + let up_c = match next_event(&ev, "node_up(node-c)") { + NodeEvent::NodeUp(info) => { + assert_eq!(info.name, "node-c"); + info + } + other => panic!("expected node_up(node-c), got {other:?}"), + }; + assert_ne!(up_b.node, up_c.node, "distinct peers get distinct ids"); + + // Commanded disconnect: down with node-b's id. + disconnect("node-b"); + assert_eq!( + next_event(&ev, "node_down(node-b)"), + NodeEvent::NodeDown { node: up_b.node } + ); + + // Peer EOF, no command: down with node-c's id. + drop(b2); + assert_eq!( + next_event(&ev, "node_down(node-c)"), + NodeEvent::NodeDown { node: up_c.node } + ); + assert_quiet(&ev); + + drop(b1); + mgr.shutdown(); + }); +} + +/// Snapshot-then-stream: a subscriber arriving after connections established +/// receives one `NodeUp` per live peer before anything else, and the view +/// call agrees with it. +#[test] +fn late_subscriber_gets_snapshot_and_view_agrees() { + run(|| { + let mgr = GenServerBuilder::new(Manager::new()) + .named(MANAGER) + .start() + .expect("manager name is free"); + + let t = TcpTransport; + let (a1, b1) = pair(&t); + let (a2, b2) = pair(&t); + spawn_established(FramedConn::new(a1), peer("node-b", 1)).expect("node-b registers"); + spawn_established(FramedConn::new(a2), peer("node-c", 1)).expect("node-c registers"); + + let ev = subscribe().expect("manager is up"); + let mut names = Vec::new(); + for _ in 0..2 { + match next_event(&ev, "a snapshot node_up") { + NodeEvent::NodeUp(info) => names.push(info.name), + other => panic!("expected a snapshot node_up, got {other:?}"), + } + } + names.sort(); + assert_eq!(names, ["node-b", "node-c"]); + assert_quiet(&ev); // the snapshot is exactly the live set + + let mut v = view().expect("manager is up"); + v.sort_by(|a, b| a.name.cmp(&b.name)); + assert_eq!(v.len(), 2); + assert_eq!(v[0].name, "node-b"); + assert_eq!(v[1].name, "node-c"); + + disconnect("node-b"); + disconnect("node-c"); + drop((b1, b2)); + // Drain the two downs so the subscription ends quiet. + let _ = next_event(&ev, "node_down"); + let _ = next_event(&ev, "node_down"); + mgr.shutdown(); + }); +} + +/// NodeId identity: a restart (same name, new incarnation) is a NEW id — the +/// ghost and its successor are distinguishable — while a reconnect blip (same +/// name, same incarnation) keeps its id. +#[test] +fn restart_gets_new_id_blip_keeps_id() { + run(|| { + let mgr = GenServerBuilder::new(Manager::new()) + .named(MANAGER) + .start() + .expect("manager name is free"); + let ev = subscribe().expect("manager is up"); + let t = TcpTransport; + + let id = |e: NodeEvent, what: &str| -> NodeId { + match e { + NodeEvent::NodeUp(info) => info.node, + other => panic!("expected node_up ({what}), got {other:?}"), + } + }; + + // Up at incarnation 1, then the peer dies (EOF). + let (a1, b1) = pair(&t); + spawn_established(FramedConn::new(a1), peer("node-b", 1)).expect("registers"); + let id1 = id(next_event(&ev, "node_up inc 1"), "inc 1"); + drop(b1); + assert_eq!( + next_event(&ev, "node_down inc 1"), + NodeEvent::NodeDown { node: id1 } + ); + + // Restart: new incarnation, new id — the ghost's id is not reused. + let (a2, b2) = pair(&t); + spawn_established(FramedConn::new(a2), peer("node-b", 2)).expect("registers"); + let id2 = id(next_event(&ev, "node_up inc 2"), "inc 2"); + assert_ne!( + id1, id2, + "a restarted node must be distinguishable from its ghost" + ); + + // Blip: the same incarnation reconnects and keeps its id. + disconnect("node-b"); + assert_eq!( + next_event(&ev, "node_down inc 2"), + NodeEvent::NodeDown { node: id2 } + ); + let (a3, b3) = pair(&t); + spawn_established(FramedConn::new(a3), peer("node-b", 2)).expect("registers"); + let id3 = id(next_event(&ev, "node_up after blip"), "blip"); + assert_eq!( + id2, id3, + "a reconnect at the same incarnation is the same node" + ); + + disconnect("node-b"); + let _ = next_event(&ev, "final node_down"); + drop((b2, b3)); + mgr.shutdown(); + }); +} + +/// A dropped subscriber is pruned on the next emit and never disturbs the +/// manager or a live subscriber. +#[test] +fn dead_subscriber_is_pruned() { + run(|| { + let mgr = GenServerBuilder::new(Manager::new()) + .named(MANAGER) + .start() + .expect("manager name is free"); + + let dead = subscribe().expect("manager is up"); + drop(dead); + let live = subscribe().expect("manager is up"); + + let t = TcpTransport; + let (a1, b1) = pair(&t); + spawn_established(FramedConn::new(a1), peer("node-b", 1)).expect("registers"); + match next_event(&live, "node_up despite a dead co-subscriber") { + NodeEvent::NodeUp(info) => assert_eq!(info.name, "node-b"), + other => panic!("expected node_up, got {other:?}"), + } + + disconnect("node-b"); + let _ = next_event(&live, "node_down"); + drop(b1); + mgr.shutdown(); + }); +}