From 4d18874909eba0d2195f26a784b287d82afd1414 Mon Sep 17 00:00:00 2001 From: Developer Date: Thu, 12 Feb 2026 13:22:35 +0000 Subject: [PATCH] feat: actor groups with pub-sub broadcast (Cycle 14) Add GroupRegistry for named actor groups. Actors join/leave groups via rt.join_group() / ctx.join_group(), and messages can be broadcast to all members via rt.publish_to() / ctx.publish(). Groups auto-create on first join, auto-delete when empty, and members are auto-removed on death. Co-Authored-By: Claude Opus 4.6 --- CLAUDE/notes/progress.md | 31 +++++- src/actor.rs | 38 +++++++ src/delivery.rs | 81 ++++++++++++++- src/runtime.rs | 52 +++++++++- src/worker.rs | 13 +++ tests/runtime_api.rs | 219 +++++++++++++++++++++++++++++++++++++++ 6 files changed, 431 insertions(+), 3 deletions(-) diff --git a/CLAUDE/notes/progress.md b/CLAUDE/notes/progress.md index 5c28cfe..0d9feec 100644 --- a/CLAUDE/notes/progress.md +++ b/CLAUDE/notes/progress.md @@ -2,7 +2,7 @@ ## Current Stage: Phase 1 — Research + First Improvement Cycle -### Status: Cycle 13 COMPLETE +### Status: Cycle 14 COMPLETE ## Plan Overview 1. **Phase 0**: Codebase audit — understand current swactor architecture, existing tests, benchmarks ✅ @@ -124,6 +124,35 @@ - `stop_nonexistent_actor_returns_error` — stop on bad address returns Err - **Result**: 82 tests pass, all workspace compiles +### Cycle 14: Actor Groups (Pub-Sub) +- **Research**: Studied group/pub-sub patterns across Erlang pg (scopes, join/leave/get_members), + Akka DistributedPubSub (mediator, topics), Ractor pg (join/leave/broadcast), Bastion (Dispatcher), + Redis pub/sub (channels, patterns) + - Common patterns: auto-cleanup on death, at-most-once delivery, string-based naming, + flat groups (not hierarchical), lazy creation/deletion + - Decision: Erlang pg-style flat groups, string keys, auto-cleanup, RwLock pattern +- **Implementation**: `GroupRegistry` in delivery.rs with forward + reverse maps + - `groups: RwLock>>` — group→members + - `memberships: RwLock>>` — actor→groups (reverse for cleanup) + - Groups auto-create on first join, auto-delete when empty + - Runtime API: `join_group(addr, name)`, `leave_group(addr, name)`, `publish_to(group, msg)`, + `group_members(group)`, `groups()` + - Ctx API: `join_group(name)`, `leave_group(name)`, `publish(group, msg)`, `group_members(group)` + - `publish` clones at the typed level (Message: Clone), sends to each member via normal routing + - Auto-cleanup: `group_registry.cleanup(&addr)` in cleanup_dead phase removes dead actor from all groups + - ContextInner extended: `join_group()`, `leave_group()`, `group_members()` (publish is Ctx-level only) +- **Tests**: 9 new behavioral tests + - `group_members_returns_joined_actors` — join + query + - `empty_group_returns_no_members` — nonexistent group → empty + - `publish_broadcasts_to_all_members` — 2 members, both receive + - `leave_group_stops_receiving_publishes` — leave → excluded from broadcast + - `dead_actor_auto_removed_from_group` — stop → removed from group + - `actor_removed_from_all_groups_on_death` — multi-group membership cleanup + - `empty_group_auto_deleted` — last member leaves → group removed from groups() + - `ctx_join_group_from_handler` — join via on_start + - `ctx_publish_broadcasts_from_handler` — publish via handler +- **Result**: 122 tests pass (115 behavioral + 7 proptest), all workspace compiles, zero warnings + ### Cycle 13: Actor Monitoring / Death Watch - **Research**: Studied monitoring across Erlang (monitor/2, DOWN messages), Akka (watch/Terminated), Ractor (link, SupervisionEvent), Actix (none), Kameo (link, on_link_died callback) diff --git a/src/actor.rs b/src/actor.rs index 3ba2008..e31ae26 100644 --- a/src/actor.rs +++ b/src/actor.rs @@ -213,6 +213,12 @@ pub trait ContextInner { fn monitor(&self, watcher: ActorAddress, target: ActorAddress) -> MonitorRef; /// Cancel a monitor subscription. fn demonitor(&self, mref: MonitorRef); + /// Add actor to a named group. + fn join_group(&self, actor: ActorAddress, group: String); + /// Remove actor from a named group. + fn leave_group(&self, actor: ActorAddress, group: &str); + /// Return all members of a named group. + fn group_members(&self, group: &str) -> Vec; } /// Actor syscall interface — passed to `ActorInterface::handle()`. @@ -321,6 +327,38 @@ impl<'a> Ctx<'a> { self.inner.demonitor(mref); } + /// Join a named group. The group is created if it doesn't exist. + /// + /// An actor can be a member of multiple groups simultaneously. + pub fn join_group(&self, group: impl Into) { + self.inner.join_group(self.self_addr, group.into()); + } + + /// Leave a named group. Empty groups are automatically deleted. + pub fn leave_group(&self, group: &str) { + self.inner.leave_group(self.self_addr, group); + } + + /// Broadcast a message to all members of a named group. + /// + /// The message is cloned for each recipient. Returns the number of + /// messages successfully enqueued. + pub fn publish(&self, group: &str, msg: M) -> usize { + let members = self.inner.group_members(group); + let mut count = 0; + for member in &members { + if self.inner.send_any(*member, Box::new(msg.clone())).is_ok() { + count += 1; + } + } + count + } + + /// Return all current members of a named group. + pub fn group_members(&self, group: &str) -> Vec { + self.inner.group_members(group) + } + /// Spawn a restartable actor. On panic, recreated via `factory` up to /// `max_restarts` times before permanent poisoning. pub fn spawn_restartable( diff --git a/src/delivery.rs b/src/delivery.rs index 027f72b..1ea5364 100644 --- a/src/delivery.rs +++ b/src/delivery.rs @@ -1,5 +1,5 @@ use std::any::Any; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; use std::sync::{Arc, OnceLock, RwLock}; use std::thread::Thread; @@ -191,6 +191,7 @@ pub(crate) struct TickContext<'a> { pub(crate) config: &'a RuntimeConfig, pub(crate) name_registry: &'a NameRegistry, pub(crate) monitor_registry: &'a MonitorRegistry, + pub(crate) group_registry: &'a GroupRegistry, pub(crate) stats_hook: Option<&'a dyn crate::stats::StatsHook>, /// Thread handles for waking parked workers on cross-worker sends. pub(crate) worker_threads: &'a [OnceLock], @@ -332,6 +333,84 @@ impl MonitorRegistry { } } +// ─── Group Registry ───────────────────────────────────────────────────────── + +/// Actor groups (pub-sub). Actors join/leave named groups; messages can be +/// broadcast to all members of a group. +/// +/// Groups are created lazily on first join and removed when empty. +pub(crate) struct GroupRegistry { + /// group_name → set of member addresses + groups: RwLock>>, + /// actor_addr → set of group names (reverse map for O(G) cleanup on death) + memberships: RwLock>>, +} + +impl GroupRegistry { + pub fn new() -> Self { + Self { + groups: RwLock::new(HashMap::new()), + memberships: RwLock::new(HashMap::new()), + } + } + + /// Add an actor to a named group. Group is created if it doesn't exist. + pub fn join(&self, group: String, addr: ActorAddress) { + self.groups.write().unwrap() + .entry(group.clone()) + .or_default() + .insert(addr); + self.memberships.write().unwrap() + .entry(addr) + .or_default() + .insert(group); + } + + /// Remove an actor from a named group. Empty groups are auto-deleted. + pub fn leave(&self, group: &str, addr: &ActorAddress) { + let mut groups = self.groups.write().unwrap(); + if let Some(members) = groups.get_mut(group) { + members.remove(addr); + if members.is_empty() { + groups.remove(group); + } + } + drop(groups); + if let Some(membership) = self.memberships.write().unwrap().get_mut(addr) { + membership.remove(group); + } + } + + /// Return all members of a group. + pub fn members(&self, group: &str) -> Vec { + self.groups.read().unwrap() + .get(group) + .map(|s| s.iter().copied().collect()) + .unwrap_or_default() + } + + /// Remove a dead actor from all its groups. + pub fn cleanup(&self, addr: &ActorAddress) { + let group_names = self.memberships.write().unwrap().remove(addr); + if let Some(names) = group_names { + let mut groups = self.groups.write().unwrap(); + for name in names { + if let Some(members) = groups.get_mut(&name) { + members.remove(addr); + if members.is_empty() { + groups.remove(&name); + } + } + } + } + } + + /// Return all active group names. + pub fn group_names(&self) -> Vec { + self.groups.read().unwrap().keys().cloned().collect() + } +} + impl<'a> TickContext<'a> { /// Route a message whose destination is not in the local address map. /// Tries inbox registry, then remote transport, then falls back to inbox error. diff --git a/src/runtime.rs b/src/runtime.rs index d1330b2..ee7a3e4 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -9,7 +9,7 @@ use crate::actor::{Actor, ActorAddress, ActorInterface, AnyActor, Message, StopS use crate::channel::{Receiver, Sender}; // Re-export config types so existing code using `runtime::RuntimeConfig` still works pub use crate::config::{BackoffPolicy, MailboxOverflow, RuntimeConfig}; -use crate::delivery::{AddressMap, Envelope, InboxRegistry, MonitorRegistry, NameRegistry, Placement, TickContext, WorkerId}; +use crate::delivery::{AddressMap, Envelope, GroupRegistry, InboxRegistry, MonitorRegistry, NameRegistry, Placement, TickContext, WorkerId}; use crate::stats::{StatsHook, WorkerStats}; // Re-export stats types so existing code using `runtime::*` still works pub use crate::stats::{RuntimeStats, WorkerInfo}; @@ -64,6 +64,7 @@ pub struct Runtime { inbox_registry: Arc, name_registry: Arc, monitor_registry: Arc, + group_registry: Arc, transfer_txs: Vec>, spawn_txs: Vec)>>, placement: Placement, @@ -123,6 +124,7 @@ impl Runtime { let inbox_registry = Arc::new(InboxRegistry::new()); let name_registry = Arc::new(NameRegistry::new()); let monitor_registry = Arc::new(MonitorRegistry::new()); + let group_registry = Arc::new(GroupRegistry::new()); let mut transfer_txs = Vec::with_capacity(num_workers); let mut spawn_txs = Vec::with_capacity(num_workers); @@ -162,6 +164,7 @@ impl Runtime { inbox_registry, name_registry, monitor_registry, + group_registry, transfer_txs, spawn_txs, placement, @@ -264,6 +267,40 @@ impl Runtime { self.name_registry.registered_names() } + /// Add an actor to a named group. The group is created if it doesn't exist. + pub fn join_group(&self, addr: ActorAddress, group: impl Into) { + self.group_registry.join(group.into(), addr); + } + + /// Remove an actor from a named group. Empty groups are auto-deleted. + pub fn leave_group(&self, addr: ActorAddress, group: &str) { + self.group_registry.leave(group, &addr); + } + + /// Broadcast a message to all members of a named group. + /// + /// Returns the number of messages successfully enqueued. + pub fn publish_to(&self, group: &str, msg: M) -> usize { + let members = self.group_registry.members(group); + let mut count = 0; + for member in &members { + if self.send_to(*member, msg.clone()).is_ok() { + count += 1; + } + } + count + } + + /// Return all current members of a named group. + pub fn group_members(&self, group: &str) -> Vec { + self.group_registry.members(group) + } + + /// Return all active group names. + pub fn groups(&self) -> Vec { + self.group_registry.group_names() + } + /// Send a message to an actor address pub fn send_to(&self, addr: ActorAddress, msg: M) -> Result<(), Error> { let result = self.send_any(addr, Box::new(msg)); @@ -296,6 +333,7 @@ impl Runtime { config: &self.config, name_registry: &self.name_registry, monitor_registry: &self.monitor_registry, + group_registry: &self.group_registry, stats_hook: self.stats_hook.as_deref(), worker_threads: &self.worker_threads, #[cfg(feature = "transport")] @@ -512,4 +550,16 @@ impl ContextInner for Runtime { fn demonitor(&self, mref: crate::actor::MonitorRef) { self.monitor_registry.deregister(mref); } + + fn join_group(&self, actor: ActorAddress, group: String) { + self.group_registry.join(group, actor); + } + + fn leave_group(&self, actor: ActorAddress, group: &str) { + self.group_registry.leave(group, &actor); + } + + fn group_members(&self, group: &str) -> Vec { + self.group_registry.members(group) + } } diff --git a/src/worker.rs b/src/worker.rs index bacbc0d..248f5e2 100644 --- a/src/worker.rs +++ b/src/worker.rs @@ -322,6 +322,7 @@ impl Worker { for &(addr, _) in &dead { tc.address_map.remove(&addr); tc.name_registry.unregister_by_addr(&addr); + tc.group_registry.cleanup(&addr); } // Re-publish num_actors after cleanup so stats reflect removal self.stats.num_actors.store(self.pool.len(), Ordering::Relaxed); @@ -461,6 +462,18 @@ impl ContextInner for WorkerContext<'_> { fn demonitor(&self, mref: crate::actor::MonitorRef) { self.tc.monitor_registry.deregister(mref); } + + fn join_group(&self, actor: ActorAddress, group: String) { + self.tc.group_registry.join(group, actor); + } + + fn leave_group(&self, actor: ActorAddress, group: &str) { + self.tc.group_registry.leave(group, &actor); + } + + fn group_members(&self, group: &str) -> Vec { + self.tc.group_registry.members(group) + } } struct ActorSlot { diff --git a/tests/runtime_api.rs b/tests/runtime_api.rs index 3c2fed2..63d0f88 100644 --- a/tests/runtime_api.rs +++ b/tests/runtime_api.rs @@ -3329,3 +3329,222 @@ fn stacked_monitors_produce_multiple_notifications() { assert!(inbox.try_recv().is_some(), "second Down"); assert!(inbox.try_recv().is_none(), "no more"); } + +// ── Actor Groups / Pub-Sub ────────────────────────────────────────────────── + +/// Given actors join a group, +/// when I query group_members, +/// then all joined actors are listed. +#[test] +fn group_members_returns_joined_actors() { + let rt = Runtime::new(RuntimeConfig::default()); + let a = rt.spawn(PingPongActor).unwrap(); + let b = rt.spawn(PingPongActor).unwrap(); + + rt.join_group(a, "workers"); + rt.join_group(b, "workers"); + + let mut members = rt.group_members("workers"); + members.sort_by_key(|addr| addr.0); + let mut expected = vec![a, b]; + expected.sort_by_key(|addr| addr.0); + assert_eq!(members, expected); +} + +/// Given no actors have joined a group, +/// when I query group_members, +/// then the result is empty. +#[test] +fn empty_group_returns_no_members() { + let rt = Runtime::new(RuntimeConfig::default()); + assert!(rt.group_members("nonexistent").is_empty()); +} + +/// Given actors in a group, +/// when a message is published to the group, +/// then all members receive the message. +#[test] +fn publish_broadcasts_to_all_members() { + let rt = Runtime::new(RuntimeConfig::default()); + let inbox1 = rt.new_inbox::().unwrap(); + let inbox2 = rt.new_inbox::().unwrap(); + + let a = rt.spawn(PingPongActor).unwrap(); + let b = rt.spawn(PingPongActor).unwrap(); + rt.join_group(a, "pongers"); + rt.join_group(b, "pongers"); + + rt.tick(); // on_start + + // Publish a Ping with different reply_to for each — but since it's cloned, + // all members get the same message. Use inbox1's addr as reply_to. + let count = rt.publish_to("pongers", Ping { reply_to: *inbox1.addr() }); + assert_eq!(count, 2, "two members, two messages sent"); + + rt.tick(); // actors handle Ping → send Pong to inbox1 + + // Both actors send to inbox1 (because the published Ping had inbox1 as reply_to) + assert!(inbox1.try_recv().is_some(), "first Pong"); + assert!(inbox1.try_recv().is_some(), "second Pong"); + assert!(inbox1.try_recv().is_none(), "no more"); + drop(inbox2); +} + +/// Given an actor leaves a group, +/// when a message is published, +/// then the leaver does not receive it. +#[test] +fn leave_group_stops_receiving_publishes() { + let rt = Runtime::new(RuntimeConfig::default()); + let inbox = rt.new_inbox::().unwrap(); + + let a = rt.spawn(PingPongActor).unwrap(); + let b = rt.spawn(PingPongActor).unwrap(); + rt.join_group(a, "pool"); + rt.join_group(b, "pool"); + rt.leave_group(b, "pool"); + + rt.tick(); // on_start + let count = rt.publish_to("pool", Ping { reply_to: *inbox.addr() }); + assert_eq!(count, 1, "only one member after leave"); + + rt.tick(); + assert!(inbox.try_recv().is_some(), "one Pong from remaining member"); + assert!(inbox.try_recv().is_none(), "no second Pong"); +} + +/// Given a group member dies, +/// when a message is published, +/// then the dead member is not included. +#[test] +fn dead_actor_auto_removed_from_group() { + let rt = Runtime::new(RuntimeConfig::default()); + let inbox = rt.new_inbox::().unwrap(); + + let a = rt.spawn(PingPongActor).unwrap(); + let b = rt.spawn(PingPongActor).unwrap(); + rt.join_group(a, "team"); + rt.join_group(b, "team"); + + rt.tick(); // on_start + rt.stop_actor(b).unwrap(); + rt.tick(); // b dies, cleaned up from group + + let count = rt.publish_to("team", Ping { reply_to: *inbox.addr() }); + assert_eq!(count, 1, "dead actor removed from group"); + + rt.tick(); + assert!(inbox.try_recv().is_some()); + assert!(inbox.try_recv().is_none()); +} + +/// Given an actor is in multiple groups, +/// when the actor dies, +/// then it is removed from all groups. +#[test] +fn actor_removed_from_all_groups_on_death() { + let rt = Runtime::new(RuntimeConfig::default()); + let actor = rt.spawn(PingPongActor).unwrap(); + rt.join_group(actor, "alpha"); + rt.join_group(actor, "beta"); + rt.join_group(actor, "gamma"); + + rt.tick(); + rt.stop_actor(actor).unwrap(); + rt.tick(); // cleanup removes from all groups + + assert!(rt.group_members("alpha").is_empty()); + assert!(rt.group_members("beta").is_empty()); + assert!(rt.group_members("gamma").is_empty()); +} + +/// Given a group becomes empty after its last member leaves, +/// then the group name disappears from the active groups list. +#[test] +fn empty_group_auto_deleted() { + let rt = Runtime::new(RuntimeConfig::default()); + let actor = rt.spawn(PingPongActor).unwrap(); + rt.join_group(actor, "temp"); + assert!(rt.groups().contains(&"temp".to_string())); + + rt.leave_group(actor, "temp"); + assert!(!rt.groups().contains(&"temp".to_string()), "empty group should be removed"); +} + +/// Given actors join groups from handlers using ctx.join_group(), +/// when group_members is queried, +/// then the joining actors are listed. +#[test] +fn ctx_join_group_from_handler() { + struct GroupJoinerActor; + + impl ActorInterface for GroupJoinerActor { + type Incoming = Ping; + type Response = (); + fn on_start(&mut self, ctx: &Ctx) { + ctx.join_group("auto-joined"); + } + fn handle(&mut self, _ctx: &Ctx, _msg: Ping) {} + } + + let rt = Runtime::new(RuntimeConfig::default()); + let a = rt.spawn(GroupJoinerActor).unwrap(); + let b = rt.spawn(GroupJoinerActor).unwrap(); + + rt.tick(); // on_start → both join "auto-joined" + + let members = rt.group_members("auto-joined"); + assert_eq!(members.len(), 2); + assert!(members.contains(&a)); + assert!(members.contains(&b)); +} + +/// Given an actor uses ctx.publish() from inside a handler, +/// when the published message is processed, +/// then all group members receive it. +#[test] +fn ctx_publish_broadcasts_from_handler() { + #[derive(Clone)] + struct BroadcastCmd { + reply_to: ActorAddress, + } + + struct BroadcasterActor; + + impl ActorInterface for BroadcasterActor { + type Incoming = BroadcastCmd; + type Response = (); + fn on_start(&mut self, ctx: &Ctx) { + ctx.join_group("broadcast-test"); + } + fn handle(&mut self, ctx: &Ctx, msg: BroadcastCmd) { + ctx.publish("broadcast-test", Ping { reply_to: msg.reply_to }); + } + } + + let rt = Runtime::new(RuntimeConfig::default()); + let inbox = rt.new_inbox::().unwrap(); + + // Spawn 3 PingPongActors and one Broadcaster, all in the same group + let _p1 = rt.spawn(PingPongActor).unwrap(); + let _p2 = rt.spawn(PingPongActor).unwrap(); + rt.join_group(_p1, "broadcast-test"); + rt.join_group(_p2, "broadcast-test"); + + let broadcaster = rt.spawn(BroadcasterActor).unwrap(); + + rt.tick(); // on_start (broadcaster joins group too) + + // Send BroadcastCmd to broadcaster + rt.send_to(broadcaster, BroadcastCmd { reply_to: *inbox.addr() }).unwrap(); + rt.tick(); // broadcaster handles → publish Ping to all 3 members (including self) + rt.tick(); // PingPong actors handle Ping → send Pong to inbox + // Broadcaster also gets the Ping but it expects BroadcastCmd, so type mismatch (silent) + + // At least 2 Pongs from the PingPongActors + let mut pong_count = 0; + while inbox.try_recv().is_some() { + pong_count += 1; + } + assert!(pong_count >= 2, "at least 2 PingPong members should reply, got {pong_count}"); +}