use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use std::time::Instant;
use tokio::sync::watch;
use super::config::MultiAgentConfig;
use super::path::AgentPath;
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum AgentStatus {
Queued,
Running,
Done,
Closed,
Paused,
}
impl AgentStatus {
pub fn name(&self) -> &'static str {
match self {
Self::Queued => "queued",
Self::Running => "running",
Self::Done => "done",
Self::Closed => "closed",
Self::Paused => "paused",
}
}
}
pub fn derive_status(registered: bool, in_flight: bool, queue_len: usize) -> AgentStatus {
match (registered, in_flight, queue_len) {
(false, _, _) => AgentStatus::Closed,
(_, true, _) => AgentStatus::Running,
(_, false, 0) => AgentStatus::Done,
(_, false, _) => AgentStatus::Queued,
}
}
#[derive(Clone, Debug)]
pub struct AgentEntry {
pub path: AgentPath,
pub depth: i32,
pub queue_len: usize,
pub in_flight: bool,
pub running_since: Option<Instant>,
pub closing: bool,
pub results_posted: usize,
pub results_handed_over: usize,
pub tool_calls: usize,
pub last_activity: Option<Instant>,
pub task: Option<String>,
}
impl AgentEntry {
pub fn status(&self) -> AgentStatus {
if self.closing {
return AgentStatus::Closed;
}
derive_status(true, self.in_flight, self.queue_len)
}
}
#[derive(Clone, Debug)]
pub struct AgentLifecycleEvent {
pub path: AgentPath,
pub from: AgentStatus,
pub to: AgentStatus,
pub reason: &'static str,
}
#[derive(Clone, Debug, Default)]
pub struct RegistrySnapshot {
pub agents: Vec<AgentSnapshot>,
}
#[derive(Clone, Debug)]
pub struct AgentSnapshot {
pub path: String,
pub status: String,
pub running_secs: Option<u64>,
pub last_activity_secs: Option<u64>,
pub tool_calls: usize,
pub task: Option<String>,
pub pending_results: usize,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum SpawnError {
MaxAgentsReached { max: usize },
DepthLimitReached { max: i32, attempted: i32 },
AlreadyExists,
}
impl std::fmt::Display for SpawnError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::MaxAgentsReached { max } => {
write!(f, "max agents reached (limit: {})", max)
}
Self::DepthLimitReached { max, attempted } => {
write!(
f,
"agent depth limit reached (max: {}, attempted: {})",
max, attempted
)
}
Self::AlreadyExists => {
write!(f, "agent with this path already exists")
}
}
}
}
pub struct AgentRegistry {
config: MultiAgentConfig,
agents: HashMap<AgentPath, AgentEntry>,
events: VecDeque<AgentLifecycleEvent>,
snapshot_tx: watch::Sender<Arc<RegistrySnapshot>>,
}
const EVENT_RING_CAP: usize = 2048;
impl AgentRegistry {
pub fn new(config: MultiAgentConfig) -> Self {
let (snapshot_tx, _) = watch::channel(Arc::new(RegistrySnapshot::default()));
Self {
config,
agents: HashMap::new(),
events: VecDeque::new(),
snapshot_tx,
}
}
pub fn can_spawn(&self, depth: i32) -> Result<(), SpawnError> {
if self.config.enabled && self.agents.len() >= self.config.max_sub_agents {
return Err(SpawnError::MaxAgentsReached {
max: self.config.max_sub_agents,
});
}
if self.config.enabled && depth > self.config.max_agent_depth {
return Err(SpawnError::DepthLimitReached {
max: self.config.max_agent_depth,
attempted: depth,
});
}
Ok(())
}
pub fn register(&mut self, path: &AgentPath, depth: i32) -> Result<(), SpawnError> {
self.can_spawn(depth)?;
if self.agents.contains_key(path) {
return Err(SpawnError::AlreadyExists);
}
self.agents.insert(
path.clone(),
AgentEntry {
path: path.clone(),
depth,
queue_len: 0,
in_flight: false,
running_since: None,
closing: false,
results_posted: 0,
results_handed_over: 0,
tool_calls: 0,
last_activity: None,
task: None,
},
);
Ok(())
}
pub fn close(&mut self, path: &AgentPath) -> Option<AgentEntry> {
let entry = self.agents.remove(path)?;
self.push_event(AgentLifecycleEvent {
path: path.clone(),
from: entry.status(),
to: AgentStatus::Closed,
reason: "unregistered",
});
self.publish_snapshot();
Some(entry)
}
fn transition(
&mut self,
path: &AgentPath,
reason: &'static str,
apply: impl FnOnce(&mut AgentEntry),
) -> bool {
let Some(entry) = self.agents.get_mut(path) else {
return false;
};
let from = entry.status();
apply(entry);
let to = entry.status();
if from != to {
self.push_event(AgentLifecycleEvent {
path: path.clone(),
from,
to,
reason,
});
self.publish_snapshot();
}
true
}
pub fn note_enqueued(&mut self, path: &AgentPath) -> bool {
self.transition(path, "task_enqueued", |e| {
e.queue_len += 1;
})
}
pub fn note_dequeued(&mut self, path: &AgentPath) -> bool {
self.transition(path, "task_dequeued", |e| {
e.queue_len = e.queue_len.saturating_sub(1);
e.in_flight = true;
e.running_since = Some(Instant::now());
})
}
pub fn note_posted(&mut self, path: &AgentPath) -> bool {
self.transition(path, "result_posted", |e| {
e.in_flight = false;
e.results_posted += 1;
e.running_since = None;
})
}
pub fn note_closing(&mut self, path: &AgentPath) -> bool {
self.transition(path, "closing", |e| {
e.closing = true;
e.in_flight = false;
e.queue_len = 0;
e.running_since = None;
})
}
pub fn note_send_failed(&mut self, path: &AgentPath) -> bool {
self.transition(path, "task_send_failed", |e| {
e.queue_len = e.queue_len.saturating_sub(1);
})
}
pub fn note_batch_handed_over(&mut self, path: &AgentPath) -> bool {
let Some(entry) = self.agents.get_mut(path) else {
return false;
};
entry.results_handed_over = entry.results_posted;
self.publish_snapshot();
true
}
pub fn set_task(&mut self, path: &AgentPath, task: String) -> bool {
match self.agents.get_mut(path) {
Some(entry) => {
if entry.task.is_none() {
entry.task = Some(task);
}
true
}
None => false,
}
}
pub fn touch(&mut self, path: &AgentPath) -> bool {
match self.agents.get_mut(path) {
Some(entry) => {
entry.last_activity = Some(Instant::now());
true
}
None => false,
}
}
pub fn record_tool_call(&mut self, path: &AgentPath) -> bool {
match self.agents.get_mut(path) {
Some(entry) => {
entry.tool_calls += 1;
entry.last_activity = Some(Instant::now());
true
}
None => false,
}
}
pub fn busy_count(&self) -> usize {
self.agents
.values()
.filter(|e| !e.closing && (e.in_flight || e.queue_len > 0))
.count()
}
pub fn quiescent(&self) -> bool {
let all_quiet = self
.agents
.values()
.all(|e| e.closing || (!e.in_flight && e.queue_len == 0 && e.results_posted >= 1));
if !all_quiet {
let blocking: Vec<String> = self
.agents
.values()
.filter(|e| !e.closing && (e.in_flight || e.queue_len > 0 || e.results_posted < 1))
.map(|e| {
format!(
"{}: in_flight={}, queue_len={}, results_posted={}",
e.path, e.in_flight, e.queue_len, e.results_posted
)
})
.collect();
tracing::info!(
blocking_agents = ?blocking,
"registry: quiescent=false"
);
}
all_quiet
}
pub fn get(&self, path: &AgentPath) -> Option<&AgentEntry> {
self.agents.get(path)
}
pub fn list(&self) -> Vec<&AgentEntry> {
let mut entries: Vec<&AgentEntry> = self.agents.values().collect();
entries.sort_by(|a, b| a.path.cmp(&b.path));
entries
}
pub fn count(&self) -> usize {
self.agents.len()
}
pub fn is_empty(&self) -> bool {
self.agents.is_empty()
}
pub fn contains(&self, path: &AgentPath) -> bool {
self.agents.contains_key(path)
}
pub fn snapshot(&self) -> RegistrySnapshot {
RegistrySnapshot {
agents: self
.list()
.into_iter()
.map(|e| AgentSnapshot {
path: e.path.to_string(),
status: e.status().name().to_string(),
running_secs: if e.status() == AgentStatus::Running {
e.running_since.map(|t| t.elapsed().as_secs())
} else {
None
},
last_activity_secs: e.last_activity.map(|t| t.elapsed().as_secs()),
tool_calls: e.tool_calls,
task: e.task.clone(),
pending_results: e.results_posted.saturating_sub(e.results_handed_over),
})
.collect(),
}
}
pub fn subscribe(&self) -> watch::Receiver<Arc<RegistrySnapshot>> {
self.snapshot_tx.subscribe()
}
pub fn latest_snapshot(&self) -> Arc<RegistrySnapshot> {
self.snapshot_tx.borrow().clone()
}
pub fn recent_events(&self, max: usize) -> Vec<AgentLifecycleEvent> {
let skip = self.events.len().saturating_sub(max);
self.events.iter().skip(skip).cloned().collect()
}
fn push_event(&mut self, event: AgentLifecycleEvent) {
if self.events.len() >= EVENT_RING_CAP {
self.events.pop_front();
}
self.events.push_back(event);
}
fn publish_snapshot(&mut self) {
let snapshot = Arc::new(self.snapshot());
self.snapshot_tx.send_replace(snapshot);
}
pub fn config(&self) -> &MultiAgentConfig {
&self.config
}
}
#[cfg(test)]
mod tests {
use super::*;
fn test_config() -> MultiAgentConfig {
MultiAgentConfig::enabled()
}
fn test_path(name: &str) -> AgentPath {
AgentPath::root().join(name)
}
fn registered(config: MultiAgentConfig, name: &str) -> (AgentRegistry, AgentPath) {
let mut reg = AgentRegistry::new(config);
let path = test_path(name);
reg.register(&path, 1).unwrap();
(reg, path)
}
#[test]
fn derive_status_is_a_pure_function() {
use AgentStatus::*;
assert_eq!(derive_status(false, true, 3), Closed);
assert_eq!(derive_status(false, false, 0), Closed);
assert_eq!(derive_status(true, true, 0), Running);
assert_eq!(derive_status(true, true, 5), Running);
assert_eq!(derive_status(true, false, 0), Done);
assert_eq!(derive_status(true, false, 1), Queued);
assert_eq!(derive_status(true, false, 9), Queued);
}
#[test]
fn status_names_are_wire_stable() {
assert_eq!(AgentStatus::Queued.name(), "queued");
assert_eq!(AgentStatus::Running.name(), "running");
assert_eq!(AgentStatus::Done.name(), "done");
assert_eq!(AgentStatus::Closed.name(), "closed");
assert_eq!(AgentStatus::Paused.name(), "paused");
}
#[test]
fn register_and_close() {
let mut reg = AgentRegistry::new(test_config());
let path = test_path("worker");
assert!(reg.register(&path, 1).is_ok());
assert_eq!(reg.count(), 1);
assert!(reg.contains(&path));
let entry = reg.get(&path).unwrap();
assert_eq!(entry.status(), AgentStatus::Done);
assert_eq!(entry.depth, 1);
assert_eq!(entry.queue_len, 0);
assert!(!entry.in_flight);
assert_eq!(entry.results_posted, 0);
assert_eq!(entry.tool_calls, 0);
assert!(entry.last_activity.is_none());
let closed = reg.close(&path).unwrap();
assert_eq!(closed.path, path);
assert_eq!(reg.count(), 0);
assert!(!reg.contains(&path));
}
#[test]
fn duplicate_register_fails() {
let mut reg = AgentRegistry::new(test_config());
let path = test_path("worker");
assert!(reg.register(&path, 1).is_ok());
assert_eq!(
reg.register(&path, 1).unwrap_err(),
SpawnError::AlreadyExists
);
}
#[test]
fn max_agents_limit() {
let config = MultiAgentConfig::with_limits(2, 1);
let mut reg = AgentRegistry::new(config);
assert!(reg.register(&test_path("a"), 1).is_ok());
assert!(reg.register(&test_path("b"), 1).is_ok());
assert_eq!(
reg.register(&test_path("c"), 1).unwrap_err(),
SpawnError::MaxAgentsReached { max: 2 }
);
}
#[test]
fn depth_limit() {
let config = MultiAgentConfig::with_limits(8, 1);
let reg = AgentRegistry::new(config);
assert!(reg.can_spawn(1).is_ok());
assert_eq!(
reg.can_spawn(2).unwrap_err(),
SpawnError::DepthLimitReached {
max: 1,
attempted: 2
}
);
}
#[test]
fn facts_drive_the_whole_lifecycle() {
let (mut reg, path) = registered(test_config(), "worker");
assert_eq!(reg.get(&path).unwrap().status(), AgentStatus::Done);
assert!(reg.note_enqueued(&path));
assert_eq!(reg.get(&path).unwrap().status(), AgentStatus::Queued);
assert!(reg.note_enqueued(&path));
assert_eq!(reg.get(&path).unwrap().status(), AgentStatus::Queued);
assert!(reg.note_dequeued(&path));
assert_eq!(reg.get(&path).unwrap().status(), AgentStatus::Running);
assert!(reg.get(&path).unwrap().running_since.is_some());
assert_eq!(reg.get(&path).unwrap().queue_len, 1);
assert!(reg.note_posted(&path));
assert_eq!(reg.get(&path).unwrap().status(), AgentStatus::Queued);
assert!(reg.get(&path).unwrap().running_since.is_none());
assert_eq!(reg.get(&path).unwrap().results_posted, 1);
assert!(reg.note_dequeued(&path));
assert_eq!(reg.get(&path).unwrap().status(), AgentStatus::Running);
assert!(reg.note_posted(&path));
assert_eq!(reg.get(&path).unwrap().status(), AgentStatus::Done);
assert_eq!(reg.get(&path).unwrap().results_posted, 2);
assert_eq!(reg.count(), 1);
}
#[test]
fn close_frees_quota() {
let config = MultiAgentConfig::with_limits(1, 1);
let mut reg = AgentRegistry::new(config);
let path = test_path("worker");
reg.register(&path, 1).unwrap();
assert_eq!(reg.count(), 1);
assert!(reg.can_spawn(1).is_err());
reg.close(&path);
assert_eq!(reg.count(), 0);
assert!(reg.can_spawn(1).is_ok());
}
#[test]
fn list_sorted() {
let mut reg = AgentRegistry::new(test_config());
reg.register(&test_path("b"), 1).unwrap();
reg.register(&test_path("a"), 1).unwrap();
let list = reg.list();
assert_eq!(list.len(), 2);
assert_eq!(list[0].path.name(), "a");
assert_eq!(list[1].path.name(), "b");
}
#[test]
fn tool_call_counter_is_monotonic() {
let (mut reg, path) = registered(test_config(), "worker");
for _ in 0..3 {
assert!(reg.record_tool_call(&path));
}
assert_eq!(reg.get(&path).unwrap().tool_calls, 3);
assert!(reg.get(&path).unwrap().last_activity.is_some());
assert!(reg.record_tool_call(&path));
assert_eq!(reg.get(&path).unwrap().tool_calls, 4);
assert!(!reg.record_tool_call(&test_path("ghost")));
}
#[test]
fn touch_marks_activity_without_counting() {
let (mut reg, path) = registered(test_config(), "worker");
assert!(reg.touch(&path));
assert!(reg.get(&path).unwrap().last_activity.is_some());
assert_eq!(reg.get(&path).unwrap().tool_calls, 0);
assert!(!reg.touch(&test_path("ghost")));
}
#[test]
fn busy_count_tracks_outstanding_work() {
let (mut reg, a) = registered(test_config(), "a");
let b = test_path("b");
reg.register(&b, 1).unwrap();
assert_eq!(reg.busy_count(), 0);
assert!(reg.note_enqueued(&a));
assert!(reg.note_dequeued(&a));
assert!(reg.note_enqueued(&b));
assert_eq!(reg.busy_count(), 2);
assert!(reg.note_posted(&a));
assert_eq!(reg.busy_count(), 1);
assert!(reg.note_dequeued(&b));
assert!(reg.note_posted(&b));
assert_eq!(reg.busy_count(), 0);
}
#[test]
fn quiescence_requires_delivery_from_every_registered_agent() {
let (mut reg, a) = registered(test_config(), "a");
let b = test_path("b");
reg.register(&b, 1).unwrap();
assert!(reg.note_enqueued(&a));
assert!(reg.note_dequeued(&a));
assert!(reg.note_posted(&a));
assert!(!reg.quiescent(), "b never delivered — spawn→send window");
assert!(reg.note_enqueued(&b));
assert!(reg.note_dequeued(&b));
assert!(!reg.quiescent(), "b is in flight");
assert!(reg.note_posted(&b));
assert!(reg.quiescent(), "everyone idle and delivered");
assert!(reg.note_enqueued(&b));
assert!(!reg.quiescent(), "queued work blocks the batch");
}
#[test]
fn batch_handover_closes_and_reopens_the_delivery_gap() {
let (mut reg, path) = registered(test_config(), "worker");
reg.note_enqueued(&path);
reg.note_dequeued(&path);
reg.note_posted(&path);
let snap = |reg: &AgentRegistry| reg.snapshot().agents[0].pending_results;
assert_eq!(snap(®), 1, "posted but no batch has fired yet");
assert!(reg.note_batch_handed_over(&path));
assert_eq!(snap(®), 0, "handed over — no longer pending");
reg.note_enqueued(&path);
reg.note_dequeued(&path);
reg.note_posted(&path);
assert_eq!(snap(®), 1, "new post is pending until the next batch");
assert!(!reg.note_batch_handed_over(&test_path("ghost")));
assert_eq!(reg.get(&path).unwrap().results_handed_over, 1);
}
#[test]
fn note_on_unknown_path_is_a_noop() {
let mut reg = AgentRegistry::new(test_config());
assert!(!reg.note_enqueued(&test_path("ghost")));
assert!(!reg.note_dequeued(&test_path("ghost")));
assert!(!reg.note_posted(&test_path("ghost")));
assert!(!reg.note_send_failed(&test_path("ghost")));
assert!(reg.recent_events(10).is_empty());
assert!(reg.latest_snapshot().agents.is_empty());
}
#[test]
fn send_failure_rolls_back_the_enqueue_fact() {
let (mut reg, path) = registered(test_config(), "worker");
assert!(reg.note_enqueued(&path));
assert_eq!(reg.get(&path).unwrap().queue_len, 1);
assert!(reg.note_send_failed(&path));
assert_eq!(reg.get(&path).unwrap().queue_len, 0);
assert_eq!(reg.get(&path).unwrap().status(), AgentStatus::Done);
assert!(!reg.quiescent());
assert!(reg.note_send_failed(&path));
assert_eq!(reg.get(&path).unwrap().queue_len, 0);
let (mut reg2, p2) = registered(test_config(), "solo");
reg2.note_enqueued(&p2);
reg2.note_dequeued(&p2);
reg2.note_posted(&p2);
reg2.note_enqueued(&p2);
assert!(!reg2.quiescent());
reg2.note_send_failed(&p2);
assert!(
reg2.quiescent(),
"rolled-back phantom queue must not wedge the batch"
);
}
#[test]
fn disabled_config_bypasses_spawn_limits() {
let config = MultiAgentConfig {
enabled: false,
..MultiAgentConfig::default()
};
let reg = AgentRegistry::new(config);
assert!(reg.can_spawn(999).is_ok());
}
#[tokio::test]
async fn running_secs_and_watch_receiver_track_running_state() {
let (mut reg, path) = registered(test_config(), "worker");
let mut rx = reg.subscribe();
assert!(
reg.snapshot().agents[0].running_secs.is_none(),
"not running yet"
);
reg.note_enqueued(&path);
reg.note_dequeued(&path);
rx.changed()
.await
.expect("watch receiver must be woken by the transition");
assert_eq!(rx.borrow().agents[0].running_secs, Some(0));
reg.note_posted(&path);
assert!(
reg.snapshot().agents[0].running_secs.is_none(),
"stale seconds must clear on post"
);
}
#[test]
fn events_and_snapshot_track_transitions() {
let (mut reg, path) = registered(test_config(), "worker");
assert!(reg.note_enqueued(&path));
assert!(reg.note_dequeued(&path));
assert!(reg.note_posted(&path));
let events = reg.recent_events(10);
assert_eq!(events.len(), 3);
assert_eq!(events[0].reason, "task_enqueued");
assert_eq!(events[0].from, AgentStatus::Done);
assert_eq!(events[0].to, AgentStatus::Queued);
assert_eq!(events[1].reason, "task_dequeued");
assert_eq!(events[1].to, AgentStatus::Running);
assert_eq!(events[2].reason, "result_posted");
assert_eq!(events[2].to, AgentStatus::Done);
reg.close(&path);
let events = reg.recent_events(10);
assert_eq!(events.last().unwrap().reason, "unregistered");
assert_eq!(events.last().unwrap().to, AgentStatus::Closed);
assert!(reg.snapshot().agents.is_empty());
}
#[test]
fn event_ring_is_bounded() {
let (mut reg, path) = registered(test_config(), "worker");
for _ in 0..(EVENT_RING_CAP + 100) {
reg.note_enqueued(&path);
reg.note_dequeued(&path);
reg.note_posted(&path);
}
assert_eq!(reg.events.len(), EVENT_RING_CAP);
assert_eq!(reg.recent_events(3).len(), 3);
}
#[test]
fn snapshot_reflects_derived_state() {
let (mut reg, path) = registered(test_config(), "worker");
reg.note_enqueued(&path);
reg.note_dequeued(&path);
reg.set_task(&path, "do the thing".to_string());
reg.record_tool_call(&path);
let snap = reg.snapshot();
assert_eq!(snap.agents.len(), 1);
let row = &snap.agents[0];
assert_eq!(row.path, "root/worker");
assert_eq!(row.status, "running");
assert!(row.running_secs.is_some());
assert_eq!(row.tool_calls, 1);
assert_eq!(row.task.as_deref(), Some("do the thing"));
let latest = reg.latest_snapshot();
assert_eq!(latest.agents[0].status, "running");
}
#[test]
fn spawn_error_display() {
assert_eq!(
SpawnError::MaxAgentsReached { max: 8 }.to_string(),
"max agents reached (limit: 8)"
);
assert_eq!(
SpawnError::DepthLimitReached {
max: 1,
attempted: 2
}
.to_string(),
"agent depth limit reached (max: 1, attempted: 2)"
);
assert_eq!(
SpawnError::AlreadyExists.to_string(),
"agent with this path already exists"
);
}
}