From 258f2fecb7a7fdf4e7052abc3caea11733870140 Mon Sep 17 00:00:00 2001 From: Zachery Aaron Shores-Chmielewski Date: Wed, 25 Feb 2026 00:23:15 +0700 Subject: [PATCH] fix: idle cpu usage and distribution CPU use on idle down from 17% to 1.5% via lazy runtime execution. Turned down SWIM gossip frequency. --- Cargo.lock | 1 - benches/mt_benchmarks.rs | 8 +- crates/bindings/python/src/lib.rs | 28 +- crates/dashboard/src/actor_detail_html.rs | 1 + crates/dashboard/src/actors_html.rs | 1 + crates/dashboard/src/dashboard_html.rs | 1 + crates/dashboard/src/topology_html.rs | 1 + crates/distribution/Cargo.toml | 4 +- crates/distribution/src/driver.rs | 496 ------------------ crates/distribution/src/iroh_driver.rs | 29 + crates/distribution/src/lib.rs | 2 - crates/distribution/src/node.rs | 22 +- crates/distribution/src/node_metadata.rs | 20 +- crates/distribution/src/snapshot.rs | 10 +- crates/distribution/src/swim/node.rs | 17 + crates/distribution/src/swim/probe.rs | 152 +++++- crates/distribution/tests/common/mod.rs | 1 + crates/distribution/tests/swim_node.rs | 1 + crates/distribution/tests/swim_probe.rs | 9 + .../distribution/tests/transport_and_codec.rs | 172 ------ crates/node/Cargo.toml | 1 - crates/node/src/main.rs | 205 ++------ crates/node/src/plugins/datastore_page.html | 1 + crates/node/src/plugins/distribution.rs | 63 ++- .../node/src/plugins/distribution_page.html | 44 +- src/channel.rs | 8 + src/config.rs | 27 - src/extension.rs | 4 + src/runtime.rs | 2 +- src/std/timer_wheel.rs | 7 + src/worker.rs | 39 +- tests/test_python.py | 4 - 32 files changed, 443 insertions(+), 938 deletions(-) delete mode 100644 crates/distribution/src/driver.rs diff --git a/Cargo.lock b/Cargo.lock index 034f52b..2de3a1b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1334,7 +1334,6 @@ dependencies = [ "serde", "serde_json", "swactor", - "swactor-transport", "tokio", ] diff --git a/benches/mt_benchmarks.rs b/benches/mt_benchmarks.rs index 12c0a1c..9d757eb 100644 --- a/benches/mt_benchmarks.rs +++ b/benches/mt_benchmarks.rs @@ -3,7 +3,7 @@ use std::time::{Duration, Instant}; use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion}; use swactor::{ actor::{ActorAddress, ActorInterface}, - config::{BackoffPolicy, RuntimeConfig}, + config::RuntimeConfig, runtime::{Ctx, Runtime}, }; @@ -16,12 +16,6 @@ fn mt_config(threads: usize, max_actors: usize, max_messages: usize) -> RuntimeC num_threads: threads, max_actors, channel_buffer_size: max_messages, - backoff_policy: BackoffPolicy { - spin_threshold: 32, - yield_threshold: 64, - sleep_increment_us: 10, - sleep_max_us: 100, - }, ..Default::default() } } diff --git a/crates/bindings/python/src/lib.rs b/crates/bindings/python/src/lib.rs index b5345b6..b46268d 100644 --- a/crates/bindings/python/src/lib.rs +++ b/crates/bindings/python/src/lib.rs @@ -5,7 +5,7 @@ use pyo3::prelude::*; use pyo3::types::PyModule; use ::swactor::actor::{Actor, ActorAddress, ActorInterface, AnyActor, Ctx, Environment, SpawnRequest}; -use ::swactor::config::{BackoffPolicy, RuntimeConfig}; +use ::swactor::config::RuntimeConfig; use ::swactor::runtime::{Inbox, Runtime, RuntimeHandle}; // ─── PyMsg newtype ─────────────────────────────────────────────────────────── @@ -231,14 +231,6 @@ pub struct PyRuntimeConfig { max_actors: usize, #[pyo3(get, set)] channel_buffer_size: usize, - #[pyo3(get, set)] - spin_threshold: u32, - #[pyo3(get, set)] - yield_threshold: u32, - #[pyo3(get, set)] - sleep_increment_us: u64, - #[pyo3(get, set)] - sleep_max_us: u64, } #[pymethods] @@ -249,28 +241,16 @@ impl PyRuntimeConfig { num_threads = 1, max_actors = 1_000, channel_buffer_size = 1_000, - spin_threshold = 64, - yield_threshold = 256, - sleep_increment_us = 50, - sleep_max_us = 1_000, ))] fn new( num_threads: usize, max_actors: usize, channel_buffer_size: usize, - spin_threshold: u32, - yield_threshold: u32, - sleep_increment_us: u64, - sleep_max_us: u64, ) -> Self { Self { num_threads, max_actors, channel_buffer_size, - spin_threshold, - yield_threshold, - sleep_increment_us, - sleep_max_us, } } } @@ -281,12 +261,6 @@ impl From for RuntimeConfig { num_threads: py.num_threads, max_actors: py.max_actors, channel_buffer_size: py.channel_buffer_size, - backoff_policy: BackoffPolicy { - spin_threshold: py.spin_threshold, - yield_threshold: py.yield_threshold, - sleep_increment_us: py.sleep_increment_us, - sleep_max_us: py.sleep_max_us, - }, ..Default::default() } } diff --git a/crates/dashboard/src/actor_detail_html.rs b/crates/dashboard/src/actor_detail_html.rs index 3d98f33..6364cee 100644 --- a/crates/dashboard/src/actor_detail_html.rs +++ b/crates/dashboard/src/actor_detail_html.rs @@ -392,6 +392,7 @@ pub const ACTOR_DETAIL_HTML: &str = r##" // ─── SSE ──────────────────────────────────────────────────── var es = new EventSource('/events'); + window.addEventListener('beforeunload', function() { es.close(); }); es.addEventListener('stats', function(e) { try { updateDetail(JSON.parse(e.data)); } catch(err) { console.error(err); } diff --git a/crates/dashboard/src/actors_html.rs b/crates/dashboard/src/actors_html.rs index 8a07027..f4fe739 100644 --- a/crates/dashboard/src/actors_html.rs +++ b/crates/dashboard/src/actors_html.rs @@ -721,6 +721,7 @@ pub const ACTORS_HTML: &str = r##" // ── SSE connection ───────────────────────────────────── var es = new EventSource('/events'); + window.addEventListener('beforeunload', function() { es.close(); }); es.addEventListener('stats', function(e) { try { diff --git a/crates/dashboard/src/dashboard_html.rs b/crates/dashboard/src/dashboard_html.rs index f777468..b3a3a1b 100644 --- a/crates/dashboard/src/dashboard_html.rs +++ b/crates/dashboard/src/dashboard_html.rs @@ -578,6 +578,7 @@ pub const DASHBOARD_HTML: &str = r##" } var es = new EventSource('/events'); + window.addEventListener('beforeunload', function() { es.close(); }); es.addEventListener('stats', function(e) { try { diff --git a/crates/dashboard/src/topology_html.rs b/crates/dashboard/src/topology_html.rs index 9f5b9da..f9f5f07 100644 --- a/crates/dashboard/src/topology_html.rs +++ b/crates/dashboard/src/topology_html.rs @@ -263,6 +263,7 @@ pub const TOPOLOGY_HTML: &str = r##" tick(); var es = new EventSource('/events'); + window.addEventListener('beforeunload', function() { es.close(); }); es.addEventListener('topology', function(e) { try { updateTopology(JSON.parse(e.data)); } catch(err) { console.error(err); } diff --git a/crates/distribution/Cargo.toml b/crates/distribution/Cargo.toml index c2567a3..d76fbd7 100644 --- a/crates/distribution/Cargo.toml +++ b/crates/distribution/Cargo.toml @@ -4,14 +4,12 @@ version = "0.1.0" edition = "2024" [features] -default = ["tcp"] -tcp = ["dep:swactor-transport", "swactor-transport/tcp"] +default = [] iroh = ["dep:iroh", "dep:tokio"] relay = ["iroh", "dep:iroh-relay"] [dependencies] swactor = { path = "../..", features = ["serde", "transport"] } -swactor-transport = { path = "../transport", optional = true } ed25519-dalek = { version = "2", features = ["rand_core"] } rand_core = { version = "0.6", features = ["getrandom"] } serde = { version = "1", features = ["derive"] } diff --git a/crates/distribution/src/driver.rs b/crates/distribution/src/driver.rs deleted file mode 100644 index 26c280f..0000000 --- a/crates/distribution/src/driver.rs +++ /dev/null @@ -1,496 +0,0 @@ -//! Network driver — bridges `DistributedNode` logic with TCP I/O. -//! -//! Translates outgoing `NodeAction`s into wire messages sent via `TcpTransport`, -//! and dispatches incoming wire messages to the appropriate `DistributedNode` -//! handler methods. -//! -//! The driver maintains a `PeerAddressBook` mapping `NodeId → SocketAddr`. -//! Address hints travel in the TCP wire frame (not in protocol messages), -//! keeping the protocol layer transport-agnostic. - -use std::collections::HashMap; -use std::net::{SocketAddr, TcpStream}; - -use serde::{Deserialize, Serialize}; -use swactor::actor::ActorAddress; -use swactor::transport::{NetworkMessage, WireEnvelope}; - -use crate::messages::*; -use crate::node::{DistributedNode, DistributedNodeConfig}; -use crate::snapshot::DistributionNodeSnapshot; -use crate::swim::node::NodeAction; -use swactor_transport::tcp::{TcpAcceptor, TcpTransport, encode_wire_envelope_with_hints}; -use crate::types::NodeId; - -/// Dummy destination address used in wire envelopes for SWIM protocol messages. -/// SWIM messages are routed by `SocketAddr`, not by `ActorAddress`, so this -/// field is unused but required by the wire format. -const SWIM_DEST: ActorAddress = ActorAddress([0u8; 32]); - -// ─── Address Hints ────────────────────────────────────────────────────────── - -/// An address hint bundled in the TCP wire frame. -/// -/// Each outgoing TCP message includes the sender's own (NodeId, SocketAddr) -/// as a hint. JoinResponse messages include all known member addresses. -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct AddressHint { - pub node_id: NodeId, - pub addr: SocketAddr, -} - -// ─── Peer Address Book ────────────────────────────────────────────────────── - -/// Maps NodeId → SocketAddr. Maintained by the TCP driver layer. -pub struct PeerAddressBook(HashMap); - -impl Default for PeerAddressBook { - fn default() -> Self { - Self::new() - } -} - -impl PeerAddressBook { - pub fn new() -> Self { - Self(HashMap::new()) - } - - /// Learn a node's address from a hint. - pub fn learn(&mut self, node_id: NodeId, addr: SocketAddr) { - self.0.insert(node_id, addr); - } - - /// Resolve a node's address. - pub fn resolve(&self, node_id: &NodeId) -> Option { - self.0.get(node_id).copied() - } - - /// Learn multiple hints at once. - pub fn bulk_learn(&mut self, hints: &[AddressHint]) { - for hint in hints { - self.learn(hint.node_id, hint.addr); - } - } -} - -// ─── NodeDriver ───────────────────────────────────────────────────────────── - -/// Network driver that owns a `DistributedNode` and performs real TCP I/O. -pub struct NodeDriver { - node: DistributedNode, - transport: TcpTransport, - acceptor: TcpAcceptor, - streams: Vec, - address_book: PeerAddressBook, - listen_addr: SocketAddr, -} - -impl NodeDriver { - /// Create a new driver. Binds a TCP listener on `listen_addr`. - pub fn new(listen_addr: SocketAddr, config: DistributedNodeConfig) -> Result { - let acceptor = TcpAcceptor::bind(listen_addr)?; - let actual_addr = acceptor.local_addr(); - let node = DistributedNode::new(config); - Ok(Self { - node, - transport: TcpTransport::pool(), - acceptor, - streams: Vec::new(), - address_book: PeerAddressBook::new(), - listen_addr: actual_addr, - }) - } - - /// Create a driver with a specific keypair for persistent identity. - pub fn with_keypair( - listen_addr: SocketAddr, - keypair: crate::crypto::Keypair, - config: DistributedNodeConfig, - ) -> Result { - let acceptor = TcpAcceptor::bind(listen_addr)?; - let actual_addr = acceptor.local_addr(); - let node = DistributedNode::with_keypair(keypair, config); - Ok(Self { - node, - transport: TcpTransport::pool(), - acceptor, - streams: Vec::new(), - address_book: PeerAddressBook::new(), - listen_addr: actual_addr, - }) - } - - /// The node's identity. - pub fn node_id(&self) -> NodeId { - self.node.node_id() - } - - /// The address this driver is listening on. - pub fn listen_addr(&self) -> SocketAddr { - self.listen_addr - } - - /// Access the underlying node (read-only). - pub fn node(&self) -> &DistributedNode { - &self.node - } - - /// Access the underlying node (mutable). - pub fn node_mut(&mut self) -> &mut DistributedNode { - &mut self.node - } - - /// Capture a snapshot of the node's state, enriched with addresses - /// from the driver's address book. - pub fn snapshot(&self) -> DistributionNodeSnapshot { - let mut snap = self.node.snapshot(); - snap.listen_addr = Some(self.listen_addr.to_string()); - - // Enrich member addresses from address book - for member in &mut snap.members { - if let Some(node_id) = parse_node_id_hex(&member.node_id) - && let Some(addr) = self.address_book.resolve(&node_id) { - member.addr = Some(addr.to_string()); - } - } - - // Enrich routing neighbor addresses from address book - for neighbor in &mut snap.routing_neighbors { - if let Some(node_id) = parse_node_id_hex(&neighbor.node_id) - && let Some(addr) = self.address_book.resolve(&node_id) { - neighbor.addr = Some(addr.to_string()); - } - } - - snap - } - - /// Join a cluster by contacting seed nodes. - /// - /// Sends `JoinRequest` messages directly to each seed over TCP. - /// On receiving a `JoinResponse`, extracts address hints and passes - /// the member records to the protocol layer. - pub fn join(&mut self, seeds: &[SocketAddr]) { - for seed_addr in seeds { - if let Err(e) = self.send_join_request(*seed_addr) { - eprintln!("driver: join send error to {seed_addr}: {e}"); - } - } - } - - fn send_join_request(&mut self, seed_addr: SocketAddr) -> Result<(), swactor::Error> { - let msg = JoinRequest { - from: self.node.node_id(), - }; - let hints = vec![AddressHint { - node_id: self.node.node_id(), - addr: self.listen_addr, - }]; - self.send_wire_with_hints::(&msg, seed_addr, &hints) - } - - /// Advance the node by one tick. - /// - /// Drives the SWIM probe cycle, sends outgoing protocol messages, - /// and handles periodic republishing. - pub fn tick(&mut self) { - let actions = self.node.tick(); - self.send_actions(&actions); - } - - /// Process incoming TCP messages. - /// - /// Reads all available wire envelopes from the acceptor, dispatches - /// each to the appropriate handler, and sends any response actions. - pub fn recv(&mut self) { - let envelopes = self.acceptor.try_recv(&mut self.streams); - for (envelope, _peer_addr, hints_bytes) in envelopes { - if !hints_bytes.is_empty() - && let Ok(hints) = serde_json::from_slice::>(&hints_bytes) { - self.learn_hints(&hints); - } - let response_actions = self.dispatch_incoming(envelope); - self.send_actions(&response_actions); - } - } - - /// Process incoming TCP messages with peer auth filtering. - /// - /// Same as `recv()`, but checks the sender's NodeId against the - /// peer allow-list before dispatching. Unauthorized messages are dropped. - pub fn recv_with_auth( - &mut self, - peer_auth: &std::sync::Arc>, - ) { - let envelopes = self.acceptor.try_recv(&mut self.streams); - for (envelope, _peer_addr, hints_bytes) in envelopes { - // Extract sender NodeId from hints - let mut sender_node_id = None; - if !hints_bytes.is_empty() - && let Ok(hints) = serde_json::from_slice::>(&hints_bytes) { - if let Some(first) = hints.first() { - sender_node_id = Some(first.node_id); - } - self.learn_hints(&hints); - } - - // Check peer auth if we know the sender - if let Some(node_id) = sender_node_id { - let allowed = peer_auth.lock().unwrap().is_allowed(&node_id); - if !allowed { - let hex: String = node_id.0[..4].iter().map(|b| format!("{b:02x}")).collect(); - eprintln!("driver: rejected message from unauthorized peer {hex}"); - continue; - } - } - - let response_actions = self.dispatch_incoming(envelope); - self.send_actions(&response_actions); - } - } - - // ─── Outgoing: NodeAction → TCP ───────────────────────────────────── - - fn send_actions(&mut self, actions: &[NodeAction]) { - for action in actions { - if let Err(e) = self.send_action(action) { - eprintln!("driver: send error: {e}"); - } - } - } - - fn send_action(&mut self, action: &NodeAction) -> Result<(), swactor::Error> { - // Standard sender hint - let sender_hint = AddressHint { - node_id: self.node.node_id(), - addr: self.listen_addr, - }; - - match action { - NodeAction::SendPing { - to, - sequence, - piggyback, - } => { - let dest = self.resolve_addr(to)?; - let msg = Ping { - from: self.node.node_id(), - sequence: *sequence, - piggyback: piggyback.clone(), - }; - self.send_wire_with_hints::(&msg, dest, &[sender_hint]) - } - - NodeAction::SendAck { - to, - sequence, - piggyback, - } => { - let dest = self.resolve_addr(to)?; - let msg = Ack { - from: self.node.node_id(), - sequence: *sequence, - piggyback: piggyback.clone(), - }; - self.send_wire_with_hints::(&msg, dest, &[sender_hint]) - } - - NodeAction::SendPingReq { - relay, - target, - sequence, - piggyback, - } => { - let dest = self.resolve_addr(relay)?; - let msg = PingReq { - from: self.node.node_id(), - target: *target, - sequence: *sequence, - piggyback: piggyback.clone(), - }; - // Include target hint so the relay can forward - let mut hints = vec![sender_hint]; - if let Some(target_addr) = self.address_book.resolve(target) { - hints.push(AddressHint { - node_id: *target, - addr: target_addr, - }); - } - self.send_wire_with_hints::(&msg, dest, &hints) - } - - NodeAction::SendJoinResponse { - to, members, - } => { - let dest = self.resolve_addr(to)?; - let msg = JoinResponse { - members: members.clone(), - }; - // Include all known member addresses as hints - let mut hints = vec![sender_hint]; - for record in members { - if let Some(addr) = self.address_book.resolve(&record.node_id) { - hints.push(AddressHint { - node_id: record.node_id, - addr, - }); - } - } - self.send_wire_with_hints::(&msg, dest, &hints) - } - - NodeAction::ForwardAck { - to, - target, - sequence, - piggyback, - } => { - let dest = self.resolve_addr(to)?; - let msg = IndirectAck { - target: *target, - sequence: *sequence, - piggyback: piggyback.clone(), - }; - self.send_wire_with_hints::(&msg, dest, &[sender_hint]) - } - - NodeAction::MembershipChanged { .. } => { - // Internal notification — no network I/O. - Ok(()) - } - } - } - - fn resolve_addr(&self, node_id: &NodeId) -> Result { - self.address_book - .resolve(node_id) - .ok_or_else(|| swactor::Error::from(format!( - "no address known for node {:?}", - node_id - ))) - } - - fn send_wire_with_hints( - &mut self, - msg: &M, - dest_addr: SocketAddr, - hints: &[AddressHint], - ) -> Result<(), swactor::Error> { - let payload = serde_json::to_vec(msg) - .map_err(|e| swactor::Error::from(format!("encode {}: {e}", M::type_tag())))?; - let hints_bytes = serde_json::to_vec(hints).unwrap_or_default(); - let envelope = WireEnvelope { - dest: SWIM_DEST, - type_tag: M::type_tag().to_string(), - payload, - }; - let buf = encode_wire_envelope_with_hints(&envelope, &hints_bytes); - let mut stream = self.transport_get_or_connect(dest_addr)?; - use std::io::Write; - match stream.write_all(&buf) { - Ok(()) => Ok(()), - Err(_) => { - self.transport.evict(dest_addr); - let mut stream = self.transport_get_or_connect(dest_addr)?; - stream - .write_all(&buf) - .map_err(|e| swactor::Error::from(format!("TCP send to {dest_addr}: {e}"))) - } - } - } - - fn transport_get_or_connect(&self, addr: SocketAddr) -> Result { - self.transport.get_or_connect(addr) - } - - // ─── Incoming: TCP → handler ──────────────────────────────────────── - - fn dispatch_incoming(&mut self, envelope: WireEnvelope) -> Vec { - match envelope.type_tag.as_str() { - "swactor_dist::Ping" => match decode::(&envelope.payload) { - Ok(msg) => { - self.node.handle_ping( - msg.from, - msg.sequence, - &msg.piggyback, - ) - } - Err(e) => { - eprintln!("driver: decode Ping: {e}"); - Vec::new() - } - }, - - "swactor_dist::Ack" => match decode::(&envelope.payload) { - Ok(msg) => self.node.handle_ack(msg.from, msg.sequence, &msg.piggyback), - Err(e) => { - eprintln!("driver: decode Ack: {e}"); - Vec::new() - } - }, - - "swactor_dist::PingReq" => match decode::(&envelope.payload) { - Ok(msg) => self.node.handle_ping_req( - msg.from, - msg.target, - msg.sequence, - &msg.piggyback, - ), - Err(e) => { - eprintln!("driver: decode PingReq: {e}"); - Vec::new() - } - }, - - "swactor_dist::JoinRequest" => match decode::(&envelope.payload) { - Ok(msg) => self.node.handle_join_request(msg.from), - Err(e) => { - eprintln!("driver: decode JoinRequest: {e}"); - Vec::new() - } - }, - - "swactor_dist::JoinResponse" => match decode::(&envelope.payload) { - Ok(msg) => self.node.handle_join_response(msg.members), - Err(e) => { - eprintln!("driver: decode JoinResponse: {e}"); - Vec::new() - } - }, - - "swactor_dist::IndirectAck" => match decode::(&envelope.payload) { - Ok(msg) => self.node.handle_indirect_ack(msg.target, msg.sequence, &msg.piggyback), - Err(e) => { - eprintln!("driver: decode IndirectAck: {e}"); - Vec::new() - } - }, - - other => { - eprintln!("driver: unknown message type: {other}"); - Vec::new() - } - } - } - - /// Process address hints extracted from incoming frames. - pub fn learn_hints(&mut self, hints: &[AddressHint]) { - self.address_book.bulk_learn(hints); - } -} - -fn decode(bytes: &[u8]) -> Result { - serde_json::from_slice(bytes).map_err(|e| e.to_string()) -} - -fn parse_node_id_hex(hex: &str) -> Option { - if hex.len() != 64 { - return None; - } - let mut bytes = [0u8; 32]; - for i in 0..32 { - bytes[i] = u8::from_str_radix(&hex[i * 2..i * 2 + 2], 16).ok()?; - } - Some(NodeId(bytes)) -} - diff --git a/crates/distribution/src/iroh_driver.rs b/crates/distribution/src/iroh_driver.rs index 14b7c7a..13eb41d 100644 --- a/crates/distribution/src/iroh_driver.rs +++ b/crates/distribution/src/iroh_driver.rs @@ -416,9 +416,24 @@ impl IrohDriver { // ─── Outgoing: NodeAction → iroh ───────────────────────────────── fn send_actions(&mut self, actions: &[NodeAction]) { + let mut failure_targets: Vec = Vec::new(); for action in actions { if let Err(e) = self.send_action(action) { eprintln!("iroh driver: send error: {e}"); + if let Some(target) = action_target(action) { + if !failure_targets.contains(&target) { + failure_targets.push(target); + } + } + } + } + for target in failure_targets { + let probe_actions = self.node.report_send_failure(target); + // Best-effort send of probe actions — no recursion on failure + for action in &probe_actions { + if let Err(e) = self.send_action(action) { + eprintln!("iroh driver: probe send error: {e}"); + } } } } @@ -765,6 +780,20 @@ async fn start_embedded_relay( Ok((server, url)) } +// ─── Helpers ───────────────────────────────────────────────────────────────── + +/// Extract the send target from a node action (if it has one). +fn action_target(action: &NodeAction) -> Option { + match action { + NodeAction::SendPing { to, .. } => Some(*to), + NodeAction::SendAck { to, .. } => Some(*to), + NodeAction::SendPingReq { relay, .. } => Some(*relay), + NodeAction::SendJoinResponse { to, .. } => Some(*to), + NodeAction::ForwardAck { to, .. } => Some(*to), + NodeAction::MembershipChanged { .. } => None, + } +} + // ─── Wire Framing Over QUIC Streams ───────────────────────────────────────── /// Write a tagged message to a QUIC send stream. diff --git a/crates/distribution/src/lib.rs b/crates/distribution/src/lib.rs index 5469829..a93901e 100644 --- a/crates/distribution/src/lib.rs +++ b/crates/distribution/src/lib.rs @@ -10,7 +10,5 @@ pub mod node; pub mod registry; pub mod node_metadata; pub mod snapshot; -#[cfg(feature = "tcp")] -pub mod driver; #[cfg(feature = "iroh")] pub mod iroh_driver; diff --git a/crates/distribution/src/node.rs b/crates/distribution/src/node.rs index aa1b1d5..58758db 100644 --- a/crates/distribution/src/node.rs +++ b/crates/distribution/src/node.rs @@ -191,6 +191,13 @@ impl DistributedNode { self.inject_piggyback(actions) } + /// Report that a send to `target` failed, triggering a reactive probe. + pub fn report_send_failure(&mut self, target: NodeId) -> Vec { + let actions = self.swim.report_send_failure(target); + self.process_membership_changes(&actions); + self.inject_piggyback(actions) + } + pub fn handle_join_request(&mut self, from: NodeId) -> Vec { let actions = self.swim.handle_join_request(from); self.maybe_update_routing_table(from); @@ -286,8 +293,16 @@ impl DistributedNode { /// Set this node's relay URL and begin gossiping it to the cluster. pub fn set_relay_url(&mut self, url: Option) { + let name = self.metadata.node_name(&self.node_id()).map(String::from); self.metadata - .set_local(self.node_id(), url, self.cluster_size()); + .set_local(self.node_id(), url, name, self.cluster_size()); + } + + /// Set this node's human-readable name and begin gossiping it to the cluster. + pub fn set_node_name(&mut self, name: String) { + let relay_url = self.metadata.relay_url(&self.node_id()).map(String::from); + self.metadata + .set_local(self.node_id(), relay_url, Some(name), self.cluster_size()); } /// Look up a node's relay URL. @@ -295,6 +310,11 @@ impl DistributedNode { self.metadata.relay_url(node_id) } + /// Look up a node's human-readable name. + pub fn node_name(&self, node_id: &NodeId) -> Option<&str> { + self.metadata.node_name(node_id) + } + /// Read-only access to the metadata disseminator. pub fn metadata(&self) -> &NodeMetadataDisseminator { &self.metadata diff --git a/crates/distribution/src/node_metadata.rs b/crates/distribution/src/node_metadata.rs index e4b4dbb..080a8af 100644 --- a/crates/distribution/src/node_metadata.rs +++ b/crates/distribution/src/node_metadata.rs @@ -15,6 +15,8 @@ use crate::types::NodeId; pub struct NodeMetadataEntry { pub node_id: NodeId, pub relay_url: Option, + #[serde(default)] + pub node_name: Option, pub generation: u64, } @@ -43,12 +45,19 @@ impl NodeMetadataDisseminator { } } - /// Set this node's relay URL and enqueue for dissemination. - pub fn set_local(&mut self, node_id: NodeId, relay_url: Option, cluster_size: usize) { + /// Set this node's metadata and enqueue for dissemination. + pub fn set_local( + &mut self, + node_id: NodeId, + relay_url: Option, + node_name: Option, + cluster_size: usize, + ) { self.local_generation += 1; let entry = NodeMetadataEntry { node_id, relay_url, + node_name, generation: self.local_generation, }; self.store.insert(node_id, entry.clone()); @@ -91,6 +100,13 @@ impl NodeMetadataDisseminator { .and_then(|e| e.relay_url.as_deref()) } + /// Look up a node's human-readable name. + pub fn node_name(&self, node_id: &NodeId) -> Option<&str> { + self.store + .get(node_id) + .and_then(|e| e.node_name.as_deref()) + } + /// Remove metadata for a dead node. pub fn remove_node(&mut self, node_id: &NodeId) { self.store.remove(node_id); diff --git a/crates/distribution/src/snapshot.rs b/crates/distribution/src/snapshot.rs index d4b3d01..9ce985a 100644 --- a/crates/distribution/src/snapshot.rs +++ b/crates/distribution/src/snapshot.rs @@ -21,6 +21,8 @@ pub struct MemberInfo { pub label: Option, #[serde(skip_serializing_if = "Option::is_none")] pub relay_url: Option, + #[serde(skip_serializing_if = "Option::is_none", default)] + pub node_name: Option, } /// Snapshot of a node in the Kademlia routing table. @@ -115,6 +117,10 @@ pub struct DistributionNodeSnapshot { /// This node's relay URL, if running an embedded relay server. #[serde(skip_serializing_if = "Option::is_none", default)] pub relay_url: Option, + + /// Build version string (e.g. "branch @ hash"). + #[serde(skip_serializing_if = "Option::is_none", default)] + pub version: Option, } fn node_id_hex(id: &NodeId) -> String { @@ -146,6 +152,7 @@ impl DistributedNode { is_authorized: None, label: None, relay_url: self.metadata().relay_url(&m.node_id).map(String::from), + node_name: self.metadata().node_name(&m.node_id).map(String::from), }) .collect(); @@ -210,9 +217,10 @@ impl DistributedNode { recent_probe_targets: recent_targets, peer_auth_mode: "open".into(), authorized_peer_count: None, - node_name: None, + node_name: self.metadata().node_name(&self.node_id()).map(String::from), invite_code: None, relay_url: self.metadata().relay_url(&self.node_id()).map(String::from), + version: None, } } } diff --git a/crates/distribution/src/swim/node.rs b/crates/distribution/src/swim/node.rs index f8baade..d1d7830 100644 --- a/crates/distribution/src/swim/node.rs +++ b/crates/distribution/src/swim/node.rs @@ -142,6 +142,15 @@ impl SwimNode { actions } + /// Report that a send to `target` failed, triggering a reactive probe. + pub fn report_send_failure(&mut self, target: NodeId) -> Vec { + let probe_actions = self.probe.step( + SwimEvent::SendFailed { to: target }, + &mut self.members, + ); + self.translate_probe_actions(probe_actions) + } + /// Handle a received indirect ack (forwarded by a relay node). pub fn handle_indirect_ack(&mut self, target: NodeId, sequence: u64, piggyback: &[u8]) -> Vec { let mut actions = self.apply_piggyback(piggyback); @@ -202,6 +211,11 @@ impl SwimNode { state: record.state, incarnation: record.incarnation, }); + // In reactive mode, probe newly discovered alive peers so they + // don't decay to dead before we ever exchange a ping/ack. + if record.state == MemberState::Alive && record.node_id != self.members.self_id() { + self.probe.enqueue_demand_probe(record.node_id); + } } } actions @@ -259,6 +273,9 @@ impl SwimNode { if changed { if update.state == MemberState::Alive { eprintln!("SWIM: alive {}", &hex_encode(&update.node_id.0)[..8]); + // In reactive mode, probe newly discovered alive peers so they + // don't decay to dead before we ever exchange a ping/ack. + self.probe.enqueue_demand_probe(update.node_id); } // Re-disseminate the update self.dissemination.enqueue( diff --git a/crates/distribution/src/swim/probe.rs b/crates/distribution/src/swim/probe.rs index 954ce22..e7334e5 100644 --- a/crates/distribution/src/swim/probe.rs +++ b/crates/distribution/src/swim/probe.rs @@ -14,10 +14,20 @@ const PROBE_HISTORY_SIZE: usize = 16; // ─── Configuration ────────────────────────────────────────────────────────── +/// Controls how the probe cycle triggers probes. +#[derive(Debug, Clone, PartialEq)] +pub enum ProbeMode { + /// Classic SWIM: probe one random member every `probe_interval` ticks. + Periodic, + /// No periodic probing. Probes triggered externally via `SendFailed`. + /// Safety sweep probes one random member every `safety_sweep_interval` ticks. + Reactive { safety_sweep_interval: u64 }, +} + /// SWIM protocol configuration. #[derive(Debug, Clone)] pub struct SwimConfig { - /// Ticks between probe cycles. + /// Ticks between probe cycles (used in Periodic mode). pub probe_interval: u64, /// Ticks to wait for a direct ack before sending indirect probes. pub probe_timeout: u64, @@ -28,6 +38,8 @@ pub struct SwimConfig { /// Ticks between dead-node reprobe attempts. 0 = disabled. /// When enabled, periodically pings dead nodes to detect partition heals. pub dead_reprobe_interval: u64, + /// Probe mode: Periodic (default) or Reactive (probe-on-failure). + pub probe_mode: ProbeMode, } impl Default for SwimConfig { @@ -38,6 +50,7 @@ impl Default for SwimConfig { indirect_probes: 3, suspicion_timeout: 30, dead_reprobe_interval: 50, + probe_mode: ProbeMode::Periodic, } } } @@ -53,6 +66,8 @@ pub enum SwimEvent { AckReceived { from: NodeId, sequence: u64 }, /// Received an indirect ack (forwarded through a relay). IndirectAckReceived { target: NodeId, sequence: u64 }, + /// A send to the given peer failed (reactive probe trigger). + SendFailed { to: NodeId }, } // ─── Actions (outputs) ────────────────────────────────────────────────────── @@ -103,6 +118,9 @@ struct SuspicionTimer { started_at: u64, } +/// Maximum demand queue size to prevent unbounded growth. +const MAX_DEMAND_QUEUE: usize = 32; + /// The SWIM probe state machine. pub struct SwimProbe { config: SwimConfig, @@ -122,6 +140,10 @@ pub struct SwimProbe { next_reprobe_tick: u64, /// Round-robin index into the dead member list for reprobe target selection. reprobe_index: usize, + /// Peers needing probes due to send failures (reactive mode). + demand_queue: VecDeque, + /// Tick at which the next safety sweep fires (reactive mode). + next_sweep_tick: u64, } impl SwimProbe { @@ -131,6 +153,10 @@ impl SwimProbe { } else { u64::MAX }; + let next_sweep = match &config.probe_mode { + ProbeMode::Reactive { safety_sweep_interval } => *safety_sweep_interval, + ProbeMode::Periodic => u64::MAX, + }; Self { next_probe_tick: config.probe_interval, next_reprobe_tick: next_reprobe, @@ -143,6 +169,8 @@ impl SwimProbe { probe_order: Vec::new(), suspicion_timers: Vec::new(), recent_targets: VecDeque::with_capacity(PROBE_HISTORY_SIZE), + demand_queue: VecDeque::new(), + next_sweep_tick: next_sweep, } } @@ -155,7 +183,15 @@ impl SwimProbe { self.tick += 1; self.check_probe_timeout(members, &mut actions); self.check_suspicion_timeouts(members, &mut actions); - self.maybe_start_probe(members, &mut actions); + match &self.config.probe_mode { + ProbeMode::Periodic => { + self.maybe_start_probe(members, &mut actions); + } + ProbeMode::Reactive { .. } => { + self.maybe_start_demand_probe(members, &mut actions); + self.maybe_safety_sweep(members, &mut actions); + } + } self.maybe_reprobe_dead(members, &mut actions); } SwimEvent::AckReceived { from, sequence } => { @@ -164,6 +200,9 @@ impl SwimProbe { SwimEvent::IndirectAckReceived { target, sequence } => { self.handle_indirect_ack(target, sequence, members, &mut actions); } + SwimEvent::SendFailed { to } => { + self.handle_send_failed(to, members, &mut actions); + } } actions @@ -390,6 +429,115 @@ impl SwimProbe { sequence: seq, }); } + + // ─── Reactive mode ───────────────────────────────────────────────── + + /// Enqueue a demand probe for a newly discovered peer (reactive mode only). + /// + /// Called when gossip or a join response introduces a new Alive member. + /// In Periodic mode this is a no-op (periodic probing covers it). + pub fn enqueue_demand_probe(&mut self, target: NodeId) { + if matches!(self.config.probe_mode, ProbeMode::Periodic) { + return; + } + if self.is_currently_probing(target) || self.demand_queue.contains(&target) { + return; + } + if self.demand_queue.len() < MAX_DEMAND_QUEUE { + self.demand_queue.push_back(target); + } + } + + /// Handle a send failure: start a probe immediately or queue it. + fn handle_send_failed(&mut self, target: NodeId, members: &MemberList, actions: &mut Vec) { + // Ignore failures for dead peers, self, or already-queued targets + if let Some(entry) = members.get(&target) { + if entry.state == MemberState::Dead { + return; + } + } else { + // Unknown peer — nothing to probe + return; + } + + // Ignore if we're already probing this target + if self.is_currently_probing(target) { + return; + } + + // Ignore if already in demand queue + if self.demand_queue.contains(&target) { + return; + } + + if matches!(self.phase, ProbePhase::Idle) { + // Start probe immediately + self.start_probe_for(target, actions); + } else { + // Queue it (capped) + if self.demand_queue.len() < MAX_DEMAND_QUEUE { + self.demand_queue.push_back(target); + } + } + } + + /// Start a directed probe to a specific target. + fn start_probe_for(&mut self, target: NodeId, actions: &mut Vec) { + if self.recent_targets.len() >= PROBE_HISTORY_SIZE { + self.recent_targets.pop_front(); + } + self.recent_targets.push_back(target); + + let seq = self.next_sequence(); + actions.push(SwimAction::SendPing { + to: target, + sequence: seq, + }); + self.phase = ProbePhase::WaitingDirectAck { + target, + sequence: seq, + sent_at: self.tick, + }; + } + + /// On each tick in reactive mode, if idle and queue non-empty, pop and probe. + fn maybe_start_demand_probe(&mut self, _members: &MemberList, actions: &mut Vec) { + if !matches!(self.phase, ProbePhase::Idle) { + return; + } + if let Some(target) = self.demand_queue.pop_front() { + self.start_probe_for(target, actions); + } + } + + /// At safety_sweep_interval, probe one random alive member. + fn maybe_safety_sweep(&mut self, members: &MemberList, actions: &mut Vec) { + if self.tick < self.next_sweep_tick { + return; + } + let interval = match &self.config.probe_mode { + ProbeMode::Reactive { safety_sweep_interval } => *safety_sweep_interval, + ProbeMode::Periodic => return, + }; + self.next_sweep_tick = self.tick + interval; + + if !matches!(self.phase, ProbePhase::Idle) { + return; + } + + if let Some(target) = self.pick_probe_target(&MemberList::clone_shallow(members)) { + self.start_probe_for(target, actions); + } + } + + /// Check if we are currently probing a specific target. + fn is_currently_probing(&self, target: NodeId) -> bool { + match &self.phase { + ProbePhase::WaitingDirectAck { target: t, .. } + | ProbePhase::WaitingIndirectAck { target: t, .. } => *t == target, + ProbePhase::Idle => false, + } + } } // Helper: we need a read-only borrow of members in pick_probe_target diff --git a/crates/distribution/tests/common/mod.rs b/crates/distribution/tests/common/mod.rs index 7e03906..067ba4f 100644 --- a/crates/distribution/tests/common/mod.rs +++ b/crates/distribution/tests/common/mod.rs @@ -21,6 +21,7 @@ pub fn test_config() -> DistributedNodeConfig { indirect_probes: 1, suspicion_timeout: 5, dead_reprobe_interval: 0, + ..SwimConfig::default() }, cache_capacity: 100, republish_interval: 50, diff --git a/crates/distribution/tests/swim_node.rs b/crates/distribution/tests/swim_node.rs index 23bb2cc..e7a9c76 100644 --- a/crates/distribution/tests/swim_node.rs +++ b/crates/distribution/tests/swim_node.rs @@ -13,6 +13,7 @@ fn fast_config() -> SwimConfig { indirect_probes: 2, suspicion_timeout: 10, dead_reprobe_interval: 0, + ..SwimConfig::default() } } diff --git a/crates/distribution/tests/swim_probe.rs b/crates/distribution/tests/swim_probe.rs index ad8fb08..ac0ccd5 100644 --- a/crates/distribution/tests/swim_probe.rs +++ b/crates/distribution/tests/swim_probe.rs @@ -95,6 +95,7 @@ fn probe_sends_ping_after_interval() { indirect_probes: 2, suspicion_timeout: 20, dead_reprobe_interval: 0, + ..SwimConfig::default() }; let mut probe = SwimProbe::new(config); let mut members = MemberList::new(node(0)); @@ -118,6 +119,7 @@ fn probe_ack_completes_cycle() { indirect_probes: 2, suspicion_timeout: 20, dead_reprobe_interval: 0, + ..SwimConfig::default() }; let mut probe = SwimProbe::new(config); let mut members = MemberList::new(node(0)); @@ -148,6 +150,7 @@ fn probe_timeout_triggers_indirect_probes() { indirect_probes: 2, suspicion_timeout: 20, dead_reprobe_interval: 0, + ..SwimConfig::default() }; let mut probe = SwimProbe::new(config); let mut members = MemberList::new(node(0)); @@ -175,6 +178,7 @@ fn no_ack_at_all_causes_suspicion() { indirect_probes: 2, suspicion_timeout: 20, dead_reprobe_interval: 0, + ..SwimConfig::default() }; let mut probe = SwimProbe::new(config); let mut members = MemberList::new(node(0)); @@ -204,6 +208,7 @@ fn suspicion_timeout_causes_death_declaration() { indirect_probes: 0, suspicion_timeout: 10, dead_reprobe_interval: 0, + ..SwimConfig::default() }; let mut probe = SwimProbe::new(config); let mut members = MemberList::new(node(0)); @@ -238,6 +243,7 @@ fn indirect_ack_rescues_suspected_node() { indirect_probes: 2, suspicion_timeout: 20, dead_reprobe_interval: 0, + ..SwimConfig::default() }; let mut probe = SwimProbe::new(config); let mut members = MemberList::new(node(0)); @@ -289,6 +295,7 @@ fn reprobe_sends_ping_to_dead_node() { indirect_probes: 0, suspicion_timeout: 10, dead_reprobe_interval: 20, + ..SwimConfig::default() }; let mut probe = SwimProbe::new(config); let mut members = MemberList::new(node(0)); @@ -314,6 +321,7 @@ fn reprobe_disabled_when_interval_is_zero() { indirect_probes: 0, suspicion_timeout: 10, dead_reprobe_interval: 0, + ..SwimConfig::default() }; let mut probe = SwimProbe::new(config); let mut members = MemberList::new(node(0)); @@ -336,6 +344,7 @@ fn reprobe_does_nothing_when_no_dead_members() { indirect_probes: 0, suspicion_timeout: 10, dead_reprobe_interval: 20, + ..SwimConfig::default() }; let mut probe = SwimProbe::new(config); let mut members = MemberList::new(node(0)); diff --git a/crates/distribution/tests/transport_and_codec.rs b/crates/distribution/tests/transport_and_codec.rs index 48f9a39..2c42bd1 100644 --- a/crates/distribution/tests/transport_and_codec.rs +++ b/crates/distribution/tests/transport_and_codec.rs @@ -74,175 +74,3 @@ fn all_message_types_registered_in_codec_registry() { } } -// ─── TCP transport tests (require "tcp" feature) ──────────────────────────── - -#[cfg(feature = "tcp")] -mod tcp_transport { - use swactor::actor::ActorAddress; - use swactor::transport::WireEnvelope; - - use distribution::codec::distribution_codec_registry; - use distribution::messages::*; - use swactor_transport::tcp::{TcpAcceptor, TcpTransport}; - use distribution::types::NodeId; - - #[test] - fn wire_envelope_roundtrips_through_tcp() { - let acceptor = TcpAcceptor::bind("127.0.0.1:0".parse().unwrap()).unwrap(); - let addr = acceptor.local_addr(); - - let original = WireEnvelope { - dest: ActorAddress::new_random(), - type_tag: "test::Msg".to_string(), - payload: vec![1, 2, 3, 4, 5], - }; - - let original_clone = original.clone(); - let sender = std::thread::spawn(move || { - let transport = TcpTransport::new(addr); - transport.send_to(addr, original_clone).unwrap(); - }); - - std::thread::sleep(std::time::Duration::from_millis(50)); - let mut streams = Vec::new(); - let envelopes = loop { - let envs = acceptor.try_recv(&mut streams); - if !envs.is_empty() { - break envs; - } - std::thread::sleep(std::time::Duration::from_millis(10)); - }; - - sender.join().unwrap(); - - assert_eq!(envelopes.len(), 1); - let (received, _peer, hints) = &envelopes[0]; - assert_eq!(received.dest, original.dest); - assert_eq!(received.type_tag, original.type_tag); - assert_eq!(received.payload, original.payload); - assert!(hints.is_empty(), "no hints expected from plain send_to"); - } - - #[test] - fn wire_envelope_minimal_roundtrips() { - let acceptor = TcpAcceptor::bind("127.0.0.1:0".parse().unwrap()).unwrap(); - let addr = acceptor.local_addr(); - - let original = WireEnvelope { - dest: ActorAddress::new_random(), - type_tag: "test::Minimal".to_string(), - payload: vec![42], - }; - - let original_clone = original.clone(); - let sender = std::thread::spawn(move || { - let transport = TcpTransport::new(addr); - transport.send_to(addr, original_clone).unwrap(); - }); - - std::thread::sleep(std::time::Duration::from_millis(50)); - let mut streams = Vec::new(); - let envelopes = loop { - let envs = acceptor.try_recv(&mut streams); - if !envs.is_empty() { - break envs; - } - std::thread::sleep(std::time::Duration::from_millis(10)); - }; - - sender.join().unwrap(); - - let (received, _, _hints) = &envelopes[0]; - assert_eq!(received.payload, vec![42]); - } - - #[test] - fn two_drivers_complete_join_handshake() { - use distribution::driver::NodeDriver; - use distribution::node::DistributedNodeConfig; - - let mut driver_a = NodeDriver::new( - "127.0.0.1:0".parse().unwrap(), - DistributedNodeConfig::default(), - ) - .unwrap(); - - let mut driver_b = NodeDriver::new( - "127.0.0.1:0".parse().unwrap(), - DistributedNodeConfig::default(), - ) - .unwrap(); - - let b_addr = driver_b.listen_addr(); - - // A sends JoinRequest to B - driver_a.join(&[b_addr]); - - // B receives the JoinRequest, learns A's address from hints, sends JoinResponse - std::thread::sleep(std::time::Duration::from_millis(50)); - driver_b.recv(); - - // A receives the JoinResponse with member list - std::thread::sleep(std::time::Duration::from_millis(50)); - driver_a.recv(); - - // A should now have members — proving the full roundtrip worked. - // If hints were broken, B couldn't resolve A's address to send the - // JoinResponse, so A's member list would stay empty. - let members = driver_a.node().members(); - assert!( - !members.is_empty(), - "join handshake should complete: A needs members from JoinResponse" - ); - } - - #[test] - fn ping_message_survives_codec_and_tcp_roundtrip() { - let codecs = distribution_codec_registry(); - - let acceptor = TcpAcceptor::bind("127.0.0.1:0".parse().unwrap()).unwrap(); - let server_addr = acceptor.local_addr(); - - let dest = ActorAddress::new_random(); - let ping = Ping { - from: NodeId([0xBB; 32]), - sequence: 99, - piggyback: vec![], - }; - - let type_id = std::any::TypeId::of::(); - let (tag, payload) = codecs.encode(type_id, Box::new(ping.clone())).unwrap(); - - let envelope = WireEnvelope { - dest, - type_tag: tag, - payload, - }; - - let envelope_clone = envelope.clone(); - let sender = std::thread::spawn(move || { - let transport = TcpTransport::new(server_addr); - transport.send_to(server_addr, envelope_clone).unwrap(); - }); - - std::thread::sleep(std::time::Duration::from_millis(50)); - let mut streams = Vec::new(); - let envelopes = loop { - let envs = acceptor.try_recv(&mut streams); - if !envs.is_empty() { - break envs; - } - std::thread::sleep(std::time::Duration::from_millis(10)); - }; - - sender.join().unwrap(); - - let (received, _, _hints) = &envelopes[0]; - - let (addr, msg_any) = codecs.receive(received.clone()).unwrap(); - assert_eq!(addr, dest); - let decoded: &Ping = msg_any.downcast_ref().unwrap(); - assert_eq!(decoded.from, NodeId([0xBB; 32])); - assert_eq!(decoded.sequence, 99); - } -} diff --git a/crates/node/Cargo.toml b/crates/node/Cargo.toml index 1fc6a68..03539d7 100644 --- a/crates/node/Cargo.toml +++ b/crates/node/Cargo.toml @@ -23,7 +23,6 @@ path = "src/lib.rs" [features] default = ["iroh", "relay"] -tcp = ["distribution/tcp"] iroh = ["distribution/iroh", "dep:iroh"] relay = ["iroh", "distribution/relay"] diff --git a/crates/node/src/main.rs b/crates/node/src/main.rs index e816508..602db31 100644 --- a/crates/node/src/main.rs +++ b/crates/node/src/main.rs @@ -73,19 +73,7 @@ struct Args { #[arg(long)] config: Option, - /// Transport to use: iroh or tcp - #[arg(long, default_value = "iroh")] - transport: String, - - /// Address to listen on for TCP transport (e.g. 10.0.1.10:7000) - #[arg(long)] - listen: Option, - - /// Seed node address to join (TCP mode: host:port) - #[arg(long)] - seed: Option, - - /// Seed node's iroh public key (iroh mode: hex-encoded 32-byte key) + /// Seed node's iroh public key (hex-encoded 32-byte key) #[arg(long)] seed_node_id: Option, @@ -272,11 +260,6 @@ fn main() { // Layer: CLI > config > defaults // For args with default values, we check if the user explicitly provided // the CLI flag; if not, we fall back to config, then to the default. - let transport = if args.transport != "iroh" { - args.transport.clone() - } else { - cfg.transport.unwrap_or_else(|| args.transport.clone()) - }; let dashboard_port = if args.dashboard_port != 9090 { args.dashboard_port } else { @@ -295,8 +278,6 @@ fn main() { let storage_path = args.storage_path.clone().or(cfg.storage_path); let peers_file = args.peers_file.clone().or(cfg.peers_file); let seed_node_id = args.seed_node_id.clone().or(cfg.seed_node_id); - #[cfg(feature = "tcp")] - let seed = args.seed.clone().or(cfg.seed); let auth_enabled = args.auth || cfg.auth.unwrap_or(false); let actors = if args.actors != 0 { args.actors @@ -331,11 +312,6 @@ fn main() { cfg.relay_bind.unwrap_or_else(|| args.relay_bind.clone()) }; let relay_hosts = cfg.relay_hosts.unwrap_or_default(); - #[cfg(feature = "tcp")] - let listen = args.listen.or_else(|| { - cfg.listen.as_ref().and_then(|s| s.parse().ok()) - }); - // Signal handler — second Ctrl+C forces immediate exit { let stop = Arc::clone(&stop); @@ -529,13 +505,16 @@ fn main() { }; // Distribution config - eprintln!("Distribution: SWIM (transport: {transport})"); + eprintln!("Distribution: SWIM (transport: iroh, mode: reactive)"); let swim_config = SwimConfig { - probe_interval: 5, - probe_timeout: 6, // 600ms — allows relay round-trip + probe_interval: 10, // unused in Reactive mode, kept for compat + probe_timeout: 15, // 1.5s — generous for relay roundtrips indirect_probes: 2, - suspicion_timeout: 40, // 4s — gives refutation time to piggyback - dead_reprobe_interval: 50, + suspicion_timeout: 80, // 8s — gives refutation time to gossip back + dead_reprobe_interval: 100, + probe_mode: distribution::swim::probe::ProbeMode::Reactive { + safety_sweep_interval: 3000, // 5 minutes at 100ms/tick + }, }; let node_config = DistributedNodeConfig { swim: swim_config, @@ -544,8 +523,9 @@ fn main() { ..Default::default() }; - // Channel for triggering SWIM joins (fed by peers plugin "Add Peer") + // Channel for triggering SWIM joins (fed by peers plugin "Add Peer" and distribution "Re-peer") let (join_tx, join_rx) = std::sync::mpsc::channel::(); + let join_tx_dist = join_tx.clone(); // Register peers plugin let peers_plugin = plugins::peers::PeersPlugin::new( @@ -554,53 +534,32 @@ fn main() { ); dash.register_plugin(Arc::new(peers_plugin)); - match transport.as_str() { - #[cfg(feature = "iroh")] - "iroh" => run_iroh( - seed_node_id, - dashboard_port, - actors, - node_config, - keypair, - Arc::clone(&peer_auth), - &handle, - &dash, - &stop, - &ds_group, - node_name, - invite_code, - join_rx, - relay_enabled, - &relay_bind, - relay_port, - relay_hosts, - ), - #[cfg(feature = "tcp")] - "tcp" => run_tcp( - listen, - seed, - dashboard_port, - actors, - node_config, - keypair, - Arc::clone(&peer_auth), - &handle, - &dash, - &stop, - &ds_group, - node_name, - invite_code, - join_rx, - ), - other => { - eprintln!("Unknown or unavailable transport: {other}"); - eprintln!("Available transports:"); - #[cfg(feature = "iroh")] - eprintln!(" iroh"); - #[cfg(feature = "tcp")] - eprintln!(" tcp"); - std::process::exit(1); - } + #[cfg(feature = "iroh")] + run_iroh( + seed_node_id, + dashboard_port, + actors, + node_config, + keypair, + Arc::clone(&peer_auth), + &handle, + &dash, + &stop, + &ds_group, + node_name, + invite_code, + join_rx, + join_tx_dist, + relay_enabled, + &relay_bind, + relay_port, + relay_hosts, + ); + + #[cfg(not(feature = "iroh"))] + { + eprintln!("iroh feature is required but not enabled"); + std::process::exit(1); } handle.shutdown(); @@ -608,89 +567,6 @@ fn main() { handle.join(); } -// ── TCP transport ──────────────────────────────────────────────────────── - -#[cfg(feature = "tcp")] -fn run_tcp( - listen: Option, - seed: Option, - dashboard_port: u16, - actors: usize, - node_config: DistributedNodeConfig, - keypair: Keypair, - peer_auth: Arc>, - handle: &swactor::runtime::RuntimeHandle, - dash: &dashboard::DashboardHandle, - stop: &Arc, - ds_group: &Option, - node_name: String, - invite_code: String, - _join_rx: std::sync::mpsc::Receiver, -) { - use distribution::driver::NodeDriver; - - let listen_addr = listen.expect("--listen is required for TCP mode"); - let dist_keypair = distribution::crypto::Keypair::from_bytes(&keypair.secret_bytes()); - let mut driver = NodeDriver::with_keypair(listen_addr, dist_keypair, node_config) - .expect("failed to create node driver"); - - eprintln!( - "Node {} listening on {} (TCP)", - hex(&driver.node_id().0[..4]), - driver.listen_addr(), - ); - - // Join seed if provided - if let Some(seed) = seed { - let seed_addr: std::net::SocketAddr = seed.parse().expect("invalid seed address"); - eprintln!("Joining cluster via seed {seed_addr}"); - driver.join(&[seed_addr]); - } - - // Spawn and register actors - let actor_addrs = spawn_actors(actors, handle, driver.node_mut()); - - // Wire distribution snapshot to dashboard via plugin - let mut snap = driver.snapshot(); - snap.node_name = Some(node_name.clone()); - snap.invite_code = Some(invite_code.clone()); - let cached_snapshot: Arc>> = - Arc::new(Mutex::new(Some(snap))); - let dist_plugin = plugins::distribution::DistributionPlugin::new( - Arc::clone(&cached_snapshot), - ); - dash.register_plugin(Arc::new(dist_plugin)); - - // Start dashboard HTTP on a standalone tokio runtime (no iroh runtime in TCP mode) - dash.start_http_standalone(); - eprintln!("Dashboard at http://0.0.0.0:{dashboard_port}"); - - // Main loop - let mut round: u64 = 0; - while !stop.load(Ordering::Relaxed) { - round += 1; - - driver.recv_with_auth(&peer_auth); - driver.tick(); - - for addr in &actor_addrs { - let _ = handle.runtime.send_to(*addr, Heartbeat); - } - - let mut snap = driver.snapshot(); - snap.node_name = Some(node_name.clone()); - snap.invite_code = Some(invite_code.clone()); - *cached_snapshot.lock().unwrap() = Some(snap); - - // Datastore ticks - if let Some(group) = ds_group { - group.tick(round); - } - - thread::sleep(Duration::from_millis(100)); - } -} - // ── iroh transport ─────────────────────────────────────────────────────── #[cfg(feature = "iroh")] @@ -708,6 +584,7 @@ fn run_iroh( node_name: String, invite_code: String, join_rx: std::sync::mpsc::Receiver, + join_tx_dist: std::sync::mpsc::Sender, relay_enabled: bool, relay_bind: &str, relay_port: u16, @@ -828,7 +705,8 @@ fn run_iroh( // Spawn and register actors let actor_addrs = spawn_actors(actors, handle, driver.node_mut()); - // Announce relay URL to cluster gossip + // Announce node name and relay URL to cluster gossip + driver.node_mut().set_node_name(node_name.clone()); let mut home_relay_set = if let Some(url) = driver.relay_url().map(|u| u.to_string()) { driver.node_mut().set_relay_url(Some(url)); true @@ -840,10 +718,12 @@ fn run_iroh( let mut snap = driver.snapshot(); snap.node_name = Some(node_name.clone()); snap.invite_code = Some(invite_code.clone()); + snap.version = Some(VERSION.to_string()); let cached_snapshot: Arc>> = Arc::new(Mutex::new(Some(snap))); let dist_plugin = plugins::distribution::DistributionPlugin::new( Arc::clone(&cached_snapshot), + Some(join_tx_dist), ); dash.register_plugin(Arc::new(dist_plugin)); @@ -930,6 +810,7 @@ fn run_iroh( let mut snap = driver.snapshot(); snap.node_name = Some(node_name.clone()); snap.invite_code = Some(invite_code.clone()); + snap.version = Some(VERSION.to_string()); *cached_snapshot.lock().unwrap() = Some(snap); // Datastore ticks diff --git a/crates/node/src/plugins/datastore_page.html b/crates/node/src/plugins/datastore_page.html index 4671a62..96cae68 100644 --- a/crates/node/src/plugins/datastore_page.html +++ b/crates/node/src/plugins/datastore_page.html @@ -543,6 +543,7 @@ // ── SSE connection ─────────────────────────────────────────────────── var es = new EventSource('/events'); + window.addEventListener('beforeunload', function() { es.close(); }); es.addEventListener('datastore', function(e) { try { diff --git a/crates/node/src/plugins/distribution.rs b/crates/node/src/plugins/distribution.rs index c21ef67..7adfea7 100644 --- a/crates/node/src/plugins/distribution.rs +++ b/crates/node/src/plugins/distribution.rs @@ -7,6 +7,7 @@ use std::collections::HashMap; use std::sync::{Arc, Mutex}; use dashboard::plugin::{DashboardPlugin, PluginResponse}; +use dashboard::JoinPeerInfo; use distribution::snapshot::DistributionNodeSnapshot; /// HTML page for the distribution plugin. @@ -15,11 +16,15 @@ const DISTRIBUTION_HTML: &str = include_str!("distribution_page.html"); /// Dashboard plugin that exposes distribution node snapshots. pub struct DistributionPlugin { cached: Arc>>, + join_sender: Option>, } impl DistributionPlugin { - pub fn new(cached: Arc>>) -> Self { - Self { cached } + pub fn new( + cached: Arc>>, + join_sender: Option>, + ) -> Self { + Self { cached, join_sender } } } @@ -39,7 +44,7 @@ impl DashboardPlugin for DistributionPlugin { method: &str, path: &str, _query: &HashMap, - _body: &[u8], + body: &[u8], ) -> PluginResponse { match (method, path) { ("GET", "") => { @@ -52,6 +57,7 @@ impl DashboardPlugin for DistributionPlugin { None => PluginResponse::json("{}".into()), } } + ("POST", "rejoin") => self.handle_rejoin(body), _ => PluginResponse::not_found(), } } @@ -60,3 +66,54 @@ impl DashboardPlugin for DistributionPlugin { Some(DISTRIBUTION_HTML) } } + +impl DistributionPlugin { + fn handle_rejoin(&self, body: &[u8]) -> PluginResponse { + let tx = match &self.join_sender { + Some(tx) => tx, + None => return PluginResponse::json(r#"{"error":"rejoin not available"}"#.into()), + }; + + // Parse { "node_id": "" } from body + let parsed: serde_json::Value = match serde_json::from_slice(body) { + Ok(v) => v, + Err(_) => return PluginResponse::json(r#"{"error":"invalid json"}"#.into()), + }; + let node_id_hex = match parsed.get("node_id").and_then(|v| v.as_str()) { + Some(s) => s, + None => return PluginResponse::json(r#"{"error":"missing node_id"}"#.into()), + }; + + // Parse hex node_id into [u8; 32] + let bytes = match parse_hex_node_id(node_id_hex) { + Some(b) => b, + None => return PluginResponse::json(r#"{"error":"invalid node_id hex"}"#.into()), + }; + + // Look up relay_url from cached snapshot + let relay_url = { + let guard = self.cached.lock().unwrap(); + guard.as_ref().and_then(|snap| { + snap.members.iter() + .find(|m| m.node_id == node_id_hex) + .and_then(|m| m.relay_url.clone()) + }) + }; + + match tx.send((bytes, relay_url)) { + Ok(()) => PluginResponse::json(r#"{"ok":true}"#.into()), + Err(_) => PluginResponse::json(r#"{"error":"channel closed"}"#.into()), + } + } +} + +fn parse_hex_node_id(hex: &str) -> Option<[u8; 32]> { + if hex.len() != 64 { + return None; + } + let mut bytes = [0u8; 32]; + for i in 0..32 { + bytes[i] = u8::from_str_radix(&hex[i * 2..i * 2 + 2], 16).ok()?; + } + Some(bytes) +} diff --git a/crates/node/src/plugins/distribution_page.html b/crates/node/src/plugins/distribution_page.html index 81032f5..1e6ba0d 100644 --- a/crates/node/src/plugins/distribution_page.html +++ b/crates/node/src/plugins/distribution_page.html @@ -80,6 +80,13 @@ .state-suspect { color: #ff9800; } .state-dead { color: #f44336; } + .repeer-btn { + background: none; border: 1px solid #6366f1; color: #6366f1; + padding: 1px 6px; font-size: 9px; font-family: inherit; + border-radius: 3px; cursor: pointer; margin-left: 4px; + } + .repeer-btn:hover { background: #6366f1; color: #fff; } + .bottom-panel { grid-column: 1 / -1; border-top: 1px solid #2a2d3e; background: #161822; display: flex; gap: 12px; padding: 12px 16px; height: 200px; @@ -171,6 +178,7 @@
+ Waiting for data...
@@ -215,7 +223,7 @@

Members

- +
StateNode IDAddressInc
StateNode IDAddressInc
@@ -632,7 +640,8 @@ ctx.globalAlpha = opacity; ctx.font = Math.round(9 / vs) + 'px monospace'; ctx.textAlign = 'center'; - var label = (i === 0 && data && data.node_name) ? data.node_name : (nodeIds[i] ? nodeIds[i].substring(0, 8) : ''); + var memberName = (i > 0 && data && data.members[i - 1]) ? data.members[i - 1].node_name : null; + var label = (i === 0 && data && data.node_name) ? data.node_name : (memberName || (nodeIds[i] ? nodeIds[i].substring(0, 8) : '')); ctx.fillText(label, posX[i], posY[i] - nr - 3/vs); ctx.globalAlpha = 1; } @@ -696,6 +705,7 @@ // Self label document.getElementById('selfLabel').textContent = 'Node: ' + (d.node_name || d.node_id.substring(0, 16) + '\u2026'); document.getElementById('nodeLabel').textContent = d.node_name || d.listen_addr; + document.getElementById('versionTag').textContent = d.version || ''; // Stats cards document.getElementById('statMembers').textContent = d.members.length; @@ -721,16 +731,22 @@ var m = d.members[i]; var cls = 'state-' + m.state; var tr = document.createElement('tr'); - var idCell = m.node_id.substring(0, 16) + '\u2026'; - if (m.label) idCell = m.label + ' (' + m.node_id.substring(0, 8) + ')'; + var displayName = m.label || m.node_name; + var idCell = displayName + ? displayName + ' (' + m.node_id.substring(0, 8) + ')' + : m.node_id.substring(0, 16) + '\u2026'; var relayBadge = m.relay_url ? ' RELAY' : ''; + var actionCell = m.state === 'dead' + ? '' + : ''; tr.innerHTML = '' + m.state + '' + '' + idCell + relayBadge + '' + '' + m.addr + '' + - '' + m.incarnation + ''; + '' + m.incarnation + '' + + '' + actionCell + ''; body.appendChild(tr); } @@ -763,9 +779,13 @@ var addr = mem ? mem.addr : '\u2014'; var cls = 'state-' + state; var tr = document.createElement('tr'); + var probeName = mem && (mem.label || mem.node_name); + var probeLabel = probeName + ? probeName + ' (' + pid.substring(0, 8) + ')' + : pid.substring(0, 12) + '\u2026'; tr.innerHTML = '' + state + '' + - '' + pid.substring(0, 12) + '\u2026' + + '' + probeLabel + '' + '' + addr + ''; probesBody.appendChild(tr); } @@ -824,6 +844,7 @@ // ── SSE connection ───────────────────────────────────────────── var es = new EventSource('/events'); + window.addEventListener('beforeunload', function() { es.close(); }); es.addEventListener('distribution', function(e) { try { @@ -912,6 +933,17 @@ }).then(function() { fetchPeers(); }); }; + window.repeerNode = function(nid) { + fetch('/api/plugin/distribution/rejoin', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ node_id: nid }) + }).then(function(r) { return r.json(); }).then(function(d) { + if (d.ok) { console.log('Re-peer triggered for ' + nid.substring(0, 8)); } + else { console.error('Re-peer failed:', d.error); } + }).catch(function(e) { console.error('Re-peer error:', e); }); + }; + // Poll peers every 5 seconds fetchPeers(); setInterval(fetchPeers, 5000); diff --git a/src/channel.rs b/src/channel.rs index 3d2a1a1..88b536f 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -28,6 +28,10 @@ impl HybridChannel { pub fn pop(&self) -> Option { self.ring.pop().or_else(|| self.overflow.pop()) } + + pub fn is_empty(&self) -> bool { + self.ring.is_empty() && self.overflow.is_empty() + } } pub(crate) struct Receiver { @@ -44,6 +48,10 @@ impl Receiver { self.queue.pop() } + pub fn is_empty(&self) -> bool { + self.queue.is_empty() + } + pub fn new_sender(&self) -> Sender { Sender { queue: self.queue.clone(), diff --git a/src/config.rs b/src/config.rs index 4772bd4..85ddee4 100644 --- a/src/config.rs +++ b/src/config.rs @@ -1,28 +1,3 @@ -/// Backoff policy for worker threads when idle. -/// -/// Workers spin → yield → sleep with increasing delay when no work is available. -pub struct BackoffPolicy { - /// Number of idle ticks before switching from spin to yield. - pub spin_threshold: u32, - /// Number of idle ticks before switching from yield to sleep. - pub yield_threshold: u32, - /// Microseconds added per tick beyond the yield threshold. - pub sleep_increment_us: u64, - /// Maximum sleep duration in microseconds. - pub sleep_max_us: u64, -} - -impl Default for BackoffPolicy { - fn default() -> Self { - Self { - spin_threshold: 64, - yield_threshold: 256, - sleep_increment_us: 50, - sleep_max_us: 1000, - } - } -} - /// What to do when a bounded mailbox is full. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum MailboxOverflow { @@ -37,7 +12,6 @@ pub struct RuntimeConfig { pub max_actors: usize, pub channel_buffer_size: usize, pub num_threads: usize, - pub backoff_policy: BackoffPolicy, /// Maximum messages processed per actor per tick. /// Prevents a single actor with a large mailbox from starving others. /// `0` means unlimited (drain entire mailbox). @@ -67,7 +41,6 @@ impl Default for RuntimeConfig { max_actors: DEFAULT_MAX_ACTORS, channel_buffer_size: DEFAULT_CHANNEL_BUFFER_SIZE, num_threads: 1, - backoff_policy: BackoffPolicy::default(), actor_message_budget: DEFAULT_ACTOR_MESSAGE_BUDGET, default_mailbox_capacity: 0, mailbox_overflow: MailboxOverflow::DropNewest, diff --git a/src/extension.rs b/src/extension.rs index a1a170d..e9ca2e4 100644 --- a/src/extension.rs +++ b/src/extension.rs @@ -53,6 +53,10 @@ pub trait RuntimeExtension: Send + Sync { /// - `handle_request`: phase 5.5 — processes deferred requests from handlers /// - `gc_dead`: after cleanup_dead — removes state for dead actors pub trait WorkerExtension: Send { + /// Returns `true` if this extension has pending work (e.g., active timers). + /// Used by the fast idle path to avoid unnecessary ticks. + fn has_pending_work(&self) -> bool { false } + /// Called each tick before tick_all. Returns messages to deliver. fn on_tick(&mut self) -> Vec<(ActorAddress, Box)>; diff --git a/src/runtime.rs b/src/runtime.rs index a15b693..e75726a 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -10,7 +10,7 @@ use crate::Instant; use crate::actor::{Actor, ActorAddress, ActorInterface, AnyActor, Environment, ExitValue, Message, ResumeSignal, SpawnRequest, StopSignal, StopWithSignal, SystemInfo}; use crate::channel::{Receiver, Sender}; // Re-export config types so existing code using `runtime::RuntimeConfig` still works -pub use crate::config::{BackoffPolicy, MailboxOverflow, RuntimeConfig}; +pub use crate::config::{MailboxOverflow, RuntimeConfig}; use crate::delivery::{AddressMap, Envelope, InboxRegistry, Placement, TickContext, WorkerId}; use crate::extension::RuntimeExtension; use crate::stats::{StatsHook, WorkerStats}; diff --git a/src/std/timer_wheel.rs b/src/std/timer_wheel.rs index b0af03a..ae4c83f 100644 --- a/src/std/timer_wheel.rs +++ b/src/std/timer_wheel.rs @@ -127,7 +127,14 @@ impl TimerWheel { } impl WorkerExtension for TimerWheel { + fn has_pending_work(&self) -> bool { + !self.once_timers.is_empty() || !self.interval_timers.is_empty() + } + fn on_tick(&mut self) -> Vec<(ActorAddress, Box)> { + if self.once_timers.is_empty() && self.interval_timers.is_empty() { + return Vec::new(); + } self.fire() } diff --git a/src/worker.rs b/src/worker.rs index e11ffa1..50ce227 100644 --- a/src/worker.rs +++ b/src/worker.rs @@ -51,6 +51,9 @@ pub(crate) struct Worker { snapshot_buf: Vec, /// Per-worker extension (e.g., timer wheel). Created by RuntimeExtension factory. pub(crate) worker_ext: Option>, + /// True if the previous tick did work — ensures one full tick follows a productive + /// tick so pending_local messages delivered to mailboxes get drained. + has_backlog: bool, } impl Worker { @@ -70,6 +73,7 @@ impl Worker { stats, snapshot_buf: Vec::new(), worker_ext: None, + has_backlog: false, } } @@ -154,6 +158,16 @@ impl Worker { #[cfg(feature = "tracing")] let _span = tracing::trace_span!("worker.tick", worker_id = self.id.0).entered(); + // Fast idle path: skip the entire tick when nothing could have changed. + // Cost: ~3 atomic loads, zero syscalls, zero actor iteration. + if !self.has_backlog + && self.spawn_rx.is_empty() + && self.transfer_rx.is_empty() + && !self.worker_ext.as_ref().map_or(false, |e| e.has_pending_work()) + { + return false; + } + let mut did_work = false; let t0 = Instant::now(); @@ -285,6 +299,7 @@ impl Worker { // 7. Clean up poisoned and stopping actors did_work |= self.cleanup_dead_actors(tc); + self.has_backlog = did_work; did_work } @@ -292,27 +307,11 @@ impl Worker { #[cfg(feature = "tracing")] let _span = tracing::info_span!("worker.run", worker_id = self.id.0).entered(); - let backoff = &tc.config.backoff_policy; - let mut idle_count: u32 = 0; while is_running.load(Ordering::Acquire) { - let did_work = self.tick_once(tc); - if did_work { - idle_count = 0; - } else { - idle_count = idle_count.saturating_add(1); - if idle_count < backoff.spin_threshold { - // Hot spin - } else if idle_count < backoff.yield_threshold { - thread::yield_now(); - } else { - let micros = std::cmp::min( - (idle_count - backoff.yield_threshold) as u64 * backoff.sleep_increment_us, - backoff.sleep_max_us, - ); - // park_timeout allows instant wakeup via Thread::unpark() - // when new work arrives (send_to/spawn notify the target worker) - thread::park_timeout(std::time::Duration::from_micros(micros)); - } + if !self.tick_once(tc) { + // Park indefinitely — woken by unpark() from send_to/spawn/stop/shutdown. + // Spurious wakes hit the fast idle path (~3 atomic loads) and park again. + thread::park(); } } } diff --git a/tests/test_python.py b/tests/test_python.py index 9e8a773..4ad7b34 100644 --- a/tests/test_python.py +++ b/tests/test_python.py @@ -38,10 +38,6 @@ class TestRuntimeConfig(unittest.TestCase): self.assertEqual(cfg.num_threads, 1) self.assertEqual(cfg.max_actors, 1000) self.assertEqual(cfg.channel_buffer_size, 1000) - self.assertEqual(cfg.spin_threshold, 64) - self.assertEqual(cfg.yield_threshold, 256) - self.assertEqual(cfg.sleep_increment_us, 50) - self.assertEqual(cfg.sleep_max_us, 1000) def test_custom(self): cfg = RuntimeConfig(num_threads=4, max_actors=500)