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