diff --git a/lib/components/fabro-agent/src/subagent.rs b/lib/components/fabro-agent/src/subagent.rs index f4d6c1b84..476094a94 100644 --- a/lib/components/fabro-agent/src/subagent.rs +++ b/lib/components/fabro-agent/src/subagent.rs @@ -89,16 +89,27 @@ pub enum SubAgentStatus { const SUBAGENT_SHUTDOWN_GRACE: Duration = Duration::from_secs(5); struct SubAgent { - status: watch::Sender, - cleanup_done: watch::Sender, - cleanup_started: bool, - monitor_task: Option>, - event_forwarder: Option>, - cleanup_task: Option>, - child_abort_handle: AbortHandle, - followup_queue: Arc>>, - cancel_token: CancellationToken, - depth: usize, + status: watch::Sender, + cleanup_done: watch::Sender, + cleanup_started: bool, + monitor_task: Option>, + event_forwarder: Option>, + cleanup_task: Option>, + child_abort_handle: AbortHandle, + followup_queue: Arc>>, + cancel_token: CancellationToken, + depth: usize, + /// Task description, set when the parent should receive this child's + /// terminal result automatically. Cleared once the result is delivered, + /// the parent retrieves it explicitly, or the agent is shut down. + /// + /// Keeping this beside the status it is delivered with means a + /// notification cannot be registered before -- or suppressed after -- the + /// state it describes: there is only one lock and one ordering. + parent_notification: Option, + /// Spawn order, so a batch is delivered oldest-first rather than in + /// whatever order the map happens to iterate. + spawn_seq: u64, } impl Drop for SubAgent { @@ -119,114 +130,8 @@ impl Drop for SubAgent { #[derive(Default)] struct SupervisorState { - agents: HashMap, -} - -#[derive(Default)] -struct ParentNotificationState { - pending: HashMap, - ready: VecDeque, -} - -struct ParentNotificationHub { - state: Mutex, - changed: watch::Sender, -} - -impl ParentNotificationHub { - fn new() -> Self { - let (changed, _) = watch::channel(0); - Self { - state: Mutex::new(ParentNotificationState::default()), - changed, - } - } - - fn register(&self, agent_id: String, description: String) { - self.state - .lock() - .expect("parent notification lock poisoned") - .pending - .insert(agent_id, description); - self.signal(); - } - - fn complete(&self, agent_id: &str, result: Result) { - { - let mut state = self - .state - .lock() - .expect("parent notification lock poisoned"); - let Some(description) = state.pending.remove(agent_id) else { - return; - }; - state.ready.push_back(SubAgentParentNotification { - agent_id: agent_id.to_string(), - description, - result, - }); - } - self.signal(); - } - - fn suppress(&self, agent_id: &str) { - let changed = { - let mut state = self - .state - .lock() - .expect("parent notification lock poisoned"); - let removed_pending = state.pending.remove(agent_id).is_some(); - let ready_len = state.ready.len(); - state - .ready - .retain(|notification| notification.agent_id != agent_id); - removed_pending || state.ready.len() != ready_len - }; - if changed { - self.signal(); - } - } - - async fn next_batch( - &self, - cancel: &CancellationToken, - ) -> Result>, Error> { - let mut changed = self.changed.subscribe(); - loop { - { - let mut state = self - .state - .lock() - .expect("parent notification lock poisoned"); - if !state.ready.is_empty() { - return Ok(Some(state.ready.drain(..).collect())); - } - if state.pending.is_empty() { - return Ok(None); - } - } - - tokio::select! { - biased; - () = cancel.cancelled() => { - return Err(Error::Interrupted(InterruptReason::Cancelled)); - } - observed = changed.changed() => { - observed.map_err(|_| { - Error::InvalidState( - "Background-agent notification observer closed unexpectedly".to_string(), - ) - })?; - } - } - } - } - - fn signal(&self) { - self.changed.send_modify(|generation| { - *generation = generation.wrapping_add(1); - }); - } + agents: HashMap, + next_spawn_seq: u64, } struct ShutdownWork { @@ -268,11 +173,20 @@ impl Drop for CleanupDoneGuard { } } +/// Wake anything parked in +/// [`SubAgentSupervisor::next_parent_notification_batch`] so it can re-evaluate +/// which children are deliverable. +fn signal_notifications(changed: &watch::Sender) { + changed.send_modify(|generation| { + *generation = generation.wrapping_add(1); + }); +} + fn spawn_result_monitor( child_task: JoinHandle>, status: watch::Sender, event_callback: Arc>>, - parent_notifications: Arc, + notifications_changed: Arc>, agent_id: String, depth: usize, ) -> JoinHandle<()> { @@ -294,6 +208,8 @@ fn spawn_result_monitor( if !committed { return; } + // The status this agent will be delivered with is now committed. + signal_notifications(¬ifications_changed); let event = match &task_result { Ok(result) => AgentEvent::SubAgentCompleted { @@ -315,7 +231,6 @@ fn spawn_result_monitor( if let Some(callback) = callback { callback(SubAgentCallbackEvent::Lifecycle(event)); } - parent_notifications.complete(&agent_id, task_result); }) } @@ -326,10 +241,10 @@ fn spawn_result_monitor( /// happen after the guard has been released. #[derive(Clone)] pub struct SubAgentSupervisor { - state: Arc>, - max_depth: usize, - event_callback: Arc>>, - parent_notifications: Arc, + state: Arc>, + max_depth: usize, + event_callback: Arc>>, + notifications_changed: Arc>, } impl SubAgentSupervisor { @@ -339,7 +254,7 @@ impl SubAgentSupervisor { state: Arc::new(Mutex::new(SupervisorState::default())), max_depth, event_callback: Arc::new(RwLock::new(None)), - parent_notifications: Arc::new(ParentNotificationHub::new()), + notifications_changed: Arc::new(watch::channel(0).0), } } @@ -474,23 +389,15 @@ impl SubAgentSupervisor { child_task, status.clone(), Arc::clone(&self.event_callback), - Arc::clone(&self.parent_notifications), + Arc::clone(&self.notifications_changed), agent_id.clone(), child_depth, ); - // Register before the agent becomes discoverable in `state.agents`. - // Once it is, a concurrent `shutdown_all` can suppress and close it; - // registering afterwards would leave a pending entry that the monitor - // never completes (it early-returns for a non-Running agent), and - // `next_batch` would then never report the queue as drained. Nothing - // can complete this registration before `start_tx.send(())` below. - if let Some(description) = parent_notification_description { - self.parent_notifications - .register(agent_id.clone(), description); - } { let mut state = self.state.lock().expect("subagent state lock poisoned"); + let spawn_seq = state.next_spawn_seq; + state.next_spawn_seq = state.next_spawn_seq.saturating_add(1); state.agents.insert(agent_id.clone(), SubAgent { status, cleanup_done, @@ -502,8 +409,11 @@ impl SubAgentSupervisor { followup_queue, cancel_token, depth: child_depth, + parent_notification: parent_notification_description, + spawn_seq, }); } + signal_notifications(&self.notifications_changed); self.emit_event(AgentEvent::SubAgentSpawned { agent_id: agent_id.clone(), @@ -587,11 +497,20 @@ impl SubAgentSupervisor { } } - /// Stop automatic delivery for an agent whose result the parent explicitly - /// retrieved. Removes a result that may already have raced into the ready - /// queue. + /// Stop automatic delivery for an agent whose result the parent retrieved + /// explicitly. pub(crate) fn suppress_parent_notification(&self, agent_id: &str) { - self.parent_notifications.suppress(agent_id); + let cleared = { + let mut state = self.state.lock().expect("subagent state lock poisoned"); + state + .agents + .get_mut(agent_id) + .and_then(|agent| agent.parent_notification.take()) + .is_some() + }; + if cleared { + signal_notifications(&self.notifications_changed); + } } /// Wait until all currently-ready background results can be delivered in @@ -616,7 +535,68 @@ impl SubAgentSupervisor { &self, cancel: &CancellationToken, ) -> Result>, Error> { - self.parent_notifications.next_batch(cancel).await + let mut changed = self.notifications_changed.subscribe(); + loop { + { + let mut state = self.state.lock().expect("subagent state lock poisoned"); + let mut ready = Vec::new(); + let mut awaiting_result = false; + for (agent_id, agent) in &state.agents { + let Some(description) = agent.parent_notification.as_ref() else { + continue; + }; + let finished = match &*agent.status.borrow() { + SubAgentStatus::Finished(result) => Some(result.clone()), + SubAgentStatus::Running => { + awaiting_result = true; + None + } + // Being torn down, so no result is coming. Ignoring + // these is what keeps a shutdown that races delivery + // from parking the parent forever. + SubAgentStatus::Closing | SubAgentStatus::Closed => None, + }; + if let Some(result) = finished { + ready.push((agent.spawn_seq, SubAgentParentNotification { + agent_id: agent_id.clone(), + description: description.clone(), + result, + })); + } + } + + if !ready.is_empty() { + ready.sort_by_key(|(spawn_seq, _)| *spawn_seq); + let batch: Vec<_> = ready + .into_iter() + .map(|(_, notification)| notification) + .collect(); + for notification in &batch { + if let Some(agent) = state.agents.get_mut(¬ification.agent_id) { + agent.parent_notification = None; + } + } + return Ok(Some(batch)); + } + if !awaiting_result { + return Ok(None); + } + } + + tokio::select! { + biased; + () = cancel.cancelled() => { + return Err(Error::Interrupted(InterruptReason::Cancelled)); + } + observed = changed.changed() => { + observed.map_err(|_| { + Error::InvalidState( + "Background-agent notification observer closed unexpectedly".to_string(), + ) + })?; + } + } + } } #[cfg(test)] @@ -666,6 +646,11 @@ impl SubAgentSupervisor { } }; + // Shutdown is committed, so this child's result will never reach the + // parent. The early returns above leave the notification intact, so a + // rejected shutdown cannot discard a result the parent is owed. + agent.parent_notification = None; + if agent.cleanup_started { return Ok(ShutdownDisposition::Follow(agent.cleanup_done.subscribe())); } @@ -772,10 +757,7 @@ impl SubAgentSupervisor { async fn ensure_closed(&self, agent_id: &str) -> Result<(), Error> { let disposition = self.begin_shutdown(agent_id, false)?; - // Only once shutdown is committed. Suppressing before `begin_shutdown` - // would also discard the result of an agent that had already finished, - // which rejects the shutdown but had a delivery pending. - self.parent_notifications.suppress(agent_id); + signal_notifications(&self.notifications_changed); let cleanup_done = match disposition { ShutdownDisposition::Lead(work) => self.spawn_shutdown(work), ShutdownDisposition::Follow(cleanup_done) => cleanup_done, @@ -788,7 +770,7 @@ impl SubAgentSupervisor { /// Strict user-facing close: only a currently running child may be closed. pub async fn close_agent(&self, agent_id: &str) -> Result<(), Error> { let disposition = self.begin_shutdown(agent_id, true)?; - self.parent_notifications.suppress(agent_id); + signal_notifications(&self.notifications_changed); let cleanup_done = match disposition { ShutdownDisposition::Lead(work) => self.spawn_shutdown(work), ShutdownDisposition::Follow(_) | ShutdownDisposition::Done => { @@ -857,7 +839,7 @@ impl SubAgentSupervisor { child_task, status.clone(), Arc::clone(&self.event_callback), - Arc::clone(&self.parent_notifications), + Arc::clone(&self.notifications_changed), agent_id.clone(), depth, ); @@ -873,6 +855,8 @@ impl SubAgentSupervisor { event_forwarder, cleanup_task: None, child_abort_handle, + parent_notification: None, + spawn_seq: 0, followup_queue: Arc::new(Mutex::new(VecDeque::new())), cancel_token, depth, @@ -1072,67 +1056,155 @@ mod tests { assert!(manager.is_empty()); } - #[tokio::test] - async fn parent_notifications_are_exactly_once_and_xml_escaped() { - let hub = ParentNotificationHub::new(); - hub.register("agent<&".to_string(), "Review & tests".to_string()); - let result = Ok(SubAgentResult { - output: "done & \"verified\"".to_string(), - success: true, - turns_used: 2, - }); - hub.complete("agent<&", result.clone()); - hub.complete("agent<&", result); + #[test] + fn parent_notification_envelope_escapes_xml() { + let envelope = format_parent_notification_batch(&[SubAgentParentNotification { + agent_id: "agent<&".to_string(), + description: "Review & tests".to_string(), + result: Ok(SubAgentResult { + output: "done & \"verified\"".to_string(), + success: true, + turns_used: 2, + }), + }]); - let notifications = hub - .next_batch(&CancellationToken::new()) - .await - .unwrap() - .unwrap(); - assert_eq!(notifications.len(), 1); - let envelope = format_parent_notification_batch(¬ifications); assert!(envelope.contains("completed")); assert!(envelope.contains("agent<&")); assert!(envelope.contains("Review <core> & tests")); assert!( envelope.contains("done <safely> & "verified"") ); - assert!( - hub.next_batch(&CancellationToken::new()) - .await - .unwrap() - .is_none() - ); } #[tokio::test] - async fn suppress_removes_pending_and_ready_parent_notifications() { - let hub = ParentNotificationHub::new(); - hub.register("pending".to_string(), "Pending".to_string()); - hub.suppress("pending"); + async fn a_finished_agent_is_delivered_to_the_parent_exactly_once() { + let supervisor = SubAgentSupervisor::new(3); + let child = make_session(vec![text_response("child result")]).await; + let agent_id = supervisor + .spawn_with_parent_notification( + child, + "task".to_string(), + "Inspect the module".to_string(), + 0, + ) + .unwrap(); + supervisor + .wait_with_cancel(&agent_id, &CancellationToken::new()) + .await + .unwrap(); + + let batch = supervisor + .next_parent_notification_batch(&CancellationToken::new()) + .await + .unwrap() + .expect("the finished child must be delivered"); + assert_eq!(batch.len(), 1); + assert_eq!(batch[0].agent_id, agent_id); + assert_eq!(batch[0].description, "Inspect the module"); + + // The status stays `Finished`, so re-delivery is prevented by clearing + // the registration rather than by consuming the result. assert!( - hub.next_batch(&CancellationToken::new()) + supervisor + .next_parent_notification_batch(&CancellationToken::new()) .await .unwrap() .is_none() ); - hub.register("ready".to_string(), "Ready".to_string()); - hub.complete( - "ready", - Ok(SubAgentResult { - output: "done".to_string(), - success: true, - turns_used: 1, - }), - ); - hub.suppress("ready"); + supervisor.shutdown_all().await; + } + + #[tokio::test] + async fn batches_are_delivered_in_spawn_order() { + let supervisor = SubAgentSupervisor::new(3); + let mut ids = Vec::new(); + for index in 0..3 { + let child = make_session(vec![text_response("done")]).await; + ids.push( + supervisor + .spawn_with_parent_notification( + child, + format!("task {index}"), + format!("Task {index}"), + 0, + ) + .unwrap(), + ); + } + for id in &ids { + supervisor + .wait_with_cancel(id, &CancellationToken::new()) + .await + .unwrap(); + } + + let batch = supervisor + .next_parent_notification_batch(&CancellationToken::new()) + .await + .unwrap() + .expect("all three children must be delivered together"); + let delivered: Vec<_> = batch.iter().map(|n| n.agent_id.clone()).collect(); + assert_eq!(delivered, ids); + + supervisor.shutdown_all().await; + } + + #[tokio::test] + async fn suppressing_before_completion_stops_delivery() { + let supervisor = SubAgentSupervisor::new(3); + let child = make_session(vec![text_response("child result")]).await; + let agent_id = supervisor + .spawn_with_parent_notification( + child, + "task".to_string(), + "Inspect the module".to_string(), + 0, + ) + .unwrap(); + + supervisor.suppress_parent_notification(&agent_id); + supervisor + .wait_with_cancel(&agent_id, &CancellationToken::new()) + .await + .unwrap(); + assert!( - hub.next_batch(&CancellationToken::new()) + supervisor + .next_parent_notification_batch(&CancellationToken::new()) .await .unwrap() .is_none() ); + + supervisor.shutdown_all().await; + } + + #[tokio::test] + async fn closing_a_running_agent_stops_delivery_without_parking_the_parent() { + let supervisor = SubAgentSupervisor::new(3); + let child = make_session(vec![text_response("child result")]).await; + let agent_id = supervisor + .spawn_with_parent_notification( + child, + "task".to_string(), + "Inspect the module".to_string(), + 0, + ) + .unwrap(); + + supervisor.close_agent(&agent_id).await.unwrap(); + + // Must resolve rather than wait for a result that will never arrive. + assert!( + supervisor + .next_parent_notification_batch(&CancellationToken::new()) + .await + .unwrap() + .is_none() + ); + + supervisor.shutdown_all().await; } #[tokio::test]