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 report_relay_failure(&mut self.relay, &event_sink, *slot, true, error.to_string());
1445 let _ = self.queue_session_metadata();
1446 return Err(error);
1447 }
1448 if let Some(sink) = &self.event_sink {
1449 sink(AgentEvent::TurnStarted { slot: *slot });
1450 }
1451 if let Some(introduced) = self.introduced.get_mut(*slot) {
1452 *introduced = true;
1453 }
1454 self.dispatches.push((*slot, prompt));
1455 let mut response = String::new();
1456 let mut emitted_text = 0usize;
1457 let completion_event = loop {
1458 if self.cancel_requested.swap(false, Ordering::AcqRel) {
1459 if let Err(error) = cancel_with_timeout(host).await {
1460 report_relay_failure(
1461 &mut self.relay,
1462 &event_sink,
1463 *slot,
1464 true,
1465 error.to_string(),
1466 );
1467 let _ = self.queue_session_metadata();
1468 return Err(error);
1469 }
1470 return Err(AdapterError::Transport("relay turn cancelled".into()));
1471 }
1472 let update = tokio::select! {
1473 update = host.next_update() => match update {
1474 Some(Ok(update)) => update,
1475 Some(Err(error)) => {
1476 report_relay_failure(
1477 &mut self.relay,
1478 &event_sink,
1479 *slot,
1480 true,
1481 error.to_string(),
1482 );
1483 let _ = self.queue_session_metadata();
1484 return Err(error);
1485 }
1486 None => {
1487 let error = AdapterError::Transport("adapter ended during turn".into());
1488 report_relay_failure(
1489 &mut self.relay,
1490 &event_sink,
1491 *slot,
1492 true,
1493 error.to_string(),
1494 );
1495 let _ = self.queue_session_metadata();
1496 return Err(error);
1497 }
1498 },
1499 _ = self.cancel_notify.notified() => {
1500 if !self.cancel_requested.swap(false, Ordering::AcqRel) {
1501 continue;
1502 }
1503 if let Err(error) = cancel_with_timeout(host).await {
1504 report_relay_failure(
1505 &mut self.relay,
1506 &event_sink,
1507 *slot,
1508 true,
1509 error.to_string(),
1510 );
1511 let _ = self.queue_session_metadata();
1512 return Err(error);
1513 }
1514 return Err(AdapterError::Transport("relay turn cancelled".into()));
1515 },
1516 permission = async {
1517 match permissions.as_mut() {
1518 Some(receiver) => receiver.recv().await,
1519 None => std::future::pending().await,
1520 }
1521 } => {
1522 let Some(permission) = permission else {
1523 permissions = None;
1524 continue;
1525 };
1526 if permission.slot != *slot {
1527 return Err(AdapterError::Transport(
1528 "permission response targets an inactive relay slot".into(),
1529 ));
1530 }
1531 host.answer_permission(permission.request_id, permission.answer).await?;
1532 continue;
1533 },
1534 };
1535 match &update.event {
1536 AgentEvent::Text { text, .. } => response.push_str(text),
1537 AgentEvent::TurnComplete { .. } => {
1538 let visible_response = response.replace(STOP_TOKEN, "");
1539 let visible_start = emitted_text.min(visible_response.len());
1540 let visible_start = floor_char_boundary(&visible_response, visible_start);
1541 if visible_start < visible_response.len()
1542 && let Some(sink) = &self.event_sink
1543 {
1544 sink(AgentEvent::Text {
1545 slot: *slot,
1546 text: visible_response[visible_start..].to_owned(),
1547 });
1548 }
1549 self.cancel_requested.store(false, Ordering::Release);
1550 break update.event.clone();
1551 }
1552 AgentEvent::Failed {
1553 started, detail, ..
1554 } => {
1555 report_relay_failure(
1556 &mut self.relay,
1557 &event_sink,
1558 *slot,
1559 *started,
1560 detail.clone(),
1561 );
1562 let _ = self.queue_session_metadata();
1563 return Err(AdapterError::Transport(detail.clone()));
1564 }
1565 _ => {}
1566 }
1567 if let AgentEvent::Text { .. } = &update.event {
1568 let visible_end = response
1571 .find(STOP_TOKEN)
1572 .unwrap_or_else(|| response.len().saturating_sub(STOP_TOKEN.len() - 1));
1573 let visible_end = floor_char_boundary(&response, visible_end);
1574 if emitted_text < visible_end {
1575 if let Some(sink) = &self.event_sink {
1576 sink(AgentEvent::Text {
1577 slot: *slot,
1578 text: response[emitted_text..visible_end].to_owned(),
1579 });
1580 }
1581 emitted_text = visible_end;
1582 }
1583 } else if let Some(sink) = &self.event_sink {
1584 sink(update.event.clone());
1585 }
1586 };
1587 let (response, requested_stop) = strip_stop_token(&response);
1588 let response = response.replace(STOP_TOKEN, "");
1589 let accepted_stop = requested_stop && *can_stop;
1590 let needs_stop_acknowledgment = accepted_stop && response.is_empty();
1591 let response = if needs_stop_acknowledgment {
1592 DEFAULT_STOP_ACKNOWLEDGMENT.to_owned()
1593 } else {
1594 response
1595 };
1596 if needs_stop_acknowledgment && let Some(sink) = &self.event_sink {
1601 sink(AgentEvent::Text {
1602 slot: *slot,
1603 text: response.clone(),
1604 });
1605 }
1606 if let Some(sink) = &self.event_sink {
1607 sink(completion_event);
1608 }
1609 if is_usage_limit_response(&response) {
1612 let detail = response.clone();
1613 let _ = self.relay.mark_limited(*slot);
1614 self.relay.finish(*slot, *direct, false);
1617 self.queue_session_metadata()?;
1618 if let Some(sink) = &self.event_sink {
1619 sink(AgentEvent::UsageLimitReached {
1620 slot: *slot,
1621 detail,
1622 });
1623 }
1624 return Ok(decision);
1625 }
1626 if !*direct && !response.is_empty() {
1627 self.relay
1628 .record_public(public_context_speaker(&speaker_name), response);
1629 }
1630 self.relay.mark_context_seen(*slot);
1631 self.relay.finish(*slot, *direct, accepted_stop);
1632 self.queue_session_metadata()?;
1633 Ok(decision)
1634 }
1635}
1636
1637fn report_relay_failure(
1638 relay: &mut Relay,
1639 event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
1640 slot: RosterSlot,
1641 started: bool,
1642 detail: String,
1643) {
1644 if is_usage_limit_response(&detail) {
1647 let _ = relay.mark_limited(slot);
1648 if let Some(sink) = event_sink {
1649 sink(AgentEvent::UsageLimitReached {
1650 slot,
1651 detail: detail.clone(),
1652 });
1653 sink(AgentEvent::Failed {
1654 slot,
1655 started,
1656 detail,
1657 });
1658 }
1659 return;
1660 }
1661 let _ = relay.tombstone(slot);
1662 if let Some(sink) = event_sink {
1663 sink(AgentEvent::Failed {
1664 slot,
1665 started,
1666 detail,
1667 });
1668 }
1669}
1670
1671#[derive(Debug)]
1673pub struct ScriptedAdapter {
1674 slot: RosterSlot,
1675 capabilities: AgentCapabilities,
1676 events: VecDeque<AdapterResult<AgentEvent>>,
1677 prompts: Vec<String>,
1678}
1679
1680impl ScriptedAdapter {
1681 pub fn new(
1682 slot: RosterSlot,
1683 capabilities: AgentCapabilities,
1684 events: impl IntoIterator<Item = AgentEvent>,
1685 ) -> Self {
1686 Self {
1687 slot,
1688 capabilities,
1689 events: events.into_iter().map(Ok).collect(),
1690 prompts: Vec::new(),
1691 }
1692 }
1693
1694 pub fn prompts(&self) -> &[String] {
1695 &self.prompts
1696 }
1697}
1698
1699#[async_trait]
1700impl AgentAdapter for ScriptedAdapter {
1701 fn slot(&self) -> RosterSlot {
1702 self.slot
1703 }
1704
1705 fn capabilities(&self) -> AgentCapabilities {
1706 self.capabilities.clone()
1707 }
1708
1709 async fn start(&mut self) -> AdapterResult<()> {
1710 Ok(())
1711 }
1712
1713 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
1714 self.prompts.push(prompt);
1715 Ok(())
1716 }
1717
1718 async fn cancel(&mut self) -> AdapterResult<bool> {
1719 Ok(self.capabilities.supports_cancel)
1720 }
1721
1722 async fn answer_permission(
1723 &mut self,
1724 _request_id: String,
1725 _answer: PermissionAnswer,
1726 ) -> AdapterResult<()> {
1727 if self.capabilities.supports_permissions {
1728 Ok(())
1729 } else {
1730 Err(AdapterError::Unsupported("permission answer"))
1731 }
1732 }
1733
1734 async fn set_mode(&mut self, _mode: String) -> AdapterResult<()> {
1735 if self.capabilities.supports_modes {
1736 Ok(())
1737 } else {
1738 Err(AdapterError::Unsupported("set_mode"))
1739 }
1740 }
1741
1742 async fn reload(&mut self) -> AdapterResult<()> {
1743 Ok(())
1744 }
1745
1746 async fn stop(&mut self) -> AdapterResult<()> {
1747 Ok(())
1748 }
1749
1750 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
1751 self.events.pop_front()
1752 }
1753}
1754
1755#[derive(Debug)]
1758pub struct AgyAdapter {
1759 slot: RosterSlot,
1760 cwd: PathBuf,
1761 command: String,
1762 mode: String,
1763 mode_policy: String,
1764 session_id: Option<String>,
1765 child: Option<Child>,
1766 sender: mpsc::Sender<AdapterResult<AgentEvent>>,
1767 receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
1768 announced_session: Arc<Mutex<Option<String>>>,
1772 cancel_requested: Arc<AtomicBool>,
1773}
1774
1775impl AgyAdapter {
1776 pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
1777 let (sender, receiver) = mpsc::channel(256);
1778 Self {
1779 slot,
1780 cwd,
1781 command: command.into(),
1782 mode: "default".into(),
1783 mode_policy: "agy:full-access".into(),
1784 session_id: None,
1785 child: None,
1786 sender,
1787 receiver,
1788 announced_session: Arc::new(Mutex::new(None)),
1789 cancel_requested: Arc::new(AtomicBool::new(false)),
1790 }
1791 }
1792
1793 pub fn with_session_id(
1794 slot: RosterSlot,
1795 cwd: PathBuf,
1796 command: impl Into<String>,
1797 session_id: impl Into<String>,
1798 ) -> Self {
1799 let mut adapter = Self::new(slot, cwd, command);
1800 adapter.session_id = Some(session_id.into());
1801 adapter
1802 }
1803
1804 fn modes() -> Vec<Mode> {
1805 vec![
1806 Mode {
1807 id: "agy:full-access".into(),
1808 label: "Auto pilot".into(),
1809 },
1810 Mode {
1811 id: "agy:manual".into(),
1812 label: "Manual".into(),
1813 },
1814 Mode {
1815 id: "accept-edits".into(),
1816 label: "Accept Edits".into(),
1817 },
1818 Mode {
1819 id: "plan".into(),
1820 label: "Plan".into(),
1821 },
1822 ]
1823 }
1824
1825 async fn emit(&self, event: AdapterResult<AgentEvent>) {
1826 let _ = self.sender.send(event).await;
1827 }
1828}
1829
1830#[async_trait]
1831impl AgentAdapter for AgyAdapter {
1832 fn slot(&self) -> RosterSlot {
1833 self.slot
1834 }
1835
1836 fn session_id(&self) -> Option<String> {
1837 self.session_id.clone()
1838 }
1839
1840 fn protocol(&self) -> &'static str {
1841 "native"
1842 }
1843
1844 fn capabilities(&self) -> AgentCapabilities {
1845 AgentCapabilities {
1846 supports_cancel: true,
1847 supports_modes: true,
1848 supports_permissions: false,
1849 supports_terminals: true,
1850 supports_session_load: true,
1851 supports_models: false,
1852 }
1853 }
1854
1855 async fn start(&mut self) -> AdapterResult<()> {
1856 self.cancel_requested.store(false, Ordering::Release);
1857 self.emit(Ok(AgentEvent::Ready {
1858 slot: self.slot,
1859 capabilities: self.capabilities(),
1860 }))
1861 .await;
1862 self.emit(Ok(AgentEvent::ModesReplaced {
1863 slot: self.slot,
1864 modes: Self::modes(),
1865 current_mode: Some(self.mode_policy.clone()),
1866 }))
1867 .await;
1868 Ok(())
1869 }
1870
1871 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
1872 if self.child.is_some() {
1873 return Err(AdapterError::Transport(
1874 "agent is already handling a turn".into(),
1875 ));
1876 }
1877 if self.session_id.is_none()
1878 && let Ok(session) = self.announced_session.lock()
1879 {
1880 self.session_id = session.clone();
1881 }
1882 self.cancel_requested.store(false, Ordering::Release);
1883 let (program, args) = parse_command_line(&self.command)
1884 .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
1885 let mut command = Command::new(program);
1886 isolate_process_group(&mut command);
1887 command
1888 .args(args)
1889 .arg("--print")
1890 .arg(prompt)
1891 .arg("--print-timeout")
1892 .arg("60m")
1893 .arg("--output-format")
1894 .arg("stream-json")
1895 .current_dir(&self.cwd)
1896 .env("CODESWARM_CWD", &self.cwd)
1897 .stdout(Stdio::piped())
1898 .stderr(Stdio::piped());
1899 if let Some(session_id) = &self.session_id {
1900 command.arg("--conversation").arg(session_id);
1901 }
1902 if self.mode != "default" {
1903 command.arg("--mode").arg(&self.mode);
1904 }
1905 let mut child = command
1906 .spawn()
1907 .map_err(|error| AdapterError::Spawn(error.to_string()))?;
1908 let stdout = match child.stdout.take() {
1909 Some(stdout) => stdout,
1910 None => {
1911 let _ = terminate_child(&mut child).await;
1912 return Err(AdapterError::Transport("agent has no stdout".into()));
1913 }
1914 };
1915 let stderr = match child.stderr.take() {
1916 Some(stderr) => stderr,
1917 None => {
1918 let _ = terminate_child(&mut child).await;
1919 return Err(AdapterError::Transport("agent has no stderr".into()));
1920 }
1921 };
1922 let sender = self.sender.clone();
1923 let slot = self.slot;
1924 let announced_session = Arc::clone(&self.announced_session);
1925 let cancel_requested = Arc::clone(&self.cancel_requested);
1926 tokio::spawn(async move {
1927 let stderr_task = tokio::spawn(async move {
1928 const MAX_STDERR: usize = 32 * 1024;
1929 let mut stderr = BufReader::new(stderr);
1930 let mut bytes = Vec::new();
1931 let mut chunk = [0_u8; 4096];
1932 while let Ok(count) = stderr.read(&mut chunk).await {
1933 if count == 0 {
1934 break;
1935 }
1936 bytes.extend_from_slice(&chunk[..count]);
1937 if bytes.len() > MAX_STDERR {
1938 let keep_from = bytes.len() - MAX_STDERR;
1939 bytes.drain(..keep_from);
1940 }
1941 }
1942 String::from_utf8_lossy(&bytes).trim().to_owned()
1943 });
1944 let mut lines = BufReader::new(stdout).lines();
1945 let mut result: Option<Value> = None;
1946 let mut streamed_response = false;
1947 while let Ok(Some(line)) = lines.next_line().await {
1948 let value = match serde_json::from_str::<Value>(&line) {
1949 Ok(value) => value,
1950 Err(_) => {
1951 continue;
1956 }
1957 };
1958 if value.get("event").and_then(Value::as_str) == Some("init")
1959 && let Some(session_id) = value
1960 .get("conversation_id")
1961 .or_else(|| value.get("conversationId"))
1962 .and_then(Value::as_str)
1963 .filter(|id| !id.is_empty())
1964 && let Ok(mut announced) = announced_session.lock()
1965 {
1966 *announced = Some(session_id.to_owned());
1967 }
1968 if value.get("event").and_then(Value::as_str) == Some("result") {
1969 result = value.get("result").cloned();
1970 }
1971 match parse_agy_value(slot, &value) {
1972 Ok(Some(event)) => {
1973 if matches!(event, AgentEvent::Text { .. }) {
1974 streamed_response = true;
1975 }
1976 if sender.send(Ok(event)).await.is_err() {
1977 break;
1978 }
1979 }
1980 Ok(None) => {}
1981 Err(error) => {
1982 let _ = sender.send(Err(error)).await;
1983 }
1984 }
1985 }
1986 let stderr = stderr_task.await.ok().unwrap_or_default();
1987 let succeeded = cancel_requested.load(Ordering::Acquire)
1988 || result
1989 .as_ref()
1990 .and_then(|result| result.get("status"))
1991 .and_then(Value::as_str)
1992 == Some("SUCCESS");
1993 if succeeded {
1994 if !streamed_response
1999 && let Some(response) = result
2000 .as_ref()
2001 .and_then(|result| result.get("response"))
2002 .and_then(Value::as_str)
2003 .filter(|response| !response.is_empty())
2004 {
2005 let _ = sender
2006 .send(Ok(AgentEvent::Text {
2007 slot,
2008 text: response.to_owned(),
2009 }))
2010 .await;
2011 }
2012 let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
2013 } else {
2014 let detail = result
2015 .as_ref()
2016 .and_then(|result| result.get("error"))
2017 .and_then(Value::as_str)
2018 .filter(|detail| !detail.is_empty())
2019 .map(str::to_owned)
2020 .or_else(|| (!stderr.is_empty()).then_some(stderr))
2021 .unwrap_or_else(|| "native stream ended before a successful result".into());
2022 let _ = sender
2023 .send(Ok(AgentEvent::Failed {
2024 slot,
2025 started: true,
2026 detail,
2027 }))
2028 .await;
2029 }
2030 });
2031 self.child = Some(child);
2032 Ok(())
2033 }
2034
2035 async fn cancel(&mut self) -> AdapterResult<bool> {
2036 self.cancel_requested.store(true, Ordering::Release);
2037 let Some(mut child) = self.child.take() else {
2038 return Ok(false);
2039 };
2040 terminate_child(&mut child).await?;
2044 let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
2045 while let Some(event) = self.receiver.recv().await {
2046 if matches!(
2047 event,
2048 Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
2049 ) {
2050 break;
2051 }
2052 }
2053 })
2054 .await;
2055 Ok(true)
2056 }
2057
2058 async fn answer_permission(
2059 &mut self,
2060 _request_id: String,
2061 _answer: PermissionAnswer,
2062 ) -> AdapterResult<()> {
2063 Err(AdapterError::Unsupported("permission answer"))
2064 }
2065
2066 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
2067 let (mode, mode_policy) = match mode.as_str() {
2068 "full-access" | "codeswarm:mode:full-access" | "auto" | "autopilot" => {
2069 ("default".to_owned(), "agy:full-access".to_owned())
2070 }
2071 "codeswarm:mode:plan" | "readonly" | "plan" => ("plan".to_owned(), "plan".to_owned()),
2072 "codeswarm:mode:accept-edits" | "acceptedits" | "accept-edits" => {
2073 ("accept-edits".to_owned(), "accept-edits".to_owned())
2074 }
2075 "codeswarm:mode:manual" | "manual" | "ask" | "default" => {
2076 ("default".to_owned(), "agy:manual".to_owned())
2077 }
2078 "agy:full-access" => ("default".to_owned(), "agy:full-access".to_owned()),
2079 "agy:manual" => ("default".to_owned(), "agy:manual".to_owned()),
2080 _ => return Err(AdapterError::Unsupported("requested Agy mode")),
2081 };
2082 self.mode = mode;
2083 self.mode_policy = mode_policy.clone();
2084 self.emit(Ok(AgentEvent::ModesReplaced {
2085 slot: self.slot,
2086 modes: Self::modes(),
2087 current_mode: Some(mode_policy),
2088 }))
2089 .await;
2090 Ok(())
2091 }
2092
2093 async fn reload(&mut self) -> AdapterResult<()> {
2094 self.stop().await?;
2095 self.start().await
2096 }
2097
2098 async fn stop(&mut self) -> AdapterResult<()> {
2099 let _ = self.cancel().await?;
2100 Ok(())
2101 }
2102
2103 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
2104 let event = self.receiver.recv().await;
2105 if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2106 && self.session_id.is_none()
2107 && let Ok(session) = self.announced_session.lock()
2108 {
2109 self.session_id = session.clone();
2110 }
2111 if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2112 && let Some(mut child) = self.child.take()
2113 {
2114 let _ = child.wait().await;
2115 }
2116 event
2117 }
2118}
2119
2120#[cfg(test)]
2121#[cfg_attr(not(test), allow(dead_code))]
2122fn parse_agy_line(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
2123 let value: Value =
2124 serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
2125 parse_agy_value(slot, &value)
2126}
2127
2128fn parse_agy_value(slot: RosterSlot, value: &Value) -> AdapterResult<Option<AgentEvent>> {
2129 let event = value.get("event").and_then(Value::as_str);
2130 if let Some(terminal) = parse_terminal_event(value, event) {
2131 return Ok(Some(AgentEvent::Terminal {
2132 slot,
2133 event: terminal,
2134 }));
2135 }
2136 match event {
2137 Some("step_update") => {
2138 let Some(update) = value.get("step_update") else {
2139 return Ok(None);
2140 };
2141 let is_response = update
2142 .get("step_type")
2143 .and_then(Value::as_str)
2144 .is_some_and(|kind| kind == "agent_response");
2145 let text = update.get("text_delta").and_then(Value::as_str);
2146 let response = is_response
2147 .then(|| text.map(str::to_owned))
2148 .flatten()
2149 .filter(|text| !text.is_empty())
2150 .map(|text| AgentEvent::Text { slot, text });
2151 Ok(response.or_else(|| parse_agy_tool(slot, value)))
2152 }
2153 _ => Ok(None),
2154 }
2155}
2156
2157fn parse_agy_tool(slot: RosterSlot, value: &Value) -> Option<AgentEvent> {
2158 let update = value.get("step_update")?;
2159 if update.get("step_type")?.as_str()? != "tool" {
2160 return None;
2161 }
2162 let step_index = update.get("step_index")?.as_i64()?;
2163 let title = update
2164 .get("tool_name")
2165 .and_then(Value::as_str)
2166 .unwrap_or("Tool call")
2167 .replace('_', " ");
2168 let status = match update.get("state").and_then(Value::as_str) {
2169 Some("DONE") => ToolStatus::Completed,
2170 Some("FAILED") => ToolStatus::Failed,
2171 Some("ACTIVE") => ToolStatus::Running,
2172 _ => ToolStatus::Pending,
2173 };
2174 let detail = update
2175 .get("tool_info")
2176 .and_then(|info| info.get("output"))
2177 .and_then(Value::as_str)
2178 .map(str::to_owned);
2179 Some(AgentEvent::Tool {
2180 slot,
2181 update: ToolUpdate {
2182 id: format!("agy-tool-{step_index}"),
2183 title,
2184 status,
2185 detail,
2186 },
2187 })
2188}
2189
2190#[derive(Debug)]
2193pub struct AcpAdapter {
2194 slot: RosterSlot,
2195 program: String,
2196 args: Vec<String>,
2197 cwd: PathBuf,
2198 child: Option<Child>,
2199 reader: Option<BufReader<ChildStdout>>,
2200 capabilities: AgentCapabilities,
2201 modes: Vec<Mode>,
2202 models: Vec<Mode>,
2203 model_config_id: Option<String>,
2204 session_id: Option<String>,
2205 next_request_id: u64,
2206 prompt_request_id: Option<u64>,
2207 queued_events: VecDeque<AdapterResult<AgentEvent>>,
2208 stderr_task: Option<tokio::task::JoinHandle<String>>,
2209 terminals: BTreeMap<String, TerminalProcess>,
2210 next_terminal_id: u64,
2211}
2212
2213impl AcpAdapter {
2214 pub fn new(
2215 slot: RosterSlot,
2216 cwd: PathBuf,
2217 program: impl Into<String>,
2218 args: Vec<String>,
2219 ) -> Self {
2220 Self {
2221 slot,
2222 program: program.into(),
2223 args,
2224 cwd,
2225 child: None,
2226 reader: None,
2227 capabilities: AgentCapabilities::default(),
2228 modes: Vec::new(),
2229 models: Vec::new(),
2230 model_config_id: None,
2231 session_id: None,
2232 next_request_id: 1,
2233 prompt_request_id: None,
2234 queued_events: VecDeque::new(),
2235 stderr_task: None,
2236 terminals: BTreeMap::new(),
2237 next_terminal_id: 1,
2238 }
2239 }
2240
2241 pub fn with_session_id(
2242 slot: RosterSlot,
2243 cwd: PathBuf,
2244 program: impl Into<String>,
2245 args: Vec<String>,
2246 session_id: impl Into<String>,
2247 ) -> Self {
2248 let mut adapter = Self::new(slot, cwd, program, args);
2249 adapter.session_id = Some(session_id.into());
2250 adapter
2251 }
2252
2253 async fn request(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2254 let request_id = self.next_request_id;
2255 self.next_request_id += 1;
2256 self.write_json(serde_json::json!({
2257 "jsonrpc": "2.0",
2258 "id": request_id,
2259 "method": method,
2260 "params": params,
2261 }))
2262 .await?;
2263 loop {
2264 let line = self.read_line().await?;
2265 let value: Value = match serde_json::from_str(&line) {
2266 Ok(value) => value,
2267 Err(_) => {
2268 continue;
2272 }
2273 };
2274 if self.reject_empty_permission_request(&value).await? {
2275 continue;
2276 }
2277 if self.handle_client_request(&value).await? {
2278 continue;
2279 }
2280 if value.get("id").and_then(Value::as_u64) == Some(request_id) {
2281 if let Some(error) = value.get("error") {
2282 return Err(AdapterError::Protocol(error.to_string()));
2283 }
2284 return value
2285 .get("result")
2286 .cloned()
2287 .ok_or_else(|| AdapterError::Protocol("response has no result".into()));
2288 }
2289 if let Some(event) = parse_acp_notification(self.slot, &line)? {
2290 self.queued_events.push_back(Ok(event));
2291 }
2292 }
2293 }
2294
2295 async fn write_json(&mut self, value: Value) -> AdapterResult<()> {
2296 let child = self
2297 .child
2298 .as_mut()
2299 .ok_or_else(|| AdapterError::Transport("ACP agent is not running".into()))?;
2300 let stdin = child
2301 .stdin
2302 .as_mut()
2303 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdin".into()))?;
2304 stdin
2305 .write_all(value.to_string().as_bytes())
2306 .await
2307 .map_err(|error| AdapterError::Transport(error.to_string()))?;
2308 stdin
2309 .write_all(b"\n")
2310 .await
2311 .map_err(|error| AdapterError::Transport(error.to_string()))
2312 }
2313
2314 async fn reject_empty_permission_request(&mut self, value: &Value) -> AdapterResult<bool> {
2319 if value.get("method").and_then(Value::as_str) != Some("session/request_permission")
2320 || value.get("id").is_none()
2321 {
2322 return Ok(false);
2323 }
2324 let valid = value
2325 .get("params")
2326 .and_then(|params| params.get("options"))
2327 .and_then(Value::as_array)
2328 .is_some_and(|options| !options.is_empty());
2329 if valid {
2330 return Ok(false);
2331 }
2332 self.write_json(serde_json::json!({
2333 "jsonrpc": "2.0",
2334 "id": value.get("id").cloned().unwrap_or(Value::Null),
2335 "error": {
2336 "code": -32602,
2337 "message": "Permission request requires at least one option",
2338 },
2339 }))
2340 .await?;
2341 Ok(true)
2342 }
2343
2344 fn workspace_path(&self, path: &str) -> Result<PathBuf, String> {
2345 let root = self
2346 .cwd
2347 .canonicalize()
2348 .map_err(|error| format!("unable to resolve workspace: {error}"))?;
2349 let requested = Path::new(path);
2350 let candidate = if requested.is_absolute() {
2351 requested.to_path_buf()
2352 } else {
2353 root.join(requested)
2354 };
2355 let resolved = if !candidate.exists() {
2356 let parent = candidate
2357 .parent()
2358 .ok_or_else(|| "file path has no parent".to_owned())?
2359 .canonicalize()
2360 .map_err(|error| format!("unable to resolve parent directory: {error}"))?;
2361 parent.join(
2362 candidate
2363 .file_name()
2364 .ok_or_else(|| "file path has no filename".to_owned())?,
2365 )
2366 } else {
2367 candidate
2368 .canonicalize()
2369 .map_err(|error| format!("unable to resolve file path: {error}"))?
2370 };
2371 if !resolved.starts_with(&root) {
2372 return Err("file path is outside the project".into());
2373 }
2374 Ok(resolved)
2375 }
2376
2377 fn read_workspace_text(
2378 &self,
2379 path: &str,
2380 line: Option<i64>,
2381 limit: Option<i64>,
2382 ) -> Result<String, String> {
2383 if line.is_some_and(|line| line < 1) {
2384 return Err("line must be positive".into());
2385 }
2386 if limit.is_some_and(|limit| limit < 0) {
2387 return Err("limit must not be negative".into());
2388 }
2389 let path = self.workspace_path(path)?;
2390 let mut bytes = Vec::new();
2391 let mut source = match std::fs::File::open(path) {
2392 Ok(source) => source,
2393 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(String::new()),
2394 Err(error) => return Err(error.to_string()),
2395 };
2396 source
2397 .by_ref()
2398 .take((MAX_FILE_READ_BYTES as u64).saturating_add(1))
2399 .read_to_end(&mut bytes)
2400 .map_err(|error| error.to_string())?;
2401 bytes.truncate(MAX_FILE_READ_BYTES);
2402 let text = String::from_utf8_lossy(&bytes);
2403 if line.is_none() && limit.is_none() {
2404 return Ok(text.into_owned());
2405 }
2406 let start = line.map_or(0, |line| line as usize - 1);
2407 let limit = limit.unwrap_or(i64::MAX) as usize;
2408 let selected = text
2409 .split_inclusive('\n')
2410 .skip(start)
2411 .take(limit)
2412 .collect::<String>();
2413 if line.is_some() {
2414 Ok(selected.trim_end_matches('\n').to_owned())
2415 } else {
2416 Ok(selected)
2417 }
2418 }
2419
2420 async fn terminal_create(&mut self, params: &Value) -> Result<Value, String> {
2421 let command = params
2422 .get("command")
2423 .and_then(Value::as_str)
2424 .filter(|command| !command.trim().is_empty())
2425 .ok_or_else(|| "terminal command is required".to_owned())?;
2426 let cwd = params.get("cwd").and_then(Value::as_str).unwrap_or(".");
2427 let cwd = self.workspace_path(cwd)?;
2428 if !cwd.is_dir() {
2429 return Err("terminal cwd is not a directory".into());
2430 }
2431 let mut process = Command::new(command);
2432 isolate_process_group(&mut process);
2433 if let Some(args) = params.get("args").and_then(Value::as_array) {
2434 process.args(args.iter().filter_map(Value::as_str));
2435 }
2436 process
2437 .current_dir(&cwd)
2438 .stdin(Stdio::null())
2439 .stdout(Stdio::piped())
2440 .stderr(Stdio::piped());
2441 if let Some(env) = params.get("env") {
2442 if let Some(entries) = env.as_array() {
2443 for entry in entries {
2444 if let (Some(name), Some(value)) = (
2445 entry.get("name").and_then(Value::as_str),
2446 entry.get("value").and_then(Value::as_str),
2447 ) {
2448 process.env(name, value);
2449 }
2450 }
2451 } else if let Some(entries) = env.as_object() {
2452 for (name, value) in entries {
2453 if let Some(value) = value.as_str() {
2454 process.env(name, value);
2455 }
2456 }
2457 }
2458 }
2459 let mut child = process.spawn().map_err(|error| error.to_string())?;
2460 let stdout = child.stdout.take();
2461 let stderr = child.stderr.take();
2462 let (Some(stdout), Some(stderr)) = (stdout, stderr) else {
2463 let _ = terminate_child(&mut child).await;
2464 return Err("terminal has no output pipes".into());
2465 };
2466 let output = Arc::new(Mutex::new(Vec::new()));
2467 let truncated = Arc::new(AtomicBool::new(false));
2468 let output_readers = Arc::new(AtomicUsize::new(2));
2469 let output_limit = params
2470 .get("outputByteLimit")
2471 .and_then(Value::as_u64)
2472 .map_or(MAX_TERMINAL_OUTPUT_BYTES, |limit| {
2473 usize::try_from(limit)
2474 .unwrap_or(MAX_TERMINAL_OUTPUT_BYTES)
2475 .min(MAX_TERMINAL_OUTPUT_BYTES)
2476 });
2477 tokio::spawn(drain_terminal_output(
2478 stdout,
2479 Arc::clone(&output),
2480 Arc::clone(&truncated),
2481 Arc::clone(&output_readers),
2482 output_limit,
2483 ));
2484 tokio::spawn(drain_terminal_output(
2485 stderr,
2486 Arc::clone(&output),
2487 Arc::clone(&truncated),
2488 Arc::clone(&output_readers),
2489 output_limit,
2490 ));
2491 let id = format!("terminal-{}", self.next_terminal_id);
2492 self.next_terminal_id = self.next_terminal_id.saturating_add(1);
2493 let state = TerminalProcess {
2494 child: Arc::new(AsyncMutex::new(Some(child))),
2495 output,
2496 truncated,
2497 output_readers,
2498 };
2499 self.terminals.insert(id.clone(), state);
2500 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2501 slot: self.slot,
2502 event: TerminalEvent::Created {
2503 id: id.clone(),
2504 command: std::iter::once(command)
2505 .chain(
2506 params
2507 .get("args")
2508 .and_then(Value::as_array)
2509 .into_iter()
2510 .flatten()
2511 .filter_map(Value::as_str),
2512 )
2513 .collect::<Vec<_>>()
2514 .join(" "),
2515 },
2516 }));
2517 Ok(serde_json::json!({"terminalId": id}))
2518 }
2519
2520 async fn terminal_output(&mut self, id: &str) -> Result<Value, String> {
2521 let terminal = self
2522 .terminals
2523 .get(id)
2524 .ok_or_else(|| "terminal not found".to_owned())?
2525 .clone();
2526 let output = terminal
2527 .output
2528 .lock()
2529 .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
2530 .unwrap_or_default();
2531 let exit_code = terminal.exit_code().await;
2532 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2533 slot: self.slot,
2534 event: TerminalEvent::Output {
2535 id: id.to_owned(),
2536 text: output.clone(),
2537 },
2538 }));
2539 let mut response = serde_json::json!({
2540 "output": output,
2541 "truncated": terminal.truncated.load(Ordering::Acquire),
2542 });
2543 if let Some(code) = exit_code {
2544 response["exitStatus"] = serde_json::json!({"exitCode": code});
2545 }
2546 Ok(response)
2547 }
2548
2549 async fn terminal_wait(&mut self, id: &str) -> Result<Value, String> {
2550 let terminal = self
2551 .terminals
2552 .get(id)
2553 .ok_or_else(|| "terminal not found".to_owned())?
2554 .clone();
2555 let exit_code = terminal.wait().await;
2556 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2557 slot: self.slot,
2558 event: TerminalEvent::Exited {
2559 id: id.to_owned(),
2560 code: exit_code.unwrap_or(-1),
2561 },
2562 }));
2563 Ok(serde_json::json!({"exitCode": exit_code, "signal": Value::Null}))
2564 }
2565
2566 async fn handle_client_request(&mut self, value: &Value) -> AdapterResult<bool> {
2570 let Some(method) = value.get("method").and_then(Value::as_str) else {
2571 return Ok(false);
2572 };
2573 let Some(id) = value.get("id").cloned() else {
2574 return Ok(false);
2575 };
2576 if method == "session/request_permission" {
2579 return Ok(false);
2580 }
2581 let params = value.get("params").cloned().unwrap_or(Value::Null);
2582 let response = match method {
2583 "fs/read_text_file" => {
2584 let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2585 let line = params.get("line").and_then(Value::as_i64);
2586 let limit = params.get("limit").and_then(Value::as_i64);
2587 match self.read_workspace_text(path, line, limit) {
2588 Ok(content) => serde_json::json!({
2589 "jsonrpc": "2.0",
2590 "id": id,
2591 "result": {"content": content},
2592 }),
2593 Err(message) => serde_json::json!({
2594 "jsonrpc": "2.0",
2595 "id": id,
2596 "error": {"code": -32602, "message": message},
2597 }),
2598 }
2599 }
2600 "fs/write_text_file" => {
2601 let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2602 let content = params.get("content").and_then(Value::as_str).unwrap_or("");
2603 let result = self.workspace_path(path).and_then(|path| {
2604 std::fs::write(path, content).map_err(|error| error.to_string())
2605 });
2606 match result {
2607 Ok(()) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}}),
2608 Err(message) => serde_json::json!({
2609 "jsonrpc": "2.0",
2610 "id": id,
2611 "error": {"code": -32602, "message": message},
2612 }),
2613 }
2614 }
2615 "terminal/create" => match self.terminal_create(¶ms).await {
2616 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2617 Err(message) => serde_json::json!({
2618 "jsonrpc": "2.0",
2619 "id": id,
2620 "error": {"code": -32602, "message": message},
2621 }),
2622 },
2623 "terminal/output" => {
2624 let terminal_id = params
2625 .get("terminalId")
2626 .and_then(Value::as_str)
2627 .unwrap_or("");
2628 match self.terminal_output(terminal_id).await {
2629 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2630 Err(message) => serde_json::json!({
2631 "jsonrpc": "2.0",
2632 "id": id,
2633 "error": {"code": -32602, "message": message},
2634 }),
2635 }
2636 }
2637 "terminal/wait_for_exit" => {
2638 let terminal_id = params
2639 .get("terminalId")
2640 .and_then(Value::as_str)
2641 .unwrap_or("");
2642 match self.terminal_wait(terminal_id).await {
2643 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2644 Err(message) => serde_json::json!({
2645 "jsonrpc": "2.0",
2646 "id": id,
2647 "error": {"code": -32602, "message": message},
2648 }),
2649 }
2650 }
2651 "terminal/kill" => {
2652 let terminal_id = params
2653 .get("terminalId")
2654 .and_then(Value::as_str)
2655 .unwrap_or("");
2656 if let Some(terminal) = self.terminals.get(terminal_id) {
2657 terminal.kill().await;
2658 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2659 } else {
2660 serde_json::json!({
2661 "jsonrpc": "2.0",
2662 "id": id,
2663 "error": {"code": -32602, "message": "terminal not found"},
2664 })
2665 }
2666 }
2667 "terminal/release" => {
2668 let terminal_id = params
2669 .get("terminalId")
2670 .and_then(Value::as_str)
2671 .unwrap_or("");
2672 if let Some(terminal) = self.terminals.remove(terminal_id) {
2673 terminal.stop().await;
2674 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2675 slot: self.slot,
2676 event: TerminalEvent::Released {
2677 id: terminal_id.to_owned(),
2678 },
2679 }));
2680 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2681 } else {
2682 serde_json::json!({
2683 "jsonrpc": "2.0",
2684 "id": id,
2685 "error": {"code": -32602, "message": "terminal not found"},
2686 })
2687 }
2688 }
2689 _ => serde_json::json!({
2690 "jsonrpc": "2.0",
2691 "id": id,
2692 "error": {"code": -32601, "message": format!("unsupported client method: {method}")},
2693 }),
2694 };
2695 self.write_json(response).await?;
2696 Ok(true)
2697 }
2698
2699 async fn read_line(&mut self) -> AdapterResult<String> {
2700 let reader = self
2701 .reader
2702 .as_mut()
2703 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdout".into()))?;
2704 read_bounded_line(reader).await
2705 }
2706
2707 async fn start(&mut self) -> AdapterResult<()> {
2708 if self.child.is_some() {
2710 self.stop().await?;
2711 }
2712 let mut command = Command::new(&self.program);
2713 isolate_process_group(&mut command);
2714 command
2715 .args(&self.args)
2716 .current_dir(&self.cwd)
2717 .stdin(Stdio::piped())
2718 .stdout(Stdio::piped())
2719 .stderr(Stdio::piped())
2720 .env("CODESWARM_CWD", &self.cwd);
2721 if self.program.to_ascii_lowercase().contains("gemini")
2722 || self
2723 .args
2724 .iter()
2725 .any(|arg| arg.to_ascii_lowercase().contains("gemini"))
2726 {
2727 command.env("GEMINI_TELEMETRY_ENABLED", "false");
2728 }
2729 let mut child = command
2730 .spawn()
2731 .map_err(|error| AdapterError::Spawn(error.to_string()))?;
2732 let stdout = match child.stdout.take() {
2733 Some(stdout) => stdout,
2734 None => {
2735 let _ = terminate_child(&mut child).await;
2736 return Err(AdapterError::Transport("ACP agent has no stdout".into()));
2737 }
2738 };
2739 let stderr = match child.stderr.take() {
2740 Some(stderr) => stderr,
2741 None => {
2742 let _ = terminate_child(&mut child).await;
2743 return Err(AdapterError::Transport("ACP agent has no stderr".into()));
2744 }
2745 };
2746 self.child = Some(child);
2747 self.reader = Some(BufReader::new(stdout));
2748 self.stderr_task = Some(tokio::spawn(drain_bounded(stderr, 32 * 1024)));
2749
2750 let initialize = match self
2751 .request(
2752 "initialize",
2753 serde_json::json!({
2754 "protocolVersion": 1,
2755 "clientCapabilities": {
2756 "fs": {"readTextFile": true, "writeTextFile": true},
2757 "terminal": true,
2758 },
2759 "clientInfo": {
2760 "name": "CodeSwarm",
2761 "title": "CodeSwarm",
2762 "version": env!("CARGO_PKG_VERSION"),
2763 },
2764 }),
2765 )
2766 .await
2767 {
2768 Ok(value) => value,
2769 Err(error) => {
2770 let _ = self.stop().await;
2771 return Err(error);
2772 }
2773 };
2774 let agent_capabilities = initialize
2775 .get("agentCapabilities")
2776 .cloned()
2777 .unwrap_or(Value::Null);
2778 self.capabilities = AgentCapabilities {
2779 supports_cancel: true,
2780 supports_modes: true,
2781 supports_permissions: true,
2782 supports_terminals: true,
2783 supports_session_load: agent_capabilities
2784 .get("loadSession")
2785 .and_then(Value::as_bool)
2786 .unwrap_or(false),
2787 supports_models: false,
2788 };
2789 let session = if let Some(session_id) = self.session_id.clone() {
2790 if !self.capabilities.supports_session_load {
2791 let _ = self.stop().await;
2792 return Err(AdapterError::Unsupported("session/load"));
2793 }
2794 match self
2795 .request(
2796 "session/load",
2797 serde_json::json!({
2798 "cwd": self.cwd,
2799 "mcpServers": [],
2800 "sessionId": session_id,
2801 }),
2802 )
2803 .await
2804 {
2805 Ok(value) => value,
2806 Err(error) => {
2807 let _ = self.stop().await;
2808 return Err(error);
2809 }
2810 }
2811 } else {
2812 let session = match self
2813 .request(
2814 "session/new",
2815 serde_json::json!({"cwd": self.cwd, "mcpServers": []}),
2816 )
2817 .await
2818 {
2819 Ok(value) => value,
2820 Err(error) => {
2821 let _ = self.stop().await;
2822 return Err(error);
2823 }
2824 };
2825 self.session_id = session
2826 .get("sessionId")
2827 .and_then(Value::as_str)
2828 .map(str::to_owned);
2829 if self.session_id.is_none() {
2830 let _ = self.stop().await;
2831 return Err(AdapterError::Protocol(
2832 "session/new returned no sessionId".into(),
2833 ));
2834 }
2835 session
2836 };
2837 self.capabilities.supports_modes = false;
2838 if let Some(modes) = session.get("modes") {
2839 let available = modes
2840 .get("availableModes")
2841 .and_then(Value::as_array)
2842 .map(|modes| {
2843 modes
2844 .iter()
2845 .filter_map(|mode| {
2846 Some(Mode {
2847 id: mode.get("id")?.as_str()?.to_owned(),
2848 label: mode.get("name")?.as_str()?.to_owned(),
2849 })
2850 })
2851 .collect::<Vec<_>>()
2852 })
2853 .unwrap_or_default();
2854 self.modes = available.clone();
2855 self.capabilities.supports_modes = !available.is_empty();
2856 self.queued_events.push_back(Ok(AgentEvent::ModesReplaced {
2857 slot: self.slot,
2858 modes: available,
2859 current_mode: modes
2860 .get("currentModeId")
2861 .and_then(Value::as_str)
2862 .map(str::to_owned),
2863 }));
2864 }
2865 self.models.clear();
2866 self.model_config_id = None;
2867 let current_model =
2868 parse_model_config(&session).and_then(|(config_id, models, current)| {
2869 self.model_config_id = Some(config_id);
2870 self.models = models;
2871 current
2872 });
2873 self.capabilities.supports_models =
2874 self.model_config_id.is_some() && !self.models.is_empty();
2875 if let Some(config_id) = self.model_config_id.clone()
2876 && !self.models.is_empty()
2877 {
2878 self.queued_events.push_back(Ok(AgentEvent::ModelsReplaced {
2879 slot: self.slot,
2880 config_id,
2881 models: self.models.clone(),
2882 current_model,
2883 }));
2884 }
2885 self.queued_events.push_back(Ok(AgentEvent::Ready {
2886 slot: self.slot,
2887 capabilities: self.capabilities(),
2888 }));
2889 Ok(())
2890 }
2891}
2892
2893fn prompt_resource_paths(prompt: &str) -> Vec<String> {
2894 let characters = prompt.chars().collect::<Vec<_>>();
2895 let mut paths = Vec::new();
2896 let mut index = 0;
2897 while index < characters.len() {
2898 if characters[index] != '@' {
2899 index += 1;
2900 continue;
2901 }
2902 index += 1;
2903 let quoted = characters.get(index) == Some(&'"');
2904 if quoted {
2905 index += 1;
2906 }
2907 let start = index;
2908 while index < characters.len()
2909 && if quoted {
2910 characters[index] != '"'
2911 } else {
2912 !characters[index].is_whitespace()
2913 }
2914 {
2915 index += 1;
2916 }
2917 if index > start {
2918 paths.push(characters[start..index].iter().collect());
2919 }
2920 if quoted && index < characters.len() {
2921 index += 1;
2922 }
2923 }
2924 paths
2925}
2926
2927fn prompt_content_blocks(cwd: &Path, prompt: &str) -> Vec<Value> {
2928 let mut blocks = vec![serde_json::json!({"type": "text", "text": prompt})];
2929 for path in prompt_resource_paths(prompt) {
2930 if path.ends_with('/') {
2931 continue;
2932 }
2933 let Ok(resource) = resources::load(cwd, &path) else {
2934 continue;
2935 };
2936 let uri = format!("file://{}", resource.path.display());
2937 let resource_value = if let Some(text) = resource.text {
2938 serde_json::json!({
2939 "uri": uri,
2940 "text": text,
2941 "mimeType": resource.mime_type,
2942 })
2943 } else if let Some(data) = resource.data {
2944 serde_json::json!({
2945 "uri": uri,
2946 "blob": BASE64.encode(data),
2947 "mimeType": resource.mime_type,
2948 })
2949 } else {
2950 continue;
2951 };
2952 blocks.push(serde_json::json!({
2953 "type": "resource",
2954 "resource": resource_value,
2955 }));
2956 }
2957 blocks
2958}
2959
2960#[async_trait]
2961impl AgentAdapter for AcpAdapter {
2962 fn slot(&self) -> RosterSlot {
2963 self.slot
2964 }
2965
2966 fn session_id(&self) -> Option<String> {
2967 self.session_id.clone()
2968 }
2969
2970 fn protocol(&self) -> &'static str {
2971 "acp"
2972 }
2973
2974 fn capabilities(&self) -> AgentCapabilities {
2975 self.capabilities.clone()
2976 }
2977
2978 async fn start(&mut self) -> AdapterResult<()> {
2979 AcpAdapter::start(self).await
2982 }
2983
2984 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
2985 let session_id = self
2986 .session_id
2987 .as_ref()
2988 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
2989 let request_id = self.next_request_id;
2990 self.next_request_id += 1;
2991 let prompt_blocks = prompt_content_blocks(&self.cwd, &prompt);
2992 self.write_json(serde_json::json!({
2993 "jsonrpc": "2.0",
2994 "id": request_id,
2995 "method": "session/prompt",
2996 "params": {
2997 "sessionId": session_id,
2998 "prompt": prompt_blocks,
2999 },
3000 }))
3001 .await?;
3002 self.prompt_request_id = Some(request_id);
3003 Ok(())
3004 }
3005
3006 async fn cancel(&mut self) -> AdapterResult<bool> {
3007 let Some(session_id) = &self.session_id else {
3008 return Ok(false);
3009 };
3010 self.write_json(serde_json::json!({
3011 "jsonrpc": "2.0",
3012 "method": "session/cancel",
3013 "params": {"sessionId": session_id, "_meta": {}},
3014 }))
3015 .await?;
3016 let settled = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
3017 loop {
3018 match <Self as AgentAdapter>::next_event(self).await {
3019 Some(Ok(AgentEvent::TurnComplete { .. })) | None => break,
3020 Some(Ok(_)) => {}
3021 Some(Err(_)) => break,
3022 }
3023 }
3024 })
3025 .await
3026 .is_ok();
3027 if !settled {
3028 self.reload().await?;
3032 }
3033 Ok(true)
3034 }
3035
3036 async fn answer_permission(
3037 &mut self,
3038 request_id: String,
3039 answer: PermissionAnswer,
3040 ) -> AdapterResult<()> {
3041 let id = request_id
3042 .parse::<u64>()
3043 .map(Value::from)
3044 .unwrap_or_else(|_| Value::String(request_id));
3045 let outcome = match answer {
3046 PermissionAnswer::Selected { option_id } => {
3047 serde_json::json!({"outcome": "selected", "optionId": option_id})
3048 }
3049 PermissionAnswer::Cancelled => serde_json::json!({"outcome": "cancelled"}),
3050 };
3051 self.write_json(serde_json::json!({
3052 "jsonrpc": "2.0",
3053 "id": id,
3054 "result": {"outcome": outcome},
3058 }))
3059 .await
3060 }
3061
3062 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
3063 let session_id = self
3064 .session_id
3065 .as_ref()
3066 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3067 let policy = match mode.as_str() {
3068 "plan" => "codeswarm:mode:plan",
3069 "default" | "manual" => "codeswarm:mode:manual",
3070 "accept-edits" => "codeswarm:mode:accept-edits",
3071 "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
3072 other => other,
3073 };
3074 let native_mode = crate::policy::resolve(policy, &self.modes)
3075 .map(|mode| mode.id)
3076 .unwrap_or(mode);
3077 let _ = self
3078 .request(
3079 "session/set_mode",
3080 serde_json::json!({"sessionId": session_id, "modeId": native_mode.clone()}),
3081 )
3082 .await?;
3083 self.queued_events.push_back(Ok(AgentEvent::ModeUpdated {
3084 slot: self.slot,
3085 current_mode: native_mode,
3086 }));
3087 Ok(())
3088 }
3089
3090 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
3091 let session_id = self
3092 .session_id
3093 .clone()
3094 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3095 let config_id = self
3096 .model_config_id
3097 .clone()
3098 .ok_or(AdapterError::Unsupported("set_model"))?;
3099 if !self.models.iter().any(|candidate| candidate.id == model) {
3100 return Err(AdapterError::Protocol(
3101 "model is not advertised by the agent".into(),
3102 ));
3103 }
3104 let _ = self
3105 .request(
3106 "session/set_config_option",
3107 serde_json::json!({
3108 "sessionId": session_id,
3109 "configId": config_id,
3110 "value": model,
3111 }),
3112 )
3113 .await?;
3114 Ok(())
3115 }
3116
3117 async fn reload(&mut self) -> AdapterResult<()> {
3118 let session_id = self
3123 .capabilities
3124 .supports_session_load
3125 .then(|| self.session_id.clone())
3126 .flatten();
3127 self.stop().await?;
3128 self.session_id = session_id.clone();
3129 let result = self.start().await;
3130 if result.is_err() {
3131 self.session_id = session_id;
3135 }
3136 result
3137 }
3138
3139 async fn stop(&mut self) -> AdapterResult<()> {
3140 let terminals = std::mem::take(&mut self.terminals);
3141 for terminal in terminals.values() {
3142 terminal.stop().await;
3143 }
3144 if let Some(mut child) = self.child.take() {
3145 terminate_child(&mut child).await?;
3146 }
3147 self.reader = None;
3148 self.session_id = None;
3149 self.prompt_request_id = None;
3150 if let Some(task) = self.stderr_task.take() {
3151 task.abort();
3152 let _ = task.await;
3153 }
3154 Ok(())
3155 }
3156
3157 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
3158 if let Some(event) = self.queued_events.pop_front() {
3159 return Some(event);
3160 }
3161 loop {
3162 let line = match self.read_line().await {
3163 Ok(line) => line,
3164 Err(error) => return Some(Err(error)),
3165 };
3166 let value: Value = match serde_json::from_str(&line) {
3167 Ok(value) => value,
3168 Err(_) => continue,
3169 };
3170 match self.reject_empty_permission_request(&value).await {
3171 Ok(true) => continue,
3172 Ok(false) => {}
3173 Err(error) => return Some(Err(error)),
3174 }
3175 match self.handle_client_request(&value).await {
3176 Ok(true) => continue,
3177 Ok(false) => {}
3178 Err(error) => return Some(Err(error)),
3179 }
3180 match parse_acp_notification(self.slot, &line) {
3181 Ok(Some(event)) => {
3182 if let AgentEvent::ModelsReplaced {
3183 config_id, models, ..
3184 } = &event
3185 {
3186 self.model_config_id = Some(config_id.clone());
3187 self.models = models.clone();
3188 self.capabilities.supports_models = !models.is_empty();
3189 }
3190 return Some(Ok(event));
3191 }
3192 Ok(None) => {}
3193 Err(error) => return Some(Err(error)),
3194 }
3195 if value.get("id").is_some_and(|id| {
3196 self.prompt_request_id
3197 .is_some_and(|expected| rpc_id_to_string(id) == expected.to_string())
3198 }) {
3199 if let Some(error) = value.get("error") {
3200 self.prompt_request_id = None;
3201 return Some(Err(AdapterError::Protocol(error.to_string())));
3202 }
3203 self.prompt_request_id = None;
3204 return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3205 }
3206 }
3207 }
3208}
3209
3210fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3211 let value: Value =
3212 serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3213 let method = value.get("method").and_then(Value::as_str);
3214 if method == Some("session/request_permission") {
3215 let params = value.get("params").cloned().unwrap_or(Value::Null);
3216 let request_id = value
3217 .get("id")
3218 .map(rpc_id_to_string)
3219 .unwrap_or_else(|| "permission".into());
3220 return Ok(parse_permission_event(
3221 slot,
3222 ¶ms,
3223 &request_id,
3224 params.get("options"),
3225 ));
3226 }
3227 if method != Some("session/update") {
3228 return Ok(None);
3229 }
3230 let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3231 return Ok(None);
3232 };
3233 let kind = update.get("sessionUpdate").and_then(Value::as_str);
3234 if kind == Some("config_option_update")
3235 && let Some((config_id, models, current_model)) = parse_model_config(update)
3236 {
3237 return Ok(Some(AgentEvent::ModelsReplaced {
3238 slot,
3239 config_id,
3240 models,
3241 current_model,
3242 }));
3243 }
3244 if kind == Some("request_permission") {
3245 let request_id = update
3246 .get("toolCall")
3247 .and_then(|tool| tool.get("toolCallId"))
3248 .and_then(Value::as_str)
3249 .unwrap_or("permission");
3250 return Ok(parse_permission_event(
3251 slot,
3252 update,
3253 request_id,
3254 update.get("options"),
3255 ));
3256 }
3257 if kind == Some("available_commands_update") {
3258 let commands = update
3259 .get("availableCommands")
3260 .and_then(Value::as_array)
3261 .map(|commands| {
3262 commands
3263 .iter()
3264 .filter_map(|command| {
3265 let name = command.get("name").and_then(Value::as_str)?.trim();
3266 (!name.is_empty()).then(|| AgentCommand {
3267 name: name.to_owned(),
3268 })
3269 })
3270 .collect::<Vec<_>>()
3271 })
3272 .unwrap_or_default();
3273 return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3274 }
3275 if kind == Some("current_mode_update") {
3276 if let Some(mode) = update
3277 .get("currentModeId")
3278 .and_then(Value::as_str)
3279 .filter(|mode| !mode.trim().is_empty())
3280 {
3281 return Ok(Some(AgentEvent::ModeUpdated {
3282 slot,
3283 current_mode: mode.to_owned(),
3284 }));
3285 }
3286 return Ok(None);
3287 }
3288 if kind == Some("usage_update") {
3289 let Some(used) = update.get("used").and_then(Value::as_u64) else {
3290 return Ok(None);
3291 };
3292 let Some(size) = update.get("size").and_then(Value::as_u64) else {
3293 return Ok(None);
3294 };
3295 return Ok(Some(AgentEvent::UsageUpdated {
3296 slot,
3297 usage: UsageUpdate { used, size },
3298 }));
3299 }
3300 if let Some(terminal) = parse_terminal_event(update, kind) {
3301 return Ok(Some(AgentEvent::Terminal {
3302 slot,
3303 event: terminal,
3304 }));
3305 }
3306 let text = update
3307 .get("content")
3308 .and_then(|content| content.get("text"))
3309 .and_then(Value::as_str)
3310 .map(str::to_owned);
3311 if kind == Some("user_message_chunk") {
3312 return Ok(text
3313 .filter(|text| !text.is_empty())
3314 .map(|text| AgentEvent::UserText { slot, text }));
3315 }
3316 if kind == Some("agent_message_chunk")
3317 && let Some(mode) = text
3318 .as_deref()
3319 .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3320 .map(str::trim)
3321 .filter(|mode| !mode.is_empty())
3322 {
3323 return Ok(Some(AgentEvent::ModesReplaced {
3328 slot,
3329 modes: vec![Mode {
3330 id: mode.to_owned(),
3331 label: mode.to_owned(),
3332 }],
3333 current_mode: Some(mode.to_owned()),
3334 }));
3335 }
3336 match (kind, text) {
3337 (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3338 Ok(Some(AgentEvent::Text { slot, text }))
3339 }
3340 (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3341 Ok(Some(AgentEvent::Thought { slot, text }))
3342 }
3343 (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3344 let id = update
3345 .get("toolCallId")
3346 .and_then(Value::as_str)
3347 .unwrap_or("unknown-tool")
3348 .to_owned();
3349 let title = update
3350 .get("title")
3351 .and_then(Value::as_str)
3352 .unwrap_or("Tool call")
3353 .to_owned();
3354 let status = match update.get("status").and_then(Value::as_str) {
3355 Some("completed") => ToolStatus::Completed,
3356 Some("failed") => ToolStatus::Failed,
3357 Some("in_progress") => ToolStatus::Running,
3358 _ => ToolStatus::Pending,
3359 };
3360 Ok(Some(AgentEvent::Tool {
3361 slot,
3362 update: ToolUpdate {
3363 id,
3364 title,
3365 status,
3366 detail: None,
3367 },
3368 }))
3369 }
3370 _ => Ok(None),
3371 }
3372}
3373
3374fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3375 let config = value
3376 .get("configOptions")?
3377 .as_array()?
3378 .iter()
3379 .find(|option| {
3380 option.get("category").and_then(Value::as_str) == Some("model")
3381 && matches!(
3382 option.get("type").and_then(Value::as_str),
3383 Some("select" | "enum")
3384 )
3385 })?;
3386 let config_id = config.get("id")?.as_str()?.to_owned();
3387 let models = config
3388 .get("options")?
3389 .as_array()?
3390 .iter()
3391 .filter_map(|option| {
3392 let id = option.get("value")?.as_str()?.to_owned();
3393 let label = option
3394 .get("name")
3395 .or_else(|| option.get("label"))
3396 .and_then(Value::as_str)
3397 .unwrap_or(&id)
3398 .to_owned();
3399 Some(Mode { id, label })
3400 })
3401 .collect::<Vec<_>>();
3402 (!models.is_empty()).then(|| {
3403 let current = config
3404 .get("currentValue")
3405 .and_then(Value::as_str)
3406 .map(str::to_owned);
3407 (config_id, models, current)
3408 })
3409}
3410
3411fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3416 let nested = value.get("terminal").unwrap_or(value);
3417 let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3418 let id = nested
3419 .get("terminalId")
3420 .or_else(|| nested.get("terminal_id"))
3421 .or_else(|| nested.get("id"))
3422 .and_then(Value::as_str)
3423 .unwrap_or("terminal")
3424 .to_owned();
3425 match kind {
3426 "terminal_created" | "terminal_create" | "terminal_started" => {
3427 let command = nested
3428 .get("command")
3429 .and_then(Value::as_str)
3430 .unwrap_or("")
3431 .to_owned();
3432 Some(TerminalEvent::Created { id, command })
3433 }
3434 "terminal_output" | "terminal_output_chunk" => {
3435 let text = nested
3436 .get("output")
3437 .or_else(|| nested.get("text"))
3438 .and_then(Value::as_str)
3439 .unwrap_or("")
3440 .to_owned();
3441 Some(TerminalEvent::Output { id, text })
3442 }
3443 "terminal_exited" | "terminal_exit" => {
3444 let code = nested
3445 .get("exitCode")
3446 .or_else(|| nested.get("exit_code"))
3447 .or_else(|| nested.get("code"))
3448 .and_then(Value::as_i64)
3449 .unwrap_or(0) as i32;
3450 Some(TerminalEvent::Exited { id, code })
3451 }
3452 "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3453 _ => None,
3454 }
3455}
3456
3457fn parse_permission_event(
3458 slot: RosterSlot,
3459 value: &Value,
3460 request_id: &str,
3461 options: Option<&Value>,
3462) -> Option<AgentEvent> {
3463 let tool = value.get("toolCall").unwrap_or(value);
3464 let title = tool
3465 .get("title")
3466 .and_then(Value::as_str)
3467 .unwrap_or("Agent requests permission")
3468 .to_owned();
3469 let (options, option_ids): (Vec<String>, Vec<String>) = options
3470 .and_then(Value::as_array)
3471 .map(|options| {
3472 options
3473 .iter()
3474 .filter_map(|option| {
3475 let label = option
3476 .get("name")
3477 .or_else(|| option.get("optionId"))
3478 .and_then(Value::as_str)?
3479 .to_owned();
3480 let option_id = option
3481 .get("optionId")
3482 .or_else(|| option.get("id"))
3483 .and_then(Value::as_str)
3484 .map(str::to_owned)
3485 .unwrap_or_else(|| label.clone());
3486 Some((label, option_id))
3487 })
3488 .unzip()
3489 })
3490 .unwrap_or_default();
3491 if options.is_empty() {
3492 return None;
3493 }
3494 Some(AgentEvent::Permission {
3495 slot,
3496 request: PermissionRequest {
3497 id: request_id.to_owned(),
3498 title,
3499 options,
3500 option_ids,
3501 },
3502 })
3503}
3504
3505fn rpc_id_to_string(value: &Value) -> String {
3506 value
3507 .as_str()
3508 .map(str::to_owned)
3509 .or_else(|| value.as_u64().map(|id| id.to_string()))
3510 .unwrap_or_else(|| value.to_string())
3511}
3512
3513#[cfg(test)]
3514mod tests {
3515 use super::{
3516 AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
3517 RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
3518 parse_model_config, prompt_content_blocks, read_bounded_line,
3519 };
3520 #[cfg(target_os = "linux")]
3521 use super::{isolate_process_group, terminate_child};
3522 use crate::TerminalEvent;
3523 use crate::{
3524 AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
3525 persistence::SessionMetadataStore,
3526 relay::{DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
3527 };
3528 use async_trait::async_trait;
3529 use serde_json::Value;
3530 use std::sync::{
3531 Arc, Mutex,
3532 atomic::{AtomicUsize, Ordering},
3533 };
3534
3535 fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
3536 let nonce = std::time::SystemTime::now()
3537 .duration_since(std::time::UNIX_EPOCH)
3538 .expect("clock")
3539 .as_nanos();
3540 std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
3541 }
3542
3543 #[cfg(target_os = "linux")]
3544 #[tokio::test]
3545 async fn termination_kills_only_the_verified_isolated_child_group() {
3546 use nix::unistd::{Pid, getpgid, getpgrp};
3547 use tokio::io::{AsyncBufReadExt, BufReader};
3548
3549 let own_group = getpgrp();
3550 let mut command = tokio::process::Command::new("sh");
3551 isolate_process_group(&mut command);
3552 command
3553 .arg("-c")
3554 .arg("sleep 60 & echo $!; wait")
3555 .stdout(std::process::Stdio::piped());
3556 let mut child = command.spawn().expect("spawn isolated shell");
3557 let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
3558 assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
3559 assert_ne!(leader, own_group);
3560
3561 let stdout = child.stdout.take().expect("child stdout");
3562 let mut lines = BufReader::new(stdout).lines();
3563 let descendant = lines
3564 .next_line()
3565 .await
3566 .expect("read descendant pid")
3567 .expect("descendant pid")
3568 .parse::<i32>()
3569 .expect("numeric descendant pid");
3570 let descendant = Pid::from_raw(descendant);
3571 assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
3572
3573 terminate_child(&mut child).await.expect("terminate group");
3574 for _ in 0..100 {
3575 if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
3576 return;
3577 }
3578 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
3579 }
3580 panic!("descendant {descendant} survived isolated group termination");
3581 }
3582
3583 #[test]
3584 fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
3585 assert_eq!(
3586 parse_command_line(r#"npx -y "@agentclientprotocol/codex-acp" --flag 'two words'"#),
3587 Ok((
3588 "npx".into(),
3589 vec![
3590 "-y".into(),
3591 "@agentclientprotocol/codex-acp".into(),
3592 "--flag".into(),
3593 "two words".into(),
3594 ]
3595 ),)
3596 );
3597 assert_eq!(
3598 parse_command_line(r#"agent "" escaped\ argument"#),
3599 Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
3600 );
3601 }
3602
3603 #[test]
3604 fn acp_prompt_expands_safe_at_path_resources() {
3605 let root = unique_test_path("codeswarm-prompt-resource", "dir");
3606 std::fs::create_dir_all(&root).expect("workspace");
3607 std::fs::write(root.join("note.md"), "resource text").expect("resource");
3608 let blocks = prompt_content_blocks(&root, "inspect @note.md");
3609 assert_eq!(blocks[0]["type"], "text");
3610 assert_eq!(blocks[0]["text"], "inspect @note.md");
3611 assert_eq!(blocks[1]["type"], "resource");
3612 assert_eq!(blocks[1]["resource"]["text"], "resource text");
3613 assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
3614 std::fs::remove_dir_all(root).expect("cleanup workspace");
3615 }
3616
3617 #[tokio::test]
3618 async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
3619 let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
3620 bytes.push(b'\n');
3621 let mut reader = tokio::io::BufReader::new(bytes.as_slice());
3622 assert!(matches!(
3623 read_bounded_line(&mut reader).await,
3624 Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
3625 ));
3626 }
3627
3628 #[test]
3629 fn rejects_malformed_configured_commands_before_spawn() {
3630 assert_eq!(
3631 parse_command_line("agent 'unfinished"),
3632 Err(super::CommandParseError::UnterminatedQuote)
3633 );
3634 assert_eq!(
3635 parse_command_line("agent\\"),
3636 Err(super::CommandParseError::TrailingEscape)
3637 );
3638 assert_eq!(
3639 parse_command_line(" \t"),
3640 Err(super::CommandParseError::Empty)
3641 );
3642 }
3643
3644 #[derive(Debug)]
3645 struct PendingAdapter {
3646 slot: usize,
3647 hang_on_cancel: bool,
3648 }
3649
3650 #[derive(Debug)]
3651 struct ConcurrentStartAdapter {
3652 slot: usize,
3653 barrier: Arc<tokio::sync::Barrier>,
3654 }
3655
3656 #[async_trait]
3657 impl AgentAdapter for ConcurrentStartAdapter {
3658 fn slot(&self) -> usize {
3659 self.slot
3660 }
3661
3662 fn capabilities(&self) -> AgentCapabilities {
3663 AgentCapabilities::default()
3664 }
3665
3666 async fn start(&mut self) -> super::AdapterResult<()> {
3667 self.barrier.wait().await;
3668 Ok(())
3669 }
3670
3671 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
3672 Ok(())
3673 }
3674
3675 async fn cancel(&mut self) -> super::AdapterResult<bool> {
3676 Ok(true)
3677 }
3678
3679 async fn answer_permission(
3680 &mut self,
3681 _request_id: String,
3682 _answer: PermissionAnswer,
3683 ) -> super::AdapterResult<()> {
3684 Ok(())
3685 }
3686
3687 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
3688 Ok(())
3689 }
3690
3691 async fn reload(&mut self) -> super::AdapterResult<()> {
3692 Ok(())
3693 }
3694
3695 async fn stop(&mut self) -> super::AdapterResult<()> {
3696 Ok(())
3697 }
3698
3699 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
3700 std::future::pending().await
3701 }
3702 }
3703
3704 #[derive(Debug)]
3705 struct PermissionBlockingAdapter {
3706 slot: usize,
3707 phase: u8,
3708 }
3709
3710 #[async_trait]
3711 impl AgentAdapter for PermissionBlockingAdapter {
3712 fn slot(&self) -> usize {
3713 self.slot
3714 }
3715
3716 fn capabilities(&self) -> AgentCapabilities {
3717 AgentCapabilities {
3718 supports_permissions: true,
3719 ..AgentCapabilities::default()
3720 }
3721 }
3722
3723 async fn start(&mut self) -> super::AdapterResult<()> {
3724 Ok(())
3725 }
3726
3727 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
3728 Ok(())
3729 }
3730
3731 async fn cancel(&mut self) -> super::AdapterResult<bool> {
3732 Ok(true)
3733 }
3734
3735 async fn answer_permission(
3736 &mut self,
3737 request_id: String,
3738 answer: PermissionAnswer,
3739 ) -> super::AdapterResult<()> {
3740 if self.phase != 1 || request_id != "permission-1" {
3741 return Err(super::AdapterError::Protocol(
3742 "unexpected permission response".into(),
3743 ));
3744 }
3745 assert!(matches!(answer, PermissionAnswer::Selected { .. }));
3746 self.phase = 2;
3747 Ok(())
3748 }
3749
3750 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
3751 Ok(())
3752 }
3753
3754 async fn reload(&mut self) -> super::AdapterResult<()> {
3755 Ok(())
3756 }
3757
3758 async fn stop(&mut self) -> super::AdapterResult<()> {
3759 Ok(())
3760 }
3761
3762 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
3763 match self.phase {
3764 0 => {
3765 self.phase = 1;
3766 Some(Ok(AgentEvent::Permission {
3767 slot: self.slot,
3768 request: crate::PermissionRequest {
3769 id: "permission-1".into(),
3770 title: "Allow?".into(),
3771 options: vec!["Allow".into()],
3772 option_ids: vec!["allow".into()],
3773 },
3774 }))
3775 }
3776 1 => std::future::pending().await,
3777 _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
3778 }
3779 }
3780 }
3781
3782 #[async_trait]
3783 impl AgentAdapter for PendingAdapter {
3784 fn slot(&self) -> usize {
3785 self.slot
3786 }
3787
3788 fn capabilities(&self) -> AgentCapabilities {
3789 AgentCapabilities {
3790 supports_cancel: true,
3791 ..AgentCapabilities::default()
3792 }
3793 }
3794
3795 async fn start(&mut self) -> super::AdapterResult<()> {
3796 Ok(())
3797 }
3798
3799 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
3800 Ok(())
3801 }
3802
3803 async fn cancel(&mut self) -> super::AdapterResult<bool> {
3804 if self.hang_on_cancel {
3805 return std::future::pending().await;
3806 }
3807 Ok(true)
3808 }
3809
3810 async fn answer_permission(
3811 &mut self,
3812 _request_id: String,
3813 _answer: PermissionAnswer,
3814 ) -> super::AdapterResult<()> {
3815 Ok(())
3816 }
3817
3818 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
3819 Ok(())
3820 }
3821
3822 async fn reload(&mut self) -> super::AdapterResult<()> {
3823 Ok(())
3824 }
3825
3826 async fn stop(&mut self) -> super::AdapterResult<()> {
3827 Ok(())
3828 }
3829
3830 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
3831 std::future::pending().await
3832 }
3833 }
3834
3835 #[derive(Debug)]
3836 struct StopTrackingAdapter {
3837 slot: usize,
3838 stopped: Arc<AtomicUsize>,
3839 fail_stop: bool,
3840 }
3841
3842 #[derive(Debug)]
3843 struct ModeOrderAdapter {
3844 slot: usize,
3845 log: Arc<Mutex<Vec<String>>>,
3846 phase: u8,
3847 }
3848
3849 #[derive(Debug)]
3850 struct StartupAcpAdapter {
3851 slot: usize,
3852 events: std::collections::VecDeque<AgentEvent>,
3853 }
3854
3855 impl StartupAcpAdapter {
3856 fn new(slot: usize) -> Self {
3857 Self {
3858 slot,
3859 events: [
3860 AgentEvent::ModesReplaced {
3861 slot,
3862 modes: vec![Mode {
3863 id: "full-access".into(),
3864 label: "Auto pilot".into(),
3865 }],
3866 current_mode: Some("full-access".into()),
3867 },
3868 AgentEvent::Ready {
3869 slot,
3870 capabilities: AgentCapabilities {
3871 supports_modes: true,
3872 ..AgentCapabilities::default()
3873 },
3874 },
3875 ]
3876 .into(),
3877 }
3878 }
3879 }
3880
3881 #[async_trait]
3882 impl AgentAdapter for StartupAcpAdapter {
3883 fn slot(&self) -> usize {
3884 self.slot
3885 }
3886
3887 fn protocol(&self) -> &'static str {
3888 "acp"
3889 }
3890
3891 fn capabilities(&self) -> AgentCapabilities {
3892 AgentCapabilities {
3893 supports_modes: true,
3894 ..AgentCapabilities::default()
3895 }
3896 }
3897
3898 async fn start(&mut self) -> super::AdapterResult<()> {
3899 Ok(())
3900 }
3901
3902 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
3903 Ok(())
3904 }
3905
3906 async fn cancel(&mut self) -> super::AdapterResult<bool> {
3907 Ok(true)
3908 }
3909
3910 async fn answer_permission(
3911 &mut self,
3912 _request_id: String,
3913 _answer: PermissionAnswer,
3914 ) -> super::AdapterResult<()> {
3915 Ok(())
3916 }
3917
3918 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
3919 Ok(())
3920 }
3921
3922 async fn reload(&mut self) -> super::AdapterResult<()> {
3923 self.events = Self::new(self.slot).events;
3924 Ok(())
3925 }
3926
3927 async fn stop(&mut self) -> super::AdapterResult<()> {
3928 Ok(())
3929 }
3930
3931 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
3932 self.events.pop_front().map(Ok)
3933 }
3934 }
3935
3936 #[async_trait]
3937 impl AgentAdapter for ModeOrderAdapter {
3938 fn slot(&self) -> usize {
3939 self.slot
3940 }
3941
3942 fn capabilities(&self) -> AgentCapabilities {
3943 AgentCapabilities {
3944 supports_modes: true,
3945 ..AgentCapabilities::default()
3946 }
3947 }
3948
3949 async fn start(&mut self) -> super::AdapterResult<()> {
3950 self.log.lock().expect("log").push("start".into());
3951 Ok(())
3952 }
3953
3954 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
3955 self.log.lock().expect("log").push("prompt".into());
3956 Ok(())
3957 }
3958
3959 async fn cancel(&mut self) -> super::AdapterResult<bool> {
3960 Ok(true)
3961 }
3962
3963 async fn answer_permission(
3964 &mut self,
3965 _request_id: String,
3966 _answer: PermissionAnswer,
3967 ) -> super::AdapterResult<()> {
3968 Ok(())
3969 }
3970
3971 async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
3972 self.log.lock().expect("log").push(format!("mode:{mode}"));
3973 Ok(())
3974 }
3975
3976 async fn reload(&mut self) -> super::AdapterResult<()> {
3977 self.log.lock().expect("log").push("reload".into());
3978 self.phase = 0;
3979 Ok(())
3980 }
3981
3982 async fn stop(&mut self) -> super::AdapterResult<()> {
3983 Ok(())
3984 }
3985
3986 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
3987 match self.phase {
3988 0 => {
3989 self.phase = 1;
3990 Some(Ok(AgentEvent::ModesReplaced {
3991 slot: self.slot,
3992 modes: vec![Mode {
3993 id: "yolo".into(),
3994 label: "YOLO".into(),
3995 }],
3996 current_mode: None,
3997 }))
3998 }
3999 1 => {
4000 self.phase = 2;
4001 Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4002 }
4003 _ => std::future::pending().await,
4004 }
4005 }
4006 }
4007
4008 #[async_trait]
4009 impl AgentAdapter for StopTrackingAdapter {
4010 fn slot(&self) -> usize {
4011 self.slot
4012 }
4013
4014 fn capabilities(&self) -> AgentCapabilities {
4015 AgentCapabilities::default()
4016 }
4017
4018 async fn start(&mut self) -> super::AdapterResult<()> {
4019 Ok(())
4020 }
4021
4022 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4023 Ok(())
4024 }
4025
4026 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4027 Ok(false)
4028 }
4029
4030 async fn answer_permission(
4031 &mut self,
4032 _request_id: String,
4033 _answer: PermissionAnswer,
4034 ) -> super::AdapterResult<()> {
4035 Ok(())
4036 }
4037
4038 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4039 Ok(())
4040 }
4041
4042 async fn reload(&mut self) -> super::AdapterResult<()> {
4043 Ok(())
4044 }
4045
4046 async fn stop(&mut self) -> super::AdapterResult<()> {
4047 self.stopped.fetch_add(1, Ordering::Relaxed);
4048 if self.fail_stop {
4049 Err(super::AdapterError::Transport("stop failed".into()))
4050 } else {
4051 Ok(())
4052 }
4053 }
4054
4055 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4056 None
4057 }
4058 }
4059
4060 #[derive(Debug)]
4064 struct FailingStartAdapter {
4065 slot: usize,
4066 stopped: Arc<AtomicUsize>,
4067 }
4068
4069 #[async_trait]
4070 impl AgentAdapter for FailingStartAdapter {
4071 fn slot(&self) -> usize {
4072 self.slot
4073 }
4074
4075 fn capabilities(&self) -> AgentCapabilities {
4076 AgentCapabilities::default()
4077 }
4078
4079 async fn start(&mut self) -> super::AdapterResult<()> {
4080 Err(super::AdapterError::Spawn("startup failed".into()))
4081 }
4082
4083 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4084 Ok(())
4085 }
4086
4087 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4088 Ok(false)
4089 }
4090
4091 async fn answer_permission(
4092 &mut self,
4093 _request_id: String,
4094 _answer: PermissionAnswer,
4095 ) -> super::AdapterResult<()> {
4096 Ok(())
4097 }
4098
4099 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4100 Ok(())
4101 }
4102
4103 async fn reload(&mut self) -> super::AdapterResult<()> {
4104 Ok(())
4105 }
4106
4107 async fn stop(&mut self) -> super::AdapterResult<()> {
4108 self.stopped.fetch_add(1, Ordering::Relaxed);
4109 Ok(())
4110 }
4111
4112 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4113 None
4114 }
4115 }
4116
4117 #[tokio::test]
4118 async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4119 let stopped = Arc::new(AtomicUsize::new(0));
4120 let relay = RelayHost::new(
4121 vec![
4122 AdapterHost::new(
4123 Box::new(StopTrackingAdapter {
4124 slot: 0,
4125 stopped: Arc::clone(&stopped),
4126 fail_stop: true,
4127 }),
4128 None,
4129 ),
4130 AdapterHost::new(
4131 Box::new(StopTrackingAdapter {
4132 slot: 1,
4133 stopped: Arc::clone(&stopped),
4134 fail_stop: false,
4135 }),
4136 None,
4137 ),
4138 ],
4139 4,
4140 )
4141 .expect("relay");
4142 let mut relay = relay;
4143
4144 let error = relay.stop().await.expect_err("first stop failure");
4145 assert!(error.to_string().contains("stop failed"));
4146 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4147 }
4148
4149 #[tokio::test]
4150 async fn relay_start_cleans_up_the_adapter_that_failed_startup() {
4151 let stopped = Arc::new(AtomicUsize::new(0));
4152 let mut relay = RelayHost::new(
4153 vec![
4154 AdapterHost::new(
4155 Box::new(StopTrackingAdapter {
4156 slot: 0,
4157 stopped: Arc::clone(&stopped),
4158 fail_stop: false,
4159 }),
4160 None,
4161 ),
4162 AdapterHost::new(
4163 Box::new(FailingStartAdapter {
4164 slot: 1,
4165 stopped: Arc::clone(&stopped),
4166 }),
4167 None,
4168 ),
4169 ],
4170 4,
4171 )
4172 .expect("relay");
4173
4174 assert!(relay.start().await.is_err());
4175 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4176 }
4177
4178 #[test]
4179 fn parses_acp_text_without_ui_dependency() {
4180 let event = parse_acp_notification(
4181 2,
4182 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4183 )
4184 .expect("valid ACP")
4185 .expect("text event");
4186 assert_eq!(
4187 event,
4188 AgentEvent::Text {
4189 slot: 2,
4190 text: "hello".into(),
4191 }
4192 );
4193 }
4194
4195 #[test]
4196 fn parses_acp_state_notifications_at_the_adapter_boundary() {
4197 let commands = parse_acp_notification(
4198 3,
4199 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
4200 )
4201 .expect("valid ACP")
4202 .expect("commands event");
4203 assert_eq!(
4204 commands,
4205 AgentEvent::CommandsReplaced {
4206 slot: 3,
4207 commands: vec![crate::AgentCommand {
4208 name: "review".into()
4209 }]
4210 }
4211 );
4212
4213 let mode = parse_acp_notification(
4214 3,
4215 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
4216 )
4217 .expect("valid ACP")
4218 .expect("mode event");
4219 assert_eq!(
4220 mode,
4221 AgentEvent::ModeUpdated {
4222 slot: 3,
4223 current_mode: "review".into()
4224 }
4225 );
4226
4227 let usage = parse_acp_notification(
4228 3,
4229 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
4230 )
4231 .expect("valid ACP")
4232 .expect("usage event");
4233 assert_eq!(
4234 usage,
4235 AgentEvent::UsageUpdated {
4236 slot: 3,
4237 usage: crate::UsageUpdate {
4238 used: 4200,
4239 size: 128000
4240 }
4241 }
4242 );
4243
4244 let models = parse_acp_notification(
4245 3,
4246 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"}]}]}}}"#,
4247 )
4248 .expect("valid ACP")
4249 .expect("models event");
4250 assert!(matches!(
4251 models,
4252 AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
4253 if models.len() == 2 && current_model.as_deref() == Some("smart")
4254 ));
4255 assert_eq!(
4256 parse_model_config(&serde_json::json!({
4257 "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
4258 })),
4259 None
4260 );
4261
4262 let user = parse_acp_notification(
4263 3,
4264 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
4265 )
4266 .expect("valid ACP")
4267 .expect("user event");
4268 assert_eq!(
4269 user,
4270 AgentEvent::UserText {
4271 slot: 3,
4272 text: "context".into()
4273 }
4274 );
4275 }
4276
4277 #[test]
4278 fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
4279 let event = parse_acp_notification(
4280 0,
4281 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
4282 )
4283 .expect("valid ACP")
4284 .expect("mode event");
4285 assert!(matches!(
4286 event,
4287 AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
4288 if mode == "yolo" && modes[0].id == "yolo"
4289 ));
4290 }
4291
4292 #[test]
4293 fn parses_native_agy_text_without_acp_bridge() {
4294 let event = parse_agy_line(
4295 1,
4296 r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
4297 )
4298 .expect("valid stream-json")
4299 .expect("text event");
4300 assert_eq!(
4301 event,
4302 AgentEvent::Text {
4303 slot: 1,
4304 text: "hello".into(),
4305 }
4306 );
4307 }
4308
4309 #[test]
4310 fn parses_tool_lifecycle_from_each_protocol() {
4311 let agy = parse_agy_line(
4312 1,
4313 r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
4314 )
4315 .expect("valid native tool")
4316 .expect("tool event");
4317 assert!(matches!(
4318 agy,
4319 AgentEvent::Tool {
4320 update: crate::ToolUpdate {
4321 status: ToolStatus::Completed,
4322 ..
4323 },
4324 ..
4325 }
4326 ));
4327
4328 let acp = parse_acp_notification(
4329 1,
4330 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
4331 )
4332 .expect("valid ACP tool")
4333 .expect("tool event");
4334 assert!(matches!(
4335 acp,
4336 AgentEvent::Tool {
4337 update: crate::ToolUpdate {
4338 status: ToolStatus::Failed,
4339 ..
4340 },
4341 ..
4342 }
4343 ));
4344 }
4345
4346 #[test]
4347 fn parses_terminal_lifecycle_from_acp_and_native_events() {
4348 let created = parse_acp_notification(
4349 0,
4350 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
4351 )
4352 .expect("valid ACP terminal")
4353 .expect("terminal event");
4354 assert_eq!(
4355 created,
4356 AgentEvent::Terminal {
4357 slot: 0,
4358 event: TerminalEvent::Created {
4359 id: "term-1".into(),
4360 command: "cargo test".into(),
4361 },
4362 }
4363 );
4364 let output = parse_agy_line(
4365 1,
4366 r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
4367 )
4368 .expect("valid native terminal")
4369 .expect("terminal event");
4370 assert_eq!(
4371 output,
4372 AgentEvent::Terminal {
4373 slot: 1,
4374 event: TerminalEvent::Output {
4375 id: "term-1".into(),
4376 text: "ok\n".into(),
4377 },
4378 }
4379 );
4380 let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
4381 .expect("valid native release")
4382 .expect("terminal event");
4383 assert!(matches!(
4384 released,
4385 AgentEvent::Terminal {
4386 event: TerminalEvent::Released { id },
4387 ..
4388 } if id == "term-1"
4389 ));
4390 }
4391
4392 #[test]
4393 fn parses_acp_permission_requests() {
4394 let event = parse_acp_notification(
4395 0,
4396 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"}]}}}"#,
4397 )
4398 .expect("valid permission")
4399 .expect("permission event");
4400 assert!(matches!(
4401 event,
4402 AgentEvent::Permission { request, .. }
4403 if request.id == "t1"
4404 && request.title == "Write file"
4405 && request.options == ["Allow once", "Reject"]
4406 && request.option_ids == ["allow-once", "reject"]
4407 ));
4408 }
4409
4410 #[test]
4411 fn parses_acp_permission_request_as_json_rpc_request() {
4412 let event = parse_acp_notification(
4413 2,
4414 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
4415 )
4416 .expect("valid permission request")
4417 .expect("permission event");
4418 assert!(matches!(
4419 event,
4420 AgentEvent::Permission { request, .. }
4421 if request.id == "17"
4422 && request.title == "Write file"
4423 && request.options == ["allow-once", "reject"]
4424 && request.option_ids == ["allow-once", "reject"]
4425 ));
4426 }
4427
4428 #[tokio::test]
4429 async fn native_adapter_explicitly_rejects_permission_answers() {
4430 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4431 assert_eq!(
4432 adapter
4433 .answer_permission(
4434 "request".into(),
4435 PermissionAnswer::Selected {
4436 option_id: "allow".into()
4437 },
4438 )
4439 .await,
4440 Err(super::AdapterError::Unsupported("permission answer"))
4441 );
4442 }
4443
4444 #[tokio::test]
4445 async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
4446 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4447 adapter
4448 .set_mode("full-access".into())
4449 .await
4450 .expect("auto-pilot alias");
4451 assert!(matches!(
4452 adapter.next_event().await,
4453 Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
4454 ));
4455 }
4456
4457 #[tokio::test]
4458 async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
4459 let script_path = unique_test_path("codeswarm-agy-session", "sh");
4460 std::fs::write(
4461 &script_path,
4462 "#!/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",
4463 )
4464 .expect("write native test script");
4465 let mut adapter = AgyAdapter::new(
4466 0,
4467 std::env::current_dir().expect("cwd"),
4468 format!("sh {}", script_path.display()),
4469 );
4470 adapter.start().await.expect("start native adapter");
4471 assert!(adapter.next_event().await.is_some());
4473 assert!(adapter.next_event().await.is_some());
4474 adapter
4475 .send_prompt("first".into())
4476 .await
4477 .expect("first prompt");
4478 while !matches!(
4479 adapter.next_event().await,
4480 Some(Ok(AgentEvent::TurnComplete { .. }))
4481 ) {}
4482 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4483 adapter
4484 .send_prompt("follow up".into())
4485 .await
4486 .expect("follow-up prompt");
4487 while !matches!(
4488 adapter.next_event().await,
4489 Some(Ok(AgentEvent::TurnComplete { .. }))
4490 ) {}
4491 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4492 adapter.stop().await.expect("stop native adapter");
4493 std::fs::remove_file(script_path).expect("cleanup native script");
4494 }
4495
4496 #[tokio::test]
4497 async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
4498 let script_path = unique_test_path("codeswarm-agy-failure", "sh");
4499 std::fs::write(
4500 &script_path,
4501 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
4502 )
4503 .expect("write native test script");
4504 let mut adapter = AgyAdapter::new(
4505 0,
4506 std::env::current_dir().expect("cwd"),
4507 format!("sh {}", script_path.display()),
4508 );
4509 adapter.start().await.expect("start native adapter");
4510 assert!(adapter.next_event().await.is_some());
4511 assert!(adapter.next_event().await.is_some());
4512 adapter.send_prompt("fail".into()).await.expect("prompt");
4513 assert!(matches!(
4514 adapter.next_event().await,
4515 Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
4516 if detail == "agent failed"
4517 ));
4518 adapter.stop().await.expect("stop native adapter");
4519 std::fs::remove_file(script_path).expect("cleanup native script");
4520 }
4521
4522 #[tokio::test]
4523 async fn acp_adapter_initializes_session_and_completes_a_prompt() {
4524 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"}}'"#;
4525 let cwd = std::env::current_dir().expect("cwd");
4526 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4527 adapter.start().await.expect("initialize");
4528 assert!(matches!(
4529 adapter.next_event().await,
4530 Some(Ok(AgentEvent::ModesReplaced { .. }))
4531 ));
4532 assert!(matches!(
4533 adapter.next_event().await,
4534 Some(Ok(AgentEvent::Ready { .. }))
4535 ));
4536 adapter.send_prompt("hello".into()).await.expect("prompt");
4537 assert!(matches!(
4538 adapter.next_event().await,
4539 Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
4540 ));
4541 assert!(matches!(
4542 adapter.next_event().await,
4543 Some(Ok(AgentEvent::TurnComplete { .. }))
4544 ));
4545 }
4546
4547 #[tokio::test]
4548 async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
4549 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"}}'"#;
4550 let cwd = std::env::current_dir().expect("cwd");
4551 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
4552 adapter.start().await.expect("initialize");
4553 assert!(adapter.next_event().await.is_some());
4554 assert!(adapter.next_event().await.is_some());
4555
4556 for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
4557 adapter.send_prompt(prompt.into()).await.expect("prompt");
4558 assert!(matches!(
4559 adapter.next_event().await,
4560 Some(Ok(AgentEvent::Text { text, .. })) if text == expected
4561 ));
4562 assert!(matches!(
4563 adapter.next_event().await,
4564 Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
4565 ));
4566 }
4567 adapter.stop().await.expect("stop");
4568 }
4569
4570 #[tokio::test]
4571 async fn empty_acp_mode_catalog_disables_mode_control() {
4572 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":[]}}}'"#;
4573 let cwd = std::env::current_dir().expect("cwd");
4574 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4575 adapter.start().await.expect("initialize");
4576 assert!(!adapter.capabilities().supports_modes);
4577 assert!(matches!(
4578 adapter.next_event().await,
4579 Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
4580 ));
4581 assert!(matches!(
4582 adapter.next_event().await,
4583 Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
4584 ));
4585 adapter.stop().await.expect("stop");
4586 }
4587
4588 #[tokio::test]
4589 async fn acp_models_are_discovered_live_and_changed_through_session_config() {
4590 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"#;
4591 let cwd = std::env::current_dir().expect("cwd");
4592 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4593 adapter.start().await.expect("initialize");
4594 assert!(adapter.capabilities().supports_models);
4595 assert!(matches!(
4596 adapter.next_event().await,
4597 Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
4598 if config_id == "model"
4599 && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
4600 && current_model.as_deref() == Some("fast")
4601 ));
4602 assert!(matches!(
4603 adapter.next_event().await,
4604 Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
4605 ));
4606 adapter.set_model("smart".into()).await.expect("set model");
4607 assert!(adapter.set_model("invented".into()).await.is_err());
4608 adapter.stop().await.expect("stop");
4609 }
4610
4611 #[tokio::test]
4612 async fn acp_mode_change_is_acknowledged_without_provider_notification() {
4613 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":{}}'"#;
4614 let cwd = std::env::current_dir().expect("cwd");
4615 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4616 adapter.start().await.expect("initialize");
4617 adapter
4618 .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
4619 .await
4620 .expect("set mode");
4621 assert!(matches!(
4622 adapter.next_event().await,
4623 Some(Ok(AgentEvent::ModesReplaced { .. }))
4624 ));
4625 assert!(matches!(
4626 adapter.next_event().await,
4627 Some(Ok(AgentEvent::Ready { .. }))
4628 ));
4629 assert!(matches!(
4630 adapter.next_event().await,
4631 Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
4632 ));
4633 adapter.stop().await.expect("stop");
4634 }
4635
4636 #[tokio::test]
4637 async fn acp_reload_preserves_a_loadable_session_id() {
4638 let cwd = std::env::current_dir().expect("cwd");
4639 let mut adapter = AcpAdapter::with_session_id(
4640 0,
4641 cwd,
4642 "__codeswarm_missing_acp_for_reload_test__",
4643 Vec::new(),
4644 "saved-session",
4645 );
4646 adapter.capabilities.supports_session_load = true;
4647 assert!(adapter.reload().await.is_err());
4651 assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
4652 }
4653
4654 #[tokio::test]
4655 async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
4656 let cwd = std::env::current_dir().expect("cwd");
4657 let mut adapter = AcpAdapter::with_session_id(
4658 0,
4659 cwd,
4660 "__codeswarm_missing_nonloadable_acp__",
4661 Vec::new(),
4662 "stale-session",
4663 );
4664 adapter.capabilities.supports_session_load = false;
4665 assert!(adapter.reload().await.is_err());
4666 assert_eq!(adapter.session_id, None);
4667 }
4668
4669 #[tokio::test]
4670 async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
4671 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"}}'"#;
4672 let cwd = std::env::current_dir().expect("cwd");
4673 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4674 adapter.start().await.expect("initialize");
4675 assert!(matches!(
4676 adapter.next_event().await,
4677 Some(Ok(AgentEvent::Ready { .. }))
4678 ));
4679 adapter.send_prompt("hello".into()).await.expect("prompt");
4680 assert!(matches!(
4681 adapter.next_event().await,
4682 Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
4683 ));
4684 assert!(matches!(
4685 adapter.next_event().await,
4686 Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
4687 ));
4688 }
4689
4690 #[tokio::test]
4691 async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
4692 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
4693 let cwd = std::env::current_dir().expect("cwd");
4694 let mut adapter = AcpAdapter::with_session_id(
4695 0,
4696 cwd,
4697 "sh",
4698 vec!["-c".into(), script.into()],
4699 "existing-session",
4700 );
4701 adapter.start().await.expect("load existing session");
4702 assert!(matches!(
4703 adapter.next_event().await,
4704 Some(Ok(AgentEvent::Ready { .. }))
4705 ));
4706 }
4707
4708 #[tokio::test]
4709 async fn acp_start_failure_reaps_transport_process() {
4710 let mut adapter = AcpAdapter::new(
4715 0,
4716 std::env::current_dir().expect("cwd"),
4717 "sh",
4718 vec!["-c".into(), "printf 'not-json\\n'".into()],
4719 );
4720 assert!(adapter.start().await.is_err());
4721 assert!(adapter.child.is_none());
4722 assert!(adapter.reader.is_none());
4723 }
4724
4725 #[tokio::test]
4726 async fn acp_adapter_answers_permission_json_rpc_requests() {
4727 let path = std::env::temp_dir().join(format!(
4728 "codeswarm-permission-answer-{}",
4729 std::process::id()
4730 ));
4731 let script = format!(
4732 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"}}}}'"#,
4733 path.display()
4734 );
4735 let mut adapter = AcpAdapter::new(
4736 0,
4737 std::env::current_dir().expect("cwd"),
4738 "sh",
4739 vec!["-c".into(), script],
4740 );
4741 adapter.start().await.expect("start ACP");
4742 assert!(matches!(
4743 adapter.next_event().await,
4744 Some(Ok(AgentEvent::Ready { .. }))
4745 ));
4746 adapter.send_prompt("do it".into()).await.expect("prompt");
4747 assert!(matches!(
4748 adapter.next_event().await,
4749 Some(Ok(AgentEvent::Permission { request, .. }))
4750 if request.id == "9"
4751 && request.options == ["Allow once"]
4752 && request.option_ids == ["allow-once"]
4753 ));
4754 adapter
4755 .answer_permission(
4756 "9".into(),
4757 PermissionAnswer::Selected {
4758 option_id: "allow-once".into(),
4759 },
4760 )
4761 .await
4762 .expect("permission answer");
4763 assert!(matches!(
4764 adapter.next_event().await,
4765 Some(Ok(AgentEvent::TurnComplete { .. }))
4766 ));
4767 let answer: Value = serde_json::from_str(
4768 &std::fs::read_to_string(&path).expect("captured permission answer"),
4769 )
4770 .expect("valid JSON-RPC answer");
4771 assert_eq!(answer["id"], 9);
4772 assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
4773 assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
4774 std::fs::remove_file(path).expect("cleanup");
4775 }
4776
4777 #[test]
4778 fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
4779 let event = parse_acp_notification(
4780 0,
4781 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
4782 )
4783 .expect("valid JSON-RPC request");
4784 assert!(event.is_none());
4785 }
4786
4787 #[tokio::test]
4788 async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
4789 let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
4790 std::fs::write(
4791 &script_path,
4792 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
4793 )
4794 .expect("write native test script");
4795 let mut adapter = AgyAdapter::new(
4796 0,
4797 std::env::current_dir().expect("cwd"),
4798 format!("sh {}", script_path.display()),
4799 );
4800 adapter.start().await.expect("start native adapter");
4801 assert!(adapter.next_event().await.is_some());
4802 assert!(adapter.next_event().await.is_some());
4803 adapter
4804 .send_prompt("continue".into())
4805 .await
4806 .expect("prompt");
4807 assert!(matches!(
4808 adapter.next_event().await,
4809 Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
4810 ));
4811 assert!(matches!(
4812 adapter.next_event().await,
4813 Some(Ok(AgentEvent::TurnComplete { .. }))
4814 ));
4815 adapter.stop().await.expect("stop native adapter");
4816 std::fs::remove_file(script_path).expect("cleanup native script");
4817 }
4818
4819 #[test]
4820 fn acp_workspace_file_access_is_root_bound_and_size_limited() {
4821 let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
4822 let _ = std::fs::remove_dir_all(&root);
4823 std::fs::create_dir_all(&root).expect("workspace");
4824 std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
4825 let outside =
4826 std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
4827 std::fs::write(&outside, "secret").expect("outside file");
4828 let link = root.join("outside-link");
4829 #[cfg(unix)]
4830 std::os::unix::fs::symlink(&outside, &link).expect("symlink");
4831 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
4832
4833 assert_eq!(
4834 adapter
4835 .read_workspace_text("inside.txt", Some(2), Some(1))
4836 .expect("read inside"),
4837 "two"
4838 );
4839 std::fs::write(
4840 root.join("large.txt"),
4841 vec![b'x'; MAX_FILE_READ_BYTES + 1024],
4842 )
4843 .expect("large file");
4844 let bounded = adapter
4845 .read_workspace_text("large.txt", None, None)
4846 .expect("bounded read");
4847 assert!(bounded.len() <= MAX_FILE_READ_BYTES);
4848 #[cfg(unix)]
4849 {
4850 std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
4851 .expect("internal symlink");
4852 assert_eq!(
4853 adapter
4854 .read_workspace_text("inside-link", None, None)
4855 .expect("read internal symlink"),
4856 "one\ntwo\nthree\n"
4857 );
4858 }
4859 assert!(adapter.workspace_path("../codeswarm-outside").is_err());
4860 assert!(
4861 adapter
4862 .workspace_path(&outside.display().to_string())
4863 .is_err()
4864 );
4865 #[cfg(unix)]
4866 assert!(adapter.workspace_path("outside-link").is_err());
4867 #[cfg(unix)]
4868 std::fs::remove_file(link).expect("cleanup symlink");
4869 #[cfg(unix)]
4870 std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
4871 std::fs::remove_file(outside).expect("cleanup outside");
4872 std::fs::remove_dir_all(root).expect("cleanup workspace");
4873 }
4874
4875 #[tokio::test]
4876 async fn running_terminal_output_omits_exit_status_until_completion() {
4877 let root = unique_test_path("codeswarm-terminal-output", "dir");
4878 std::fs::create_dir_all(&root).expect("workspace");
4879 let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
4880 let result = adapter
4881 .terminal_create(&serde_json::json!({
4882 "command": "sh",
4883 "args": ["-c", "sleep 0.2; printf done"],
4884 "cwd": ".",
4885 }))
4886 .await
4887 .expect("terminal create");
4888 let id = result["terminalId"].as_str().expect("terminal id");
4889 let output = adapter.terminal_output(id).await.expect("terminal output");
4890 assert!(output.get("exitStatus").is_none());
4891 if let Some(terminal) = adapter.terminals.remove(id) {
4892 terminal.stop().await;
4893 }
4894 std::fs::remove_dir_all(root).expect("cleanup workspace");
4895 }
4896
4897 #[tokio::test]
4898 async fn acp_adapter_answers_workspace_read_requests() {
4899 let root =
4900 std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
4901 let _ = std::fs::remove_dir_all(&root);
4902 std::fs::create_dir_all(&root).expect("workspace");
4903 let source = root.join("inside.txt");
4904 let answer = root.join("answer.json");
4905 std::fs::write(&source, "workspace content").expect("source");
4906 let script = format!(
4907 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"}}}}'"#,
4908 source.display(),
4909 answer.display(),
4910 );
4911 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
4912 adapter.start().await.expect("start ACP");
4913 assert!(matches!(
4914 adapter.next_event().await,
4915 Some(Ok(AgentEvent::Ready { .. }))
4916 ));
4917 adapter.send_prompt("read it".into()).await.expect("prompt");
4918 assert!(matches!(
4919 adapter.next_event().await,
4920 Some(Ok(AgentEvent::TurnComplete { .. }))
4921 ));
4922 let response: Value =
4923 serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
4924 .expect("response JSON");
4925 assert_eq!(response["id"], 9);
4926 assert_eq!(response["result"]["content"], "workspace content");
4927 adapter.stop().await.expect("stop ACP");
4928 std::fs::remove_dir_all(root).expect("cleanup workspace");
4929 }
4930
4931 #[tokio::test]
4932 async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
4933 let root =
4934 std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
4935 let _ = std::fs::remove_dir_all(&root);
4936 std::fs::create_dir_all(&root).expect("workspace");
4937 let create_request = serde_json::json!({
4938 "jsonrpc": "2.0",
4939 "id": 9,
4940 "method": "terminal/create",
4941 "params": {
4942 "sessionId": "s1",
4943 "command": "sh",
4944 "args": ["-c", "sleep 0.1; printf terminal-ok"],
4945 "cwd": ".",
4946 },
4947 });
4948 let wait_request = serde_json::json!({
4949 "jsonrpc": "2.0",
4950 "id": 10,
4951 "method": "terminal/wait_for_exit",
4952 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
4953 });
4954 let output_request = serde_json::json!({
4955 "jsonrpc": "2.0",
4956 "id": 11,
4957 "method": "terminal/output",
4958 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
4959 });
4960 let create_answer = root.join("create-answer.json");
4961 let wait_answer = root.join("wait-answer.json");
4962 let output_answer = root.join("output-answer.json");
4963 let script = format!(
4964 "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\"}}}}'",
4965 create_request,
4966 create_answer.display(),
4967 wait_request,
4968 wait_answer.display(),
4969 output_request,
4970 output_answer.display(),
4971 );
4972 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
4973 adapter.start().await.expect("start ACP");
4974 assert!(matches!(
4975 adapter.next_event().await,
4976 Some(Ok(AgentEvent::Ready { .. }))
4977 ));
4978 adapter
4979 .send_prompt("run terminal".into())
4980 .await
4981 .expect("prompt");
4982 let mut saw_complete = false;
4983 for _ in 0..6 {
4984 match adapter.next_event().await {
4985 Some(Ok(AgentEvent::TurnComplete { .. })) => {
4986 saw_complete = true;
4987 break;
4988 }
4989 Some(_) => {}
4990 None => break,
4991 }
4992 }
4993 assert!(saw_complete, "terminal requests should not stall ACP");
4994 let create: Value = serde_json::from_str(
4995 &std::fs::read_to_string(&create_answer).expect("captured create response"),
4996 )
4997 .expect("create JSON");
4998 assert_eq!(create["result"]["terminalId"], "terminal-1");
4999 let output: Value = serde_json::from_str(
5000 &std::fs::read_to_string(&output_answer).expect("captured output response"),
5001 )
5002 .expect("output JSON");
5003 assert!(
5004 output["result"]["output"]
5005 .as_str()
5006 .unwrap_or_default()
5007 .contains("terminal-ok"),
5008 "output response: {output}"
5009 );
5010 adapter.stop().await.expect("stop ACP");
5011 std::fs::remove_dir_all(root).expect("cleanup workspace");
5012 }
5013
5014 #[tokio::test]
5015 async fn host_reduces_and_persists_adapter_events() {
5016 let path =
5017 std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
5018 let adapter = ScriptedAdapter::new(
5019 0,
5020 AgentCapabilities::default(),
5021 [AgentEvent::Text {
5022 slot: 0,
5023 text: "hello".into(),
5024 }],
5025 );
5026 let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
5027 host.start().await.expect("start");
5028 host.next_effects()
5029 .await
5030 .expect("event")
5031 .expect("valid event");
5032 assert_eq!(host.state.public_text[0].1, "hello");
5033 assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
5034 std::fs::remove_file(path).expect("cleanup");
5035 }
5036
5037 #[tokio::test]
5038 async fn relay_applies_default_policy_before_the_first_prompt() {
5039 let first_log = Arc::new(Mutex::new(Vec::new()));
5040 let second_log = Arc::new(Mutex::new(Vec::new()));
5041 let first = AdapterHost::new(
5042 Box::new(ModeOrderAdapter {
5043 slot: 0,
5044 log: Arc::clone(&first_log),
5045 phase: 0,
5046 }),
5047 None,
5048 );
5049 let second = AdapterHost::new(
5050 Box::new(ModeOrderAdapter {
5051 slot: 1,
5052 log: Arc::clone(&second_log),
5053 phase: 0,
5054 }),
5055 None,
5056 );
5057 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
5058 relay.start().await.expect("start and synchronize policy");
5059 relay.run_turn("task", 0).await.expect("first turn");
5060 {
5061 let log = first_log.lock().expect("log");
5062 assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
5063 }
5064 assert_eq!(
5065 second_log.lock().expect("log").as_slice(),
5066 ["start", "mode:yolo"]
5067 );
5068 let added_log = Arc::new(Mutex::new(Vec::new()));
5069 relay
5070 .add_agent(
5071 AdapterHost::new(
5072 Box::new(ModeOrderAdapter {
5073 slot: 2,
5074 log: Arc::clone(&added_log),
5075 phase: 0,
5076 }),
5077 None,
5078 ),
5079 "Added",
5080 "added.example",
5081 "added-agent",
5082 )
5083 .await
5084 .expect("add with synchronized policy");
5085 assert_eq!(
5086 added_log.lock().expect("log").as_slice(),
5087 ["start", "mode:yolo"]
5088 );
5089 relay.drop_agent(2).await.expect("drop added agent");
5090 added_log.lock().expect("log").clear();
5091 relay
5092 .reload(2)
5093 .await
5094 .expect("reload with synchronized policy");
5095 assert_eq!(
5096 added_log.lock().expect("log").as_slice(),
5097 ["reload", "mode:yolo"]
5098 );
5099 }
5100
5101 #[tokio::test]
5102 async fn acp_roster_is_ready_before_any_prompt_is_sent() {
5103 let hosts = (0..2)
5104 .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
5105 .collect::<Vec<_>>();
5106 let startup_events = Arc::new(Mutex::new(Vec::new()));
5107 let captured = Arc::clone(&startup_events);
5108 let mut relay = RelayHost::new(hosts, 4).expect("relay");
5109 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5110
5111 relay.start().await.expect("complete startup handshake");
5112
5113 assert!(relay.dispatches().is_empty());
5114 let ready_slots = startup_events
5115 .lock()
5116 .expect("events")
5117 .iter()
5118 .filter_map(|event| match event {
5119 AgentEvent::Ready { slot, .. } => Some(*slot),
5120 _ => None,
5121 })
5122 .collect::<Vec<_>>();
5123 assert_eq!(ready_slots, vec![0, 1]);
5124 }
5125
5126 #[tokio::test]
5127 async fn independent_roster_adapters_start_concurrently() {
5128 let barrier = Arc::new(tokio::sync::Barrier::new(2));
5129 let hosts = (0..2)
5130 .map(|slot| {
5131 AdapterHost::new(
5132 Box::new(ConcurrentStartAdapter {
5133 slot,
5134 barrier: Arc::clone(&barrier),
5135 }),
5136 None,
5137 )
5138 })
5139 .collect::<Vec<_>>();
5140 let mut relay = RelayHost::new(hosts, 4).expect("relay");
5141 tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
5142 .await
5143 .expect("startup should not serialize barrier participants")
5144 .expect("startup succeeds");
5145 }
5146
5147 #[tokio::test]
5148 async fn relay_host_dispatches_turns_sequentially() {
5149 let capabilities = AgentCapabilities {
5150 supports_cancel: true,
5151 ..AgentCapabilities::default()
5152 };
5153 let first = ScriptedAdapter::new(
5154 0,
5155 capabilities.clone(),
5156 [
5157 AgentEvent::Text {
5158 slot: 0,
5159 text: "first".into(),
5160 },
5161 AgentEvent::TurnComplete { slot: 0 },
5162 ],
5163 );
5164 let second = ScriptedAdapter::new(
5165 1,
5166 capabilities,
5167 [
5168 AgentEvent::Text {
5169 slot: 1,
5170 text: "review".into(),
5171 },
5172 AgentEvent::TurnComplete { slot: 1 },
5173 ],
5174 );
5175 let hosts = vec![
5176 AdapterHost::new(Box::new(first), None),
5177 AdapterHost::new(Box::new(second), None),
5178 ];
5179 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5180 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
5181 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5182 let captured = std::sync::Arc::clone(&events);
5183 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5184 relay.start().await.expect("start");
5185 events.lock().expect("events").clear();
5186 assert!(matches!(
5187 relay.run_turn("task", 0).await.expect("first turn"),
5188 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5189 ));
5190 assert!(matches!(
5191 relay.run_turn("first", 0).await.expect("second turn"),
5192 crate::relay::RelayDecision::Dispatch {
5193 slot: 1,
5194 can_stop: true,
5195 ..
5196 }
5197 ));
5198 assert_eq!(
5199 relay
5200 .dispatches()
5201 .iter()
5202 .map(|(slot, _)| *slot)
5203 .collect::<Vec<_>>(),
5204 [0, 1]
5205 );
5206 assert!(relay.dispatches()[0].1.contains("You are Claude"));
5207 assert!(
5208 relay.dispatches()[0]
5209 .1
5210 .contains("CodeSwarm roster (ordered)")
5211 );
5212 assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
5213 assert!(relay.dispatches()[0].1.contains("2. Codex"));
5214 assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
5215 assert!(relay.dispatches()[0].1.contains("Do not use"));
5216 let lifecycle = events.lock().expect("events");
5217 let positions = lifecycle
5218 .iter()
5219 .filter_map(|event| match event {
5220 AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
5221 AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
5222 _ => None,
5223 })
5224 .collect::<Vec<_>>();
5225 assert_eq!(
5226 positions,
5227 [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
5228 );
5229 }
5230
5231 #[tokio::test]
5232 async fn relay_host_routes_around_a_usage_limited_agent() {
5233 let capabilities = AgentCapabilities::default();
5234 let first = ScriptedAdapter::new(
5235 0,
5236 capabilities.clone(),
5237 [
5238 AgentEvent::Text {
5239 slot: 0,
5240 text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
5241 credits or try again later."
5242 .into(),
5243 },
5244 AgentEvent::TurnComplete { slot: 0 },
5245 ],
5246 );
5247 let second = ScriptedAdapter::new(
5248 1,
5249 capabilities,
5250 [
5251 AgentEvent::Text {
5252 slot: 1,
5253 text: "review done".into(),
5254 },
5255 AgentEvent::TurnComplete { slot: 1 },
5256 ],
5257 );
5258 let hosts = vec![
5259 AdapterHost::new(Box::new(first), None),
5260 AdapterHost::new(Box::new(second), None),
5261 ];
5262 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5263 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5264 let captured = std::sync::Arc::clone(&events);
5265 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5266 relay.start().await.expect("start");
5267 events.lock().expect("events").clear();
5268 assert!(matches!(
5269 relay.run_turn("task", 0).await.expect("limited turn"),
5270 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5271 ));
5272 assert!(
5273 events
5274 .lock()
5275 .expect("events")
5276 .iter()
5277 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
5278 );
5279 assert!(matches!(
5281 relay.run_turn("", 0).await.expect("next turn"),
5282 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
5283 ));
5284 assert!(relay.relay().is_limited(0));
5285 relay.reload(0).await.expect("reload");
5289 assert!(!relay.relay().is_limited(0));
5290 }
5291
5292 #[tokio::test]
5293 async fn relay_failure_is_tombstoned_and_reported_to_the_ui_sink() {
5294 let failed = ScriptedAdapter::new(
5295 0,
5296 AgentCapabilities::default(),
5297 [AgentEvent::Failed {
5298 slot: 0,
5299 started: true,
5300 detail: "connection lost".into(),
5301 }],
5302 );
5303 let healthy = ScriptedAdapter::new(
5304 1,
5305 AgentCapabilities::default(),
5306 [AgentEvent::TurnComplete { slot: 1 }],
5307 );
5308 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5309 let captured = std::sync::Arc::clone(&events);
5310 let mut relay = RelayHost::new(
5311 vec![
5312 AdapterHost::new(Box::new(failed), None),
5313 AdapterHost::new(Box::new(healthy), None),
5314 ],
5315 4,
5316 )
5317 .expect("relay");
5318 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5319 relay.start().await.expect("start");
5320
5321 let error = relay.run_turn("task", 0).await.expect_err("failure");
5322 assert!(error.to_string().contains("connection lost"));
5323 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
5324 assert!(events.lock().expect("lock").iter().any(|event| {
5325 matches!(
5326 event,
5327 AgentEvent::Failed {
5328 slot: 0,
5329 started: true,
5330 ..
5331 }
5332 )
5333 }));
5334 }
5335
5336 #[tokio::test]
5337 async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
5338 let first = ScriptedAdapter::new(
5339 0,
5340 AgentCapabilities::default(),
5341 [
5342 AgentEvent::Text {
5343 slot: 0,
5344 text: "done".into(),
5345 },
5346 AgentEvent::TurnComplete { slot: 0 },
5347 ],
5348 );
5349 let reviewer = ScriptedAdapter::new(
5350 1,
5351 AgentCapabilities::default(),
5352 [
5353 AgentEvent::Text {
5354 slot: 1,
5355 text: STOP_TOKEN.into(),
5356 },
5357 AgentEvent::TurnComplete { slot: 1 },
5358 ],
5359 );
5360 let mut relay = RelayHost::new(
5361 vec![
5362 AdapterHost::new(Box::new(first), None),
5363 AdapterHost::new(Box::new(reviewer), None),
5364 ],
5365 10,
5366 )
5367 .expect("relay");
5368 relay.start().await.expect("start");
5369 let first_decision = relay.run_turn("task", 0).await.expect("first");
5370 assert!(matches!(
5371 first_decision,
5372 RelayDecision::Dispatch { slot: 0, .. }
5373 ));
5374 let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
5375 assert!(matches!(
5376 reviewer_decision,
5377 RelayDecision::Dispatch {
5378 slot: 1,
5379 can_stop: true,
5380 ..
5381 }
5382 ));
5383 assert_eq!(
5384 relay.run_turn("", 0).await.expect("complete"),
5385 RelayDecision::Complete
5386 );
5387 }
5388
5389 #[tokio::test]
5390 async fn stop_token_is_filtered_from_streamed_ui_events() {
5391 let first = ScriptedAdapter::new(
5392 0,
5393 AgentCapabilities::default(),
5394 [
5395 AgentEvent::Text {
5396 slot: 0,
5397 text: format!("visible {STOP_TOKEN} trailing"),
5398 },
5399 AgentEvent::TurnComplete { slot: 0 },
5400 ],
5401 );
5402 let reviewer = ScriptedAdapter::new(
5403 1,
5404 AgentCapabilities::default(),
5405 [AgentEvent::TurnComplete { slot: 1 }],
5406 );
5407 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5408 let captured = std::sync::Arc::clone(&events);
5409 let mut relay = RelayHost::new(
5410 vec![
5411 AdapterHost::new(Box::new(first), None),
5412 AdapterHost::new(Box::new(reviewer), None),
5413 ],
5414 2,
5415 )
5416 .expect("relay");
5417 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5418 relay.start().await.expect("start");
5419 relay.run_turn("task", 0).await.expect("turn");
5420 let captured = events.lock().expect("lock");
5421 assert!(captured.iter().all(|event| match event {
5422 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5423 _ => true,
5424 }));
5425 let visible = captured
5426 .iter()
5427 .filter_map(|event| match event {
5428 AgentEvent::Text { text, .. } => Some(text.as_str()),
5429 _ => None,
5430 })
5431 .collect::<String>();
5432 assert_eq!(visible, "visible trailing");
5433 }
5434
5435 #[tokio::test]
5436 async fn token_only_reviewer_response_emits_visible_acknowledgment() {
5437 let first = ScriptedAdapter::new(
5438 0,
5439 AgentCapabilities::default(),
5440 [
5441 AgentEvent::Text {
5442 slot: 0,
5443 text: "done".into(),
5444 },
5445 AgentEvent::TurnComplete { slot: 0 },
5446 ],
5447 );
5448 let reviewer = ScriptedAdapter::new(
5449 1,
5450 AgentCapabilities::default(),
5451 [
5452 AgentEvent::Text {
5453 slot: 1,
5454 text: STOP_TOKEN.into(),
5455 },
5456 AgentEvent::TurnComplete { slot: 1 },
5457 ],
5458 );
5459 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5460 let captured = std::sync::Arc::clone(&events);
5461 let mut relay = RelayHost::new(
5462 vec![
5463 AdapterHost::new(Box::new(first), None),
5464 AdapterHost::new(Box::new(reviewer), None),
5465 ],
5466 4,
5467 )
5468 .expect("relay");
5469 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5470 relay.start().await.expect("start");
5471 relay.run_turn("task", 0).await.expect("first turn");
5472 relay.run_turn("", 0).await.expect("review turn");
5473 let captured = events.lock().expect("lock");
5474 assert!(captured.iter().any(|event| {
5475 matches!(
5476 event,
5477 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
5478 )
5479 }));
5480 assert!(captured.iter().all(|event| match event {
5481 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5482 _ => true,
5483 }));
5484 let acknowledgment = captured
5485 .iter()
5486 .position(|event| {
5487 matches!(
5488 event,
5489 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
5490 )
5491 })
5492 .expect("visible acknowledgment");
5493 let completion = captured
5494 .iter()
5495 .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
5496 .expect("reviewer completion");
5497 assert!(acknowledgment < completion);
5498 }
5499
5500 #[tokio::test]
5501 async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
5502 let first = ScriptedAdapter::new(
5503 0,
5504 AgentCapabilities::default(),
5505 [
5506 AgentEvent::Text {
5507 slot: 0,
5508 text: "done".into(),
5509 },
5510 AgentEvent::TurnComplete { slot: 0 },
5511 ],
5512 );
5513 let reviewer = ScriptedAdapter::new(
5514 1,
5515 AgentCapabilities::default(),
5516 [
5517 AgentEvent::Text {
5518 slot: 1,
5519 text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
5520 },
5521 AgentEvent::TurnComplete { slot: 1 },
5522 ],
5523 );
5524 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5525 let captured = std::sync::Arc::clone(&events);
5526 let mut relay = RelayHost::new(
5527 vec![
5528 AdapterHost::new(Box::new(first), None),
5529 AdapterHost::new(Box::new(reviewer), None),
5530 ],
5531 4,
5532 )
5533 .expect("relay");
5534 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5535 relay.start().await.expect("start");
5536 relay.run_turn("task", 0).await.expect("first turn");
5537 relay.run_turn("", 0).await.expect("review turn");
5538
5539 let visible = events
5540 .lock()
5541 .expect("lock")
5542 .iter()
5543 .filter_map(|event| match event {
5544 AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
5545 _ => None,
5546 })
5547 .collect::<String>();
5548 assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
5549 assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
5550 }
5551
5552 #[tokio::test]
5553 async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
5554 let first = AdapterHost::new(
5555 Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
5556 None,
5557 );
5558 let second = AdapterHost::new(
5559 Box::new(ScriptedAdapter::new(
5560 1,
5561 AgentCapabilities::default(),
5562 [AgentEvent::TurnComplete { slot: 1 }],
5563 )),
5564 None,
5565 );
5566 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
5567 let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
5568 relay.set_event_sink(move |event| {
5569 if matches!(event, AgentEvent::Permission { .. }) {
5570 let _ = seen_sender.send(());
5571 }
5572 });
5573 relay.start().await.expect("start");
5574 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
5575 let answer = async move {
5576 seen_receiver.recv().await.expect("permission request");
5577 sender
5578 .send(super::RelayPermissionAnswer {
5579 slot: 0,
5580 request_id: "permission-1".into(),
5581 answer: PermissionAnswer::Selected {
5582 option_id: "allow".into(),
5583 },
5584 })
5585 .expect("queue permission answer");
5586 };
5587 tokio::time::timeout(std::time::Duration::from_millis(100), async {
5588 let ((), result) = tokio::join!(
5589 answer,
5590 relay.run_turn_with_permissions("task", 0, &mut receiver)
5591 );
5592 result
5593 })
5594 .await
5595 .expect("permission-gated turn should not deadlock")
5596 .expect("turn completes");
5597 }
5598
5599 #[tokio::test]
5600 async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
5601 let first = AdapterHost::new(
5602 Box::new(PendingAdapter {
5603 slot: 0,
5604 hang_on_cancel: false,
5605 }),
5606 None,
5607 );
5608 let second = AdapterHost::new(
5609 Box::new(ScriptedAdapter::new(
5610 1,
5611 AgentCapabilities::default(),
5612 [AgentEvent::TurnComplete { slot: 1 }],
5613 )),
5614 None,
5615 );
5616 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
5617 relay.start().await.expect("start");
5618 let cancellation = relay.cancellation();
5619 let error = {
5620 let turn = relay.run_turn("task", 0);
5621 tokio::pin!(turn);
5622 cancellation.request();
5623 turn.await.expect_err("cancellation should stop turn")
5624 };
5625 assert!(error.to_string().contains("relay turn cancelled"));
5626
5627 assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
5628 relay
5629 .run_turn("", 1)
5630 .await
5631 .expect("replacement job reaches the selected peer");
5632 let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
5633 assert!(replacement.contains("replacement job"));
5634 assert!(replacement.contains("User "));
5635 assert!(replacement.contains(":\ntask"));
5636 let owner_updates = relay.relay_mut().unseen_context(0);
5637 assert!(owner_updates.contains("User "));
5638 assert!(owner_updates.contains(":\ntask"));
5639 assert!(owner_updates.contains(":\nreplacement job"));
5640 }
5641
5642 #[tokio::test]
5643 async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
5644 let first = AdapterHost::new(
5645 Box::new(PendingAdapter {
5646 slot: 0,
5647 hang_on_cancel: true,
5648 }),
5649 None,
5650 );
5651 let second = AdapterHost::new(
5652 Box::new(PendingAdapter {
5653 slot: 1,
5654 hang_on_cancel: false,
5655 }),
5656 None,
5657 );
5658 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
5659 relay.start().await.expect("start");
5660 let cancellation = relay.cancellation();
5661 let turn = relay.run_turn("task", 0);
5662 tokio::pin!(turn);
5663 cancellation.request();
5664 let error = turn.await.expect_err("cancellation should stop turn");
5665 assert!(error.to_string().contains("timed out"));
5666 }
5667
5668 #[tokio::test]
5669 async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
5670 let event = [AgentEvent::TurnComplete { slot: 0 }];
5671 let first = AdapterHost::new(
5672 Box::new(ScriptedAdapter::new(
5673 0,
5674 AgentCapabilities::default(),
5675 event.clone(),
5676 )),
5677 None,
5678 );
5679 let second = AdapterHost::new(
5680 Box::new(ScriptedAdapter::new(
5681 1,
5682 AgentCapabilities::default(),
5683 [AgentEvent::TurnComplete { slot: 1 }],
5684 )),
5685 None,
5686 );
5687 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
5688 relay.start().await.expect("start");
5689
5690 relay.pause();
5691 assert_eq!(
5692 relay.run_turn("paused", 0).await.expect("paused turn"),
5693 crate::relay::RelayDecision::Paused
5694 );
5695 assert!(relay.dispatches().is_empty());
5696
5697 relay.resume();
5698 relay.relay_mut().drop_agent(1).expect("drop reviewer");
5699 assert!(matches!(
5700 relay
5701 .run_turn("solo follow-up", 0)
5702 .await
5703 .expect("solo turn"),
5704 crate::relay::RelayDecision::Dispatch {
5705 slot: 0,
5706 can_stop: false,
5707 ..
5708 }
5709 ));
5710 assert_eq!(relay.dispatches().len(), 1);
5711 }
5712
5713 #[tokio::test]
5714 async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
5715 let first = AdapterHost::new(
5716 Box::new(ScriptedAdapter::new(
5717 0,
5718 AgentCapabilities::default(),
5719 [AgentEvent::TurnComplete { slot: 0 }],
5720 )),
5721 None,
5722 );
5723 let second = AdapterHost::new(
5724 Box::new(ScriptedAdapter::new(
5725 1,
5726 AgentCapabilities::default(),
5727 [AgentEvent::TurnComplete { slot: 1 }],
5728 )),
5729 None,
5730 );
5731 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
5732 relay.set_roster_names(vec!["First".into(), "Second".into()]);
5733 relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
5734 relay.set_roster_launch_specs(vec![
5735 ("custom".into(), "owner".into()),
5736 ("custom".into(), "peer".into()),
5737 ]);
5738 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5739 let captured = std::sync::Arc::clone(&events);
5740 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5741 relay.start().await.expect("start");
5742 let slot = relay
5743 .add_agent(
5744 AdapterHost::new(
5745 Box::new(ScriptedAdapter::new(
5746 2,
5747 AgentCapabilities::default(),
5748 [AgentEvent::TurnComplete { slot: 2 }],
5749 )),
5750 None,
5751 ),
5752 "Reviewer",
5753 "reviewer.example",
5754 "reviewer --acp",
5755 )
5756 .await
5757 .expect("append agent");
5758 assert_eq!(slot, 2);
5759 assert_eq!(
5760 relay.relay().active_slots().collect::<Vec<_>>(),
5761 vec![0, 1, 2]
5762 );
5763 assert_eq!(
5764 relay
5765 .session_metadata()
5766 .get("agents")
5767 .and_then(|value| value.as_array())
5768 .map(Vec::len),
5769 Some(3)
5770 );
5771 relay.drop_agent(1).await.expect("drop middle peer");
5772 let metadata = relay.session_metadata();
5773 assert_eq!(
5774 metadata.get("agents"),
5775 Some(&serde_json::json!([
5776 {"name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
5777 {"name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
5778 ]))
5779 );
5780 assert!(
5781 events
5782 .lock()
5783 .expect("lock")
5784 .iter()
5785 .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
5786 );
5787 }
5788
5789 #[tokio::test]
5790 async fn relay_host_persists_coordinator_owned_runtime_metadata() {
5791 let path = unique_test_path("codeswarm-session-metadata", "json");
5792 let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
5793 let writer = metadata_store.buffered().expect("metadata writer");
5794 let first = AdapterHost::new(
5795 Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
5796 None,
5797 );
5798 let second = AdapterHost::new(
5799 Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
5800 None,
5801 );
5802 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
5803 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
5804 relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
5805 relay.set_roster_launch_specs(vec![
5806 ("custom".into(), "claude".into()),
5807 ("custom".into(), "codex".into()),
5808 ]);
5809 relay.set_session_metadata_writer(writer);
5810 relay.start().await.expect("start");
5811 relay.drop_agent(0).await.expect("drop first agent");
5812 relay.stop().await.expect("stop");
5813
5814 let loaded = metadata_store
5815 .read()
5816 .expect("read metadata")
5817 .expect("metadata snapshot");
5818 assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
5819 assert_eq!(
5820 loaded.get("agents"),
5821 Some(&serde_json::json!([{
5822 "name": "Codex", "identity": "openai.com", "protocol": "custom",
5823 "command": "codex", "supports_load_session": false
5824 }]))
5825 );
5826 assert!(loaded.get("owner").is_none());
5827 let _ = std::fs::remove_file(path);
5828 }
5829
5830 #[tokio::test]
5831 async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
5832 let first = AdapterHost::new(
5833 Box::new(ScriptedAdapter::new(
5834 0,
5835 AgentCapabilities::default(),
5836 [
5837 AgentEvent::Text {
5838 slot: 0,
5839 text: "owner stream".into(),
5840 },
5841 AgentEvent::TurnComplete { slot: 0 },
5842 ],
5843 )),
5844 None,
5845 );
5846 let second = AdapterHost::new(
5847 Box::new(ScriptedAdapter::new(
5848 1,
5849 AgentCapabilities::default(),
5850 [
5851 AgentEvent::Text {
5852 slot: 1,
5853 text: "peer stream".into(),
5854 },
5855 AgentEvent::TurnComplete { slot: 1 },
5856 ],
5857 )),
5858 None,
5859 );
5860 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
5861 relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
5862 relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
5863 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5864 let captured = std::sync::Arc::clone(&events);
5865 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5866 relay.start().await.expect("start");
5867
5868 relay.swap_agents(0, 1).expect("swap peers");
5869 assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
5870 assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
5871 relay.run_turn("task", 0).await.expect("swapped turn");
5872 let events = events.lock().expect("events");
5873 assert!(events.iter().any(|event| {
5874 matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
5875 }));
5876 assert!(relay.dispatches()[0].1.contains("You are Peer"));
5877 }
5878
5879 #[tokio::test]
5880 async fn relay_host_persists_all_active_agent_metadata_off_thread() {
5881 let path = unique_test_path("codeswarm-session-metadata", "json");
5882 let first = AdapterHost::new(
5883 Box::new(ScriptedAdapter::new(
5884 0,
5885 AgentCapabilities::default(),
5886 [AgentEvent::TurnComplete { slot: 0 }],
5887 )),
5888 None,
5889 );
5890 let second = AdapterHost::new(
5891 Box::new(ScriptedAdapter::new(
5892 1,
5893 AgentCapabilities::default(),
5894 [AgentEvent::TurnComplete { slot: 1 }],
5895 )),
5896 None,
5897 );
5898 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
5899 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
5900 relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
5901 relay.set_roster_launch_specs(vec![
5902 ("custom".into(), "claude".into()),
5903 ("custom".into(), "codex".into()),
5904 ]);
5905 let writer = SessionMetadataStore::open(&path)
5906 .buffered()
5907 .expect("metadata writer");
5908 relay.set_session_metadata_writer(writer);
5909 relay.start().await.expect("start");
5910 relay.stop().await.expect("stop");
5911 let loaded = SessionMetadataStore::open(&path)
5912 .read()
5913 .expect("read metadata")
5914 .expect("metadata snapshot");
5915 let agents = loaded
5916 .get("agents")
5917 .and_then(|value| value.as_array())
5918 .expect("agents");
5919 assert_eq!(agents.len(), 2);
5920 assert_eq!(agents[0]["identity"], "claude.com");
5921 assert_eq!(agents[1]["identity"], "openai.com");
5922 let _ = std::fs::remove_file(path);
5923 }
5924
5925 #[tokio::test]
5926 async fn relay_host_routes_unseen_public_context_to_next_agent() {
5927 let first = AdapterHost::new(
5928 Box::new(ScriptedAdapter::new(
5929 0,
5930 AgentCapabilities::default(),
5931 [
5932 AgentEvent::Text {
5933 slot: 0,
5934 text: "implemented the fix".into(),
5935 },
5936 AgentEvent::TurnComplete { slot: 0 },
5937 ],
5938 )),
5939 None,
5940 );
5941 let second = AdapterHost::new(
5942 Box::new(ScriptedAdapter::new(
5943 1,
5944 AgentCapabilities::default(),
5945 [AgentEvent::TurnComplete { slot: 1 }],
5946 )),
5947 None,
5948 );
5949 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
5950 relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
5951 relay.start().await.expect("start");
5952 relay.run_turn("task", 0).await.expect("first turn");
5953 relay.run_turn("review this", 0).await.expect("review turn");
5954
5955 assert_eq!(relay.dispatches().len(), 2);
5956 assert_eq!(relay.dispatches()[0].0, 0);
5957 assert!(relay.dispatches()[0].1.contains("task"));
5958 assert!(relay.dispatches()[0].1.contains("You are Codex"));
5959 assert!(relay.dispatches()[0].1.contains("2. Qwen"));
5960 assert_eq!(relay.dispatches()[1].0, 1);
5961 assert!(relay.dispatches()[1].1.contains("review this"));
5962 let public = relay.dispatches()[1]
5963 .1
5964 .split_once("Public updates:\n")
5965 .map(|(_, updates)| updates)
5966 .expect("review receives public context");
5967 let header = public
5968 .lines()
5969 .find(|line| line.starts_with("Codex "))
5970 .expect("named previous agent");
5971 let timestamp = header
5972 .strip_prefix("Codex ")
5973 .and_then(|value| value.strip_suffix(':'))
5974 .expect("timestamped header");
5975 assert_eq!(timestamp.len(), 5);
5976 assert_eq!(timestamp.as_bytes()[2], b':');
5977 assert!(
5978 timestamp
5979 .bytes()
5980 .enumerate()
5981 .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
5982 );
5983 assert!(public.contains("implemented the fix"));
5984 assert!(!public.contains("Agent 0"));
5985 }
5986}