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