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,
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
30#[derive(Debug, Default)]
31struct ParserState {
32    tools: BTreeMap<String, ToolUpdate>,
33}
34
35fn content_text(value: &Value) -> Option<String> {
36    match value {
37        Value::String(text) => (!text.is_empty()).then(|| text.to_owned()),
38        Value::Array(content) => {
39            let text = content
40                .iter()
41                .filter_map(content_text)
42                .collect::<Vec<_>>()
43                .join("\n");
44            (!text.is_empty()).then_some(text)
45        }
46        Value::Object(content) => content
47            .get("text")
48            .and_then(Value::as_str)
49            .filter(|text| !text.is_empty())
50            .map(str::to_owned)
51            .or_else(|| content.get("content").and_then(content_text)),
52        _ => None,
53    }
54}
55
56fn parse_stream_event(
57    slot: RosterSlot,
58    value: &Value,
59    state: &mut ParserState,
60) -> Option<AgentEvent> {
61    if value.get("type").and_then(Value::as_str) != Some("stream_event") {
62        return None;
63    }
64    let event = value.get("event")?;
65    let index = event.get("index").and_then(Value::as_u64);
66    match event.get("type").and_then(Value::as_str)? {
67        "content_block_delta" => {
68            let delta = event.get("delta")?;
69            match delta.get("type").and_then(Value::as_str)? {
70                "text_delta" => delta
71                    .get("text")
72                    .and_then(Value::as_str)
73                    .filter(|text| !text.is_empty())
74                    .map(|text| AgentEvent::Text {
75                        slot,
76                        text: text.to_owned(),
77                    }),
78                "thinking_delta" => delta
79                    .get("thinking")
80                    .and_then(Value::as_str)
81                    .filter(|text| !text.is_empty())
82                    .map(|text| AgentEvent::Thought {
83                        slot,
84                        text: text.to_owned(),
85                    }),
86                _ => None,
87            }
88        }
89        "content_block_start" => {
90            let index = index?;
91            let block = event.get("content_block")?;
92            if block.get("type").and_then(Value::as_str) != Some("tool_use") {
93                return None;
94            }
95            let id = block
96                .get("id")
97                .and_then(Value::as_str)
98                .filter(|id| !id.is_empty())
99                .map_or_else(|| format!("claude-tool-{index}"), str::to_owned);
100            let title = block
101                .get("name")
102                .and_then(Value::as_str)
103                .unwrap_or("Tool call")
104                .replace('_', " ");
105            state.tools.insert(
106                id.clone(),
107                ToolUpdate {
108                    id: id.clone(),
109                    title: title.clone(),
110                    status: ToolStatus::Running,
111                    detail: None,
112                },
113            );
114            Some(AgentEvent::Tool {
115                slot,
116                update: ToolUpdate {
117                    id,
118                    title,
119                    status: ToolStatus::Running,
120                    detail: None,
121                },
122            })
123        }
124        "content_block_stop" => {
125            // This only closes Claude's streamed `tool_use` input block. The
126            // tool is still executing; its later `tool_result` user message
127            // carries the actual completion status and output.
128            None
129        }
130        _ => None,
131    }
132}
133
134fn parse_tool_results(slot: RosterSlot, value: &Value, state: &mut ParserState) -> Vec<AgentEvent> {
135    if value.get("type").and_then(Value::as_str) != Some("user") {
136        return Vec::new();
137    }
138    let Some(content) = value
139        .get("message")
140        .and_then(|message| message.get("content"))
141        .and_then(Value::as_array)
142    else {
143        return Vec::new();
144    };
145    content
146        .iter()
147        .filter(|block| block.get("type").and_then(Value::as_str) == Some("tool_result"))
148        .filter_map(|block| {
149            let id = block
150                .get("tool_use_id")
151                .and_then(Value::as_str)
152                .filter(|id| !id.is_empty())?;
153            let tool = state
154                .tools
155                .entry(id.to_owned())
156                .or_insert_with(|| ToolUpdate {
157                    id: id.to_owned(),
158                    title: "Tool call".into(),
159                    status: ToolStatus::Running,
160                    detail: None,
161                });
162            tool.status = if block
163                .get("is_error")
164                .and_then(Value::as_bool)
165                .unwrap_or(false)
166            {
167                ToolStatus::Failed
168            } else {
169                ToolStatus::Completed
170            };
171            tool.detail = block.get("content").and_then(content_text);
172            Some(AgentEvent::Tool {
173                slot,
174                update: tool.clone(),
175            })
176        })
177        .collect()
178}
179
180fn result_text(value: &Value) -> Option<String> {
181    value
182        .get("result")
183        .and_then(Value::as_str)
184        .filter(|text| !text.is_empty())
185        .map(str::to_owned)
186}
187
188fn session_id(value: &Value) -> Option<String> {
189    value
190        .get("session_id")
191        .or_else(|| value.get("sessionId"))
192        .and_then(Value::as_str)
193        .filter(|id| !id.is_empty())
194        .map(str::to_owned)
195}
196
197#[derive(Debug)]
198pub struct ClaudeAdapter {
199    slot: RosterSlot,
200    cwd: PathBuf,
201    command: String,
202    mode: String,
203    model: Option<String>,
204    session_id: Option<String>,
205    child: Option<Child>,
206    sender: mpsc::Sender<AdapterResult<AgentEvent>>,
207    receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
208    announced_session: Arc<Mutex<Option<String>>>,
209    cancel_requested: Arc<AtomicBool>,
210}
211
212impl ClaudeAdapter {
213    pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
214        let (sender, receiver) = mpsc::channel(256);
215        Self {
216            slot,
217            cwd,
218            command: command.into(),
219            mode: "bypassPermissions".into(),
220            model: None,
221            session_id: None,
222            child: None,
223            sender,
224            receiver,
225            announced_session: Arc::new(Mutex::new(None)),
226            cancel_requested: Arc::new(AtomicBool::new(false)),
227        }
228    }
229
230    pub fn with_session_id(
231        slot: RosterSlot,
232        cwd: PathBuf,
233        command: impl Into<String>,
234        session_id: impl Into<String>,
235    ) -> Self {
236        let mut adapter = Self::new(slot, cwd, command);
237        adapter.session_id = Some(session_id.into());
238        adapter
239    }
240
241    fn modes() -> Vec<Mode> {
242        [("plan", "Plan"), ("bypassPermissions", "Full Access")]
243            .into_iter()
244            .map(|(id, label)| Mode {
245                id: id.into(),
246                label: label.into(),
247            })
248            .collect()
249    }
250
251    fn models(&self) -> Vec<Mode> {
252        let mut models = [
253            ("fable", "Fable"),
254            ("opus", "Opus"),
255            ("sonnet", "Sonnet"),
256            ("haiku", "Haiku"),
257        ]
258        .into_iter()
259        .map(|(id, label)| Mode {
260            id: id.into(),
261            label: label.into(),
262        })
263        .collect::<Vec<_>>();
264        if let Some(model) = &self.model
265            && !models.iter().any(|candidate| candidate.id == *model)
266        {
267            // Claude Code has no model-discovery command. Keep its documented
268            // aliases static, but retain an explicitly configured full model
269            // ID instead of pretending the alias list is exhaustive.
270            models.push(Mode {
271                id: model.clone(),
272                label: model.clone(),
273            });
274        }
275        models
276    }
277
278    async fn emit(&self, event: AdapterResult<AgentEvent>) {
279        let _ = self.sender.send(event).await;
280    }
281}
282
283#[async_trait]
284impl AgentAdapter for ClaudeAdapter {
285    fn slot(&self) -> RosterSlot {
286        self.slot
287    }
288
289    fn display_name(&self) -> String {
290        "Claude".into()
291    }
292
293    fn session_id(&self) -> Option<String> {
294        self.session_id.clone()
295    }
296
297    fn protocol(&self) -> &'static str {
298        "native"
299    }
300
301    fn capabilities(&self) -> AgentCapabilities {
302        AgentCapabilities {
303            supports_cancel: true,
304            supports_modes: true,
305            supports_permissions: false,
306            supports_terminals: false,
307            supports_session_load: true,
308            supports_models: true,
309        }
310    }
311
312    async fn start(&mut self) -> AdapterResult<()> {
313        self.cancel_requested.store(false, Ordering::Release);
314        self.emit(Ok(AgentEvent::ModesReplaced {
315            slot: self.slot,
316            modes: Self::modes(),
317            current_mode: Some(self.mode.clone()),
318        }))
319        .await;
320        self.emit(Ok(AgentEvent::ModelsReplaced {
321            slot: self.slot,
322            config_id: "claude:model".into(),
323            models: self.models(),
324            current_model: self.model.clone(),
325        }))
326        .await;
327        self.emit(Ok(AgentEvent::Ready {
328            slot: self.slot,
329            capabilities: self.capabilities(),
330        }))
331        .await;
332        Ok(())
333    }
334
335    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
336        if self.child.is_some() {
337            return Err(AdapterError::Transport(
338                "agent is already handling a turn".into(),
339            ));
340        }
341        self.cancel_requested.store(false, Ordering::Release);
342        let (program, args) = parse_command_line(&self.command)
343            .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
344        let mut command = Command::new(program);
345        isolate_process_group(&mut command);
346        command
347            .args(args)
348            .arg("--print")
349            .arg("--output-format")
350            .arg("stream-json")
351            .arg("--verbose")
352            .arg("--include-partial-messages")
353            .arg("--permission-prompts")
354            .arg("none")
355            .arg("--permission-mode")
356            .arg(&self.mode)
357            .current_dir(&self.cwd)
358            .env("CODESWARM_CWD", &self.cwd);
359        if self.mode == "bypassPermissions" {
360            command.arg("--allow-dangerously-skip-permissions");
361        }
362        if let Some(model) = &self.model {
363            command.arg("--model").arg(model);
364        }
365        if let Some(session_id) = &self.session_id {
366            command.arg("--resume").arg(session_id);
367        }
368        let NativeTurn {
369            child,
370            stdout,
371            stderr,
372        } = spawn_native_turn(command, prompt).await?;
373        let sender = self.sender.clone();
374        let slot = self.slot;
375        let announced = Arc::clone(&self.announced_session);
376        let cancelled = Arc::clone(&self.cancel_requested);
377        tokio::spawn(async move {
378            let stderr_task = tokio::spawn(drain_bounded(stderr, 32 * 1024));
379            let mut lines = BufReader::new(stdout).lines();
380            let mut state = ParserState::default();
381            let mut result = None;
382            let mut streamed = false;
383            while let Ok(Some(line)) = lines.next_line().await {
384                let Ok(value) = serde_json::from_str::<Value>(&line) else {
385                    continue;
386                };
387                if let Some(id) = session_id(&value)
388                    && let Ok(mut current) = announced.lock()
389                {
390                    *current = Some(id);
391                }
392                if value.get("type").and_then(Value::as_str) == Some("result") {
393                    result = Some(value.clone());
394                }
395                if let Some(event) = parse_stream_event(slot, &value, &mut state) {
396                    streamed |= matches!(event, AgentEvent::Text { .. });
397                    if sender.send(Ok(event)).await.is_err() {
398                        break;
399                    }
400                }
401                for event in parse_tool_results(slot, &value, &mut state) {
402                    if sender.send(Ok(event)).await.is_err() {
403                        break;
404                    }
405                }
406            }
407            let stderr = stderr_task.await.ok().unwrap_or_default();
408            let succeeded = cancelled.load(Ordering::Acquire)
409                || result.as_ref().is_some_and(|value| {
410                    value.get("subtype").and_then(Value::as_str) == Some("success")
411                        && !value
412                            .get("is_error")
413                            .and_then(Value::as_bool)
414                            .unwrap_or(false)
415                });
416            if succeeded {
417                if !streamed && let Some(text) = result.as_ref().and_then(result_text) {
418                    let _ = sender.send(Ok(AgentEvent::Text { slot, text })).await;
419                }
420                let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
421            } else {
422                let detail = result
423                    .as_ref()
424                    .and_then(result_text)
425                    .or_else(|| {
426                        result
427                            .as_ref()
428                            .and_then(|value| value.get("subtype"))
429                            .and_then(Value::as_str)
430                            .map(str::to_owned)
431                    })
432                    .or_else(|| (!stderr.is_empty()).then_some(stderr))
433                    .unwrap_or_else(|| "Claude stream ended before a successful result".into());
434                let _ = sender
435                    .send(Ok(AgentEvent::Failed {
436                        slot,
437                        started: true,
438                        detail,
439                    }))
440                    .await;
441            }
442        });
443        self.child = Some(child);
444        Ok(())
445    }
446
447    async fn cancel(&mut self) -> AdapterResult<bool> {
448        self.cancel_requested.store(true, Ordering::Release);
449        let Some(mut child) = self.child.take() else {
450            return Ok(false);
451        };
452        terminate_child(&mut child).await?;
453        let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
454            while let Some(event) = self.receiver.recv().await {
455                if matches!(
456                    event,
457                    Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
458                ) {
459                    break;
460                }
461            }
462        })
463        .await;
464        Ok(true)
465    }
466
467    async fn answer_permission(
468        &mut self,
469        _request_id: String,
470        _answer: PermissionAnswer,
471    ) -> AdapterResult<()> {
472        Err(AdapterError::Unsupported("permission answer"))
473    }
474
475    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
476        self.mode = match mode.as_str() {
477            "codeswarm:mode:full-access"
478            | "full-access"
479            | "auto"
480            | "autopilot"
481            | "bypassPermissions" => "bypassPermissions",
482            "codeswarm:mode:plan" | "readonly" | "plan" => "plan",
483            _ => return Err(AdapterError::Unsupported("requested Claude mode")),
484        }
485        .into();
486        self.emit(Ok(AgentEvent::ModesReplaced {
487            slot: self.slot,
488            modes: Self::modes(),
489            current_mode: Some(self.mode.clone()),
490        }))
491        .await;
492        Ok(())
493    }
494
495    async fn set_model(&mut self, model: String) -> AdapterResult<()> {
496        let documented_alias = self.models().iter().any(|candidate| candidate.id == model);
497        let full_model_id = model.strip_prefix("claude-").is_some_and(|suffix| {
498            !suffix.is_empty()
499                && suffix
500                    .bytes()
501                    .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-')
502        });
503        if !documented_alias && !full_model_id {
504            return Err(AdapterError::Protocol(
505                "model must be a documented Claude Code alias or a full claude-* model ID".into(),
506            ));
507        }
508        self.model = Some(model.clone());
509        self.emit(Ok(AgentEvent::ModelUpdated {
510            slot: self.slot,
511            current_model: model,
512        }))
513        .await;
514        Ok(())
515    }
516
517    async fn reload(&mut self) -> AdapterResult<()> {
518        self.stop().await?;
519        self.start().await
520    }
521
522    async fn stop(&mut self) -> AdapterResult<()> {
523        let _ = self.cancel().await?;
524        Ok(())
525    }
526
527    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
528        let event = self.receiver.recv().await;
529        if matches!(
530            event.as_ref(),
531            Some(Ok(
532                AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. }
533            ))
534        ) {
535            if self.session_id.is_none()
536                && let Ok(session) = self.announced_session.lock()
537            {
538                self.session_id = session.clone();
539            }
540            if let Some(mut child) = self.child.take() {
541                let _ = child.wait().await;
542            }
543        }
544        event
545    }
546}
547
548#[cfg(test)]
549mod tests {
550    use super::{ClaudeAdapter, ParserState, parse_stream_event, parse_tool_results};
551    use crate::{AgentAdapter, AgentEvent, ToolStatus};
552    use serde_json::json;
553
554    #[test]
555    fn parses_claude_text_thought_and_tool_events() {
556        let mut state = ParserState::default();
557        assert!(matches!(
558            parse_stream_event(2, &json!({"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"text_delta","text":"hello"}}}), &mut state),
559            Some(AgentEvent::Text { slot: 2, text }) if text == "hello"
560        ));
561        assert!(matches!(
562            parse_stream_event(2, &json!({"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"thinking_delta","thinking":"check"}}}), &mut state),
563            Some(AgentEvent::Thought { text, .. }) if text == "check"
564        ));
565        assert!(matches!(
566            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),
567            Some(AgentEvent::Tool { update, .. }) if update.status == ToolStatus::Running && update.title == "Read"
568        ));
569        assert_eq!(
570            parse_stream_event(
571                2,
572                &json!({"type":"stream_event","event":{"type":"content_block_stop","index":1}}),
573                &mut state
574            ),
575            None
576        );
577
578        let completed = parse_tool_results(
579            2,
580            &json!({
581                "type": "user",
582                "message": {
583                    "content": [{
584                        "type": "tool_result",
585                        "tool_use_id": "tool-1",
586                        "content": [{"type": "text", "text": "file contents"}]
587                    }]
588                }
589            }),
590            &mut state,
591        );
592        assert!(matches!(
593            completed.as_slice(),
594            [AgentEvent::Tool { update, .. }]
595                if update.status == ToolStatus::Completed
596                    && update.title == "Read"
597                    && update.detail.as_deref() == Some("file contents")
598        ));
599    }
600
601    #[test]
602    fn failed_claude_tool_result_preserves_error_detail() {
603        let mut state = ParserState::default();
604        let failed = parse_tool_results(
605            4,
606            &json!({
607                "type": "user",
608                "message": {
609                    "content": [{
610                        "type": "tool_result",
611                        "tool_use_id": "tool-without-partial-start",
612                        "is_error": true,
613                        "content": "permission denied"
614                    }]
615                }
616            }),
617            &mut state,
618        );
619        assert!(matches!(
620            failed.as_slice(),
621            [AgentEvent::Tool { slot: 4, update }]
622                if update.status == ToolStatus::Failed
623                    && update.id == "tool-without-partial-start"
624                    && update.detail.as_deref() == Some("permission denied")
625        ));
626    }
627
628    #[tokio::test]
629    async fn advertises_only_noninteractive_modes_and_accepts_full_model_ids() {
630        let mut adapter = ClaudeAdapter::new(0, std::env::current_dir().unwrap(), "claude");
631        assert_eq!(
632            ClaudeAdapter::modes()
633                .into_iter()
634                .map(|mode| mode.id)
635                .collect::<Vec<_>>(),
636            ["plan", "bypassPermissions"]
637        );
638        assert_eq!(
639            adapter
640                .models()
641                .into_iter()
642                .map(|model| model.id)
643                .collect::<Vec<_>>(),
644            ["fable", "opus", "sonnet", "haiku"]
645        );
646        assert!(adapter.set_mode("manual".into()).await.is_err());
647        adapter
648            .set_model("claude-sonnet-4-5-20250929".into())
649            .await
650            .unwrap();
651        assert!(
652            adapter
653                .models()
654                .iter()
655                .any(|model| model.id == "claude-sonnet-4-5-20250929")
656        );
657        assert!(adapter.set_model("made-up-alias".into()).await.is_err());
658    }
659
660    #[tokio::test]
661    async fn native_claude_process_captures_session_and_resumes() {
662        let args_path =
663            std::env::temp_dir().join(format!("codeswarm-claude-args-{}", std::process::id()));
664        let stdin_path =
665            std::env::temp_dir().join(format!("codeswarm-claude-stdin-{}", std::process::id()));
666        let script_path =
667            std::env::temp_dir().join(format!("codeswarm-claude-script-{}", std::process::id()));
668        std::fs::write(
669            &script_path,
670            format!(
671                "printf '%s\\n' \"$*\" >> '{}'\nsed -n 'p' >> '{}'\nprintf '\\n' >> '{}'\nprintf '%s\\n' '{{\"type\":\"system\",\"session_id\":\"session-native\"}}' '{{\"type\":\"result\",\"subtype\":\"success\",\"result\":\"hello\",\"session_id\":\"session-native\"}}'\n",
672                args_path.display(),
673                stdin_path.display(),
674                stdin_path.display()
675            ),
676        )
677        .unwrap();
678        let mut adapter = ClaudeAdapter::new(
679            0,
680            std::env::current_dir().unwrap(),
681            format!("sh {}", script_path.display()),
682        );
683        adapter.start().await.unwrap();
684        assert!(matches!(
685            adapter.next_event().await,
686            Some(Ok(AgentEvent::ModesReplaced { .. }))
687        ));
688        assert!(matches!(
689            adapter.next_event().await,
690            Some(Ok(AgentEvent::ModelsReplaced { .. }))
691        ));
692        assert!(matches!(
693            adapter.next_event().await,
694            Some(Ok(AgentEvent::Ready { .. }))
695        ));
696        adapter.send_prompt("-first prompt".into()).await.unwrap();
697        assert!(
698            matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "hello")
699        );
700        assert!(matches!(
701            adapter.next_event().await,
702            Some(Ok(AgentEvent::TurnComplete { .. }))
703        ));
704        assert_eq!(adapter.session_id(), Some("session-native".into()));
705        adapter.send_prompt("second prompt".into()).await.unwrap();
706        assert!(matches!(
707            adapter.next_event().await,
708            Some(Ok(AgentEvent::Text { .. }))
709        ));
710        assert!(matches!(
711            adapter.next_event().await,
712            Some(Ok(AgentEvent::TurnComplete { .. }))
713        ));
714        let args = std::fs::read_to_string(&args_path).unwrap();
715        assert!(args.contains("--resume session-native"), "{args}");
716        assert!(!args.contains("first prompt"), "{args}");
717        assert!(!args.contains("second prompt"), "{args}");
718        assert_eq!(
719            std::fs::read_to_string(&stdin_path).unwrap(),
720            "-first prompt\nsecond prompt\n"
721        );
722        adapter.stop().await.unwrap();
723        let _ = std::fs::remove_file(args_path);
724        let _ = std::fs::remove_file(stdin_path);
725        let _ = std::fs::remove_file(script_path);
726    }
727}