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