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