From 60e5e2fad1b4304a31fbc25522a2d5744540ea81 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 23 Mar 2026 13:18:28 -0400 Subject: [PATCH] Clean up subagents before emitting SessionEnded in Session.close() MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Matches spec shutdown order: cleanup subagents → emit SESSION_END → transition to CLOSED. Session now holds an optional SubAgentManager reference and calls close_all() during shutdown. Co-Authored-By: Claude Opus 4.6 (1M context) --- lib/crates/fabro-agent/src/session.rs | 65 ++++++++++++++++++++++++++- 1 file changed, 64 insertions(+), 1 deletion(-) diff --git a/lib/crates/fabro-agent/src/session.rs b/lib/crates/fabro-agent/src/session.rs index b2a3823e7..26f2f7f8f 100644 --- a/lib/crates/fabro-agent/src/session.rs +++ b/lib/crates/fabro-agent/src/session.rs @@ -44,6 +44,7 @@ pub struct Session { system_prompt: String, file_tracker: FileTracker, tool_env: Option>, + subagent_manager: Option>>, } impl Session { @@ -73,6 +74,7 @@ impl Session { system_prompt: String::new(), file_tracker: FileTracker::default(), tool_env: None, + subagent_manager: None, } } @@ -80,6 +82,13 @@ impl Session { self.tool_env = Some(env); } + pub fn set_subagent_manager( + &mut self, + manager: Arc>, + ) { + self.subagent_manager = Some(manager); + } + /// Initialize session by discovering project docs and capturing environment context. /// Call before `process_input`. pub async fn initialize(&mut self) { @@ -449,9 +458,15 @@ impl Session { pub fn close(&mut self) { if self.state != SessionState::Closed { - self.state = SessionState::Closed; + // Clean up subagents before emitting SessionEnded + if let Some(ref manager) = self.subagent_manager { + if let Ok(mut mgr) = manager.try_lock() { + mgr.close_all(); + } + } self.event_emitter .emit(self.id.clone(), AgentEvent::SessionEnded); + self.state = SessionState::Closed; } } @@ -2630,4 +2645,52 @@ mod tests { assert_eq!(turns.len(), 2); assert!(matches!(&turns[1], Turn::Assistant { content, .. } if content == "Fast response")); } + + #[tokio::test] + async fn close_cleans_up_subagents_before_emitting_session_ended() { + use crate::subagent::SubAgentManager; + + let mut session = make_session(vec![text_response("done")]).await; + let manager = Arc::new(tokio::sync::Mutex::new(SubAgentManager::new(3))); + + // Wire the manager's event callback to the session's emitter + manager + .lock() + .await + .set_event_callback(session.event_callback()); + + // Spawn a subagent + let child = make_session(vec![text_response("child done")]).await; + let agent_id = manager.lock().await.spawn(child, "task".into(), 0).unwrap(); + + session.set_subagent_manager(manager.clone()); + + // Collect events + let mut rx = session.subscribe(); + session.close(); + + // The subagent should have been cleaned up + assert!(manager.lock().await.get(&agent_id).is_none()); + + // Verify event ordering: SubAgentClosed before SessionEnded + let mut events = Vec::new(); + while let Ok(envelope) = rx.try_recv() { + events.push(envelope.event); + } + let closed_idx = events + .iter() + .position(|e| matches!(e, AgentEvent::SubAgentClosed { .. })); + let ended_idx = events + .iter() + .position(|e| matches!(e, AgentEvent::SessionEnded)); + assert!( + closed_idx.is_some(), + "SubAgentClosed event should be emitted" + ); + assert!(ended_idx.is_some(), "SessionEnded event should be emitted"); + assert!( + closed_idx.unwrap() < ended_idx.unwrap(), + "SubAgentClosed must come before SessionEnded" + ); + } }