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("60m")
1971 .arg("--output-format")
1972 .arg("stream-json")
1973 .current_dir(&self.cwd)
1974 .env("CODESWARM_CWD", &self.cwd)
1975 .stdout(Stdio::piped())
1976 .stderr(Stdio::piped());
1977 if let Some(session_id) = &self.session_id {
1978 command.arg("--conversation").arg(session_id);
1979 }
1980 if self.mode != "default" {
1981 command.arg("--mode").arg(&self.mode);
1982 }
1983 let mut child = command
1984 .spawn()
1985 .map_err(|error| AdapterError::Spawn(error.to_string()))?;
1986 let stdout = match child.stdout.take() {
1987 Some(stdout) => stdout,
1988 None => {
1989 let _ = terminate_child(&mut child).await;
1990 return Err(AdapterError::Transport("agent has no stdout".into()));
1991 }
1992 };
1993 let stderr = match child.stderr.take() {
1994 Some(stderr) => stderr,
1995 None => {
1996 let _ = terminate_child(&mut child).await;
1997 return Err(AdapterError::Transport("agent has no stderr".into()));
1998 }
1999 };
2000 let sender = self.sender.clone();
2001 let slot = self.slot;
2002 let announced_session = Arc::clone(&self.announced_session);
2003 let cancel_requested = Arc::clone(&self.cancel_requested);
2004 tokio::spawn(async move {
2005 let stderr_task = tokio::spawn(async move {
2006 const MAX_STDERR: usize = 32 * 1024;
2007 let mut stderr = BufReader::new(stderr);
2008 let mut bytes = Vec::new();
2009 let mut chunk = [0_u8; 4096];
2010 while let Ok(count) = stderr.read(&mut chunk).await {
2011 if count == 0 {
2012 break;
2013 }
2014 bytes.extend_from_slice(&chunk[..count]);
2015 if bytes.len() > MAX_STDERR {
2016 let keep_from = bytes.len() - MAX_STDERR;
2017 bytes.drain(..keep_from);
2018 }
2019 }
2020 String::from_utf8_lossy(&bytes).trim().to_owned()
2021 });
2022 let mut lines = BufReader::new(stdout).lines();
2023 let mut result: Option<Value> = None;
2024 let mut streamed_response = false;
2025 while let Ok(Some(line)) = lines.next_line().await {
2026 let value = match serde_json::from_str::<Value>(&line) {
2027 Ok(value) => value,
2028 Err(_) => {
2029 continue;
2034 }
2035 };
2036 if value.get("event").and_then(Value::as_str) == Some("init")
2037 && let Some(session_id) = value
2038 .get("conversation_id")
2039 .or_else(|| value.get("conversationId"))
2040 .and_then(Value::as_str)
2041 .filter(|id| !id.is_empty())
2042 && let Ok(mut announced) = announced_session.lock()
2043 {
2044 *announced = Some(session_id.to_owned());
2045 }
2046 if value.get("event").and_then(Value::as_str) == Some("result") {
2047 result = value.get("result").cloned();
2048 }
2049 match parse_agy_value(slot, &value) {
2050 Ok(Some(event)) => {
2051 if matches!(event, AgentEvent::Text { .. }) {
2052 streamed_response = true;
2053 }
2054 if sender.send(Ok(event)).await.is_err() {
2055 break;
2056 }
2057 }
2058 Ok(None) => {}
2059 Err(error) => {
2060 let _ = sender.send(Err(error)).await;
2061 }
2062 }
2063 }
2064 let stderr = stderr_task.await.ok().unwrap_or_default();
2065 let succeeded = cancel_requested.load(Ordering::Acquire)
2066 || result
2067 .as_ref()
2068 .and_then(|result| result.get("status"))
2069 .and_then(Value::as_str)
2070 == Some("SUCCESS");
2071 if succeeded {
2072 if !streamed_response
2077 && let Some(response) = result
2078 .as_ref()
2079 .and_then(|result| result.get("response"))
2080 .and_then(Value::as_str)
2081 .filter(|response| !response.is_empty())
2082 {
2083 let _ = sender
2084 .send(Ok(AgentEvent::Text {
2085 slot,
2086 text: response.to_owned(),
2087 }))
2088 .await;
2089 }
2090 let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
2091 } else {
2092 let detail = result
2093 .as_ref()
2094 .and_then(|result| result.get("error"))
2095 .and_then(Value::as_str)
2096 .filter(|detail| !detail.is_empty())
2097 .map(str::to_owned)
2098 .or_else(|| (!stderr.is_empty()).then_some(stderr))
2099 .unwrap_or_else(|| "native stream ended before a successful result".into());
2100 let _ = sender
2101 .send(Ok(AgentEvent::Failed {
2102 slot,
2103 started: true,
2104 detail,
2105 }))
2106 .await;
2107 }
2108 });
2109 self.child = Some(child);
2110 Ok(())
2111 }
2112
2113 async fn cancel(&mut self) -> AdapterResult<bool> {
2114 self.cancel_requested.store(true, Ordering::Release);
2115 let Some(mut child) = self.child.take() else {
2116 return Ok(false);
2117 };
2118 terminate_child(&mut child).await?;
2122 let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
2123 while let Some(event) = self.receiver.recv().await {
2124 if matches!(
2125 event,
2126 Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
2127 ) {
2128 break;
2129 }
2130 }
2131 })
2132 .await;
2133 Ok(true)
2134 }
2135
2136 async fn answer_permission(
2137 &mut self,
2138 _request_id: String,
2139 _answer: PermissionAnswer,
2140 ) -> AdapterResult<()> {
2141 Err(AdapterError::Unsupported("permission answer"))
2142 }
2143
2144 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
2145 let (mode, mode_policy) = match mode.as_str() {
2146 "full-access" | "codeswarm:mode:full-access" | "auto" | "autopilot" => {
2147 ("default".to_owned(), "agy:full-access".to_owned())
2148 }
2149 "codeswarm:mode:plan" | "readonly" | "plan" => ("plan".to_owned(), "plan".to_owned()),
2150 "codeswarm:mode:accept-edits" | "acceptedits" | "accept-edits" => {
2151 ("accept-edits".to_owned(), "accept-edits".to_owned())
2152 }
2153 "codeswarm:mode:manual" | "manual" | "ask" | "default" => {
2154 ("default".to_owned(), "agy:manual".to_owned())
2155 }
2156 "agy:full-access" => ("default".to_owned(), "agy:full-access".to_owned()),
2157 "agy:manual" => ("default".to_owned(), "agy:manual".to_owned()),
2158 _ => return Err(AdapterError::Unsupported("requested Agy mode")),
2159 };
2160 self.mode = mode;
2161 self.mode_policy = mode_policy.clone();
2162 self.emit(Ok(AgentEvent::ModesReplaced {
2163 slot: self.slot,
2164 modes: Self::modes(),
2165 current_mode: Some(mode_policy),
2166 }))
2167 .await;
2168 Ok(())
2169 }
2170
2171 async fn reload(&mut self) -> AdapterResult<()> {
2172 self.stop().await?;
2173 self.start().await
2174 }
2175
2176 async fn stop(&mut self) -> AdapterResult<()> {
2177 let _ = self.cancel().await?;
2178 Ok(())
2179 }
2180
2181 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
2182 let event = self.receiver.recv().await;
2183 if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2184 && self.session_id.is_none()
2185 && let Ok(session) = self.announced_session.lock()
2186 {
2187 self.session_id = session.clone();
2188 }
2189 if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2190 && let Some(mut child) = self.child.take()
2191 {
2192 let _ = child.wait().await;
2193 }
2194 event
2195 }
2196}
2197
2198#[cfg(test)]
2199#[cfg_attr(not(test), allow(dead_code))]
2200fn parse_agy_line(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
2201 let value: Value =
2202 serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
2203 parse_agy_value(slot, &value)
2204}
2205
2206fn parse_agy_value(slot: RosterSlot, value: &Value) -> AdapterResult<Option<AgentEvent>> {
2207 let event = value.get("event").and_then(Value::as_str);
2208 if let Some(terminal) = parse_terminal_event(value, event) {
2209 return Ok(Some(AgentEvent::Terminal {
2210 slot,
2211 event: terminal,
2212 }));
2213 }
2214 match event {
2215 Some("step_update") => {
2216 let Some(update) = value.get("step_update") else {
2217 return Ok(None);
2218 };
2219 let is_response = update
2220 .get("step_type")
2221 .and_then(Value::as_str)
2222 .is_some_and(|kind| kind == "agent_response");
2223 let text = update.get("text_delta").and_then(Value::as_str);
2224 let response = is_response
2225 .then(|| text.map(str::to_owned))
2226 .flatten()
2227 .filter(|text| !text.is_empty())
2228 .map(|text| AgentEvent::Text { slot, text });
2229 Ok(response.or_else(|| parse_agy_tool(slot, value)))
2230 }
2231 _ => Ok(None),
2232 }
2233}
2234
2235fn parse_agy_tool(slot: RosterSlot, value: &Value) -> Option<AgentEvent> {
2236 let update = value.get("step_update")?;
2237 if update.get("step_type")?.as_str()? != "tool" {
2238 return None;
2239 }
2240 let step_index = update.get("step_index")?.as_i64()?;
2241 let title = update
2242 .get("tool_name")
2243 .and_then(Value::as_str)
2244 .unwrap_or("Tool call")
2245 .replace('_', " ");
2246 let status = match update.get("state").and_then(Value::as_str) {
2247 Some("DONE") => ToolStatus::Completed,
2248 Some("FAILED") => ToolStatus::Failed,
2249 Some("ACTIVE") => ToolStatus::Running,
2250 _ => ToolStatus::Pending,
2251 };
2252 let detail = update
2253 .get("tool_info")
2254 .and_then(|info| info.get("output"))
2255 .and_then(Value::as_str)
2256 .map(str::to_owned);
2257 Some(AgentEvent::Tool {
2258 slot,
2259 update: ToolUpdate {
2260 id: format!("agy-tool-{step_index}"),
2261 title,
2262 status,
2263 detail,
2264 },
2265 })
2266}
2267
2268#[derive(Debug)]
2271pub struct AcpAdapter {
2272 slot: RosterSlot,
2273 program: String,
2274 args: Vec<String>,
2275 cwd: PathBuf,
2276 child: Option<Child>,
2277 reader: Option<BufReader<ChildStdout>>,
2278 capabilities: AgentCapabilities,
2279 modes: Vec<Mode>,
2280 models: Vec<Mode>,
2281 model_config_id: Option<String>,
2282 session_id: Option<String>,
2283 next_request_id: u64,
2284 prompt_request_id: Option<u64>,
2285 queued_events: VecDeque<AdapterResult<AgentEvent>>,
2286 stderr_task: Option<tokio::task::JoinHandle<String>>,
2287 terminals: BTreeMap<String, TerminalProcess>,
2288 next_terminal_id: u64,
2289}
2290
2291impl AcpAdapter {
2292 pub fn new(
2293 slot: RosterSlot,
2294 cwd: PathBuf,
2295 program: impl Into<String>,
2296 args: Vec<String>,
2297 ) -> Self {
2298 Self {
2299 slot,
2300 program: program.into(),
2301 args,
2302 cwd,
2303 child: None,
2304 reader: None,
2305 capabilities: AgentCapabilities::default(),
2306 modes: Vec::new(),
2307 models: Vec::new(),
2308 model_config_id: None,
2309 session_id: None,
2310 next_request_id: 1,
2311 prompt_request_id: None,
2312 queued_events: VecDeque::new(),
2313 stderr_task: None,
2314 terminals: BTreeMap::new(),
2315 next_terminal_id: 1,
2316 }
2317 }
2318
2319 pub fn with_session_id(
2320 slot: RosterSlot,
2321 cwd: PathBuf,
2322 program: impl Into<String>,
2323 args: Vec<String>,
2324 session_id: impl Into<String>,
2325 ) -> Self {
2326 let mut adapter = Self::new(slot, cwd, program, args);
2327 adapter.session_id = Some(session_id.into());
2328 adapter
2329 }
2330
2331 async fn request(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2332 self.request_with_timeout(method, params, std::time::Duration::from_secs(30))
2333 .await
2334 }
2335
2336 async fn request_with_timeout(
2337 &mut self,
2338 method: &str,
2339 params: Value,
2340 deadline: std::time::Duration,
2341 ) -> AdapterResult<Value> {
2342 tokio::time::timeout(deadline, self.request_inner(method, params))
2343 .await
2344 .map_err(|_| {
2345 AdapterError::Transport(format!(
2346 "ACP {method} timed out; reload the agent to retry"
2347 ))
2348 })?
2349 }
2350
2351 async fn request_inner(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2352 let request_id = self.next_request_id;
2353 self.next_request_id += 1;
2354 self.write_json(serde_json::json!({
2355 "jsonrpc": "2.0",
2356 "id": request_id,
2357 "method": method,
2358 "params": params,
2359 }))
2360 .await?;
2361 loop {
2362 let line = self.read_line().await?;
2363 let value: Value = match serde_json::from_str(&line) {
2364 Ok(value) => value,
2365 Err(_) => {
2366 continue;
2370 }
2371 };
2372 if self.reject_empty_permission_request(&value).await? {
2373 continue;
2374 }
2375 if self.handle_client_request(&value).await? {
2376 continue;
2377 }
2378 if value
2379 .get("id")
2380 .is_some_and(|id| rpc_id_to_string(id) == request_id.to_string())
2381 {
2382 if let Some(error) = value.get("error") {
2383 return Err(AdapterError::Protocol(error.to_string()));
2384 }
2385 return value
2386 .get("result")
2387 .cloned()
2388 .ok_or_else(|| AdapterError::Protocol("response has no result".into()));
2389 }
2390 if let Some(event) = parse_acp_notification(self.slot, &line)? {
2391 let event = if method == "session/load" {
2392 restored_history_event(event)
2393 } else {
2394 event
2395 };
2396 self.queued_events.push_back(Ok(event));
2397 }
2398 }
2399 }
2400
2401 async fn write_json(&mut self, value: Value) -> AdapterResult<()> {
2402 let child = self
2403 .child
2404 .as_mut()
2405 .ok_or_else(|| AdapterError::Transport("ACP agent is not running".into()))?;
2406 let stdin = child
2407 .stdin
2408 .as_mut()
2409 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdin".into()))?;
2410 stdin
2411 .write_all(value.to_string().as_bytes())
2412 .await
2413 .map_err(|error| AdapterError::Transport(error.to_string()))?;
2414 stdin
2415 .write_all(b"\n")
2416 .await
2417 .map_err(|error| AdapterError::Transport(error.to_string()))
2418 }
2419
2420 async fn reject_empty_permission_request(&mut self, value: &Value) -> AdapterResult<bool> {
2425 if value.get("method").and_then(Value::as_str) != Some("session/request_permission")
2426 || value.get("id").is_none()
2427 {
2428 return Ok(false);
2429 }
2430 let valid = value
2431 .get("params")
2432 .and_then(|params| params.get("options"))
2433 .and_then(Value::as_array)
2434 .is_some_and(|options| !options.is_empty());
2435 if valid {
2436 return Ok(false);
2437 }
2438 self.write_json(serde_json::json!({
2439 "jsonrpc": "2.0",
2440 "id": value.get("id").cloned().unwrap_or(Value::Null),
2441 "error": {
2442 "code": -32602,
2443 "message": "Permission request requires at least one option",
2444 },
2445 }))
2446 .await?;
2447 Ok(true)
2448 }
2449
2450 fn workspace_path(&self, path: &str) -> Result<PathBuf, String> {
2451 let root = self
2452 .cwd
2453 .canonicalize()
2454 .map_err(|error| format!("unable to resolve workspace: {error}"))?;
2455 let requested = Path::new(path);
2456 let candidate = if requested.is_absolute() {
2457 requested.to_path_buf()
2458 } else {
2459 root.join(requested)
2460 };
2461 let resolved = if !candidate.exists() {
2462 let parent = candidate
2463 .parent()
2464 .ok_or_else(|| "file path has no parent".to_owned())?
2465 .canonicalize()
2466 .map_err(|error| format!("unable to resolve parent directory: {error}"))?;
2467 parent.join(
2468 candidate
2469 .file_name()
2470 .ok_or_else(|| "file path has no filename".to_owned())?,
2471 )
2472 } else {
2473 candidate
2474 .canonicalize()
2475 .map_err(|error| format!("unable to resolve file path: {error}"))?
2476 };
2477 if !resolved.starts_with(&root) {
2478 return Err("file path is outside the project".into());
2479 }
2480 Ok(resolved)
2481 }
2482
2483 fn read_workspace_text(
2484 &self,
2485 path: &str,
2486 line: Option<i64>,
2487 limit: Option<i64>,
2488 ) -> Result<String, String> {
2489 if line.is_some_and(|line| line < 1) {
2490 return Err("line must be positive".into());
2491 }
2492 if limit.is_some_and(|limit| limit < 0) {
2493 return Err("limit must not be negative".into());
2494 }
2495 let path = self.workspace_path(path)?;
2496 let mut bytes = Vec::new();
2497 let mut source = match std::fs::File::open(path) {
2498 Ok(source) => source,
2499 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(String::new()),
2500 Err(error) => return Err(error.to_string()),
2501 };
2502 source
2503 .by_ref()
2504 .take((MAX_FILE_READ_BYTES as u64).saturating_add(1))
2505 .read_to_end(&mut bytes)
2506 .map_err(|error| error.to_string())?;
2507 bytes.truncate(MAX_FILE_READ_BYTES);
2508 let text = String::from_utf8_lossy(&bytes);
2509 if line.is_none() && limit.is_none() {
2510 return Ok(text.into_owned());
2511 }
2512 let start = line.map_or(0, |line| line as usize - 1);
2513 let limit = limit.unwrap_or(i64::MAX) as usize;
2514 let selected = text
2515 .split_inclusive('\n')
2516 .skip(start)
2517 .take(limit)
2518 .collect::<String>();
2519 if line.is_some() {
2520 Ok(selected.trim_end_matches('\n').to_owned())
2521 } else {
2522 Ok(selected)
2523 }
2524 }
2525
2526 fn write_workspace_text(&self, params: &Value) -> Result<(), String> {
2527 let path = params
2528 .get("path")
2529 .and_then(Value::as_str)
2530 .filter(|path| !path.is_empty())
2531 .ok_or("path must be a non-empty string")?;
2532 let content = params
2533 .get("content")
2534 .and_then(Value::as_str)
2535 .ok_or("content must be a string")?;
2536 let path = self.workspace_path(path)?;
2537 std::fs::write(path, content).map_err(|error| error.to_string())
2538 }
2539
2540 async fn terminal_create(&mut self, params: &Value) -> Result<Value, String> {
2541 let command = params
2542 .get("command")
2543 .and_then(Value::as_str)
2544 .filter(|command| !command.trim().is_empty())
2545 .ok_or_else(|| "terminal command is required".to_owned())?;
2546 let cwd = params.get("cwd").and_then(Value::as_str).unwrap_or(".");
2547 let cwd = self.workspace_path(cwd)?;
2548 if !cwd.is_dir() {
2549 return Err("terminal cwd is not a directory".into());
2550 }
2551 let mut process = Command::new(command);
2552 isolate_process_group(&mut process);
2553 if let Some(args) = params.get("args").and_then(Value::as_array) {
2554 process.args(args.iter().filter_map(Value::as_str));
2555 }
2556 process
2557 .current_dir(&cwd)
2558 .stdin(Stdio::null())
2559 .stdout(Stdio::piped())
2560 .stderr(Stdio::piped());
2561 if let Some(env) = params.get("env") {
2562 if let Some(entries) = env.as_array() {
2563 for entry in entries {
2564 if let (Some(name), Some(value)) = (
2565 entry.get("name").and_then(Value::as_str),
2566 entry.get("value").and_then(Value::as_str),
2567 ) {
2568 process.env(name, value);
2569 }
2570 }
2571 } else if let Some(entries) = env.as_object() {
2572 for (name, value) in entries {
2573 if let Some(value) = value.as_str() {
2574 process.env(name, value);
2575 }
2576 }
2577 }
2578 }
2579 let mut child = process.spawn().map_err(|error| error.to_string())?;
2580 let stdout = child.stdout.take();
2581 let stderr = child.stderr.take();
2582 let (Some(stdout), Some(stderr)) = (stdout, stderr) else {
2583 let _ = terminate_child(&mut child).await;
2584 return Err("terminal has no output pipes".into());
2585 };
2586 let output = Arc::new(Mutex::new(Vec::new()));
2587 let truncated = Arc::new(AtomicBool::new(false));
2588 let output_readers = Arc::new(AtomicUsize::new(2));
2589 let output_limit = params
2590 .get("outputByteLimit")
2591 .and_then(Value::as_u64)
2592 .map_or(MAX_TERMINAL_OUTPUT_BYTES, |limit| {
2593 usize::try_from(limit)
2594 .unwrap_or(MAX_TERMINAL_OUTPUT_BYTES)
2595 .min(MAX_TERMINAL_OUTPUT_BYTES)
2596 });
2597 tokio::spawn(drain_terminal_output(
2598 stdout,
2599 Arc::clone(&output),
2600 Arc::clone(&truncated),
2601 Arc::clone(&output_readers),
2602 output_limit,
2603 ));
2604 tokio::spawn(drain_terminal_output(
2605 stderr,
2606 Arc::clone(&output),
2607 Arc::clone(&truncated),
2608 Arc::clone(&output_readers),
2609 output_limit,
2610 ));
2611 let id = format!("terminal-{}", self.next_terminal_id);
2612 self.next_terminal_id = self.next_terminal_id.saturating_add(1);
2613 let state = TerminalProcess {
2614 child: Arc::new(AsyncMutex::new(Some(child))),
2615 output,
2616 truncated,
2617 output_readers,
2618 };
2619 self.terminals.insert(id.clone(), state);
2620 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2621 slot: self.slot,
2622 event: TerminalEvent::Created {
2623 id: id.clone(),
2624 command: std::iter::once(command)
2625 .chain(
2626 params
2627 .get("args")
2628 .and_then(Value::as_array)
2629 .into_iter()
2630 .flatten()
2631 .filter_map(Value::as_str),
2632 )
2633 .collect::<Vec<_>>()
2634 .join(" "),
2635 },
2636 }));
2637 Ok(serde_json::json!({"terminalId": id}))
2638 }
2639
2640 async fn terminal_output(&mut self, id: &str) -> Result<Value, String> {
2641 let terminal = self
2642 .terminals
2643 .get(id)
2644 .ok_or_else(|| "terminal not found".to_owned())?
2645 .clone();
2646 let output = terminal
2647 .output
2648 .lock()
2649 .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
2650 .unwrap_or_default();
2651 let exit_code = terminal.exit_code().await;
2652 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2653 slot: self.slot,
2654 event: TerminalEvent::Output {
2655 id: id.to_owned(),
2656 text: output.clone(),
2657 },
2658 }));
2659 let mut response = serde_json::json!({
2660 "output": output,
2661 "truncated": terminal.truncated.load(Ordering::Acquire),
2662 });
2663 if let Some(code) = exit_code {
2664 response["exitStatus"] = serde_json::json!({"exitCode": code});
2665 }
2666 Ok(response)
2667 }
2668
2669 async fn terminal_wait(&mut self, id: &str) -> Result<Value, String> {
2670 let terminal = self
2671 .terminals
2672 .get(id)
2673 .ok_or_else(|| "terminal not found".to_owned())?
2674 .clone();
2675 let exit_code = terminal.wait().await;
2676 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2677 slot: self.slot,
2678 event: TerminalEvent::Exited {
2679 id: id.to_owned(),
2680 code: exit_code.unwrap_or(-1),
2681 },
2682 }));
2683 Ok(serde_json::json!({"exitCode": exit_code, "signal": Value::Null}))
2684 }
2685
2686 async fn handle_client_request(&mut self, value: &Value) -> AdapterResult<bool> {
2690 let Some(method) = value.get("method").and_then(Value::as_str) else {
2691 return Ok(false);
2692 };
2693 let Some(id) = value.get("id").cloned() else {
2694 return Ok(false);
2695 };
2696 if method == "session/request_permission" {
2699 return Ok(false);
2700 }
2701 let params = value.get("params").cloned().unwrap_or(Value::Null);
2702 let response = match method {
2703 "fs/read_text_file" => {
2704 let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2705 let line = params.get("line").and_then(Value::as_i64);
2706 let limit = params.get("limit").and_then(Value::as_i64);
2707 match self.read_workspace_text(path, line, limit) {
2708 Ok(content) => serde_json::json!({
2709 "jsonrpc": "2.0",
2710 "id": id,
2711 "result": {"content": content},
2712 }),
2713 Err(message) => serde_json::json!({
2714 "jsonrpc": "2.0",
2715 "id": id,
2716 "error": {"code": -32602, "message": message},
2717 }),
2718 }
2719 }
2720 "fs/write_text_file" => {
2721 let result = self.write_workspace_text(¶ms);
2722 match result {
2723 Ok(()) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}}),
2724 Err(message) => serde_json::json!({
2725 "jsonrpc": "2.0",
2726 "id": id,
2727 "error": {"code": -32602, "message": message},
2728 }),
2729 }
2730 }
2731 "terminal/create" => match self.terminal_create(¶ms).await {
2732 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2733 Err(message) => serde_json::json!({
2734 "jsonrpc": "2.0",
2735 "id": id,
2736 "error": {"code": -32602, "message": message},
2737 }),
2738 },
2739 "terminal/output" => {
2740 let terminal_id = params
2741 .get("terminalId")
2742 .and_then(Value::as_str)
2743 .unwrap_or("");
2744 match self.terminal_output(terminal_id).await {
2745 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2746 Err(message) => serde_json::json!({
2747 "jsonrpc": "2.0",
2748 "id": id,
2749 "error": {"code": -32602, "message": message},
2750 }),
2751 }
2752 }
2753 "terminal/wait_for_exit" => {
2754 let terminal_id = params
2755 .get("terminalId")
2756 .and_then(Value::as_str)
2757 .unwrap_or("");
2758 match self.terminal_wait(terminal_id).await {
2759 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2760 Err(message) => serde_json::json!({
2761 "jsonrpc": "2.0",
2762 "id": id,
2763 "error": {"code": -32602, "message": message},
2764 }),
2765 }
2766 }
2767 "terminal/kill" => {
2768 let terminal_id = params
2769 .get("terminalId")
2770 .and_then(Value::as_str)
2771 .unwrap_or("");
2772 if let Some(terminal) = self.terminals.get(terminal_id) {
2773 terminal.kill().await;
2774 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2775 } else {
2776 serde_json::json!({
2777 "jsonrpc": "2.0",
2778 "id": id,
2779 "error": {"code": -32602, "message": "terminal not found"},
2780 })
2781 }
2782 }
2783 "terminal/release" => {
2784 let terminal_id = params
2785 .get("terminalId")
2786 .and_then(Value::as_str)
2787 .unwrap_or("");
2788 if let Some(terminal) = self.terminals.remove(terminal_id) {
2789 terminal.stop().await;
2790 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2791 slot: self.slot,
2792 event: TerminalEvent::Released {
2793 id: terminal_id.to_owned(),
2794 },
2795 }));
2796 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2797 } else {
2798 serde_json::json!({
2799 "jsonrpc": "2.0",
2800 "id": id,
2801 "error": {"code": -32602, "message": "terminal not found"},
2802 })
2803 }
2804 }
2805 _ => serde_json::json!({
2806 "jsonrpc": "2.0",
2807 "id": id,
2808 "error": {"code": -32601, "message": format!("unsupported client method: {method}")},
2809 }),
2810 };
2811 self.write_json(response).await?;
2812 Ok(true)
2813 }
2814
2815 async fn read_line(&mut self) -> AdapterResult<String> {
2816 let reader = self
2817 .reader
2818 .as_mut()
2819 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdout".into()))?;
2820 read_bounded_line(reader).await
2821 }
2822
2823 async fn start(&mut self) -> AdapterResult<()> {
2824 if self.child.is_some() {
2826 self.stop().await?;
2827 }
2828 let mut command = Command::new(&self.program);
2829 isolate_process_group(&mut command);
2830 command
2831 .args(&self.args)
2832 .current_dir(&self.cwd)
2833 .stdin(Stdio::piped())
2834 .stdout(Stdio::piped())
2835 .stderr(Stdio::piped())
2836 .env("CODESWARM_CWD", &self.cwd);
2837 if self.program.to_ascii_lowercase().contains("gemini")
2838 || self
2839 .args
2840 .iter()
2841 .any(|arg| arg.to_ascii_lowercase().contains("gemini"))
2842 {
2843 command.env("GEMINI_TELEMETRY_ENABLED", "false");
2844 }
2845 let mut child = command
2846 .spawn()
2847 .map_err(|error| AdapterError::Spawn(error.to_string()))?;
2848 let stdout = match child.stdout.take() {
2849 Some(stdout) => stdout,
2850 None => {
2851 let _ = terminate_child(&mut child).await;
2852 return Err(AdapterError::Transport("ACP agent has no stdout".into()));
2853 }
2854 };
2855 let stderr = match child.stderr.take() {
2856 Some(stderr) => stderr,
2857 None => {
2858 let _ = terminate_child(&mut child).await;
2859 return Err(AdapterError::Transport("ACP agent has no stderr".into()));
2860 }
2861 };
2862 self.child = Some(child);
2863 self.reader = Some(BufReader::new(stdout));
2864 self.stderr_task = Some(tokio::spawn(drain_bounded(stderr, 32 * 1024)));
2865
2866 let initialize = match self
2867 .request(
2868 "initialize",
2869 serde_json::json!({
2870 "protocolVersion": 1,
2871 "clientCapabilities": {
2872 "fs": {"readTextFile": true, "writeTextFile": true},
2873 "terminal": true,
2874 },
2875 "clientInfo": {
2876 "name": "CodeSwarm",
2877 "title": "CodeSwarm",
2878 "version": env!("CARGO_PKG_VERSION"),
2879 },
2880 }),
2881 )
2882 .await
2883 {
2884 Ok(value) => value,
2885 Err(error) => {
2886 let _ = self.stop().await;
2887 return Err(error);
2888 }
2889 };
2890 let agent_capabilities = initialize
2891 .get("agentCapabilities")
2892 .cloned()
2893 .unwrap_or(Value::Null);
2894 self.capabilities = AgentCapabilities {
2895 supports_cancel: true,
2896 supports_modes: true,
2897 supports_permissions: true,
2898 supports_terminals: true,
2899 supports_session_load: agent_capabilities
2900 .get("loadSession")
2901 .and_then(Value::as_bool)
2902 .unwrap_or(false),
2903 supports_models: false,
2904 };
2905 let session = if let Some(session_id) = self.session_id.clone() {
2906 if !self.capabilities.supports_session_load {
2907 let _ = self.stop().await;
2908 return Err(AdapterError::Unsupported("session/load"));
2909 }
2910 match self
2911 .request(
2912 "session/load",
2913 serde_json::json!({
2914 "cwd": self.cwd,
2915 "mcpServers": [],
2916 "sessionId": session_id,
2917 }),
2918 )
2919 .await
2920 {
2921 Ok(value) => value,
2922 Err(error) => {
2923 let _ = self.stop().await;
2924 return Err(error);
2925 }
2926 }
2927 } else {
2928 let session = match self
2929 .request(
2930 "session/new",
2931 serde_json::json!({"cwd": self.cwd, "mcpServers": []}),
2932 )
2933 .await
2934 {
2935 Ok(value) => value,
2936 Err(error) => {
2937 let _ = self.stop().await;
2938 return Err(error);
2939 }
2940 };
2941 self.session_id = session
2942 .get("sessionId")
2943 .and_then(Value::as_str)
2944 .map(str::to_owned);
2945 if self.session_id.is_none() {
2946 let _ = self.stop().await;
2947 return Err(AdapterError::Protocol(
2948 "session/new returned no sessionId".into(),
2949 ));
2950 }
2951 session
2952 };
2953 self.capabilities.supports_modes = false;
2954 if let Some(modes) = session.get("modes") {
2955 let available = modes
2956 .get("availableModes")
2957 .and_then(Value::as_array)
2958 .map(|modes| {
2959 modes
2960 .iter()
2961 .filter_map(|mode| {
2962 Some(Mode {
2963 id: mode.get("id")?.as_str()?.to_owned(),
2964 label: mode.get("name")?.as_str()?.to_owned(),
2965 })
2966 })
2967 .collect::<Vec<_>>()
2968 })
2969 .unwrap_or_default();
2970 self.modes = available.clone();
2971 self.capabilities.supports_modes = !available.is_empty();
2972 self.queued_events.push_back(Ok(AgentEvent::ModesReplaced {
2973 slot: self.slot,
2974 modes: available,
2975 current_mode: modes
2976 .get("currentModeId")
2977 .and_then(Value::as_str)
2978 .map(str::to_owned),
2979 }));
2980 }
2981 self.models.clear();
2982 self.model_config_id = None;
2983 let current_model =
2984 parse_model_config(&session).and_then(|(config_id, models, current)| {
2985 self.model_config_id = Some(config_id);
2986 self.models = models;
2987 current
2988 });
2989 self.capabilities.supports_models =
2990 self.model_config_id.is_some() && !self.models.is_empty();
2991 if let Some(config_id) = self.model_config_id.clone()
2992 && !self.models.is_empty()
2993 {
2994 self.queued_events.push_back(Ok(AgentEvent::ModelsReplaced {
2995 slot: self.slot,
2996 config_id,
2997 models: self.models.clone(),
2998 current_model,
2999 }));
3000 }
3001 self.queued_events.push_back(Ok(AgentEvent::Ready {
3002 slot: self.slot,
3003 capabilities: self.capabilities(),
3004 }));
3005 Ok(())
3006 }
3007}
3008
3009fn prompt_resource_paths(prompt: &str) -> Vec<String> {
3010 let characters = prompt.chars().collect::<Vec<_>>();
3011 let mut paths = Vec::new();
3012 let mut index = 0;
3013 while index < characters.len() {
3014 if characters[index] != '@' {
3015 index += 1;
3016 continue;
3017 }
3018 index += 1;
3019 let quoted = characters.get(index) == Some(&'"');
3020 if quoted {
3021 index += 1;
3022 }
3023 let start = index;
3024 while index < characters.len()
3025 && if quoted {
3026 characters[index] != '"'
3027 } else {
3028 !characters[index].is_whitespace()
3029 }
3030 {
3031 index += 1;
3032 }
3033 if index > start {
3034 paths.push(characters[start..index].iter().collect());
3035 }
3036 if quoted && index < characters.len() {
3037 index += 1;
3038 }
3039 }
3040 paths
3041}
3042
3043fn prompt_content_blocks(cwd: &Path, prompt: &str) -> Vec<Value> {
3044 let mut blocks = vec![serde_json::json!({"type": "text", "text": prompt})];
3045 for path in prompt_resource_paths(prompt) {
3046 if path.ends_with('/') {
3047 continue;
3048 }
3049 let Ok(resource) = resources::load(cwd, &path) else {
3050 continue;
3051 };
3052 let uri = format!("file://{}", resource.path.display());
3053 let resource_value = if let Some(text) = resource.text {
3054 serde_json::json!({
3055 "uri": uri,
3056 "text": text,
3057 "mimeType": resource.mime_type,
3058 })
3059 } else if let Some(data) = resource.data {
3060 serde_json::json!({
3061 "uri": uri,
3062 "blob": BASE64.encode(data),
3063 "mimeType": resource.mime_type,
3064 })
3065 } else {
3066 continue;
3067 };
3068 blocks.push(serde_json::json!({
3069 "type": "resource",
3070 "resource": resource_value,
3071 }));
3072 }
3073 blocks
3074}
3075
3076#[async_trait]
3077impl AgentAdapter for AcpAdapter {
3078 fn slot(&self) -> RosterSlot {
3079 self.slot
3080 }
3081
3082 fn session_id(&self) -> Option<String> {
3083 self.session_id.clone()
3084 }
3085
3086 fn protocol(&self) -> &'static str {
3087 "acp"
3088 }
3089
3090 fn capabilities(&self) -> AgentCapabilities {
3091 self.capabilities.clone()
3092 }
3093
3094 async fn start(&mut self) -> AdapterResult<()> {
3095 AcpAdapter::start(self).await
3098 }
3099
3100 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
3101 let session_id = self
3102 .session_id
3103 .as_ref()
3104 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3105 let request_id = self.next_request_id;
3106 self.next_request_id += 1;
3107 let prompt_blocks = prompt_content_blocks(&self.cwd, &prompt);
3108 self.write_json(serde_json::json!({
3109 "jsonrpc": "2.0",
3110 "id": request_id,
3111 "method": "session/prompt",
3112 "params": {
3113 "sessionId": session_id,
3114 "prompt": prompt_blocks,
3115 },
3116 }))
3117 .await?;
3118 self.prompt_request_id = Some(request_id);
3119 Ok(())
3120 }
3121
3122 async fn cancel(&mut self) -> AdapterResult<bool> {
3123 let Some(session_id) = &self.session_id else {
3124 return Ok(false);
3125 };
3126 self.write_json(serde_json::json!({
3127 "jsonrpc": "2.0",
3128 "method": "session/cancel",
3129 "params": {"sessionId": session_id, "_meta": {}},
3130 }))
3131 .await?;
3132 let settled = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
3133 loop {
3134 match <Self as AgentAdapter>::next_event(self).await {
3135 Some(Ok(AgentEvent::TurnComplete { .. })) | None => break,
3136 Some(Ok(_)) => {}
3137 Some(Err(_)) => break,
3138 }
3139 }
3140 })
3141 .await
3142 .is_ok();
3143 if !settled {
3144 self.reload().await?;
3148 }
3149 Ok(true)
3150 }
3151
3152 async fn answer_permission(
3153 &mut self,
3154 request_id: String,
3155 answer: PermissionAnswer,
3156 ) -> AdapterResult<()> {
3157 let id = request_id
3158 .parse::<u64>()
3159 .map(Value::from)
3160 .unwrap_or_else(|_| Value::String(request_id));
3161 let outcome = match answer {
3162 PermissionAnswer::Selected { option_id } => {
3163 serde_json::json!({"outcome": "selected", "optionId": option_id})
3164 }
3165 PermissionAnswer::Cancelled => serde_json::json!({"outcome": "cancelled"}),
3166 };
3167 self.write_json(serde_json::json!({
3168 "jsonrpc": "2.0",
3169 "id": id,
3170 "result": {"outcome": outcome},
3174 }))
3175 .await
3176 }
3177
3178 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
3179 let session_id = self
3180 .session_id
3181 .as_ref()
3182 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3183 let policy = match mode.as_str() {
3184 "plan" => "codeswarm:mode:plan",
3185 "default" | "manual" => "codeswarm:mode:manual",
3186 "accept-edits" => "codeswarm:mode:accept-edits",
3187 "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
3188 other => other,
3189 };
3190 let native_mode = crate::policy::resolve(policy, &self.modes)
3191 .map(|mode| mode.id)
3192 .unwrap_or(mode);
3193 let _ = self
3194 .request(
3195 "session/set_mode",
3196 serde_json::json!({"sessionId": session_id, "modeId": native_mode.clone()}),
3197 )
3198 .await?;
3199 self.queued_events.push_back(Ok(AgentEvent::ModeUpdated {
3200 slot: self.slot,
3201 current_mode: native_mode,
3202 }));
3203 Ok(())
3204 }
3205
3206 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
3207 let session_id = self
3208 .session_id
3209 .clone()
3210 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3211 let config_id = self
3212 .model_config_id
3213 .clone()
3214 .ok_or(AdapterError::Unsupported("set_model"))?;
3215 if !self.models.iter().any(|candidate| candidate.id == model) {
3216 return Err(AdapterError::Protocol(
3217 "model is not advertised by the agent".into(),
3218 ));
3219 }
3220 let _ = self
3221 .request(
3222 "session/set_config_option",
3223 serde_json::json!({
3224 "sessionId": session_id,
3225 "configId": config_id,
3226 "value": model,
3227 }),
3228 )
3229 .await?;
3230 Ok(())
3231 }
3232
3233 async fn reload(&mut self) -> AdapterResult<()> {
3234 let session_id = self
3239 .capabilities
3240 .supports_session_load
3241 .then(|| self.session_id.clone())
3242 .flatten();
3243 self.stop().await?;
3244 self.session_id = session_id.clone();
3245 let result = self.start().await;
3246 if result.is_err() {
3247 self.session_id = session_id;
3251 }
3252 result
3253 }
3254
3255 async fn stop(&mut self) -> AdapterResult<()> {
3256 let terminals = std::mem::take(&mut self.terminals);
3257 for terminal in terminals.values() {
3258 terminal.stop().await;
3259 }
3260 if let Some(mut child) = self.child.take() {
3261 terminate_child(&mut child).await?;
3262 }
3263 self.reader = None;
3264 self.session_id = None;
3265 self.prompt_request_id = None;
3266 if let Some(task) = self.stderr_task.take() {
3267 task.abort();
3268 let _ = task.await;
3269 }
3270 Ok(())
3271 }
3272
3273 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
3274 if let Some(event) = self.queued_events.pop_front() {
3275 return Some(event);
3276 }
3277 loop {
3278 let line = match self.read_line().await {
3279 Ok(line) => line,
3280 Err(error) => return Some(Err(error)),
3281 };
3282 let value: Value = match serde_json::from_str(&line) {
3283 Ok(value) => value,
3284 Err(_) => continue,
3285 };
3286 match self.reject_empty_permission_request(&value).await {
3287 Ok(true) => continue,
3288 Ok(false) => {}
3289 Err(error) => return Some(Err(error)),
3290 }
3291 match self.handle_client_request(&value).await {
3292 Ok(true) => continue,
3293 Ok(false) => {}
3294 Err(error) => return Some(Err(error)),
3295 }
3296 match parse_acp_notification(self.slot, &line) {
3297 Ok(Some(event)) => {
3298 if let AgentEvent::ModelsReplaced {
3299 config_id, models, ..
3300 } = &event
3301 {
3302 self.model_config_id = Some(config_id.clone());
3303 self.models = models.clone();
3304 self.capabilities.supports_models = !models.is_empty();
3305 }
3306 return Some(Ok(event));
3307 }
3308 Ok(None) => {}
3309 Err(error) => return Some(Err(error)),
3310 }
3311 if value.get("id").is_some_and(|id| {
3312 self.prompt_request_id
3313 .is_some_and(|expected| rpc_id_to_string(id) == expected.to_string())
3314 }) {
3315 if let Some(error) = value.get("error") {
3316 self.prompt_request_id = None;
3317 return Some(Err(AdapterError::Protocol(error.to_string())));
3318 }
3319 self.prompt_request_id = None;
3320 return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3321 }
3322 }
3323 }
3324}
3325
3326fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3327 let value: Value =
3328 serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3329 let method = value.get("method").and_then(Value::as_str);
3330 if method == Some("session/request_permission") {
3331 let params = value.get("params").cloned().unwrap_or(Value::Null);
3332 let request_id = value
3333 .get("id")
3334 .map(rpc_id_to_string)
3335 .unwrap_or_else(|| "permission".into());
3336 return Ok(parse_permission_event(
3337 slot,
3338 ¶ms,
3339 &request_id,
3340 params.get("options"),
3341 ));
3342 }
3343 if method != Some("session/update") {
3344 return Ok(None);
3345 }
3346 let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3347 return Ok(None);
3348 };
3349 let kind = update.get("sessionUpdate").and_then(Value::as_str);
3350 if kind == Some("config_option_update")
3351 && let Some((config_id, models, current_model)) = parse_model_config(update)
3352 {
3353 return Ok(Some(AgentEvent::ModelsReplaced {
3354 slot,
3355 config_id,
3356 models,
3357 current_model,
3358 }));
3359 }
3360 if kind == Some("request_permission") {
3361 let request_id = update
3362 .get("toolCall")
3363 .and_then(|tool| tool.get("toolCallId"))
3364 .and_then(Value::as_str)
3365 .unwrap_or("permission");
3366 return Ok(parse_permission_event(
3367 slot,
3368 update,
3369 request_id,
3370 update.get("options"),
3371 ));
3372 }
3373 if kind == Some("available_commands_update") {
3374 let commands = update
3375 .get("availableCommands")
3376 .and_then(Value::as_array)
3377 .map(|commands| {
3378 commands
3379 .iter()
3380 .filter_map(|command| {
3381 let name = command.get("name").and_then(Value::as_str)?.trim();
3382 (!name.is_empty()).then(|| AgentCommand {
3383 name: name.to_owned(),
3384 })
3385 })
3386 .collect::<Vec<_>>()
3387 })
3388 .unwrap_or_default();
3389 return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3390 }
3391 if kind == Some("current_mode_update") {
3392 if let Some(mode) = update
3393 .get("currentModeId")
3394 .and_then(Value::as_str)
3395 .filter(|mode| !mode.trim().is_empty())
3396 {
3397 return Ok(Some(AgentEvent::ModeUpdated {
3398 slot,
3399 current_mode: mode.to_owned(),
3400 }));
3401 }
3402 return Ok(None);
3403 }
3404 if kind == Some("usage_update") {
3405 let Some(used) = update.get("used").and_then(Value::as_u64) else {
3406 return Ok(None);
3407 };
3408 let Some(size) = update.get("size").and_then(Value::as_u64) else {
3409 return Ok(None);
3410 };
3411 return Ok(Some(AgentEvent::UsageUpdated {
3412 slot,
3413 usage: UsageUpdate { used, size },
3414 }));
3415 }
3416 if let Some(terminal) = parse_terminal_event(update, kind) {
3417 return Ok(Some(AgentEvent::Terminal {
3418 slot,
3419 event: terminal,
3420 }));
3421 }
3422 let text = update
3423 .get("content")
3424 .and_then(|content| content.get("text"))
3425 .and_then(Value::as_str)
3426 .map(str::to_owned);
3427 if kind == Some("user_message_chunk") {
3428 return Ok(text
3429 .filter(|text| !text.is_empty())
3430 .map(|text| AgentEvent::UserText { slot, text }));
3431 }
3432 if kind == Some("agent_message_chunk")
3433 && let Some(mode) = text
3434 .as_deref()
3435 .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3436 .map(str::trim)
3437 .filter(|mode| !mode.is_empty())
3438 {
3439 return Ok(Some(AgentEvent::ModesReplaced {
3444 slot,
3445 modes: vec![Mode {
3446 id: mode.to_owned(),
3447 label: mode.to_owned(),
3448 }],
3449 current_mode: Some(mode.to_owned()),
3450 }));
3451 }
3452 match (kind, text) {
3453 (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3454 Ok(Some(AgentEvent::Text { slot, text }))
3455 }
3456 (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3457 Ok(Some(AgentEvent::Thought { slot, text }))
3458 }
3459 (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3460 let id = update
3461 .get("toolCallId")
3462 .and_then(Value::as_str)
3463 .unwrap_or("unknown-tool")
3464 .to_owned();
3465 let title = update
3466 .get("title")
3467 .and_then(Value::as_str)
3468 .unwrap_or("Tool call")
3469 .to_owned();
3470 let status = match update.get("status").and_then(Value::as_str) {
3471 Some("completed") => ToolStatus::Completed,
3472 Some("failed") => ToolStatus::Failed,
3473 Some("in_progress") => ToolStatus::Running,
3474 _ => ToolStatus::Pending,
3475 };
3476 Ok(Some(AgentEvent::Tool {
3477 slot,
3478 update: ToolUpdate {
3479 id,
3480 title,
3481 status,
3482 detail: None,
3483 },
3484 }))
3485 }
3486 _ => Ok(None),
3487 }
3488}
3489
3490fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3491 let config = value
3492 .get("configOptions")?
3493 .as_array()?
3494 .iter()
3495 .find(|option| {
3496 option.get("category").and_then(Value::as_str) == Some("model")
3497 && matches!(
3498 option.get("type").and_then(Value::as_str),
3499 Some("select" | "enum")
3500 )
3501 })?;
3502 let config_id = config.get("id")?.as_str()?.to_owned();
3503 let models = config
3504 .get("options")?
3505 .as_array()?
3506 .iter()
3507 .filter_map(|option| {
3508 let id = option.get("value")?.as_str()?.to_owned();
3509 let label = option
3510 .get("name")
3511 .or_else(|| option.get("label"))
3512 .and_then(Value::as_str)
3513 .unwrap_or(&id)
3514 .to_owned();
3515 Some(Mode { id, label })
3516 })
3517 .collect::<Vec<_>>();
3518 (!models.is_empty()).then(|| {
3519 let current = config
3520 .get("currentValue")
3521 .and_then(Value::as_str)
3522 .map(str::to_owned);
3523 (config_id, models, current)
3524 })
3525}
3526
3527fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3532 let nested = value.get("terminal").unwrap_or(value);
3533 let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3534 let id = nested
3535 .get("terminalId")
3536 .or_else(|| nested.get("terminal_id"))
3537 .or_else(|| nested.get("id"))
3538 .and_then(Value::as_str)
3539 .unwrap_or("terminal")
3540 .to_owned();
3541 match kind {
3542 "terminal_created" | "terminal_create" | "terminal_started" => {
3543 let command = nested
3544 .get("command")
3545 .and_then(Value::as_str)
3546 .unwrap_or("")
3547 .to_owned();
3548 Some(TerminalEvent::Created { id, command })
3549 }
3550 "terminal_output" | "terminal_output_chunk" => {
3551 let text = nested
3552 .get("output")
3553 .or_else(|| nested.get("text"))
3554 .and_then(Value::as_str)
3555 .unwrap_or("")
3556 .to_owned();
3557 Some(TerminalEvent::Output { id, text })
3558 }
3559 "terminal_exited" | "terminal_exit" => {
3560 let code = nested
3561 .get("exitCode")
3562 .or_else(|| nested.get("exit_code"))
3563 .or_else(|| nested.get("code"))
3564 .and_then(Value::as_i64)
3565 .unwrap_or(0) as i32;
3566 Some(TerminalEvent::Exited { id, code })
3567 }
3568 "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3569 _ => None,
3570 }
3571}
3572
3573fn parse_permission_event(
3574 slot: RosterSlot,
3575 value: &Value,
3576 request_id: &str,
3577 options: Option<&Value>,
3578) -> Option<AgentEvent> {
3579 let tool = value.get("toolCall").unwrap_or(value);
3580 let title = tool
3581 .get("title")
3582 .and_then(Value::as_str)
3583 .unwrap_or("Agent requests permission")
3584 .to_owned();
3585 let (options, option_ids): (Vec<String>, Vec<String>) = options
3586 .and_then(Value::as_array)
3587 .map(|options| {
3588 options
3589 .iter()
3590 .filter_map(|option| {
3591 let label = option
3592 .get("name")
3593 .or_else(|| option.get("optionId"))
3594 .and_then(Value::as_str)?
3595 .to_owned();
3596 let option_id = option
3597 .get("optionId")
3598 .or_else(|| option.get("id"))
3599 .and_then(Value::as_str)
3600 .map(str::to_owned)
3601 .unwrap_or_else(|| label.clone());
3602 Some((label, option_id))
3603 })
3604 .unzip()
3605 })
3606 .unwrap_or_default();
3607 if options.is_empty() {
3608 return None;
3609 }
3610 Some(AgentEvent::Permission {
3611 slot,
3612 request: PermissionRequest {
3613 id: request_id.to_owned(),
3614 title,
3615 options,
3616 option_ids,
3617 },
3618 })
3619}
3620
3621fn rpc_id_to_string(value: &Value) -> String {
3622 value
3623 .as_str()
3624 .map(str::to_owned)
3625 .or_else(|| value.as_u64().map(|id| id.to_string()))
3626 .unwrap_or_else(|| value.to_string())
3627}
3628
3629#[cfg(test)]
3630mod tests {
3631 use super::{
3632 AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
3633 RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
3634 parse_model_config, prompt_content_blocks, read_bounded_line,
3635 };
3636 #[cfg(target_os = "linux")]
3637 use super::{isolate_process_group, terminate_child};
3638 use crate::TerminalEvent;
3639 use crate::{
3640 AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
3641 persistence::SessionMetadataStore,
3642 relay::{DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
3643 };
3644 use async_trait::async_trait;
3645 use serde_json::Value;
3646 use std::sync::{
3647 Arc, Mutex,
3648 atomic::{AtomicUsize, Ordering},
3649 };
3650
3651 fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
3652 let nonce = std::time::SystemTime::now()
3653 .duration_since(std::time::UNIX_EPOCH)
3654 .expect("clock")
3655 .as_nanos();
3656 std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
3657 }
3658
3659 #[test]
3660 fn malformed_file_writes_preserve_existing_content() {
3661 let root = unique_test_path("codeswarm-write-validation", "dir");
3662 std::fs::create_dir_all(&root).unwrap();
3663 let file = root.join("keep.txt");
3664 std::fs::write(&file, "valuable content").unwrap();
3665 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
3666 for content in [
3667 Value::Null,
3668 serde_json::json!(false),
3669 serde_json::json!(42),
3670 serde_json::json!([]),
3671 ] {
3672 assert!(
3673 adapter
3674 .write_workspace_text(
3675 &serde_json::json!({"path":"keep.txt", "content": content})
3676 )
3677 .is_err()
3678 );
3679 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3680 }
3681 assert!(
3682 adapter
3683 .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
3684 .is_err()
3685 );
3686 assert!(
3687 adapter
3688 .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
3689 .is_err()
3690 );
3691 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3692 adapter
3693 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
3694 .unwrap();
3695 assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
3696 adapter
3697 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
3698 .unwrap();
3699 assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
3700 std::fs::remove_dir_all(root).unwrap();
3701 }
3702
3703 #[tokio::test]
3704 async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
3705 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
3706 let mut adapter = AcpAdapter::new(
3707 0,
3708 std::env::current_dir().unwrap(),
3709 "sh",
3710 vec!["-c".into(), script.into()],
3711 );
3712 adapter.start().await.unwrap();
3713 let error = adapter
3714 .request_with_timeout(
3715 "session/set_mode",
3716 serde_json::json!({}),
3717 std::time::Duration::from_millis(10),
3718 )
3719 .await
3720 .unwrap_err();
3721 assert!(error.to_string().contains("session/set_mode timed out"));
3722 adapter.stop().await.unwrap();
3723 assert!(adapter.child.is_none());
3724 }
3725
3726 #[tokio::test]
3727 async fn goals_reach_every_roster_slot_without_native_goal_support() {
3728 use crate::goal::GoalCommand;
3729 let hosts = (0..3)
3730 .map(|slot| {
3731 AdapterHost::new(
3732 Box::new(ScriptedAdapter::new(
3733 slot,
3734 AgentCapabilities::default(),
3735 [
3736 AgentEvent::TurnComplete { slot },
3737 AgentEvent::TurnComplete { slot },
3738 ],
3739 )),
3740 None,
3741 )
3742 })
3743 .collect();
3744 let mut relay = RelayHost::new(hosts, 10).unwrap();
3745 relay.start().await.unwrap();
3746 let task = relay
3747 .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
3748 .unwrap()
3749 .unwrap();
3750 relay.relay_mut().enqueue_human(task, Some(0));
3751 for slot in 0..3 {
3752 relay.run_turn("", 0).await.unwrap();
3753 let (actual, prompt) = relay.dispatches().last().unwrap();
3754 assert_eq!(*actual, slot);
3755 assert!(prompt.contains("Active shared goal: Ship the settings screen"));
3756 }
3757 let snapshot = relay.session_metadata();
3758 let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
3759 assert!(restored.is_some());
3760 relay.restore_goal(restored);
3761 relay.reload(0).await.unwrap();
3762 relay.run_turn("", 0).await.unwrap();
3763 assert!(
3764 relay
3765 .dispatches()
3766 .last()
3767 .unwrap()
3768 .1
3769 .contains("Active shared goal: Ship the settings screen")
3770 );
3771 relay.apply_goal(GoalCommand::Done).unwrap();
3772 relay.run_turn("", 0).await.unwrap();
3773 assert!(
3774 relay
3775 .dispatches()
3776 .last()
3777 .unwrap()
3778 .1
3779 .contains("No active shared goal")
3780 );
3781 relay.apply_goal(GoalCommand::Clear).unwrap();
3782 assert!(relay.session_metadata().get("goal").unwrap().is_null());
3783 }
3784
3785 #[tokio::test]
3786 async fn replacement_agent_receives_task_after_public_journal_pruning() {
3787 let hosts = (0..2)
3788 .map(|slot| {
3789 AdapterHost::new(
3790 Box::new(ScriptedAdapter::new(
3791 slot,
3792 AgentCapabilities::default(),
3793 [
3794 AgentEvent::Text {
3795 slot,
3796 text: "progress".into(),
3797 },
3798 AgentEvent::TurnComplete { slot },
3799 AgentEvent::Text {
3800 slot,
3801 text: "more progress".into(),
3802 },
3803 AgentEvent::TurnComplete { slot },
3804 ],
3805 )),
3806 None,
3807 )
3808 })
3809 .collect();
3810 let mut relay = RelayHost::new(hosts, 10).unwrap();
3811 relay.start().await.unwrap();
3812 relay
3813 .relay_mut()
3814 .enqueue_human("Fix the login bug", Some(0));
3815 relay.run_turn("", 0).await.unwrap();
3816 relay.run_turn("", 0).await.unwrap();
3817 relay.run_turn("", 0).await.unwrap();
3818 relay.reload(1).await.unwrap();
3819 assert!(
3820 !relay
3821 .relay_mut()
3822 .unseen_context(1)
3823 .contains("Fix the login bug")
3824 );
3825 relay.run_turn("", 0).await.unwrap();
3826 assert!(
3827 relay
3828 .dispatches()
3829 .last()
3830 .unwrap()
3831 .1
3832 .contains("Shared task:\nFix the login bug")
3833 );
3834 }
3835
3836 #[cfg(target_os = "linux")]
3837 #[tokio::test]
3838 async fn termination_kills_only_the_verified_isolated_child_group() {
3839 use nix::unistd::{Pid, getpgid, getpgrp};
3840 use tokio::io::{AsyncBufReadExt, BufReader};
3841
3842 let own_group = getpgrp();
3843 let mut command = tokio::process::Command::new("sh");
3844 isolate_process_group(&mut command);
3845 command
3846 .arg("-c")
3847 .arg("sleep 60 & echo $!; wait")
3848 .stdout(std::process::Stdio::piped());
3849 let mut child = command.spawn().expect("spawn isolated shell");
3850 let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
3851 assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
3852 assert_ne!(leader, own_group);
3853
3854 let stdout = child.stdout.take().expect("child stdout");
3855 let mut lines = BufReader::new(stdout).lines();
3856 let descendant = lines
3857 .next_line()
3858 .await
3859 .expect("read descendant pid")
3860 .expect("descendant pid")
3861 .parse::<i32>()
3862 .expect("numeric descendant pid");
3863 let descendant = Pid::from_raw(descendant);
3864 assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
3865
3866 terminate_child(&mut child).await.expect("terminate group");
3867 for _ in 0..100 {
3868 if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
3869 return;
3870 }
3871 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
3872 }
3873 panic!("descendant {descendant} survived isolated group termination");
3874 }
3875
3876 #[test]
3877 fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
3878 assert_eq!(
3879 parse_command_line(r#"npx -y "@agentclientprotocol/codex-acp" --flag 'two words'"#),
3880 Ok((
3881 "npx".into(),
3882 vec![
3883 "-y".into(),
3884 "@agentclientprotocol/codex-acp".into(),
3885 "--flag".into(),
3886 "two words".into(),
3887 ]
3888 ),)
3889 );
3890 assert_eq!(
3891 parse_command_line(r#"agent "" escaped\ argument"#),
3892 Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
3893 );
3894 }
3895
3896 #[test]
3897 fn acp_prompt_expands_safe_at_path_resources() {
3898 let root = unique_test_path("codeswarm-prompt-resource", "dir");
3899 std::fs::create_dir_all(&root).expect("workspace");
3900 std::fs::write(root.join("note.md"), "resource text").expect("resource");
3901 let blocks = prompt_content_blocks(&root, "inspect @note.md");
3902 assert_eq!(blocks[0]["type"], "text");
3903 assert_eq!(blocks[0]["text"], "inspect @note.md");
3904 assert_eq!(blocks[1]["type"], "resource");
3905 assert_eq!(blocks[1]["resource"]["text"], "resource text");
3906 assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
3907 std::fs::remove_dir_all(root).expect("cleanup workspace");
3908 }
3909
3910 #[tokio::test]
3911 async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
3912 let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
3913 bytes.push(b'\n');
3914 let mut reader = tokio::io::BufReader::new(bytes.as_slice());
3915 assert!(matches!(
3916 read_bounded_line(&mut reader).await,
3917 Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
3918 ));
3919 }
3920
3921 #[test]
3922 fn rejects_malformed_configured_commands_before_spawn() {
3923 assert_eq!(
3924 parse_command_line("agent 'unfinished"),
3925 Err(super::CommandParseError::UnterminatedQuote)
3926 );
3927 assert_eq!(
3928 parse_command_line("agent\\"),
3929 Err(super::CommandParseError::TrailingEscape)
3930 );
3931 assert_eq!(
3932 parse_command_line(" \t"),
3933 Err(super::CommandParseError::Empty)
3934 );
3935 }
3936
3937 #[derive(Debug)]
3938 struct PendingAdapter {
3939 slot: usize,
3940 hang_on_cancel: bool,
3941 }
3942
3943 #[derive(Debug)]
3944 struct ConcurrentStartAdapter {
3945 slot: usize,
3946 barrier: Arc<tokio::sync::Barrier>,
3947 }
3948
3949 #[async_trait]
3950 impl AgentAdapter for ConcurrentStartAdapter {
3951 fn slot(&self) -> usize {
3952 self.slot
3953 }
3954
3955 fn capabilities(&self) -> AgentCapabilities {
3956 AgentCapabilities::default()
3957 }
3958
3959 async fn start(&mut self) -> super::AdapterResult<()> {
3960 self.barrier.wait().await;
3961 Ok(())
3962 }
3963
3964 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
3965 Ok(())
3966 }
3967
3968 async fn cancel(&mut self) -> super::AdapterResult<bool> {
3969 Ok(true)
3970 }
3971
3972 async fn answer_permission(
3973 &mut self,
3974 _request_id: String,
3975 _answer: PermissionAnswer,
3976 ) -> super::AdapterResult<()> {
3977 Ok(())
3978 }
3979
3980 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
3981 Ok(())
3982 }
3983
3984 async fn reload(&mut self) -> super::AdapterResult<()> {
3985 Ok(())
3986 }
3987
3988 async fn stop(&mut self) -> super::AdapterResult<()> {
3989 Ok(())
3990 }
3991
3992 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
3993 std::future::pending().await
3994 }
3995 }
3996
3997 #[derive(Debug)]
3998 struct PermissionBlockingAdapter {
3999 slot: usize,
4000 phase: u8,
4001 }
4002
4003 #[async_trait]
4004 impl AgentAdapter for PermissionBlockingAdapter {
4005 fn slot(&self) -> usize {
4006 self.slot
4007 }
4008
4009 fn capabilities(&self) -> AgentCapabilities {
4010 AgentCapabilities {
4011 supports_permissions: true,
4012 ..AgentCapabilities::default()
4013 }
4014 }
4015
4016 async fn start(&mut self) -> super::AdapterResult<()> {
4017 Ok(())
4018 }
4019
4020 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4021 Ok(())
4022 }
4023
4024 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4025 Ok(true)
4026 }
4027
4028 async fn answer_permission(
4029 &mut self,
4030 request_id: String,
4031 answer: PermissionAnswer,
4032 ) -> super::AdapterResult<()> {
4033 if self.phase != 1 || request_id != "permission-1" {
4034 return Err(super::AdapterError::Protocol(
4035 "unexpected permission response".into(),
4036 ));
4037 }
4038 assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4039 self.phase = 2;
4040 Ok(())
4041 }
4042
4043 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4044 Ok(())
4045 }
4046
4047 async fn reload(&mut self) -> super::AdapterResult<()> {
4048 Ok(())
4049 }
4050
4051 async fn stop(&mut self) -> super::AdapterResult<()> {
4052 Ok(())
4053 }
4054
4055 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4056 match self.phase {
4057 0 => {
4058 self.phase = 1;
4059 Some(Ok(AgentEvent::Permission {
4060 slot: self.slot,
4061 request: crate::PermissionRequest {
4062 id: "permission-1".into(),
4063 title: "Allow?".into(),
4064 options: vec!["Allow".into()],
4065 option_ids: vec!["allow".into()],
4066 },
4067 }))
4068 }
4069 1 => std::future::pending().await,
4070 _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4071 }
4072 }
4073 }
4074
4075 #[async_trait]
4076 impl AgentAdapter for PendingAdapter {
4077 fn slot(&self) -> usize {
4078 self.slot
4079 }
4080
4081 fn capabilities(&self) -> AgentCapabilities {
4082 AgentCapabilities {
4083 supports_cancel: true,
4084 ..AgentCapabilities::default()
4085 }
4086 }
4087
4088 async fn start(&mut self) -> super::AdapterResult<()> {
4089 Ok(())
4090 }
4091
4092 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4093 Ok(())
4094 }
4095
4096 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4097 if self.hang_on_cancel {
4098 return std::future::pending().await;
4099 }
4100 Ok(true)
4101 }
4102
4103 async fn answer_permission(
4104 &mut self,
4105 _request_id: String,
4106 _answer: PermissionAnswer,
4107 ) -> super::AdapterResult<()> {
4108 Ok(())
4109 }
4110
4111 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4112 Ok(())
4113 }
4114
4115 async fn reload(&mut self) -> super::AdapterResult<()> {
4116 Ok(())
4117 }
4118
4119 async fn stop(&mut self) -> super::AdapterResult<()> {
4120 Ok(())
4121 }
4122
4123 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4124 std::future::pending().await
4125 }
4126 }
4127
4128 #[derive(Debug)]
4129 struct StopTrackingAdapter {
4130 slot: usize,
4131 stopped: Arc<AtomicUsize>,
4132 fail_stop: bool,
4133 }
4134
4135 #[derive(Debug)]
4136 struct ModeOrderAdapter {
4137 slot: usize,
4138 log: Arc<Mutex<Vec<String>>>,
4139 phase: u8,
4140 }
4141
4142 #[derive(Debug)]
4143 struct StartupAcpAdapter {
4144 slot: usize,
4145 events: std::collections::VecDeque<AgentEvent>,
4146 }
4147
4148 impl StartupAcpAdapter {
4149 fn new(slot: usize) -> Self {
4150 Self {
4151 slot,
4152 events: [
4153 AgentEvent::ModesReplaced {
4154 slot,
4155 modes: vec![Mode {
4156 id: "full-access".into(),
4157 label: "Auto pilot".into(),
4158 }],
4159 current_mode: Some("full-access".into()),
4160 },
4161 AgentEvent::Ready {
4162 slot,
4163 capabilities: AgentCapabilities {
4164 supports_modes: true,
4165 ..AgentCapabilities::default()
4166 },
4167 },
4168 ]
4169 .into(),
4170 }
4171 }
4172 }
4173
4174 #[async_trait]
4175 impl AgentAdapter for StartupAcpAdapter {
4176 fn slot(&self) -> usize {
4177 self.slot
4178 }
4179
4180 fn protocol(&self) -> &'static str {
4181 "acp"
4182 }
4183
4184 fn capabilities(&self) -> AgentCapabilities {
4185 AgentCapabilities {
4186 supports_modes: true,
4187 ..AgentCapabilities::default()
4188 }
4189 }
4190
4191 async fn start(&mut self) -> super::AdapterResult<()> {
4192 Ok(())
4193 }
4194
4195 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4196 Ok(())
4197 }
4198
4199 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4200 Ok(true)
4201 }
4202
4203 async fn answer_permission(
4204 &mut self,
4205 _request_id: String,
4206 _answer: PermissionAnswer,
4207 ) -> super::AdapterResult<()> {
4208 Ok(())
4209 }
4210
4211 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4212 Ok(())
4213 }
4214
4215 async fn reload(&mut self) -> super::AdapterResult<()> {
4216 self.events = Self::new(self.slot).events;
4217 Ok(())
4218 }
4219
4220 async fn stop(&mut self) -> super::AdapterResult<()> {
4221 Ok(())
4222 }
4223
4224 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4225 self.events.pop_front().map(Ok)
4226 }
4227 }
4228
4229 #[async_trait]
4230 impl AgentAdapter for ModeOrderAdapter {
4231 fn slot(&self) -> usize {
4232 self.slot
4233 }
4234
4235 fn capabilities(&self) -> AgentCapabilities {
4236 AgentCapabilities {
4237 supports_modes: true,
4238 ..AgentCapabilities::default()
4239 }
4240 }
4241
4242 async fn start(&mut self) -> super::AdapterResult<()> {
4243 self.log.lock().expect("log").push("start".into());
4244 Ok(())
4245 }
4246
4247 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4248 self.log.lock().expect("log").push("prompt".into());
4249 Ok(())
4250 }
4251
4252 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4253 Ok(true)
4254 }
4255
4256 async fn answer_permission(
4257 &mut self,
4258 _request_id: String,
4259 _answer: PermissionAnswer,
4260 ) -> super::AdapterResult<()> {
4261 Ok(())
4262 }
4263
4264 async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4265 self.log.lock().expect("log").push(format!("mode:{mode}"));
4266 Ok(())
4267 }
4268
4269 async fn reload(&mut self) -> super::AdapterResult<()> {
4270 self.log.lock().expect("log").push("reload".into());
4271 self.phase = 0;
4272 Ok(())
4273 }
4274
4275 async fn stop(&mut self) -> super::AdapterResult<()> {
4276 Ok(())
4277 }
4278
4279 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4280 match self.phase {
4281 0 => {
4282 self.phase = 1;
4283 Some(Ok(AgentEvent::ModesReplaced {
4284 slot: self.slot,
4285 modes: vec![Mode {
4286 id: "yolo".into(),
4287 label: "YOLO".into(),
4288 }],
4289 current_mode: None,
4290 }))
4291 }
4292 1 => {
4293 self.phase = 2;
4294 Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4295 }
4296 _ => std::future::pending().await,
4297 }
4298 }
4299 }
4300
4301 #[async_trait]
4302 impl AgentAdapter for StopTrackingAdapter {
4303 fn slot(&self) -> usize {
4304 self.slot
4305 }
4306
4307 fn capabilities(&self) -> AgentCapabilities {
4308 AgentCapabilities::default()
4309 }
4310
4311 async fn start(&mut self) -> super::AdapterResult<()> {
4312 Ok(())
4313 }
4314
4315 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4316 Ok(())
4317 }
4318
4319 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4320 Ok(false)
4321 }
4322
4323 async fn answer_permission(
4324 &mut self,
4325 _request_id: String,
4326 _answer: PermissionAnswer,
4327 ) -> super::AdapterResult<()> {
4328 Ok(())
4329 }
4330
4331 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4332 Ok(())
4333 }
4334
4335 async fn reload(&mut self) -> super::AdapterResult<()> {
4336 Ok(())
4337 }
4338
4339 async fn stop(&mut self) -> super::AdapterResult<()> {
4340 self.stopped.fetch_add(1, Ordering::Relaxed);
4341 if self.fail_stop {
4342 Err(super::AdapterError::Transport("stop failed".into()))
4343 } else {
4344 Ok(())
4345 }
4346 }
4347
4348 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4349 None
4350 }
4351 }
4352
4353 #[derive(Debug)]
4357 struct FailingStartAdapter {
4358 slot: usize,
4359 stopped: Arc<AtomicUsize>,
4360 }
4361
4362 #[async_trait]
4363 impl AgentAdapter for FailingStartAdapter {
4364 fn slot(&self) -> usize {
4365 self.slot
4366 }
4367
4368 fn capabilities(&self) -> AgentCapabilities {
4369 AgentCapabilities::default()
4370 }
4371
4372 async fn start(&mut self) -> super::AdapterResult<()> {
4373 Err(super::AdapterError::Spawn("startup failed".into()))
4374 }
4375
4376 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4377 Ok(())
4378 }
4379
4380 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4381 Ok(false)
4382 }
4383
4384 async fn answer_permission(
4385 &mut self,
4386 _request_id: String,
4387 _answer: PermissionAnswer,
4388 ) -> super::AdapterResult<()> {
4389 Ok(())
4390 }
4391
4392 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4393 Ok(())
4394 }
4395
4396 async fn reload(&mut self) -> super::AdapterResult<()> {
4397 Ok(())
4398 }
4399
4400 async fn stop(&mut self) -> super::AdapterResult<()> {
4401 self.stopped.fetch_add(1, Ordering::Relaxed);
4402 Ok(())
4403 }
4404
4405 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4406 None
4407 }
4408 }
4409
4410 #[tokio::test]
4411 async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4412 let stopped = Arc::new(AtomicUsize::new(0));
4413 let relay = RelayHost::new(
4414 vec![
4415 AdapterHost::new(
4416 Box::new(StopTrackingAdapter {
4417 slot: 0,
4418 stopped: Arc::clone(&stopped),
4419 fail_stop: true,
4420 }),
4421 None,
4422 ),
4423 AdapterHost::new(
4424 Box::new(StopTrackingAdapter {
4425 slot: 1,
4426 stopped: Arc::clone(&stopped),
4427 fail_stop: false,
4428 }),
4429 None,
4430 ),
4431 ],
4432 4,
4433 )
4434 .expect("relay");
4435 let mut relay = relay;
4436
4437 let error = relay.stop().await.expect_err("first stop failure");
4438 assert!(error.to_string().contains("stop failed"));
4439 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4440 }
4441
4442 #[tokio::test]
4443 async fn relay_start_cleans_up_the_adapter_that_failed_startup() {
4444 let stopped = Arc::new(AtomicUsize::new(0));
4445 let mut relay = RelayHost::new(
4446 vec![
4447 AdapterHost::new(
4448 Box::new(StopTrackingAdapter {
4449 slot: 0,
4450 stopped: Arc::clone(&stopped),
4451 fail_stop: false,
4452 }),
4453 None,
4454 ),
4455 AdapterHost::new(
4456 Box::new(FailingStartAdapter {
4457 slot: 1,
4458 stopped: Arc::clone(&stopped),
4459 }),
4460 None,
4461 ),
4462 ],
4463 4,
4464 )
4465 .expect("relay");
4466
4467 assert!(relay.start().await.is_err());
4468 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4469 }
4470
4471 #[test]
4472 fn parses_acp_text_without_ui_dependency() {
4473 let event = parse_acp_notification(
4474 2,
4475 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4476 )
4477 .expect("valid ACP")
4478 .expect("text event");
4479 assert_eq!(
4480 event,
4481 AgentEvent::Text {
4482 slot: 2,
4483 text: "hello".into(),
4484 }
4485 );
4486 }
4487
4488 #[test]
4489 fn parses_acp_state_notifications_at_the_adapter_boundary() {
4490 let commands = parse_acp_notification(
4491 3,
4492 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
4493 )
4494 .expect("valid ACP")
4495 .expect("commands event");
4496 assert_eq!(
4497 commands,
4498 AgentEvent::CommandsReplaced {
4499 slot: 3,
4500 commands: vec![crate::AgentCommand {
4501 name: "review".into()
4502 }]
4503 }
4504 );
4505
4506 let mode = parse_acp_notification(
4507 3,
4508 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
4509 )
4510 .expect("valid ACP")
4511 .expect("mode event");
4512 assert_eq!(
4513 mode,
4514 AgentEvent::ModeUpdated {
4515 slot: 3,
4516 current_mode: "review".into()
4517 }
4518 );
4519
4520 let usage = parse_acp_notification(
4521 3,
4522 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
4523 )
4524 .expect("valid ACP")
4525 .expect("usage event");
4526 assert_eq!(
4527 usage,
4528 AgentEvent::UsageUpdated {
4529 slot: 3,
4530 usage: crate::UsageUpdate {
4531 used: 4200,
4532 size: 128000
4533 }
4534 }
4535 );
4536
4537 let models = parse_acp_notification(
4538 3,
4539 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"}]}]}}}"#,
4540 )
4541 .expect("valid ACP")
4542 .expect("models event");
4543 assert!(matches!(
4544 models,
4545 AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
4546 if models.len() == 2 && current_model.as_deref() == Some("smart")
4547 ));
4548 assert_eq!(
4549 parse_model_config(&serde_json::json!({
4550 "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
4551 })),
4552 None
4553 );
4554
4555 let user = parse_acp_notification(
4556 3,
4557 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
4558 )
4559 .expect("valid ACP")
4560 .expect("user event");
4561 assert_eq!(
4562 user,
4563 AgentEvent::UserText {
4564 slot: 3,
4565 text: "context".into()
4566 }
4567 );
4568 }
4569
4570 #[test]
4571 fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
4572 let event = parse_acp_notification(
4573 0,
4574 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
4575 )
4576 .expect("valid ACP")
4577 .expect("mode event");
4578 assert!(matches!(
4579 event,
4580 AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
4581 if mode == "yolo" && modes[0].id == "yolo"
4582 ));
4583 }
4584
4585 #[test]
4586 fn parses_native_agy_text_without_acp_bridge() {
4587 let event = parse_agy_line(
4588 1,
4589 r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
4590 )
4591 .expect("valid stream-json")
4592 .expect("text event");
4593 assert_eq!(
4594 event,
4595 AgentEvent::Text {
4596 slot: 1,
4597 text: "hello".into(),
4598 }
4599 );
4600 }
4601
4602 #[test]
4603 fn parses_tool_lifecycle_from_each_protocol() {
4604 let agy = parse_agy_line(
4605 1,
4606 r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
4607 )
4608 .expect("valid native tool")
4609 .expect("tool event");
4610 assert!(matches!(
4611 agy,
4612 AgentEvent::Tool {
4613 update: crate::ToolUpdate {
4614 status: ToolStatus::Completed,
4615 ..
4616 },
4617 ..
4618 }
4619 ));
4620
4621 let acp = parse_acp_notification(
4622 1,
4623 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
4624 )
4625 .expect("valid ACP tool")
4626 .expect("tool event");
4627 assert!(matches!(
4628 acp,
4629 AgentEvent::Tool {
4630 update: crate::ToolUpdate {
4631 status: ToolStatus::Failed,
4632 ..
4633 },
4634 ..
4635 }
4636 ));
4637 }
4638
4639 #[test]
4640 fn parses_terminal_lifecycle_from_acp_and_native_events() {
4641 let created = parse_acp_notification(
4642 0,
4643 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
4644 )
4645 .expect("valid ACP terminal")
4646 .expect("terminal event");
4647 assert_eq!(
4648 created,
4649 AgentEvent::Terminal {
4650 slot: 0,
4651 event: TerminalEvent::Created {
4652 id: "term-1".into(),
4653 command: "cargo test".into(),
4654 },
4655 }
4656 );
4657 let output = parse_agy_line(
4658 1,
4659 r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
4660 )
4661 .expect("valid native terminal")
4662 .expect("terminal event");
4663 assert_eq!(
4664 output,
4665 AgentEvent::Terminal {
4666 slot: 1,
4667 event: TerminalEvent::Output {
4668 id: "term-1".into(),
4669 text: "ok\n".into(),
4670 },
4671 }
4672 );
4673 let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
4674 .expect("valid native release")
4675 .expect("terminal event");
4676 assert!(matches!(
4677 released,
4678 AgentEvent::Terminal {
4679 event: TerminalEvent::Released { id },
4680 ..
4681 } if id == "term-1"
4682 ));
4683 }
4684
4685 #[test]
4686 fn parses_acp_permission_requests() {
4687 let event = parse_acp_notification(
4688 0,
4689 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"}]}}}"#,
4690 )
4691 .expect("valid permission")
4692 .expect("permission event");
4693 assert!(matches!(
4694 event,
4695 AgentEvent::Permission { request, .. }
4696 if request.id == "t1"
4697 && request.title == "Write file"
4698 && request.options == ["Allow once", "Reject"]
4699 && request.option_ids == ["allow-once", "reject"]
4700 ));
4701 }
4702
4703 #[test]
4704 fn parses_acp_permission_request_as_json_rpc_request() {
4705 let event = parse_acp_notification(
4706 2,
4707 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
4708 )
4709 .expect("valid permission request")
4710 .expect("permission event");
4711 assert!(matches!(
4712 event,
4713 AgentEvent::Permission { request, .. }
4714 if request.id == "17"
4715 && request.title == "Write file"
4716 && request.options == ["allow-once", "reject"]
4717 && request.option_ids == ["allow-once", "reject"]
4718 ));
4719 }
4720
4721 #[tokio::test]
4722 async fn native_adapter_explicitly_rejects_permission_answers() {
4723 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4724 assert_eq!(
4725 adapter
4726 .answer_permission(
4727 "request".into(),
4728 PermissionAnswer::Selected {
4729 option_id: "allow".into()
4730 },
4731 )
4732 .await,
4733 Err(super::AdapterError::Unsupported("permission answer"))
4734 );
4735 }
4736
4737 #[tokio::test]
4738 async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
4739 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4740 adapter
4741 .set_mode("full-access".into())
4742 .await
4743 .expect("auto-pilot alias");
4744 assert!(matches!(
4745 adapter.next_event().await,
4746 Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
4747 ));
4748 }
4749
4750 #[tokio::test]
4751 async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
4752 let script_path = unique_test_path("codeswarm-agy-session", "sh");
4753 std::fs::write(
4754 &script_path,
4755 "#!/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",
4756 )
4757 .expect("write native test script");
4758 let mut adapter = AgyAdapter::new(
4759 0,
4760 std::env::current_dir().expect("cwd"),
4761 format!("sh {}", script_path.display()),
4762 );
4763 adapter.start().await.expect("start native adapter");
4764 assert!(adapter.next_event().await.is_some());
4766 assert!(adapter.next_event().await.is_some());
4767 adapter
4768 .send_prompt("first".into())
4769 .await
4770 .expect("first prompt");
4771 while !matches!(
4772 adapter.next_event().await,
4773 Some(Ok(AgentEvent::TurnComplete { .. }))
4774 ) {}
4775 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4776 adapter
4777 .send_prompt("follow up".into())
4778 .await
4779 .expect("follow-up prompt");
4780 while !matches!(
4781 adapter.next_event().await,
4782 Some(Ok(AgentEvent::TurnComplete { .. }))
4783 ) {}
4784 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4785 adapter.stop().await.expect("stop native adapter");
4786 std::fs::remove_file(script_path).expect("cleanup native script");
4787 }
4788
4789 #[tokio::test]
4790 async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
4791 let script_path = unique_test_path("codeswarm-agy-failure", "sh");
4792 std::fs::write(
4793 &script_path,
4794 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
4795 )
4796 .expect("write native test script");
4797 let mut adapter = AgyAdapter::new(
4798 0,
4799 std::env::current_dir().expect("cwd"),
4800 format!("sh {}", script_path.display()),
4801 );
4802 adapter.start().await.expect("start native adapter");
4803 assert!(adapter.next_event().await.is_some());
4804 assert!(adapter.next_event().await.is_some());
4805 adapter.send_prompt("fail".into()).await.expect("prompt");
4806 assert!(matches!(
4807 adapter.next_event().await,
4808 Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
4809 if detail == "agent failed"
4810 ));
4811 adapter.stop().await.expect("stop native adapter");
4812 std::fs::remove_file(script_path).expect("cleanup native script");
4813 }
4814
4815 #[tokio::test]
4816 async fn acp_adapter_initializes_session_and_completes_a_prompt() {
4817 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"}}'"#;
4818 let cwd = std::env::current_dir().expect("cwd");
4819 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4820 adapter.start().await.expect("initialize");
4821 assert!(matches!(
4822 adapter.next_event().await,
4823 Some(Ok(AgentEvent::ModesReplaced { .. }))
4824 ));
4825 assert!(matches!(
4826 adapter.next_event().await,
4827 Some(Ok(AgentEvent::Ready { .. }))
4828 ));
4829 adapter.send_prompt("hello".into()).await.expect("prompt");
4830 assert!(matches!(
4831 adapter.next_event().await,
4832 Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
4833 ));
4834 assert!(matches!(
4835 adapter.next_event().await,
4836 Some(Ok(AgentEvent::TurnComplete { .. }))
4837 ));
4838 }
4839
4840 #[tokio::test]
4841 async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
4842 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"}}'"#;
4843 let cwd = std::env::current_dir().expect("cwd");
4844 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
4845 adapter.start().await.expect("initialize");
4846 assert!(adapter.next_event().await.is_some());
4847 assert!(adapter.next_event().await.is_some());
4848
4849 for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
4850 adapter.send_prompt(prompt.into()).await.expect("prompt");
4851 assert!(matches!(
4852 adapter.next_event().await,
4853 Some(Ok(AgentEvent::Text { text, .. })) if text == expected
4854 ));
4855 assert!(matches!(
4856 adapter.next_event().await,
4857 Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
4858 ));
4859 }
4860 adapter.stop().await.expect("stop");
4861 }
4862
4863 #[tokio::test]
4864 async fn empty_acp_mode_catalog_disables_mode_control() {
4865 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":[]}}}'"#;
4866 let cwd = std::env::current_dir().expect("cwd");
4867 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4868 adapter.start().await.expect("initialize");
4869 assert!(!adapter.capabilities().supports_modes);
4870 assert!(matches!(
4871 adapter.next_event().await,
4872 Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
4873 ));
4874 assert!(matches!(
4875 adapter.next_event().await,
4876 Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
4877 ));
4878 adapter.stop().await.expect("stop");
4879 }
4880
4881 #[tokio::test]
4882 async fn acp_models_are_discovered_live_and_changed_through_session_config() {
4883 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"#;
4884 let cwd = std::env::current_dir().expect("cwd");
4885 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4886 adapter.start().await.expect("initialize");
4887 assert!(adapter.capabilities().supports_models);
4888 assert!(matches!(
4889 adapter.next_event().await,
4890 Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
4891 if config_id == "model"
4892 && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
4893 && current_model.as_deref() == Some("fast")
4894 ));
4895 assert!(matches!(
4896 adapter.next_event().await,
4897 Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
4898 ));
4899 adapter.set_model("smart".into()).await.expect("set model");
4900 assert!(adapter.set_model("invented".into()).await.is_err());
4901 adapter.stop().await.expect("stop");
4902 }
4903
4904 #[tokio::test]
4905 async fn acp_mode_change_is_acknowledged_without_provider_notification() {
4906 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":{}}'"#;
4907 let cwd = std::env::current_dir().expect("cwd");
4908 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4909 adapter.start().await.expect("initialize");
4910 adapter
4911 .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
4912 .await
4913 .expect("set mode");
4914 assert!(matches!(
4915 adapter.next_event().await,
4916 Some(Ok(AgentEvent::ModesReplaced { .. }))
4917 ));
4918 assert!(matches!(
4919 adapter.next_event().await,
4920 Some(Ok(AgentEvent::Ready { .. }))
4921 ));
4922 assert!(matches!(
4923 adapter.next_event().await,
4924 Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
4925 ));
4926 adapter.stop().await.expect("stop");
4927 }
4928
4929 #[tokio::test]
4930 async fn acp_reload_preserves_a_loadable_session_id() {
4931 let cwd = std::env::current_dir().expect("cwd");
4932 let mut adapter = AcpAdapter::with_session_id(
4933 0,
4934 cwd,
4935 "__codeswarm_missing_acp_for_reload_test__",
4936 Vec::new(),
4937 "saved-session",
4938 );
4939 adapter.capabilities.supports_session_load = true;
4940 assert!(adapter.reload().await.is_err());
4944 assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
4945 }
4946
4947 #[tokio::test]
4948 async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
4949 let cwd = std::env::current_dir().expect("cwd");
4950 let mut adapter = AcpAdapter::with_session_id(
4951 0,
4952 cwd,
4953 "__codeswarm_missing_nonloadable_acp__",
4954 Vec::new(),
4955 "stale-session",
4956 );
4957 adapter.capabilities.supports_session_load = false;
4958 assert!(adapter.reload().await.is_err());
4959 assert_eq!(adapter.session_id, None);
4960 }
4961
4962 #[tokio::test]
4963 async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
4964 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"}}'"#;
4965 let cwd = std::env::current_dir().expect("cwd");
4966 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4967 adapter.start().await.expect("initialize");
4968 assert!(matches!(
4969 adapter.next_event().await,
4970 Some(Ok(AgentEvent::Ready { .. }))
4971 ));
4972 adapter.send_prompt("hello".into()).await.expect("prompt");
4973 assert!(matches!(
4974 adapter.next_event().await,
4975 Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
4976 ));
4977 assert!(matches!(
4978 adapter.next_event().await,
4979 Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
4980 ));
4981 }
4982
4983 #[tokio::test]
4984 async fn acp_load_replays_history_without_starting_a_turn() {
4985 let script = r#"
4986read _
4987echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'
4988read request
4989case "$request" in *session/load*) ;; *) exit 2;; esac
4990echo '{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"text":"old question"}}}}'
4991echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"old answer"}}}}'
4992echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_thought_chunk","content":{"text":"old reasoning"}}}}'
4993echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call","toolCallId":"old-tool","title":"Read","status":"in_progress"}}}'
4994echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"old-tool","title":"Read","status":"completed"}}}'
4995echo '{"jsonrpc":"2.0","id":2,"result":{}}'
4996read request
4997case "$request" in *session/prompt*) ;; *) exit 3;; esac
4998echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"new answer"}}}}'
4999echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
5000"#;
5001 let mut adapter = AcpAdapter::with_session_id(
5002 2,
5003 std::env::current_dir().unwrap(),
5004 "sh",
5005 vec!["-c".into(), script.into()],
5006 "saved",
5007 );
5008 adapter.start().await.unwrap();
5009 let mut state = crate::SessionState::new(3);
5010 for _ in 0..5 {
5011 let event = adapter.next_event().await.unwrap().unwrap();
5012 assert!(matches!(&event, AgentEvent::History { slot: 2, .. }));
5013 crate::reduce(&mut state, event);
5014 assert_eq!(state.active_slot, None);
5015 }
5016 assert!(matches!(
5017 adapter.next_event().await,
5018 Some(Ok(AgentEvent::Ready { slot: 2, .. }))
5019 ));
5020 adapter.send_prompt("new question".into()).await.unwrap();
5021 assert!(
5022 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "new answer")
5023 );
5024 assert!(matches!(
5025 adapter.next_event().await,
5026 Some(Ok(AgentEvent::TurnComplete { slot: 2 }))
5027 ));
5028 adapter.stop().await.unwrap();
5029 }
5030
5031 #[tokio::test]
5032 async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
5033 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
5034 let cwd = std::env::current_dir().expect("cwd");
5035 let mut adapter = AcpAdapter::with_session_id(
5036 0,
5037 cwd,
5038 "sh",
5039 vec!["-c".into(), script.into()],
5040 "existing-session",
5041 );
5042 adapter.start().await.expect("load existing session");
5043 assert!(matches!(
5044 adapter.next_event().await,
5045 Some(Ok(AgentEvent::Ready { .. }))
5046 ));
5047 }
5048
5049 #[tokio::test]
5050 async fn acp_start_failure_reaps_transport_process() {
5051 let mut adapter = AcpAdapter::new(
5056 0,
5057 std::env::current_dir().expect("cwd"),
5058 "sh",
5059 vec!["-c".into(), "printf 'not-json\\n'".into()],
5060 );
5061 assert!(adapter.start().await.is_err());
5062 assert!(adapter.child.is_none());
5063 assert!(adapter.reader.is_none());
5064 }
5065
5066 #[tokio::test]
5067 async fn acp_adapter_answers_permission_json_rpc_requests() {
5068 let path = std::env::temp_dir().join(format!(
5069 "codeswarm-permission-answer-{}",
5070 std::process::id()
5071 ));
5072 let script = format!(
5073 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"}}}}'"#,
5074 path.display()
5075 );
5076 let mut adapter = AcpAdapter::new(
5077 0,
5078 std::env::current_dir().expect("cwd"),
5079 "sh",
5080 vec!["-c".into(), script],
5081 );
5082 adapter.start().await.expect("start ACP");
5083 assert!(matches!(
5084 adapter.next_event().await,
5085 Some(Ok(AgentEvent::Ready { .. }))
5086 ));
5087 adapter.send_prompt("do it".into()).await.expect("prompt");
5088 assert!(matches!(
5089 adapter.next_event().await,
5090 Some(Ok(AgentEvent::Permission { request, .. }))
5091 if request.id == "9"
5092 && request.options == ["Allow once"]
5093 && request.option_ids == ["allow-once"]
5094 ));
5095 adapter
5096 .answer_permission(
5097 "9".into(),
5098 PermissionAnswer::Selected {
5099 option_id: "allow-once".into(),
5100 },
5101 )
5102 .await
5103 .expect("permission answer");
5104 assert!(matches!(
5105 adapter.next_event().await,
5106 Some(Ok(AgentEvent::TurnComplete { .. }))
5107 ));
5108 let answer: Value = serde_json::from_str(
5109 &std::fs::read_to_string(&path).expect("captured permission answer"),
5110 )
5111 .expect("valid JSON-RPC answer");
5112 assert_eq!(answer["id"], 9);
5113 assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
5114 assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
5115 std::fs::remove_file(path).expect("cleanup");
5116 }
5117
5118 #[test]
5119 fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
5120 let event = parse_acp_notification(
5121 0,
5122 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
5123 )
5124 .expect("valid JSON-RPC request");
5125 assert!(event.is_none());
5126 }
5127
5128 #[tokio::test]
5129 async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
5130 let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
5131 std::fs::write(
5132 &script_path,
5133 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
5134 )
5135 .expect("write native test script");
5136 let mut adapter = AgyAdapter::new(
5137 0,
5138 std::env::current_dir().expect("cwd"),
5139 format!("sh {}", script_path.display()),
5140 );
5141 adapter.start().await.expect("start native adapter");
5142 assert!(adapter.next_event().await.is_some());
5143 assert!(adapter.next_event().await.is_some());
5144 adapter
5145 .send_prompt("continue".into())
5146 .await
5147 .expect("prompt");
5148 assert!(matches!(
5149 adapter.next_event().await,
5150 Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
5151 ));
5152 assert!(matches!(
5153 adapter.next_event().await,
5154 Some(Ok(AgentEvent::TurnComplete { .. }))
5155 ));
5156 adapter.stop().await.expect("stop native adapter");
5157 std::fs::remove_file(script_path).expect("cleanup native script");
5158 }
5159
5160 #[test]
5161 fn acp_workspace_file_access_is_root_bound_and_size_limited() {
5162 let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
5163 let _ = std::fs::remove_dir_all(&root);
5164 std::fs::create_dir_all(&root).expect("workspace");
5165 std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
5166 let outside =
5167 std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
5168 std::fs::write(&outside, "secret").expect("outside file");
5169 let link = root.join("outside-link");
5170 #[cfg(unix)]
5171 std::os::unix::fs::symlink(&outside, &link).expect("symlink");
5172 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5173
5174 assert_eq!(
5175 adapter
5176 .read_workspace_text("inside.txt", Some(2), Some(1))
5177 .expect("read inside"),
5178 "two"
5179 );
5180 std::fs::write(
5181 root.join("large.txt"),
5182 vec![b'x'; MAX_FILE_READ_BYTES + 1024],
5183 )
5184 .expect("large file");
5185 let bounded = adapter
5186 .read_workspace_text("large.txt", None, None)
5187 .expect("bounded read");
5188 assert!(bounded.len() <= MAX_FILE_READ_BYTES);
5189 #[cfg(unix)]
5190 {
5191 std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
5192 .expect("internal symlink");
5193 assert_eq!(
5194 adapter
5195 .read_workspace_text("inside-link", None, None)
5196 .expect("read internal symlink"),
5197 "one\ntwo\nthree\n"
5198 );
5199 }
5200 assert!(adapter.workspace_path("../codeswarm-outside").is_err());
5201 assert!(
5202 adapter
5203 .workspace_path(&outside.display().to_string())
5204 .is_err()
5205 );
5206 #[cfg(unix)]
5207 assert!(adapter.workspace_path("outside-link").is_err());
5208 #[cfg(unix)]
5209 std::fs::remove_file(link).expect("cleanup symlink");
5210 #[cfg(unix)]
5211 std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
5212 std::fs::remove_file(outside).expect("cleanup outside");
5213 std::fs::remove_dir_all(root).expect("cleanup workspace");
5214 }
5215
5216 #[tokio::test]
5217 async fn running_terminal_output_omits_exit_status_until_completion() {
5218 let root = unique_test_path("codeswarm-terminal-output", "dir");
5219 std::fs::create_dir_all(&root).expect("workspace");
5220 let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5221 let result = adapter
5222 .terminal_create(&serde_json::json!({
5223 "command": "sh",
5224 "args": ["-c", "sleep 0.2; printf done"],
5225 "cwd": ".",
5226 }))
5227 .await
5228 .expect("terminal create");
5229 let id = result["terminalId"].as_str().expect("terminal id");
5230 let output = adapter.terminal_output(id).await.expect("terminal output");
5231 assert!(output.get("exitStatus").is_none());
5232 if let Some(terminal) = adapter.terminals.remove(id) {
5233 terminal.stop().await;
5234 }
5235 std::fs::remove_dir_all(root).expect("cleanup workspace");
5236 }
5237
5238 #[tokio::test]
5239 async fn acp_adapter_answers_workspace_read_requests() {
5240 let root =
5241 std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
5242 let _ = std::fs::remove_dir_all(&root);
5243 std::fs::create_dir_all(&root).expect("workspace");
5244 let source = root.join("inside.txt");
5245 let answer = root.join("answer.json");
5246 std::fs::write(&source, "workspace content").expect("source");
5247 let script = format!(
5248 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"}}}}'"#,
5249 source.display(),
5250 answer.display(),
5251 );
5252 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
5253 adapter.start().await.expect("start ACP");
5254 assert!(matches!(
5255 adapter.next_event().await,
5256 Some(Ok(AgentEvent::Ready { .. }))
5257 ));
5258 adapter.send_prompt("read it".into()).await.expect("prompt");
5259 assert!(matches!(
5260 adapter.next_event().await,
5261 Some(Ok(AgentEvent::TurnComplete { .. }))
5262 ));
5263 let response: Value =
5264 serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
5265 .expect("response JSON");
5266 assert_eq!(response["id"], 9);
5267 assert_eq!(response["result"]["content"], "workspace content");
5268 adapter.stop().await.expect("stop ACP");
5269 std::fs::remove_dir_all(root).expect("cleanup workspace");
5270 }
5271
5272 #[tokio::test]
5273 async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
5274 let root =
5275 std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
5276 let _ = std::fs::remove_dir_all(&root);
5277 std::fs::create_dir_all(&root).expect("workspace");
5278 let create_request = serde_json::json!({
5279 "jsonrpc": "2.0",
5280 "id": 9,
5281 "method": "terminal/create",
5282 "params": {
5283 "sessionId": "s1",
5284 "command": "sh",
5285 "args": ["-c", "sleep 0.1; printf terminal-ok"],
5286 "cwd": ".",
5287 },
5288 });
5289 let wait_request = serde_json::json!({
5290 "jsonrpc": "2.0",
5291 "id": 10,
5292 "method": "terminal/wait_for_exit",
5293 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
5294 });
5295 let output_request = serde_json::json!({
5296 "jsonrpc": "2.0",
5297 "id": 11,
5298 "method": "terminal/output",
5299 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
5300 });
5301 let create_answer = root.join("create-answer.json");
5302 let wait_answer = root.join("wait-answer.json");
5303 let output_answer = root.join("output-answer.json");
5304 let script = format!(
5305 "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\"}}}}'",
5306 create_request,
5307 create_answer.display(),
5308 wait_request,
5309 wait_answer.display(),
5310 output_request,
5311 output_answer.display(),
5312 );
5313 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
5314 adapter.start().await.expect("start ACP");
5315 assert!(matches!(
5316 adapter.next_event().await,
5317 Some(Ok(AgentEvent::Ready { .. }))
5318 ));
5319 adapter
5320 .send_prompt("run terminal".into())
5321 .await
5322 .expect("prompt");
5323 let mut saw_complete = false;
5324 for _ in 0..6 {
5325 match adapter.next_event().await {
5326 Some(Ok(AgentEvent::TurnComplete { .. })) => {
5327 saw_complete = true;
5328 break;
5329 }
5330 Some(_) => {}
5331 None => break,
5332 }
5333 }
5334 assert!(saw_complete, "terminal requests should not stall ACP");
5335 let create: Value = serde_json::from_str(
5336 &std::fs::read_to_string(&create_answer).expect("captured create response"),
5337 )
5338 .expect("create JSON");
5339 assert_eq!(create["result"]["terminalId"], "terminal-1");
5340 let output: Value = serde_json::from_str(
5341 &std::fs::read_to_string(&output_answer).expect("captured output response"),
5342 )
5343 .expect("output JSON");
5344 assert!(
5345 output["result"]["output"]
5346 .as_str()
5347 .unwrap_or_default()
5348 .contains("terminal-ok"),
5349 "output response: {output}"
5350 );
5351 adapter.stop().await.expect("stop ACP");
5352 std::fs::remove_dir_all(root).expect("cleanup workspace");
5353 }
5354
5355 #[tokio::test]
5356 async fn host_reduces_and_persists_adapter_events() {
5357 let path =
5358 std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
5359 let adapter = ScriptedAdapter::new(
5360 0,
5361 AgentCapabilities::default(),
5362 [AgentEvent::Text {
5363 slot: 0,
5364 text: "hello".into(),
5365 }],
5366 );
5367 let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
5368 host.start().await.expect("start");
5369 host.next_effects()
5370 .await
5371 .expect("event")
5372 .expect("valid event");
5373 assert_eq!(host.state.public_text[0].1, "hello");
5374 assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
5375 std::fs::remove_file(path).expect("cleanup");
5376 }
5377
5378 #[tokio::test]
5379 async fn relay_applies_default_policy_before_the_first_prompt() {
5380 let first_log = Arc::new(Mutex::new(Vec::new()));
5381 let second_log = Arc::new(Mutex::new(Vec::new()));
5382 let first = AdapterHost::new(
5383 Box::new(ModeOrderAdapter {
5384 slot: 0,
5385 log: Arc::clone(&first_log),
5386 phase: 0,
5387 }),
5388 None,
5389 );
5390 let second = AdapterHost::new(
5391 Box::new(ModeOrderAdapter {
5392 slot: 1,
5393 log: Arc::clone(&second_log),
5394 phase: 0,
5395 }),
5396 None,
5397 );
5398 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
5399 relay.start().await.expect("start and synchronize policy");
5400 relay.run_turn("task", 0).await.expect("first turn");
5401 {
5402 let log = first_log.lock().expect("log");
5403 assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
5404 }
5405 assert_eq!(
5406 second_log.lock().expect("log").as_slice(),
5407 ["start", "mode:yolo"]
5408 );
5409 let added_log = Arc::new(Mutex::new(Vec::new()));
5410 relay
5411 .add_agent(
5412 AdapterHost::new(
5413 Box::new(ModeOrderAdapter {
5414 slot: 2,
5415 log: Arc::clone(&added_log),
5416 phase: 0,
5417 }),
5418 None,
5419 ),
5420 "Added",
5421 "added.example",
5422 "added-agent",
5423 )
5424 .await
5425 .expect("add with synchronized policy");
5426 assert_eq!(
5427 added_log.lock().expect("log").as_slice(),
5428 ["start", "mode:yolo"]
5429 );
5430 relay.drop_agent(2).await.expect("drop added agent");
5431 added_log.lock().expect("log").clear();
5432 relay
5433 .reload(2)
5434 .await
5435 .expect("reload with synchronized policy");
5436 assert_eq!(
5437 added_log.lock().expect("log").as_slice(),
5438 ["reload", "mode:yolo"]
5439 );
5440 }
5441
5442 #[tokio::test]
5443 async fn acp_roster_is_ready_before_any_prompt_is_sent() {
5444 let hosts = (0..2)
5445 .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
5446 .collect::<Vec<_>>();
5447 let startup_events = Arc::new(Mutex::new(Vec::new()));
5448 let captured = Arc::clone(&startup_events);
5449 let mut relay = RelayHost::new(hosts, 4).expect("relay");
5450 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5451
5452 relay.start().await.expect("complete startup handshake");
5453
5454 assert!(relay.dispatches().is_empty());
5455 let ready_slots = startup_events
5456 .lock()
5457 .expect("events")
5458 .iter()
5459 .filter_map(|event| match event {
5460 AgentEvent::Ready { slot, .. } => Some(*slot),
5461 _ => None,
5462 })
5463 .collect::<Vec<_>>();
5464 assert_eq!(ready_slots, vec![0, 1]);
5465 }
5466
5467 #[tokio::test]
5468 async fn independent_roster_adapters_start_concurrently() {
5469 let barrier = Arc::new(tokio::sync::Barrier::new(2));
5470 let hosts = (0..2)
5471 .map(|slot| {
5472 AdapterHost::new(
5473 Box::new(ConcurrentStartAdapter {
5474 slot,
5475 barrier: Arc::clone(&barrier),
5476 }),
5477 None,
5478 )
5479 })
5480 .collect::<Vec<_>>();
5481 let mut relay = RelayHost::new(hosts, 4).expect("relay");
5482 tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
5483 .await
5484 .expect("startup should not serialize barrier participants")
5485 .expect("startup succeeds");
5486 }
5487
5488 #[tokio::test]
5489 async fn relay_host_dispatches_turns_sequentially() {
5490 let capabilities = AgentCapabilities {
5491 supports_cancel: true,
5492 ..AgentCapabilities::default()
5493 };
5494 let first = ScriptedAdapter::new(
5495 0,
5496 capabilities.clone(),
5497 [
5498 AgentEvent::Text {
5499 slot: 0,
5500 text: "first".into(),
5501 },
5502 AgentEvent::TurnComplete { slot: 0 },
5503 ],
5504 );
5505 let second = ScriptedAdapter::new(
5506 1,
5507 capabilities,
5508 [
5509 AgentEvent::Text {
5510 slot: 1,
5511 text: "review".into(),
5512 },
5513 AgentEvent::TurnComplete { slot: 1 },
5514 ],
5515 );
5516 let hosts = vec![
5517 AdapterHost::new(Box::new(first), None),
5518 AdapterHost::new(Box::new(second), None),
5519 ];
5520 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5521 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
5522 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5523 let captured = std::sync::Arc::clone(&events);
5524 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5525 relay.start().await.expect("start");
5526 events.lock().expect("events").clear();
5527 assert!(matches!(
5528 relay.run_turn("task", 0).await.expect("first turn"),
5529 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5530 ));
5531 assert!(matches!(
5532 relay.run_turn("first", 0).await.expect("second turn"),
5533 crate::relay::RelayDecision::Dispatch {
5534 slot: 1,
5535 can_stop: true,
5536 ..
5537 }
5538 ));
5539 assert_eq!(
5540 relay
5541 .dispatches()
5542 .iter()
5543 .map(|(slot, _)| *slot)
5544 .collect::<Vec<_>>(),
5545 [0, 1]
5546 );
5547 assert!(relay.dispatches()[0].1.contains("You are Claude"));
5548 assert!(
5549 relay.dispatches()[0]
5550 .1
5551 .contains("CodeSwarm roster (ordered)")
5552 );
5553 assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
5554 assert!(relay.dispatches()[0].1.contains("2. Codex"));
5555 assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
5556 assert!(relay.dispatches()[0].1.contains("Do not use"));
5557 let lifecycle = events.lock().expect("events");
5558 let positions = lifecycle
5559 .iter()
5560 .filter_map(|event| match event {
5561 AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
5562 AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
5563 _ => None,
5564 })
5565 .collect::<Vec<_>>();
5566 assert_eq!(
5567 positions,
5568 [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
5569 );
5570 }
5571
5572 #[tokio::test]
5573 async fn relay_host_routes_around_a_usage_limited_agent() {
5574 let capabilities = AgentCapabilities::default();
5575 let first = ScriptedAdapter::new(
5576 0,
5577 capabilities.clone(),
5578 [
5579 AgentEvent::Text {
5580 slot: 0,
5581 text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
5582 credits or try again later."
5583 .into(),
5584 },
5585 AgentEvent::TurnComplete { slot: 0 },
5586 ],
5587 );
5588 let second = ScriptedAdapter::new(
5589 1,
5590 capabilities,
5591 [
5592 AgentEvent::Text {
5593 slot: 1,
5594 text: "review done".into(),
5595 },
5596 AgentEvent::TurnComplete { slot: 1 },
5597 ],
5598 );
5599 let hosts = vec![
5600 AdapterHost::new(Box::new(first), None),
5601 AdapterHost::new(Box::new(second), None),
5602 ];
5603 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5604 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5605 let captured = std::sync::Arc::clone(&events);
5606 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5607 relay.start().await.expect("start");
5608 events.lock().expect("events").clear();
5609 assert!(matches!(
5610 relay.run_turn("task", 0).await.expect("limited turn"),
5611 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5612 ));
5613 assert!(
5614 events
5615 .lock()
5616 .expect("events")
5617 .iter()
5618 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
5619 );
5620 assert!(matches!(
5622 relay.run_turn("", 0).await.expect("next turn"),
5623 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
5624 ));
5625 assert!(relay.relay().is_limited(0));
5626 relay.reload(0).await.expect("reload");
5630 assert!(!relay.relay().is_limited(0));
5631 }
5632
5633 #[tokio::test]
5634 async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
5635 let limited = ScriptedAdapter::new(
5636 0,
5637 AgentCapabilities::default(),
5638 [AgentEvent::Failed {
5639 slot: 0,
5640 started: true,
5641 detail: "request failed: insufficient_quota".into(),
5642 }],
5643 );
5644 let healthy = ScriptedAdapter::new(
5645 1,
5646 AgentCapabilities::default(),
5647 [AgentEvent::TurnComplete { slot: 1 }],
5648 );
5649 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5650 let captured = std::sync::Arc::clone(&events);
5651 let mut relay = RelayHost::new(
5652 vec![
5653 AdapterHost::new(Box::new(limited), None),
5654 AdapterHost::new(Box::new(healthy), None),
5655 ],
5656 4,
5657 )
5658 .expect("relay");
5659 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5660 relay.start().await.expect("start");
5661 events.lock().expect("events").clear();
5662
5663 assert!(matches!(
5664 relay.run_turn("task", 0).await.expect("limited failure"),
5665 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5666 ));
5667 assert!(relay.relay().is_limited(0));
5668 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
5669 {
5670 let events = events.lock().expect("events");
5671 assert!(
5672 events
5673 .iter()
5674 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
5675 );
5676 assert!(
5677 !events
5678 .iter()
5679 .any(|event| matches!(event, AgentEvent::Failed { .. }))
5680 );
5681 }
5682
5683 assert!(matches!(
5684 relay.run_turn("", 0).await.expect("healthy peer"),
5685 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
5686 ));
5687 }
5688
5689 #[tokio::test]
5690 async fn relay_failure_is_tombstoned_and_reported_to_the_ui_sink() {
5691 let failed = ScriptedAdapter::new(
5692 0,
5693 AgentCapabilities::default(),
5694 [AgentEvent::Failed {
5695 slot: 0,
5696 started: true,
5697 detail: "connection lost".into(),
5698 }],
5699 );
5700 let healthy = ScriptedAdapter::new(
5701 1,
5702 AgentCapabilities::default(),
5703 [AgentEvent::TurnComplete { slot: 1 }],
5704 );
5705 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5706 let captured = std::sync::Arc::clone(&events);
5707 let mut relay = RelayHost::new(
5708 vec![
5709 AdapterHost::new(Box::new(failed), None),
5710 AdapterHost::new(Box::new(healthy), None),
5711 ],
5712 4,
5713 )
5714 .expect("relay");
5715 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5716 relay.start().await.expect("start");
5717
5718 let error = relay.run_turn("task", 0).await.expect_err("failure");
5719 assert!(error.to_string().contains("connection lost"));
5720 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
5721 assert!(events.lock().expect("lock").iter().any(|event| {
5722 matches!(
5723 event,
5724 AgentEvent::Failed {
5725 slot: 0,
5726 started: true,
5727 ..
5728 }
5729 )
5730 }));
5731 }
5732
5733 #[tokio::test]
5734 async fn codex_stop_does_not_skip_later_roster_reviewers() {
5735 let hosts = (0..3)
5736 .map(|slot| {
5737 AdapterHost::new(
5738 Box::new(ScriptedAdapter::new(
5739 slot,
5740 AgentCapabilities::default(),
5741 [
5742 AgentEvent::Text {
5743 slot,
5744 text: STOP_TOKEN.into(),
5745 },
5746 AgentEvent::TurnComplete { slot },
5747 ],
5748 )),
5749 None,
5750 )
5751 })
5752 .collect();
5753 let mut relay = RelayHost::new(hosts, 10).expect("relay");
5754 relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
5755 relay.start().await.expect("start");
5756 for expected in 0..3 {
5757 assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
5758 RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
5759 }
5760 assert_eq!(
5761 relay.run_turn("", 0).await.expect("complete"),
5762 RelayDecision::Complete
5763 );
5764 }
5765
5766 #[tokio::test]
5767 async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
5768 let first = ScriptedAdapter::new(
5769 0,
5770 AgentCapabilities::default(),
5771 [
5772 AgentEvent::Text {
5773 slot: 0,
5774 text: "done".into(),
5775 },
5776 AgentEvent::TurnComplete { slot: 0 },
5777 ],
5778 );
5779 let reviewer = ScriptedAdapter::new(
5780 1,
5781 AgentCapabilities::default(),
5782 [
5783 AgentEvent::Text {
5784 slot: 1,
5785 text: STOP_TOKEN.into(),
5786 },
5787 AgentEvent::TurnComplete { slot: 1 },
5788 ],
5789 );
5790 let mut relay = RelayHost::new(
5791 vec![
5792 AdapterHost::new(Box::new(first), None),
5793 AdapterHost::new(Box::new(reviewer), None),
5794 ],
5795 10,
5796 )
5797 .expect("relay");
5798 relay.start().await.expect("start");
5799 let first_decision = relay.run_turn("task", 0).await.expect("first");
5800 assert!(matches!(
5801 first_decision,
5802 RelayDecision::Dispatch { slot: 0, .. }
5803 ));
5804 let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
5805 assert!(matches!(
5806 reviewer_decision,
5807 RelayDecision::Dispatch {
5808 slot: 1,
5809 can_stop: true,
5810 ..
5811 }
5812 ));
5813 assert_eq!(
5814 relay.run_turn("", 0).await.expect("complete"),
5815 RelayDecision::Complete
5816 );
5817 }
5818
5819 #[tokio::test]
5820 async fn relay_stream_emits_text_and_thought_endings_before_tools() {
5821 let tool = AgentEvent::Tool {
5822 slot: 0,
5823 update: crate::ToolUpdate {
5824 id: "read".into(),
5825 title: "Read file".into(),
5826 status: ToolStatus::Running,
5827 detail: None,
5828 },
5829 };
5830 let updates = vec![
5831 AgentEvent::Thought {
5832 slot: 0,
5833 text: "Check the buffer. ✈".into(),
5834 },
5835 AgentEvent::Text {
5836 slot: 0,
5837 text: "Let me check.".into(),
5838 },
5839 tool.clone(),
5840 AgentEvent::Text {
5841 slot: 0,
5842 text: "[CODE".into(),
5843 },
5844 AgentEvent::Text {
5845 slot: 0,
5846 text: " is ordinary.".into(),
5847 },
5848 AgentEvent::Text {
5849 slot: 0,
5850 text: "[CODESWARM:".into(),
5851 },
5852 AgentEvent::Text {
5853 slot: 0,
5854 text: "STOP] Done.".into(),
5855 },
5856 tool,
5857 AgentEvent::TurnComplete { slot: 0 },
5858 ];
5859 let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
5860 let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
5861 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5862 let captured = std::sync::Arc::clone(&events);
5863 let mut relay = RelayHost::new(
5864 vec![
5865 AdapterHost::new(Box::new(first), None),
5866 AdapterHost::new(Box::new(reviewer), None),
5867 ],
5868 2,
5869 )
5870 .expect("relay");
5871 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
5872 relay.start().await.unwrap();
5873 relay.run_turn("task", 0).await.unwrap();
5874 let captured = events.lock().unwrap();
5875 let visible: Vec<_> = captured
5876 .iter()
5877 .filter(|event| {
5878 matches!(
5879 event,
5880 AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
5881 )
5882 })
5883 .cloned()
5884 .collect();
5885 assert_eq!(
5886 visible,
5887 vec![
5888 updates[0].clone(),
5889 updates[1].clone(),
5890 updates[2].clone(),
5891 AgentEvent::Text {
5892 slot: 0,
5893 text: "[CODE is ordinary.".into()
5894 },
5895 AgentEvent::Text {
5896 slot: 0,
5897 text: " Done.".into()
5898 },
5899 updates[7].clone(),
5900 ]
5901 );
5902 }
5903
5904 #[tokio::test]
5905 async fn stop_token_is_filtered_from_streamed_ui_events() {
5906 let first = ScriptedAdapter::new(
5907 0,
5908 AgentCapabilities::default(),
5909 [
5910 AgentEvent::Text {
5911 slot: 0,
5912 text: format!("visible {STOP_TOKEN} trailing"),
5913 },
5914 AgentEvent::TurnComplete { slot: 0 },
5915 ],
5916 );
5917 let reviewer = ScriptedAdapter::new(
5918 1,
5919 AgentCapabilities::default(),
5920 [AgentEvent::TurnComplete { slot: 1 }],
5921 );
5922 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5923 let captured = std::sync::Arc::clone(&events);
5924 let mut relay = RelayHost::new(
5925 vec![
5926 AdapterHost::new(Box::new(first), None),
5927 AdapterHost::new(Box::new(reviewer), None),
5928 ],
5929 2,
5930 )
5931 .expect("relay");
5932 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5933 relay.start().await.expect("start");
5934 relay.run_turn("task", 0).await.expect("turn");
5935 let captured = events.lock().expect("lock");
5936 assert!(captured.iter().all(|event| match event {
5937 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5938 _ => true,
5939 }));
5940 let visible = captured
5941 .iter()
5942 .filter_map(|event| match event {
5943 AgentEvent::Text { text, .. } => Some(text.as_str()),
5944 _ => None,
5945 })
5946 .collect::<String>();
5947 assert_eq!(visible, "visible trailing");
5948 }
5949
5950 #[tokio::test]
5951 async fn token_only_reviewer_response_emits_visible_acknowledgment() {
5952 let first = ScriptedAdapter::new(
5953 0,
5954 AgentCapabilities::default(),
5955 [
5956 AgentEvent::Text {
5957 slot: 0,
5958 text: "done".into(),
5959 },
5960 AgentEvent::TurnComplete { slot: 0 },
5961 ],
5962 );
5963 let reviewer = ScriptedAdapter::new(
5964 1,
5965 AgentCapabilities::default(),
5966 [
5967 AgentEvent::Text {
5968 slot: 1,
5969 text: STOP_TOKEN.into(),
5970 },
5971 AgentEvent::TurnComplete { slot: 1 },
5972 ],
5973 );
5974 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5975 let captured = std::sync::Arc::clone(&events);
5976 let mut relay = RelayHost::new(
5977 vec![
5978 AdapterHost::new(Box::new(first), None),
5979 AdapterHost::new(Box::new(reviewer), None),
5980 ],
5981 4,
5982 )
5983 .expect("relay");
5984 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5985 relay.start().await.expect("start");
5986 relay.run_turn("task", 0).await.expect("first turn");
5987 relay.run_turn("", 0).await.expect("review turn");
5988 let captured = events.lock().expect("lock");
5989 assert!(captured.iter().any(|event| {
5990 matches!(
5991 event,
5992 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
5993 )
5994 }));
5995 assert!(captured.iter().all(|event| match event {
5996 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5997 _ => true,
5998 }));
5999 let acknowledgment = captured
6000 .iter()
6001 .position(|event| {
6002 matches!(
6003 event,
6004 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
6005 )
6006 })
6007 .expect("visible acknowledgment");
6008 let completion = captured
6009 .iter()
6010 .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
6011 .expect("reviewer completion");
6012 assert!(acknowledgment < completion);
6013 }
6014
6015 #[tokio::test]
6016 async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
6017 let first = ScriptedAdapter::new(
6018 0,
6019 AgentCapabilities::default(),
6020 [
6021 AgentEvent::Text {
6022 slot: 0,
6023 text: "done".into(),
6024 },
6025 AgentEvent::TurnComplete { slot: 0 },
6026 ],
6027 );
6028 let reviewer = ScriptedAdapter::new(
6029 1,
6030 AgentCapabilities::default(),
6031 [
6032 AgentEvent::Text {
6033 slot: 1,
6034 text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
6035 },
6036 AgentEvent::TurnComplete { slot: 1 },
6037 ],
6038 );
6039 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6040 let captured = std::sync::Arc::clone(&events);
6041 let mut relay = RelayHost::new(
6042 vec![
6043 AdapterHost::new(Box::new(first), None),
6044 AdapterHost::new(Box::new(reviewer), None),
6045 ],
6046 4,
6047 )
6048 .expect("relay");
6049 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6050 relay.start().await.expect("start");
6051 relay.run_turn("task", 0).await.expect("first turn");
6052 relay.run_turn("", 0).await.expect("review turn");
6053
6054 let visible = events
6055 .lock()
6056 .expect("lock")
6057 .iter()
6058 .filter_map(|event| match event {
6059 AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
6060 _ => None,
6061 })
6062 .collect::<String>();
6063 assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
6064 assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
6065 }
6066
6067 #[tokio::test]
6068 async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
6069 let first = AdapterHost::new(
6070 Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
6071 None,
6072 );
6073 let second = AdapterHost::new(
6074 Box::new(ScriptedAdapter::new(
6075 1,
6076 AgentCapabilities::default(),
6077 [AgentEvent::TurnComplete { slot: 1 }],
6078 )),
6079 None,
6080 );
6081 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6082 let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
6083 relay.set_event_sink(move |event| {
6084 if matches!(event, AgentEvent::Permission { .. }) {
6085 let _ = seen_sender.send(());
6086 }
6087 });
6088 relay.start().await.expect("start");
6089 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
6090 let answer = async move {
6091 seen_receiver.recv().await.expect("permission request");
6092 sender
6093 .send(super::RelayPermissionAnswer {
6094 slot: 0,
6095 request_id: "permission-1".into(),
6096 answer: PermissionAnswer::Selected {
6097 option_id: "allow".into(),
6098 },
6099 })
6100 .expect("queue permission answer");
6101 };
6102 tokio::time::timeout(std::time::Duration::from_millis(100), async {
6103 let ((), result) = tokio::join!(
6104 answer,
6105 relay.run_turn_with_permissions("task", 0, &mut receiver)
6106 );
6107 result
6108 })
6109 .await
6110 .expect("permission-gated turn should not deadlock")
6111 .expect("turn completes");
6112 }
6113
6114 #[tokio::test]
6115 async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
6116 let first = AdapterHost::new(
6117 Box::new(PendingAdapter {
6118 slot: 0,
6119 hang_on_cancel: false,
6120 }),
6121 None,
6122 );
6123 let second = AdapterHost::new(
6124 Box::new(ScriptedAdapter::new(
6125 1,
6126 AgentCapabilities::default(),
6127 [AgentEvent::TurnComplete { slot: 1 }],
6128 )),
6129 None,
6130 );
6131 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6132 relay.start().await.expect("start");
6133 let cancellation = relay.cancellation();
6134 let error = {
6135 let turn = relay.run_turn("task", 0);
6136 tokio::pin!(turn);
6137 cancellation.request();
6138 turn.await.expect_err("cancellation should stop turn")
6139 };
6140 assert!(error.to_string().contains("relay turn cancelled"));
6141
6142 assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
6143 relay
6144 .run_turn("", 1)
6145 .await
6146 .expect("replacement job reaches the selected peer");
6147 let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
6148 assert!(replacement.contains("replacement job"));
6149 assert!(replacement.contains("User "));
6150 assert!(replacement.contains(":\ntask"));
6151 let owner_updates = relay.relay_mut().unseen_context(0);
6152 assert!(owner_updates.contains("User "));
6153 assert!(owner_updates.contains(":\ntask"));
6154 assert!(owner_updates.contains(":\nreplacement job"));
6155 }
6156
6157 #[tokio::test]
6158 async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
6159 let first = AdapterHost::new(
6160 Box::new(PendingAdapter {
6161 slot: 0,
6162 hang_on_cancel: true,
6163 }),
6164 None,
6165 );
6166 let second = AdapterHost::new(
6167 Box::new(PendingAdapter {
6168 slot: 1,
6169 hang_on_cancel: false,
6170 }),
6171 None,
6172 );
6173 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6174 relay.start().await.expect("start");
6175 let cancellation = relay.cancellation();
6176 let turn = relay.run_turn("task", 0);
6177 tokio::pin!(turn);
6178 cancellation.request();
6179 let error = turn.await.expect_err("cancellation should stop turn");
6180 assert!(error.to_string().contains("timed out"));
6181 }
6182
6183 #[tokio::test]
6184 async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
6185 let event = [AgentEvent::TurnComplete { slot: 0 }];
6186 let first = AdapterHost::new(
6187 Box::new(ScriptedAdapter::new(
6188 0,
6189 AgentCapabilities::default(),
6190 event.clone(),
6191 )),
6192 None,
6193 );
6194 let second = AdapterHost::new(
6195 Box::new(ScriptedAdapter::new(
6196 1,
6197 AgentCapabilities::default(),
6198 [AgentEvent::TurnComplete { slot: 1 }],
6199 )),
6200 None,
6201 );
6202 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6203 relay.start().await.expect("start");
6204
6205 relay.pause();
6206 assert_eq!(
6207 relay.run_turn("paused", 0).await.expect("paused turn"),
6208 crate::relay::RelayDecision::Paused
6209 );
6210 assert!(relay.dispatches().is_empty());
6211
6212 relay.resume();
6213 relay.relay_mut().drop_agent(1).expect("drop reviewer");
6214 assert!(matches!(
6215 relay
6216 .run_turn("solo follow-up", 0)
6217 .await
6218 .expect("solo turn"),
6219 crate::relay::RelayDecision::Dispatch {
6220 slot: 0,
6221 can_stop: false,
6222 ..
6223 }
6224 ));
6225 assert_eq!(relay.dispatches().len(), 1);
6226 }
6227
6228 #[tokio::test]
6229 async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
6230 let first = AdapterHost::new(
6231 Box::new(ScriptedAdapter::new(
6232 0,
6233 AgentCapabilities::default(),
6234 [AgentEvent::TurnComplete { slot: 0 }],
6235 )),
6236 None,
6237 );
6238 let second = AdapterHost::new(
6239 Box::new(ScriptedAdapter::new(
6240 1,
6241 AgentCapabilities::default(),
6242 [AgentEvent::TurnComplete { slot: 1 }],
6243 )),
6244 None,
6245 );
6246 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6247 relay.set_roster_names(vec!["First".into(), "Second".into()]);
6248 relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
6249 relay.set_roster_launch_specs(vec![
6250 ("custom".into(), "owner".into()),
6251 ("custom".into(), "peer".into()),
6252 ]);
6253 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6254 let captured = std::sync::Arc::clone(&events);
6255 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6256 relay.start().await.expect("start");
6257 let slot = relay
6258 .add_agent(
6259 AdapterHost::new(
6260 Box::new(ScriptedAdapter::new(
6261 2,
6262 AgentCapabilities::default(),
6263 [AgentEvent::TurnComplete { slot: 2 }],
6264 )),
6265 None,
6266 ),
6267 "Reviewer",
6268 "reviewer.example",
6269 "reviewer --acp",
6270 )
6271 .await
6272 .expect("append agent");
6273 assert_eq!(slot, 2);
6274 assert_eq!(
6275 relay.relay().active_slots().collect::<Vec<_>>(),
6276 vec![0, 1, 2]
6277 );
6278 assert_eq!(
6279 relay
6280 .session_metadata()
6281 .get("agents")
6282 .and_then(|value| value.as_array())
6283 .map(Vec::len),
6284 Some(3)
6285 );
6286 relay.drop_agent(1).await.expect("drop middle peer");
6287 let metadata = relay.session_metadata();
6288 assert_eq!(
6289 metadata.get("agents"),
6290 Some(&serde_json::json!([
6291 {"name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
6292 {"name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
6293 ]))
6294 );
6295 assert!(
6296 events
6297 .lock()
6298 .expect("lock")
6299 .iter()
6300 .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
6301 );
6302 }
6303
6304 #[tokio::test]
6305 async fn relay_host_persists_coordinator_owned_runtime_metadata() {
6306 let path = unique_test_path("codeswarm-session-metadata", "json");
6307 let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
6308 let writer = metadata_store.buffered().expect("metadata writer");
6309 let first = AdapterHost::new(
6310 Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
6311 None,
6312 );
6313 let second = AdapterHost::new(
6314 Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
6315 None,
6316 );
6317 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6318 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6319 relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
6320 relay.set_roster_launch_specs(vec![
6321 ("custom".into(), "claude".into()),
6322 ("custom".into(), "codex".into()),
6323 ]);
6324 relay.set_session_metadata_writer(writer);
6325 relay.start().await.expect("start");
6326 relay.drop_agent(0).await.expect("drop first agent");
6327 relay.stop().await.expect("stop");
6328
6329 let loaded = metadata_store
6330 .read()
6331 .expect("read metadata")
6332 .expect("metadata snapshot");
6333 assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
6334 assert_eq!(
6335 loaded.get("agents"),
6336 Some(&serde_json::json!([{
6337 "name": "Codex", "identity": "openai.com", "protocol": "custom",
6338 "command": "codex", "supports_load_session": false
6339 }]))
6340 );
6341 assert!(loaded.get("owner").is_none());
6342 let _ = std::fs::remove_file(path);
6343 }
6344
6345 #[tokio::test]
6346 async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
6347 let first = AdapterHost::new(
6348 Box::new(ScriptedAdapter::new(
6349 0,
6350 AgentCapabilities::default(),
6351 [
6352 AgentEvent::Text {
6353 slot: 0,
6354 text: "owner stream".into(),
6355 },
6356 AgentEvent::TurnComplete { slot: 0 },
6357 ],
6358 )),
6359 None,
6360 );
6361 let second = AdapterHost::new(
6362 Box::new(ScriptedAdapter::new(
6363 1,
6364 AgentCapabilities::default(),
6365 [
6366 AgentEvent::Text {
6367 slot: 1,
6368 text: "peer stream".into(),
6369 },
6370 AgentEvent::TurnComplete { slot: 1 },
6371 ],
6372 )),
6373 None,
6374 );
6375 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6376 relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
6377 relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
6378 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6379 let captured = std::sync::Arc::clone(&events);
6380 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6381 relay.start().await.expect("start");
6382
6383 relay.swap_agents(0, 1).expect("swap peers");
6384 assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
6385 assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
6386 relay.run_turn("task", 0).await.expect("swapped turn");
6387 let events = events.lock().expect("events");
6388 assert!(events.iter().any(|event| {
6389 matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
6390 }));
6391 assert!(relay.dispatches()[0].1.contains("You are Peer"));
6392 }
6393
6394 #[tokio::test]
6395 async fn relay_host_persists_all_active_agent_metadata_off_thread() {
6396 let path = unique_test_path("codeswarm-session-metadata", "json");
6397 let first = AdapterHost::new(
6398 Box::new(ScriptedAdapter::new(
6399 0,
6400 AgentCapabilities::default(),
6401 [AgentEvent::TurnComplete { slot: 0 }],
6402 )),
6403 None,
6404 );
6405 let second = AdapterHost::new(
6406 Box::new(ScriptedAdapter::new(
6407 1,
6408 AgentCapabilities::default(),
6409 [AgentEvent::TurnComplete { slot: 1 }],
6410 )),
6411 None,
6412 );
6413 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6414 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6415 relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
6416 relay.set_roster_launch_specs(vec![
6417 ("custom".into(), "claude".into()),
6418 ("custom".into(), "codex".into()),
6419 ]);
6420 let writer = SessionMetadataStore::open(&path)
6421 .buffered()
6422 .expect("metadata writer");
6423 relay.set_session_metadata_writer(writer);
6424 relay.start().await.expect("start");
6425 relay.stop().await.expect("stop");
6426 let loaded = SessionMetadataStore::open(&path)
6427 .read()
6428 .expect("read metadata")
6429 .expect("metadata snapshot");
6430 let agents = loaded
6431 .get("agents")
6432 .and_then(|value| value.as_array())
6433 .expect("agents");
6434 assert_eq!(agents.len(), 2);
6435 assert_eq!(agents[0]["identity"], "claude.com");
6436 assert_eq!(agents[1]["identity"], "openai.com");
6437 let _ = std::fs::remove_file(path);
6438 }
6439
6440 #[tokio::test]
6441 async fn relay_host_routes_unseen_public_context_to_next_agent() {
6442 let first = AdapterHost::new(
6443 Box::new(ScriptedAdapter::new(
6444 0,
6445 AgentCapabilities::default(),
6446 [
6447 AgentEvent::Text {
6448 slot: 0,
6449 text: "implemented the fix".into(),
6450 },
6451 AgentEvent::TurnComplete { slot: 0 },
6452 ],
6453 )),
6454 None,
6455 );
6456 let second = AdapterHost::new(
6457 Box::new(ScriptedAdapter::new(
6458 1,
6459 AgentCapabilities::default(),
6460 [AgentEvent::TurnComplete { slot: 1 }],
6461 )),
6462 None,
6463 );
6464 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6465 relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
6466 relay.start().await.expect("start");
6467 relay.run_turn("task", 0).await.expect("first turn");
6468 relay.run_turn("review this", 0).await.expect("review turn");
6469
6470 assert_eq!(relay.dispatches().len(), 2);
6471 assert_eq!(relay.dispatches()[0].0, 0);
6472 assert!(relay.dispatches()[0].1.contains("task"));
6473 assert!(relay.dispatches()[0].1.contains("You are Codex"));
6474 assert!(relay.dispatches()[0].1.contains("2. Qwen"));
6475 assert_eq!(relay.dispatches()[1].0, 1);
6476 assert!(relay.dispatches()[1].1.contains("review this"));
6477 let public = relay.dispatches()[1]
6478 .1
6479 .split_once("Public updates:\n")
6480 .map(|(_, updates)| updates)
6481 .expect("review receives public context");
6482 let header = public
6483 .lines()
6484 .find(|line| line.starts_with("Codex "))
6485 .expect("named previous agent");
6486 let timestamp = header
6487 .strip_prefix("Codex ")
6488 .and_then(|value| value.strip_suffix(':'))
6489 .expect("timestamped header");
6490 assert_eq!(timestamp.len(), 5);
6491 assert_eq!(timestamp.as_bytes()[2], b':');
6492 assert!(
6493 timestamp
6494 .bytes()
6495 .enumerate()
6496 .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
6497 );
6498 assert!(public.contains("implemented the fix"));
6499 assert!(!public.contains("Agent 0"));
6500 }
6501}