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