//! RFC 010 — `Discovery::Withdrawn`: a strategy retracts a candidate and the //! connector stops dialing it. //! //! Cross-process: a plain *server* node, and a *client* whose strategy is a //! script: announce a decoy `(ghost, addr)` where `addr` is a raw //! `TcpListener` the client itself holds (an OS thread accepts and //! immediately closes, so every dial fails at handshake and the connector //! keeps retrying on backoff — the accept count is the dial count); after a //! beat, withdraw the decoy and announce the real server. The client waits //! for the server's `node_up` — which is *after* the withdrawal in the //! strategy's own stream — then watches the decoy's accept count stay flat //! across a window longer than the pending backoff. Before withdrawal it //! must have been climbing (≥ 1), or the negative proves nothing. #![cfg(feature = "cluster")] mod common; use common::{maybe_child, spawn_node}; use smarm::channel::Sender; use smarm::cluster::discovery::{Discovery, Strategy}; use smarm::cluster::envelope::NodeMeta; use smarm::cluster::membership::{subscribe, NodeEvent}; use smarm::cluster::{start, Config, StaticSeeds, Timing}; use std::net::TcpListener; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; use std::time::{Duration, Instant}; const ROLES: &[(&str, fn())] = &[("server", role_server), ("client", role_client)]; fn meta() -> NodeMeta { NodeMeta { role: "withdraw".into(), region: "local".into(), } } fn role_server() { smarm::run(|| { let cluster = start(Config { node_name: "server".into(), meta: meta(), listen_addr: std::env::var("SMARM_LISTEN_ADDR") .unwrap_or_else(|_| "127.0.0.1:0".into()), strategy: Box::new(StaticSeeds::new(Vec::<(String, String)>::new())), timing: Timing::default(), }) .expect("binds"); println!("LISTENING {}", cluster.local_addr()); loop { smarm::sleep(Duration::from_secs(3600)); } }); } /// Scripted strategy: decoy, pause, withdraw decoy, real server, done. struct Script { decoy: String, server: String, } impl Strategy for Script { fn run(self: Box, out: Sender) { let _ = out.send(Discovery::Candidate { name: "ghost".into(), addr: self.decoy.clone(), }); // Long enough for the 250ms/500ms retries to land: ≥ 3 dials. smarm::sleep(Duration::from_millis(1100)); let _ = out.send(Discovery::Withdrawn { name: "ghost".into(), addr: self.decoy, }); let _ = out.send(Discovery::Candidate { name: "server".into(), addr: self.server, }); } } fn role_client() { let server_addr = std::env::var("SMARM_SERVER_ADDR").expect("SMARM_SERVER_ADDR"); // The decoy: accept-and-close on an OS thread; count every accept. let decoy = TcpListener::bind("127.0.0.1:0").unwrap(); let decoy_addr = decoy.local_addr().unwrap().to_string(); let dials = Arc::new(AtomicUsize::new(0)); let counter = dials.clone(); std::thread::spawn(move || { for conn in decoy.incoming() { counter.fetch_add(1, Ordering::SeqCst); drop(conn); } }); smarm::run(move || { let _cluster = start(Config { node_name: "client".into(), meta: meta(), listen_addr: "127.0.0.1:0".into(), strategy: Box::new(Script { decoy: decoy_addr, server: server_addr, }), timing: Timing::default(), }) .expect("binds"); let ev = subscribe().unwrap(); loop { match ev.rx.recv() { Ok(NodeEvent::NodeUp(i)) if i.name == "server" => break, Ok(_) => continue, Err(_) => panic!("manager gone"), } } // The withdrawal preceded the server candidate in the strategy's // stream, so it has been applied. Any dial that started before it // is bounded by the connect+handshake deadlines; let it drain, then // hold the count flat across a window longer than the pending // backoff would be (1s at this point, 2s next). let before = dials.load(Ordering::SeqCst); smarm::sleep(Duration::from_millis(500)); let settled = dials.load(Ordering::SeqCst); let t0 = Instant::now(); while t0.elapsed() < Duration::from_millis(3000) { smarm::sleep(Duration::from_millis(100)); } let after = dials.load(Ordering::SeqCst); println!("WITHDRAWN before={before} settled={settled} after={after}"); println!("CLIENT DONE"); loop { smarm::sleep(Duration::from_secs(3600)); } }); } #[test] fn withdrawn_candidate_is_no_longer_dialed() { maybe_child(ROLES); let mut server = spawn_node("server", &[]); let saddr = server.wait_listening(); let mut client = spawn_node("client", &[("SMARM_SERVER_ADDR", &saddr)]); let line = client.wait_line("WITHDRAWN", |l| l.starts_with("WITHDRAWN ")); let mut nums = line .split_whitespace() .skip(1) .map(|kv| kv.split_once('=').unwrap().1.parse::().unwrap()); let (before, settled, after) = ( nums.next().unwrap(), nums.next().unwrap(), nums.next().unwrap(), ); assert!( before >= 1, "decoy was never dialed; the negative proves nothing: {line}" ); assert_eq!( settled, after, "connector kept dialing a withdrawn candidate: {line}" ); client.wait_line("CLIENT DONE", |l| l == "CLIENT DONE"); }