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