1use std::collections::VecDeque;
19
20use serde_json::Value;
21
22use super::events::{
23 PlanEntry, PlanEntryPriority, PlanEntryStatus, SubagentEvent, ToolCallStatus, ToolDiffSummary,
24 ToolKind, ToolOutputBlock,
25};
26
27const RESULT_TAIL_BYTES: usize = 4_000;
33
34#[derive(Debug, Default)]
37pub struct InternalEventBridge {
38 issued: usize,
40 outstanding: VecDeque<(String, String)>,
43}
44
45impl InternalEventBridge {
46 #[must_use]
48 pub fn new() -> Self {
49 Self::default()
50 }
51
52 #[must_use]
55 pub fn message(text: String) -> SubagentEvent {
56 SubagentEvent::Message { text }
57 }
58
59 #[must_use]
61 pub fn thought(text: String) -> SubagentEvent {
62 SubagentEvent::Thought { text }
63 }
64
65 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 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 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 #[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 _ => PlanEntryStatus::Pending,
162 },
163 })
164 .collect();
165 (!entries.is_empty()).then_some(SubagentEvent::Plan { entries })
166 }
167
168 #[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#[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 _ => ToolKind::Unknown,
212 }
213}
214
215fn 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
228fn 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}