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