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