diff --git a/crates/distribution/src/directory_actor.rs b/crates/distribution/src/directory_actor.rs index 215ae62..5669f54 100644 --- a/crates/distribution/src/directory_actor.rs +++ b/crates/distribution/src/directory_actor.rs @@ -23,7 +23,7 @@ //! [`RegistryActor`]: crate::registry_actor::RegistryActor //! [`MetadataActor`]: crate::node_metadata_actor::MetadataActor -use std::collections::{BTreeSet, HashMap}; +use std::collections::{BTreeSet, HashMap, HashSet}; use std::sync::Arc; use swactor::actor::{ActorAddress, ActorInterface}; @@ -95,6 +95,10 @@ pub struct DirectoryActor { /// Registers a route so the runtime can deliver an app message addressed to a /// remotely-hosted actor (the §5 egress seam). Called from [`Self::republish`]. route_binder: Arc, + /// Remotely-hosted actors bound by the previous [`Self::republish`]; the + /// diff against the next view drives `remove_route` so departed actors + /// do not accumulate in the binder and router forever. + bound_remote: HashSet, /// Routes owned outside SWIM membership (for example an exec child actor /// runtime attached to this host). Directory republishing always preserves /// these entries. @@ -133,6 +137,7 @@ impl DirectoryActor { peer_directory, route_view, route_binder, + bound_remote: HashSet::new(), pinned_routes, } } @@ -225,13 +230,10 @@ impl DirectoryActor { /// Republish the §5 route view: every actor whose host is reachable right now /// (self, or an alive peer). A dead host's actors are omitted, so the egress - /// never routes to a host SWIM has buried; the claim remains cached for recovery. - /// - /// Builds the new view locally, swaps it in under the lock, then registers a - /// route for each remotely-hosted actor *outside* the lock — so the egress - /// (which takes the view's read lock on every send) never contends with the - /// route registration. - fn republish(&self) { + /// never routes to a host SWIM has buried; the claim remains cached for + /// recovery. Actors that left the remotely-routed set are unbound so the + /// binder and transport router do not grow without bound. + fn republish(&mut self) { let pinned = self .pinned_routes .read() @@ -240,7 +242,7 @@ impl DirectoryActor { let mut remote = pinned .iter() .filter_map(|(actor, node)| (*node != self.self_id).then_some(*actor)) - .collect::>(); + .collect::>(); drop(pinned); for (actor, claim) in &self.map { if view.contains_key(actor) { @@ -250,13 +252,22 @@ impl DirectoryActor { view.insert(*actor, claim.node_id); } else if self.alive.contains(&claim.node_id) { view.insert(*actor, claim.node_id); - remote.push(*actor); + remote.insert(*actor); } } *self.route_view.write().expect("route view poisoned") = view; - for actor in remote { + // Unbind actors that fell out of the remotely-routed set (host left + // the cluster or claim superseded): without this diff the binder and + // router retain every ever-seen actor forever. + let departed: Vec = self.bound_remote.difference(&remote).copied().collect(); + for actor in departed { + self.route_binder.remove_route(&actor); + } + let joined: Vec = remote.difference(&self.bound_remote).copied().collect(); + for actor in joined { self.route_binder.ensure_routable(actor); } + self.bound_remote = remote; } /// Take up to `limit` armed claims for this tick's batch, spending one unit of diff --git a/crates/distribution/tests/directory_unbind.rs b/crates/distribution/tests/directory_unbind.rs new file mode 100644 index 0000000..340b13c --- /dev/null +++ b/crates/distribution/tests/directory_unbind.rs @@ -0,0 +1,110 @@ +//! Directory route-binding lifecycle: republish must unbind actors whose host +//! left the cluster, or the binder and transport router grow without bound +//! under membership churn. + +use std::collections::HashMap; +use std::sync::Arc; + +use parking_lot::Mutex; +use swactor::actor::ActorAddress; +use swactor::runtime::{RuntimeConfig, RuntimeParts}; +use swactor::std::StdExtension; +use swactor_engine::{Engine, SteppingBackend}; + +use distribution::crypto::{Keypair, KeypairExt}; +use distribution::directory_actor::{DirectoryActor, DirectoryIn}; +use distribution::swim::actor::{MembershipChanged, SharedPeerDirectory}; +use distribution::transport_bridge::{RouteBinder, RouteView}; +use distribution::types::{MemberState, NodeId}; + +/// Records every bind/unbind so a test can assert the republish diff. +struct RecordingBinder { + events: Mutex>, +} + +impl RouteBinder for RecordingBinder { + fn ensure_routable(&self, actor: ActorAddress) { + self.events.lock().push((actor, true)); + } + + fn remove_route(&self, actor: &ActorAddress) { + self.events.lock().push((*actor, false)); + } +} + +fn membership(node: NodeId, state: MemberState) -> DirectoryIn { + DirectoryIn::Membership(MembershipChanged { + node_id: node, + state, + incarnation: 0, + }) +} + +#[test] +fn republish_unbinds_routes_of_departed_hosts() { + let parts = + RuntimeParts::new(RuntimeConfig::default()).with_extension(Arc::new(StdExtension::new())); + let rt = parts.runtime().clone(); + let backend = SteppingBackend::new(); + let _engine = Engine::new(parts, backend.clone()).expect("create engine"); + + let self_key = Keypair::generate(); + let self_id = self_key.node_id(); + let peer_key = Keypair::generate(); + let peer_id = peer_key.node_id(); + + let route_view: RouteView = Arc::new(std::sync::RwLock::new(HashMap::new())); + let binder = Arc::new(RecordingBinder { + events: Mutex::new(Vec::new()), + }); + let directory = rt + .spawn(DirectoryActor::new( + self_id, + Arc::new(SharedPeerDirectory::new()), + route_view.clone(), + binder.clone(), + )) + .expect("spawn DirectoryActor"); + + // A remotely-hosted actor claim plus its host being alive. + let actor = ActorAddress::new_random(); + let claim = peer_key.sign_directory_entry(actor, 1); + rt.send_to(directory, DirectoryIn::Register(claim)).unwrap(); + rt.send_to(directory, membership(peer_id, MemberState::Alive)) + .unwrap(); + backend.step(); + + let events = binder.events.lock().clone(); + assert!( + events.contains(&(actor, true)), + "alive peer's actor must be bound: {events:?}" + ); + + // Host dies: its actor leaves the route view and must be unbound. + rt.send_to(directory, membership(peer_id, MemberState::Dead)) + .unwrap(); + backend.step(); + let events = binder.events.lock().clone(); + assert!( + events.contains(&(actor, false)), + "dead host's actor must be unbound: {events:?}" + ); + assert!( + route_view.read().unwrap().get(&actor).is_none(), + "dead host's actor must leave the route view" + ); + + // Host returns: the cached claim re-enters the view, so bind again. + rt.send_to(directory, membership(peer_id, MemberState::Alive)) + .unwrap(); + backend.step(); + let events = binder.events.lock().clone(); + assert_eq!( + events + .iter() + .filter(|(bound_actor, bound)| *bound_actor == actor && *bound) + .count(), + 2, + "returning host's actor must be re-bound: {events:?}" + ); +}