Skip to main content

codeswarm_adapters/
claude.rs

1//! Native adapter for Claude Code's print-mode JSONL stream.
2
3use 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            // This only closes Claude's streamed `tool_use` input block. The
250            // tool is still executing; its later `tool_result` user message
251            // carries the actual completion status and output.
252            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            // Claude Code has no model-discovery command. Keep its documented
493            // aliases static, but retain an explicitly configured full model
494            // ID instead of pretending the alias list is exhaustive.
495            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}