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