1use std::collections::{BTreeMap, VecDeque};
7use std::io::Read;
8use std::path::{Path, PathBuf};
9use std::process::Stdio;
10use std::sync::{
11 Arc, Mutex,
12 atomic::{AtomicBool, AtomicUsize, Ordering},
13};
14
15use crate::{
16 AgentCapabilities, AgentCommand, AgentEvent, Effect, EventLog, Mode, PermissionAnswer,
17 PermissionRequest, RosterSlot, RosterUpdate, SessionState, TerminalEvent, ToolStatus,
18 ToolUpdate, UsageUpdate,
19 persistence::{BufferedSessionMetadataStore, SessionMetadata},
20 reduce,
21 relay::{
22 CollaborationStrategy, DEFAULT_STOP_ACKNOWLEDGMENT, Relay, RelayDecision, STOP_TOKEN,
23 is_usage_limit_response, strip_stop_token,
24 },
25 resources,
26};
27use async_trait::async_trait;
28use base64::{Engine, engine::general_purpose::STANDARD as BASE64};
29use serde_json::Value;
30use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWriteExt, BufReader};
31use tokio::process::{Child, ChildStdout, Command};
32use tokio::sync::{Mutex as AsyncMutex, Notify, mpsc};
33
34pub type AdapterResult<T> = Result<T, AdapterError>;
35
36const MAX_ACP_LINE_BYTES: usize = 10 * 1024 * 1024;
40const MAX_FILE_READ_BYTES: usize = 4 * 1024 * 1024;
41const MAX_TERMINAL_OUTPUT_BYTES: usize = 1024 * 1024;
42
43#[derive(Clone, Debug)]
44struct TerminalProcess {
45 child: Arc<AsyncMutex<Option<Child>>>,
46 output: Arc<Mutex<Vec<u8>>>,
47 truncated: Arc<AtomicBool>,
48 output_readers: Arc<AtomicUsize>,
49}
50
51impl TerminalProcess {
52 async fn kill(&self) {
53 if let Some(child) = self.child.lock().await.as_mut() {
54 #[cfg(unix)]
55 if signal_isolated_process_group(child, nix::sys::signal::Signal::SIGTERM) {
56 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
57 signal_isolated_process_group(child, nix::sys::signal::Signal::SIGKILL);
58 }
59 let _ = child.start_kill();
60 }
61 }
62
63 async fn stop(&self) {
64 if let Some(mut child) = self.child.lock().await.take() {
65 let _ = terminate_child(&mut child).await;
66 }
67 }
68
69 async fn wait(&self) -> Option<i32> {
70 loop {
71 let code = {
72 let mut child = self.child.lock().await;
73 match child.as_mut() {
74 None => Some(-1),
75 Some(child) => match child.try_wait() {
76 Ok(Some(status)) => Some(status.code().unwrap_or(-1)),
77 Ok(None) => None,
78 Err(_) => Some(-1),
83 },
84 }
85 };
86 if code.is_some() {
87 while self.output_readers.load(Ordering::Acquire) != 0 {
88 tokio::task::yield_now().await;
89 }
90 return code;
91 }
92 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
93 }
94 }
95
96 async fn exit_code(&self) -> Option<i32> {
97 let mut child = self.child.lock().await;
98 child
99 .as_mut()
100 .and_then(|child| child.try_wait().ok().flatten())
101 .map(|status| status.code().unwrap_or(-1))
102 }
103}
104
105#[derive(Clone, Debug)]
106pub struct HostUpdate {
107 pub event: AgentEvent,
108 pub effects: Vec<Effect>,
109}
110
111#[derive(Clone, Debug, Eq, PartialEq)]
112pub enum AdapterError {
113 Unsupported(&'static str),
114 Spawn(String),
115 Transport(String),
116 Protocol(String),
117}
118
119impl std::fmt::Display for AdapterError {
120 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
121 match self {
122 Self::Unsupported(operation) => write!(formatter, "unsupported operation: {operation}"),
123 Self::Spawn(error) => write!(formatter, "unable to launch agent: {error}"),
124 Self::Transport(error) => write!(formatter, "agent transport error: {error}"),
125 Self::Protocol(error) => write!(formatter, "agent protocol error: {error}"),
126 }
127 }
128}
129
130impl std::error::Error for AdapterError {}
131
132fn floor_char_boundary(text: &str, index: usize) -> usize {
133 let mut index = index.min(text.len());
134 while index > 0 && !text.is_char_boundary(index) {
135 index -= 1;
136 }
137 index
138}
139
140#[derive(Clone, Debug, Eq, PartialEq)]
148pub enum CommandParseError {
149 Empty,
150 UnterminatedQuote,
151 TrailingEscape,
152}
153
154impl std::fmt::Display for CommandParseError {
155 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
156 let message = match self {
157 Self::Empty => "command is empty",
158 Self::UnterminatedQuote => "command contains an unterminated quote",
159 Self::TrailingEscape => "command ends with an incomplete escape",
160 };
161 formatter.write_str(message)
162 }
163}
164
165impl std::error::Error for CommandParseError {}
166
167pub fn parse_command_line(command: &str) -> Result<(String, Vec<String>), CommandParseError> {
170 let mut argv = Vec::new();
171 let mut argument = String::new();
172 let mut quoted = None;
173 let mut escaped = false;
174 let mut started = false;
175
176 for character in command.chars() {
177 if escaped {
178 argument.push(character);
179 escaped = false;
180 started = true;
181 continue;
182 }
183 match (quoted, character) {
184 (_, '\\') if quoted != Some('\'') => {
185 escaped = true;
186 started = true;
187 }
188 (None, '\'' | '"') => {
189 quoted = Some(character);
190 started = true;
191 }
192 (Some(quote), character) if character == quote => quoted = None,
193 (None, character) if character.is_whitespace() => {
194 if started {
195 argv.push(std::mem::take(&mut argument));
196 started = false;
197 }
198 }
199 (_, character) => {
200 argument.push(character);
201 started = true;
202 }
203 }
204 }
205
206 if escaped {
207 return Err(CommandParseError::TrailingEscape);
208 }
209 if quoted.is_some() {
210 return Err(CommandParseError::UnterminatedQuote);
211 }
212 if started {
213 argv.push(argument);
214 }
215 let Some((program, args)) = argv.split_first() else {
216 return Err(CommandParseError::Empty);
217 };
218 Ok((program.clone(), args.to_vec()))
219}
220
221async fn terminate_child(child: &mut Child) -> AdapterResult<()> {
225 #[cfg(unix)]
226 if signal_isolated_process_group(child, nix::sys::signal::Signal::SIGTERM) {
227 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
228 signal_isolated_process_group(child, nix::sys::signal::Signal::SIGKILL);
229 }
230 let kill_error = child.start_kill().err();
231 let wait_error = child.wait().await.err();
232 if let Some(error) = kill_error.or(wait_error) {
233 return Err(AdapterError::Transport(error.to_string()));
234 }
235 Ok(())
236}
237
238#[cfg(unix)]
239fn signal_isolated_process_group(child: &Child, signal: nix::sys::signal::Signal) -> bool {
240 use nix::{
241 sys::signal::killpg,
242 unistd::{Pid, getpgid, getpgrp},
243 };
244
245 let Some(raw_pid) = child.id().and_then(|pid| i32::try_from(pid).ok()) else {
246 return false;
247 };
248 let pid = Pid::from_raw(raw_pid);
249 if getpgid(Some(pid)).ok() == Some(pid) && pid != getpgrp() {
254 let _ = killpg(pid, signal);
255 true
256 } else {
257 false
258 }
259}
260
261fn isolate_process_group(command: &mut Command) {
262 #[cfg(unix)]
263 command.process_group(0);
264}
265
266async fn drain_bounded<R>(mut reader: R, limit: usize) -> String
269where
270 R: AsyncRead + Unpin,
271{
272 let mut bytes = Vec::new();
273 let mut chunk = [0_u8; 4096];
274 while let Ok(count) = reader.read(&mut chunk).await {
275 if count == 0 {
276 break;
277 }
278 bytes.extend_from_slice(&chunk[..count]);
279 if bytes.len() > limit {
280 let keep_from = bytes.len() - limit;
281 bytes.drain(..keep_from);
282 }
283 }
284 String::from_utf8_lossy(&bytes).trim().to_owned()
285}
286
287async fn read_bounded_line<R>(reader: &mut R) -> AdapterResult<String>
292where
293 R: AsyncBufRead + Unpin,
294{
295 let mut bytes = Vec::with_capacity(4096);
296 loop {
297 let buffer = reader
298 .fill_buf()
299 .await
300 .map_err(|error| AdapterError::Transport(error.to_string()))?;
301 if buffer.is_empty() {
302 if bytes.is_empty() {
303 return Err(AdapterError::Transport("ACP stream closed".into()));
304 }
305 break;
306 }
307 let newline = buffer.iter().position(|byte| *byte == b'\n');
308 let available = newline.map_or(buffer.len(), |index| index + 1);
309 let remaining = MAX_ACP_LINE_BYTES
310 .saturating_add(1)
311 .saturating_sub(bytes.len());
312 if available > remaining {
313 reader.consume(remaining);
314 return Err(AdapterError::Protocol(format!(
315 "ACP protocol line exceeds {MAX_ACP_LINE_BYTES} bytes"
316 )));
317 }
318 bytes.extend_from_slice(&buffer[..available]);
319 reader.consume(available);
320 if newline.is_some() {
321 break;
322 }
323 }
324 if bytes.len() > MAX_ACP_LINE_BYTES {
325 return Err(AdapterError::Protocol(format!(
326 "ACP protocol line exceeds {MAX_ACP_LINE_BYTES} bytes"
327 )));
328 }
329 String::from_utf8(bytes).map_err(|error| AdapterError::Protocol(error.to_string()))
330}
331
332async fn drain_terminal_output<R>(
333 mut reader: R,
334 output: Arc<Mutex<Vec<u8>>>,
335 truncated: Arc<AtomicBool>,
336 output_readers: Arc<AtomicUsize>,
337 limit: usize,
338) where
339 R: AsyncRead + Unpin,
340{
341 let mut chunk = [0_u8; 4096];
342 while let Ok(count) = reader.read(&mut chunk).await {
343 if count == 0 {
344 break;
345 }
346 if let Ok(mut bytes) = output.lock() {
347 let remaining = limit.saturating_sub(bytes.len());
348 if count > remaining {
349 bytes.extend_from_slice(&chunk[..remaining]);
350 truncated.store(true, Ordering::Release);
351 } else {
352 bytes.extend_from_slice(&chunk[..count]);
353 }
354 }
355 }
356 output_readers.fetch_sub(1, Ordering::AcqRel);
357}
358
359#[async_trait]
361pub trait AgentAdapter: Send {
362 fn slot(&self) -> RosterSlot;
363 fn display_name(&self) -> String {
367 format!("Agent {}", self.slot().saturating_add(1))
368 }
369 fn session_id(&self) -> Option<String> {
373 None
374 }
375 fn protocol(&self) -> &'static str {
379 "custom"
380 }
381 fn capabilities(&self) -> AgentCapabilities;
382 async fn start(&mut self) -> AdapterResult<()>;
383 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()>;
384 async fn cancel(&mut self) -> AdapterResult<bool>;
385 async fn answer_permission(
386 &mut self,
387 request_id: String,
388 answer: PermissionAnswer,
389 ) -> AdapterResult<()>;
390 async fn set_mode(&mut self, mode: String) -> AdapterResult<()>;
391 async fn set_model(&mut self, _model: String) -> AdapterResult<()> {
392 Err(AdapterError::Unsupported("set_model"))
393 }
394 async fn reload(&mut self) -> AdapterResult<()>;
395 async fn stop(&mut self) -> AdapterResult<()>;
396 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>>;
397}
398
399struct SlotMappedAdapter {
404 logical_slot: RosterSlot,
405 inner: Box<dyn AgentAdapter>,
406}
407
408impl std::fmt::Debug for SlotMappedAdapter {
409 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
410 formatter
411 .debug_struct("SlotMappedAdapter")
412 .field("logical_slot", &self.logical_slot)
413 .field("inner_slot", &self.inner.slot())
414 .finish_non_exhaustive()
415 }
416}
417
418fn map_event_slot(event: AgentEvent, slot: RosterSlot) -> AgentEvent {
419 match event {
420 AgentEvent::GoalUpdated { goal } => AgentEvent::GoalUpdated { goal },
421 AgentEvent::RosterUpdated { update } => AgentEvent::RosterUpdated { update },
422 AgentEvent::Ready { capabilities, .. } => AgentEvent::Ready { slot, capabilities },
423 AgentEvent::TurnStarted { .. } => AgentEvent::TurnStarted { slot },
424 AgentEvent::ModesReplaced {
425 modes,
426 current_mode,
427 ..
428 } => AgentEvent::ModesReplaced {
429 slot,
430 modes,
431 current_mode,
432 },
433 AgentEvent::ModeUpdated { current_mode, .. } => {
434 AgentEvent::ModeUpdated { slot, current_mode }
435 }
436 AgentEvent::ModelsReplaced {
437 config_id,
438 models,
439 current_model,
440 ..
441 } => AgentEvent::ModelsReplaced {
442 slot,
443 config_id,
444 models,
445 current_model,
446 },
447 AgentEvent::ModelUpdated { current_model, .. } => AgentEvent::ModelUpdated {
448 slot,
449 current_model,
450 },
451 AgentEvent::UserText { text, .. } => AgentEvent::UserText { slot, text },
452 AgentEvent::CommandsReplaced { commands, .. } => {
453 AgentEvent::CommandsReplaced { slot, commands }
454 }
455 AgentEvent::UsageUpdated { usage, .. } => AgentEvent::UsageUpdated { slot, usage },
456 AgentEvent::Text { text, .. } => AgentEvent::Text { slot, text },
457 AgentEvent::Thought { text, .. } => AgentEvent::Thought { slot, text },
458 AgentEvent::Tool { update, .. } => AgentEvent::Tool { slot, update },
459 AgentEvent::Permission { request, .. } => AgentEvent::Permission { slot, request },
460 AgentEvent::Terminal { event, .. } => AgentEvent::Terminal { slot, event },
461 AgentEvent::TurnComplete { .. } => AgentEvent::TurnComplete { slot },
462 AgentEvent::UsageLimitReached { detail, .. } => {
463 AgentEvent::UsageLimitReached { slot, detail }
464 }
465 AgentEvent::Failed {
466 started, detail, ..
467 } => AgentEvent::Failed {
468 slot,
469 started,
470 detail,
471 },
472 }
473}
474
475#[async_trait]
476impl AgentAdapter for SlotMappedAdapter {
477 fn slot(&self) -> RosterSlot {
478 self.logical_slot
479 }
480
481 fn display_name(&self) -> String {
482 self.inner.display_name()
483 }
484
485 fn session_id(&self) -> Option<String> {
486 self.inner.session_id()
487 }
488
489 fn protocol(&self) -> &'static str {
490 self.inner.protocol()
491 }
492
493 fn capabilities(&self) -> AgentCapabilities {
494 self.inner.capabilities()
495 }
496
497 async fn start(&mut self) -> AdapterResult<()> {
498 self.inner.start().await
499 }
500
501 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
502 self.inner.send_prompt(prompt).await
503 }
504
505 async fn cancel(&mut self) -> AdapterResult<bool> {
506 self.inner.cancel().await
507 }
508
509 async fn answer_permission(
510 &mut self,
511 request_id: String,
512 answer: PermissionAnswer,
513 ) -> AdapterResult<()> {
514 self.inner.answer_permission(request_id, answer).await
515 }
516
517 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
518 self.inner.set_mode(mode).await
519 }
520
521 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
522 self.inner.set_model(model).await
523 }
524
525 async fn reload(&mut self) -> AdapterResult<()> {
526 self.inner.reload().await
527 }
528
529 async fn stop(&mut self) -> AdapterResult<()> {
530 self.inner.stop().await
531 }
532
533 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
534 let slot = self.logical_slot;
535 self.inner
536 .next_event()
537 .await
538 .map(|result| result.map(|event| map_event_slot(event, slot)))
539 }
540}
541
542pub struct AdapterHost {
545 adapter: Box<dyn AgentAdapter>,
546 pub state: SessionState,
547 pub last_error: Option<String>,
548 event_log: Option<EventLog>,
549}
550
551impl std::fmt::Debug for AdapterHost {
552 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
553 formatter
554 .debug_struct("AdapterHost")
555 .field("state", &self.state)
556 .field("last_error", &self.last_error)
557 .field("event_log", &self.event_log)
558 .finish_non_exhaustive()
559 }
560}
561
562impl AdapterHost {
563 pub fn new(adapter: Box<dyn AgentAdapter>, event_log: Option<EventLog>) -> Self {
564 let slot = adapter.slot();
565 Self {
566 adapter,
567 state: SessionState::new(slot.saturating_add(1)),
568 last_error: None,
569 event_log,
570 }
571 }
572
573 pub async fn start(&mut self) -> AdapterResult<()> {
574 self.adapter.start().await
575 }
576
577 pub async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
578 self.adapter.send_prompt(prompt).await
579 }
580
581 pub async fn cancel(&mut self) -> AdapterResult<bool> {
582 self.adapter.cancel().await
583 }
584
585 pub async fn answer_permission(
586 &mut self,
587 request_id: String,
588 answer: PermissionAnswer,
589 ) -> AdapterResult<()> {
590 self.adapter.answer_permission(request_id, answer).await
591 }
592
593 pub async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
594 self.adapter.set_mode(mode).await
595 }
596
597 pub async fn set_model(&mut self, model: String) -> AdapterResult<()> {
598 self.adapter.set_model(model).await
599 }
600
601 pub async fn reload(&mut self) -> AdapterResult<()> {
602 self.adapter.reload().await?;
603 let slot = self.adapter.slot();
604 if let Some(agent) = self.state.slots.get_mut(slot) {
605 agent.active = true;
606 agent.capabilities = self.adapter.capabilities();
607 }
608 self.last_error = None;
609 Ok(())
610 }
611
612 pub async fn stop(&mut self) -> AdapterResult<()> {
613 self.adapter.stop().await
614 }
615
616 pub async fn next_effects(&mut self) -> Option<AdapterResult<Vec<Effect>>> {
617 Some(self.next_update().await?.map(|update| update.effects))
618 }
619
620 pub async fn next_update(&mut self) -> Option<AdapterResult<HostUpdate>> {
621 let event = match self.adapter.next_event().await {
622 None => return None,
623 Some(Err(error)) => {
624 self.last_error = Some(error.to_string());
625 let slot = self.adapter.slot();
626 let failure = AgentEvent::Failed {
627 slot,
628 started: true,
629 detail: error.to_string(),
630 };
631 let effects = reduce(&mut self.state, failure.clone());
632 return Some(Ok(HostUpdate {
633 event: failure,
634 effects,
635 }));
636 }
637 Some(Ok(event)) => event,
638 };
639 if let Some(log) = &self.event_log
640 && let Err(error) = log.append(&event)
641 {
642 return Some(Err(AdapterError::Transport(error.to_string())));
643 }
644 let effects = reduce(&mut self.state, event.clone());
645 Some(Ok(HostUpdate { event, effects }))
646 }
647
648 pub fn adapter(&self) -> &dyn AgentAdapter {
649 &*self.adapter
650 }
651
652 pub fn session_id(&self) -> Option<String> {
653 self.adapter.session_id()
654 }
655
656 fn remap(self, logical_slot: RosterSlot) -> Self {
659 let old_slot = self.adapter.slot();
660 if old_slot == logical_slot {
661 return self;
662 }
663
664 let mut state = self.state;
665 if state.slots.len() <= logical_slot {
666 state
667 .slots
668 .resize(logical_slot.saturating_add(1), Default::default());
669 }
670 if let Some(agent) = state.slots.get(old_slot).cloned() {
671 state.slots[logical_slot] = agent;
672 }
673 if state.active_slot == Some(old_slot) {
674 state.active_slot = Some(logical_slot);
675 }
676 for (slot, _) in &mut state.queued_prompts {
677 if *slot == old_slot {
678 *slot = logical_slot;
679 }
680 }
681 for (slot, _) in &mut state.public_text {
682 if *slot == old_slot {
683 *slot = logical_slot;
684 }
685 }
686
687 Self {
688 adapter: Box::new(SlotMappedAdapter {
689 logical_slot,
690 inner: self.adapter,
691 }),
692 state,
693 last_error: self.last_error,
694 event_log: self.event_log,
695 }
696 }
697}
698
699pub struct RelayHost {
702 goal: Option<crate::goal::Goal>,
703 hosts: Vec<AdapterHost>,
704 relay: Relay,
705 introduced: Vec<bool>,
706 roster_names: Vec<String>,
707 roster_identities: Vec<String>,
708 roster_launch_specs: Vec<(String, String)>,
709 desired_policy: String,
710 metadata_writer: Option<BufferedSessionMetadataStore>,
711 metadata_workspace: Option<String>,
712 dispatches: Vec<(RosterSlot, String)>,
713 event_sink: Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
714 cancel_requested: Arc<AtomicBool>,
715 cancel_notify: Arc<Notify>,
716}
717
718#[derive(Clone, Debug)]
721pub struct RelayCancellation {
722 requested: Arc<AtomicBool>,
723 notify: Arc<Notify>,
724}
725
726#[derive(Debug)]
731pub struct RelayPermissionAnswer {
732 pub slot: RosterSlot,
733 pub request_id: String,
734 pub answer: PermissionAnswer,
735}
736
737#[cfg(test)]
738const CANCEL_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(250);
739#[cfg(not(test))]
740const CANCEL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
741
742#[cfg(test)]
743const CANCEL_SETTLE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(20);
744#[cfg(not(test))]
745const CANCEL_SETTLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
746
747async fn cancel_with_timeout(host: &mut AdapterHost) -> AdapterResult<bool> {
752 tokio::time::timeout(CANCEL_TIMEOUT, host.cancel())
753 .await
754 .map_err(|_| AdapterError::Transport("adapter cancellation timed out".into()))?
755}
756
757fn canonical_policy_id(policy: &str) -> &str {
758 match policy {
759 "plan" => "codeswarm:mode:plan",
760 "default" | "manual" => "codeswarm:mode:manual",
761 "accept-edits" => "codeswarm:mode:accept-edits",
762 "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
763 other => other,
764 }
765}
766
767async fn apply_policy_to_host(host: &mut AdapterHost, policy: &str) -> AdapterResult<()> {
768 if !host.adapter().capabilities().supports_modes {
769 return Ok(());
770 }
771 let policy_id = canonical_policy_id(policy);
772 let slot = host.adapter().slot();
773 let advertised = host
774 .state
775 .slots
776 .get(slot)
777 .map(|agent| agent.modes.as_slice())
778 .unwrap_or_default();
779 let native = if advertised.is_empty() {
780 match policy_id {
781 "codeswarm:mode:plan" => "plan".into(),
782 "codeswarm:mode:manual" => "default".into(),
783 "codeswarm:mode:accept-edits" => "accept-edits".into(),
784 "codeswarm:mode:full-access" => "full-access".into(),
785 other => other.into(),
786 }
787 } else {
788 crate::policy::resolve(policy_id, advertised)
789 .map(|mode| mode.id)
790 .ok_or(AdapterError::Unsupported(
791 "desired policy is unavailable for adapter",
792 ))?
793 };
794 host.set_mode(native).await
795}
796
797async fn refresh_mode_catalog(
798 host: &mut AdapterHost,
799 event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
800) -> AdapterResult<bool> {
801 if !host.adapter().capabilities().supports_modes {
802 return Ok(false);
803 }
804 tokio::time::timeout(std::time::Duration::from_secs(2), async {
805 let mut ready_seen = false;
806 loop {
807 let update = host.next_update().await.ok_or_else(|| {
808 AdapterError::Transport("adapter ended before advertising modes".into())
809 })??;
810 let catalog_ready = matches!(update.event, AgentEvent::ModesReplaced { .. });
811 ready_seen |= matches!(update.event, AgentEvent::Ready { .. });
812 if let Some(sink) = event_sink {
813 sink(update.event);
814 }
815 if catalog_ready {
816 return Ok(ready_seen);
817 }
818 }
819 })
820 .await
821 .map_err(|_| AdapterError::Transport("adapter mode catalog timed out".into()))?
822}
823
824async fn refresh_adapter_startup(
825 host: &mut AdapterHost,
826 event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
827) -> AdapterResult<()> {
828 let ready_seen = refresh_mode_catalog(host, event_sink).await?;
829 if ready_seen || !matches!(host.adapter().protocol(), "native" | "acp") {
830 return Ok(());
831 }
832 tokio::time::timeout(std::time::Duration::from_secs(2), async {
833 loop {
834 let update = host.next_update().await.ok_or_else(|| {
835 AdapterError::Transport("adapter ended before becoming ready".into())
836 })??;
837 let ready = matches!(update.event, AgentEvent::Ready { .. });
838 if let Some(sink) = event_sink {
839 sink(update.event);
840 }
841 if ready {
842 return Ok(());
843 }
844 }
845 })
846 .await
847 .map_err(|_| AdapterError::Transport("adapter ready handshake timed out".into()))?
848}
849
850fn public_context_speaker(name: &str) -> String {
851 let now = time::OffsetDateTime::now_local().unwrap_or_else(|_| time::OffsetDateTime::now_utc());
852 format!("{name} {:02}:{:02}", now.hour(), now.minute())
853}
854
855impl RelayCancellation {
856 pub fn request(&self) {
857 self.requested.store(true, Ordering::Release);
858 self.notify.notify_one();
859 }
860}
861
862impl std::fmt::Debug for RelayHost {
863 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
864 formatter
865 .debug_struct("RelayHost")
866 .field("hosts", &self.hosts)
867 .field("relay", &self.relay)
868 .field("dispatches", &self.dispatches)
869 .field("event_sink", &self.event_sink.is_some())
870 .field(
871 "cancel_requested",
872 &self.cancel_requested.load(Ordering::Acquire),
873 )
874 .finish()
875 }
876}
877
878impl RelayHost {
879 pub fn new(hosts: Vec<AdapterHost>, max_rounds: usize) -> Result<Self, AdapterError> {
880 if hosts.is_empty() {
881 return Err(AdapterError::Unsupported("relay requires an adapter"));
882 }
883 Ok(Self {
884 relay: Relay::new(hosts.len(), max_rounds),
885 introduced: vec![false; hosts.len()],
886 roster_names: hosts
887 .iter()
888 .map(|host| host.adapter().display_name())
889 .collect(),
890 roster_identities: hosts
891 .iter()
892 .map(|host| host.adapter().display_name())
893 .collect(),
894 roster_launch_specs: Vec::new(),
895 desired_policy: crate::policy::DEFAULT_POLICY_ID.into(),
896 metadata_writer: None,
897 metadata_workspace: None,
898 hosts,
899 dispatches: Vec::new(),
900 goal: None,
901 event_sink: None,
902 cancel_requested: Arc::new(AtomicBool::new(false)),
903 cancel_notify: Arc::new(Notify::new()),
904 })
905 }
906
907 pub fn set_event_sink<F>(&mut self, sink: F)
911 where
912 F: Fn(AgentEvent) + Send + Sync + 'static,
913 {
914 self.event_sink = Some(Arc::new(sink));
915 }
916
917 pub fn set_roster_names(&mut self, names: Vec<String>) {
920 if names.len() == self.hosts.len() {
921 self.roster_names = names;
922 }
923 }
924
925 pub fn set_roster_identities(&mut self, identities: Vec<String>) {
929 if identities.len() == self.hosts.len() {
930 self.roster_identities = identities;
931 }
932 }
933
934 pub fn set_roster_launch_specs(&mut self, specs: Vec<(String, String)>) {
935 if specs.len() == self.hosts.len() {
936 self.roster_launch_specs = specs;
937 }
938 }
939
940 pub fn set_session_metadata_writer(&mut self, writer: BufferedSessionMetadataStore) {
943 self.metadata_writer = Some(writer);
944 }
945
946 pub fn set_session_metadata_workspace(&mut self, workspace: impl Into<String>) {
947 self.metadata_workspace = Some(workspace.into());
948 }
949
950 pub fn session_metadata(&self) -> SessionMetadata {
952 let active = self.relay.active_slots().collect::<Vec<_>>();
953 let mut data = serde_json::Map::new();
954 data.insert(
955 "goal".into(),
956 serde_json::to_value(&self.goal).expect("goal serializes"),
957 );
958 if let Some(workspace) = &self.metadata_workspace {
959 data.insert("cwd".into(), serde_json::Value::String(workspace.clone()));
960 }
961 data.insert(
962 "title".into(),
963 serde_json::Value::String("CodeSwarm".into()),
964 );
965 data.insert(
966 "agents".into(),
967 serde_json::Value::Array(
968 active
969 .into_iter()
970 .filter_map(|slot| {
971 let host = self.hosts.get(slot)?;
972 let (protocol, command) = self.roster_launch_specs.get(slot)?;
973 let mut agent = serde_json::Map::new();
974 agent.insert(
975 "name".into(),
976 serde_json::Value::String(
977 self.roster_names
978 .get(slot)
979 .cloned()
980 .unwrap_or_else(|| host.adapter().display_name()),
981 ),
982 );
983 agent.insert(
984 "identity".into(),
985 serde_json::Value::String(
986 self.roster_identities
987 .get(slot)
988 .cloned()
989 .unwrap_or_else(|| host.adapter().display_name()),
990 ),
991 );
992 agent.insert(
993 "protocol".into(),
994 serde_json::Value::String(protocol.clone()),
995 );
996 agent.insert("command".into(), serde_json::Value::String(command.clone()));
997 agent.insert(
998 "supports_load_session".into(),
999 serde_json::Value::Bool(
1000 host.adapter().capabilities().supports_session_load,
1001 ),
1002 );
1003 if let Some(session_id) = host.session_id() {
1004 agent
1005 .insert("session_id".into(), serde_json::Value::String(session_id));
1006 }
1007 Some(serde_json::Value::Object(agent))
1008 })
1009 .collect(),
1010 ),
1011 );
1012 SessionMetadata::new(data)
1013 }
1014
1015 fn queue_session_metadata(&self) -> AdapterResult<()> {
1016 if let Some(writer) = &self.metadata_writer {
1017 writer
1018 .write(self.session_metadata())
1019 .map_err(|error| AdapterError::Transport(error.to_string()))?;
1020 }
1021 Ok(())
1022 }
1023
1024 pub fn restore_goal(&mut self, goal: Option<crate::goal::Goal>) {
1025 self.goal = goal;
1026 if let Some(sink) = &self.event_sink {
1027 sink(AgentEvent::GoalUpdated {
1028 goal: self.goal.clone(),
1029 });
1030 }
1031 }
1032
1033 pub fn apply_goal(
1034 &mut self,
1035 command: crate::goal::GoalCommand,
1036 ) -> Result<Option<String>, String> {
1037 let task = crate::goal::apply(&mut self.goal, command)?;
1038 if let Some(sink) = &self.event_sink {
1039 sink(AgentEvent::GoalUpdated {
1040 goal: self.goal.clone(),
1041 });
1042 }
1043 self.queue_session_metadata()
1044 .map_err(|error| error.to_string())?;
1045 Ok(task)
1046 }
1047
1048 pub fn roster_names(&self) -> &[String] {
1049 &self.roster_names
1050 }
1051
1052 pub fn session_ids(&self) -> Vec<Option<String>> {
1053 self.hosts.iter().map(AdapterHost::session_id).collect()
1054 }
1055
1056 pub async fn start(&mut self) -> AdapterResult<()> {
1057 let event_sink = self.event_sink.clone();
1058 let startups = self.hosts.iter_mut().map(|host| {
1059 let event_sink = event_sink.clone();
1060 async move {
1061 host.start().await?;
1062 refresh_adapter_startup(host, &event_sink).await
1063 }
1064 });
1065 let results = futures::future::join_all(startups).await;
1066 if let Some(error) = results.into_iter().find_map(Result::err) {
1067 for host in &mut self.hosts {
1070 let _ = host.stop().await;
1071 }
1072 return Err(error);
1073 }
1074 if let Err(error) = self
1075 .set_policy(crate::policy::DEFAULT_POLICY_ID.into())
1076 .await
1077 {
1078 for host in &mut self.hosts {
1079 let _ = host.stop().await;
1080 }
1081 return Err(error);
1082 }
1083 let _ = self.queue_session_metadata();
1084 Ok(())
1085 }
1086
1087 pub async fn stop(&mut self) -> AdapterResult<()> {
1088 let mut first_error = None;
1093 for host in &mut self.hosts {
1094 if let Err(error) = host.stop().await
1095 && first_error.is_none()
1096 {
1097 first_error = Some(error);
1098 }
1099 }
1100 if let Some(writer) = &self.metadata_writer
1101 && let Err(error) = writer.flush()
1102 && first_error.is_none()
1103 {
1104 first_error = Some(AdapterError::Transport(error.to_string()));
1105 }
1106 first_error.map_or(Ok(()), Err)
1107 }
1108
1109 pub async fn answer_permission(
1112 &mut self,
1113 slot: RosterSlot,
1114 request_id: String,
1115 answer: PermissionAnswer,
1116 ) -> AdapterResult<()> {
1117 let host = self
1118 .hosts
1119 .get_mut(slot)
1120 .ok_or_else(|| AdapterError::Transport("permission target is missing".into()))?;
1121 host.answer_permission(request_id, answer).await
1122 }
1123
1124 pub fn pause(&mut self) {
1125 self.relay.pause();
1126 }
1127
1128 pub fn resume(&mut self) {
1129 self.relay.resume();
1130 }
1131
1132 pub fn set_strategy(&mut self, strategy: CollaborationStrategy) {
1135 self.relay.set_strategy(strategy);
1136 }
1137
1138 pub fn strategy(&self) -> CollaborationStrategy {
1139 self.relay.strategy()
1140 }
1141
1142 pub fn roster_identity(&self, slot: RosterSlot) -> Option<&str> {
1143 self.roster_identities.get(slot).map(String::as_str)
1144 }
1145
1146 pub fn active_slot_for_identity(&self, identity: &str) -> Option<RosterSlot> {
1147 self.relay.active_slots().find(|slot| {
1148 self.roster_identity(*slot)
1149 .is_some_and(|candidate| candidate.eq_ignore_ascii_case(identity))
1150 })
1151 }
1152
1153 pub async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
1156 let active = self.relay.active_slots().collect::<Vec<_>>();
1157 for slot in active {
1158 let Some(host) = self.hosts.get_mut(slot) else {
1159 continue;
1160 };
1161 if host.adapter().capabilities().supports_modes {
1162 host.set_mode(mode.clone()).await?;
1163 }
1164 }
1165 Ok(())
1166 }
1167
1168 pub async fn set_model(&mut self, slot: RosterSlot, model: String) -> AdapterResult<()> {
1169 let host = self
1170 .hosts
1171 .get_mut(slot)
1172 .ok_or_else(|| AdapterError::Transport("model target is missing".into()))?;
1173 if !host.adapter().capabilities().supports_models {
1174 return Err(AdapterError::Unsupported("set_model"));
1175 }
1176 host.set_model(model.clone()).await?;
1177 if let Some(sink) = &self.event_sink {
1178 sink(AgentEvent::ModelUpdated {
1179 slot,
1180 current_model: model,
1181 });
1182 }
1183 Ok(())
1184 }
1185
1186 pub async fn set_policy(&mut self, policy: String) -> AdapterResult<()> {
1190 self.desired_policy = canonical_policy_id(&policy).to_owned();
1191 let desired_policy = self.desired_policy.clone();
1192 let mut first_error = None;
1193 let active = self.relay.active_slots().collect::<Vec<_>>();
1194 for active_slot in active {
1195 let Some(host) = self.hosts.get_mut(active_slot) else {
1196 continue;
1197 };
1198 if let Err(error) = apply_policy_to_host(host, &desired_policy).await
1199 && first_error.is_none()
1200 {
1201 first_error = Some(error);
1202 }
1203 }
1204 first_error.map_or(Ok(()), Err)
1205 }
1206
1207 pub async fn reload(&mut self, slot: RosterSlot) -> AdapterResult<()> {
1208 let desired_policy = self.desired_policy.clone();
1209 let event_sink = self.event_sink.clone();
1210 let host = self
1211 .hosts
1212 .get_mut(slot)
1213 .ok_or_else(|| AdapterError::Transport("reload target is missing".into()))?;
1214 host.reload().await?;
1215 refresh_adapter_startup(host, &event_sink).await?;
1216 apply_policy_to_host(host, &desired_policy).await?;
1217 if let Some(introduced) = self.introduced.get_mut(slot) {
1218 *introduced = false;
1219 }
1220 self.relay
1221 .reactivate(slot)
1222 .map_err(|error| AdapterError::Transport(error.into()))?;
1223 let _ = self.queue_session_metadata();
1224 if let Some(sink) = &self.event_sink {
1225 sink(AgentEvent::RosterUpdated {
1226 update: RosterUpdate::Reloaded { slot },
1227 });
1228 }
1229 let _ = self.relay.clear_limited(slot);
1230 Ok(())
1231 }
1232
1233 pub async fn drop_agent(&mut self, slot: RosterSlot) -> AdapterResult<()> {
1235 self.relay
1236 .drop_agent(slot)
1237 .map_err(|error| AdapterError::Transport(error.into()))?;
1238 let _stop_result = if let Some(host) = self.hosts.get_mut(slot) {
1239 host.stop().await
1240 } else {
1241 Ok(())
1242 };
1243 let _ = self.queue_session_metadata();
1247 if let Some(sink) = &self.event_sink {
1248 sink(AgentEvent::RosterUpdated {
1249 update: RosterUpdate::Dropped { slot },
1250 });
1251 }
1252 Ok(())
1253 }
1254
1255 pub async fn add_agent(
1259 &mut self,
1260 mut host: AdapterHost,
1261 name: impl Into<String>,
1262 identity: impl Into<String>,
1263 command: impl Into<String>,
1264 ) -> AdapterResult<RosterSlot> {
1265 let slot = self.hosts.len();
1266 if host.adapter().slot() != slot {
1267 return Err(AdapterError::Transport(
1268 "new adapter slot must append after the existing roster".into(),
1269 ));
1270 }
1271 if let Err(error) = host.start().await {
1272 let _ = host.stop().await;
1273 return Err(error);
1274 }
1275 if let Err(error) = refresh_adapter_startup(&mut host, &self.event_sink).await {
1276 let _ = host.stop().await;
1277 return Err(error);
1278 }
1279 if let Err(error) = apply_policy_to_host(&mut host, &self.desired_policy).await {
1280 let _ = host.stop().await;
1281 return Err(error);
1282 }
1283 let capabilities = host.adapter().capabilities();
1284 self.hosts.push(host);
1285 self.relay.add_agent();
1286 self.introduced.push(false);
1287 let name = name.into();
1288 let identity = identity.into();
1289 self.roster_names.push(name.clone());
1290 self.roster_identities.push(identity.clone());
1291 self.roster_launch_specs
1292 .push((self.hosts[slot].adapter().protocol().into(), command.into()));
1293 let _ = self.queue_session_metadata();
1294 if let Some(sink) = &self.event_sink {
1295 sink(AgentEvent::RosterUpdated {
1296 update: RosterUpdate::Added {
1297 slot,
1298 name,
1299 identity,
1300 },
1301 });
1302 sink(AgentEvent::Ready { slot, capabilities });
1303 }
1304 Ok(slot)
1305 }
1306
1307 pub fn swap_agents(&mut self, first: RosterSlot, second: RosterSlot) -> AdapterResult<()> {
1311 if first == second {
1312 return Ok(());
1313 }
1314 if first >= self.hosts.len() || second >= self.hosts.len() {
1315 return Err(AdapterError::Transport("roster slot out of range".into()));
1316 }
1317 self.relay
1318 .swap_agents(first, second)
1319 .map_err(|error| AdapterError::Transport(error.into()))?;
1320 let low = first.min(second);
1321 let high = first.max(second);
1322 let high_host = self.hosts.remove(high);
1323 let low_host = self.hosts.remove(low);
1324 self.hosts.insert(low, high_host.remap(low));
1325 self.hosts.insert(high, low_host.remap(high));
1326 self.roster_names.swap(first, second);
1327 self.roster_identities.swap(first, second);
1328 if self.roster_launch_specs.len() == self.hosts.len() {
1329 self.roster_launch_specs.swap(first, second);
1330 }
1331 self.introduced.swap(first, second);
1332 if let Some(sink) = &self.event_sink {
1333 sink(AgentEvent::RosterUpdated {
1334 update: RosterUpdate::Swapped { first, second },
1335 });
1336 sink(AgentEvent::Ready {
1337 slot: first,
1338 capabilities: self.hosts[first].adapter().capabilities(),
1339 });
1340 sink(AgentEvent::Ready {
1341 slot: second,
1342 capabilities: self.hosts[second].adapter().capabilities(),
1343 });
1344 }
1345 let _ = self.queue_session_metadata();
1346 Ok(())
1347 }
1348
1349 pub fn relay(&self) -> &Relay {
1350 &self.relay
1351 }
1352
1353 pub fn next_slot(&self) -> RosterSlot {
1354 self.hosts.len()
1355 }
1356
1357 pub fn relay_mut(&mut self) -> &mut Relay {
1358 &mut self.relay
1359 }
1360
1361 pub fn cancellation(&self) -> RelayCancellation {
1362 RelayCancellation {
1363 requested: Arc::clone(&self.cancel_requested),
1364 notify: Arc::clone(&self.cancel_notify),
1365 }
1366 }
1367
1368 pub fn dispatches(&self) -> &[(RosterSlot, String)] {
1372 &self.dispatches
1373 }
1374
1375 pub async fn run_turn(
1376 &mut self,
1377 task: impl Into<String>,
1378 first_slot: RosterSlot,
1379 ) -> AdapterResult<RelayDecision> {
1380 self.run_turn_inner(task.into(), first_slot, None).await
1381 }
1382
1383 pub async fn run_turn_with_permissions(
1384 &mut self,
1385 task: impl Into<String>,
1386 first_slot: RosterSlot,
1387 permissions: &mut tokio::sync::mpsc::UnboundedReceiver<RelayPermissionAnswer>,
1388 ) -> AdapterResult<RelayDecision> {
1389 self.run_turn_inner(task.into(), first_slot, Some(permissions))
1390 .await
1391 }
1392
1393 async fn run_turn_inner(
1394 &mut self,
1395 task: String,
1396 first_slot: RosterSlot,
1397 mut permissions: Option<&mut tokio::sync::mpsc::UnboundedReceiver<RelayPermissionAnswer>>,
1398 ) -> AdapterResult<RelayDecision> {
1399 let decision = self.relay.begin(task, first_slot);
1400 let RelayDecision::Dispatch {
1401 slot,
1402 prompt,
1403 direct,
1404 can_stop,
1405 } = &decision
1406 else {
1407 return Ok(decision);
1408 };
1409 let speaker_name = self
1410 .roster_names
1411 .get(*slot)
1412 .cloned()
1413 .unwrap_or_else(|| self.hosts[*slot].adapter().display_name());
1414 let unseen = self.relay.unseen_context(*slot);
1415 if !*direct && !prompt.trim().is_empty() {
1422 if self.relay.shared_task().is_none() {
1423 self.relay.set_shared_task(prompt.clone());
1424 }
1425 self.relay
1426 .record_public(public_context_speaker("User"), prompt.clone());
1427 }
1428 let prompt = if unseen.is_empty() {
1429 prompt.clone()
1430 } else {
1431 format!("{prompt}\n\nPublic updates:\n{unseen}")
1432 };
1433 let introduction = if !self.introduced.get(*slot).copied().unwrap_or(false) {
1434 let self_name = speaker_name.clone();
1435 let roster = self
1436 .relay
1437 .active_slots()
1438 .map(|candidate| {
1439 let name = self
1440 .roster_names
1441 .get(candidate)
1442 .cloned()
1443 .unwrap_or_else(|| self.hosts[candidate].adapter().display_name());
1444 let role = if candidate == *slot { " — you" } else { "" };
1445 format!("{}. {name}{role}", candidate.saturating_add(1))
1446 })
1447 .collect::<Vec<_>>();
1448 let shared_task = self
1449 .relay
1450 .shared_task()
1451 .filter(|task| *task != prompt)
1452 .map(|task| format!("\n\nShared task:\n{task}"))
1453 .unwrap_or_default();
1454 format!(
1455 "You are {self_name}.\nCodeSwarm roster (ordered):\n{}\n\
1456 Turns relay sequentially through this roster. Treat the user request as the shared task; use timestamped public updates as conversation context.{shared_task}",
1457 roster.join("\n")
1458 )
1459 } else {
1460 String::new()
1461 };
1462 let prompt = format!(
1463 "{introduction}{separator}{prompt}\n\n{}",
1464 if *can_stop {
1465 format!(
1466 "You are reviewing another agent. If no meaningful correction is needed,\nreply with an optional emoji followed by {STOP_TOKEN} on the final line.\nCodeSwarm hides the token and ends this review batch."
1467 )
1468 } else {
1469 format!(
1470 "Do not use {STOP_TOKEN} on this turn. Your response must be offered to another agent for review."
1471 )
1472 },
1473 separator = if introduction.is_empty() { "" } else { "\n\n" },
1474 );
1475 let event_sink = self.event_sink.clone();
1476 let host = self
1477 .hosts
1478 .get_mut(*slot)
1479 .ok_or_else(|| AdapterError::Transport("relay selected missing adapter".into()))?;
1480 let prompt = crate::goal::prompt(self.goal.as_ref(), &prompt);
1481 if let Err(error) = host.send_prompt(prompt.clone()).await {
1482 let limited =
1483 report_relay_failure(&mut self.relay, &event_sink, *slot, true, error.to_string());
1484 if limited {
1485 self.relay.finish(*slot, *direct, false);
1486 }
1487 let _ = self.queue_session_metadata();
1488 if limited {
1489 return Ok(decision);
1490 }
1491 return Err(error);
1492 }
1493 if let Some(sink) = &self.event_sink {
1494 sink(AgentEvent::TurnStarted { slot: *slot });
1495 }
1496 if let Some(introduced) = self.introduced.get_mut(*slot) {
1497 *introduced = true;
1498 }
1499 self.dispatches.push((*slot, prompt));
1500 let mut response = String::new();
1501 let mut emitted_text = 0usize;
1502 let completion_event = loop {
1503 if self.cancel_requested.swap(false, Ordering::AcqRel) {
1504 if let Err(error) = cancel_with_timeout(host).await {
1505 let limited = report_relay_failure(
1506 &mut self.relay,
1507 &event_sink,
1508 *slot,
1509 true,
1510 error.to_string(),
1511 );
1512 if limited {
1513 self.relay.finish(*slot, *direct, false);
1514 }
1515 let _ = self.queue_session_metadata();
1516 if limited {
1517 return Ok(decision);
1518 }
1519 return Err(error);
1520 }
1521 return Err(AdapterError::Transport("relay turn cancelled".into()));
1522 }
1523 let update = tokio::select! {
1524 update = host.next_update() => match update {
1525 Some(Ok(update)) => update,
1526 Some(Err(error)) => {
1527 let limited = report_relay_failure(
1528 &mut self.relay,
1529 &event_sink,
1530 *slot,
1531 true,
1532 error.to_string(),
1533 );
1534 if limited {
1535 self.relay.finish(*slot, *direct, false);
1536 }
1537 let _ = self.queue_session_metadata();
1538 if limited {
1539 return Ok(decision);
1540 }
1541 return Err(error);
1542 }
1543 None => {
1544 let error = AdapterError::Transport("adapter ended during turn".into());
1545 let limited = report_relay_failure(
1546 &mut self.relay,
1547 &event_sink,
1548 *slot,
1549 true,
1550 error.to_string(),
1551 );
1552 if limited {
1553 self.relay.finish(*slot, *direct, false);
1554 }
1555 let _ = self.queue_session_metadata();
1556 if limited {
1557 return Ok(decision);
1558 }
1559 return Err(error);
1560 }
1561 },
1562 _ = self.cancel_notify.notified() => {
1563 if !self.cancel_requested.swap(false, Ordering::AcqRel) {
1564 continue;
1565 }
1566 if let Err(error) = cancel_with_timeout(host).await {
1567 let limited = report_relay_failure(
1568 &mut self.relay,
1569 &event_sink,
1570 *slot,
1571 true,
1572 error.to_string(),
1573 );
1574 if limited {
1575 self.relay.finish(*slot, *direct, false);
1576 }
1577 let _ = self.queue_session_metadata();
1578 if limited {
1579 return Ok(decision);
1580 }
1581 return Err(error);
1582 }
1583 return Err(AdapterError::Transport("relay turn cancelled".into()));
1584 },
1585 permission = async {
1586 match permissions.as_mut() {
1587 Some(receiver) => receiver.recv().await,
1588 None => std::future::pending().await,
1589 }
1590 } => {
1591 let Some(permission) = permission else {
1592 permissions = None;
1593 continue;
1594 };
1595 if permission.slot != *slot {
1596 return Err(AdapterError::Transport(
1597 "permission response targets an inactive relay slot".into(),
1598 ));
1599 }
1600 host.answer_permission(permission.request_id, permission.answer).await?;
1601 continue;
1602 },
1603 };
1604 match &update.event {
1605 AgentEvent::Text { text, .. } => response.push_str(text),
1606 AgentEvent::TurnComplete { .. } => {
1607 let visible_response = response.replace(STOP_TOKEN, "");
1608 let visible_start = emitted_text.min(visible_response.len());
1609 let visible_start = floor_char_boundary(&visible_response, visible_start);
1610 if visible_start < visible_response.len()
1611 && let Some(sink) = &self.event_sink
1612 {
1613 sink(AgentEvent::Text {
1614 slot: *slot,
1615 text: visible_response[visible_start..].to_owned(),
1616 });
1617 }
1618 self.cancel_requested.store(false, Ordering::Release);
1619 break update.event.clone();
1620 }
1621 AgentEvent::Failed {
1622 started, detail, ..
1623 } => {
1624 let limited = report_relay_failure(
1625 &mut self.relay,
1626 &event_sink,
1627 *slot,
1628 *started,
1629 detail.clone(),
1630 );
1631 if limited {
1632 self.relay.finish(*slot, *direct, false);
1633 }
1634 let _ = self.queue_session_metadata();
1635 if limited {
1636 return Ok(decision);
1637 }
1638 return Err(AdapterError::Transport(detail.clone()));
1639 }
1640 _ => {}
1641 }
1642 if let AgentEvent::Text { .. } = &update.event {
1643 let visible_end = response
1646 .find(STOP_TOKEN)
1647 .unwrap_or_else(|| response.len().saturating_sub(STOP_TOKEN.len() - 1));
1648 let visible_end = floor_char_boundary(&response, visible_end);
1649 if emitted_text < visible_end {
1650 if let Some(sink) = &self.event_sink {
1651 sink(AgentEvent::Text {
1652 slot: *slot,
1653 text: response[emitted_text..visible_end].to_owned(),
1654 });
1655 }
1656 emitted_text = visible_end;
1657 }
1658 } else if let Some(sink) = &self.event_sink {
1659 sink(update.event.clone());
1660 }
1661 };
1662 let (response, requested_stop) = strip_stop_token(&response);
1663 let response = response.replace(STOP_TOKEN, "");
1664 let accepted_stop = requested_stop && *can_stop;
1665 let needs_stop_acknowledgment = accepted_stop && response.is_empty();
1666 let response = if needs_stop_acknowledgment {
1667 DEFAULT_STOP_ACKNOWLEDGMENT.to_owned()
1668 } else {
1669 response
1670 };
1671 if needs_stop_acknowledgment && let Some(sink) = &self.event_sink {
1676 sink(AgentEvent::Text {
1677 slot: *slot,
1678 text: response.clone(),
1679 });
1680 }
1681 if let Some(sink) = &self.event_sink {
1682 sink(completion_event);
1683 }
1684 if is_usage_limit_response(&response) {
1687 let detail = response.clone();
1688 let _ = self.relay.mark_limited(*slot);
1689 self.relay.finish(*slot, *direct, false);
1692 self.queue_session_metadata()?;
1693 if let Some(sink) = &self.event_sink {
1694 sink(AgentEvent::UsageLimitReached {
1695 slot: *slot,
1696 detail,
1697 });
1698 }
1699 return Ok(decision);
1700 }
1701 if !*direct && !response.is_empty() {
1702 self.relay
1703 .record_public(public_context_speaker(&speaker_name), response);
1704 }
1705 self.relay.mark_context_seen(*slot);
1706 self.relay.finish(*slot, *direct, accepted_stop);
1707 self.queue_session_metadata()?;
1708 Ok(decision)
1709 }
1710}
1711
1712fn report_relay_failure(
1713 relay: &mut Relay,
1714 event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
1715 slot: RosterSlot,
1716 started: bool,
1717 detail: String,
1718) -> bool {
1719 if is_usage_limit_response(&detail) {
1722 let _ = relay.mark_limited(slot);
1723 if let Some(sink) = event_sink {
1724 sink(AgentEvent::UsageLimitReached { slot, detail });
1725 }
1726 return true;
1727 }
1728 let _ = relay.tombstone(slot);
1729 if let Some(sink) = event_sink {
1730 sink(AgentEvent::Failed {
1731 slot,
1732 started,
1733 detail,
1734 });
1735 }
1736 false
1737}
1738
1739#[derive(Debug)]
1741pub struct ScriptedAdapter {
1742 slot: RosterSlot,
1743 capabilities: AgentCapabilities,
1744 events: VecDeque<AdapterResult<AgentEvent>>,
1745 prompts: Vec<String>,
1746}
1747
1748impl ScriptedAdapter {
1749 pub fn new(
1750 slot: RosterSlot,
1751 capabilities: AgentCapabilities,
1752 events: impl IntoIterator<Item = AgentEvent>,
1753 ) -> Self {
1754 Self {
1755 slot,
1756 capabilities,
1757 events: events.into_iter().map(Ok).collect(),
1758 prompts: Vec::new(),
1759 }
1760 }
1761
1762 pub fn prompts(&self) -> &[String] {
1763 &self.prompts
1764 }
1765}
1766
1767#[async_trait]
1768impl AgentAdapter for ScriptedAdapter {
1769 fn slot(&self) -> RosterSlot {
1770 self.slot
1771 }
1772
1773 fn capabilities(&self) -> AgentCapabilities {
1774 self.capabilities.clone()
1775 }
1776
1777 async fn start(&mut self) -> AdapterResult<()> {
1778 Ok(())
1779 }
1780
1781 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
1782 self.prompts.push(prompt);
1783 Ok(())
1784 }
1785
1786 async fn cancel(&mut self) -> AdapterResult<bool> {
1787 Ok(self.capabilities.supports_cancel)
1788 }
1789
1790 async fn answer_permission(
1791 &mut self,
1792 _request_id: String,
1793 _answer: PermissionAnswer,
1794 ) -> AdapterResult<()> {
1795 if self.capabilities.supports_permissions {
1796 Ok(())
1797 } else {
1798 Err(AdapterError::Unsupported("permission answer"))
1799 }
1800 }
1801
1802 async fn set_mode(&mut self, _mode: String) -> AdapterResult<()> {
1803 if self.capabilities.supports_modes {
1804 Ok(())
1805 } else {
1806 Err(AdapterError::Unsupported("set_mode"))
1807 }
1808 }
1809
1810 async fn reload(&mut self) -> AdapterResult<()> {
1811 Ok(())
1812 }
1813
1814 async fn stop(&mut self) -> AdapterResult<()> {
1815 Ok(())
1816 }
1817
1818 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
1819 self.events.pop_front()
1820 }
1821}
1822
1823#[derive(Debug)]
1826pub struct AgyAdapter {
1827 slot: RosterSlot,
1828 cwd: PathBuf,
1829 command: String,
1830 mode: String,
1831 mode_policy: String,
1832 session_id: Option<String>,
1833 child: Option<Child>,
1834 sender: mpsc::Sender<AdapterResult<AgentEvent>>,
1835 receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
1836 announced_session: Arc<Mutex<Option<String>>>,
1840 cancel_requested: Arc<AtomicBool>,
1841}
1842
1843impl AgyAdapter {
1844 pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
1845 let (sender, receiver) = mpsc::channel(256);
1846 Self {
1847 slot,
1848 cwd,
1849 command: command.into(),
1850 mode: "default".into(),
1851 mode_policy: "agy:full-access".into(),
1852 session_id: None,
1853 child: None,
1854 sender,
1855 receiver,
1856 announced_session: Arc::new(Mutex::new(None)),
1857 cancel_requested: Arc::new(AtomicBool::new(false)),
1858 }
1859 }
1860
1861 pub fn with_session_id(
1862 slot: RosterSlot,
1863 cwd: PathBuf,
1864 command: impl Into<String>,
1865 session_id: impl Into<String>,
1866 ) -> Self {
1867 let mut adapter = Self::new(slot, cwd, command);
1868 adapter.session_id = Some(session_id.into());
1869 adapter
1870 }
1871
1872 fn modes() -> Vec<Mode> {
1873 vec![
1874 Mode {
1875 id: "agy:full-access".into(),
1876 label: "Auto pilot".into(),
1877 },
1878 Mode {
1879 id: "agy:manual".into(),
1880 label: "Manual".into(),
1881 },
1882 Mode {
1883 id: "accept-edits".into(),
1884 label: "Accept Edits".into(),
1885 },
1886 Mode {
1887 id: "plan".into(),
1888 label: "Plan".into(),
1889 },
1890 ]
1891 }
1892
1893 async fn emit(&self, event: AdapterResult<AgentEvent>) {
1894 let _ = self.sender.send(event).await;
1895 }
1896}
1897
1898#[async_trait]
1899impl AgentAdapter for AgyAdapter {
1900 fn slot(&self) -> RosterSlot {
1901 self.slot
1902 }
1903
1904 fn session_id(&self) -> Option<String> {
1905 self.session_id.clone()
1906 }
1907
1908 fn protocol(&self) -> &'static str {
1909 "native"
1910 }
1911
1912 fn capabilities(&self) -> AgentCapabilities {
1913 AgentCapabilities {
1914 supports_cancel: true,
1915 supports_modes: true,
1916 supports_permissions: false,
1917 supports_terminals: true,
1918 supports_session_load: true,
1919 supports_models: false,
1920 }
1921 }
1922
1923 async fn start(&mut self) -> AdapterResult<()> {
1924 self.cancel_requested.store(false, Ordering::Release);
1925 self.emit(Ok(AgentEvent::Ready {
1926 slot: self.slot,
1927 capabilities: self.capabilities(),
1928 }))
1929 .await;
1930 self.emit(Ok(AgentEvent::ModesReplaced {
1931 slot: self.slot,
1932 modes: Self::modes(),
1933 current_mode: Some(self.mode_policy.clone()),
1934 }))
1935 .await;
1936 Ok(())
1937 }
1938
1939 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
1940 if self.child.is_some() {
1941 return Err(AdapterError::Transport(
1942 "agent is already handling a turn".into(),
1943 ));
1944 }
1945 if self.session_id.is_none()
1946 && let Ok(session) = self.announced_session.lock()
1947 {
1948 self.session_id = session.clone();
1949 }
1950 self.cancel_requested.store(false, Ordering::Release);
1951 let (program, args) = parse_command_line(&self.command)
1952 .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
1953 let mut command = Command::new(program);
1954 isolate_process_group(&mut command);
1955 command
1956 .args(args)
1957 .arg("--print")
1958 .arg(prompt)
1959 .arg("--print-timeout")
1960 .arg("60m")
1961 .arg("--output-format")
1962 .arg("stream-json")
1963 .current_dir(&self.cwd)
1964 .env("CODESWARM_CWD", &self.cwd)
1965 .stdout(Stdio::piped())
1966 .stderr(Stdio::piped());
1967 if let Some(session_id) = &self.session_id {
1968 command.arg("--conversation").arg(session_id);
1969 }
1970 if self.mode != "default" {
1971 command.arg("--mode").arg(&self.mode);
1972 }
1973 let mut child = command
1974 .spawn()
1975 .map_err(|error| AdapterError::Spawn(error.to_string()))?;
1976 let stdout = match child.stdout.take() {
1977 Some(stdout) => stdout,
1978 None => {
1979 let _ = terminate_child(&mut child).await;
1980 return Err(AdapterError::Transport("agent has no stdout".into()));
1981 }
1982 };
1983 let stderr = match child.stderr.take() {
1984 Some(stderr) => stderr,
1985 None => {
1986 let _ = terminate_child(&mut child).await;
1987 return Err(AdapterError::Transport("agent has no stderr".into()));
1988 }
1989 };
1990 let sender = self.sender.clone();
1991 let slot = self.slot;
1992 let announced_session = Arc::clone(&self.announced_session);
1993 let cancel_requested = Arc::clone(&self.cancel_requested);
1994 tokio::spawn(async move {
1995 let stderr_task = tokio::spawn(async move {
1996 const MAX_STDERR: usize = 32 * 1024;
1997 let mut stderr = BufReader::new(stderr);
1998 let mut bytes = Vec::new();
1999 let mut chunk = [0_u8; 4096];
2000 while let Ok(count) = stderr.read(&mut chunk).await {
2001 if count == 0 {
2002 break;
2003 }
2004 bytes.extend_from_slice(&chunk[..count]);
2005 if bytes.len() > MAX_STDERR {
2006 let keep_from = bytes.len() - MAX_STDERR;
2007 bytes.drain(..keep_from);
2008 }
2009 }
2010 String::from_utf8_lossy(&bytes).trim().to_owned()
2011 });
2012 let mut lines = BufReader::new(stdout).lines();
2013 let mut result: Option<Value> = None;
2014 let mut streamed_response = false;
2015 while let Ok(Some(line)) = lines.next_line().await {
2016 let value = match serde_json::from_str::<Value>(&line) {
2017 Ok(value) => value,
2018 Err(_) => {
2019 continue;
2024 }
2025 };
2026 if value.get("event").and_then(Value::as_str) == Some("init")
2027 && let Some(session_id) = value
2028 .get("conversation_id")
2029 .or_else(|| value.get("conversationId"))
2030 .and_then(Value::as_str)
2031 .filter(|id| !id.is_empty())
2032 && let Ok(mut announced) = announced_session.lock()
2033 {
2034 *announced = Some(session_id.to_owned());
2035 }
2036 if value.get("event").and_then(Value::as_str) == Some("result") {
2037 result = value.get("result").cloned();
2038 }
2039 match parse_agy_value(slot, &value) {
2040 Ok(Some(event)) => {
2041 if matches!(event, AgentEvent::Text { .. }) {
2042 streamed_response = true;
2043 }
2044 if sender.send(Ok(event)).await.is_err() {
2045 break;
2046 }
2047 }
2048 Ok(None) => {}
2049 Err(error) => {
2050 let _ = sender.send(Err(error)).await;
2051 }
2052 }
2053 }
2054 let stderr = stderr_task.await.ok().unwrap_or_default();
2055 let succeeded = cancel_requested.load(Ordering::Acquire)
2056 || result
2057 .as_ref()
2058 .and_then(|result| result.get("status"))
2059 .and_then(Value::as_str)
2060 == Some("SUCCESS");
2061 if succeeded {
2062 if !streamed_response
2067 && let Some(response) = result
2068 .as_ref()
2069 .and_then(|result| result.get("response"))
2070 .and_then(Value::as_str)
2071 .filter(|response| !response.is_empty())
2072 {
2073 let _ = sender
2074 .send(Ok(AgentEvent::Text {
2075 slot,
2076 text: response.to_owned(),
2077 }))
2078 .await;
2079 }
2080 let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
2081 } else {
2082 let detail = result
2083 .as_ref()
2084 .and_then(|result| result.get("error"))
2085 .and_then(Value::as_str)
2086 .filter(|detail| !detail.is_empty())
2087 .map(str::to_owned)
2088 .or_else(|| (!stderr.is_empty()).then_some(stderr))
2089 .unwrap_or_else(|| "native stream ended before a successful result".into());
2090 let _ = sender
2091 .send(Ok(AgentEvent::Failed {
2092 slot,
2093 started: true,
2094 detail,
2095 }))
2096 .await;
2097 }
2098 });
2099 self.child = Some(child);
2100 Ok(())
2101 }
2102
2103 async fn cancel(&mut self) -> AdapterResult<bool> {
2104 self.cancel_requested.store(true, Ordering::Release);
2105 let Some(mut child) = self.child.take() else {
2106 return Ok(false);
2107 };
2108 terminate_child(&mut child).await?;
2112 let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
2113 while let Some(event) = self.receiver.recv().await {
2114 if matches!(
2115 event,
2116 Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
2117 ) {
2118 break;
2119 }
2120 }
2121 })
2122 .await;
2123 Ok(true)
2124 }
2125
2126 async fn answer_permission(
2127 &mut self,
2128 _request_id: String,
2129 _answer: PermissionAnswer,
2130 ) -> AdapterResult<()> {
2131 Err(AdapterError::Unsupported("permission answer"))
2132 }
2133
2134 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
2135 let (mode, mode_policy) = match mode.as_str() {
2136 "full-access" | "codeswarm:mode:full-access" | "auto" | "autopilot" => {
2137 ("default".to_owned(), "agy:full-access".to_owned())
2138 }
2139 "codeswarm:mode:plan" | "readonly" | "plan" => ("plan".to_owned(), "plan".to_owned()),
2140 "codeswarm:mode:accept-edits" | "acceptedits" | "accept-edits" => {
2141 ("accept-edits".to_owned(), "accept-edits".to_owned())
2142 }
2143 "codeswarm:mode:manual" | "manual" | "ask" | "default" => {
2144 ("default".to_owned(), "agy:manual".to_owned())
2145 }
2146 "agy:full-access" => ("default".to_owned(), "agy:full-access".to_owned()),
2147 "agy:manual" => ("default".to_owned(), "agy:manual".to_owned()),
2148 _ => return Err(AdapterError::Unsupported("requested Agy mode")),
2149 };
2150 self.mode = mode;
2151 self.mode_policy = mode_policy.clone();
2152 self.emit(Ok(AgentEvent::ModesReplaced {
2153 slot: self.slot,
2154 modes: Self::modes(),
2155 current_mode: Some(mode_policy),
2156 }))
2157 .await;
2158 Ok(())
2159 }
2160
2161 async fn reload(&mut self) -> AdapterResult<()> {
2162 self.stop().await?;
2163 self.start().await
2164 }
2165
2166 async fn stop(&mut self) -> AdapterResult<()> {
2167 let _ = self.cancel().await?;
2168 Ok(())
2169 }
2170
2171 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
2172 let event = self.receiver.recv().await;
2173 if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2174 && self.session_id.is_none()
2175 && let Ok(session) = self.announced_session.lock()
2176 {
2177 self.session_id = session.clone();
2178 }
2179 if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2180 && let Some(mut child) = self.child.take()
2181 {
2182 let _ = child.wait().await;
2183 }
2184 event
2185 }
2186}
2187
2188#[cfg(test)]
2189#[cfg_attr(not(test), allow(dead_code))]
2190fn parse_agy_line(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
2191 let value: Value =
2192 serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
2193 parse_agy_value(slot, &value)
2194}
2195
2196fn parse_agy_value(slot: RosterSlot, value: &Value) -> AdapterResult<Option<AgentEvent>> {
2197 let event = value.get("event").and_then(Value::as_str);
2198 if let Some(terminal) = parse_terminal_event(value, event) {
2199 return Ok(Some(AgentEvent::Terminal {
2200 slot,
2201 event: terminal,
2202 }));
2203 }
2204 match event {
2205 Some("step_update") => {
2206 let Some(update) = value.get("step_update") else {
2207 return Ok(None);
2208 };
2209 let is_response = update
2210 .get("step_type")
2211 .and_then(Value::as_str)
2212 .is_some_and(|kind| kind == "agent_response");
2213 let text = update.get("text_delta").and_then(Value::as_str);
2214 let response = is_response
2215 .then(|| text.map(str::to_owned))
2216 .flatten()
2217 .filter(|text| !text.is_empty())
2218 .map(|text| AgentEvent::Text { slot, text });
2219 Ok(response.or_else(|| parse_agy_tool(slot, value)))
2220 }
2221 _ => Ok(None),
2222 }
2223}
2224
2225fn parse_agy_tool(slot: RosterSlot, value: &Value) -> Option<AgentEvent> {
2226 let update = value.get("step_update")?;
2227 if update.get("step_type")?.as_str()? != "tool" {
2228 return None;
2229 }
2230 let step_index = update.get("step_index")?.as_i64()?;
2231 let title = update
2232 .get("tool_name")
2233 .and_then(Value::as_str)
2234 .unwrap_or("Tool call")
2235 .replace('_', " ");
2236 let status = match update.get("state").and_then(Value::as_str) {
2237 Some("DONE") => ToolStatus::Completed,
2238 Some("FAILED") => ToolStatus::Failed,
2239 Some("ACTIVE") => ToolStatus::Running,
2240 _ => ToolStatus::Pending,
2241 };
2242 let detail = update
2243 .get("tool_info")
2244 .and_then(|info| info.get("output"))
2245 .and_then(Value::as_str)
2246 .map(str::to_owned);
2247 Some(AgentEvent::Tool {
2248 slot,
2249 update: ToolUpdate {
2250 id: format!("agy-tool-{step_index}"),
2251 title,
2252 status,
2253 detail,
2254 },
2255 })
2256}
2257
2258#[derive(Debug)]
2261pub struct AcpAdapter {
2262 slot: RosterSlot,
2263 program: String,
2264 args: Vec<String>,
2265 cwd: PathBuf,
2266 child: Option<Child>,
2267 reader: Option<BufReader<ChildStdout>>,
2268 capabilities: AgentCapabilities,
2269 modes: Vec<Mode>,
2270 models: Vec<Mode>,
2271 model_config_id: Option<String>,
2272 session_id: Option<String>,
2273 next_request_id: u64,
2274 prompt_request_id: Option<u64>,
2275 queued_events: VecDeque<AdapterResult<AgentEvent>>,
2276 stderr_task: Option<tokio::task::JoinHandle<String>>,
2277 terminals: BTreeMap<String, TerminalProcess>,
2278 next_terminal_id: u64,
2279}
2280
2281impl AcpAdapter {
2282 pub fn new(
2283 slot: RosterSlot,
2284 cwd: PathBuf,
2285 program: impl Into<String>,
2286 args: Vec<String>,
2287 ) -> Self {
2288 Self {
2289 slot,
2290 program: program.into(),
2291 args,
2292 cwd,
2293 child: None,
2294 reader: None,
2295 capabilities: AgentCapabilities::default(),
2296 modes: Vec::new(),
2297 models: Vec::new(),
2298 model_config_id: None,
2299 session_id: None,
2300 next_request_id: 1,
2301 prompt_request_id: None,
2302 queued_events: VecDeque::new(),
2303 stderr_task: None,
2304 terminals: BTreeMap::new(),
2305 next_terminal_id: 1,
2306 }
2307 }
2308
2309 pub fn with_session_id(
2310 slot: RosterSlot,
2311 cwd: PathBuf,
2312 program: impl Into<String>,
2313 args: Vec<String>,
2314 session_id: impl Into<String>,
2315 ) -> Self {
2316 let mut adapter = Self::new(slot, cwd, program, args);
2317 adapter.session_id = Some(session_id.into());
2318 adapter
2319 }
2320
2321 async fn request(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2322 self.request_with_timeout(method, params, std::time::Duration::from_secs(30))
2323 .await
2324 }
2325
2326 async fn request_with_timeout(
2327 &mut self,
2328 method: &str,
2329 params: Value,
2330 deadline: std::time::Duration,
2331 ) -> AdapterResult<Value> {
2332 tokio::time::timeout(deadline, self.request_inner(method, params))
2333 .await
2334 .map_err(|_| {
2335 AdapterError::Transport(format!(
2336 "ACP {method} timed out; reload the agent to retry"
2337 ))
2338 })?
2339 }
2340
2341 async fn request_inner(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2342 let request_id = self.next_request_id;
2343 self.next_request_id += 1;
2344 self.write_json(serde_json::json!({
2345 "jsonrpc": "2.0",
2346 "id": request_id,
2347 "method": method,
2348 "params": params,
2349 }))
2350 .await?;
2351 loop {
2352 let line = self.read_line().await?;
2353 let value: Value = match serde_json::from_str(&line) {
2354 Ok(value) => value,
2355 Err(_) => {
2356 continue;
2360 }
2361 };
2362 if self.reject_empty_permission_request(&value).await? {
2363 continue;
2364 }
2365 if self.handle_client_request(&value).await? {
2366 continue;
2367 }
2368 if value
2369 .get("id")
2370 .is_some_and(|id| rpc_id_to_string(id) == request_id.to_string())
2371 {
2372 if let Some(error) = value.get("error") {
2373 return Err(AdapterError::Protocol(error.to_string()));
2374 }
2375 return value
2376 .get("result")
2377 .cloned()
2378 .ok_or_else(|| AdapterError::Protocol("response has no result".into()));
2379 }
2380 if let Some(event) = parse_acp_notification(self.slot, &line)? {
2381 self.queued_events.push_back(Ok(event));
2382 }
2383 }
2384 }
2385
2386 async fn write_json(&mut self, value: Value) -> AdapterResult<()> {
2387 let child = self
2388 .child
2389 .as_mut()
2390 .ok_or_else(|| AdapterError::Transport("ACP agent is not running".into()))?;
2391 let stdin = child
2392 .stdin
2393 .as_mut()
2394 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdin".into()))?;
2395 stdin
2396 .write_all(value.to_string().as_bytes())
2397 .await
2398 .map_err(|error| AdapterError::Transport(error.to_string()))?;
2399 stdin
2400 .write_all(b"\n")
2401 .await
2402 .map_err(|error| AdapterError::Transport(error.to_string()))
2403 }
2404
2405 async fn reject_empty_permission_request(&mut self, value: &Value) -> AdapterResult<bool> {
2410 if value.get("method").and_then(Value::as_str) != Some("session/request_permission")
2411 || value.get("id").is_none()
2412 {
2413 return Ok(false);
2414 }
2415 let valid = value
2416 .get("params")
2417 .and_then(|params| params.get("options"))
2418 .and_then(Value::as_array)
2419 .is_some_and(|options| !options.is_empty());
2420 if valid {
2421 return Ok(false);
2422 }
2423 self.write_json(serde_json::json!({
2424 "jsonrpc": "2.0",
2425 "id": value.get("id").cloned().unwrap_or(Value::Null),
2426 "error": {
2427 "code": -32602,
2428 "message": "Permission request requires at least one option",
2429 },
2430 }))
2431 .await?;
2432 Ok(true)
2433 }
2434
2435 fn workspace_path(&self, path: &str) -> Result<PathBuf, String> {
2436 let root = self
2437 .cwd
2438 .canonicalize()
2439 .map_err(|error| format!("unable to resolve workspace: {error}"))?;
2440 let requested = Path::new(path);
2441 let candidate = if requested.is_absolute() {
2442 requested.to_path_buf()
2443 } else {
2444 root.join(requested)
2445 };
2446 let resolved = if !candidate.exists() {
2447 let parent = candidate
2448 .parent()
2449 .ok_or_else(|| "file path has no parent".to_owned())?
2450 .canonicalize()
2451 .map_err(|error| format!("unable to resolve parent directory: {error}"))?;
2452 parent.join(
2453 candidate
2454 .file_name()
2455 .ok_or_else(|| "file path has no filename".to_owned())?,
2456 )
2457 } else {
2458 candidate
2459 .canonicalize()
2460 .map_err(|error| format!("unable to resolve file path: {error}"))?
2461 };
2462 if !resolved.starts_with(&root) {
2463 return Err("file path is outside the project".into());
2464 }
2465 Ok(resolved)
2466 }
2467
2468 fn read_workspace_text(
2469 &self,
2470 path: &str,
2471 line: Option<i64>,
2472 limit: Option<i64>,
2473 ) -> Result<String, String> {
2474 if line.is_some_and(|line| line < 1) {
2475 return Err("line must be positive".into());
2476 }
2477 if limit.is_some_and(|limit| limit < 0) {
2478 return Err("limit must not be negative".into());
2479 }
2480 let path = self.workspace_path(path)?;
2481 let mut bytes = Vec::new();
2482 let mut source = match std::fs::File::open(path) {
2483 Ok(source) => source,
2484 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(String::new()),
2485 Err(error) => return Err(error.to_string()),
2486 };
2487 source
2488 .by_ref()
2489 .take((MAX_FILE_READ_BYTES as u64).saturating_add(1))
2490 .read_to_end(&mut bytes)
2491 .map_err(|error| error.to_string())?;
2492 bytes.truncate(MAX_FILE_READ_BYTES);
2493 let text = String::from_utf8_lossy(&bytes);
2494 if line.is_none() && limit.is_none() {
2495 return Ok(text.into_owned());
2496 }
2497 let start = line.map_or(0, |line| line as usize - 1);
2498 let limit = limit.unwrap_or(i64::MAX) as usize;
2499 let selected = text
2500 .split_inclusive('\n')
2501 .skip(start)
2502 .take(limit)
2503 .collect::<String>();
2504 if line.is_some() {
2505 Ok(selected.trim_end_matches('\n').to_owned())
2506 } else {
2507 Ok(selected)
2508 }
2509 }
2510
2511 fn write_workspace_text(&self, params: &Value) -> Result<(), String> {
2512 let path = params
2513 .get("path")
2514 .and_then(Value::as_str)
2515 .filter(|path| !path.is_empty())
2516 .ok_or("path must be a non-empty string")?;
2517 let content = params
2518 .get("content")
2519 .and_then(Value::as_str)
2520 .ok_or("content must be a string")?;
2521 let path = self.workspace_path(path)?;
2522 std::fs::write(path, content).map_err(|error| error.to_string())
2523 }
2524
2525 async fn terminal_create(&mut self, params: &Value) -> Result<Value, String> {
2526 let command = params
2527 .get("command")
2528 .and_then(Value::as_str)
2529 .filter(|command| !command.trim().is_empty())
2530 .ok_or_else(|| "terminal command is required".to_owned())?;
2531 let cwd = params.get("cwd").and_then(Value::as_str).unwrap_or(".");
2532 let cwd = self.workspace_path(cwd)?;
2533 if !cwd.is_dir() {
2534 return Err("terminal cwd is not a directory".into());
2535 }
2536 let mut process = Command::new(command);
2537 isolate_process_group(&mut process);
2538 if let Some(args) = params.get("args").and_then(Value::as_array) {
2539 process.args(args.iter().filter_map(Value::as_str));
2540 }
2541 process
2542 .current_dir(&cwd)
2543 .stdin(Stdio::null())
2544 .stdout(Stdio::piped())
2545 .stderr(Stdio::piped());
2546 if let Some(env) = params.get("env") {
2547 if let Some(entries) = env.as_array() {
2548 for entry in entries {
2549 if let (Some(name), Some(value)) = (
2550 entry.get("name").and_then(Value::as_str),
2551 entry.get("value").and_then(Value::as_str),
2552 ) {
2553 process.env(name, value);
2554 }
2555 }
2556 } else if let Some(entries) = env.as_object() {
2557 for (name, value) in entries {
2558 if let Some(value) = value.as_str() {
2559 process.env(name, value);
2560 }
2561 }
2562 }
2563 }
2564 let mut child = process.spawn().map_err(|error| error.to_string())?;
2565 let stdout = child.stdout.take();
2566 let stderr = child.stderr.take();
2567 let (Some(stdout), Some(stderr)) = (stdout, stderr) else {
2568 let _ = terminate_child(&mut child).await;
2569 return Err("terminal has no output pipes".into());
2570 };
2571 let output = Arc::new(Mutex::new(Vec::new()));
2572 let truncated = Arc::new(AtomicBool::new(false));
2573 let output_readers = Arc::new(AtomicUsize::new(2));
2574 let output_limit = params
2575 .get("outputByteLimit")
2576 .and_then(Value::as_u64)
2577 .map_or(MAX_TERMINAL_OUTPUT_BYTES, |limit| {
2578 usize::try_from(limit)
2579 .unwrap_or(MAX_TERMINAL_OUTPUT_BYTES)
2580 .min(MAX_TERMINAL_OUTPUT_BYTES)
2581 });
2582 tokio::spawn(drain_terminal_output(
2583 stdout,
2584 Arc::clone(&output),
2585 Arc::clone(&truncated),
2586 Arc::clone(&output_readers),
2587 output_limit,
2588 ));
2589 tokio::spawn(drain_terminal_output(
2590 stderr,
2591 Arc::clone(&output),
2592 Arc::clone(&truncated),
2593 Arc::clone(&output_readers),
2594 output_limit,
2595 ));
2596 let id = format!("terminal-{}", self.next_terminal_id);
2597 self.next_terminal_id = self.next_terminal_id.saturating_add(1);
2598 let state = TerminalProcess {
2599 child: Arc::new(AsyncMutex::new(Some(child))),
2600 output,
2601 truncated,
2602 output_readers,
2603 };
2604 self.terminals.insert(id.clone(), state);
2605 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2606 slot: self.slot,
2607 event: TerminalEvent::Created {
2608 id: id.clone(),
2609 command: std::iter::once(command)
2610 .chain(
2611 params
2612 .get("args")
2613 .and_then(Value::as_array)
2614 .into_iter()
2615 .flatten()
2616 .filter_map(Value::as_str),
2617 )
2618 .collect::<Vec<_>>()
2619 .join(" "),
2620 },
2621 }));
2622 Ok(serde_json::json!({"terminalId": id}))
2623 }
2624
2625 async fn terminal_output(&mut self, id: &str) -> Result<Value, String> {
2626 let terminal = self
2627 .terminals
2628 .get(id)
2629 .ok_or_else(|| "terminal not found".to_owned())?
2630 .clone();
2631 let output = terminal
2632 .output
2633 .lock()
2634 .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
2635 .unwrap_or_default();
2636 let exit_code = terminal.exit_code().await;
2637 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2638 slot: self.slot,
2639 event: TerminalEvent::Output {
2640 id: id.to_owned(),
2641 text: output.clone(),
2642 },
2643 }));
2644 let mut response = serde_json::json!({
2645 "output": output,
2646 "truncated": terminal.truncated.load(Ordering::Acquire),
2647 });
2648 if let Some(code) = exit_code {
2649 response["exitStatus"] = serde_json::json!({"exitCode": code});
2650 }
2651 Ok(response)
2652 }
2653
2654 async fn terminal_wait(&mut self, id: &str) -> Result<Value, String> {
2655 let terminal = self
2656 .terminals
2657 .get(id)
2658 .ok_or_else(|| "terminal not found".to_owned())?
2659 .clone();
2660 let exit_code = terminal.wait().await;
2661 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2662 slot: self.slot,
2663 event: TerminalEvent::Exited {
2664 id: id.to_owned(),
2665 code: exit_code.unwrap_or(-1),
2666 },
2667 }));
2668 Ok(serde_json::json!({"exitCode": exit_code, "signal": Value::Null}))
2669 }
2670
2671 async fn handle_client_request(&mut self, value: &Value) -> AdapterResult<bool> {
2675 let Some(method) = value.get("method").and_then(Value::as_str) else {
2676 return Ok(false);
2677 };
2678 let Some(id) = value.get("id").cloned() else {
2679 return Ok(false);
2680 };
2681 if method == "session/request_permission" {
2684 return Ok(false);
2685 }
2686 let params = value.get("params").cloned().unwrap_or(Value::Null);
2687 let response = match method {
2688 "fs/read_text_file" => {
2689 let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2690 let line = params.get("line").and_then(Value::as_i64);
2691 let limit = params.get("limit").and_then(Value::as_i64);
2692 match self.read_workspace_text(path, line, limit) {
2693 Ok(content) => serde_json::json!({
2694 "jsonrpc": "2.0",
2695 "id": id,
2696 "result": {"content": content},
2697 }),
2698 Err(message) => serde_json::json!({
2699 "jsonrpc": "2.0",
2700 "id": id,
2701 "error": {"code": -32602, "message": message},
2702 }),
2703 }
2704 }
2705 "fs/write_text_file" => {
2706 let result = self.write_workspace_text(¶ms);
2707 match result {
2708 Ok(()) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}}),
2709 Err(message) => serde_json::json!({
2710 "jsonrpc": "2.0",
2711 "id": id,
2712 "error": {"code": -32602, "message": message},
2713 }),
2714 }
2715 }
2716 "terminal/create" => match self.terminal_create(¶ms).await {
2717 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2718 Err(message) => serde_json::json!({
2719 "jsonrpc": "2.0",
2720 "id": id,
2721 "error": {"code": -32602, "message": message},
2722 }),
2723 },
2724 "terminal/output" => {
2725 let terminal_id = params
2726 .get("terminalId")
2727 .and_then(Value::as_str)
2728 .unwrap_or("");
2729 match self.terminal_output(terminal_id).await {
2730 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2731 Err(message) => serde_json::json!({
2732 "jsonrpc": "2.0",
2733 "id": id,
2734 "error": {"code": -32602, "message": message},
2735 }),
2736 }
2737 }
2738 "terminal/wait_for_exit" => {
2739 let terminal_id = params
2740 .get("terminalId")
2741 .and_then(Value::as_str)
2742 .unwrap_or("");
2743 match self.terminal_wait(terminal_id).await {
2744 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2745 Err(message) => serde_json::json!({
2746 "jsonrpc": "2.0",
2747 "id": id,
2748 "error": {"code": -32602, "message": message},
2749 }),
2750 }
2751 }
2752 "terminal/kill" => {
2753 let terminal_id = params
2754 .get("terminalId")
2755 .and_then(Value::as_str)
2756 .unwrap_or("");
2757 if let Some(terminal) = self.terminals.get(terminal_id) {
2758 terminal.kill().await;
2759 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2760 } else {
2761 serde_json::json!({
2762 "jsonrpc": "2.0",
2763 "id": id,
2764 "error": {"code": -32602, "message": "terminal not found"},
2765 })
2766 }
2767 }
2768 "terminal/release" => {
2769 let terminal_id = params
2770 .get("terminalId")
2771 .and_then(Value::as_str)
2772 .unwrap_or("");
2773 if let Some(terminal) = self.terminals.remove(terminal_id) {
2774 terminal.stop().await;
2775 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2776 slot: self.slot,
2777 event: TerminalEvent::Released {
2778 id: terminal_id.to_owned(),
2779 },
2780 }));
2781 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2782 } else {
2783 serde_json::json!({
2784 "jsonrpc": "2.0",
2785 "id": id,
2786 "error": {"code": -32602, "message": "terminal not found"},
2787 })
2788 }
2789 }
2790 _ => serde_json::json!({
2791 "jsonrpc": "2.0",
2792 "id": id,
2793 "error": {"code": -32601, "message": format!("unsupported client method: {method}")},
2794 }),
2795 };
2796 self.write_json(response).await?;
2797 Ok(true)
2798 }
2799
2800 async fn read_line(&mut self) -> AdapterResult<String> {
2801 let reader = self
2802 .reader
2803 .as_mut()
2804 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdout".into()))?;
2805 read_bounded_line(reader).await
2806 }
2807
2808 async fn start(&mut self) -> AdapterResult<()> {
2809 if self.child.is_some() {
2811 self.stop().await?;
2812 }
2813 let mut command = Command::new(&self.program);
2814 isolate_process_group(&mut command);
2815 command
2816 .args(&self.args)
2817 .current_dir(&self.cwd)
2818 .stdin(Stdio::piped())
2819 .stdout(Stdio::piped())
2820 .stderr(Stdio::piped())
2821 .env("CODESWARM_CWD", &self.cwd);
2822 if self.program.to_ascii_lowercase().contains("gemini")
2823 || self
2824 .args
2825 .iter()
2826 .any(|arg| arg.to_ascii_lowercase().contains("gemini"))
2827 {
2828 command.env("GEMINI_TELEMETRY_ENABLED", "false");
2829 }
2830 let mut child = command
2831 .spawn()
2832 .map_err(|error| AdapterError::Spawn(error.to_string()))?;
2833 let stdout = match child.stdout.take() {
2834 Some(stdout) => stdout,
2835 None => {
2836 let _ = terminate_child(&mut child).await;
2837 return Err(AdapterError::Transport("ACP agent has no stdout".into()));
2838 }
2839 };
2840 let stderr = match child.stderr.take() {
2841 Some(stderr) => stderr,
2842 None => {
2843 let _ = terminate_child(&mut child).await;
2844 return Err(AdapterError::Transport("ACP agent has no stderr".into()));
2845 }
2846 };
2847 self.child = Some(child);
2848 self.reader = Some(BufReader::new(stdout));
2849 self.stderr_task = Some(tokio::spawn(drain_bounded(stderr, 32 * 1024)));
2850
2851 let initialize = match self
2852 .request(
2853 "initialize",
2854 serde_json::json!({
2855 "protocolVersion": 1,
2856 "clientCapabilities": {
2857 "fs": {"readTextFile": true, "writeTextFile": true},
2858 "terminal": true,
2859 },
2860 "clientInfo": {
2861 "name": "CodeSwarm",
2862 "title": "CodeSwarm",
2863 "version": env!("CARGO_PKG_VERSION"),
2864 },
2865 }),
2866 )
2867 .await
2868 {
2869 Ok(value) => value,
2870 Err(error) => {
2871 let _ = self.stop().await;
2872 return Err(error);
2873 }
2874 };
2875 let agent_capabilities = initialize
2876 .get("agentCapabilities")
2877 .cloned()
2878 .unwrap_or(Value::Null);
2879 self.capabilities = AgentCapabilities {
2880 supports_cancel: true,
2881 supports_modes: true,
2882 supports_permissions: true,
2883 supports_terminals: true,
2884 supports_session_load: agent_capabilities
2885 .get("loadSession")
2886 .and_then(Value::as_bool)
2887 .unwrap_or(false),
2888 supports_models: false,
2889 };
2890 let session = if let Some(session_id) = self.session_id.clone() {
2891 if !self.capabilities.supports_session_load {
2892 let _ = self.stop().await;
2893 return Err(AdapterError::Unsupported("session/load"));
2894 }
2895 match self
2896 .request(
2897 "session/load",
2898 serde_json::json!({
2899 "cwd": self.cwd,
2900 "mcpServers": [],
2901 "sessionId": session_id,
2902 }),
2903 )
2904 .await
2905 {
2906 Ok(value) => value,
2907 Err(error) => {
2908 let _ = self.stop().await;
2909 return Err(error);
2910 }
2911 }
2912 } else {
2913 let session = match self
2914 .request(
2915 "session/new",
2916 serde_json::json!({"cwd": self.cwd, "mcpServers": []}),
2917 )
2918 .await
2919 {
2920 Ok(value) => value,
2921 Err(error) => {
2922 let _ = self.stop().await;
2923 return Err(error);
2924 }
2925 };
2926 self.session_id = session
2927 .get("sessionId")
2928 .and_then(Value::as_str)
2929 .map(str::to_owned);
2930 if self.session_id.is_none() {
2931 let _ = self.stop().await;
2932 return Err(AdapterError::Protocol(
2933 "session/new returned no sessionId".into(),
2934 ));
2935 }
2936 session
2937 };
2938 self.capabilities.supports_modes = false;
2939 if let Some(modes) = session.get("modes") {
2940 let available = modes
2941 .get("availableModes")
2942 .and_then(Value::as_array)
2943 .map(|modes| {
2944 modes
2945 .iter()
2946 .filter_map(|mode| {
2947 Some(Mode {
2948 id: mode.get("id")?.as_str()?.to_owned(),
2949 label: mode.get("name")?.as_str()?.to_owned(),
2950 })
2951 })
2952 .collect::<Vec<_>>()
2953 })
2954 .unwrap_or_default();
2955 self.modes = available.clone();
2956 self.capabilities.supports_modes = !available.is_empty();
2957 self.queued_events.push_back(Ok(AgentEvent::ModesReplaced {
2958 slot: self.slot,
2959 modes: available,
2960 current_mode: modes
2961 .get("currentModeId")
2962 .and_then(Value::as_str)
2963 .map(str::to_owned),
2964 }));
2965 }
2966 self.models.clear();
2967 self.model_config_id = None;
2968 let current_model =
2969 parse_model_config(&session).and_then(|(config_id, models, current)| {
2970 self.model_config_id = Some(config_id);
2971 self.models = models;
2972 current
2973 });
2974 self.capabilities.supports_models =
2975 self.model_config_id.is_some() && !self.models.is_empty();
2976 if let Some(config_id) = self.model_config_id.clone()
2977 && !self.models.is_empty()
2978 {
2979 self.queued_events.push_back(Ok(AgentEvent::ModelsReplaced {
2980 slot: self.slot,
2981 config_id,
2982 models: self.models.clone(),
2983 current_model,
2984 }));
2985 }
2986 self.queued_events.push_back(Ok(AgentEvent::Ready {
2987 slot: self.slot,
2988 capabilities: self.capabilities(),
2989 }));
2990 Ok(())
2991 }
2992}
2993
2994fn prompt_resource_paths(prompt: &str) -> Vec<String> {
2995 let characters = prompt.chars().collect::<Vec<_>>();
2996 let mut paths = Vec::new();
2997 let mut index = 0;
2998 while index < characters.len() {
2999 if characters[index] != '@' {
3000 index += 1;
3001 continue;
3002 }
3003 index += 1;
3004 let quoted = characters.get(index) == Some(&'"');
3005 if quoted {
3006 index += 1;
3007 }
3008 let start = index;
3009 while index < characters.len()
3010 && if quoted {
3011 characters[index] != '"'
3012 } else {
3013 !characters[index].is_whitespace()
3014 }
3015 {
3016 index += 1;
3017 }
3018 if index > start {
3019 paths.push(characters[start..index].iter().collect());
3020 }
3021 if quoted && index < characters.len() {
3022 index += 1;
3023 }
3024 }
3025 paths
3026}
3027
3028fn prompt_content_blocks(cwd: &Path, prompt: &str) -> Vec<Value> {
3029 let mut blocks = vec![serde_json::json!({"type": "text", "text": prompt})];
3030 for path in prompt_resource_paths(prompt) {
3031 if path.ends_with('/') {
3032 continue;
3033 }
3034 let Ok(resource) = resources::load(cwd, &path) else {
3035 continue;
3036 };
3037 let uri = format!("file://{}", resource.path.display());
3038 let resource_value = if let Some(text) = resource.text {
3039 serde_json::json!({
3040 "uri": uri,
3041 "text": text,
3042 "mimeType": resource.mime_type,
3043 })
3044 } else if let Some(data) = resource.data {
3045 serde_json::json!({
3046 "uri": uri,
3047 "blob": BASE64.encode(data),
3048 "mimeType": resource.mime_type,
3049 })
3050 } else {
3051 continue;
3052 };
3053 blocks.push(serde_json::json!({
3054 "type": "resource",
3055 "resource": resource_value,
3056 }));
3057 }
3058 blocks
3059}
3060
3061#[async_trait]
3062impl AgentAdapter for AcpAdapter {
3063 fn slot(&self) -> RosterSlot {
3064 self.slot
3065 }
3066
3067 fn session_id(&self) -> Option<String> {
3068 self.session_id.clone()
3069 }
3070
3071 fn protocol(&self) -> &'static str {
3072 "acp"
3073 }
3074
3075 fn capabilities(&self) -> AgentCapabilities {
3076 self.capabilities.clone()
3077 }
3078
3079 async fn start(&mut self) -> AdapterResult<()> {
3080 AcpAdapter::start(self).await
3083 }
3084
3085 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
3086 let session_id = self
3087 .session_id
3088 .as_ref()
3089 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3090 let request_id = self.next_request_id;
3091 self.next_request_id += 1;
3092 let prompt_blocks = prompt_content_blocks(&self.cwd, &prompt);
3093 self.write_json(serde_json::json!({
3094 "jsonrpc": "2.0",
3095 "id": request_id,
3096 "method": "session/prompt",
3097 "params": {
3098 "sessionId": session_id,
3099 "prompt": prompt_blocks,
3100 },
3101 }))
3102 .await?;
3103 self.prompt_request_id = Some(request_id);
3104 Ok(())
3105 }
3106
3107 async fn cancel(&mut self) -> AdapterResult<bool> {
3108 let Some(session_id) = &self.session_id else {
3109 return Ok(false);
3110 };
3111 self.write_json(serde_json::json!({
3112 "jsonrpc": "2.0",
3113 "method": "session/cancel",
3114 "params": {"sessionId": session_id, "_meta": {}},
3115 }))
3116 .await?;
3117 let settled = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
3118 loop {
3119 match <Self as AgentAdapter>::next_event(self).await {
3120 Some(Ok(AgentEvent::TurnComplete { .. })) | None => break,
3121 Some(Ok(_)) => {}
3122 Some(Err(_)) => break,
3123 }
3124 }
3125 })
3126 .await
3127 .is_ok();
3128 if !settled {
3129 self.reload().await?;
3133 }
3134 Ok(true)
3135 }
3136
3137 async fn answer_permission(
3138 &mut self,
3139 request_id: String,
3140 answer: PermissionAnswer,
3141 ) -> AdapterResult<()> {
3142 let id = request_id
3143 .parse::<u64>()
3144 .map(Value::from)
3145 .unwrap_or_else(|_| Value::String(request_id));
3146 let outcome = match answer {
3147 PermissionAnswer::Selected { option_id } => {
3148 serde_json::json!({"outcome": "selected", "optionId": option_id})
3149 }
3150 PermissionAnswer::Cancelled => serde_json::json!({"outcome": "cancelled"}),
3151 };
3152 self.write_json(serde_json::json!({
3153 "jsonrpc": "2.0",
3154 "id": id,
3155 "result": {"outcome": outcome},
3159 }))
3160 .await
3161 }
3162
3163 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
3164 let session_id = self
3165 .session_id
3166 .as_ref()
3167 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3168 let policy = match mode.as_str() {
3169 "plan" => "codeswarm:mode:plan",
3170 "default" | "manual" => "codeswarm:mode:manual",
3171 "accept-edits" => "codeswarm:mode:accept-edits",
3172 "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
3173 other => other,
3174 };
3175 let native_mode = crate::policy::resolve(policy, &self.modes)
3176 .map(|mode| mode.id)
3177 .unwrap_or(mode);
3178 let _ = self
3179 .request(
3180 "session/set_mode",
3181 serde_json::json!({"sessionId": session_id, "modeId": native_mode.clone()}),
3182 )
3183 .await?;
3184 self.queued_events.push_back(Ok(AgentEvent::ModeUpdated {
3185 slot: self.slot,
3186 current_mode: native_mode,
3187 }));
3188 Ok(())
3189 }
3190
3191 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
3192 let session_id = self
3193 .session_id
3194 .clone()
3195 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3196 let config_id = self
3197 .model_config_id
3198 .clone()
3199 .ok_or(AdapterError::Unsupported("set_model"))?;
3200 if !self.models.iter().any(|candidate| candidate.id == model) {
3201 return Err(AdapterError::Protocol(
3202 "model is not advertised by the agent".into(),
3203 ));
3204 }
3205 let _ = self
3206 .request(
3207 "session/set_config_option",
3208 serde_json::json!({
3209 "sessionId": session_id,
3210 "configId": config_id,
3211 "value": model,
3212 }),
3213 )
3214 .await?;
3215 Ok(())
3216 }
3217
3218 async fn reload(&mut self) -> AdapterResult<()> {
3219 let session_id = self
3224 .capabilities
3225 .supports_session_load
3226 .then(|| self.session_id.clone())
3227 .flatten();
3228 self.stop().await?;
3229 self.session_id = session_id.clone();
3230 let result = self.start().await;
3231 if result.is_err() {
3232 self.session_id = session_id;
3236 }
3237 result
3238 }
3239
3240 async fn stop(&mut self) -> AdapterResult<()> {
3241 let terminals = std::mem::take(&mut self.terminals);
3242 for terminal in terminals.values() {
3243 terminal.stop().await;
3244 }
3245 if let Some(mut child) = self.child.take() {
3246 terminate_child(&mut child).await?;
3247 }
3248 self.reader = None;
3249 self.session_id = None;
3250 self.prompt_request_id = None;
3251 if let Some(task) = self.stderr_task.take() {
3252 task.abort();
3253 let _ = task.await;
3254 }
3255 Ok(())
3256 }
3257
3258 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
3259 if let Some(event) = self.queued_events.pop_front() {
3260 return Some(event);
3261 }
3262 loop {
3263 let line = match self.read_line().await {
3264 Ok(line) => line,
3265 Err(error) => return Some(Err(error)),
3266 };
3267 let value: Value = match serde_json::from_str(&line) {
3268 Ok(value) => value,
3269 Err(_) => continue,
3270 };
3271 match self.reject_empty_permission_request(&value).await {
3272 Ok(true) => continue,
3273 Ok(false) => {}
3274 Err(error) => return Some(Err(error)),
3275 }
3276 match self.handle_client_request(&value).await {
3277 Ok(true) => continue,
3278 Ok(false) => {}
3279 Err(error) => return Some(Err(error)),
3280 }
3281 match parse_acp_notification(self.slot, &line) {
3282 Ok(Some(event)) => {
3283 if let AgentEvent::ModelsReplaced {
3284 config_id, models, ..
3285 } = &event
3286 {
3287 self.model_config_id = Some(config_id.clone());
3288 self.models = models.clone();
3289 self.capabilities.supports_models = !models.is_empty();
3290 }
3291 return Some(Ok(event));
3292 }
3293 Ok(None) => {}
3294 Err(error) => return Some(Err(error)),
3295 }
3296 if value.get("id").is_some_and(|id| {
3297 self.prompt_request_id
3298 .is_some_and(|expected| rpc_id_to_string(id) == expected.to_string())
3299 }) {
3300 if let Some(error) = value.get("error") {
3301 self.prompt_request_id = None;
3302 return Some(Err(AdapterError::Protocol(error.to_string())));
3303 }
3304 self.prompt_request_id = None;
3305 return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3306 }
3307 }
3308 }
3309}
3310
3311fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3312 let value: Value =
3313 serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3314 let method = value.get("method").and_then(Value::as_str);
3315 if method == Some("session/request_permission") {
3316 let params = value.get("params").cloned().unwrap_or(Value::Null);
3317 let request_id = value
3318 .get("id")
3319 .map(rpc_id_to_string)
3320 .unwrap_or_else(|| "permission".into());
3321 return Ok(parse_permission_event(
3322 slot,
3323 ¶ms,
3324 &request_id,
3325 params.get("options"),
3326 ));
3327 }
3328 if method != Some("session/update") {
3329 return Ok(None);
3330 }
3331 let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3332 return Ok(None);
3333 };
3334 let kind = update.get("sessionUpdate").and_then(Value::as_str);
3335 if kind == Some("config_option_update")
3336 && let Some((config_id, models, current_model)) = parse_model_config(update)
3337 {
3338 return Ok(Some(AgentEvent::ModelsReplaced {
3339 slot,
3340 config_id,
3341 models,
3342 current_model,
3343 }));
3344 }
3345 if kind == Some("request_permission") {
3346 let request_id = update
3347 .get("toolCall")
3348 .and_then(|tool| tool.get("toolCallId"))
3349 .and_then(Value::as_str)
3350 .unwrap_or("permission");
3351 return Ok(parse_permission_event(
3352 slot,
3353 update,
3354 request_id,
3355 update.get("options"),
3356 ));
3357 }
3358 if kind == Some("available_commands_update") {
3359 let commands = update
3360 .get("availableCommands")
3361 .and_then(Value::as_array)
3362 .map(|commands| {
3363 commands
3364 .iter()
3365 .filter_map(|command| {
3366 let name = command.get("name").and_then(Value::as_str)?.trim();
3367 (!name.is_empty()).then(|| AgentCommand {
3368 name: name.to_owned(),
3369 })
3370 })
3371 .collect::<Vec<_>>()
3372 })
3373 .unwrap_or_default();
3374 return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3375 }
3376 if kind == Some("current_mode_update") {
3377 if let Some(mode) = update
3378 .get("currentModeId")
3379 .and_then(Value::as_str)
3380 .filter(|mode| !mode.trim().is_empty())
3381 {
3382 return Ok(Some(AgentEvent::ModeUpdated {
3383 slot,
3384 current_mode: mode.to_owned(),
3385 }));
3386 }
3387 return Ok(None);
3388 }
3389 if kind == Some("usage_update") {
3390 let Some(used) = update.get("used").and_then(Value::as_u64) else {
3391 return Ok(None);
3392 };
3393 let Some(size) = update.get("size").and_then(Value::as_u64) else {
3394 return Ok(None);
3395 };
3396 return Ok(Some(AgentEvent::UsageUpdated {
3397 slot,
3398 usage: UsageUpdate { used, size },
3399 }));
3400 }
3401 if let Some(terminal) = parse_terminal_event(update, kind) {
3402 return Ok(Some(AgentEvent::Terminal {
3403 slot,
3404 event: terminal,
3405 }));
3406 }
3407 let text = update
3408 .get("content")
3409 .and_then(|content| content.get("text"))
3410 .and_then(Value::as_str)
3411 .map(str::to_owned);
3412 if kind == Some("user_message_chunk") {
3413 return Ok(text
3414 .filter(|text| !text.is_empty())
3415 .map(|text| AgentEvent::UserText { slot, text }));
3416 }
3417 if kind == Some("agent_message_chunk")
3418 && let Some(mode) = text
3419 .as_deref()
3420 .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3421 .map(str::trim)
3422 .filter(|mode| !mode.is_empty())
3423 {
3424 return Ok(Some(AgentEvent::ModesReplaced {
3429 slot,
3430 modes: vec![Mode {
3431 id: mode.to_owned(),
3432 label: mode.to_owned(),
3433 }],
3434 current_mode: Some(mode.to_owned()),
3435 }));
3436 }
3437 match (kind, text) {
3438 (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3439 Ok(Some(AgentEvent::Text { slot, text }))
3440 }
3441 (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3442 Ok(Some(AgentEvent::Thought { slot, text }))
3443 }
3444 (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3445 let id = update
3446 .get("toolCallId")
3447 .and_then(Value::as_str)
3448 .unwrap_or("unknown-tool")
3449 .to_owned();
3450 let title = update
3451 .get("title")
3452 .and_then(Value::as_str)
3453 .unwrap_or("Tool call")
3454 .to_owned();
3455 let status = match update.get("status").and_then(Value::as_str) {
3456 Some("completed") => ToolStatus::Completed,
3457 Some("failed") => ToolStatus::Failed,
3458 Some("in_progress") => ToolStatus::Running,
3459 _ => ToolStatus::Pending,
3460 };
3461 Ok(Some(AgentEvent::Tool {
3462 slot,
3463 update: ToolUpdate {
3464 id,
3465 title,
3466 status,
3467 detail: None,
3468 },
3469 }))
3470 }
3471 _ => Ok(None),
3472 }
3473}
3474
3475fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3476 let config = value
3477 .get("configOptions")?
3478 .as_array()?
3479 .iter()
3480 .find(|option| {
3481 option.get("category").and_then(Value::as_str) == Some("model")
3482 && matches!(
3483 option.get("type").and_then(Value::as_str),
3484 Some("select" | "enum")
3485 )
3486 })?;
3487 let config_id = config.get("id")?.as_str()?.to_owned();
3488 let models = config
3489 .get("options")?
3490 .as_array()?
3491 .iter()
3492 .filter_map(|option| {
3493 let id = option.get("value")?.as_str()?.to_owned();
3494 let label = option
3495 .get("name")
3496 .or_else(|| option.get("label"))
3497 .and_then(Value::as_str)
3498 .unwrap_or(&id)
3499 .to_owned();
3500 Some(Mode { id, label })
3501 })
3502 .collect::<Vec<_>>();
3503 (!models.is_empty()).then(|| {
3504 let current = config
3505 .get("currentValue")
3506 .and_then(Value::as_str)
3507 .map(str::to_owned);
3508 (config_id, models, current)
3509 })
3510}
3511
3512fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3517 let nested = value.get("terminal").unwrap_or(value);
3518 let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3519 let id = nested
3520 .get("terminalId")
3521 .or_else(|| nested.get("terminal_id"))
3522 .or_else(|| nested.get("id"))
3523 .and_then(Value::as_str)
3524 .unwrap_or("terminal")
3525 .to_owned();
3526 match kind {
3527 "terminal_created" | "terminal_create" | "terminal_started" => {
3528 let command = nested
3529 .get("command")
3530 .and_then(Value::as_str)
3531 .unwrap_or("")
3532 .to_owned();
3533 Some(TerminalEvent::Created { id, command })
3534 }
3535 "terminal_output" | "terminal_output_chunk" => {
3536 let text = nested
3537 .get("output")
3538 .or_else(|| nested.get("text"))
3539 .and_then(Value::as_str)
3540 .unwrap_or("")
3541 .to_owned();
3542 Some(TerminalEvent::Output { id, text })
3543 }
3544 "terminal_exited" | "terminal_exit" => {
3545 let code = nested
3546 .get("exitCode")
3547 .or_else(|| nested.get("exit_code"))
3548 .or_else(|| nested.get("code"))
3549 .and_then(Value::as_i64)
3550 .unwrap_or(0) as i32;
3551 Some(TerminalEvent::Exited { id, code })
3552 }
3553 "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3554 _ => None,
3555 }
3556}
3557
3558fn parse_permission_event(
3559 slot: RosterSlot,
3560 value: &Value,
3561 request_id: &str,
3562 options: Option<&Value>,
3563) -> Option<AgentEvent> {
3564 let tool = value.get("toolCall").unwrap_or(value);
3565 let title = tool
3566 .get("title")
3567 .and_then(Value::as_str)
3568 .unwrap_or("Agent requests permission")
3569 .to_owned();
3570 let (options, option_ids): (Vec<String>, Vec<String>) = options
3571 .and_then(Value::as_array)
3572 .map(|options| {
3573 options
3574 .iter()
3575 .filter_map(|option| {
3576 let label = option
3577 .get("name")
3578 .or_else(|| option.get("optionId"))
3579 .and_then(Value::as_str)?
3580 .to_owned();
3581 let option_id = option
3582 .get("optionId")
3583 .or_else(|| option.get("id"))
3584 .and_then(Value::as_str)
3585 .map(str::to_owned)
3586 .unwrap_or_else(|| label.clone());
3587 Some((label, option_id))
3588 })
3589 .unzip()
3590 })
3591 .unwrap_or_default();
3592 if options.is_empty() {
3593 return None;
3594 }
3595 Some(AgentEvent::Permission {
3596 slot,
3597 request: PermissionRequest {
3598 id: request_id.to_owned(),
3599 title,
3600 options,
3601 option_ids,
3602 },
3603 })
3604}
3605
3606fn rpc_id_to_string(value: &Value) -> String {
3607 value
3608 .as_str()
3609 .map(str::to_owned)
3610 .or_else(|| value.as_u64().map(|id| id.to_string()))
3611 .unwrap_or_else(|| value.to_string())
3612}
3613
3614#[cfg(test)]
3615mod tests {
3616 use super::{
3617 AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
3618 RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
3619 parse_model_config, prompt_content_blocks, read_bounded_line,
3620 };
3621 #[cfg(target_os = "linux")]
3622 use super::{isolate_process_group, terminate_child};
3623 use crate::TerminalEvent;
3624 use crate::{
3625 AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
3626 persistence::SessionMetadataStore,
3627 relay::{DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
3628 };
3629 use async_trait::async_trait;
3630 use serde_json::Value;
3631 use std::sync::{
3632 Arc, Mutex,
3633 atomic::{AtomicUsize, Ordering},
3634 };
3635
3636 fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
3637 let nonce = std::time::SystemTime::now()
3638 .duration_since(std::time::UNIX_EPOCH)
3639 .expect("clock")
3640 .as_nanos();
3641 std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
3642 }
3643
3644 #[test]
3645 fn malformed_file_writes_preserve_existing_content() {
3646 let root = unique_test_path("codeswarm-write-validation", "dir");
3647 std::fs::create_dir_all(&root).unwrap();
3648 let file = root.join("keep.txt");
3649 std::fs::write(&file, "valuable content").unwrap();
3650 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
3651 for content in [
3652 Value::Null,
3653 serde_json::json!(false),
3654 serde_json::json!(42),
3655 serde_json::json!([]),
3656 ] {
3657 assert!(
3658 adapter
3659 .write_workspace_text(
3660 &serde_json::json!({"path":"keep.txt", "content": content})
3661 )
3662 .is_err()
3663 );
3664 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3665 }
3666 assert!(
3667 adapter
3668 .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
3669 .is_err()
3670 );
3671 assert!(
3672 adapter
3673 .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
3674 .is_err()
3675 );
3676 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3677 adapter
3678 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
3679 .unwrap();
3680 assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
3681 adapter
3682 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
3683 .unwrap();
3684 assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
3685 std::fs::remove_dir_all(root).unwrap();
3686 }
3687
3688 #[tokio::test]
3689 async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
3690 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
3691 let mut adapter = AcpAdapter::new(
3692 0,
3693 std::env::current_dir().unwrap(),
3694 "sh",
3695 vec!["-c".into(), script.into()],
3696 );
3697 adapter.start().await.unwrap();
3698 let error = adapter
3699 .request_with_timeout(
3700 "session/set_mode",
3701 serde_json::json!({}),
3702 std::time::Duration::from_millis(10),
3703 )
3704 .await
3705 .unwrap_err();
3706 assert!(error.to_string().contains("session/set_mode timed out"));
3707 adapter.stop().await.unwrap();
3708 assert!(adapter.child.is_none());
3709 }
3710
3711 #[tokio::test]
3712 async fn goals_reach_every_roster_slot_without_native_goal_support() {
3713 use crate::goal::GoalCommand;
3714 let hosts = (0..3)
3715 .map(|slot| {
3716 AdapterHost::new(
3717 Box::new(ScriptedAdapter::new(
3718 slot,
3719 AgentCapabilities::default(),
3720 [
3721 AgentEvent::TurnComplete { slot },
3722 AgentEvent::TurnComplete { slot },
3723 ],
3724 )),
3725 None,
3726 )
3727 })
3728 .collect();
3729 let mut relay = RelayHost::new(hosts, 10).unwrap();
3730 relay.start().await.unwrap();
3731 let task = relay
3732 .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
3733 .unwrap()
3734 .unwrap();
3735 relay.relay_mut().enqueue_human(task, Some(0));
3736 for slot in 0..3 {
3737 relay.run_turn("", 0).await.unwrap();
3738 let (actual, prompt) = relay.dispatches().last().unwrap();
3739 assert_eq!(*actual, slot);
3740 assert!(prompt.contains("Active shared goal: Ship the settings screen"));
3741 }
3742 let snapshot = relay.session_metadata();
3743 let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
3744 assert!(restored.is_some());
3745 relay.restore_goal(restored);
3746 relay.reload(0).await.unwrap();
3747 relay.run_turn("", 0).await.unwrap();
3748 assert!(
3749 relay
3750 .dispatches()
3751 .last()
3752 .unwrap()
3753 .1
3754 .contains("Active shared goal: Ship the settings screen")
3755 );
3756 relay.apply_goal(GoalCommand::Done).unwrap();
3757 relay.run_turn("", 0).await.unwrap();
3758 assert!(
3759 relay
3760 .dispatches()
3761 .last()
3762 .unwrap()
3763 .1
3764 .contains("No active shared goal")
3765 );
3766 relay.apply_goal(GoalCommand::Clear).unwrap();
3767 assert!(relay.session_metadata().get("goal").unwrap().is_null());
3768 }
3769
3770 #[tokio::test]
3771 async fn replacement_agent_receives_task_after_public_journal_pruning() {
3772 let hosts = (0..2)
3773 .map(|slot| {
3774 AdapterHost::new(
3775 Box::new(ScriptedAdapter::new(
3776 slot,
3777 AgentCapabilities::default(),
3778 [
3779 AgentEvent::Text {
3780 slot,
3781 text: "progress".into(),
3782 },
3783 AgentEvent::TurnComplete { slot },
3784 AgentEvent::Text {
3785 slot,
3786 text: "more progress".into(),
3787 },
3788 AgentEvent::TurnComplete { slot },
3789 ],
3790 )),
3791 None,
3792 )
3793 })
3794 .collect();
3795 let mut relay = RelayHost::new(hosts, 10).unwrap();
3796 relay.start().await.unwrap();
3797 relay
3798 .relay_mut()
3799 .enqueue_human("Fix the login bug", Some(0));
3800 relay.run_turn("", 0).await.unwrap();
3801 relay.run_turn("", 0).await.unwrap();
3802 relay.run_turn("", 0).await.unwrap();
3803 relay.reload(1).await.unwrap();
3804 assert!(
3805 !relay
3806 .relay_mut()
3807 .unseen_context(1)
3808 .contains("Fix the login bug")
3809 );
3810 relay.run_turn("", 0).await.unwrap();
3811 assert!(
3812 relay
3813 .dispatches()
3814 .last()
3815 .unwrap()
3816 .1
3817 .contains("Shared task:\nFix the login bug")
3818 );
3819 }
3820
3821 #[cfg(target_os = "linux")]
3822 #[tokio::test]
3823 async fn termination_kills_only_the_verified_isolated_child_group() {
3824 use nix::unistd::{Pid, getpgid, getpgrp};
3825 use tokio::io::{AsyncBufReadExt, BufReader};
3826
3827 let own_group = getpgrp();
3828 let mut command = tokio::process::Command::new("sh");
3829 isolate_process_group(&mut command);
3830 command
3831 .arg("-c")
3832 .arg("sleep 60 & echo $!; wait")
3833 .stdout(std::process::Stdio::piped());
3834 let mut child = command.spawn().expect("spawn isolated shell");
3835 let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
3836 assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
3837 assert_ne!(leader, own_group);
3838
3839 let stdout = child.stdout.take().expect("child stdout");
3840 let mut lines = BufReader::new(stdout).lines();
3841 let descendant = lines
3842 .next_line()
3843 .await
3844 .expect("read descendant pid")
3845 .expect("descendant pid")
3846 .parse::<i32>()
3847 .expect("numeric descendant pid");
3848 let descendant = Pid::from_raw(descendant);
3849 assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
3850
3851 terminate_child(&mut child).await.expect("terminate group");
3852 for _ in 0..100 {
3853 if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
3854 return;
3855 }
3856 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
3857 }
3858 panic!("descendant {descendant} survived isolated group termination");
3859 }
3860
3861 #[test]
3862 fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
3863 assert_eq!(
3864 parse_command_line(r#"npx -y "@agentclientprotocol/codex-acp" --flag 'two words'"#),
3865 Ok((
3866 "npx".into(),
3867 vec![
3868 "-y".into(),
3869 "@agentclientprotocol/codex-acp".into(),
3870 "--flag".into(),
3871 "two words".into(),
3872 ]
3873 ),)
3874 );
3875 assert_eq!(
3876 parse_command_line(r#"agent "" escaped\ argument"#),
3877 Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
3878 );
3879 }
3880
3881 #[test]
3882 fn acp_prompt_expands_safe_at_path_resources() {
3883 let root = unique_test_path("codeswarm-prompt-resource", "dir");
3884 std::fs::create_dir_all(&root).expect("workspace");
3885 std::fs::write(root.join("note.md"), "resource text").expect("resource");
3886 let blocks = prompt_content_blocks(&root, "inspect @note.md");
3887 assert_eq!(blocks[0]["type"], "text");
3888 assert_eq!(blocks[0]["text"], "inspect @note.md");
3889 assert_eq!(blocks[1]["type"], "resource");
3890 assert_eq!(blocks[1]["resource"]["text"], "resource text");
3891 assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
3892 std::fs::remove_dir_all(root).expect("cleanup workspace");
3893 }
3894
3895 #[tokio::test]
3896 async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
3897 let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
3898 bytes.push(b'\n');
3899 let mut reader = tokio::io::BufReader::new(bytes.as_slice());
3900 assert!(matches!(
3901 read_bounded_line(&mut reader).await,
3902 Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
3903 ));
3904 }
3905
3906 #[test]
3907 fn rejects_malformed_configured_commands_before_spawn() {
3908 assert_eq!(
3909 parse_command_line("agent 'unfinished"),
3910 Err(super::CommandParseError::UnterminatedQuote)
3911 );
3912 assert_eq!(
3913 parse_command_line("agent\\"),
3914 Err(super::CommandParseError::TrailingEscape)
3915 );
3916 assert_eq!(
3917 parse_command_line(" \t"),
3918 Err(super::CommandParseError::Empty)
3919 );
3920 }
3921
3922 #[derive(Debug)]
3923 struct PendingAdapter {
3924 slot: usize,
3925 hang_on_cancel: bool,
3926 }
3927
3928 #[derive(Debug)]
3929 struct ConcurrentStartAdapter {
3930 slot: usize,
3931 barrier: Arc<tokio::sync::Barrier>,
3932 }
3933
3934 #[async_trait]
3935 impl AgentAdapter for ConcurrentStartAdapter {
3936 fn slot(&self) -> usize {
3937 self.slot
3938 }
3939
3940 fn capabilities(&self) -> AgentCapabilities {
3941 AgentCapabilities::default()
3942 }
3943
3944 async fn start(&mut self) -> super::AdapterResult<()> {
3945 self.barrier.wait().await;
3946 Ok(())
3947 }
3948
3949 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
3950 Ok(())
3951 }
3952
3953 async fn cancel(&mut self) -> super::AdapterResult<bool> {
3954 Ok(true)
3955 }
3956
3957 async fn answer_permission(
3958 &mut self,
3959 _request_id: String,
3960 _answer: PermissionAnswer,
3961 ) -> super::AdapterResult<()> {
3962 Ok(())
3963 }
3964
3965 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
3966 Ok(())
3967 }
3968
3969 async fn reload(&mut self) -> super::AdapterResult<()> {
3970 Ok(())
3971 }
3972
3973 async fn stop(&mut self) -> super::AdapterResult<()> {
3974 Ok(())
3975 }
3976
3977 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
3978 std::future::pending().await
3979 }
3980 }
3981
3982 #[derive(Debug)]
3983 struct PermissionBlockingAdapter {
3984 slot: usize,
3985 phase: u8,
3986 }
3987
3988 #[async_trait]
3989 impl AgentAdapter for PermissionBlockingAdapter {
3990 fn slot(&self) -> usize {
3991 self.slot
3992 }
3993
3994 fn capabilities(&self) -> AgentCapabilities {
3995 AgentCapabilities {
3996 supports_permissions: true,
3997 ..AgentCapabilities::default()
3998 }
3999 }
4000
4001 async fn start(&mut self) -> super::AdapterResult<()> {
4002 Ok(())
4003 }
4004
4005 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4006 Ok(())
4007 }
4008
4009 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4010 Ok(true)
4011 }
4012
4013 async fn answer_permission(
4014 &mut self,
4015 request_id: String,
4016 answer: PermissionAnswer,
4017 ) -> super::AdapterResult<()> {
4018 if self.phase != 1 || request_id != "permission-1" {
4019 return Err(super::AdapterError::Protocol(
4020 "unexpected permission response".into(),
4021 ));
4022 }
4023 assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4024 self.phase = 2;
4025 Ok(())
4026 }
4027
4028 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4029 Ok(())
4030 }
4031
4032 async fn reload(&mut self) -> super::AdapterResult<()> {
4033 Ok(())
4034 }
4035
4036 async fn stop(&mut self) -> super::AdapterResult<()> {
4037 Ok(())
4038 }
4039
4040 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4041 match self.phase {
4042 0 => {
4043 self.phase = 1;
4044 Some(Ok(AgentEvent::Permission {
4045 slot: self.slot,
4046 request: crate::PermissionRequest {
4047 id: "permission-1".into(),
4048 title: "Allow?".into(),
4049 options: vec!["Allow".into()],
4050 option_ids: vec!["allow".into()],
4051 },
4052 }))
4053 }
4054 1 => std::future::pending().await,
4055 _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4056 }
4057 }
4058 }
4059
4060 #[async_trait]
4061 impl AgentAdapter for PendingAdapter {
4062 fn slot(&self) -> usize {
4063 self.slot
4064 }
4065
4066 fn capabilities(&self) -> AgentCapabilities {
4067 AgentCapabilities {
4068 supports_cancel: true,
4069 ..AgentCapabilities::default()
4070 }
4071 }
4072
4073 async fn start(&mut self) -> super::AdapterResult<()> {
4074 Ok(())
4075 }
4076
4077 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4078 Ok(())
4079 }
4080
4081 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4082 if self.hang_on_cancel {
4083 return std::future::pending().await;
4084 }
4085 Ok(true)
4086 }
4087
4088 async fn answer_permission(
4089 &mut self,
4090 _request_id: String,
4091 _answer: PermissionAnswer,
4092 ) -> super::AdapterResult<()> {
4093 Ok(())
4094 }
4095
4096 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4097 Ok(())
4098 }
4099
4100 async fn reload(&mut self) -> super::AdapterResult<()> {
4101 Ok(())
4102 }
4103
4104 async fn stop(&mut self) -> super::AdapterResult<()> {
4105 Ok(())
4106 }
4107
4108 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4109 std::future::pending().await
4110 }
4111 }
4112
4113 #[derive(Debug)]
4114 struct StopTrackingAdapter {
4115 slot: usize,
4116 stopped: Arc<AtomicUsize>,
4117 fail_stop: bool,
4118 }
4119
4120 #[derive(Debug)]
4121 struct ModeOrderAdapter {
4122 slot: usize,
4123 log: Arc<Mutex<Vec<String>>>,
4124 phase: u8,
4125 }
4126
4127 #[derive(Debug)]
4128 struct StartupAcpAdapter {
4129 slot: usize,
4130 events: std::collections::VecDeque<AgentEvent>,
4131 }
4132
4133 impl StartupAcpAdapter {
4134 fn new(slot: usize) -> Self {
4135 Self {
4136 slot,
4137 events: [
4138 AgentEvent::ModesReplaced {
4139 slot,
4140 modes: vec![Mode {
4141 id: "full-access".into(),
4142 label: "Auto pilot".into(),
4143 }],
4144 current_mode: Some("full-access".into()),
4145 },
4146 AgentEvent::Ready {
4147 slot,
4148 capabilities: AgentCapabilities {
4149 supports_modes: true,
4150 ..AgentCapabilities::default()
4151 },
4152 },
4153 ]
4154 .into(),
4155 }
4156 }
4157 }
4158
4159 #[async_trait]
4160 impl AgentAdapter for StartupAcpAdapter {
4161 fn slot(&self) -> usize {
4162 self.slot
4163 }
4164
4165 fn protocol(&self) -> &'static str {
4166 "acp"
4167 }
4168
4169 fn capabilities(&self) -> AgentCapabilities {
4170 AgentCapabilities {
4171 supports_modes: true,
4172 ..AgentCapabilities::default()
4173 }
4174 }
4175
4176 async fn start(&mut self) -> super::AdapterResult<()> {
4177 Ok(())
4178 }
4179
4180 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4181 Ok(())
4182 }
4183
4184 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4185 Ok(true)
4186 }
4187
4188 async fn answer_permission(
4189 &mut self,
4190 _request_id: String,
4191 _answer: PermissionAnswer,
4192 ) -> super::AdapterResult<()> {
4193 Ok(())
4194 }
4195
4196 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4197 Ok(())
4198 }
4199
4200 async fn reload(&mut self) -> super::AdapterResult<()> {
4201 self.events = Self::new(self.slot).events;
4202 Ok(())
4203 }
4204
4205 async fn stop(&mut self) -> super::AdapterResult<()> {
4206 Ok(())
4207 }
4208
4209 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4210 self.events.pop_front().map(Ok)
4211 }
4212 }
4213
4214 #[async_trait]
4215 impl AgentAdapter for ModeOrderAdapter {
4216 fn slot(&self) -> usize {
4217 self.slot
4218 }
4219
4220 fn capabilities(&self) -> AgentCapabilities {
4221 AgentCapabilities {
4222 supports_modes: true,
4223 ..AgentCapabilities::default()
4224 }
4225 }
4226
4227 async fn start(&mut self) -> super::AdapterResult<()> {
4228 self.log.lock().expect("log").push("start".into());
4229 Ok(())
4230 }
4231
4232 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4233 self.log.lock().expect("log").push("prompt".into());
4234 Ok(())
4235 }
4236
4237 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4238 Ok(true)
4239 }
4240
4241 async fn answer_permission(
4242 &mut self,
4243 _request_id: String,
4244 _answer: PermissionAnswer,
4245 ) -> super::AdapterResult<()> {
4246 Ok(())
4247 }
4248
4249 async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4250 self.log.lock().expect("log").push(format!("mode:{mode}"));
4251 Ok(())
4252 }
4253
4254 async fn reload(&mut self) -> super::AdapterResult<()> {
4255 self.log.lock().expect("log").push("reload".into());
4256 self.phase = 0;
4257 Ok(())
4258 }
4259
4260 async fn stop(&mut self) -> super::AdapterResult<()> {
4261 Ok(())
4262 }
4263
4264 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4265 match self.phase {
4266 0 => {
4267 self.phase = 1;
4268 Some(Ok(AgentEvent::ModesReplaced {
4269 slot: self.slot,
4270 modes: vec![Mode {
4271 id: "yolo".into(),
4272 label: "YOLO".into(),
4273 }],
4274 current_mode: None,
4275 }))
4276 }
4277 1 => {
4278 self.phase = 2;
4279 Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4280 }
4281 _ => std::future::pending().await,
4282 }
4283 }
4284 }
4285
4286 #[async_trait]
4287 impl AgentAdapter for StopTrackingAdapter {
4288 fn slot(&self) -> usize {
4289 self.slot
4290 }
4291
4292 fn capabilities(&self) -> AgentCapabilities {
4293 AgentCapabilities::default()
4294 }
4295
4296 async fn start(&mut self) -> super::AdapterResult<()> {
4297 Ok(())
4298 }
4299
4300 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4301 Ok(())
4302 }
4303
4304 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4305 Ok(false)
4306 }
4307
4308 async fn answer_permission(
4309 &mut self,
4310 _request_id: String,
4311 _answer: PermissionAnswer,
4312 ) -> super::AdapterResult<()> {
4313 Ok(())
4314 }
4315
4316 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4317 Ok(())
4318 }
4319
4320 async fn reload(&mut self) -> super::AdapterResult<()> {
4321 Ok(())
4322 }
4323
4324 async fn stop(&mut self) -> super::AdapterResult<()> {
4325 self.stopped.fetch_add(1, Ordering::Relaxed);
4326 if self.fail_stop {
4327 Err(super::AdapterError::Transport("stop failed".into()))
4328 } else {
4329 Ok(())
4330 }
4331 }
4332
4333 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4334 None
4335 }
4336 }
4337
4338 #[derive(Debug)]
4342 struct FailingStartAdapter {
4343 slot: usize,
4344 stopped: Arc<AtomicUsize>,
4345 }
4346
4347 #[async_trait]
4348 impl AgentAdapter for FailingStartAdapter {
4349 fn slot(&self) -> usize {
4350 self.slot
4351 }
4352
4353 fn capabilities(&self) -> AgentCapabilities {
4354 AgentCapabilities::default()
4355 }
4356
4357 async fn start(&mut self) -> super::AdapterResult<()> {
4358 Err(super::AdapterError::Spawn("startup failed".into()))
4359 }
4360
4361 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4362 Ok(())
4363 }
4364
4365 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4366 Ok(false)
4367 }
4368
4369 async fn answer_permission(
4370 &mut self,
4371 _request_id: String,
4372 _answer: PermissionAnswer,
4373 ) -> super::AdapterResult<()> {
4374 Ok(())
4375 }
4376
4377 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4378 Ok(())
4379 }
4380
4381 async fn reload(&mut self) -> super::AdapterResult<()> {
4382 Ok(())
4383 }
4384
4385 async fn stop(&mut self) -> super::AdapterResult<()> {
4386 self.stopped.fetch_add(1, Ordering::Relaxed);
4387 Ok(())
4388 }
4389
4390 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4391 None
4392 }
4393 }
4394
4395 #[tokio::test]
4396 async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4397 let stopped = Arc::new(AtomicUsize::new(0));
4398 let relay = RelayHost::new(
4399 vec![
4400 AdapterHost::new(
4401 Box::new(StopTrackingAdapter {
4402 slot: 0,
4403 stopped: Arc::clone(&stopped),
4404 fail_stop: true,
4405 }),
4406 None,
4407 ),
4408 AdapterHost::new(
4409 Box::new(StopTrackingAdapter {
4410 slot: 1,
4411 stopped: Arc::clone(&stopped),
4412 fail_stop: false,
4413 }),
4414 None,
4415 ),
4416 ],
4417 4,
4418 )
4419 .expect("relay");
4420 let mut relay = relay;
4421
4422 let error = relay.stop().await.expect_err("first stop failure");
4423 assert!(error.to_string().contains("stop failed"));
4424 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4425 }
4426
4427 #[tokio::test]
4428 async fn relay_start_cleans_up_the_adapter_that_failed_startup() {
4429 let stopped = Arc::new(AtomicUsize::new(0));
4430 let mut relay = RelayHost::new(
4431 vec![
4432 AdapterHost::new(
4433 Box::new(StopTrackingAdapter {
4434 slot: 0,
4435 stopped: Arc::clone(&stopped),
4436 fail_stop: false,
4437 }),
4438 None,
4439 ),
4440 AdapterHost::new(
4441 Box::new(FailingStartAdapter {
4442 slot: 1,
4443 stopped: Arc::clone(&stopped),
4444 }),
4445 None,
4446 ),
4447 ],
4448 4,
4449 )
4450 .expect("relay");
4451
4452 assert!(relay.start().await.is_err());
4453 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4454 }
4455
4456 #[test]
4457 fn parses_acp_text_without_ui_dependency() {
4458 let event = parse_acp_notification(
4459 2,
4460 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4461 )
4462 .expect("valid ACP")
4463 .expect("text event");
4464 assert_eq!(
4465 event,
4466 AgentEvent::Text {
4467 slot: 2,
4468 text: "hello".into(),
4469 }
4470 );
4471 }
4472
4473 #[test]
4474 fn parses_acp_state_notifications_at_the_adapter_boundary() {
4475 let commands = parse_acp_notification(
4476 3,
4477 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
4478 )
4479 .expect("valid ACP")
4480 .expect("commands event");
4481 assert_eq!(
4482 commands,
4483 AgentEvent::CommandsReplaced {
4484 slot: 3,
4485 commands: vec![crate::AgentCommand {
4486 name: "review".into()
4487 }]
4488 }
4489 );
4490
4491 let mode = parse_acp_notification(
4492 3,
4493 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
4494 )
4495 .expect("valid ACP")
4496 .expect("mode event");
4497 assert_eq!(
4498 mode,
4499 AgentEvent::ModeUpdated {
4500 slot: 3,
4501 current_mode: "review".into()
4502 }
4503 );
4504
4505 let usage = parse_acp_notification(
4506 3,
4507 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
4508 )
4509 .expect("valid ACP")
4510 .expect("usage event");
4511 assert_eq!(
4512 usage,
4513 AgentEvent::UsageUpdated {
4514 slot: 3,
4515 usage: crate::UsageUpdate {
4516 used: 4200,
4517 size: 128000
4518 }
4519 }
4520 );
4521
4522 let models = parse_acp_notification(
4523 3,
4524 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"config_option_update","configOptions":[{"id":"model","category":"model","type":"select","currentValue":"smart","options":[{"value":"fast","name":"Fast"},{"value":"smart","name":"Smart"}]}]}}}"#,
4525 )
4526 .expect("valid ACP")
4527 .expect("models event");
4528 assert!(matches!(
4529 models,
4530 AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
4531 if models.len() == 2 && current_model.as_deref() == Some("smart")
4532 ));
4533 assert_eq!(
4534 parse_model_config(&serde_json::json!({
4535 "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
4536 })),
4537 None
4538 );
4539
4540 let user = parse_acp_notification(
4541 3,
4542 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
4543 )
4544 .expect("valid ACP")
4545 .expect("user event");
4546 assert_eq!(
4547 user,
4548 AgentEvent::UserText {
4549 slot: 3,
4550 text: "context".into()
4551 }
4552 );
4553 }
4554
4555 #[test]
4556 fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
4557 let event = parse_acp_notification(
4558 0,
4559 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
4560 )
4561 .expect("valid ACP")
4562 .expect("mode event");
4563 assert!(matches!(
4564 event,
4565 AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
4566 if mode == "yolo" && modes[0].id == "yolo"
4567 ));
4568 }
4569
4570 #[test]
4571 fn parses_native_agy_text_without_acp_bridge() {
4572 let event = parse_agy_line(
4573 1,
4574 r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
4575 )
4576 .expect("valid stream-json")
4577 .expect("text event");
4578 assert_eq!(
4579 event,
4580 AgentEvent::Text {
4581 slot: 1,
4582 text: "hello".into(),
4583 }
4584 );
4585 }
4586
4587 #[test]
4588 fn parses_tool_lifecycle_from_each_protocol() {
4589 let agy = parse_agy_line(
4590 1,
4591 r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
4592 )
4593 .expect("valid native tool")
4594 .expect("tool event");
4595 assert!(matches!(
4596 agy,
4597 AgentEvent::Tool {
4598 update: crate::ToolUpdate {
4599 status: ToolStatus::Completed,
4600 ..
4601 },
4602 ..
4603 }
4604 ));
4605
4606 let acp = parse_acp_notification(
4607 1,
4608 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
4609 )
4610 .expect("valid ACP tool")
4611 .expect("tool event");
4612 assert!(matches!(
4613 acp,
4614 AgentEvent::Tool {
4615 update: crate::ToolUpdate {
4616 status: ToolStatus::Failed,
4617 ..
4618 },
4619 ..
4620 }
4621 ));
4622 }
4623
4624 #[test]
4625 fn parses_terminal_lifecycle_from_acp_and_native_events() {
4626 let created = parse_acp_notification(
4627 0,
4628 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
4629 )
4630 .expect("valid ACP terminal")
4631 .expect("terminal event");
4632 assert_eq!(
4633 created,
4634 AgentEvent::Terminal {
4635 slot: 0,
4636 event: TerminalEvent::Created {
4637 id: "term-1".into(),
4638 command: "cargo test".into(),
4639 },
4640 }
4641 );
4642 let output = parse_agy_line(
4643 1,
4644 r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
4645 )
4646 .expect("valid native terminal")
4647 .expect("terminal event");
4648 assert_eq!(
4649 output,
4650 AgentEvent::Terminal {
4651 slot: 1,
4652 event: TerminalEvent::Output {
4653 id: "term-1".into(),
4654 text: "ok\n".into(),
4655 },
4656 }
4657 );
4658 let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
4659 .expect("valid native release")
4660 .expect("terminal event");
4661 assert!(matches!(
4662 released,
4663 AgentEvent::Terminal {
4664 event: TerminalEvent::Released { id },
4665 ..
4666 } if id == "term-1"
4667 ));
4668 }
4669
4670 #[test]
4671 fn parses_acp_permission_requests() {
4672 let event = parse_acp_notification(
4673 0,
4674 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"request_permission","toolCall":{"toolCallId":"t1","title":"Write file"},"options":[{"name":"Allow once","optionId":"allow-once"},{"name":"Reject","optionId":"reject"}]}}}"#,
4675 )
4676 .expect("valid permission")
4677 .expect("permission event");
4678 assert!(matches!(
4679 event,
4680 AgentEvent::Permission { request, .. }
4681 if request.id == "t1"
4682 && request.title == "Write file"
4683 && request.options == ["Allow once", "Reject"]
4684 && request.option_ids == ["allow-once", "reject"]
4685 ));
4686 }
4687
4688 #[test]
4689 fn parses_acp_permission_request_as_json_rpc_request() {
4690 let event = parse_acp_notification(
4691 2,
4692 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
4693 )
4694 .expect("valid permission request")
4695 .expect("permission event");
4696 assert!(matches!(
4697 event,
4698 AgentEvent::Permission { request, .. }
4699 if request.id == "17"
4700 && request.title == "Write file"
4701 && request.options == ["allow-once", "reject"]
4702 && request.option_ids == ["allow-once", "reject"]
4703 ));
4704 }
4705
4706 #[tokio::test]
4707 async fn native_adapter_explicitly_rejects_permission_answers() {
4708 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4709 assert_eq!(
4710 adapter
4711 .answer_permission(
4712 "request".into(),
4713 PermissionAnswer::Selected {
4714 option_id: "allow".into()
4715 },
4716 )
4717 .await,
4718 Err(super::AdapterError::Unsupported("permission answer"))
4719 );
4720 }
4721
4722 #[tokio::test]
4723 async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
4724 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4725 adapter
4726 .set_mode("full-access".into())
4727 .await
4728 .expect("auto-pilot alias");
4729 assert!(matches!(
4730 adapter.next_event().await,
4731 Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
4732 ));
4733 }
4734
4735 #[tokio::test]
4736 async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
4737 let script_path = unique_test_path("codeswarm-agy-session", "sh");
4738 std::fs::write(
4739 &script_path,
4740 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"init\",\"conversation_id\":\"native-session\"}' '{\"event\":\"step_update\",\"step_update\":{\"step_type\":\"agent_response\",\"text_delta\":\"ok\"}}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"ok\"}}'\n",
4741 )
4742 .expect("write native test script");
4743 let mut adapter = AgyAdapter::new(
4744 0,
4745 std::env::current_dir().expect("cwd"),
4746 format!("sh {}", script_path.display()),
4747 );
4748 adapter.start().await.expect("start native adapter");
4749 assert!(adapter.next_event().await.is_some());
4751 assert!(adapter.next_event().await.is_some());
4752 adapter
4753 .send_prompt("first".into())
4754 .await
4755 .expect("first prompt");
4756 while !matches!(
4757 adapter.next_event().await,
4758 Some(Ok(AgentEvent::TurnComplete { .. }))
4759 ) {}
4760 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4761 adapter
4762 .send_prompt("follow up".into())
4763 .await
4764 .expect("follow-up prompt");
4765 while !matches!(
4766 adapter.next_event().await,
4767 Some(Ok(AgentEvent::TurnComplete { .. }))
4768 ) {}
4769 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4770 adapter.stop().await.expect("stop native adapter");
4771 std::fs::remove_file(script_path).expect("cleanup native script");
4772 }
4773
4774 #[tokio::test]
4775 async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
4776 let script_path = unique_test_path("codeswarm-agy-failure", "sh");
4777 std::fs::write(
4778 &script_path,
4779 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
4780 )
4781 .expect("write native test script");
4782 let mut adapter = AgyAdapter::new(
4783 0,
4784 std::env::current_dir().expect("cwd"),
4785 format!("sh {}", script_path.display()),
4786 );
4787 adapter.start().await.expect("start native adapter");
4788 assert!(adapter.next_event().await.is_some());
4789 assert!(adapter.next_event().await.is_some());
4790 adapter.send_prompt("fail".into()).await.expect("prompt");
4791 assert!(matches!(
4792 adapter.next_event().await,
4793 Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
4794 if detail == "agent failed"
4795 ));
4796 adapter.stop().await.expect("stop native adapter");
4797 std::fs::remove_file(script_path).expect("cleanup native script");
4798 }
4799
4800 #[tokio::test]
4801 async fn acp_adapter_initializes_session_and_completes_a_prompt() {
4802 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","modes":{"currentModeId":"plan","availableModes":[{"id":"plan","name":"Plan"}]}}}'; read _; echo '{"jsonrpc":"2.0","method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}'; echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'"#;
4803 let cwd = std::env::current_dir().expect("cwd");
4804 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4805 adapter.start().await.expect("initialize");
4806 assert!(matches!(
4807 adapter.next_event().await,
4808 Some(Ok(AgentEvent::ModesReplaced { .. }))
4809 ));
4810 assert!(matches!(
4811 adapter.next_event().await,
4812 Some(Ok(AgentEvent::Ready { .. }))
4813 ));
4814 adapter.send_prompt("hello".into()).await.expect("prompt");
4815 assert!(matches!(
4816 adapter.next_event().await,
4817 Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
4818 ));
4819 assert!(matches!(
4820 adapter.next_event().await,
4821 Some(Ok(AgentEvent::TurnComplete { .. }))
4822 ));
4823 }
4824
4825 #[tokio::test]
4826 async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
4827 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","modes":{"currentModeId":"plan","availableModes":[{"id":"plan","name":"Plan"}]}}}'; read _; echo '{"jsonrpc":"2.0","method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"first"}}}}'; echo '{"jsonrpc":"2.0","id":"3","result":{"stopReason":"end_turn"}}'; read _; echo '{"jsonrpc":"2.0","method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"second"}}}}'; echo '{"jsonrpc":"2.0","id":"4","result":{"stopReason":"end_turn"}}'"#;
4828 let cwd = std::env::current_dir().expect("cwd");
4829 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
4830 adapter.start().await.expect("initialize");
4831 assert!(adapter.next_event().await.is_some());
4832 assert!(adapter.next_event().await.is_some());
4833
4834 for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
4835 adapter.send_prompt(prompt.into()).await.expect("prompt");
4836 assert!(matches!(
4837 adapter.next_event().await,
4838 Some(Ok(AgentEvent::Text { text, .. })) if text == expected
4839 ));
4840 assert!(matches!(
4841 adapter.next_event().await,
4842 Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
4843 ));
4844 }
4845 adapter.stop().await.expect("stop");
4846 }
4847
4848 #[tokio::test]
4849 async fn empty_acp_mode_catalog_disables_mode_control() {
4850 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","modes":{"availableModes":[]}}}'"#;
4851 let cwd = std::env::current_dir().expect("cwd");
4852 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4853 adapter.start().await.expect("initialize");
4854 assert!(!adapter.capabilities().supports_modes);
4855 assert!(matches!(
4856 adapter.next_event().await,
4857 Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
4858 ));
4859 assert!(matches!(
4860 adapter.next_event().await,
4861 Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
4862 ));
4863 adapter.stop().await.expect("stop");
4864 }
4865
4866 #[tokio::test]
4867 async fn acp_models_are_discovered_live_and_changed_through_session_config() {
4868 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","configOptions":[{"id":"model","category":"model","type":"select","currentValue":"fast","options":[{"value":"fast","name":"Fast"},{"value":"smart","name":"Smart"}]}]}}'; read request; case "$request" in *session/set_config_option*\"value\":\"smart\"*) echo '{"jsonrpc":"2.0","id":3,"result":{}}';; *) echo '{"jsonrpc":"2.0","id":3,"error":{"code":-32602,"message":"wrong model request"}}';; esac"#;
4869 let cwd = std::env::current_dir().expect("cwd");
4870 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4871 adapter.start().await.expect("initialize");
4872 assert!(adapter.capabilities().supports_models);
4873 assert!(matches!(
4874 adapter.next_event().await,
4875 Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
4876 if config_id == "model"
4877 && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
4878 && current_model.as_deref() == Some("fast")
4879 ));
4880 assert!(matches!(
4881 adapter.next_event().await,
4882 Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
4883 ));
4884 adapter.set_model("smart".into()).await.expect("set model");
4885 assert!(adapter.set_model("invented".into()).await.is_err());
4886 adapter.stop().await.expect("stop");
4887 }
4888
4889 #[tokio::test]
4890 async fn acp_mode_change_is_acknowledged_without_provider_notification() {
4891 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","modes":{"currentModeId":"plan","availableModes":[{"id":"plan","name":"Plan"},{"id":"yolo","name":"YOLO"}]}}}'; read _; echo '{"jsonrpc":"2.0","id":3,"result":{}}'"#;
4892 let cwd = std::env::current_dir().expect("cwd");
4893 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4894 adapter.start().await.expect("initialize");
4895 adapter
4896 .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
4897 .await
4898 .expect("set mode");
4899 assert!(matches!(
4900 adapter.next_event().await,
4901 Some(Ok(AgentEvent::ModesReplaced { .. }))
4902 ));
4903 assert!(matches!(
4904 adapter.next_event().await,
4905 Some(Ok(AgentEvent::Ready { .. }))
4906 ));
4907 assert!(matches!(
4908 adapter.next_event().await,
4909 Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
4910 ));
4911 adapter.stop().await.expect("stop");
4912 }
4913
4914 #[tokio::test]
4915 async fn acp_reload_preserves_a_loadable_session_id() {
4916 let cwd = std::env::current_dir().expect("cwd");
4917 let mut adapter = AcpAdapter::with_session_id(
4918 0,
4919 cwd,
4920 "__codeswarm_missing_acp_for_reload_test__",
4921 Vec::new(),
4922 "saved-session",
4923 );
4924 adapter.capabilities.supports_session_load = true;
4925 assert!(adapter.reload().await.is_err());
4929 assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
4930 }
4931
4932 #[tokio::test]
4933 async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
4934 let cwd = std::env::current_dir().expect("cwd");
4935 let mut adapter = AcpAdapter::with_session_id(
4936 0,
4937 cwd,
4938 "__codeswarm_missing_nonloadable_acp__",
4939 Vec::new(),
4940 "stale-session",
4941 );
4942 adapter.capabilities.supports_session_load = false;
4943 assert!(adapter.reload().await.is_err());
4944 assert_eq!(adapter.session_id, None);
4945 }
4946
4947 #[tokio::test]
4948 async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
4949 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"s1"}}'; read _; echo 'diagnostic from wrapper'; echo '{"jsonrpc":"2.0","method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"partial"}}}}'; echo '{"jsonrpc":"2.0","id":3,"error":{"code":-32000,"message":"capacity"}}'"#;
4950 let cwd = std::env::current_dir().expect("cwd");
4951 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4952 adapter.start().await.expect("initialize");
4953 assert!(matches!(
4954 adapter.next_event().await,
4955 Some(Ok(AgentEvent::Ready { .. }))
4956 ));
4957 adapter.send_prompt("hello".into()).await.expect("prompt");
4958 assert!(matches!(
4959 adapter.next_event().await,
4960 Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
4961 ));
4962 assert!(matches!(
4963 adapter.next_event().await,
4964 Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
4965 ));
4966 }
4967
4968 #[tokio::test]
4969 async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
4970 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
4971 let cwd = std::env::current_dir().expect("cwd");
4972 let mut adapter = AcpAdapter::with_session_id(
4973 0,
4974 cwd,
4975 "sh",
4976 vec!["-c".into(), script.into()],
4977 "existing-session",
4978 );
4979 adapter.start().await.expect("load existing session");
4980 assert!(matches!(
4981 adapter.next_event().await,
4982 Some(Ok(AgentEvent::Ready { .. }))
4983 ));
4984 }
4985
4986 #[tokio::test]
4987 async fn acp_start_failure_reaps_transport_process() {
4988 let mut adapter = AcpAdapter::new(
4993 0,
4994 std::env::current_dir().expect("cwd"),
4995 "sh",
4996 vec!["-c".into(), "printf 'not-json\\n'".into()],
4997 );
4998 assert!(adapter.start().await.is_err());
4999 assert!(adapter.child.is_none());
5000 assert!(adapter.reader.is_none());
5001 }
5002
5003 #[tokio::test]
5004 async fn acp_adapter_answers_permission_json_rpc_requests() {
5005 let path = std::env::temp_dir().join(format!(
5006 "codeswarm-permission-answer-{}",
5007 std::process::id()
5008 ));
5009 let script = format!(
5010 r#"read _; echo '{{"jsonrpc":"2.0","id":1,"result":{{"agentCapabilities":{{}}}}}}'; read _; echo '{{"jsonrpc":"2.0","id":2,"result":{{"sessionId":"s1"}}}}'; read _; echo '{{"jsonrpc":"2.0","id":9,"method":"session/request_permission","params":{{"toolCall":{{"title":"Write file"}},"options":[{{"optionId":"allow-once","name":"Allow once"}}]}}}}'; read answer; printf '%s' "$answer" > '{}'; echo '{{"jsonrpc":"2.0","id":3,"result":{{"stopReason":"end_turn"}}}}'"#,
5011 path.display()
5012 );
5013 let mut adapter = AcpAdapter::new(
5014 0,
5015 std::env::current_dir().expect("cwd"),
5016 "sh",
5017 vec!["-c".into(), script],
5018 );
5019 adapter.start().await.expect("start ACP");
5020 assert!(matches!(
5021 adapter.next_event().await,
5022 Some(Ok(AgentEvent::Ready { .. }))
5023 ));
5024 adapter.send_prompt("do it".into()).await.expect("prompt");
5025 assert!(matches!(
5026 adapter.next_event().await,
5027 Some(Ok(AgentEvent::Permission { request, .. }))
5028 if request.id == "9"
5029 && request.options == ["Allow once"]
5030 && request.option_ids == ["allow-once"]
5031 ));
5032 adapter
5033 .answer_permission(
5034 "9".into(),
5035 PermissionAnswer::Selected {
5036 option_id: "allow-once".into(),
5037 },
5038 )
5039 .await
5040 .expect("permission answer");
5041 assert!(matches!(
5042 adapter.next_event().await,
5043 Some(Ok(AgentEvent::TurnComplete { .. }))
5044 ));
5045 let answer: Value = serde_json::from_str(
5046 &std::fs::read_to_string(&path).expect("captured permission answer"),
5047 )
5048 .expect("valid JSON-RPC answer");
5049 assert_eq!(answer["id"], 9);
5050 assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
5051 assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
5052 std::fs::remove_file(path).expect("cleanup");
5053 }
5054
5055 #[test]
5056 fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
5057 let event = parse_acp_notification(
5058 0,
5059 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
5060 )
5061 .expect("valid JSON-RPC request");
5062 assert!(event.is_none());
5063 }
5064
5065 #[tokio::test]
5066 async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
5067 let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
5068 std::fs::write(
5069 &script_path,
5070 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
5071 )
5072 .expect("write native test script");
5073 let mut adapter = AgyAdapter::new(
5074 0,
5075 std::env::current_dir().expect("cwd"),
5076 format!("sh {}", script_path.display()),
5077 );
5078 adapter.start().await.expect("start native adapter");
5079 assert!(adapter.next_event().await.is_some());
5080 assert!(adapter.next_event().await.is_some());
5081 adapter
5082 .send_prompt("continue".into())
5083 .await
5084 .expect("prompt");
5085 assert!(matches!(
5086 adapter.next_event().await,
5087 Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
5088 ));
5089 assert!(matches!(
5090 adapter.next_event().await,
5091 Some(Ok(AgentEvent::TurnComplete { .. }))
5092 ));
5093 adapter.stop().await.expect("stop native adapter");
5094 std::fs::remove_file(script_path).expect("cleanup native script");
5095 }
5096
5097 #[test]
5098 fn acp_workspace_file_access_is_root_bound_and_size_limited() {
5099 let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
5100 let _ = std::fs::remove_dir_all(&root);
5101 std::fs::create_dir_all(&root).expect("workspace");
5102 std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
5103 let outside =
5104 std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
5105 std::fs::write(&outside, "secret").expect("outside file");
5106 let link = root.join("outside-link");
5107 #[cfg(unix)]
5108 std::os::unix::fs::symlink(&outside, &link).expect("symlink");
5109 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5110
5111 assert_eq!(
5112 adapter
5113 .read_workspace_text("inside.txt", Some(2), Some(1))
5114 .expect("read inside"),
5115 "two"
5116 );
5117 std::fs::write(
5118 root.join("large.txt"),
5119 vec![b'x'; MAX_FILE_READ_BYTES + 1024],
5120 )
5121 .expect("large file");
5122 let bounded = adapter
5123 .read_workspace_text("large.txt", None, None)
5124 .expect("bounded read");
5125 assert!(bounded.len() <= MAX_FILE_READ_BYTES);
5126 #[cfg(unix)]
5127 {
5128 std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
5129 .expect("internal symlink");
5130 assert_eq!(
5131 adapter
5132 .read_workspace_text("inside-link", None, None)
5133 .expect("read internal symlink"),
5134 "one\ntwo\nthree\n"
5135 );
5136 }
5137 assert!(adapter.workspace_path("../codeswarm-outside").is_err());
5138 assert!(
5139 adapter
5140 .workspace_path(&outside.display().to_string())
5141 .is_err()
5142 );
5143 #[cfg(unix)]
5144 assert!(adapter.workspace_path("outside-link").is_err());
5145 #[cfg(unix)]
5146 std::fs::remove_file(link).expect("cleanup symlink");
5147 #[cfg(unix)]
5148 std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
5149 std::fs::remove_file(outside).expect("cleanup outside");
5150 std::fs::remove_dir_all(root).expect("cleanup workspace");
5151 }
5152
5153 #[tokio::test]
5154 async fn running_terminal_output_omits_exit_status_until_completion() {
5155 let root = unique_test_path("codeswarm-terminal-output", "dir");
5156 std::fs::create_dir_all(&root).expect("workspace");
5157 let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5158 let result = adapter
5159 .terminal_create(&serde_json::json!({
5160 "command": "sh",
5161 "args": ["-c", "sleep 0.2; printf done"],
5162 "cwd": ".",
5163 }))
5164 .await
5165 .expect("terminal create");
5166 let id = result["terminalId"].as_str().expect("terminal id");
5167 let output = adapter.terminal_output(id).await.expect("terminal output");
5168 assert!(output.get("exitStatus").is_none());
5169 if let Some(terminal) = adapter.terminals.remove(id) {
5170 terminal.stop().await;
5171 }
5172 std::fs::remove_dir_all(root).expect("cleanup workspace");
5173 }
5174
5175 #[tokio::test]
5176 async fn acp_adapter_answers_workspace_read_requests() {
5177 let root =
5178 std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
5179 let _ = std::fs::remove_dir_all(&root);
5180 std::fs::create_dir_all(&root).expect("workspace");
5181 let source = root.join("inside.txt");
5182 let answer = root.join("answer.json");
5183 std::fs::write(&source, "workspace content").expect("source");
5184 let script = format!(
5185 r#"read _; echo '{{"jsonrpc":"2.0","id":1,"result":{{"agentCapabilities":{{}}}}}}'; read _; echo '{{"jsonrpc":"2.0","id":2,"result":{{"sessionId":"s1"}}}}'; read _; echo '{{"jsonrpc":"2.0","id":9,"method":"fs/read_text_file","params":{{"sessionId":"s1","path":"{}"}}}}'; read response; printf '%s' "$response" > '{}'; echo '{{"jsonrpc":"2.0","id":3,"result":{{"stopReason":"end_turn"}}}}'"#,
5186 source.display(),
5187 answer.display(),
5188 );
5189 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
5190 adapter.start().await.expect("start ACP");
5191 assert!(matches!(
5192 adapter.next_event().await,
5193 Some(Ok(AgentEvent::Ready { .. }))
5194 ));
5195 adapter.send_prompt("read it".into()).await.expect("prompt");
5196 assert!(matches!(
5197 adapter.next_event().await,
5198 Some(Ok(AgentEvent::TurnComplete { .. }))
5199 ));
5200 let response: Value =
5201 serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
5202 .expect("response JSON");
5203 assert_eq!(response["id"], 9);
5204 assert_eq!(response["result"]["content"], "workspace content");
5205 adapter.stop().await.expect("stop ACP");
5206 std::fs::remove_dir_all(root).expect("cleanup workspace");
5207 }
5208
5209 #[tokio::test]
5210 async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
5211 let root =
5212 std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
5213 let _ = std::fs::remove_dir_all(&root);
5214 std::fs::create_dir_all(&root).expect("workspace");
5215 let create_request = serde_json::json!({
5216 "jsonrpc": "2.0",
5217 "id": 9,
5218 "method": "terminal/create",
5219 "params": {
5220 "sessionId": "s1",
5221 "command": "sh",
5222 "args": ["-c", "sleep 0.1; printf terminal-ok"],
5223 "cwd": ".",
5224 },
5225 });
5226 let wait_request = serde_json::json!({
5227 "jsonrpc": "2.0",
5228 "id": 10,
5229 "method": "terminal/wait_for_exit",
5230 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
5231 });
5232 let output_request = serde_json::json!({
5233 "jsonrpc": "2.0",
5234 "id": 11,
5235 "method": "terminal/output",
5236 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
5237 });
5238 let create_answer = root.join("create-answer.json");
5239 let wait_answer = root.join("wait-answer.json");
5240 let output_answer = root.join("output-answer.json");
5241 let script = format!(
5242 "read _; echo '{{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":{{\"agentCapabilities\":{{}}}}}}'; read _; echo '{{\"jsonrpc\":\"2.0\",\"id\":2,\"result\":{{\"sessionId\":\"s1\"}}}}'; read _; echo '{}'; read response; printf '%s' \"$response\" > '{}'; echo '{}'; read response; printf '%s' \"$response\" > '{}'; echo '{}'; read response; printf '%s' \"$response\" > '{}'; echo '{{\"jsonrpc\":\"2.0\",\"id\":3,\"result\":{{\"stopReason\":\"end_turn\"}}}}'",
5243 create_request,
5244 create_answer.display(),
5245 wait_request,
5246 wait_answer.display(),
5247 output_request,
5248 output_answer.display(),
5249 );
5250 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
5251 adapter.start().await.expect("start ACP");
5252 assert!(matches!(
5253 adapter.next_event().await,
5254 Some(Ok(AgentEvent::Ready { .. }))
5255 ));
5256 adapter
5257 .send_prompt("run terminal".into())
5258 .await
5259 .expect("prompt");
5260 let mut saw_complete = false;
5261 for _ in 0..6 {
5262 match adapter.next_event().await {
5263 Some(Ok(AgentEvent::TurnComplete { .. })) => {
5264 saw_complete = true;
5265 break;
5266 }
5267 Some(_) => {}
5268 None => break,
5269 }
5270 }
5271 assert!(saw_complete, "terminal requests should not stall ACP");
5272 let create: Value = serde_json::from_str(
5273 &std::fs::read_to_string(&create_answer).expect("captured create response"),
5274 )
5275 .expect("create JSON");
5276 assert_eq!(create["result"]["terminalId"], "terminal-1");
5277 let output: Value = serde_json::from_str(
5278 &std::fs::read_to_string(&output_answer).expect("captured output response"),
5279 )
5280 .expect("output JSON");
5281 assert!(
5282 output["result"]["output"]
5283 .as_str()
5284 .unwrap_or_default()
5285 .contains("terminal-ok"),
5286 "output response: {output}"
5287 );
5288 adapter.stop().await.expect("stop ACP");
5289 std::fs::remove_dir_all(root).expect("cleanup workspace");
5290 }
5291
5292 #[tokio::test]
5293 async fn host_reduces_and_persists_adapter_events() {
5294 let path =
5295 std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
5296 let adapter = ScriptedAdapter::new(
5297 0,
5298 AgentCapabilities::default(),
5299 [AgentEvent::Text {
5300 slot: 0,
5301 text: "hello".into(),
5302 }],
5303 );
5304 let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
5305 host.start().await.expect("start");
5306 host.next_effects()
5307 .await
5308 .expect("event")
5309 .expect("valid event");
5310 assert_eq!(host.state.public_text[0].1, "hello");
5311 assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
5312 std::fs::remove_file(path).expect("cleanup");
5313 }
5314
5315 #[tokio::test]
5316 async fn relay_applies_default_policy_before_the_first_prompt() {
5317 let first_log = Arc::new(Mutex::new(Vec::new()));
5318 let second_log = Arc::new(Mutex::new(Vec::new()));
5319 let first = AdapterHost::new(
5320 Box::new(ModeOrderAdapter {
5321 slot: 0,
5322 log: Arc::clone(&first_log),
5323 phase: 0,
5324 }),
5325 None,
5326 );
5327 let second = AdapterHost::new(
5328 Box::new(ModeOrderAdapter {
5329 slot: 1,
5330 log: Arc::clone(&second_log),
5331 phase: 0,
5332 }),
5333 None,
5334 );
5335 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
5336 relay.start().await.expect("start and synchronize policy");
5337 relay.run_turn("task", 0).await.expect("first turn");
5338 {
5339 let log = first_log.lock().expect("log");
5340 assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
5341 }
5342 assert_eq!(
5343 second_log.lock().expect("log").as_slice(),
5344 ["start", "mode:yolo"]
5345 );
5346 let added_log = Arc::new(Mutex::new(Vec::new()));
5347 relay
5348 .add_agent(
5349 AdapterHost::new(
5350 Box::new(ModeOrderAdapter {
5351 slot: 2,
5352 log: Arc::clone(&added_log),
5353 phase: 0,
5354 }),
5355 None,
5356 ),
5357 "Added",
5358 "added.example",
5359 "added-agent",
5360 )
5361 .await
5362 .expect("add with synchronized policy");
5363 assert_eq!(
5364 added_log.lock().expect("log").as_slice(),
5365 ["start", "mode:yolo"]
5366 );
5367 relay.drop_agent(2).await.expect("drop added agent");
5368 added_log.lock().expect("log").clear();
5369 relay
5370 .reload(2)
5371 .await
5372 .expect("reload with synchronized policy");
5373 assert_eq!(
5374 added_log.lock().expect("log").as_slice(),
5375 ["reload", "mode:yolo"]
5376 );
5377 }
5378
5379 #[tokio::test]
5380 async fn acp_roster_is_ready_before_any_prompt_is_sent() {
5381 let hosts = (0..2)
5382 .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
5383 .collect::<Vec<_>>();
5384 let startup_events = Arc::new(Mutex::new(Vec::new()));
5385 let captured = Arc::clone(&startup_events);
5386 let mut relay = RelayHost::new(hosts, 4).expect("relay");
5387 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5388
5389 relay.start().await.expect("complete startup handshake");
5390
5391 assert!(relay.dispatches().is_empty());
5392 let ready_slots = startup_events
5393 .lock()
5394 .expect("events")
5395 .iter()
5396 .filter_map(|event| match event {
5397 AgentEvent::Ready { slot, .. } => Some(*slot),
5398 _ => None,
5399 })
5400 .collect::<Vec<_>>();
5401 assert_eq!(ready_slots, vec![0, 1]);
5402 }
5403
5404 #[tokio::test]
5405 async fn independent_roster_adapters_start_concurrently() {
5406 let barrier = Arc::new(tokio::sync::Barrier::new(2));
5407 let hosts = (0..2)
5408 .map(|slot| {
5409 AdapterHost::new(
5410 Box::new(ConcurrentStartAdapter {
5411 slot,
5412 barrier: Arc::clone(&barrier),
5413 }),
5414 None,
5415 )
5416 })
5417 .collect::<Vec<_>>();
5418 let mut relay = RelayHost::new(hosts, 4).expect("relay");
5419 tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
5420 .await
5421 .expect("startup should not serialize barrier participants")
5422 .expect("startup succeeds");
5423 }
5424
5425 #[tokio::test]
5426 async fn relay_host_dispatches_turns_sequentially() {
5427 let capabilities = AgentCapabilities {
5428 supports_cancel: true,
5429 ..AgentCapabilities::default()
5430 };
5431 let first = ScriptedAdapter::new(
5432 0,
5433 capabilities.clone(),
5434 [
5435 AgentEvent::Text {
5436 slot: 0,
5437 text: "first".into(),
5438 },
5439 AgentEvent::TurnComplete { slot: 0 },
5440 ],
5441 );
5442 let second = ScriptedAdapter::new(
5443 1,
5444 capabilities,
5445 [
5446 AgentEvent::Text {
5447 slot: 1,
5448 text: "review".into(),
5449 },
5450 AgentEvent::TurnComplete { slot: 1 },
5451 ],
5452 );
5453 let hosts = vec![
5454 AdapterHost::new(Box::new(first), None),
5455 AdapterHost::new(Box::new(second), None),
5456 ];
5457 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5458 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
5459 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5460 let captured = std::sync::Arc::clone(&events);
5461 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5462 relay.start().await.expect("start");
5463 events.lock().expect("events").clear();
5464 assert!(matches!(
5465 relay.run_turn("task", 0).await.expect("first turn"),
5466 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5467 ));
5468 assert!(matches!(
5469 relay.run_turn("first", 0).await.expect("second turn"),
5470 crate::relay::RelayDecision::Dispatch {
5471 slot: 1,
5472 can_stop: true,
5473 ..
5474 }
5475 ));
5476 assert_eq!(
5477 relay
5478 .dispatches()
5479 .iter()
5480 .map(|(slot, _)| *slot)
5481 .collect::<Vec<_>>(),
5482 [0, 1]
5483 );
5484 assert!(relay.dispatches()[0].1.contains("You are Claude"));
5485 assert!(
5486 relay.dispatches()[0]
5487 .1
5488 .contains("CodeSwarm roster (ordered)")
5489 );
5490 assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
5491 assert!(relay.dispatches()[0].1.contains("2. Codex"));
5492 assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
5493 assert!(relay.dispatches()[0].1.contains("Do not use"));
5494 let lifecycle = events.lock().expect("events");
5495 let positions = lifecycle
5496 .iter()
5497 .filter_map(|event| match event {
5498 AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
5499 AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
5500 _ => None,
5501 })
5502 .collect::<Vec<_>>();
5503 assert_eq!(
5504 positions,
5505 [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
5506 );
5507 }
5508
5509 #[tokio::test]
5510 async fn relay_host_routes_around_a_usage_limited_agent() {
5511 let capabilities = AgentCapabilities::default();
5512 let first = ScriptedAdapter::new(
5513 0,
5514 capabilities.clone(),
5515 [
5516 AgentEvent::Text {
5517 slot: 0,
5518 text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
5519 credits or try again later."
5520 .into(),
5521 },
5522 AgentEvent::TurnComplete { slot: 0 },
5523 ],
5524 );
5525 let second = ScriptedAdapter::new(
5526 1,
5527 capabilities,
5528 [
5529 AgentEvent::Text {
5530 slot: 1,
5531 text: "review done".into(),
5532 },
5533 AgentEvent::TurnComplete { slot: 1 },
5534 ],
5535 );
5536 let hosts = vec![
5537 AdapterHost::new(Box::new(first), None),
5538 AdapterHost::new(Box::new(second), None),
5539 ];
5540 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5541 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5542 let captured = std::sync::Arc::clone(&events);
5543 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5544 relay.start().await.expect("start");
5545 events.lock().expect("events").clear();
5546 assert!(matches!(
5547 relay.run_turn("task", 0).await.expect("limited turn"),
5548 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5549 ));
5550 assert!(
5551 events
5552 .lock()
5553 .expect("events")
5554 .iter()
5555 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
5556 );
5557 assert!(matches!(
5559 relay.run_turn("", 0).await.expect("next turn"),
5560 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
5561 ));
5562 assert!(relay.relay().is_limited(0));
5563 relay.reload(0).await.expect("reload");
5567 assert!(!relay.relay().is_limited(0));
5568 }
5569
5570 #[tokio::test]
5571 async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
5572 let limited = ScriptedAdapter::new(
5573 0,
5574 AgentCapabilities::default(),
5575 [AgentEvent::Failed {
5576 slot: 0,
5577 started: true,
5578 detail: "request failed: insufficient_quota".into(),
5579 }],
5580 );
5581 let healthy = ScriptedAdapter::new(
5582 1,
5583 AgentCapabilities::default(),
5584 [AgentEvent::TurnComplete { slot: 1 }],
5585 );
5586 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5587 let captured = std::sync::Arc::clone(&events);
5588 let mut relay = RelayHost::new(
5589 vec![
5590 AdapterHost::new(Box::new(limited), None),
5591 AdapterHost::new(Box::new(healthy), None),
5592 ],
5593 4,
5594 )
5595 .expect("relay");
5596 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5597 relay.start().await.expect("start");
5598 events.lock().expect("events").clear();
5599
5600 assert!(matches!(
5601 relay.run_turn("task", 0).await.expect("limited failure"),
5602 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5603 ));
5604 assert!(relay.relay().is_limited(0));
5605 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
5606 {
5607 let events = events.lock().expect("events");
5608 assert!(
5609 events
5610 .iter()
5611 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
5612 );
5613 assert!(
5614 !events
5615 .iter()
5616 .any(|event| matches!(event, AgentEvent::Failed { .. }))
5617 );
5618 }
5619
5620 assert!(matches!(
5621 relay.run_turn("", 0).await.expect("healthy peer"),
5622 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
5623 ));
5624 }
5625
5626 #[tokio::test]
5627 async fn relay_failure_is_tombstoned_and_reported_to_the_ui_sink() {
5628 let failed = ScriptedAdapter::new(
5629 0,
5630 AgentCapabilities::default(),
5631 [AgentEvent::Failed {
5632 slot: 0,
5633 started: true,
5634 detail: "connection lost".into(),
5635 }],
5636 );
5637 let healthy = ScriptedAdapter::new(
5638 1,
5639 AgentCapabilities::default(),
5640 [AgentEvent::TurnComplete { slot: 1 }],
5641 );
5642 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5643 let captured = std::sync::Arc::clone(&events);
5644 let mut relay = RelayHost::new(
5645 vec![
5646 AdapterHost::new(Box::new(failed), None),
5647 AdapterHost::new(Box::new(healthy), None),
5648 ],
5649 4,
5650 )
5651 .expect("relay");
5652 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5653 relay.start().await.expect("start");
5654
5655 let error = relay.run_turn("task", 0).await.expect_err("failure");
5656 assert!(error.to_string().contains("connection lost"));
5657 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
5658 assert!(events.lock().expect("lock").iter().any(|event| {
5659 matches!(
5660 event,
5661 AgentEvent::Failed {
5662 slot: 0,
5663 started: true,
5664 ..
5665 }
5666 )
5667 }));
5668 }
5669
5670 #[tokio::test]
5671 async fn codex_stop_does_not_skip_later_roster_reviewers() {
5672 let hosts = (0..3)
5673 .map(|slot| {
5674 AdapterHost::new(
5675 Box::new(ScriptedAdapter::new(
5676 slot,
5677 AgentCapabilities::default(),
5678 [
5679 AgentEvent::Text {
5680 slot,
5681 text: STOP_TOKEN.into(),
5682 },
5683 AgentEvent::TurnComplete { slot },
5684 ],
5685 )),
5686 None,
5687 )
5688 })
5689 .collect();
5690 let mut relay = RelayHost::new(hosts, 10).expect("relay");
5691 relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
5692 relay.start().await.expect("start");
5693 for expected in 0..3 {
5694 assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
5695 RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
5696 }
5697 assert_eq!(
5698 relay.run_turn("", 0).await.expect("complete"),
5699 RelayDecision::Complete
5700 );
5701 }
5702
5703 #[tokio::test]
5704 async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
5705 let first = ScriptedAdapter::new(
5706 0,
5707 AgentCapabilities::default(),
5708 [
5709 AgentEvent::Text {
5710 slot: 0,
5711 text: "done".into(),
5712 },
5713 AgentEvent::TurnComplete { slot: 0 },
5714 ],
5715 );
5716 let reviewer = ScriptedAdapter::new(
5717 1,
5718 AgentCapabilities::default(),
5719 [
5720 AgentEvent::Text {
5721 slot: 1,
5722 text: STOP_TOKEN.into(),
5723 },
5724 AgentEvent::TurnComplete { slot: 1 },
5725 ],
5726 );
5727 let mut relay = RelayHost::new(
5728 vec![
5729 AdapterHost::new(Box::new(first), None),
5730 AdapterHost::new(Box::new(reviewer), None),
5731 ],
5732 10,
5733 )
5734 .expect("relay");
5735 relay.start().await.expect("start");
5736 let first_decision = relay.run_turn("task", 0).await.expect("first");
5737 assert!(matches!(
5738 first_decision,
5739 RelayDecision::Dispatch { slot: 0, .. }
5740 ));
5741 let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
5742 assert!(matches!(
5743 reviewer_decision,
5744 RelayDecision::Dispatch {
5745 slot: 1,
5746 can_stop: true,
5747 ..
5748 }
5749 ));
5750 assert_eq!(
5751 relay.run_turn("", 0).await.expect("complete"),
5752 RelayDecision::Complete
5753 );
5754 }
5755
5756 #[tokio::test]
5757 async fn stop_token_is_filtered_from_streamed_ui_events() {
5758 let first = ScriptedAdapter::new(
5759 0,
5760 AgentCapabilities::default(),
5761 [
5762 AgentEvent::Text {
5763 slot: 0,
5764 text: format!("visible {STOP_TOKEN} trailing"),
5765 },
5766 AgentEvent::TurnComplete { slot: 0 },
5767 ],
5768 );
5769 let reviewer = ScriptedAdapter::new(
5770 1,
5771 AgentCapabilities::default(),
5772 [AgentEvent::TurnComplete { slot: 1 }],
5773 );
5774 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5775 let captured = std::sync::Arc::clone(&events);
5776 let mut relay = RelayHost::new(
5777 vec![
5778 AdapterHost::new(Box::new(first), None),
5779 AdapterHost::new(Box::new(reviewer), None),
5780 ],
5781 2,
5782 )
5783 .expect("relay");
5784 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5785 relay.start().await.expect("start");
5786 relay.run_turn("task", 0).await.expect("turn");
5787 let captured = events.lock().expect("lock");
5788 assert!(captured.iter().all(|event| match event {
5789 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5790 _ => true,
5791 }));
5792 let visible = captured
5793 .iter()
5794 .filter_map(|event| match event {
5795 AgentEvent::Text { text, .. } => Some(text.as_str()),
5796 _ => None,
5797 })
5798 .collect::<String>();
5799 assert_eq!(visible, "visible trailing");
5800 }
5801
5802 #[tokio::test]
5803 async fn token_only_reviewer_response_emits_visible_acknowledgment() {
5804 let first = ScriptedAdapter::new(
5805 0,
5806 AgentCapabilities::default(),
5807 [
5808 AgentEvent::Text {
5809 slot: 0,
5810 text: "done".into(),
5811 },
5812 AgentEvent::TurnComplete { slot: 0 },
5813 ],
5814 );
5815 let reviewer = ScriptedAdapter::new(
5816 1,
5817 AgentCapabilities::default(),
5818 [
5819 AgentEvent::Text {
5820 slot: 1,
5821 text: STOP_TOKEN.into(),
5822 },
5823 AgentEvent::TurnComplete { slot: 1 },
5824 ],
5825 );
5826 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5827 let captured = std::sync::Arc::clone(&events);
5828 let mut relay = RelayHost::new(
5829 vec![
5830 AdapterHost::new(Box::new(first), None),
5831 AdapterHost::new(Box::new(reviewer), None),
5832 ],
5833 4,
5834 )
5835 .expect("relay");
5836 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5837 relay.start().await.expect("start");
5838 relay.run_turn("task", 0).await.expect("first turn");
5839 relay.run_turn("", 0).await.expect("review turn");
5840 let captured = events.lock().expect("lock");
5841 assert!(captured.iter().any(|event| {
5842 matches!(
5843 event,
5844 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
5845 )
5846 }));
5847 assert!(captured.iter().all(|event| match event {
5848 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5849 _ => true,
5850 }));
5851 let acknowledgment = captured
5852 .iter()
5853 .position(|event| {
5854 matches!(
5855 event,
5856 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
5857 )
5858 })
5859 .expect("visible acknowledgment");
5860 let completion = captured
5861 .iter()
5862 .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
5863 .expect("reviewer completion");
5864 assert!(acknowledgment < completion);
5865 }
5866
5867 #[tokio::test]
5868 async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
5869 let first = ScriptedAdapter::new(
5870 0,
5871 AgentCapabilities::default(),
5872 [
5873 AgentEvent::Text {
5874 slot: 0,
5875 text: "done".into(),
5876 },
5877 AgentEvent::TurnComplete { slot: 0 },
5878 ],
5879 );
5880 let reviewer = ScriptedAdapter::new(
5881 1,
5882 AgentCapabilities::default(),
5883 [
5884 AgentEvent::Text {
5885 slot: 1,
5886 text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
5887 },
5888 AgentEvent::TurnComplete { slot: 1 },
5889 ],
5890 );
5891 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5892 let captured = std::sync::Arc::clone(&events);
5893 let mut relay = RelayHost::new(
5894 vec![
5895 AdapterHost::new(Box::new(first), None),
5896 AdapterHost::new(Box::new(reviewer), None),
5897 ],
5898 4,
5899 )
5900 .expect("relay");
5901 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5902 relay.start().await.expect("start");
5903 relay.run_turn("task", 0).await.expect("first turn");
5904 relay.run_turn("", 0).await.expect("review turn");
5905
5906 let visible = events
5907 .lock()
5908 .expect("lock")
5909 .iter()
5910 .filter_map(|event| match event {
5911 AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
5912 _ => None,
5913 })
5914 .collect::<String>();
5915 assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
5916 assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
5917 }
5918
5919 #[tokio::test]
5920 async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
5921 let first = AdapterHost::new(
5922 Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
5923 None,
5924 );
5925 let second = AdapterHost::new(
5926 Box::new(ScriptedAdapter::new(
5927 1,
5928 AgentCapabilities::default(),
5929 [AgentEvent::TurnComplete { slot: 1 }],
5930 )),
5931 None,
5932 );
5933 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
5934 let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
5935 relay.set_event_sink(move |event| {
5936 if matches!(event, AgentEvent::Permission { .. }) {
5937 let _ = seen_sender.send(());
5938 }
5939 });
5940 relay.start().await.expect("start");
5941 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
5942 let answer = async move {
5943 seen_receiver.recv().await.expect("permission request");
5944 sender
5945 .send(super::RelayPermissionAnswer {
5946 slot: 0,
5947 request_id: "permission-1".into(),
5948 answer: PermissionAnswer::Selected {
5949 option_id: "allow".into(),
5950 },
5951 })
5952 .expect("queue permission answer");
5953 };
5954 tokio::time::timeout(std::time::Duration::from_millis(100), async {
5955 let ((), result) = tokio::join!(
5956 answer,
5957 relay.run_turn_with_permissions("task", 0, &mut receiver)
5958 );
5959 result
5960 })
5961 .await
5962 .expect("permission-gated turn should not deadlock")
5963 .expect("turn completes");
5964 }
5965
5966 #[tokio::test]
5967 async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
5968 let first = AdapterHost::new(
5969 Box::new(PendingAdapter {
5970 slot: 0,
5971 hang_on_cancel: false,
5972 }),
5973 None,
5974 );
5975 let second = AdapterHost::new(
5976 Box::new(ScriptedAdapter::new(
5977 1,
5978 AgentCapabilities::default(),
5979 [AgentEvent::TurnComplete { slot: 1 }],
5980 )),
5981 None,
5982 );
5983 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
5984 relay.start().await.expect("start");
5985 let cancellation = relay.cancellation();
5986 let error = {
5987 let turn = relay.run_turn("task", 0);
5988 tokio::pin!(turn);
5989 cancellation.request();
5990 turn.await.expect_err("cancellation should stop turn")
5991 };
5992 assert!(error.to_string().contains("relay turn cancelled"));
5993
5994 assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
5995 relay
5996 .run_turn("", 1)
5997 .await
5998 .expect("replacement job reaches the selected peer");
5999 let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
6000 assert!(replacement.contains("replacement job"));
6001 assert!(replacement.contains("User "));
6002 assert!(replacement.contains(":\ntask"));
6003 let owner_updates = relay.relay_mut().unseen_context(0);
6004 assert!(owner_updates.contains("User "));
6005 assert!(owner_updates.contains(":\ntask"));
6006 assert!(owner_updates.contains(":\nreplacement job"));
6007 }
6008
6009 #[tokio::test]
6010 async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
6011 let first = AdapterHost::new(
6012 Box::new(PendingAdapter {
6013 slot: 0,
6014 hang_on_cancel: true,
6015 }),
6016 None,
6017 );
6018 let second = AdapterHost::new(
6019 Box::new(PendingAdapter {
6020 slot: 1,
6021 hang_on_cancel: false,
6022 }),
6023 None,
6024 );
6025 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6026 relay.start().await.expect("start");
6027 let cancellation = relay.cancellation();
6028 let turn = relay.run_turn("task", 0);
6029 tokio::pin!(turn);
6030 cancellation.request();
6031 let error = turn.await.expect_err("cancellation should stop turn");
6032 assert!(error.to_string().contains("timed out"));
6033 }
6034
6035 #[tokio::test]
6036 async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
6037 let event = [AgentEvent::TurnComplete { slot: 0 }];
6038 let first = AdapterHost::new(
6039 Box::new(ScriptedAdapter::new(
6040 0,
6041 AgentCapabilities::default(),
6042 event.clone(),
6043 )),
6044 None,
6045 );
6046 let second = AdapterHost::new(
6047 Box::new(ScriptedAdapter::new(
6048 1,
6049 AgentCapabilities::default(),
6050 [AgentEvent::TurnComplete { slot: 1 }],
6051 )),
6052 None,
6053 );
6054 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6055 relay.start().await.expect("start");
6056
6057 relay.pause();
6058 assert_eq!(
6059 relay.run_turn("paused", 0).await.expect("paused turn"),
6060 crate::relay::RelayDecision::Paused
6061 );
6062 assert!(relay.dispatches().is_empty());
6063
6064 relay.resume();
6065 relay.relay_mut().drop_agent(1).expect("drop reviewer");
6066 assert!(matches!(
6067 relay
6068 .run_turn("solo follow-up", 0)
6069 .await
6070 .expect("solo turn"),
6071 crate::relay::RelayDecision::Dispatch {
6072 slot: 0,
6073 can_stop: false,
6074 ..
6075 }
6076 ));
6077 assert_eq!(relay.dispatches().len(), 1);
6078 }
6079
6080 #[tokio::test]
6081 async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
6082 let first = AdapterHost::new(
6083 Box::new(ScriptedAdapter::new(
6084 0,
6085 AgentCapabilities::default(),
6086 [AgentEvent::TurnComplete { slot: 0 }],
6087 )),
6088 None,
6089 );
6090 let second = AdapterHost::new(
6091 Box::new(ScriptedAdapter::new(
6092 1,
6093 AgentCapabilities::default(),
6094 [AgentEvent::TurnComplete { slot: 1 }],
6095 )),
6096 None,
6097 );
6098 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6099 relay.set_roster_names(vec!["First".into(), "Second".into()]);
6100 relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
6101 relay.set_roster_launch_specs(vec![
6102 ("custom".into(), "owner".into()),
6103 ("custom".into(), "peer".into()),
6104 ]);
6105 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6106 let captured = std::sync::Arc::clone(&events);
6107 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6108 relay.start().await.expect("start");
6109 let slot = relay
6110 .add_agent(
6111 AdapterHost::new(
6112 Box::new(ScriptedAdapter::new(
6113 2,
6114 AgentCapabilities::default(),
6115 [AgentEvent::TurnComplete { slot: 2 }],
6116 )),
6117 None,
6118 ),
6119 "Reviewer",
6120 "reviewer.example",
6121 "reviewer --acp",
6122 )
6123 .await
6124 .expect("append agent");
6125 assert_eq!(slot, 2);
6126 assert_eq!(
6127 relay.relay().active_slots().collect::<Vec<_>>(),
6128 vec![0, 1, 2]
6129 );
6130 assert_eq!(
6131 relay
6132 .session_metadata()
6133 .get("agents")
6134 .and_then(|value| value.as_array())
6135 .map(Vec::len),
6136 Some(3)
6137 );
6138 relay.drop_agent(1).await.expect("drop middle peer");
6139 let metadata = relay.session_metadata();
6140 assert_eq!(
6141 metadata.get("agents"),
6142 Some(&serde_json::json!([
6143 {"name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
6144 {"name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
6145 ]))
6146 );
6147 assert!(
6148 events
6149 .lock()
6150 .expect("lock")
6151 .iter()
6152 .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
6153 );
6154 }
6155
6156 #[tokio::test]
6157 async fn relay_host_persists_coordinator_owned_runtime_metadata() {
6158 let path = unique_test_path("codeswarm-session-metadata", "json");
6159 let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
6160 let writer = metadata_store.buffered().expect("metadata writer");
6161 let first = AdapterHost::new(
6162 Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
6163 None,
6164 );
6165 let second = AdapterHost::new(
6166 Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
6167 None,
6168 );
6169 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6170 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6171 relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
6172 relay.set_roster_launch_specs(vec![
6173 ("custom".into(), "claude".into()),
6174 ("custom".into(), "codex".into()),
6175 ]);
6176 relay.set_session_metadata_writer(writer);
6177 relay.start().await.expect("start");
6178 relay.drop_agent(0).await.expect("drop first agent");
6179 relay.stop().await.expect("stop");
6180
6181 let loaded = metadata_store
6182 .read()
6183 .expect("read metadata")
6184 .expect("metadata snapshot");
6185 assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
6186 assert_eq!(
6187 loaded.get("agents"),
6188 Some(&serde_json::json!([{
6189 "name": "Codex", "identity": "openai.com", "protocol": "custom",
6190 "command": "codex", "supports_load_session": false
6191 }]))
6192 );
6193 assert!(loaded.get("owner").is_none());
6194 let _ = std::fs::remove_file(path);
6195 }
6196
6197 #[tokio::test]
6198 async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
6199 let first = AdapterHost::new(
6200 Box::new(ScriptedAdapter::new(
6201 0,
6202 AgentCapabilities::default(),
6203 [
6204 AgentEvent::Text {
6205 slot: 0,
6206 text: "owner stream".into(),
6207 },
6208 AgentEvent::TurnComplete { slot: 0 },
6209 ],
6210 )),
6211 None,
6212 );
6213 let second = AdapterHost::new(
6214 Box::new(ScriptedAdapter::new(
6215 1,
6216 AgentCapabilities::default(),
6217 [
6218 AgentEvent::Text {
6219 slot: 1,
6220 text: "peer stream".into(),
6221 },
6222 AgentEvent::TurnComplete { slot: 1 },
6223 ],
6224 )),
6225 None,
6226 );
6227 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6228 relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
6229 relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
6230 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6231 let captured = std::sync::Arc::clone(&events);
6232 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6233 relay.start().await.expect("start");
6234
6235 relay.swap_agents(0, 1).expect("swap peers");
6236 assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
6237 assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
6238 relay.run_turn("task", 0).await.expect("swapped turn");
6239 let events = events.lock().expect("events");
6240 assert!(events.iter().any(|event| {
6241 matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
6242 }));
6243 assert!(relay.dispatches()[0].1.contains("You are Peer"));
6244 }
6245
6246 #[tokio::test]
6247 async fn relay_host_persists_all_active_agent_metadata_off_thread() {
6248 let path = unique_test_path("codeswarm-session-metadata", "json");
6249 let first = AdapterHost::new(
6250 Box::new(ScriptedAdapter::new(
6251 0,
6252 AgentCapabilities::default(),
6253 [AgentEvent::TurnComplete { slot: 0 }],
6254 )),
6255 None,
6256 );
6257 let second = AdapterHost::new(
6258 Box::new(ScriptedAdapter::new(
6259 1,
6260 AgentCapabilities::default(),
6261 [AgentEvent::TurnComplete { slot: 1 }],
6262 )),
6263 None,
6264 );
6265 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6266 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6267 relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
6268 relay.set_roster_launch_specs(vec![
6269 ("custom".into(), "claude".into()),
6270 ("custom".into(), "codex".into()),
6271 ]);
6272 let writer = SessionMetadataStore::open(&path)
6273 .buffered()
6274 .expect("metadata writer");
6275 relay.set_session_metadata_writer(writer);
6276 relay.start().await.expect("start");
6277 relay.stop().await.expect("stop");
6278 let loaded = SessionMetadataStore::open(&path)
6279 .read()
6280 .expect("read metadata")
6281 .expect("metadata snapshot");
6282 let agents = loaded
6283 .get("agents")
6284 .and_then(|value| value.as_array())
6285 .expect("agents");
6286 assert_eq!(agents.len(), 2);
6287 assert_eq!(agents[0]["identity"], "claude.com");
6288 assert_eq!(agents[1]["identity"], "openai.com");
6289 let _ = std::fs::remove_file(path);
6290 }
6291
6292 #[tokio::test]
6293 async fn relay_host_routes_unseen_public_context_to_next_agent() {
6294 let first = AdapterHost::new(
6295 Box::new(ScriptedAdapter::new(
6296 0,
6297 AgentCapabilities::default(),
6298 [
6299 AgentEvent::Text {
6300 slot: 0,
6301 text: "implemented the fix".into(),
6302 },
6303 AgentEvent::TurnComplete { slot: 0 },
6304 ],
6305 )),
6306 None,
6307 );
6308 let second = AdapterHost::new(
6309 Box::new(ScriptedAdapter::new(
6310 1,
6311 AgentCapabilities::default(),
6312 [AgentEvent::TurnComplete { slot: 1 }],
6313 )),
6314 None,
6315 );
6316 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6317 relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
6318 relay.start().await.expect("start");
6319 relay.run_turn("task", 0).await.expect("first turn");
6320 relay.run_turn("review this", 0).await.expect("review turn");
6321
6322 assert_eq!(relay.dispatches().len(), 2);
6323 assert_eq!(relay.dispatches()[0].0, 0);
6324 assert!(relay.dispatches()[0].1.contains("task"));
6325 assert!(relay.dispatches()[0].1.contains("You are Codex"));
6326 assert!(relay.dispatches()[0].1.contains("2. Qwen"));
6327 assert_eq!(relay.dispatches()[1].0, 1);
6328 assert!(relay.dispatches()[1].1.contains("review this"));
6329 let public = relay.dispatches()[1]
6330 .1
6331 .split_once("Public updates:\n")
6332 .map(|(_, updates)| updates)
6333 .expect("review receives public context");
6334 let header = public
6335 .lines()
6336 .find(|line| line.starts_with("Codex "))
6337 .expect("named previous agent");
6338 let timestamp = header
6339 .strip_prefix("Codex ")
6340 .and_then(|value| value.strip_suffix(':'))
6341 .expect("timestamped header");
6342 assert_eq!(timestamp.len(), 5);
6343 assert_eq!(timestamp.as_bytes()[2], b':');
6344 assert!(
6345 timestamp
6346 .bytes()
6347 .enumerate()
6348 .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
6349 );
6350 assert!(public.contains("implemented the fix"));
6351 assert!(!public.contains("Agent 0"));
6352 }
6353}