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