From 8c19b1e94a3e83f47cb2e51ab5c47352505d0a37 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 1 Aug 2026 10:13:27 -0400 Subject: [PATCH] fix: keep forwarding a reused child's events after a broadcast lag The subagent event forwarder left its loop on any `recv` error, including `Lagged`. A lagged broadcast receiver stays usable, so one transient lag silenced the child for the rest of its life while the task completed normally and shutdown joined it without noticing. Session reuse widens that window from a single turn to the whole parent session. Also borrow each result's output when rendering a parent notification instead of cloning it. Co-Authored-By: Claude Opus 5 (1M context) --- lib/components/fabro-agent/src/subagent.rs | 24 +++++++++++++++++----- 1 file changed, 19 insertions(+), 5 deletions(-) diff --git a/lib/components/fabro-agent/src/subagent.rs b/lib/components/fabro-agent/src/subagent.rs index d6f29f4de..2a90384fc 100644 --- a/lib/components/fabro-agent/src/subagent.rs +++ b/lib/components/fabro-agent/src/subagent.rs @@ -1,3 +1,4 @@ +use std::borrow::Cow; use std::collections::{HashMap, VecDeque}; use std::sync::{Arc, Mutex, RwLock, Weak}; use std::time::Duration; @@ -6,7 +7,7 @@ use fabro_llm::types::ToolDefinition; use fabro_types::INITIAL_SUBAGENT_GENERATION; use fabro_util::error as util_error; use futures::future; -use tokio::sync::{mpsc, oneshot, watch}; +use tokio::sync::{broadcast, mpsc, oneshot, watch}; use tokio::task::{AbortHandle, JoinHandle}; use tokio::time::{Instant, timeout_at}; use tokio_util::sync::CancellationToken; @@ -48,9 +49,14 @@ fn format_parent_notification_batch(notifications: &[SubAgentParentNotification] .iter() .map(|notification| { let (status, result) = match ¬ification.result { - Ok(result) if result.success => ("completed", result.output.clone()), - Ok(result) => ("failed", result.output.clone()), - Err(error) => ("failed", util_error::collect_chain(error).join(": ")), + Ok(result) if result.success => { + ("completed", Cow::Borrowed(result.output.as_str())) + } + Ok(result) => ("failed", Cow::Borrowed(result.output.as_str())), + Err(error) => ( + "failed", + Cow::Owned(util_error::collect_chain(error).join(": ")), + ), }; format!( "\n {}\n {status}\n \ @@ -621,7 +627,15 @@ impl SubAgentSupervisor { let mut rx = session.subscribe(); let callback = Arc::clone(&self.event_callback); Some(tokio::spawn(async move { - while let Ok(event) = rx.recv().await { + loop { + let event = match rx.recv().await { + Ok(event) => event, + // A lagged receiver stays usable, and a reused child + // forwards for the whole parent session. Giving up here + // would silence the child for the rest of its life. + Err(broadcast::error::RecvError::Lagged(_)) => continue, + Err(broadcast::error::RecvError::Closed) => break, + }; // Skip streaming / noise events if event.event.is_streaming_noise() || matches!(