1use 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 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 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}