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