From e78f5ff957896d224a9c7fd55fcc2b876834f0db Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 23 Jul 2026 18:04:58 -0400 Subject: [PATCH 1/4] Fix cancellation during subagent waits --- lib/crates/fabro-agent/src/session.rs | 161 +++++++++++++++++++- lib/crates/fabro-agent/src/subagent.rs | 201 ++++++++++++++++++------- 2 files changed, 308 insertions(+), 54 deletions(-) diff --git a/lib/crates/fabro-agent/src/session.rs b/lib/crates/fabro-agent/src/session.rs index 34986a29d..0d0861c11 100644 --- a/lib/crates/fabro-agent/src/session.rs +++ b/lib/crates/fabro-agent/src/session.rs @@ -2073,7 +2073,7 @@ mod tests { use super::*; use crate::config::{ToolAccess, ToolAccessPolicy, ToolApprovalAdapter, ToolExposureMode}; use crate::skills::{Skill, make_use_skill_tool}; - use crate::subagent::SubAgentStatus; + use crate::subagent::{SubAgentStatus, make_wait_tool}; use crate::test_support::*; use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry, ToolSource}; @@ -4703,6 +4703,165 @@ mod tests { ); } + async fn make_parent_waiting_on_blocked_subagent() -> ( + Session, + Arc>, + String, + CancellationToken, + ) { + let block_until_cancelled = RegisteredTool { + definition: ToolDefinition { + name: "block_until_cancelled".into(), + description: "Waits until cancelled".into(), + parameters: serde_json::json!({"type": "object"}), + }, + executor: Arc::new(|_args, ctx| { + Box::pin(async move { + ctx.cancel.cancelled().await; + Ok("cancelled".to_string()) + }) + }), + source: ToolSource::Native, + }; + let mut child_registry = ToolRegistry::new(); + child_registry.register(block_until_cancelled); + let child = make_session_with_tools( + vec![tool_call_response( + "block_until_cancelled", + "child_call", + serde_json::json!({}), + )], + child_registry, + ) + .await; + let child_cancel = child.cancel_token(); + + let manager = Arc::new(AsyncMutex::new(SubAgentManager::new(3))); + let agent_id = manager + .lock() + .await + .spawn(child, "block until cancelled".into(), 0) + .unwrap(); + + let mut parent_registry = ToolRegistry::new(); + parent_registry.register(make_wait_tool(manager.clone())); + let parent_provider = Arc::new(ScriptedStreamProvider::new(vec![ + ScriptedStreamCall::Response(Box::new(tool_call_response( + "wait", + "parent_wait_call", + serde_json::json!({ "agent_id": agent_id }), + ))), + ScriptedStreamCall::Response(Box::new(text_response("resumed"))), + ])); + let client = make_client(parent_provider).await; + let profile = Arc::new(TestProfile::with_tools(parent_registry)); + let env = Arc::new(MockSandbox::default()); + let session = Session::new( + client, + profile, + env, + SessionOptions::default(), + Some(manager.clone()), + ); + manager + .lock() + .await + .set_event_callback(session.sub_agent_event_callback()); + + (session, manager, agent_id, child_cancel) + } + + async fn wait_for_agent_event( + rx: &mut broadcast::Receiver, + predicate: impl Fn(&AgentEvent) -> bool, + ) { + loop { + let event = rx + .recv() + .await + .expect("session event stream should remain open"); + if predicate(&event.event) { + return; + } + } + } + + #[tokio::test] + async fn control_interrupt_during_subagent_wait_closes_child_and_resumes_after_steer() { + let (mut session, manager, agent_id, child_cancel) = + make_parent_waiting_on_blocked_subagent().await; + let control = session.control_handle(); + let mut events = session.subscribe(); + let controller = tokio::spawn(async move { + wait_for_agent_event(&mut events, |event| { + matches!( + event, + AgentEvent::ToolCallStarted { tool_name, .. } if tool_name == "wait" + ) + }) + .await; + control.interrupt(None); + wait_for_agent_event(&mut events, |event| { + matches!(event, AgentEvent::SubAgentClosed { .. }) + }) + .await; + control.steer("resume after interrupt".into(), None); + }); + + timeout( + Duration::from_secs(1), + session.process_input("wait for the child"), + ) + .await + .expect("interrupt should unblock the subagent wait") + .unwrap(); + controller.await.unwrap(); + + assert_eq!(session.state(), SessionState::Idle); + assert!(child_cancel.is_cancelled()); + assert!(matches!( + manager.lock().await.status(&agent_id), + Some(SubAgentStatus::Closed) + )); + } + + #[tokio::test] + async fn terminal_cancel_during_subagent_wait_closes_child_and_session() { + let (mut session, manager, agent_id, child_cancel) = + make_parent_waiting_on_blocked_subagent().await; + let cancel = session.cancel_token(); + let mut events = session.subscribe(); + let controller = tokio::spawn(async move { + wait_for_agent_event(&mut events, |event| { + matches!( + event, + AgentEvent::ToolCallStarted { tool_name, .. } if tool_name == "wait" + ) + }) + .await; + cancel.cancel(); + }); + + let result = timeout( + Duration::from_secs(1), + session.process_input("wait for the child"), + ) + .await + .expect("terminal cancellation should unblock the subagent wait"); + controller.await.unwrap(); + + assert!(matches!( + result, + Err(Error::Interrupted(InterruptReason::Cancelled)) + )); + assert_eq!(session.state(), SessionState::Closed); + assert!(child_cancel.is_cancelled()); + assert!(matches!( + manager.lock().await.status(&agent_id), + Some(SubAgentStatus::Closed) + )); + } + #[tokio::test] async fn close_cleans_up_subagents_before_emitting_session_ended() { use crate::subagent::SubAgentManager; diff --git a/lib/crates/fabro-agent/src/subagent.rs b/lib/crates/fabro-agent/src/subagent.rs index 52cec77c6..ed156dfed 100644 --- a/lib/crates/fabro-agent/src/subagent.rs +++ b/lib/crates/fabro-agent/src/subagent.rs @@ -2,8 +2,9 @@ use std::collections::{HashMap, VecDeque}; use std::sync::{Arc, Mutex}; use fabro_llm::types::ToolDefinition; +use futures::future::{BoxFuture, FutureExt, Shared}; use tokio::sync::Mutex as AsyncMutex; -use tokio::task::JoinHandle; +use tokio::task::{AbortHandle, JoinHandle}; use tokio_util::sync::CancellationToken; use crate::error::Error; @@ -29,6 +30,23 @@ pub struct SubAgentResult { pub turns_used: usize, } +type SubAgentTask = Shared>>; + +fn shared_subagent_task( + task: JoinHandle>, +) -> (SubAgentTask, AbortHandle) { + let abort_handle = task.abort_handle(); + let task = async move { + match task.await { + Ok(result) => result, + Err(e) => Err(Error::InvalidState(format!("Agent task panicked: {e}"))), + } + } + .boxed() + .shared(); + (task, abort_handle) +} + #[derive(Debug, Clone)] pub enum SubAgentStatus { Running, @@ -37,7 +55,8 @@ pub enum SubAgentStatus { } pub struct SubAgent { - task: Option>>, + task: SubAgentTask, + abort_handle: AbortHandle, followup_queue: Arc>>, cancel_token: CancellationToken, depth: usize, @@ -124,9 +143,11 @@ impl SubAgentManager { turns_used: turns.len(), }) }); + let (task, abort_handle) = shared_subagent_task(task); self.agents.insert(agent_id.clone(), SubAgent { - task: Some(task), + task, + abort_handle, followup_queue, cancel_token, depth: depth + 1, @@ -168,45 +189,52 @@ impl SubAgentManager { } pub async fn wait(&mut self, agent_id: &str) -> Result { - // Phase 1: Check existence and current status - let agent = self.agents.get(agent_id); - let depth = match agent { - None => { - return Err(Error::InvalidState(format!( - "No agent found with id: {agent_id} (it was never spawned)" - ))); - } - Some(a) => a.depth, - }; + let task = self.result_future(agent_id)?; + let task_result = task.await; + self.record_result(agent_id, task_result) + } - match &self.agents[agent_id].status { - SubAgentStatus::Closed => { - return Err(Error::InvalidState(format!( - "Agent {agent_id} has been closed" - ))); - } - SubAgentStatus::Finished(result) => { - return result.clone(); - } - SubAgentStatus::Running => {} + fn result_future(&self, agent_id: &str) -> Result { + let agent = self.agents.get(agent_id).ok_or_else(|| { + Error::InvalidState(format!( + "No agent found with id: {agent_id} (it was never spawned)" + )) + })?; + + match &agent.status { + SubAgentStatus::Closed => Err(Error::InvalidState(format!( + "Agent {agent_id} has been closed" + ))), + SubAgentStatus::Running | SubAgentStatus::Finished(_) => Ok(agent.task.clone()), } + } - // Phase 2: Take the JoinHandle (brief mutable borrow, no await) - let join_handle = self - .agents - .get_mut(agent_id) - .expect("agent should still exist after status check") - .task - .take() - .ok_or_else(|| Error::InvalidState(format!("Agent {agent_id} has no running task")))?; + fn record_result( + &mut self, + agent_id: &str, + task_result: Result, + ) -> Result { + let depth = { + let agent = self.agents.get_mut(agent_id).ok_or_else(|| { + Error::InvalidState(format!( + "No agent found with id: {agent_id} (it was never spawned)" + )) + })?; - // Phase 3: Await the task (no borrow held) - let task_result = match join_handle.await { - Ok(result) => result, - Err(e) => Err(Error::InvalidState(format!("Agent task panicked: {e}"))), + match &agent.status { + SubAgentStatus::Closed => { + return Err(Error::InvalidState(format!( + "Agent {agent_id} has been closed" + ))); + } + SubAgentStatus::Finished(result) => return result.clone(), + SubAgentStatus::Running => {} + } + + agent.status = SubAgentStatus::Finished(task_result.clone()); + agent.depth }; - // Phase 4: Emit event match &task_result { Ok(result) => { self.emit_event(AgentEvent::SubAgentCompleted { @@ -225,17 +253,7 @@ impl SubAgentManager { } } - // Phase 5: Store result in status and return clone - let agent = self - .agents - .get_mut(agent_id) - .expect("agent should still exist when storing task result"); - agent.status = SubAgentStatus::Finished(task_result); - - match &agent.status { - SubAgentStatus::Finished(result) => result.clone(), - _ => unreachable!("agent status was just assigned to Finished on the line above"), - } + task_result } pub fn close(&mut self, agent_id: &str) -> Result<(), Error> { @@ -253,9 +271,7 @@ impl SubAgentManager { } SubAgentStatus::Running => { agent.cancel_token.cancel(); - if let Some(join_handle) = agent.task.take() { - join_handle.abort(); - } + agent.abort_handle.abort(); } SubAgentStatus::Finished(_) => { // No task to cancel, just transition status @@ -415,13 +431,28 @@ pub fn make_wait_tool(manager: Arc>) -> RegisteredTo "required": ["agent_id"] }), }, - executor: Arc::new(move |args, _ctx| { + executor: Arc::new(move |args, ctx| { let manager = manager.clone(); Box::pin(async move { let agent_id = required_str(&args, "agent_id")?; - let mut mgr = manager.lock().await; - let result = mgr.wait(agent_id).await.map_err(|e| e.to_string())?; + let task = { + let mgr = manager.lock().await; + mgr.result_future(agent_id).map_err(|e| e.to_string())? + }; + let task_result = tokio::select! { + biased; + () = ctx.cancel.cancelled() => { + let _ = manager.lock().await.close(agent_id); + return Err("Cancelled".to_string()); + } + result = task => result, + }; + let result = manager + .lock() + .await + .record_result(agent_id, task_result) + .map_err(|e| e.to_string())?; Ok(format!( "Agent completed (success: {}, turns: {})\n\n{}", result.success, result.turns_used, result.output @@ -471,6 +502,7 @@ mod tests { use super::*; use crate::config::SessionOptions; use crate::test_support::*; + use crate::tool_registry::ToolContext; // --- Tests --- @@ -610,6 +642,69 @@ mod tests { )); } + #[tokio::test] + async fn wait_tool_returns_when_context_is_cancelled() { + let child_cancel = CancellationToken::new(); + let child_cancel_probe = child_cancel.clone(); + let task_cancel = child_cancel.clone(); + let task = tokio::spawn(async move { + task_cancel.cancelled().await; + Ok(SubAgentResult { + output: String::new(), + success: false, + turns_used: 0, + }) + }); + let (task, abort_handle) = shared_subagent_task(task); + + let agent_id = "blocked-agent".to_string(); + let manager = Arc::new(AsyncMutex::new(SubAgentManager::new(3))); + manager + .lock() + .await + .agents + .insert(agent_id.clone(), SubAgent { + task, + abort_handle, + followup_queue: Arc::new(Mutex::new(VecDeque::new())), + cancel_token: child_cancel, + depth: 1, + status: SubAgentStatus::Running, + }); + + let tool = make_wait_tool(manager.clone()); + let tool_cancel = CancellationToken::new(); + let ctx = ToolContext { + env: Arc::new(MockSandbox::default()), + cancel: tool_cancel.clone(), + tool_env_provider: None, + session_id: None, + root_session_id: None, + tool_call_id: None, + agent_event_emitter: None, + }; + let mut wait = (tool.executor)(serde_json::json!({ "agent_id": agent_id }), ctx); + + assert!( + futures::poll!(wait.as_mut()).is_pending(), + "blocked subagent should leave the wait tool pending" + ); + tool_cancel.cancel(); + + let result = time::timeout(std::time::Duration::from_millis(100), wait.as_mut()).await; + drop(wait); + manager.lock().await.close_all(); + + let result = + result.expect("wait tool should return promptly when its context is cancelled"); + assert_eq!(result, Err("Cancelled".to_string())); + assert!(child_cancel_probe.is_cancelled()); + assert!(matches!( + manager.lock().await.status(&agent_id), + Some(SubAgentStatus::Closed) + )); + } + #[test] fn tool_definitions_correct() { let manager = Arc::new(AsyncMutex::new(SubAgentManager::new(3))); From 5cf1c7d1837a9c56d54034d4542618ccf12ad61d Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 23 Jul 2026 20:40:22 -0400 Subject: [PATCH 2/4] Harden cancellation and interrupt lifecycles --- .../components/runs-list/row-actions-menu.tsx | 39 +- .../app/components/steer-bar.test.tsx | 22 + apps/fabro-web/app/components/steer-bar.tsx | 33 +- apps/fabro-web/app/data/runs.test.ts | 2 + apps/fabro-web/app/data/runs.ts | 3 + apps/fabro-web/app/lib/board-events.test.tsx | 1 + apps/fabro-web/app/lib/board-events.ts | 3 + apps/fabro-web/app/lib/mutations.ts | 5 +- apps/fabro-web/app/lib/query-keys.test.ts | 1 + apps/fabro-web/app/lib/run-actions.test.ts | 36 + apps/fabro-web/app/lib/run-actions.ts | 32 +- apps/fabro-web/app/lib/run-events.test.tsx | 17 + apps/fabro-web/app/lib/run-events.ts | 16 + apps/fabro-web/app/routes/run-detail.test.ts | 49 +- apps/fabro-web/app/routes/run-detail.tsx | 13 +- .../app/routes/run-detail/docked-controls.tsx | 8 +- .../app/routes/run-detail/lifecycle-toasts.ts | 4 +- apps/fabro-web/app/routes/run-stages.test.ts | 25 + apps/fabro-web/app/routes/run-stages.tsx | 7 + apps/fabro-web/app/routes/runs.test.tsx | 1 + docs/public/api-reference/fabro-api.yaml | 26 +- lib/crates/fabro-agent/src/agent_profile.rs | 15 +- lib/crates/fabro-agent/src/cli.rs | 44 +- lib/crates/fabro-agent/src/lib.rs | 8 +- .../fabro-agent/src/profiles/anthropic.rs | 12 +- lib/crates/fabro-agent/src/profiles/gemini.rs | 8 +- lib/crates/fabro-agent/src/profiles/openai.rs | 8 +- lib/crates/fabro-agent/src/session.rs | 445 ++++++-- lib/crates/fabro-agent/src/subagent.rs | 1013 ++++++++++++----- lib/crates/fabro-agent/src/types.rs | 8 + .../fabro-agent/tests/it/parity_matrix.rs | 9 +- lib/crates/fabro-api/build.rs | 1 + lib/crates/fabro-api/src/lib.rs | 15 +- .../tests/run_projection_round_trip.rs | 1 + .../tests/stage_projection_round_trip.rs | 19 +- lib/crates/fabro-cli/tests/it/cmd/runner.rs | 2 +- lib/crates/fabro-server/src/server.rs | 15 +- .../src/server/handler/lifecycle.rs | 134 ++- lib/crates/fabro-server/src/server/tests.rs | 247 +++- lib/crates/fabro-server/src/worker_runtime.rs | 8 + .../tests/it/scenario/lifecycle.rs | 4 +- lib/crates/fabro-store/src/run_state.rs | 139 ++- lib/crates/fabro-types/src/lib.rs | 11 +- lib/crates/fabro-types/src/run_event/agent.rs | 6 + lib/crates/fabro-types/src/run_event/mod.rs | 27 + lib/crates/fabro-types/src/run_projection.rs | 31 +- .../fabro-workflow/src/event/convert.rs | 30 + lib/crates/fabro-workflow/src/event/names.rs | 12 + .../fabro-workflow/src/handler/llm/api.rs | 260 ++++- .../src/.openapi-generator/FILES | 1 + .../fabro-api-client/src/api/runs-api.ts | 8 +- .../src/models/agent-control-state.ts | 26 + .../fabro-api-client/src/models/index.ts | 1 + .../src/models/stage-projection.ts | 7 + 54 files changed, 2346 insertions(+), 572 deletions(-) create mode 100644 apps/fabro-web/app/components/steer-bar.test.tsx create mode 100644 lib/packages/fabro-api-client/src/models/agent-control-state.ts diff --git a/apps/fabro-web/app/components/runs-list/row-actions-menu.tsx b/apps/fabro-web/app/components/runs-list/row-actions-menu.tsx index b1582dbef..49bd39ace 100644 --- a/apps/fabro-web/app/components/runs-list/row-actions-menu.tsx +++ b/apps/fabro-web/app/components/runs-list/row-actions-menu.tsx @@ -12,9 +12,12 @@ import { canCancel, canDelete, canUnarchive, + cancellationActionLabel, + cancellationSuccessMessage, cancelRun, deleteRun, denyRun, + isCancellationPendingState, mapError, retryRun, unarchiveRun, @@ -32,7 +35,9 @@ const MENU_ITEM_DANGER_CLASS = export function RowActionsMenu({ run }: { run: RunWithStatus }) { const { mutate } = useSWRConfig(); const { push } = useToast(); - const [pending, setPending] = useState(false); + const [pendingAction, setPendingAction] = useState(null); + const [optimisticCancellationRunId, setOptimisticCancellationRunId] = + useState(null); const [deleteDialogOpen, setDeleteDialogOpen] = useState(false); const [idCopied, setIdCopied] = useState(false); @@ -44,6 +49,12 @@ export function RowActionsMenu({ run }: { run: RunWithStatus }) { const showUnarchive = canUnarchive(status); const showCancel = canCancel(status); const showDelete = canDelete(status); + const cancellationPending = isCancellationPendingState( + status, + run.pendingControl, + pendingAction === "cancel" || optimisticCancellationRunId === run.id, + ); + const pending = pendingAction !== null || cancellationPending; const hasLifecycle = showRetry || showArchive || showUnarchive; const hasDestructive = showDeny || showCancel || showDelete; @@ -51,17 +62,25 @@ export function RowActionsMenu({ run }: { run: RunWithStatus }) { async function runAction( label: LifecycleAction, action: () => Promise, - successMessage: string, + successMessage: string | ((result: T) => string), ) { if (pending) return; - setPending(true); + setPendingAction(label); try { - await action(); - push({ message: successMessage }); + const result = await action(); + if (label === "cancel") { + setOptimisticCancellationRunId(run.id); + } + push({ + message: + typeof successMessage === "function" + ? successMessage(result) + : successMessage, + }); } catch (error) { push({ message: mapError(error, label), tone: "error" }); } finally { - setPending(false); + setPendingAction(null); mutateRunListCaches(mutate); } } @@ -81,7 +100,7 @@ export function RowActionsMenu({ run }: { run: RunWithStatus }) { async function handleDeleteConfirm() { if (pending) return; - setPending(true); + setPendingAction("delete"); try { await deleteRun(run.id); push({ message: "Deleted run." }); @@ -89,7 +108,7 @@ export function RowActionsMenu({ run }: { run: RunWithStatus }) { // deleteRun throws LifecycleActionError shapes via lifecycleActionErrorFromError push({ message: mapError(error, "archive"), tone: "error" }); } finally { - setPending(false); + setPendingAction(null); setDeleteDialogOpen(false); mutateRunListCaches(mutate); } @@ -208,12 +227,12 @@ export function RowActionsMenu({ run }: { run: RunWithStatus }) { )} diff --git a/apps/fabro-web/app/components/steer-bar.test.tsx b/apps/fabro-web/app/components/steer-bar.test.tsx new file mode 100644 index 000000000..33fefa556 --- /dev/null +++ b/apps/fabro-web/app/components/steer-bar.test.tsx @@ -0,0 +1,22 @@ +import { describe, expect, test } from "bun:test"; +import { createElement } from "react"; +import { renderToStaticMarkup } from "react-dom/server"; + +import { + isInterruptDisabled, + SteerWaitingStatus, +} from "./steer-bar"; + +describe("SteerBar", () => { + test("shows durable waiting state and prevents a second interrupt", () => { + expect(isInterruptDisabled(true, false)).toBe(true); + expect(isInterruptDisabled(false, true)).toBe(true); + expect(isInterruptDisabled(false, false)).toBe(false); + + const html = renderToStaticMarkup( + createElement(SteerWaitingStatus, { waitingForSteer: true }), + ); + expect(html).toContain('role="status"'); + expect(html).toContain("Interrupted — waiting for steering"); + }); +}); diff --git a/apps/fabro-web/app/components/steer-bar.tsx b/apps/fabro-web/app/components/steer-bar.tsx index 6c169dc9a..b78a62f92 100644 --- a/apps/fabro-web/app/components/steer-bar.tsx +++ b/apps/fabro-web/app/components/steer-bar.tsx @@ -13,6 +13,7 @@ import { ErrorMessage } from "./ui"; export interface SteerBarProps { runId: string; + waitingForSteer?: boolean; ref?: Ref; } @@ -20,13 +21,38 @@ export interface SteerBarHandle { focus(): void; } -export function SteerBar({ runId, ref }: SteerBarProps) { +export function isInterruptDisabled( + waitingForSteer: boolean, + mutationPending: boolean, +): boolean { + return waitingForSteer || mutationPending; +} + +export function SteerWaitingStatus({ + waitingForSteer, +}: { + waitingForSteer: boolean; +}) { + if (!waitingForSteer) return null; + return ( +

+ Interrupted — waiting for steering +

+ ); +} + +export function SteerBar({ + runId, + waitingForSteer = false, + ref, +}: SteerBarProps) { const [text, setText] = useState(""); const [errorMessage, setErrorMessage] = useState(null); const textareaRef = useRef(null); const steer = useSteerRun(runId); const interrupt = useInterruptRun(runId); const pending = steer.isMutating || interrupt.isMutating; + const interruptDisabled = isInterruptDisabled(waitingForSteer, pending); useImperativeHandle(ref, () => ({ focus() { @@ -49,7 +75,7 @@ export function SteerBar({ runId, ref }: SteerBarProps) { } async function fireInterrupt() { - if (pending) return; + if (interruptDisabled) return; setErrorMessage(null); try { await interrupt.trigger(); @@ -91,7 +117,7 @@ export function SteerBar({ runId, ref }: SteerBarProps) {