1use std::{
4 collections::{BTreeMap, BTreeSet},
5 path::PathBuf,
6 sync::{
7 Arc, Mutex,
8 atomic::{AtomicBool, Ordering},
9 },
10};
11
12use async_trait::async_trait;
13use serde_json::Value;
14use tokio::{
15 io::{AsyncBufReadExt, BufReader},
16 process::{Child, Command},
17 sync::mpsc,
18};
19
20use crate::{
21 AgentCapabilities, AgentEvent, Mode, PermissionAnswer, RosterSlot, ToolStatus, ToolUpdate,
22};
23
24use super::native::{NativeTurn, spawn_native_turn};
25use super::{
26 AdapterError, AdapterResult, AgentAdapter, CANCEL_SETTLE_TIMEOUT, drain_bounded,
27 isolate_process_group, parse_command_line, terminate_child,
28};
29
30const MODEL_CONFIG_ID: &str = "claude:model";
31
32fn model_label(model: &str) -> String {
33 match model {
34 "default" => "Default".into(),
35 "best" => "Best".into(),
36 "fable" => "Fable".into(),
37 "opus" => "Opus".into(),
38 "sonnet" => "Sonnet".into(),
39 "haiku" => "Haiku".into(),
40 "sonnet[1m]" => "Sonnet (1M)".into(),
41 "opus[1m]" => "Opus (1M)".into(),
42 "opusplan" => "Opus plan".into(),
43 _ => model.to_owned(),
44 }
45}
46
47#[derive(Debug, Default)]
48struct ParserState {
49 tools: BTreeMap<String, ToolUpdate>,
50 finished_tools: BTreeSet<String>,
51 streamed_thoughts: BTreeMap<u64, String>,
52}
53
54fn content_text(value: &Value) -> Option<String> {
55 match value {
56 Value::String(text) => (!text.is_empty()).then(|| text.to_owned()),
57 Value::Array(content) => {
58 let text = content
59 .iter()
60 .filter_map(content_text)
61 .collect::<Vec<_>>()
62 .join("\n");
63 (!text.is_empty()).then_some(text)
64 }
65 Value::Object(content) => {
66 let direct = [
67 "text",
68 "stdout",
69 "stderr",
70 "error",
71 "error_code",
72 "error_message",
73 "message",
74 ]
75 .into_iter()
76 .filter_map(|key| content.get(key).and_then(Value::as_str))
77 .filter(|text| !text.is_empty())
78 .collect::<Vec<_>>()
79 .join("\n");
80 (!direct.is_empty())
81 .then_some(direct)
82 .or_else(|| content.get("content").and_then(content_text))
83 }
84 _ => None,
85 }
86}
87
88fn is_tool_use_type(kind: &str) -> bool {
89 matches!(kind, "tool_use" | "server_tool_use" | "mcp_tool_use")
90}
91
92fn is_tool_result_type(kind: &str) -> bool {
93 matches!(
94 kind,
95 "tool_result"
96 | "tool_search_tool_result"
97 | "web_fetch_tool_result"
98 | "web_search_tool_result"
99 | "code_execution_tool_result"
100 | "bash_code_execution_tool_result"
101 | "text_editor_code_execution_tool_result"
102 | "mcp_tool_result"
103 )
104}
105
106fn tool_title(block: &Value) -> String {
107 let name = block
108 .get("name")
109 .and_then(Value::as_str)
110 .filter(|name| !name.is_empty())
111 .unwrap_or("Tool call")
112 .replace('_', " ");
113 let description = block
114 .get("input")
115 .and_then(|input| input.get("description"))
116 .and_then(Value::as_str)
117 .filter(|description| !description.is_empty());
118 description.map_or(name.clone(), |description| format!("{name}: {description}"))
119}
120
121fn parse_tool_uses(slot: RosterSlot, value: &Value, state: &mut ParserState) -> Vec<AgentEvent> {
122 if value.get("type").and_then(Value::as_str) != Some("assistant") {
123 return Vec::new();
124 }
125 let Some(content) = value
126 .get("message")
127 .and_then(|message| message.get("content"))
128 .and_then(Value::as_array)
129 else {
130 return Vec::new();
131 };
132 content
133 .iter()
134 .filter(|block| {
135 block
136 .get("type")
137 .and_then(Value::as_str)
138 .is_some_and(is_tool_use_type)
139 })
140 .filter_map(|block| {
141 let id = block
142 .get("id")
143 .and_then(Value::as_str)
144 .filter(|id| !id.is_empty())?;
145 let title = tool_title(block);
146 if state.finished_tools.contains(id) {
147 return None;
148 }
149 let update = state
150 .tools
151 .entry(id.to_owned())
152 .or_insert_with(|| ToolUpdate {
153 id: id.to_owned(),
154 title: title.clone(),
155 status: ToolStatus::Running,
156 detail: None,
157 });
158 update.title = title;
159 Some(AgentEvent::Tool {
160 slot,
161 update: update.clone(),
162 })
163 })
164 .collect()
165}
166
167fn parse_stream_event(
168 slot: RosterSlot,
169 value: &Value,
170 state: &mut ParserState,
171) -> Option<AgentEvent> {
172 if value.get("type").and_then(Value::as_str) != Some("stream_event") {
173 return None;
174 }
175 let event = value.get("event")?;
176 let event_type = event.get("type").and_then(Value::as_str)?;
177 if event_type == "message_start" {
178 state.streamed_thoughts.clear();
179 return None;
180 }
181 let index = event.get("index").and_then(Value::as_u64);
182 match event_type {
183 "content_block_delta" => {
184 let delta = event.get("delta")?;
185 match delta.get("type").and_then(Value::as_str)? {
186 "text_delta" => delta
187 .get("text")
188 .and_then(Value::as_str)
189 .filter(|text| !text.is_empty())
190 .map(|text| AgentEvent::Text {
191 slot,
192 text: text.to_owned(),
193 }),
194 "thinking_delta" => {
195 let text = delta
196 .get("thinking")
197 .and_then(Value::as_str)
198 .filter(|text| !text.is_empty())?;
199 state
200 .streamed_thoughts
201 .entry(index?)
202 .or_default()
203 .push_str(text);
204 Some(AgentEvent::Thought {
205 slot,
206 text: text.to_owned(),
207 })
208 }
209 _ => None,
210 }
211 }
212 "content_block_start" => {
213 let index = index?;
214 let block = event.get("content_block")?;
215 if !block
216 .get("type")
217 .and_then(Value::as_str)
218 .is_some_and(is_tool_use_type)
219 {
220 return None;
221 }
222 let id = block
223 .get("id")
224 .and_then(Value::as_str)
225 .filter(|id| !id.is_empty())
226 .map_or_else(|| format!("claude-tool-{index}"), str::to_owned);
227 let title = tool_title(block);
228 state.finished_tools.remove(&id);
229 state.tools.insert(
230 id.clone(),
231 ToolUpdate {
232 id: id.clone(),
233 title: title.clone(),
234 status: ToolStatus::Running,
235 detail: None,
236 },
237 );
238 Some(AgentEvent::Tool {
239 slot,
240 update: ToolUpdate {
241 id,
242 title,
243 status: ToolStatus::Running,
244 detail: None,
245 },
246 })
247 }
248 "content_block_stop" => {
249 None
253 }
254 _ => None,
255 }
256}
257
258fn parse_consolidated_thoughts(
259 slot: RosterSlot,
260 value: &Value,
261 state: &mut ParserState,
262) -> Vec<AgentEvent> {
263 if value.get("type").and_then(Value::as_str) != Some("assistant") {
264 return Vec::new();
265 }
266 let Some(content) = value
267 .get("message")
268 .and_then(|message| message.get("content"))
269 .and_then(Value::as_array)
270 else {
271 return Vec::new();
272 };
273 let events = content
274 .iter()
275 .enumerate()
276 .filter_map(|(index, block)| {
277 if block.get("type").and_then(Value::as_str) != Some("thinking") {
278 return None;
279 }
280 let text = block
281 .get("thinking")
282 .and_then(Value::as_str)
283 .filter(|text| !text.is_empty())?;
284 let streamed = state
285 .streamed_thoughts
286 .get(&(index as u64))
287 .map(String::as_str)
288 .unwrap_or_default();
289 let remainder = text.strip_prefix(streamed).unwrap_or(text);
290 (!remainder.is_empty()).then(|| AgentEvent::Thought {
291 slot,
292 text: remainder.to_owned(),
293 })
294 })
295 .collect();
296 state.streamed_thoughts.clear();
297 events
298}
299
300fn parse_tool_results(slot: RosterSlot, value: &Value, state: &mut ParserState) -> Vec<AgentEvent> {
301 if value.get("type").and_then(Value::as_str) != Some("user") {
302 return Vec::new();
303 }
304 let Some(content) = value
305 .get("message")
306 .and_then(|message| message.get("content"))
307 .and_then(Value::as_array)
308 else {
309 return Vec::new();
310 };
311 content
312 .iter()
313 .filter(|block| {
314 block
315 .get("type")
316 .and_then(Value::as_str)
317 .is_some_and(is_tool_result_type)
318 })
319 .filter_map(|block| {
320 let id = block
321 .get("tool_use_id")
322 .and_then(Value::as_str)
323 .filter(|id| !id.is_empty())?;
324 let tool = state
325 .tools
326 .entry(id.to_owned())
327 .or_insert_with(|| ToolUpdate {
328 id: id.to_owned(),
329 title: "Tool call".into(),
330 status: ToolStatus::Running,
331 detail: None,
332 });
333 let result_type = block
334 .get("type")
335 .and_then(Value::as_str)
336 .unwrap_or_default();
337 let result = block.get("content");
338 let structured_result_type = result
339 .and_then(|result| result.get("type"))
340 .and_then(Value::as_str)
341 .unwrap_or_default();
342 let nonzero_exit = result.is_some_and(|result| {
343 ["return_code", "exit_code"]
344 .into_iter()
345 .filter_map(|key| result.get(key).and_then(Value::as_i64))
346 .any(|code| code != 0)
347 });
348 tool.status = if block
349 .get("is_error")
350 .and_then(Value::as_bool)
351 .unwrap_or(false)
352 || result_type.ends_with("_error")
353 || structured_result_type.ends_with("_error")
354 || nonzero_exit
355 {
356 ToolStatus::Failed
357 } else {
358 ToolStatus::Completed
359 };
360 tool.detail = result.and_then(content_text);
361 state.finished_tools.insert(id.to_owned());
362 Some(AgentEvent::Tool {
363 slot,
364 update: tool.clone(),
365 })
366 })
367 .collect()
368}
369
370fn parse_tool_progress(
371 slot: RosterSlot,
372 value: &Value,
373 state: &mut ParserState,
374) -> Option<AgentEvent> {
375 if value.get("type").and_then(Value::as_str) != Some("tool_progress") {
376 return None;
377 }
378 let reported = value.get("tool_use_id").and_then(Value::as_str);
379 let parent = value.get("parent_tool_use_id").and_then(Value::as_str);
380 let id = reported
381 .filter(|id| state.tools.contains_key(*id) && !state.finished_tools.contains(*id))
382 .or_else(|| {
383 parent.filter(|id| state.tools.contains_key(*id) && !state.finished_tools.contains(*id))
384 })?;
385 let update = state.tools.get_mut(id)?;
386 update.status = ToolStatus::Running;
387 if let Some(seconds) = value.get("elapsed_time_seconds").and_then(Value::as_u64) {
388 update.detail = Some(format!("running for {seconds}s"));
389 }
390 Some(AgentEvent::Tool {
391 slot,
392 update: update.clone(),
393 })
394}
395
396fn result_text(value: &Value) -> Option<String> {
397 value
398 .get("result")
399 .and_then(Value::as_str)
400 .filter(|text| !text.is_empty())
401 .map(str::to_owned)
402}
403
404fn session_id(value: &Value) -> Option<String> {
405 value
406 .get("session_id")
407 .or_else(|| value.get("sessionId"))
408 .and_then(Value::as_str)
409 .filter(|id| !id.is_empty())
410 .map(str::to_owned)
411}
412
413#[derive(Debug)]
414pub struct ClaudeAdapter {
415 slot: RosterSlot,
416 cwd: PathBuf,
417 command: String,
418 mode: String,
419 model: Option<String>,
420 session_id: Option<String>,
421 child: Option<Child>,
422 sender: mpsc::Sender<AdapterResult<AgentEvent>>,
423 receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
424 announced_session: Arc<Mutex<Option<String>>>,
425 cancel_requested: Arc<AtomicBool>,
426}
427
428impl ClaudeAdapter {
429 pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
430 let (sender, receiver) = mpsc::channel(256);
431 Self {
432 slot,
433 cwd,
434 command: command.into(),
435 mode: "bypassPermissions".into(),
436 model: None,
437 session_id: None,
438 child: None,
439 sender,
440 receiver,
441 announced_session: Arc::new(Mutex::new(None)),
442 cancel_requested: Arc::new(AtomicBool::new(false)),
443 }
444 }
445
446 pub fn with_session_id(
447 slot: RosterSlot,
448 cwd: PathBuf,
449 command: impl Into<String>,
450 session_id: impl Into<String>,
451 ) -> Self {
452 let mut adapter = Self::new(slot, cwd, command);
453 adapter.session_id = Some(session_id.into());
454 adapter
455 }
456
457 fn modes() -> Vec<Mode> {
458 [("plan", "Plan"), ("bypassPermissions", "Full Access")]
459 .into_iter()
460 .map(|(id, label)| Mode {
461 id: id.into(),
462 label: label.into(),
463 })
464 .collect()
465 }
466
467 fn models(&self) -> Vec<Mode> {
468 let ids = [
469 "best",
470 "opus",
471 "sonnet",
472 "haiku",
473 "sonnet[1m]",
474 "opus[1m]",
475 "opusplan",
476 ];
477 let mut models = vec![Mode {
478 id: "default".into(),
479 label: "Default".into(),
480 }];
481 for id in ids.map(str::to_owned) {
482 if !models.iter().any(|candidate| candidate.id == id) {
483 models.push(Mode {
484 label: model_label(&id),
485 id,
486 });
487 }
488 }
489 if let Some(model) = &self.model
490 && !models.iter().any(|candidate| candidate.id == *model)
491 {
492 models.push(Mode {
496 id: model.clone(),
497 label: model.clone(),
498 });
499 }
500 models
501 }
502
503 async fn emit(&self, event: AdapterResult<AgentEvent>) {
504 let _ = self.sender.send(event).await;
505 }
506}
507
508#[async_trait]
509impl AgentAdapter for ClaudeAdapter {
510 fn slot(&self) -> RosterSlot {
511 self.slot
512 }
513
514 fn display_name(&self) -> String {
515 "Claude".into()
516 }
517
518 fn session_id(&self) -> Option<String> {
519 self.session_id.clone()
520 }
521
522 fn protocol(&self) -> &'static str {
523 "native"
524 }
525
526 fn capabilities(&self) -> AgentCapabilities {
527 AgentCapabilities {
528 supports_cancel: true,
529 supports_modes: true,
530 supports_permissions: false,
531 supports_terminals: false,
532 supports_session_load: true,
533 supports_models: true,
534 }
535 }
536
537 async fn start(&mut self) -> AdapterResult<()> {
538 self.cancel_requested.store(false, Ordering::Release);
539 self.emit(Ok(AgentEvent::ModesReplaced {
540 slot: self.slot,
541 modes: Self::modes(),
542 current_mode: Some(self.mode.clone()),
543 }))
544 .await;
545 self.emit(Ok(AgentEvent::ModelsReplaced {
546 slot: self.slot,
547 config_id: MODEL_CONFIG_ID.into(),
548 models: self.models(),
549 current_model: self.model.clone(),
550 }))
551 .await;
552 self.emit(Ok(AgentEvent::Ready {
553 slot: self.slot,
554 capabilities: self.capabilities(),
555 }))
556 .await;
557 Ok(())
558 }
559
560 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
561 if self.child.is_some() {
562 return Err(AdapterError::Transport(
563 "agent is already handling a turn".into(),
564 ));
565 }
566 self.cancel_requested.store(false, Ordering::Release);
567 let (program, args) = parse_command_line(&self.command)
568 .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
569 let mut command = Command::new(program);
570 isolate_process_group(&mut command);
571 command
572 .args(args)
573 .arg("--print")
574 .arg("--output-format")
575 .arg("stream-json")
576 .arg("--verbose")
577 .arg("--include-partial-messages")
578 .arg("--permission-prompts")
579 .arg("none")
580 .arg("--permission-mode")
581 .arg(&self.mode)
582 .current_dir(&self.cwd)
583 .env("CODESWARM_CWD", &self.cwd);
584 if self.mode == "bypassPermissions" {
585 command.arg("--allow-dangerously-skip-permissions");
586 }
587 if let Some(model) = &self.model {
588 command.arg("--model").arg(model);
589 }
590 if let Some(session_id) = &self.session_id {
591 command.arg("--resume").arg(session_id);
592 }
593 let NativeTurn {
594 child,
595 stdout,
596 stderr,
597 } = spawn_native_turn(command, prompt).await?;
598 let sender = self.sender.clone();
599 let slot = self.slot;
600 let announced = Arc::clone(&self.announced_session);
601 let cancelled = Arc::clone(&self.cancel_requested);
602 tokio::spawn(async move {
603 let stderr_task = tokio::spawn(drain_bounded(stderr, 32 * 1024));
604 let mut lines = BufReader::new(stdout).lines();
605 let mut state = ParserState::default();
606 let mut result = None;
607 let mut streamed = false;
608 while let Ok(Some(line)) = lines.next_line().await {
609 let Ok(value) = serde_json::from_str::<Value>(&line) else {
610 continue;
611 };
612 if let Some(id) = session_id(&value)
613 && let Ok(mut current) = announced.lock()
614 {
615 *current = Some(id);
616 }
617 if value.get("type").and_then(Value::as_str) == Some("result") {
618 result = Some(value.clone());
619 }
620 if let Some(event) = parse_stream_event(slot, &value, &mut state) {
621 streamed |= matches!(event, AgentEvent::Text { .. });
622 if sender.send(Ok(event)).await.is_err() {
623 break;
624 }
625 }
626 for event in parse_consolidated_thoughts(slot, &value, &mut state) {
627 if sender.send(Ok(event)).await.is_err() {
628 break;
629 }
630 }
631 for event in parse_tool_uses(slot, &value, &mut state) {
632 if sender.send(Ok(event)).await.is_err() {
633 break;
634 }
635 }
636 if let Some(event) = parse_tool_progress(slot, &value, &mut state)
637 && sender.send(Ok(event)).await.is_err()
638 {
639 break;
640 }
641 for event in parse_tool_results(slot, &value, &mut state) {
642 if sender.send(Ok(event)).await.is_err() {
643 break;
644 }
645 }
646 }
647 let stderr = stderr_task.await.ok().unwrap_or_default();
648 let succeeded = cancelled.load(Ordering::Acquire)
649 || result.as_ref().is_some_and(|value| {
650 value.get("subtype").and_then(Value::as_str) == Some("success")
651 && !value
652 .get("is_error")
653 .and_then(Value::as_bool)
654 .unwrap_or(false)
655 });
656 if succeeded {
657 if !streamed && let Some(text) = result.as_ref().and_then(result_text) {
658 let _ = sender.send(Ok(AgentEvent::Text { slot, text })).await;
659 }
660 let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
661 } else {
662 let detail = result
663 .as_ref()
664 .and_then(result_text)
665 .or_else(|| {
666 result
667 .as_ref()
668 .and_then(|value| value.get("subtype"))
669 .and_then(Value::as_str)
670 .map(str::to_owned)
671 })
672 .or_else(|| (!stderr.is_empty()).then_some(stderr))
673 .unwrap_or_else(|| "Claude stream ended before a successful result".into());
674 let _ = sender
675 .send(Ok(AgentEvent::Failed {
676 slot,
677 started: true,
678 detail,
679 }))
680 .await;
681 }
682 });
683 self.child = Some(child);
684 Ok(())
685 }
686
687 async fn cancel(&mut self) -> AdapterResult<bool> {
688 self.cancel_requested.store(true, Ordering::Release);
689 let Some(mut child) = self.child.take() else {
690 return Ok(false);
691 };
692 terminate_child(&mut child).await?;
693 let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
694 while let Some(event) = self.receiver.recv().await {
695 if matches!(
696 event,
697 Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
698 ) {
699 break;
700 }
701 }
702 })
703 .await;
704 Ok(true)
705 }
706
707 async fn answer_permission(
708 &mut self,
709 _request_id: String,
710 _answer: PermissionAnswer,
711 ) -> AdapterResult<()> {
712 Err(AdapterError::Unsupported("permission answer"))
713 }
714
715 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
716 self.mode = match mode.as_str() {
717 "codeswarm:mode:full-access"
718 | "full-access"
719 | "auto"
720 | "autopilot"
721 | "bypassPermissions" => "bypassPermissions",
722 "codeswarm:mode:plan" | "readonly" | "plan" => "plan",
723 _ => return Err(AdapterError::Unsupported("requested Claude mode")),
724 }
725 .into();
726 self.emit(Ok(AgentEvent::ModesReplaced {
727 slot: self.slot,
728 modes: Self::modes(),
729 current_mode: Some(self.mode.clone()),
730 }))
731 .await;
732 Ok(())
733 }
734
735 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
736 let model = model.trim();
737 let listed = self.models().iter().any(|candidate| candidate.id == model);
738 let full_model_id = model.starts_with("claude-")
739 && !model.contains(char::is_whitespace)
740 && !model.contains('\0');
741 if !listed && !full_model_id {
742 return Err(AdapterError::Protocol(
743 "model must be a listed Claude alias or a full claude-* model ID".into(),
744 ));
745 }
746 self.model = Some(model.to_owned());
747 self.emit(Ok(AgentEvent::ModelsReplaced {
748 slot: self.slot,
749 config_id: MODEL_CONFIG_ID.into(),
750 models: self.models(),
751 current_model: self.model.clone(),
752 }))
753 .await;
754 Ok(())
755 }
756
757 async fn reload(&mut self) -> AdapterResult<()> {
758 self.stop().await?;
759 self.start().await
760 }
761
762 async fn stop(&mut self) -> AdapterResult<()> {
763 let _ = self.cancel().await?;
764 Ok(())
765 }
766
767 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
768 let event = self.receiver.recv().await;
769 if matches!(
770 event.as_ref(),
771 Some(Ok(
772 AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. }
773 ))
774 ) {
775 if self.session_id.is_none()
776 && let Ok(session) = self.announced_session.lock()
777 {
778 self.session_id = session.clone();
779 }
780 if let Some(mut child) = self.child.take() {
781 let _ = child.wait().await;
782 }
783 }
784 event
785 }
786}
787
788#[cfg(test)]
789mod tests {
790 use super::{
791 ClaudeAdapter, ParserState, parse_consolidated_thoughts, parse_stream_event,
792 parse_tool_progress, parse_tool_results, parse_tool_uses,
793 };
794 use crate::{AgentAdapter, AgentEvent, ToolStatus};
795 use serde_json::json;
796
797 #[test]
798 fn parses_claude_text_thought_and_tool_events() {
799 let mut state = ParserState::default();
800 assert!(matches!(
801 parse_stream_event(2, &json!({"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"text_delta","text":"hello"}}}), &mut state),
802 Some(AgentEvent::Text { slot: 2, text }) if text == "hello"
803 ));
804 assert!(matches!(
805 parse_stream_event(2, &json!({"type":"stream_event","event":{"type":"content_block_delta","index":0,"delta":{"type":"thinking_delta","thinking":"check"}}}), &mut state),
806 Some(AgentEvent::Thought { text, .. }) if text == "check"
807 ));
808 assert!(matches!(
809 parse_stream_event(2, &json!({"type":"stream_event","event":{"type":"content_block_start","index":1,"content_block":{"type":"tool_use","id":"tool-1","name":"Read"}}}), &mut state),
810 Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Running && update.title == "Read"
811 ));
812 assert_eq!(
813 parse_stream_event(
814 2,
815 &json!({"type":"stream_event","event":{"type":"content_block_stop","index":1}}),
816 &mut state
817 ),
818 None
819 );
820
821 let completed = parse_tool_results(
822 2,
823 &json!({
824 "type": "user",
825 "message": {
826 "content": [{
827 "type": "tool_result",
828 "tool_use_id": "tool-1",
829 "content": [{"type": "text", "text": "file contents"}]
830 }]
831 }
832 }),
833 &mut state,
834 );
835 assert!(matches!(
836 completed.as_slice(),
837 [AgentEvent::Tool { update, .. }]
838 if update.status == ToolStatus::Completed
839 && update.title == "Read"
840 && update.detail.as_deref() == Some("file contents")
841 ));
842 }
843
844 #[test]
845 fn failed_claude_tool_result_preserves_error_detail() {
846 let mut state = ParserState::default();
847 let failed = parse_tool_results(
848 4,
849 &json!({
850 "type": "user",
851 "message": {
852 "content": [{
853 "type": "tool_result",
854 "tool_use_id": "tool-without-partial-start",
855 "is_error": true,
856 "content": "permission denied"
857 }]
858 }
859 }),
860 &mut state,
861 );
862 assert!(matches!(
863 failed.as_slice(),
864 [AgentEvent::Tool { slot: 4, update }]
865 if update.status == ToolStatus::Failed
866 && update.id == "tool-without-partial-start"
867 && update.detail.as_deref() == Some("permission denied")
868 ));
869 }
870
871 #[test]
872 fn consolidated_tool_use_and_sdk_result_variants_match_claude_acp() {
873 let mut state = ParserState::default();
874 let refined = parse_tool_uses(
875 1,
876 &json!({
877 "type": "assistant",
878 "message": {"content": [{
879 "type": "server_tool_use",
880 "id": "server-tool",
881 "name": "bash_code_execution",
882 "input": {"description": "Compile the project"}
883 }]}
884 }),
885 &mut state,
886 );
887 assert!(matches!(
888 refined.as_slice(),
889 [AgentEvent::Tool { update, .. }]
890 if update.status == ToolStatus::Running
891 && update.title == "bash code execution: Compile the project"
892 ));
893
894 let completed = parse_tool_results(
895 1,
896 &json!({
897 "type": "user",
898 "message": {"content": [{
899 "type": "bash_code_execution_tool_result",
900 "tool_use_id": "server-tool",
901 "content": {
902 "type": "bash_code_execution_result",
903 "stdout": "partial output",
904 "stderr": "compiler error",
905 "return_code": 2
906 }
907 }]}
908 }),
909 &mut state,
910 );
911 assert!(matches!(
912 completed.as_slice(),
913 [AgentEvent::Tool { update, .. }]
914 if update.status == ToolStatus::Failed
915 && update.title == "bash code execution: Compile the project"
916 && update.detail.as_deref() == Some("partial output\ncompiler error")
917 ));
918 assert!(
919 parse_tool_progress(
920 1,
921 &json!({
922 "type": "tool_progress",
923 "tool_use_id": "server-tool-heartbeat-1",
924 "parent_tool_use_id": "server-tool",
925 "tool_name": "bash_code_execution",
926 "elapsed_time_seconds": 30
927 }),
928 &mut state,
929 )
930 .is_none(),
931 "a late heartbeat must not reopen a completed tool"
932 );
933 }
934
935 #[test]
936 fn tool_progress_uses_the_real_parent_id_without_creating_phantom_tools() {
937 let mut state = ParserState::default();
938 let _ = parse_stream_event(
939 3,
940 &json!({
941 "type": "stream_event",
942 "event": {
943 "type": "content_block_start",
944 "index": 0,
945 "content_block": {"type": "tool_use", "id": "tool-3", "name": "Bash"}
946 }
947 }),
948 &mut state,
949 );
950 assert!(matches!(
951 parse_tool_progress(
952 3,
953 &json!({
954 "type": "tool_progress",
955 "tool_use_id": "tool-3-heartbeat-1",
956 "parent_tool_use_id": "tool-3",
957 "tool_name": "Bash",
958 "elapsed_time_seconds": 30
959 }),
960 &mut state,
961 ),
962 Some(AgentEvent::Tool { update, .. })
963 if update.id == "tool-3"
964 && update.status == ToolStatus::Running
965 && update.detail.as_deref() == Some("running for 30s")
966 ));
967 assert_eq!(state.tools.len(), 1);
968 }
969
970 #[test]
971 fn structured_sdk_error_results_are_failed_with_their_details() {
972 for (outer, inner) in [
973 ("tool_search_tool_result", "tool_search_tool_result_error"),
974 ("web_fetch_tool_result", "web_fetch_tool_result_error"),
975 ("web_search_tool_result", "web_search_tool_result_error"),
976 (
977 "code_execution_tool_result",
978 "code_execution_tool_result_error",
979 ),
980 (
981 "bash_code_execution_tool_result",
982 "bash_code_execution_tool_result_error",
983 ),
984 (
985 "text_editor_code_execution_tool_result",
986 "text_editor_code_execution_tool_result_error",
987 ),
988 ] {
989 let mut state = ParserState::default();
990 let events = parse_tool_results(
991 5,
992 &json!({
993 "type": "user",
994 "message": {"content": [{
995 "type": outer,
996 "tool_use_id": "failed-tool",
997 "content": {
998 "type": inner,
999 "error_code": "execution_failed",
1000 "error_message": "provider rejected the tool"
1001 }
1002 }]}
1003 }),
1004 &mut state,
1005 );
1006 assert!(matches!(
1007 events.as_slice(),
1008 [AgentEvent::Tool { update, .. }]
1009 if update.status == ToolStatus::Failed
1010 && update.detail.as_deref()
1011 == Some("execution_failed\nprovider rejected the tool")
1012 ));
1013 }
1014 }
1015
1016 #[test]
1017 fn consolidated_thoughts_fill_missing_deltas_without_duplication() {
1018 let mut state = ParserState::default();
1019 assert!(
1020 parse_stream_event(
1021 0,
1022 &json!({"type":"stream_event","event":{"type":"message_start","message":{}}}),
1023 &mut state,
1024 )
1025 .is_none()
1026 );
1027 assert!(matches!(
1028 parse_stream_event(
1029 0,
1030 &json!({"type":"stream_event","event":{"type":"content_block_delta","index":0,"delta":{"type":"thinking_delta","thinking":"first "}}}),
1031 &mut state,
1032 ),
1033 Some(AgentEvent::Thought { text, .. }) if text == "first "
1034 ));
1035 let remainder = parse_consolidated_thoughts(
1036 0,
1037 &json!({
1038 "type": "assistant",
1039 "message": {"content": [{"type": "thinking", "thinking": "first second"}]}
1040 }),
1041 &mut state,
1042 );
1043 assert!(matches!(
1044 remainder.as_slice(),
1045 [AgentEvent::Thought { text, .. }] if text == "second"
1046 ));
1047
1048 let fallback = parse_consolidated_thoughts(
1049 0,
1050 &json!({
1051 "type": "assistant",
1052 "message": {"content": [{"type": "thinking", "thinking": "gateway-only thought"}]}
1053 }),
1054 &mut state,
1055 );
1056 assert!(matches!(
1057 fallback.as_slice(),
1058 [AgentEvent::Thought { text, .. }] if text == "gateway-only thought"
1059 ));
1060 }
1061
1062 #[tokio::test]
1063 async fn advertises_current_aliases_and_accepts_full_model_ids() {
1064 let mut adapter = ClaudeAdapter::new(0, std::env::current_dir().unwrap(), "claude");
1065 assert_eq!(
1066 ClaudeAdapter::modes()
1067 .into_iter()
1068 .map(|mode| mode.id)
1069 .collect::<Vec<_>>(),
1070 ["plan", "bypassPermissions"]
1071 );
1072 assert_eq!(
1073 adapter
1074 .models()
1075 .into_iter()
1076 .map(|model| model.id)
1077 .collect::<Vec<_>>(),
1078 [
1079 "default",
1080 "best",
1081 "opus",
1082 "sonnet",
1083 "haiku",
1084 "sonnet[1m]",
1085 "opus[1m]",
1086 "opusplan"
1087 ]
1088 );
1089 assert!(adapter.set_mode("manual".into()).await.is_err());
1090 adapter
1091 .set_model("claude-sonnet-4-5-20250929".into())
1092 .await
1093 .unwrap();
1094 assert!(
1095 adapter
1096 .models()
1097 .iter()
1098 .any(|model| model.id == "claude-sonnet-4-5-20250929")
1099 );
1100 assert!(
1101 adapter
1102 .set_model("provider/model:latest".into())
1103 .await
1104 .is_err()
1105 );
1106 assert!(adapter.set_model(" ".into()).await.is_err());
1107 assert!(adapter.set_model("bad\0model".into()).await.is_err());
1108 }
1109
1110 #[tokio::test]
1111 async fn native_claude_process_captures_session_and_resumes() {
1112 let args_path =
1113 std::env::temp_dir().join(format!("codeswarm-claude-args-{}", std::process::id()));
1114 let stdin_path =
1115 std::env::temp_dir().join(format!("codeswarm-claude-stdin-{}", std::process::id()));
1116 let script_path =
1117 std::env::temp_dir().join(format!("codeswarm-claude-script-{}", std::process::id()));
1118 std::fs::write(
1119 &script_path,
1120 format!(
1121 "printf '%s\\n' \"$*\" >> '{}'\nsed -n 'p' >> '{}'\nprintf '\\n' >> '{}'\nprintf '%s\\n' '{{\"type\":\"system\",\"subtype\":\"init\",\"model\":\"provider-runtime-model\",\"session_id\":\"session-native\"}}' '{{\"type\":\"assistant\",\"message\":{{\"content\":[{{\"type\":\"thinking\",\"thinking\":\"checked context\"}}]}}}}' '{{\"type\":\"result\",\"subtype\":\"success\",\"result\":\"hello\",\"session_id\":\"session-native\"}}'\n",
1122 args_path.display(),
1123 stdin_path.display(),
1124 stdin_path.display()
1125 ),
1126 )
1127 .unwrap();
1128 let mut adapter = ClaudeAdapter::new(
1129 0,
1130 std::env::current_dir().unwrap(),
1131 format!("sh {}", script_path.display()),
1132 );
1133 adapter.start().await.unwrap();
1134 assert!(matches!(
1135 adapter.next_event().await,
1136 Some(Ok(AgentEvent::ModesReplaced { .. }))
1137 ));
1138 assert!(matches!(
1139 adapter.next_event().await,
1140 Some(Ok(AgentEvent::ModelsReplaced {
1141 current_model: None,
1142 ..
1143 }))
1144 ));
1145 assert!(matches!(
1146 adapter.next_event().await,
1147 Some(Ok(AgentEvent::Ready { .. }))
1148 ));
1149 adapter.send_prompt("-first prompt".into()).await.unwrap();
1150 assert!(matches!(
1151 adapter.next_event().await,
1152 Some(Ok(AgentEvent::Thought { text, .. })) if text == "checked context"
1153 ));
1154 assert!(
1155 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "hello")
1156 );
1157 assert!(matches!(
1158 adapter.next_event().await,
1159 Some(Ok(AgentEvent::TurnComplete { .. }))
1160 ));
1161 assert_eq!(adapter.session_id(), Some("session-native".into()));
1162 adapter.set_model("default".into()).await.unwrap();
1163 assert!(matches!(
1164 adapter.next_event().await,
1165 Some(Ok(AgentEvent::ModelsReplaced { current_model, .. }))
1166 if current_model.as_deref() == Some("default")
1167 ));
1168 adapter.send_prompt("second prompt".into()).await.unwrap();
1169 assert!(matches!(
1170 adapter.next_event().await,
1171 Some(Ok(AgentEvent::Thought { .. }))
1172 ));
1173 assert!(matches!(
1174 adapter.next_event().await,
1175 Some(Ok(AgentEvent::Text { .. }))
1176 ));
1177 assert!(matches!(
1178 adapter.next_event().await,
1179 Some(Ok(AgentEvent::TurnComplete { .. }))
1180 ));
1181 let args = std::fs::read_to_string(&args_path).unwrap();
1182 let mut argument_lines = args.lines();
1183 let first_args = argument_lines.next().unwrap_or_default();
1184 let second_args = argument_lines.next().unwrap_or_default();
1185 assert!(!first_args.contains("--model"), "{args}");
1186 assert!(second_args.contains("--model default"), "{args}");
1187 assert!(second_args.contains("--resume session-native"), "{args}");
1188 assert!(!args.contains("first prompt"), "{args}");
1189 assert!(!args.contains("second prompt"), "{args}");
1190 assert_eq!(
1191 std::fs::read_to_string(&stdin_path).unwrap(),
1192 "-first prompt\nsecond prompt\n"
1193 );
1194 adapter.stop().await.unwrap();
1195 let _ = std::fs::remove_file(args_path);
1196 let _ = std::fs::remove_file(stdin_path);
1197 let _ = std::fs::remove_file(script_path);
1198 }
1199
1200 #[tokio::test]
1201 async fn native_claude_reload_reaps_a_silent_turn() {
1202 let mut adapter =
1203 ClaudeAdapter::new(0, std::env::current_dir().unwrap(), "sh -c 'sleep 10'");
1204 adapter.start().await.unwrap();
1205 for _ in 0..3 {
1206 assert!(adapter.next_event().await.is_some());
1207 }
1208 adapter.send_prompt("stuck".into()).await.unwrap();
1209 tokio::time::timeout(std::time::Duration::from_secs(5), adapter.reload())
1210 .await
1211 .expect("reload should not hang")
1212 .expect("reload should succeed");
1213 assert!(adapter.child.is_none());
1214 adapter.stop().await.unwrap();
1215 }
1216}