From 4c4ced14e1f5ab8ea6f58e14ab65b83dcb6218b9 Mon Sep 17 00:00:00 2001 From: Tryanks Date: Mon, 21 Sep 2026 16:29:45 +0800 Subject: [PATCH] Close nested Claude subagent mirrors instead of leaving nameless ones running MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A Claude subagent that calls the Agent tool produces a grandchild whose stdout lines carry the grandchild's own tool_use id as parent_tool_use_id, while the spawning tool_use only appears inside the child's stream. The child transcript mapper turned every tool_use into a plain ToolCall, so no Subagent item ever existed for the grandchild: the runtime created a placeholder mirror titled "subagent", opened its turn, and nothing closed it — task events found no tool item, and after a host restart the open turn sat in the persisted log for good. The transcript mapper now emits Agent/Task calls as Subagent items keyed by their bare tool_use id and parented to the child. The Claude mapper registers them in tool_items, so the grandchild's task_started tail, task_notification, model observation and terminal gating work as for a top-level spawn; a foreground nested spawn settles with the child's tool_result, a background one (run_in_background, async_launched, or the launch acknowledgement text that subagent transcripts carry without a toolUseResult) with its notification. Child items from the stdout feed and the transcript tail are deduplicated per lifecycle stage rather than only after completion. The runtime keeps a Subagent item recorded inside a mirror as that mirror's content and, once a grandchild item names it as parent, opens a titled mirror nested under the child mirror; its snapshots then update and close that mirror. Mirror lookup, SessionClosed interruption and orchestrator child teardown follow the descendant chain. Loading a mirror whose last persisted turn is still open while no live parent tracks it synthesizes an interrupted TurnCompleted, so existing zombie mirrors settle on upgrade. is_agent_tool now matches only the spawning tools (agent, task), so ListAgents no longer produces a "subagent: " mirror. --- crates/agent/src/claude.rs | 383 +++++++++++++++--- crates/agent/src/subagent_tail.rs | 126 +++++- .../claude/subagent_nested_trace.jsonl | 15 + crates/runtime/src/app/lifecycle.rs | 9 +- crates/runtime/src/app/mod.rs | 4 + crates/runtime/src/app/sessions.rs | 3 +- crates/runtime/src/app/subagents.rs | 265 ++++++++---- crates/runtime/src/app/tests.rs | 353 ++++++++++++++++ 8 files changed, 1010 insertions(+), 148 deletions(-) create mode 100644 crates/agent/tests/fixtures/claude/subagent_nested_trace.jsonl diff --git a/crates/agent/src/claude.rs b/crates/agent/src/claude.rs index 160d4b62..a2043c06 100644 --- a/crates/agent/src/claude.rs +++ b/crates/agent/src/claude.rs @@ -1124,9 +1124,36 @@ enum ToolItem { summary: Option, model: Option, effort: Option, + /// The subagent this one was spawned from, for nested spawns. + parent_item_id: Option, }, } +impl ToolItem { + fn subagent(content: &ItemContent, parent_item_id: Option) -> Option { + let ItemContent::Subagent { + agent_type, + description, + status, + summary, + model, + effort, + } = content + else { + return None; + }; + Some(ToolItem::Subagent { + agent_type: agent_type.clone(), + description: description.clone(), + status: *status, + summary: summary.clone(), + model: model.clone(), + effort: effort.clone(), + parent_item_id, + }) + } +} + enum TailRequest { Start { parent_id: String, @@ -1153,6 +1180,14 @@ struct PendingApproval { suggestions: Option, } +/// Where a child item's lifecycle stands, as far as the canonical stream has +/// been told. An update keeps its snapshot so only a changed one passes. +enum ChildItemStage { + Started, + Updated(ItemContent), + Completed, +} + struct PendingRewind { checkpoint_id: String, mode: RewindMode, @@ -1193,9 +1228,11 @@ pub(crate) struct Mapper { /// then, so the mirror's turn closes after the tail's final child items. tailed: HashSet, pending_subagent_terminals: HashMap, - /// Child items already completed by either feed (stdout `parent_tool_use_id` - /// lines or the transcript tail); a stale start must not reopen them. - finished_child_items: HashSet, + /// The last lifecycle stage each child item reached through either feed + /// (stdout `parent_tool_use_id` lines or the transcript tail); the other + /// feed's copy of the same stage is dropped and a stale start cannot + /// reopen a completed item. + child_items: HashMap, pending_approvals: HashMap, /// Pending `AskUserQuestion` prompts: control request_id → the original /// `questions` array, echoed back verbatim in the allow response. @@ -1295,7 +1332,7 @@ impl Mapper { tail_requests: Vec::new(), tailed: HashSet::new(), pending_subagent_terminals: HashMap::new(), - finished_child_items: HashSet::new(), + child_items: HashMap::new(), pending_approvals: HashMap::new(), pending_user_input: HashMap::new(), approval_mode, @@ -1552,7 +1589,7 @@ impl Mapper { .entry(parent_id.to_owned()) .or_insert_with(|| crate::subagent_tail::TranscriptMapper::new(parent_id)) .map_value(&msg); - events.extend(self.dedupe_child_events(child)); + events.extend(self.fold_child_events(child)); return events; } match msg.get("type").and_then(Value::as_str) { @@ -1882,6 +1919,7 @@ impl Mapper { summary: saved_summary, model, effort, + parent_item_id, }) = self.tool_items.get_mut(tool_use_id) else { return Vec::new(); @@ -1892,7 +1930,7 @@ impl Mapper { *saved_status = status; let event = AgentEvent::ItemUpdated(ThreadItem { id: tool_use_id.to_owned(), - parent_item_id: None, + parent_item_id: parent_item_id.clone(), content: ItemContent::Subagent { agent_type: agent_type.clone(), description: description.clone(), @@ -1921,17 +1959,64 @@ impl Mapper { } } - fn dedupe_child_events(&mut self, mut events: Vec) -> Vec { - let finished = &mut self.finished_child_items; - events.retain(|event| match event { - AgentEvent::ItemStarted(item) => !finished.contains(&item.id), - AgentEvent::ItemCompleted(item) => { - finished.insert(item.id.clone()); + /// Fold one feed's child items into the canonical stream: drop the stage + /// the other feed already delivered, and give a nested spawn the same + /// `tool_items` lifecycle as a top-level one so its task events and tail + /// settle it. + fn fold_child_events(&mut self, events: Vec) -> Vec { + let mut out = Vec::new(); + for event in events { + if !self.advance_child_item(&event) { + continue; + } + match &event { + AgentEvent::ItemStarted(item) + if matches!(item.content, ItemContent::Subagent { .. }) => + { + if let Some(tool) = + ToolItem::subagent(&item.content, item.parent_item_id.clone()) + { + self.tool_items.entry(item.id.clone()).or_insert(tool); + } + out.push(event); + } + AgentEvent::ItemCompleted(item) + if matches!(item.content, ItemContent::Subagent { .. }) + && matches!( + self.tool_items.get(&item.id), + Some(ToolItem::Subagent { .. }) + ) => + { + let id = item.id.clone(); + self.tool_items.remove(&id); + out.extend(self.gate_subagent_terminal(&id, event)); + } + _ => out.push(event), + } + } + out + } + + /// Record a child item's stage; false when this stage was already delivered. + fn advance_child_item(&mut self, event: &AgentEvent) -> bool { + let (item, stage) = match event { + AgentEvent::ItemStarted(item) => (item, ChildItemStage::Started), + AgentEvent::ItemUpdated(item) => (item, ChildItemStage::Updated(item.content.clone())), + AgentEvent::ItemCompleted(item) => (item, ChildItemStage::Completed), + _ => return true, + }; + match (self.child_items.get(&item.id), &stage) { + (Some(ChildItemStage::Completed), _) | (Some(_), ChildItemStage::Started) => false, + (Some(ChildItemStage::Updated(previous)), ChildItemStage::Updated(content)) + if previous == content => + { + false + } + _ => { + self.child_items.insert(item.id.clone(), stage); true } - _ => true, - }); - events + } } /// Fold one batch from a subagent transcript tail: its child items, the @@ -1942,7 +2027,7 @@ impl Mapper { Some((model, effort)) => self.note_subagent_model(¬ice.parent_id, model, effort), None => Vec::new(), }; - events.extend(self.dedupe_child_events(notice.events)); + events.extend(self.fold_child_events(notice.events)); if notice.stopped { self.tailed.remove(¬ice.parent_id); events.extend(self.pending_subagent_terminals.remove(¬ice.parent_id)); @@ -2293,38 +2378,9 @@ impl Mapper { } let (item, content) = if is_agent_tool(&name.to_lowercase()) { - let agent_type = input - .get("subagent_type") - .and_then(Value::as_str) - .unwrap_or("subagent") - .to_owned(); - let description = subagent_description(&input); - // The Agent tool only carries a model alias when the caller picked - // one; the resolved model and effort arrive with the child's first - // assistant message. - let model = input - .get("model") - .and_then(Value::as_str) - .filter(|model| !model.is_empty()) - .map(str::to_owned); - ( - ToolItem::Subagent { - agent_type: agent_type.clone(), - description: description.clone(), - status: ItemStatus::InProgress, - summary: None, - model: model.clone(), - effort: None, - }, - ItemContent::Subagent { - agent_type, - description, - status: ItemStatus::InProgress, - summary: None, - model, - effort: None, - }, - ) + let content = spawned_subagent(&input); + let item = ToolItem::subagent(&content, None).expect("a spawn snapshot is a Subagent"); + (item, content) } else if name == "Bash" { let command = input .get("command") @@ -2439,6 +2495,7 @@ impl Mapper { } else { ItemStatus::Completed }; + let mut parent_item_id = None; let content = match item { ToolItem::Command { command, .. } if background_task_id.is_some() => { self.tool_items.insert( @@ -2496,16 +2553,21 @@ impl Mapper { summary, model, effort, + parent_item_id: parent, .. - } => ItemContent::Subagent { - agent_type, - description, - status, - summary: summary - .or_else(|| (!output.trim().is_empty()).then(|| one_line_summary(&output))), - model, - effort, - }, + } => { + parent_item_id = parent; + ItemContent::Subagent { + agent_type, + description, + status, + summary: summary.or_else(|| { + (!output.trim().is_empty()).then(|| one_line_summary(&output)) + }), + model, + effort, + } + } }; let event = if matches!( &content, @@ -2521,7 +2583,7 @@ impl Mapper { let is_subagent = matches!(content, ItemContent::Subagent { .. }); let event = event(ThreadItem { id: tool_use_id.clone(), - parent_item_id: None, + parent_item_id, content, }); if is_subagent { @@ -2865,7 +2927,7 @@ fn subagent_status(status: &str) -> ItemStatus { } } -fn one_line_summary(text: &str) -> String { +pub(crate) fn one_line_summary(text: &str) -> String { text.split_whitespace().collect::>().join(" ") } @@ -2878,8 +2940,32 @@ enum ClaudeRequestType { ToolUse, } -fn is_agent_tool(normalized: &str) -> bool { - normalized.contains("agent") || normalized == "task" +/// The tools that spawn a subagent: `Agent` and its former name `Task`. +/// Coordination tools such as `ListAgents` and `SendMessage` do not. +pub(crate) fn is_agent_tool(normalized: &str) -> bool { + matches!(normalized, "agent" | "task") +} + +/// The in-progress Subagent snapshot for an Agent tool call. The call only +/// carries a model alias when the caller picked one; the resolved model and +/// effort arrive with the child's first assistant message. +pub(crate) fn spawned_subagent(input: &Value) -> ItemContent { + ItemContent::Subagent { + agent_type: input + .get("subagent_type") + .and_then(Value::as_str) + .unwrap_or("subagent") + .to_owned(), + description: subagent_description(input), + status: ItemStatus::InProgress, + summary: None, + model: input + .get("model") + .and_then(Value::as_str) + .filter(|model| !model.is_empty()) + .map(str::to_owned), + effort: None, + } } fn subagent_description(input: &Value) -> String { @@ -5417,6 +5503,181 @@ mod tests { assert!(mapper.take_pending_subagent_terminals().is_empty()); } + /// A background subagent that itself spawns a background Agent: the + /// nested spawn is a Subagent item parented to the child, the + /// grandchild's transcript is parented to the nested spawn, and the + /// grandchild's task_notification settles it after its tail's final + /// flush. Both feeds carry every child record; each stage is emitted once. + #[test] + fn nested_background_subagent_settles_through_its_own_task_notification() { + let trace = include_str!("../tests/fixtures/claude/subagent_nested_trace.jsonl"); + let mut mapper = Mapper::new(); + let mut events = Vec::new(); + for line in trace.lines() { + events.extend(feed(&mut mapper, line)); + } + let subagent_snapshots = |events: &[AgentEvent], id: &str| -> Vec { + events + .iter() + .filter_map(|event| match event { + AgentEvent::ItemStarted(item) + | AgentEvent::ItemUpdated(item) + | AgentEvent::ItemCompleted(item) + if item.id == id + && matches!(item.content, ItemContent::Subagent { .. }) => + { + Some(item.clone()) + } + _ => None, + }) + .collect() + }; + + let grandchild = subagent_snapshots(&events, "toolu_grandchild"); + assert!(matches!( + &grandchild[0], + ThreadItem { parent_item_id: Some(parent), content: ItemContent::Subagent { agent_type, description, status: ItemStatus::InProgress, .. }, .. } + if parent == "toolu_child" && agent_type == "general-purpose" && description == "Research Workers limits" + )); + assert_eq!( + grandchild + .iter() + .filter(|item| matches!( + item.content, + ItemContent::Subagent { + status: ItemStatus::InProgress, + .. + } + )) + .count(), + grandchild.len(), + "launch acknowledgement and task_notification wait for the grandchild's tail" + ); + assert!( + grandchild + .iter() + .all(|item| item.parent_item_id.as_deref() == Some("toolu_child")) + ); + assert!(grandchild.iter().any(|item| matches!( + &item.content, + ItemContent::Subagent { model: Some(model), effort: Some(effort), .. } + if model == "claude-sonnet-5" && effort == "high" + ))); + assert!(events.iter().any(|event| matches!( + event, + AgentEvent::ItemCompleted(ThreadItem { id, parent_item_id: Some(parent), content: ItemContent::ToolCall { name, status: ItemStatus::Completed, .. }, .. }) + if id == "toolu_grandchild:toolu_gc_fetch" && parent == "toolu_grandchild" && name == "WebFetch" + ))); + assert_eq!( + events + .iter() + .filter(|event| matches!(event, AgentEvent::ItemStarted(ThreadItem { id, .. }) if id == "toolu_grandchild")) + .count(), + 1 + ); + let requests = mapper.take_tail_requests(); + assert!(requests.iter().any(|request| matches!(request, TailRequest::Start { parent_id, .. } if parent_id == "toolu_grandchild"))); + assert!(requests.iter().any(|request| matches!(request, TailRequest::Stop { parent_id } if parent_id == "toolu_grandchild"))); + + // The child's tail replays the nested spawn and the grandchild's tail + // replays its transcript; neither reopens or repeats a delivered stage. + let mut child_tail = crate::subagent_tail::TranscriptMapper::new("toolu_child"); + let mut grandchild_tail = crate::subagent_tail::TranscriptMapper::new("toolu_grandchild"); + let mut child_events = Vec::new(); + let mut grandchild_events = Vec::new(); + for line in trace.lines() { + let value: Value = serde_json::from_str(line).unwrap(); + match value.get("parent_tool_use_id").and_then(Value::as_str) { + Some("toolu_child") => child_events.extend(child_tail.map_value(&value)), + Some("toolu_grandchild") => { + grandchild_events.extend(grandchild_tail.map_value(&value)) + } + _ => {} + } + } + let flushed = mapper.on_tail_notice(SubagentTailNotice { + parent_id: "toolu_grandchild".into(), + events: grandchild_events, + model: None, + stopped: true, + }); + assert!( + matches!( + flushed.as_slice(), + [AgentEvent::ItemUpdated(ThreadItem { id, parent_item_id: Some(parent), content: ItemContent::Subagent { status: ItemStatus::Completed, summary: Some(summary), model: Some(model), effort: Some(effort), .. }, .. })] + if id == "toolu_grandchild" && parent == "toolu_child" && summary == "Agent \"Research Workers limits\" finished" && model == "claude-sonnet-5" && effort == "high" + ), + "{flushed:?}" + ); + let flushed = mapper.on_tail_notice(SubagentTailNotice { + parent_id: "toolu_child".into(), + events: child_events, + model: None, + stopped: true, + }); + assert!( + matches!( + flushed.as_slice(), + [AgentEvent::ItemUpdated(ThreadItem { id, parent_item_id: None, content: ItemContent::Subagent { status: ItemStatus::Completed, summary: Some(summary), .. }, .. })] + if id == "toolu_child" && summary == "Hosting researched" + ), + "{flushed:?}" + ); + assert!(mapper.take_pending_subagent_terminals().is_empty()); + } + + /// A nested Agent that runs in the foreground settles with the child's + /// tool_result; a coordination tool is not a spawn. + #[test] + fn nested_foreground_subagent_completes_on_the_childs_tool_result() { + let mut mapper = Mapper::new(); + let mut events = feed( + &mut mapper, + r#"{"type":"assistant","message":{"id":"msg-spawn","content":[{"type":"tool_use","id":"toolu_child","name":"Agent","input":{"description":"Audit","prompt":"Audit routing.","subagent_type":"Explore"}}]}}"#, + ); + events.extend(feed( + &mut mapper, + r#"{"type":"assistant","parent_tool_use_id":"toolu_child","message":{"id":"msg-child-1","content":[{"type":"tool_use","id":"toolu_list","name":"ListAgents","input":{}},{"type":"tool_use","id":"toolu_nested","name":"Agent","input":{"description":"Check tests","prompt":"Check the tests.","subagent_type":"Explore"}}]}}"#, + )); + events.extend(feed( + &mut mapper, + r#"{"type":"user","parent_tool_use_id":"toolu_child","message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"toolu_list","content":"none"},{"type":"tool_result","tool_use_id":"toolu_nested","content":[{"type":"text","text":"Tests are\ngreen."}]}]}}"#, + )); + assert!(events.iter().any(|event| matches!( + event, + AgentEvent::ItemStarted(ThreadItem { id, parent_item_id: Some(parent), content: ItemContent::ToolCall { name, .. }, .. }) + if id == "toolu_child:toolu_list" && parent == "toolu_child" && name == "ListAgents" + ))); + assert!( + matches!( + events.last(), + Some(AgentEvent::ItemCompleted(ThreadItem { id, parent_item_id: Some(parent), content: ItemContent::Subagent { agent_type, status: ItemStatus::Completed, summary: Some(summary), .. }, .. })) + if id == "toolu_nested" && parent == "toolu_child" && agent_type == "Explore" && summary == "Tests are green." + ), + "{events:?}" + ); + assert!(mapper.take_pending_subagent_terminals().is_empty()); + } + + #[test] + fn only_spawning_tools_are_agent_tools() { + for name in ["agent", "task"] { + assert!(is_agent_tool(name), "{name}"); + } + for name in ["listagents", "sendmessage", "taskcreate", "subagent_run"] { + assert!(!is_agent_tool(name), "{name}"); + } + let mut mapper = Mapper::new(); + let events = feed( + &mut mapper, + r#"{"type":"assistant","message":{"id":"msg-list","content":[{"type":"tool_use","id":"toolu_list","name":"ListAgents","input":{}}]}}"#, + ); + assert!(matches!( + events.as_slice(), + [AgentEvent::ItemStarted(ThreadItem { content: ItemContent::ToolCall { name, .. }, .. })] if name == "ListAgents" + )); + } + /// The tail keeps the transcript discovery found; the notification's /// `output_file` only seeds a reader when discovery found nothing. #[test] diff --git a/crates/agent/src/subagent_tail.rs b/crates/agent/src/subagent_tail.rs index 46365bab..1a1c188a 100644 --- a/crates/agent/src/subagent_tail.rs +++ b/crates/agent/src/subagent_tail.rs @@ -7,13 +7,27 @@ use std::path::{Path, PathBuf}; use serde_json::Value; +use crate::claude::{is_agent_tool, one_line_summary, spawned_subagent}; use crate::{AgentEvent, ItemContent, ItemStatus, ThreadItem}; +enum ChildTool { + Tool { + name: String, + input: Value, + }, + /// An `Agent` call inside the subagent: a nested subagent whose own + /// transcript arrives keyed by this tool_use id. + Subagent { + content: ItemContent, + background: bool, + }, +} + /// Stateful mapper for one subagent transcript. Tool calls are retained until /// their matching result arrives so completion snapshots keep the original input. pub(crate) struct TranscriptMapper { parent_id: String, - tools: HashMap, + tools: HashMap, next_user_id: u64, /// Latest model/effort seen on an assistant record, until taken. model: Option<(Option, Option)>, @@ -90,8 +104,29 @@ impl TranscriptMapper { .unwrap_or("tool") .to_owned(); let input = block.get("input").cloned().unwrap_or(Value::Null); - self.tools - .insert(id.to_owned(), (name.clone(), input.clone())); + if is_agent_tool(&name.to_lowercase()) { + let content = spawned_subagent(&input); + let background = input + .get("run_in_background") + .and_then(Value::as_bool) + .unwrap_or(false); + self.tools.insert( + id.to_owned(), + ChildTool::Subagent { + content: content.clone(), + background, + }, + ); + events.push(AgentEvent::ItemStarted(self.nested_subagent(id, content))); + continue; + } + self.tools.insert( + id.to_owned(), + ChildTool::Tool { + name: name.clone(), + input: input.clone(), + }, + ); events.push(AgentEvent::ItemStarted(self.item( id, ItemContent::ToolCall { @@ -143,26 +178,63 @@ impl TranscriptMapper { let Some(id) = block.get("tool_use_id").and_then(Value::as_str) else { continue; }; - let Some((name, input)) = self.tools.remove(id) else { + let Some(tool) = self.tools.remove(id) else { continue; }; let failed = block .get("is_error") .and_then(Value::as_bool) .unwrap_or(false); - events.push(AgentEvent::ItemCompleted(self.item( - id, - ItemContent::ToolCall { - name, - input, - output: Some(content_text(block.get("content"))), - status: if failed { - ItemStatus::Failed - } else { - ItemStatus::Completed - }, - }, - ))); + let output = content_text(block.get("content")); + let status = if failed { + ItemStatus::Failed + } else { + ItemStatus::Completed + }; + match tool { + ChildTool::Tool { name, input } => { + events.push(AgentEvent::ItemCompleted(self.item( + id, + ItemContent::ToolCall { + name, + input, + output: Some(output), + status, + }, + ))); + } + ChildTool::Subagent { + content, + background, + } => { + // A background launch only acknowledges the spawn; the + // nested subagent settles through its own task_notification. + if background || launch_acknowledged(value, &output) { + continue; + } + let ItemContent::Subagent { + agent_type, + description, + model, + effort, + .. + } = content + else { + continue; + }; + events.push(AgentEvent::ItemCompleted(self.nested_subagent( + id, + ItemContent::Subagent { + agent_type, + description, + status, + summary: (!output.trim().is_empty()).then(|| one_line_summary(&output)), + model, + effort, + }, + ))); + } + } } events } @@ -174,6 +246,26 @@ impl TranscriptMapper { content, } } + + /// A nested spawn keeps its bare tool_use id: the grandchild's transcript + /// lines name it as their `parent_tool_use_id`. + fn nested_subagent(&self, tool_use_id: &str, content: ItemContent) -> ThreadItem { + ThreadItem { + id: tool_use_id.to_owned(), + parent_item_id: Some(self.parent_id.clone()), + content, + } + } +} + +/// Whether a nested Agent result only reports an asynchronous launch. Claude +/// writes `toolUseResult.status` for top-level agents but omits it from +/// subagent transcripts, where the acknowledgement text is the only marker. +fn launch_acknowledged(value: &Value, output: &str) -> bool { + ["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/tool_use_result/status", "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/toolUseResult/status"] + .iter() + .any(|pointer| value.pointer(pointer).and_then(Value::as_str) == Some("async_launched")) + || output.starts_with("Async agent launched successfully") } /// Incremental reader used both by the polling task and deterministic tests. diff --git a/crates/agent/tests/fixtures/claude/subagent_nested_trace.jsonl b/crates/agent/tests/fixtures/claude/subagent_nested_trace.jsonl new file mode 100644 index 00000000..10a73a87 --- /dev/null +++ b/crates/agent/tests/fixtures/claude/subagent_nested_trace.jsonl @@ -0,0 +1,15 @@ +{"type":"assistant","message":{"id":"msg-spawn","content":[{"type":"tool_use","id":"toolu_child","name":"Agent","input":{"description":"Research hosting","prompt":"Research hosting options and report.","subagent_type":"general-purpose","run_in_background":true}}]}} +{"type":"system","subtype":"task_started","task_id":"task-child","tool_use_id":"toolu_child","description":"Research hosting","subagent_type":"general-purpose","task_type":"local_agent","prompt":"Research hosting options and report.","session_id":"session-nested"} +{"type":"user","message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"toolu_child","content":[{"type":"text","text":"Async agent launched successfully. (This tool result is internal metadata.)\nagentId: a-child\nThe agent is working in the background. You will be notified automatically when it completes."}],"is_error":false}]},"tool_use_result":{"isAsync":true,"status":"async_launched","agentId":"a-child","description":"Research hosting","resolvedModel":"claude-sonnet-5"}} +{"type":"result","subtype":"success","is_error":false,"result":"Started the research in the background.","usage":{"input_tokens":10,"output_tokens":5}} +{"type":"user","parent_tool_use_id":"toolu_child","message":{"role":"user","content":"Research hosting options and report."}} +{"type":"assistant","parent_tool_use_id":"toolu_child","effort":"high","message":{"id":"msg-child-1","model":"claude-sonnet-5","role":"assistant","content":[{"type":"tool_use","id":"toolu_grandchild","name":"Agent","input":{"description":"Research Workers limits","subagent_type":"general-purpose","run_in_background":true,"prompt":"Fetch the Workers limits page and report the free-plan request cap."}}]}} +{"type":"system","subtype":"task_started","task_id":"task-grandchild","tool_use_id":"toolu_grandchild","description":"Research Workers limits","subagent_type":"general-purpose","task_type":"local_agent","prompt":"Fetch the Workers limits page and report the free-plan request cap.","session_id":"session-nested"} +{"type":"user","parent_tool_use_id":"toolu_child","message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"toolu_grandchild","content":[{"type":"text","text":"Async agent launched successfully. (This tool result is internal metadata.)\nagentId: a-grandchild\nThe agent is working in the background. You will be notified automatically when it completes."}],"is_error":false}]}} +{"type":"user","parent_tool_use_id":"toolu_grandchild","message":{"role":"user","content":"Fetch the Workers limits page and report the free-plan request cap."}} +{"type":"assistant","parent_tool_use_id":"toolu_grandchild","effort":"high","message":{"id":"msg-grandchild-1","model":"claude-sonnet-5","role":"assistant","content":[{"type":"tool_use","id":"toolu_gc_fetch","name":"WebFetch","input":{"url":"https://developers.cloudflare.com/workers/platform/limits/","prompt":"free plan request cap"}}]}} +{"type":"user","parent_tool_use_id":"toolu_grandchild","message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"toolu_gc_fetch","content":"Free plan: 100,000 requests per day.","is_error":false}]}} +{"type":"assistant","parent_tool_use_id":"toolu_grandchild","message":{"id":"msg-grandchild-2","model":"claude-sonnet-5","role":"assistant","content":[{"type":"text","text":"Workers Free allows 100,000 requests per day."}]}} +{"type":"system","subtype":"task_notification","task_id":"task-grandchild","tool_use_id":"toolu_grandchild","status":"completed","output_file":"/tmp/tcode-synthetic/tasks/task-grandchild.output","summary":"Agent \"Research Workers limits\" finished","usage":{"total_tokens":12,"tool_uses":1,"duration_ms":34}} +{"type":"assistant","parent_tool_use_id":"toolu_child","message":{"id":"msg-child-2","model":"claude-sonnet-5","role":"assistant","content":[{"type":"text","text":"Hosting research done."}]}} +{"type":"system","subtype":"task_notification","task_id":"task-child","tool_use_id":"toolu_child","status":"completed","output_file":"/tmp/tcode-synthetic/tasks/task-child.output","summary":"Hosting researched","usage":{"total_tokens":40,"tool_uses":2,"duration_ms":90}} diff --git a/crates/runtime/src/app/lifecycle.rs b/crates/runtime/src/app/lifecycle.rs index e800aa81..e9531c27 100644 --- a/crates/runtime/src/app/lifecycle.rs +++ b/crates/runtime/src/app/lifecycle.rs @@ -1,3 +1,4 @@ +use super::sessions::descendant_session_ids; use super::*; impl AppState { @@ -520,10 +521,16 @@ impl AppState { } pub(super) fn close_orchestrator_children(&mut self, parent_id: &str, cx: &mut HostCx) { + // Nested subagent mirrors hang under other mirrors, yet they belong to + // the one provider process that just closed. + let descendants = descendant_session_ids(&self.sessions, parent_id); let child_ids: Vec<_> = self .sessions .iter() - .filter(|meta| meta.parent_session_id.as_deref() == Some(parent_id)) + .filter(|meta| { + meta.parent_session_id.as_deref() == Some(parent_id) + || (meta.native_subagent.is_some() && descendants.contains(&meta.id)) + }) .map(|meta| meta.id.clone()) .collect(); for child_id in child_ids { diff --git a/crates/runtime/src/app/mod.rs b/crates/runtime/src/app/mod.rs index 78cfeae2..2a707b09 100644 --- a/crates/runtime/src/app/mod.rs +++ b/crates/runtime/src/app/mod.rs @@ -339,6 +339,9 @@ pub struct AppState { native_subagent_sessions: HashMap<(String, String), String>, /// Synthetic turn state survives eviction; false remembers a finished child. native_subagent_turns: HashMap, + /// Subagent items recorded inside mirrors, by (session, item id): the + /// spawns a nested mirror is created from once a grandchild item arrives. + nested_subagent_spawns: HashMap<(String, String), subagents::NestedSpawn>, pub settings: Settings, pub providers: ProviderCatalog, terminal_preferences_path: PathBuf, @@ -509,6 +512,7 @@ impl AppState { pending_native_rewinds: HashMap::new(), native_subagent_sessions: HashMap::new(), native_subagent_turns: HashMap::new(), + nested_subagent_spawns: HashMap::new(), settings, providers: ProviderCatalog::new(model_catalogs, provider_secret_names), terminal_preferences_path, diff --git a/crates/runtime/src/app/sessions.rs b/crates/runtime/src/app/sessions.rs index c407b9fd..1b79c590 100644 --- a/crates/runtime/src/app/sessions.rs +++ b/crates/runtime/src/app/sessions.rs @@ -1467,7 +1467,7 @@ impl AppState { let git_branch = load_branch.then(|| read_git_branch(&cwd)); (timeline, folded, loaded, git_branch) }; - host_cx.enqueue(move |state, _cx| { + host_cx.enqueue(move |state, cx| { let generation_matches = state.timeline_load_generations.get(&session_id).copied() == Some(generation); let target_matches = match target { @@ -1502,6 +1502,7 @@ impl AppState { session.git_branch = git_branch; } } + state.repair_orphaned_mirror_turn(&session_id, cx); }); }); } diff --git a/crates/runtime/src/app/subagents.rs b/crates/runtime/src/app/subagents.rs index 0b19f4fa..dbc852a8 100644 --- a/crates/runtime/src/app/subagents.rs +++ b/crates/runtime/src/app/subagents.rs @@ -1,5 +1,15 @@ +use super::sessions::descendant_session_ids; use super::*; +/// A Subagent item recorded inside a mirror: where it was recorded and its +/// latest snapshot, so the grandchild's mirror can be titled and settled once +/// the grandchild's own transcript names it. +#[derive(Clone)] +pub(super) struct NestedSpawn { + mirror_id: String, + item: ThreadItem, +} + impl AppState { /// Reroute provider-native subagent transcript items into read-only mirror /// sessions. Returns true when the event was consumed instead of entering @@ -14,11 +24,11 @@ impl AppState { return false; }; - // Check child ownership before content: Codex child transcript items can - // themselves be Subagent-shaped, but they remain plain mirror content. + // Check child ownership before content: a Subagent item inside a mirror + // is content of that mirror, and only becomes a nested mirror of its + // own once a grandchild item names it as parent. if let Some(parent_item_id) = item.parent_item_id.as_deref() { - let mirror_id = - self.ensure_native_subagent_mirror(parent_session_id, parent_item_id, None, cx); + let mirror_id = self.child_item_mirror(parent_session_id, parent_item_id, cx); let Some(mirror_id) = mirror_id else { log::warn!( "dropping native subagent child item {}: parent session {} is missing", @@ -37,15 +47,26 @@ impl AppState { self.sync_mirror_turn(&mirror_id, true, parent_item_id, TurnStatus::Completed, cx); } self.record_event(&mirror_id, &strip_parent_item_id(event), cx); + if matches!(item.content, ItemContent::Subagent { .. }) { + self.nested_subagent_spawns.insert( + (parent_session_id.to_string(), item.id.clone()), + NestedSpawn { + mirror_id, + item: item.clone(), + }, + ); + if let Some(grandchild_mirror_id) = + self.find_native_subagent_mirror(parent_session_id, &item.id, cx) + { + self.apply_subagent_snapshot(&grandchild_mirror_id, item, cx); + } + } return true; } let ItemContent::Subagent { agent_type, description, - status, - model, - effort, .. } = &item.content else { @@ -55,58 +76,111 @@ impl AppState { if let Some(mirror_id) = self.ensure_native_subagent_mirror(parent_session_id, &item.id, Some(details), cx) { - let in_progress = matches!(status, ItemStatus::InProgress); - let title = mirror_title(agent_type, description); - let meta = self.resident_mut(&mirror_id).map(|mirror| { - let title_changed = mirror.meta.title == "subagent" && mirror.meta.title != title; - if title_changed { - mirror.meta.title = title; - } - let mut settings_changed = false; - if let Some(model) = model - && mirror.meta.model.as_ref() != Some(model) + self.apply_subagent_snapshot(&mirror_id, item, cx); + } + false + } + + /// The mirror that receives items parented to `parent_item_id`: the + /// subagent's own mirror, a nested spawn's mirror under the mirror it was + /// spawned in, or — for a subagent this session never announced — a + /// placeholder that closes with the parent process. + fn child_item_mirror( + &mut self, + session_id: &str, + parent_item_id: &str, + cx: &mut HostCx, + ) -> Option { + if let Some(id) = self.find_native_subagent_mirror(session_id, parent_item_id, cx) { + return Some(id); + } + let key = (session_id.to_string(), parent_item_id.to_string()); + let Some(spawn) = self.nested_subagent_spawns.get(&key).cloned() else { + return self.ensure_native_subagent_mirror(session_id, parent_item_id, None, cx); + }; + let ItemContent::Subagent { + agent_type, + description, + .. + } = &spawn.item.content + else { + return None; + }; + let mirror_id = self.ensure_native_subagent_mirror( + &spawn.mirror_id, + parent_item_id, + Some((agent_type, description)), + cx, + )?; + self.native_subagent_sessions.insert(key, mirror_id.clone()); + self.apply_subagent_snapshot(&mirror_id, &spawn.item, cx); + Some(mirror_id) + } + + /// Fold a Subagent snapshot into its mirror: title, model and effort on + /// the metadata, and the synthesized turn from its status. + fn apply_subagent_snapshot(&mut self, mirror_id: &str, item: &ThreadItem, cx: &mut HostCx) { + let ItemContent::Subagent { + agent_type, + description, + status, + model, + effort, + .. + } = &item.content + else { + return; + }; + let in_progress = matches!(status, ItemStatus::InProgress); + let title = mirror_title(agent_type, description); + let meta = self.resident_mut(mirror_id).map(|mirror| { + let title_changed = mirror.meta.title == "subagent" && mirror.meta.title != title; + if title_changed { + mirror.meta.title = title; + } + let mut settings_changed = false; + if let Some(model) = model + && mirror.meta.model.as_ref() != Some(model) + { + mirror.meta.model = Some(model.clone()); + settings_changed = true; + } + if let Some(effort) = effort { + let value = serde_json::Value::String(effort.clone()); + if let Some(selection) = mirror + .meta + .option_selections + .iter_mut() + .find(|selection| selection.id == "reasoningEffort") { - mirror.meta.model = Some(model.clone()); - settings_changed = true; - } - if let Some(effort) = effort { - let value = serde_json::Value::String(effort.clone()); - if let Some(selection) = mirror - .meta - .option_selections - .iter_mut() - .find(|selection| selection.id == "reasoningEffort") - { - if selection.value != value { - selection.value = value; - settings_changed = true; - } - } else { - mirror.meta.option_selections.push(OptionSelection { - id: "reasoningEffort".into(), - value, - }); + if selection.value != value { + selection.value = value; settings_changed = true; } - } - if !in_progress || title_changed || settings_changed { - mirror.meta.updated_at = now_secs(); - Some(mirror.meta.clone()) } else { - None + mirror.meta.option_selections.push(OptionSelection { + id: "reasoningEffort".into(), + value, + }); + settings_changed = true; } - }); - if let Some(meta) = meta.flatten() { - self.persist_meta(&meta, cx); } - let turn_status = match status { - ItemStatus::Failed | ItemStatus::Declined => TurnStatus::Failed, - ItemStatus::Interrupted => TurnStatus::Interrupted, - _ => TurnStatus::Completed, - }; - self.sync_mirror_turn(&mirror_id, in_progress, &item.id, turn_status, cx); + if !in_progress || title_changed || settings_changed { + mirror.meta.updated_at = now_secs(); + Some(mirror.meta.clone()) + } else { + None + } + }); + if let Some(meta) = meta.flatten() { + self.persist_meta(&meta, cx); } - false + let turn_status = match status { + ItemStatus::Failed | ItemStatus::Declined => TurnStatus::Failed, + ItemStatus::Interrupted => TurnStatus::Interrupted, + _ => TurnStatus::Completed, + }; + self.sync_mirror_turn(mirror_id, in_progress, &item.id, turn_status, cx); } /// Mirrors never receive provider Turn events, but the chat view derives @@ -150,34 +224,50 @@ impl AppState { } } - fn ensure_native_subagent_mirror( + /// The existing mirror for `subagent_item_id` anywhere under `session_id`: + /// a direct child, or a nested spawn's mirror under another mirror. + fn find_native_subagent_mirror( &mut self, - parent_session_id: &str, + session_id: &str, subagent_item_id: &str, - details: Option<(&str, &str)>, cx: &mut HostCx, ) -> Option { - let key = (parent_session_id.to_string(), subagent_item_id.to_string()); + let key = (session_id.to_string(), subagent_item_id.to_string()); if let Some(id) = self.native_subagent_sessions.get(&key).cloned() { return Some(id); } - - if let Some(meta) = self + let descendants = descendant_session_ids(&self.sessions, session_id); + let meta = self .sessions .iter() .find(|meta| { - meta.parent_session_id.as_deref() == Some(parent_session_id) - && meta.native_subagent.as_deref() == Some(subagent_item_id) + meta.native_subagent.as_deref() == Some(subagent_item_id) + && meta + .parent_session_id + .as_ref() + .is_some_and(|parent| descendants.contains(parent)) }) - .cloned() + .cloned()?; + let id = meta.id.clone(); + if self.resident(&id).is_none() { + self.load_background_session(meta, cx); + } + self.native_subagent_sessions.insert(key, id.clone()); + Some(id) + } + + fn ensure_native_subagent_mirror( + &mut self, + parent_session_id: &str, + subagent_item_id: &str, + details: Option<(&str, &str)>, + cx: &mut HostCx, + ) -> Option { + if let Some(id) = self.find_native_subagent_mirror(parent_session_id, subagent_item_id, cx) { - let id = meta.id.clone(); - if self.resident(&id).is_none() { - self.load_background_session(meta, cx); - } - self.native_subagent_sessions.insert(key, id.clone()); return Some(id); } + let key = (parent_session_id.to_string(), subagent_item_id.to_string()); let parent = self.find_meta(parent_session_id)?.clone(); let mut meta = SessionMeta::new(parent.provider, parent.cwd.clone(), parent.model.clone()); @@ -223,17 +313,20 @@ impl AppState { } /// The parent process is gone, so no Subagent status will ever close the - /// mirrors it was still running. + /// mirrors it was still running, nested ones included. pub(super) fn interrupt_native_subagent_work( &mut self, parent_session_id: &str, cx: &mut HostCx, ) { + self.nested_subagent_spawns + .retain(|(session_id, _), _| session_id != parent_session_id); + let descendants = descendant_session_ids(&self.sessions, parent_session_id); let running: Vec<_> = self .sessions .iter() .filter(|meta| { - meta.parent_session_id.as_deref() == Some(parent_session_id) + descendants.contains(&meta.id) && self.native_subagent_turns.get(&meta.id) == Some(&true) }) .filter_map(|meta| Some((meta.id.clone(), meta.native_subagent.clone()?))) @@ -255,6 +348,42 @@ impl AppState { } } } + + /// A mirror loaded with its last turn still open, while no live parent is + /// tracking it as running, was orphaned by a host restart or a parent that + /// left without a close: end the turn so it stops reporting work nothing + /// can finish. + pub(super) fn repair_orphaned_mirror_turn(&mut self, mirror_id: &str, cx: &mut HostCx) { + if self.native_subagent_turns.get(mirror_id) == Some(&true) { + return; + } + let Some(mirror) = self.resident(mirror_id) else { + return; + }; + let Some(subagent_item_id) = mirror.meta.native_subagent.clone() else { + return; + }; + let Some(turn) = mirror + .timeline + .turns + .last() + .filter(|turn| turn.status.is_none()) + else { + return; + }; + let turn_id = turn.provider_turn_id.clone().unwrap_or(subagent_item_id); + self.native_subagent_turns + .insert(mirror_id.to_string(), false); + self.record_event( + mirror_id, + &AgentEvent::TurnCompleted { + turn_id, + status: TurnStatus::Interrupted, + usage: None, + }, + cx, + ); + } } fn lifecycle_item(event: &AgentEvent) -> Option<&ThreadItem> { diff --git a/crates/runtime/src/app/tests.rs b/crates/runtime/src/app/tests.rs index cbf3fbec..a8a3f797 100644 --- a/crates/runtime/src/app/tests.rs +++ b/crates/runtime/src/app/tests.rs @@ -424,6 +424,359 @@ fn native_mirror_keeps_one_turn_across_residency_late_items_and_parent_completio } } } + +fn subagent_item( + id: &str, + parent_item_id: Option<&str>, + agent_type: &str, + description: &str, + status: ItemStatus, +) -> ThreadItem { + ThreadItem { + id: id.into(), + parent_item_id: parent_item_id.map(str::to_owned), + content: ItemContent::Subagent { + agent_type: agent_type.into(), + description: description.into(), + status, + summary: None, + model: None, + effort: None, + }, + } +} + +/// A subagent's transcript can itself spawn a subagent. The nested spawn is +/// content of the child mirror, and once the grandchild's own transcript +/// arrives it gets a titled mirror under the child mirror — never a nameless +/// one on the root — whose turn closes with the nested spawn's terminal status. +#[test] +fn nested_subagent_items_open_a_titled_mirror_under_the_child_mirror_and_close_it() { + let cx = &mut TestAppContext::default(); + let test_store = TestStore::new("tcode-nested-native-mirror"); + let state = cx.new_entity(TestClientState::new((*test_store).clone())); + let (child_id, grandchild_id, second_id) = state.update(cx, |state, cx| { + let mut meta = SessionMeta::new(ProviderKind::ClaudeCode, PathBuf::from("/tmp"), None); + meta.id = "parent".into(); + state.sessions.push(meta.clone()); + state.install_selected(ActiveSession::new(meta, false, Vec::new())); + state.on_event( + "parent", + AgentEvent::ItemStarted(subagent_item( + "toolu_child", + None, + "general-purpose", + "Research hosting", + ItemStatus::InProgress, + )), + cx, + ); + state.on_event( + "parent", + AgentEvent::ItemStarted(subagent_item( + "toolu_grandchild", + Some("toolu_child"), + "Explore", + "Research Workers limits\nwith sources", + ItemStatus::InProgress, + )), + cx, + ); + let child_id = state + .sessions + .iter() + .find(|meta| meta.native_subagent.as_deref() == Some("toolu_child")) + .unwrap() + .id + .clone(); + assert!( + state + .sessions + .iter() + .all(|meta| meta.native_subagent.as_deref() != Some("toolu_grandchild")), + "a spawn item alone is content of the child mirror, not a mirror" + ); + assert!( + state + .resident(&child_id) + .unwrap() + .timeline + .entries + .iter() + .any(|entry| { + entry.id == "toolu_grandchild" + && matches!( + entry.content, + EntryContent::Item(ItemContent::Subagent { .. }) + ) + }) + ); + + state.on_event( + "parent", + AgentEvent::ItemCompleted(ThreadItem { + id: "toolu_grandchild:user-1".into(), + parent_item_id: Some("toolu_grandchild".into()), + content: ItemContent::UserMessage { + text: "Fetch the limits page.".into(), + context_len: None, + attachments: Vec::new(), + }, + }), + cx, + ); + let grandchild = state + .sessions + .iter() + .find(|meta| meta.native_subagent.as_deref() == Some("toolu_grandchild")) + .cloned() + .expect("grandchild mirror"); + assert_eq!( + grandchild.parent_session_id.as_deref(), + Some(child_id.as_str()) + ); + assert_eq!(grandchild.title, "Explore: Research Workers limits"); + assert!(state.resident(&grandchild.id).unwrap().has_work()); + assert!(state.resident(&child_id).unwrap().has_work()); + assert_eq!( + state + .sessions + .iter() + .filter(|meta| meta.parent_session_id.as_deref() == Some("parent")) + .count(), + 1, + "the root session owns only the child it spawned" + ); + + // The grandchild's terminal snapshot travels as child content too. + state.on_event( + "parent", + AgentEvent::ItemUpdated(subagent_item( + "toolu_grandchild", + Some("toolu_child"), + "Explore", + "Research Workers limits", + ItemStatus::Completed, + )), + cx, + ); + assert!(!state.resident(&grandchild.id).unwrap().has_work()); + assert!(state.resident(&child_id).unwrap().has_work()); + state.on_event( + "parent", + AgentEvent::ItemCompleted(subagent_item( + "toolu_child", + None, + "general-purpose", + "Research hosting", + ItemStatus::Completed, + )), + cx, + ); + assert!(!state.resident(&child_id).unwrap().has_work()); + assert_eq!(state.resident("parent").unwrap().timeline.entries.len(), 1); + + // A second nested spawn is still running when the root process closes: + // it belongs to that process, so it ends with it. + for event in [ + AgentEvent::ItemStarted(subagent_item( + "toolu_grandchild_2", + Some("toolu_child"), + "Explore", + "Check pricing", + ItemStatus::InProgress, + )), + AgentEvent::ItemCompleted(ThreadItem { + id: "toolu_grandchild_2:msg:0".into(), + parent_item_id: Some("toolu_grandchild_2".into()), + content: ItemContent::AssistantMessage { + text: "Fetching pricing.".into(), + }, + }), + ] { + state.on_event("parent", event, cx); + } + let second = state + .sessions + .iter() + .find(|meta| meta.native_subagent.as_deref() == Some("toolu_grandchild_2")) + .cloned() + .expect("second grandchild mirror"); + assert_eq!(second.parent_session_id.as_deref(), Some(child_id.as_str())); + assert!(state.resident(&second.id).unwrap().has_work()); + state.on_event("parent", AgentEvent::SessionClosed { reason: None }, cx); + assert!( + !state + .resident(&second.id) + .is_some_and(ActiveSession::has_work) + ); + (child_id, grandchild.id, second.id) + }); + cx.run_until_parked(); + state.update(cx, |state, _| { + let events = state.store.read_events(&grandchild_id); + assert!(matches!( + &events.first().unwrap().event, + AgentEvent::TurnStarted { turn_id } if turn_id == "toolu_grandchild" + )); + assert!(matches!( + &events.last().unwrap().event, + AgentEvent::TurnCompleted { + status: TurnStatus::Completed, + .. + } + )); + assert!(matches!( + &state.store.read_events(&second_id).last().unwrap().event, + AgentEvent::TurnCompleted { turn_id, status: TurnStatus::Interrupted, .. } + if turn_id == "toolu_grandchild_2" + )); + assert!( + state + .store + .read_events(&child_id) + .iter() + .all(|stored| match &stored.event { + AgentEvent::ItemStarted(item) + | AgentEvent::ItemUpdated(item) + | AgentEvent::ItemCompleted(item) => item.parent_item_id.is_none(), + _ => true, + }) + ); + }); +} + +/// A child item whose spawn this session never announced still gets a mirror, +/// and the parent process closing ends it like every other running mirror. +#[test] +fn unknown_parent_mirror_closes_when_the_parent_session_closes() { + let cx = &mut TestAppContext::default(); + let test_store = TestStore::new("tcode-unknown-parent-native-mirror"); + let state = cx.new_entity(TestClientState::new((*test_store).clone())); + let mirror_id = state.update(cx, |state, cx| { + let mut meta = SessionMeta::new(ProviderKind::ClaudeCode, PathBuf::from("/tmp"), None); + meta.id = "parent".into(); + state.sessions.push(meta.clone()); + state.install_selected(ActiveSession::new(meta, false, Vec::new())); + state.on_event( + "parent", + AgentEvent::ItemCompleted(ThreadItem { + id: "toolu_orphan:msg:0".into(), + parent_item_id: Some("toolu_orphan".into()), + content: ItemContent::AssistantMessage { + text: "working".into(), + }, + }), + cx, + ); + let mirror = state + .sessions + .iter() + .find(|meta| meta.native_subagent.as_deref() == Some("toolu_orphan")) + .cloned() + .expect("placeholder mirror"); + assert_eq!(mirror.title, "subagent"); + assert!(state.resident(&mirror.id).unwrap().has_work()); + state.on_event("parent", AgentEvent::SessionClosed { reason: None }, cx); + assert!( + !state + .resident(&mirror.id) + .is_some_and(ActiveSession::has_work) + ); + mirror.id + }); + cx.run_until_parked(); + state.update(cx, |state, _| { + let events = state.store.read_events(&mirror_id); + assert!(matches!( + &events.last().unwrap().event, + AgentEvent::TurnCompleted { turn_id, status: TurnStatus::Interrupted, .. } + if turn_id == "toolu_orphan" + )); + }); +} + +/// A mirror persisted with an open turn by a host that stopped before the +/// subagent settled: loading it after a restart ends the turn, because no +/// live parent is tracking it and nothing else ever will. +#[test] +fn loading_a_mirror_with_an_open_turn_and_no_live_parent_ends_it() { + let cx = &mut TestAppContext::default(); + let test_store = TestStore::new("tcode-orphaned-native-mirror-repair"); + let mut parent = SessionMeta::new(ProviderKind::ClaudeCode, PathBuf::from("/tmp"), None); + parent.id = "parent".into(); + let mut mirror = SessionMeta::new(ProviderKind::ClaudeCode, PathBuf::from("/tmp"), None); + mirror.id = "mirror".into(); + mirror.title = "subagent".into(); + mirror.parent_session_id = Some("parent".into()); + mirror.native_subagent = Some("toolu_zombie".into()); + test_store.upsert_meta(&parent).unwrap(); + test_store.upsert_meta(&mirror).unwrap(); + let stored = [ + AgentEvent::TurnStarted { + turn_id: "toolu_zombie".into(), + }, + AgentEvent::ItemCompleted(ThreadItem { + id: "toolu_zombie:msg:0".into(), + parent_item_id: None, + content: ItemContent::AssistantMessage { + text: "Research complete.".into(), + }, + }), + ]; + for (i, event) in stored.iter().enumerate() { + test_store + .append_event("mirror", 1_000 + i as u64, event) + .unwrap(); + } + + let state = cx.new_entity(TestClientState::new((*test_store).clone())); + state.update(cx, |state, cx| { + assert_eq!(state.sessions.len(), 2); + state.select_session("mirror", cx); + }); + cx.run_until(|state| { + state + .resident("mirror") + .is_some_and(|mirror| mirror.timeline.turns.len() == 1) + }); + cx.run_until_parked(); + state.update(cx, |state, _| { + let mirror = state.resident("mirror").unwrap(); + assert!(!mirror.has_work()); + assert!(!mirror.timeline.turn_running); + assert_eq!( + mirror.timeline.turns[0].status, + Some(TurnStatus::Interrupted) + ); + assert_eq!( + mirror.timeline.last_turn_status, + Some(TurnStatus::Interrupted) + ); + let events = state.store.read_events("mirror"); + assert_eq!(events.len(), 3); + assert!(matches!( + &events[2].event, + AgentEvent::TurnCompleted { turn_id, status: TurnStatus::Interrupted, .. } + if turn_id == "toolu_zombie" + )); + assert_eq!(state.native_subagent_turns.get("mirror"), Some(&false)); + }); + + // Reloading the repaired mirror leaves the closed turn alone. + let state = cx.new_entity(TestClientState::new((*test_store).clone())); + state.update(cx, |state, cx| state.select_session("mirror", cx)); + cx.run_until(|state| { + state + .resident("mirror") + .is_some_and(|mirror| mirror.timeline.turns.len() == 1) + }); + cx.run_until_parked(); + state.update(cx, |state, _| { + assert_eq!(state.store.read_events("mirror").len(), 3); + }); +} + #[test] fn settings_patches_preserve_top_level_and_nested_siblings_over_the_pipe() { let cx = &mut TestAppContext::default();