Skip to main content

cosh_tools/subagent/
internal.rs

1//! Synthesis of typed [`SubagentEvent`]s for the INTERNAL sub-agent.
2//!
3//! The external ACP harnesses emit their activity natively: the `session/update`
4//! stream is mapped into [`SubagentEvent`]s by [`super::events`]. The internal
5//! sub-agent is a nested cosh harness and speaks no ACP — its loop emits plain
6//! [`crate::harness`-level] events (tool calls, results, context info). This
7//! module bridges that gap: it turns the nested loop's event stream into the
8//! SAME typed events the TUI sub-agent box already renders, so both sub-agent
9//! flavors look and behave identically (same tool rows, same spinner lifecycle,
10//! same plan mini-chat, same context usage footer).
11//!
12//! The nested loop has no call ids in its events, so ids are synthesized
13//! per announcement and paired with results FIFO: the harness guarantees the
14//! Nth `ToolResult`/`ToolError` matches the Nth `ToolCall` (sequential
15//! ordering, see the harness event docs). A result arriving without a
16//! matching announcement is dropped rather than mis-attributed.
17
18use std::collections::VecDeque;
19
20use serde_json::Value;
21
22use super::events::{
23    PlanEntry, PlanEntryPriority, PlanEntryStatus, SubagentEvent, ToolCallStatus, ToolDiffSummary,
24    ToolKind, ToolOutputBlock,
25};
26
27/// Byte cap of one tool result's displayed tail. Mirrors the TUI's own
28/// `SUBAGENT_TOOL_TAIL_BYTES`: the box keeps the TAIL of a tool output, so
29/// the bridge pre-caps the event payload the same way (the TUI re-bounds,
30/// the two never disagree — and a multi-megabyte bash dump never crosses
31/// the event channel whole).
32const RESULT_TAIL_BYTES: usize = 4_000;
33
34/// Stateful synthesizer of typed sub-agent events from a nested harness's
35/// plain event stream. One instance per internal sub-agent turn.
36#[derive(Debug, Default)]
37pub struct InternalEventBridge {
38    /// Monotonic counter backing the synthesized call ids.
39    issued: usize,
40    /// Announced calls awaiting their result, FIFO: `(id, tool name)`.
41    /// The harness guarantees the Nth result matches the Nth call.
42    outstanding: VecDeque<(String, String)>,
43}
44
45impl InternalEventBridge {
46    /// A fresh bridge for one internal sub-agent turn.
47    #[must_use]
48    pub fn new() -> Self {
49        Self::default()
50    }
51
52    /// The nested loop streamed an assistant text chunk: the box's message
53    /// entry (coalescing handled by the TUI, same as the ACP path).
54    #[must_use]
55    pub fn message(text: String) -> SubagentEvent {
56        SubagentEvent::Message { text }
57    }
58
59    /// The nested loop streamed a reasoning chunk: the box's thought entry.
60    #[must_use]
61    pub fn thought(text: String) -> SubagentEvent {
62        SubagentEvent::Thought { text }
63    }
64
65    /// A tool call was announced by the nested loop: a `ToolCall` event with
66    /// a synthesized id, the tool-derived kind and title, and the raw input
67    /// (the TUI extracts the same detail keys — path, query, command — it
68    /// shows for ACP calls).
69    pub fn tool_call(&mut self, tool: &str, input: &Value) -> SubagentEvent {
70        let id = format!("internal-{}", self.issued);
71        self.issued += 1;
72        self.outstanding.push_back((id.clone(), tool.to_string()));
73        SubagentEvent::ToolCall {
74            id,
75            title: tool.to_string(),
76            kind: tool_kind(tool),
77            status: ToolCallStatus::InProgress,
78            raw_input: Some(input.clone()),
79        }
80    }
81
82    /// A tool call completed with a serialized output: the FIFO-oldest
83    /// outstanding call is patched to `Completed`, carrying the output tail
84    /// and — for `fs_edit`-family results — a one-line diff summary parsed
85    /// from the result's unified diff, exactly what an ACP `ToolCallUpdate`
86    /// with a diff block renders.
87    ///
88    /// `None` when no call is outstanding (an unpaired result cannot be
89    /// attributed to a row).
90    pub fn tool_result(&mut self, output: &str) -> Option<SubagentEvent> {
91        let (id, tool) = self.outstanding.pop_front()?;
92        Some(SubagentEvent::ToolCallUpdate {
93            id,
94            status: Some(ToolCallStatus::Completed),
95            title: None,
96            raw_output: None,
97            content: ToolOutputBlock {
98                text: output_tail(output),
99                skipped: 0,
100                diff: diff_summary_from_result(tool, output),
101            },
102        })
103    }
104
105    /// A tool call failed during dispatch: the FIFO-oldest outstanding call
106    /// is patched to `Failed`, carrying the error message as the row's tail
107    /// (an ACP failed update renders its output text the same way).
108    /// `None` when no call is outstanding.
109    pub fn tool_error(&mut self, error: &str) -> Option<SubagentEvent> {
110        let (id, _) = self.outstanding.pop_front()?;
111        Some(SubagentEvent::ToolCallUpdate {
112            id,
113            status: Some(ToolCallStatus::Failed),
114            title: None,
115            raw_output: None,
116            content: ToolOutputBlock {
117                text: output_tail(error),
118                skipped: 0,
119                diff: None,
120            },
121        })
122    }
123
124    /// A `plan_todo_write` call's input mapped to the box's plan entry —
125    /// the same mini-chat TODO the ACP agents emit natively. `None` when
126    /// the input does not deserialize as a todo write (a malformed call is
127    /// still announced as a plain tool row by [`Self::tool_call`]).
128    #[must_use]
129    pub fn plan_from_todo_write(input: &Value) -> Option<SubagentEvent> {
130        #[derive(serde::Deserialize)]
131        struct TodoItem {
132            description: String,
133            #[serde(default)]
134            status: TodoStatus,
135        }
136        #[derive(serde::Deserialize, Default)]
137        #[serde(rename_all = "snake_case")]
138        enum TodoStatus {
139            #[default]
140            Pending,
141            InProgress,
142            Completed,
143            Cancelled,
144        }
145
146        let todos = input.get("todos").cloned()?;
147        let parsed: Vec<TodoItem> = serde_json::from_value(todos).ok()?;
148        let entries: Vec<PlanEntry> = parsed
149            .into_iter()
150            .map(|item| PlanEntry {
151                content: item.description,
152                priority: match item.status {
153                    TodoStatus::InProgress => PlanEntryPriority::High,
154                    _ => PlanEntryPriority::Medium,
155                },
156                status: match item.status {
157                    TodoStatus::InProgress => PlanEntryStatus::InProgress,
158                    TodoStatus::Completed => PlanEntryStatus::Completed,
159                    // The plan mirror has no "cancelled" state; a cancelled
160                    // task renders as the neutral not-started glyph.
161                    _ => PlanEntryStatus::Pending,
162                },
163            })
164            .collect();
165        (!entries.is_empty()).then_some(SubagentEvent::Plan { entries })
166    }
167
168    /// The nested harness's context snapshot mapped to the box's usage
169    /// footer (`context_window` = the budget, `tokens_in_context` = the
170    /// estimated tokens held) — the same numbers the ACP `UsageUpdate`
171    /// carries.
172    #[must_use]
173    pub fn usage(total_tokens: usize, max_tokens: usize) -> SubagentEvent {
174        SubagentEvent::Usage {
175            context_window: u64::try_from(max_tokens).unwrap_or(u64::MAX),
176            tokens_in_context: u64::try_from(total_tokens).unwrap_or(u64::MAX),
177        }
178    }
179}
180
181/// Display category of an internal tool name (the mirror the ACP kinds
182/// render through). File reads and lookups search, mutations edit, shell
183/// work executes, web work fetches.
184#[must_use]
185pub fn tool_kind(tool: &str) -> ToolKind {
186    match tool {
187        "fs_read" | "lsp_hover" => ToolKind::Read,
188        "fs_write" | "fs_edit" | "fs_edit_lines" | "fs_ast_edit" | "fs_rollback" | "lsp_rename"
189        | "lsp_code_actions" | "lsp_restart" => ToolKind::Edit,
190        "bash_run"
191        | "subagent_call"
192        | "computer_apps"
193        | "computer_snapshot"
194        | "computer_wait"
195        | "computer_screenshot"
196        | "computer_act"
197        | "computer_control" => ToolKind::Execute,
198        "find_glob"
199        | "find_grep"
200        | "recall_search"
201        | "skills_match_skills"
202        | "lsp_definitions"
203        | "lsp_references"
204        | "lsp_symbols"
205        | "lsp_workspace_symbols"
206        | "lsp_call_hierarchy" => ToolKind::Search,
207        "web_fetch" | "web_search" => ToolKind::Fetch,
208        "plan_todo_write" | "skills_list" | "skills_read" | "skills_read_asset" => ToolKind::Think,
209        // `ask_questions`/`stop_agent_loop` are blocklisted inside a
210        // sub-agent; anything else unknown renders as the generic tool row.
211        _ => ToolKind::Unknown,
212    }
213}
214
215/// The displayed tail of a tool output: the LAST `RESULT_TAIL_BYTES` bytes,
216/// cut at a character boundary (mirrors the TUI's bounded tail).
217fn output_tail(output: &str) -> String {
218    if output.len() <= RESULT_TAIL_BYTES {
219        return output.to_string();
220    }
221    let mut start = output.len() - RESULT_TAIL_BYTES;
222    while !output.is_char_boundary(start) {
223        start += 1;
224    }
225    output[start..].to_string()
226}
227
228/// One-line diff summary parsed from a serialized `fs_edit`-family result
229/// (a JSON array of per-target objects whose `diff` field is a unified
230/// diff). The LAST diff wins — one summary line per tool call, the same
231/// rule the ACP content extraction applies to multiple diff blocks.
232///
233/// `None` for anything that is not an edit result with a diff (bash text,
234/// grep listings, read output — those are not file modifications).
235fn diff_summary_from_result(tool: String, output: &str) -> Option<ToolDiffSummary> {
236    if !matches!(tool.as_str(), "fs_edit" | "fs_edit_lines" | "fs_ast_edit") {
237        return None;
238    }
239    let results: Vec<Value> = serde_json::from_str(output).ok()?;
240    results
241        .iter()
242        .filter_map(|entry| {
243            let diff = entry.get("diff")?.as_str()?;
244            let path = entry.get("path")?.as_str()?;
245            let mut added = 0_u32;
246            let mut removed = 0_u32;
247            for line in diff.lines() {
248                if line.starts_with('+') && !line.starts_with("+++") {
249                    added += 1;
250                } else if line.starts_with('-') && !line.starts_with("---") {
251                    removed += 1;
252                }
253            }
254            Some(ToolDiffSummary {
255                path: path.to_string(),
256                added,
257                removed,
258            })
259        })
260        .next_back()
261}
262
263#[cfg(test)]
264mod tests {
265    use super::*;
266
267    #[test]
268    fn message_and_thought_map_to_the_typed_stream_variants() {
269        assert!(matches!(
270            InternalEventBridge::message("hi".to_string()),
271            SubagentEvent::Message { text } if text == "hi"
272        ));
273        assert!(matches!(
274            InternalEventBridge::thought("pondering".to_string()),
275            SubagentEvent::Thought { text } if text == "pondering"
276        ));
277    }
278
279    #[test]
280    fn tool_call_announces_in_progress_with_a_synthesized_id() {
281        let mut bridge = InternalEventBridge::new();
282        let event = bridge.tool_call("fs_read", &serde_json::json!({"path": "src/a.rs"}));
283        match event {
284            SubagentEvent::ToolCall {
285                id,
286                title,
287                kind,
288                status,
289                raw_input,
290            } => {
291                assert_eq!(id, "internal-0");
292                assert_eq!(title, "fs_read");
293                assert_eq!(kind, ToolKind::Read);
294                assert_eq!(status, ToolCallStatus::InProgress);
295                assert_eq!(raw_input, Some(serde_json::json!({"path": "src/a.rs"})));
296            }
297            other => panic!("expected ToolCall, got {other:?}"),
298        }
299    }
300
301    #[test]
302    fn ids_are_unique_across_calls() {
303        let mut bridge = InternalEventBridge::new();
304        let first = bridge.tool_call("bash_run", &serde_json::json!({}));
305        let second = bridge.tool_call("bash_run", &serde_json::json!({}));
306        let id_of = |e: SubagentEvent| match e {
307            SubagentEvent::ToolCall { id, .. } => id,
308            other => panic!("expected ToolCall, got {other:?}"),
309        };
310        assert_ne!(id_of(first), id_of(second));
311    }
312
313    #[test]
314    fn result_pairs_with_the_oldest_outstanding_call() {
315        let mut bridge = InternalEventBridge::new();
316        bridge.tool_call("bash_run", &serde_json::json!({}));
317        bridge.tool_call("fs_read", &serde_json::json!({}));
318        let update = bridge.tool_result("ok").expect("first result pairs");
319        match update {
320            SubagentEvent::ToolCallUpdate {
321                id,
322                status,
323                content,
324                ..
325            } => {
326                assert_eq!(id, "internal-0");
327                assert_eq!(status, Some(ToolCallStatus::Completed));
328                assert_eq!(content.text, "ok");
329            }
330            other => panic!("expected ToolCallUpdate, got {other:?}"),
331        }
332        let update = bridge.tool_result("data").expect("second result pairs");
333        assert!(matches!(
334            update,
335            SubagentEvent::ToolCallUpdate { ref id, .. } if id == "internal-1"
336        ));
337    }
338
339    #[test]
340    fn unpaired_result_is_dropped() {
341        let mut bridge = InternalEventBridge::new();
342        assert!(bridge.tool_result("stray").is_none());
343        assert!(bridge.tool_error("stray").is_none());
344    }
345
346    #[test]
347    fn error_marks_the_call_failed() {
348        let mut bridge = InternalEventBridge::new();
349        bridge.tool_call("fs_write", &serde_json::json!({}));
350        let update = bridge.tool_error("permission denied").expect("error pairs");
351        assert!(matches!(
352            update,
353            SubagentEvent::ToolCallUpdate {
354                status: Some(ToolCallStatus::Failed),
355                ..
356            }
357        ));
358    }
359
360    #[test]
361    fn long_output_is_tail_capped_at_a_char_boundary() {
362        let long = format!("{}é", "x".repeat(RESULT_TAIL_BYTES));
363        let mut bridge = InternalEventBridge::new();
364        bridge.tool_call("bash_run", &serde_json::json!({}));
365        let update = bridge.tool_result(&long).expect("pairs");
366        match update {
367            SubagentEvent::ToolCallUpdate { content, .. } => {
368                assert!(content.text.len() <= RESULT_TAIL_BYTES + 2);
369                assert!(content.text.ends_with('é'));
370            }
371            other => panic!("expected ToolCallUpdate, got {other:?}"),
372        }
373    }
374
375    #[test]
376    fn edit_result_yields_a_diff_summary_last_diff_wins() {
377        let output = serde_json::json!([
378            {"path": "a.rs", "diff": "+one\n+two\n-three\n"},
379            {"path": "b.rs", "diff": "+only\n"}
380        ])
381        .to_string();
382        let mut bridge = InternalEventBridge::new();
383        bridge.tool_call("fs_edit", &serde_json::json!({}));
384        let update = bridge.tool_result(&output).expect("pairs");
385        match update {
386            SubagentEvent::ToolCallUpdate { content, .. } => {
387                let diff = content.diff.expect("edit carries a diff summary");
388                assert_eq!(diff.path, "b.rs");
389                assert_eq!(diff.added, 1);
390                assert_eq!(diff.removed, 0);
391            }
392            other => panic!("expected ToolCallUpdate, got {other:?}"),
393        }
394    }
395
396    #[test]
397    fn non_edit_results_carry_no_diff_summary() {
398        let mut bridge = InternalEventBridge::new();
399        bridge.tool_call("bash_run", &serde_json::json!({}));
400        let update = bridge
401            .tool_result("[{\"path\": \"a.rs\", \"diff\": \"+x\"}]")
402            .expect("pairs");
403        match update {
404            SubagentEvent::ToolCallUpdate { content, .. } => assert!(content.diff.is_none()),
405            other => panic!("expected ToolCallUpdate, got {other:?}"),
406        }
407    }
408
409    #[test]
410    fn malformed_edit_output_carries_no_diff_summary() {
411        let mut bridge = InternalEventBridge::new();
412        bridge.tool_call("fs_edit", &serde_json::json!({}));
413        let update = bridge.tool_result("not json").expect("pairs");
414        match update {
415            SubagentEvent::ToolCallUpdate { content, .. } => assert!(content.diff.is_none()),
416            other => panic!("expected ToolCallUpdate, got {other:?}"),
417        }
418    }
419
420    #[test]
421    fn todo_write_input_maps_to_a_plan_entry_per_todo() {
422        let input = serde_json::json!({"todos": [
423            {"description": "first", "status": "in_progress"},
424            {"description": "second", "status": "completed"},
425            {"description": "third", "status": "cancelled"}
426        ]});
427        let event = InternalEventBridge::plan_from_todo_write(&input).expect("plan");
428        match event {
429            SubagentEvent::Plan { entries } => {
430                assert_eq!(entries.len(), 3);
431                assert_eq!(entries[0].content, "first");
432                assert_eq!(entries[0].status, PlanEntryStatus::InProgress);
433                assert_eq!(entries[0].priority, PlanEntryPriority::High);
434                assert_eq!(entries[1].status, PlanEntryStatus::Completed);
435                assert_eq!(entries[2].status, PlanEntryStatus::Pending);
436            }
437            other => panic!("expected Plan, got {other:?}"),
438        }
439    }
440
441    #[test]
442    fn todo_write_without_todos_yields_no_plan() {
443        assert!(InternalEventBridge::plan_from_todo_write(&serde_json::json!({})).is_none());
444        assert!(
445            InternalEventBridge::plan_from_todo_write(&serde_json::json!({"todos": []})).is_none()
446        );
447        assert!(
448            InternalEventBridge::plan_from_todo_write(&serde_json::json!({"todos": "no"}))
449                .is_none()
450        );
451    }
452
453    #[test]
454    fn context_info_maps_to_the_usage_footer() {
455        let event = InternalEventBridge::usage(1_234, 10_000);
456        assert_eq!(
457            event,
458            SubagentEvent::Usage {
459                context_window: 10_000,
460                tokens_in_context: 1_234,
461            }
462        );
463    }
464
465    #[test]
466    fn tool_kinds_cover_the_internal_dispatch_table() {
467        assert_eq!(tool_kind("fs_read"), ToolKind::Read);
468        assert_eq!(tool_kind("fs_edit_lines"), ToolKind::Edit);
469        assert_eq!(tool_kind("bash_run"), ToolKind::Execute);
470        assert_eq!(tool_kind("find_grep"), ToolKind::Search);
471        assert_eq!(tool_kind("web_fetch"), ToolKind::Fetch);
472        assert_eq!(tool_kind("plan_todo_write"), ToolKind::Think);
473        assert_eq!(tool_kind("mystery"), ToolKind::Unknown);
474    }
475}