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 control_token_visible_end, is_usage_limit_response, requested_next_slot,
24 strip_control_tokens, strip_stop_token,
25 },
26 resources,
27};
28use async_trait::async_trait;
29use base64::{Engine, engine::general_purpose::STANDARD as BASE64};
30use serde_json::Value;
31use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWriteExt, BufReader};
32use tokio::process::{Child, ChildStdout, Command};
33use tokio::sync::{Mutex as AsyncMutex, Notify, mpsc};
34
35#[path = "codex.rs"]
36mod codex;
37#[path = "native.rs"]
38mod native;
39pub use codex::CodexAdapter;
40#[path = "claude.rs"]
41mod claude;
42pub use claude::ClaudeAdapter;
43
44pub type AdapterResult<T> = Result<T, AdapterError>;
45
46const MAX_ACP_LINE_BYTES: usize = 10 * 1024 * 1024;
50const MAX_FILE_READ_BYTES: usize = 4 * 1024 * 1024;
51const MAX_TERMINAL_OUTPUT_BYTES: usize = 1024 * 1024;
52
53#[derive(Clone, Debug)]
54struct TerminalProcess {
55 child: Arc<AsyncMutex<Option<Child>>>,
56 output: Arc<Mutex<Vec<u8>>>,
57 truncated: Arc<AtomicBool>,
58 output_readers: Arc<AtomicUsize>,
59}
60
61impl TerminalProcess {
62 async fn kill(&self) {
63 if let Some(child) = self.child.lock().await.as_mut() {
64 #[cfg(unix)]
65 if signal_isolated_process_group(child, nix::sys::signal::Signal::SIGTERM) {
66 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
67 signal_isolated_process_group(child, nix::sys::signal::Signal::SIGKILL);
68 }
69 let _ = child.start_kill();
70 }
71 }
72
73 async fn stop(&self) {
74 if let Some(mut child) = self.child.lock().await.take() {
75 let _ = terminate_child(&mut child).await;
76 }
77 }
78
79 async fn wait(&self) -> Option<i32> {
80 loop {
81 let code = {
82 let mut child = self.child.lock().await;
83 match child.as_mut() {
84 None => Some(-1),
85 Some(child) => match child.try_wait() {
86 Ok(Some(status)) => Some(status.code().unwrap_or(-1)),
87 Ok(None) => None,
88 Err(_) => Some(-1),
93 },
94 }
95 };
96 if code.is_some() {
97 while self.output_readers.load(Ordering::Acquire) != 0 {
98 tokio::task::yield_now().await;
99 }
100 return code;
101 }
102 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
103 }
104 }
105
106 async fn exit_code(&self) -> Option<i32> {
107 let mut child = self.child.lock().await;
108 child
109 .as_mut()
110 .and_then(|child| child.try_wait().ok().flatten())
111 .map(|status| status.code().unwrap_or(-1))
112 }
113}
114
115#[derive(Clone, Debug)]
116pub struct HostUpdate {
117 pub event: AgentEvent,
118 pub effects: Vec<Effect>,
119}
120
121#[derive(Clone, Debug, Eq, PartialEq)]
122pub enum AdapterError {
123 Unsupported(&'static str),
124 Spawn(String),
125 Transport(String),
126 Protocol(String),
127}
128
129impl std::fmt::Display for AdapterError {
130 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
131 match self {
132 Self::Unsupported(operation) => write!(formatter, "unsupported operation: {operation}"),
133 Self::Spawn(error) => write!(formatter, "unable to launch agent: {error}"),
134 Self::Transport(error) => write!(formatter, "agent transport error: {error}"),
135 Self::Protocol(error) => write!(formatter, "agent protocol error: {error}"),
136 }
137 }
138}
139
140impl std::error::Error for AdapterError {}
141
142fn floor_char_boundary(text: &str, index: usize) -> usize {
143 let mut index = index.min(text.len());
144 while index > 0 && !text.is_char_boundary(index) {
145 index -= 1;
146 }
147 index
148}
149
150#[derive(Clone, Debug, Eq, PartialEq)]
158pub enum CommandParseError {
159 Empty,
160 UnterminatedQuote,
161 TrailingEscape,
162}
163
164impl std::fmt::Display for CommandParseError {
165 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
166 let message = match self {
167 Self::Empty => "command is empty",
168 Self::UnterminatedQuote => "command contains an unterminated quote",
169 Self::TrailingEscape => "command ends with an incomplete escape",
170 };
171 formatter.write_str(message)
172 }
173}
174
175impl std::error::Error for CommandParseError {}
176
177pub fn parse_command_line(command: &str) -> Result<(String, Vec<String>), CommandParseError> {
180 let mut argv = Vec::new();
181 let mut argument = String::new();
182 let mut quoted = None;
183 let mut escaped = false;
184 let mut started = false;
185
186 for character in command.chars() {
187 if escaped {
188 argument.push(character);
189 escaped = false;
190 started = true;
191 continue;
192 }
193 match (quoted, character) {
194 (_, '\\') if quoted != Some('\'') => {
195 escaped = true;
196 started = true;
197 }
198 (None, '\'' | '"') => {
199 quoted = Some(character);
200 started = true;
201 }
202 (Some(quote), character) if character == quote => quoted = None,
203 (None, character) if character.is_whitespace() => {
204 if started {
205 argv.push(std::mem::take(&mut argument));
206 started = false;
207 }
208 }
209 (_, character) => {
210 argument.push(character);
211 started = true;
212 }
213 }
214 }
215
216 if escaped {
217 return Err(CommandParseError::TrailingEscape);
218 }
219 if quoted.is_some() {
220 return Err(CommandParseError::UnterminatedQuote);
221 }
222 if started {
223 argv.push(argument);
224 }
225 let Some((program, args)) = argv.split_first() else {
226 return Err(CommandParseError::Empty);
227 };
228 Ok((program.clone(), args.to_vec()))
229}
230
231async fn terminate_child(child: &mut Child) -> AdapterResult<()> {
235 #[cfg(unix)]
236 if signal_isolated_process_group(child, nix::sys::signal::Signal::SIGTERM) {
237 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
238 signal_isolated_process_group(child, nix::sys::signal::Signal::SIGKILL);
239 }
240 let kill_error = child.start_kill().err();
241 let wait_error = child.wait().await.err();
242 if let Some(error) = kill_error.or(wait_error) {
243 return Err(AdapterError::Transport(error.to_string()));
244 }
245 Ok(())
246}
247
248#[cfg(unix)]
249fn signal_isolated_process_group(child: &Child, signal: nix::sys::signal::Signal) -> bool {
250 use nix::{
251 sys::signal::killpg,
252 unistd::{Pid, getpgid, getpgrp},
253 };
254
255 let Some(raw_pid) = child.id().and_then(|pid| i32::try_from(pid).ok()) else {
256 return false;
257 };
258 let pid = Pid::from_raw(raw_pid);
259 if getpgid(Some(pid)).ok() == Some(pid) && pid != getpgrp() {
264 let _ = killpg(pid, signal);
265 true
266 } else {
267 false
268 }
269}
270
271fn isolate_process_group(command: &mut Command) {
272 #[cfg(unix)]
273 command.process_group(0);
274}
275
276async fn drain_bounded<R>(mut reader: R, limit: usize) -> String
279where
280 R: AsyncRead + Unpin,
281{
282 let mut bytes = Vec::new();
283 let mut chunk = [0_u8; 4096];
284 while let Ok(count) = reader.read(&mut chunk).await {
285 if count == 0 {
286 break;
287 }
288 bytes.extend_from_slice(&chunk[..count]);
289 if bytes.len() > limit {
290 let keep_from = bytes.len() - limit;
291 bytes.drain(..keep_from);
292 }
293 }
294 String::from_utf8_lossy(&bytes).trim().to_owned()
295}
296
297async fn read_bounded_line<R>(reader: &mut R) -> AdapterResult<String>
302where
303 R: AsyncBufRead + Unpin,
304{
305 let mut bytes = Vec::with_capacity(4096);
306 loop {
307 let buffer = reader
308 .fill_buf()
309 .await
310 .map_err(|error| AdapterError::Transport(error.to_string()))?;
311 if buffer.is_empty() {
312 if bytes.is_empty() {
313 return Err(AdapterError::Transport("ACP stream closed".into()));
314 }
315 break;
316 }
317 let newline = buffer.iter().position(|byte| *byte == b'\n');
318 let available = newline.map_or(buffer.len(), |index| index + 1);
319 let remaining = MAX_ACP_LINE_BYTES
320 .saturating_add(1)
321 .saturating_sub(bytes.len());
322 if available > remaining {
323 reader.consume(remaining);
324 return Err(AdapterError::Protocol(format!(
325 "ACP protocol line exceeds {MAX_ACP_LINE_BYTES} bytes"
326 )));
327 }
328 bytes.extend_from_slice(&buffer[..available]);
329 reader.consume(available);
330 if newline.is_some() {
331 break;
332 }
333 }
334 if bytes.len() > MAX_ACP_LINE_BYTES {
335 return Err(AdapterError::Protocol(format!(
336 "ACP protocol line exceeds {MAX_ACP_LINE_BYTES} bytes"
337 )));
338 }
339 String::from_utf8(bytes).map_err(|error| AdapterError::Protocol(error.to_string()))
340}
341
342async fn drain_terminal_output<R>(
343 mut reader: R,
344 output: Arc<Mutex<Vec<u8>>>,
345 truncated: Arc<AtomicBool>,
346 output_readers: Arc<AtomicUsize>,
347 limit: usize,
348) where
349 R: AsyncRead + Unpin,
350{
351 let mut chunk = [0_u8; 4096];
352 while let Ok(count) = reader.read(&mut chunk).await {
353 if count == 0 {
354 break;
355 }
356 if let Ok(mut bytes) = output.lock() {
357 let remaining = limit.saturating_sub(bytes.len());
358 if count > remaining {
359 bytes.extend_from_slice(&chunk[..remaining]);
360 truncated.store(true, Ordering::Release);
361 } else {
362 bytes.extend_from_slice(&chunk[..count]);
363 }
364 }
365 }
366 output_readers.fetch_sub(1, Ordering::AcqRel);
367}
368
369#[async_trait]
371pub trait AgentAdapter: Send {
372 fn slot(&self) -> RosterSlot;
373 fn display_name(&self) -> String {
377 format!("Agent {}", self.slot().saturating_add(1))
378 }
379 fn session_id(&self) -> Option<String> {
383 None
384 }
385 fn protocol(&self) -> &'static str {
389 "custom"
390 }
391 fn needs_restart(&self) -> bool {
396 false
397 }
398 fn capabilities(&self) -> AgentCapabilities;
399 async fn start(&mut self) -> AdapterResult<()>;
400 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()>;
401 async fn cancel(&mut self) -> AdapterResult<bool>;
402 async fn answer_permission(
403 &mut self,
404 request_id: String,
405 answer: PermissionAnswer,
406 ) -> AdapterResult<()>;
407 async fn set_mode(&mut self, mode: String) -> AdapterResult<()>;
408 async fn set_model(&mut self, _model: String) -> AdapterResult<()> {
409 Err(AdapterError::Unsupported("set_model"))
410 }
411 async fn reload(&mut self) -> AdapterResult<()>;
412 async fn stop(&mut self) -> AdapterResult<()>;
413 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>>;
414}
415
416struct SlotMappedAdapter {
421 logical_slot: RosterSlot,
422 inner: Box<dyn AgentAdapter>,
423}
424
425impl std::fmt::Debug for SlotMappedAdapter {
426 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
427 formatter
428 .debug_struct("SlotMappedAdapter")
429 .field("logical_slot", &self.logical_slot)
430 .field("inner_slot", &self.inner.slot())
431 .finish_non_exhaustive()
432 }
433}
434
435fn restored_history_event(event: AgentEvent) -> AgentEvent {
436 use crate::HistoryContent;
437 let (slot, content) = match event {
438 AgentEvent::UserText { slot, text } => (slot, HistoryContent::UserText(text)),
439 AgentEvent::Text { slot, text } => (slot, HistoryContent::Text(text)),
440 AgentEvent::Thought { slot, text } => (slot, HistoryContent::Thought(text)),
441 AgentEvent::Tool { slot, update } => (slot, HistoryContent::Tool(update)),
442 other => return other,
443 };
444 AgentEvent::History { slot, content }
445}
446
447fn map_event_slot(event: AgentEvent, slot: RosterSlot) -> AgentEvent {
448 match event {
449 AgentEvent::SessionMetadataUpdated { metadata } => {
450 AgentEvent::SessionMetadataUpdated { metadata }
451 }
452 AgentEvent::History { content, .. } => AgentEvent::History { slot, content },
453 AgentEvent::GoalUpdated { goal } => AgentEvent::GoalUpdated { goal },
454 AgentEvent::RosterUpdated { update } => AgentEvent::RosterUpdated { update },
455 AgentEvent::Ready { capabilities, .. } => AgentEvent::Ready { slot, capabilities },
456 AgentEvent::TurnStarted { .. } => AgentEvent::TurnStarted { slot },
457 AgentEvent::ModesReplaced {
458 modes,
459 current_mode,
460 ..
461 } => AgentEvent::ModesReplaced {
462 slot,
463 modes,
464 current_mode,
465 },
466 AgentEvent::ModeUpdated { current_mode, .. } => {
467 AgentEvent::ModeUpdated { slot, current_mode }
468 }
469 AgentEvent::ModelsReplaced {
470 config_id,
471 models,
472 current_model,
473 ..
474 } => AgentEvent::ModelsReplaced {
475 slot,
476 config_id,
477 models,
478 current_model,
479 },
480 AgentEvent::ModelUpdated { current_model, .. } => AgentEvent::ModelUpdated {
481 slot,
482 current_model,
483 },
484 AgentEvent::UserText { text, .. } => AgentEvent::UserText { slot, text },
485 AgentEvent::CommandsReplaced { commands, .. } => {
486 AgentEvent::CommandsReplaced { slot, commands }
487 }
488 AgentEvent::UsageUpdated { usage, .. } => AgentEvent::UsageUpdated { slot, usage },
489 AgentEvent::Text { text, .. } => AgentEvent::Text { slot, text },
490 AgentEvent::Thought { text, .. } => AgentEvent::Thought { slot, text },
491 AgentEvent::Tool { update, .. } => AgentEvent::Tool { slot, update },
492 AgentEvent::Permission { request, .. } => AgentEvent::Permission { slot, request },
493 AgentEvent::Terminal { event, .. } => AgentEvent::Terminal { slot, event },
494 AgentEvent::TurnComplete { .. } => AgentEvent::TurnComplete { slot },
495 AgentEvent::BatchComplete { elapsed } => AgentEvent::BatchComplete { elapsed },
496 AgentEvent::UsageLimitReached { detail, .. } => {
497 AgentEvent::UsageLimitReached { slot, detail }
498 }
499 AgentEvent::Failed {
500 started, detail, ..
501 } => AgentEvent::Failed {
502 slot,
503 started,
504 detail,
505 },
506 }
507}
508
509#[async_trait]
510impl AgentAdapter for SlotMappedAdapter {
511 fn slot(&self) -> RosterSlot {
512 self.logical_slot
513 }
514
515 fn display_name(&self) -> String {
516 self.inner.display_name()
517 }
518
519 fn session_id(&self) -> Option<String> {
520 self.inner.session_id()
521 }
522
523 fn protocol(&self) -> &'static str {
524 self.inner.protocol()
525 }
526
527 fn capabilities(&self) -> AgentCapabilities {
528 self.inner.capabilities()
529 }
530
531 fn needs_restart(&self) -> bool {
532 self.inner.needs_restart()
533 }
534
535 async fn start(&mut self) -> AdapterResult<()> {
536 self.inner.start().await
537 }
538
539 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
540 self.inner.send_prompt(prompt).await
541 }
542
543 async fn cancel(&mut self) -> AdapterResult<bool> {
544 self.inner.cancel().await
545 }
546
547 async fn answer_permission(
548 &mut self,
549 request_id: String,
550 answer: PermissionAnswer,
551 ) -> AdapterResult<()> {
552 self.inner.answer_permission(request_id, answer).await
553 }
554
555 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
556 self.inner.set_mode(mode).await
557 }
558
559 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
560 self.inner.set_model(model).await
561 }
562
563 async fn reload(&mut self) -> AdapterResult<()> {
564 self.inner.reload().await
565 }
566
567 async fn stop(&mut self) -> AdapterResult<()> {
568 self.inner.stop().await
569 }
570
571 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
572 let slot = self.logical_slot;
573 self.inner
574 .next_event()
575 .await
576 .map(|result| result.map(|event| map_event_slot(event, slot)))
577 }
578}
579
580pub struct AdapterHost {
583 adapter: Box<dyn AgentAdapter>,
584 pub state: SessionState,
585 pub last_error: Option<String>,
586 event_log: Option<EventLog>,
587}
588
589impl std::fmt::Debug for AdapterHost {
590 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
591 formatter
592 .debug_struct("AdapterHost")
593 .field("state", &self.state)
594 .field("last_error", &self.last_error)
595 .field("event_log", &self.event_log)
596 .finish_non_exhaustive()
597 }
598}
599
600impl AdapterHost {
601 pub fn new(adapter: Box<dyn AgentAdapter>, event_log: Option<EventLog>) -> Self {
602 let slot = adapter.slot();
603 Self {
604 adapter,
605 state: SessionState::new(slot.saturating_add(1)),
606 last_error: None,
607 event_log,
608 }
609 }
610
611 pub async fn start(&mut self) -> AdapterResult<()> {
612 self.adapter.start().await
613 }
614
615 pub async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
616 self.adapter.send_prompt(prompt).await
617 }
618
619 pub async fn cancel(&mut self) -> AdapterResult<bool> {
620 self.adapter.cancel().await
621 }
622
623 pub async fn answer_permission(
624 &mut self,
625 request_id: String,
626 answer: PermissionAnswer,
627 ) -> AdapterResult<()> {
628 self.adapter.answer_permission(request_id, answer).await
629 }
630
631 pub async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
632 self.adapter.set_mode(mode).await
633 }
634
635 pub async fn set_model(&mut self, model: String) -> AdapterResult<()> {
636 self.adapter.set_model(model).await
637 }
638
639 pub async fn reload(&mut self) -> AdapterResult<()> {
640 self.adapter.reload().await?;
641 let slot = self.adapter.slot();
642 if let Some(agent) = self.state.slots.get_mut(slot) {
643 agent.active = true;
644 agent.capabilities = self.adapter.capabilities();
645 }
646 self.last_error = None;
647 Ok(())
648 }
649
650 pub async fn stop(&mut self) -> AdapterResult<()> {
651 self.adapter.stop().await
652 }
653
654 pub async fn next_effects(&mut self) -> Option<AdapterResult<Vec<Effect>>> {
655 Some(self.next_update().await?.map(|update| update.effects))
656 }
657
658 pub async fn next_update(&mut self) -> Option<AdapterResult<HostUpdate>> {
659 let event = match self.adapter.next_event().await {
660 None => return None,
661 Some(Err(error)) => {
662 self.last_error = Some(error.to_string());
663 let slot = self.adapter.slot();
664 let failure = AgentEvent::Failed {
665 slot,
666 started: true,
667 detail: error.to_string(),
668 };
669 let effects = reduce(&mut self.state, failure.clone());
670 return Some(Ok(HostUpdate {
671 event: failure,
672 effects,
673 }));
674 }
675 Some(Ok(event)) => event,
676 };
677 if let Some(log) = &self.event_log
678 && let Err(error) = log.append(&event)
679 {
680 return Some(Err(AdapterError::Transport(error.to_string())));
681 }
682 let effects = reduce(&mut self.state, event.clone());
683 Some(Ok(HostUpdate { event, effects }))
684 }
685
686 pub fn adapter(&self) -> &dyn AgentAdapter {
687 &*self.adapter
688 }
689
690 pub fn session_id(&self) -> Option<String> {
691 self.adapter.session_id()
692 }
693
694 fn remap(self, logical_slot: RosterSlot) -> Self {
697 let old_slot = self.adapter.slot();
698 if old_slot == logical_slot {
699 return self;
700 }
701
702 let mut state = self.state;
703 if state.slots.len() <= logical_slot {
704 state
705 .slots
706 .resize(logical_slot.saturating_add(1), Default::default());
707 }
708 if let Some(agent) = state.slots.get(old_slot).cloned() {
709 state.slots[logical_slot] = agent;
710 }
711 if state.active_slot == Some(old_slot) {
712 state.active_slot = Some(logical_slot);
713 }
714 for (slot, _) in &mut state.queued_prompts {
715 if *slot == old_slot {
716 *slot = logical_slot;
717 }
718 }
719 for (slot, _) in &mut state.public_text {
720 if *slot == old_slot {
721 *slot = logical_slot;
722 }
723 }
724
725 Self {
726 adapter: Box::new(SlotMappedAdapter {
727 logical_slot,
728 inner: self.adapter,
729 }),
730 state,
731 last_error: self.last_error,
732 event_log: self.event_log,
733 }
734 }
735}
736
737pub struct RelayHost {
740 goal: Option<crate::goal::Goal>,
741 hosts: Vec<AdapterHost>,
742 relay: Relay,
743 introduced: Vec<bool>,
744 roster_names: Vec<String>,
745 roster_identities: Vec<String>,
746 roster_launch_specs: Vec<(String, String)>,
747 desired_policy: String,
748 metadata_writer: Option<BufferedSessionMetadataStore>,
749 metadata_workspace: Option<String>,
750 dispatches: Vec<(RosterSlot, String)>,
751 last_public_dispatch: Option<RosterSlot>,
752 pair_implementer: Option<RosterSlot>,
753 event_sink: Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
754 cancel_requested: Arc<AtomicBool>,
755 active_turn_slot: Arc<AtomicUsize>,
756 cancel_notify: Arc<Notify>,
757}
758
759#[derive(Clone, Debug)]
762pub struct RelayCancellation {
763 requested: Arc<AtomicBool>,
764 active_turn_slot: Arc<AtomicUsize>,
765 notify: Arc<Notify>,
766}
767
768const NO_ACTIVE_TURN: usize = usize::MAX;
769
770struct ActiveTurnGuard(Arc<AtomicUsize>);
771
772impl Drop for ActiveTurnGuard {
773 fn drop(&mut self) {
774 self.0.store(NO_ACTIVE_TURN, Ordering::Release);
775 }
776}
777
778#[derive(Debug)]
783pub struct RelayPermissionAnswer {
784 pub slot: RosterSlot,
785 pub request_id: String,
786 pub answer: PermissionAnswer,
787}
788
789#[cfg(test)]
790const CANCEL_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(250);
791#[cfg(not(test))]
792const CANCEL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
793
794#[cfg(test)]
795const CANCEL_SETTLE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(20);
796#[cfg(not(test))]
797const CANCEL_SETTLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
798
799async fn cancel_with_timeout(host: &mut AdapterHost) -> AdapterResult<bool> {
804 tokio::time::timeout(CANCEL_TIMEOUT, host.cancel())
805 .await
806 .map_err(|_| AdapterError::Transport("adapter cancellation timed out".into()))?
807}
808
809fn canonical_policy_id(policy: &str) -> &str {
810 match policy {
811 "plan" => "codeswarm:mode:plan",
812 "default" | "manual" => "codeswarm:mode:manual",
813 "accept-edits" => "codeswarm:mode:accept-edits",
814 "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
815 other => other,
816 }
817}
818
819async fn apply_policy_to_host(host: &mut AdapterHost, policy: &str) -> AdapterResult<()> {
820 if !host.adapter().capabilities().supports_modes {
821 return Ok(());
822 }
823 let policy_id = canonical_policy_id(policy);
824 let slot = host.adapter().slot();
825 let advertised = host
826 .state
827 .slots
828 .get(slot)
829 .map(|agent| agent.modes.as_slice())
830 .unwrap_or_default();
831 let native = if advertised.is_empty() {
832 match policy_id {
833 "codeswarm:mode:plan" => "plan".into(),
834 "codeswarm:mode:manual" => "default".into(),
835 "codeswarm:mode:accept-edits" => "accept-edits".into(),
836 "codeswarm:mode:full-access" => "full-access".into(),
837 other => other.into(),
838 }
839 } else {
840 crate::policy::resolve(policy_id, advertised)
841 .map(|mode| mode.id)
842 .ok_or(AdapterError::Unsupported(
843 "desired policy is unavailable for adapter",
844 ))?
845 };
846 host.set_mode(native).await
847}
848
849async fn refresh_mode_catalog(
850 host: &mut AdapterHost,
851 event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
852) -> AdapterResult<bool> {
853 if !host.adapter().capabilities().supports_modes {
854 return Ok(false);
855 }
856 tokio::time::timeout(std::time::Duration::from_secs(2), async {
857 let mut ready_seen = false;
858 loop {
859 let update = host.next_update().await.ok_or_else(|| {
860 AdapterError::Transport("adapter ended before advertising modes".into())
861 })??;
862 let catalog_ready = matches!(update.event, AgentEvent::ModesReplaced { .. });
863 ready_seen |= matches!(update.event, AgentEvent::Ready { .. });
864 if let Some(sink) = event_sink {
865 sink(update.event);
866 }
867 if catalog_ready {
868 return Ok(ready_seen);
869 }
870 }
871 })
872 .await
873 .map_err(|_| AdapterError::Transport("adapter mode catalog timed out".into()))?
874}
875
876async fn refresh_adapter_startup(
877 host: &mut AdapterHost,
878 event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
879) -> AdapterResult<()> {
880 let ready_seen = refresh_mode_catalog(host, event_sink).await?;
881 if ready_seen || !matches!(host.adapter().protocol(), "native" | "acp") {
882 return Ok(());
883 }
884 tokio::time::timeout(std::time::Duration::from_secs(2), async {
885 loop {
886 let update = host.next_update().await.ok_or_else(|| {
887 AdapterError::Transport("adapter ended before becoming ready".into())
888 })??;
889 let ready = matches!(update.event, AgentEvent::Ready { .. });
890 if let Some(sink) = event_sink {
891 sink(update.event);
892 }
893 if ready {
894 return Ok(());
895 }
896 }
897 })
898 .await
899 .map_err(|_| AdapterError::Transport("adapter ready handshake timed out".into()))?
900}
901
902fn public_context_speaker(name: &str) -> String {
903 let now = time::OffsetDateTime::now_local().unwrap_or_else(|_| time::OffsetDateTime::now_utc());
904 format!("{name} {:02}:{:02}", now.hour(), now.minute())
905}
906
907impl RelayCancellation {
908 pub fn request(&self) {
909 self.requested.store(true, Ordering::Release);
910 self.notify.notify_one();
911 }
912
913 pub fn request_if_active(&self, slot: RosterSlot) -> bool {
917 if self.active_turn_slot.load(Ordering::Acquire) != slot {
918 return false;
919 }
920 self.request();
921 true
922 }
923}
924
925impl std::fmt::Debug for RelayHost {
926 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
927 formatter
928 .debug_struct("RelayHost")
929 .field("hosts", &self.hosts)
930 .field("relay", &self.relay)
931 .field("dispatches", &self.dispatches)
932 .field("event_sink", &self.event_sink.is_some())
933 .field(
934 "cancel_requested",
935 &self.cancel_requested.load(Ordering::Acquire),
936 )
937 .finish()
938 }
939}
940
941impl RelayHost {
942 pub fn new(hosts: Vec<AdapterHost>, max_rounds: usize) -> Result<Self, AdapterError> {
943 if hosts.is_empty() {
944 return Err(AdapterError::Unsupported("relay requires an adapter"));
945 }
946 Ok(Self {
947 relay: Relay::new(hosts.len(), max_rounds),
948 introduced: vec![false; hosts.len()],
949 roster_names: hosts
950 .iter()
951 .map(|host| host.adapter().display_name())
952 .collect(),
953 roster_identities: hosts
954 .iter()
955 .map(|host| host.adapter().display_name())
956 .collect(),
957 roster_launch_specs: Vec::new(),
958 desired_policy: crate::policy::DEFAULT_POLICY_ID.into(),
959 metadata_writer: None,
960 metadata_workspace: None,
961 hosts,
962 dispatches: Vec::new(),
963 last_public_dispatch: None,
964 pair_implementer: None,
965 goal: None,
966 event_sink: None,
967 cancel_requested: Arc::new(AtomicBool::new(false)),
968 active_turn_slot: Arc::new(AtomicUsize::new(NO_ACTIVE_TURN)),
969 cancel_notify: Arc::new(Notify::new()),
970 })
971 }
972
973 pub fn set_event_sink<F>(&mut self, sink: F)
977 where
978 F: Fn(AgentEvent) + Send + Sync + 'static,
979 {
980 self.event_sink = Some(Arc::new(sink));
981 }
982
983 pub fn set_roster_names(&mut self, names: Vec<String>) {
986 if names.len() == self.hosts.len() {
987 self.roster_names = names;
988 }
989 }
990
991 pub fn set_roster_identities(&mut self, identities: Vec<String>) {
995 if identities.len() == self.hosts.len() {
996 self.roster_identities = identities;
997 }
998 }
999
1000 pub fn set_roster_launch_specs(&mut self, specs: Vec<(String, String)>) {
1001 if specs.len() == self.hosts.len() {
1002 self.roster_launch_specs = specs;
1003 }
1004 }
1005
1006 pub fn set_session_metadata_writer(&mut self, writer: BufferedSessionMetadataStore) {
1009 self.metadata_writer = Some(writer);
1010 }
1011
1012 pub fn set_session_metadata_workspace(&mut self, workspace: impl Into<String>) {
1013 self.metadata_workspace = Some(workspace.into());
1014 }
1015
1016 pub fn session_metadata(&self) -> SessionMetadata {
1018 let active = self.relay.active_slots().collect::<Vec<_>>();
1019 let mut data = serde_json::Map::new();
1020 data.insert(
1021 "goal".into(),
1022 serde_json::to_value(&self.goal).expect("goal serializes"),
1023 );
1024 if let Some(workspace) = &self.metadata_workspace {
1025 data.insert("cwd".into(), serde_json::Value::String(workspace.clone()));
1026 }
1027 data.insert(
1028 "title".into(),
1029 serde_json::Value::String("CodeSwarm".into()),
1030 );
1031 data.insert(
1032 "agents".into(),
1033 serde_json::Value::Array(
1034 active
1035 .into_iter()
1036 .filter_map(|slot| {
1037 let host = self.hosts.get(slot)?;
1038 let (protocol, command) = self.roster_launch_specs.get(slot)?;
1039 let mut agent = serde_json::Map::new();
1040 agent.insert("slot".into(), serde_json::json!(slot));
1041 agent.insert(
1042 "name".into(),
1043 serde_json::Value::String(
1044 self.roster_names
1045 .get(slot)
1046 .cloned()
1047 .unwrap_or_else(|| host.adapter().display_name()),
1048 ),
1049 );
1050 agent.insert(
1051 "identity".into(),
1052 serde_json::Value::String(
1053 self.roster_identities
1054 .get(slot)
1055 .cloned()
1056 .unwrap_or_else(|| host.adapter().display_name()),
1057 ),
1058 );
1059 agent.insert(
1060 "protocol".into(),
1061 serde_json::Value::String(protocol.clone()),
1062 );
1063 agent.insert("command".into(), serde_json::Value::String(command.clone()));
1064 agent.insert(
1065 "supports_load_session".into(),
1066 serde_json::Value::Bool(
1067 host.adapter().capabilities().supports_session_load,
1068 ),
1069 );
1070 if let Some(session_id) = host.session_id() {
1071 agent
1072 .insert("session_id".into(), serde_json::Value::String(session_id));
1073 }
1074 Some(serde_json::Value::Object(agent))
1075 })
1076 .collect(),
1077 ),
1078 );
1079 SessionMetadata::new(data)
1080 }
1081
1082 fn queue_session_metadata(&self) -> AdapterResult<()> {
1083 let metadata = self.session_metadata();
1084 if let Some(sink) = &self.event_sink {
1085 sink(AgentEvent::SessionMetadataUpdated {
1086 metadata: metadata.to_value(),
1087 });
1088 }
1089 if let Some(writer) = &self.metadata_writer {
1090 writer
1091 .write(metadata)
1092 .map_err(|error| AdapterError::Transport(error.to_string()))?;
1093 }
1094 Ok(())
1095 }
1096
1097 pub fn restore_goal(&mut self, goal: Option<crate::goal::Goal>) {
1098 self.goal = goal;
1099 if let Some(sink) = &self.event_sink {
1100 sink(AgentEvent::GoalUpdated {
1101 goal: self.goal.clone(),
1102 });
1103 }
1104 }
1105
1106 pub fn apply_goal(
1107 &mut self,
1108 command: crate::goal::GoalCommand,
1109 ) -> Result<Option<String>, String> {
1110 let task = crate::goal::apply(&mut self.goal, command)?;
1111 if let Some(sink) = &self.event_sink {
1112 sink(AgentEvent::GoalUpdated {
1113 goal: self.goal.clone(),
1114 });
1115 }
1116 self.queue_session_metadata()
1117 .map_err(|error| error.to_string())?;
1118 Ok(task)
1119 }
1120
1121 pub fn roster_names(&self) -> &[String] {
1122 &self.roster_names
1123 }
1124
1125 pub fn session_ids(&self) -> Vec<Option<String>> {
1126 self.hosts.iter().map(AdapterHost::session_id).collect()
1127 }
1128
1129 pub async fn start(&mut self) -> AdapterResult<()> {
1130 self.start_isolating_failures("agent could not start").await
1131 }
1132
1133 pub async fn start_resuming(&mut self) -> AdapterResult<()> {
1136 self.start_isolating_failures("saved session could not be restored")
1137 .await
1138 }
1139
1140 async fn start_isolating_failures(&mut self, failure_prefix: &str) -> AdapterResult<()> {
1141 let event_sink = self.event_sink.clone();
1142 let policy = self.desired_policy.clone();
1143 let starts = self.hosts.iter_mut().map(|host| {
1144 let sink = event_sink.clone();
1145 let policy = policy.clone();
1146 async move {
1147 host.start().await?;
1148 refresh_adapter_startup(host, &sink).await?;
1149 apply_policy_to_host(host, &policy).await
1150 }
1151 });
1152 let results = futures::future::join_all(starts).await;
1153 let mut first_error = None;
1154 for (slot, result) in results.into_iter().enumerate() {
1155 if let Err(error) = result {
1156 let _ = self.relay.tombstone(slot);
1157 let _ = self.hosts[slot].stop().await;
1158 if let Some(sink) = &event_sink {
1159 sink(AgentEvent::Failed {
1160 slot,
1161 started: false,
1162 detail: format!("{failure_prefix}: {error}"),
1163 });
1164 }
1165 if first_error.is_none() {
1166 first_error = Some(error);
1167 }
1168 }
1169 }
1170 if self.relay.active_slots().next().is_none() {
1171 return Err(first_error
1172 .unwrap_or_else(|| AdapterError::Transport("no agents could be started".into())));
1173 }
1174 let _ = self.queue_session_metadata();
1175 Ok(())
1176 }
1177
1178 pub async fn stop(&mut self) -> AdapterResult<()> {
1179 let mut first_error = None;
1184 let active = self.relay.active_slots().collect::<Vec<_>>();
1185 for slot in active {
1186 let host = &mut self.hosts[slot];
1187 if let Err(error) = host.stop().await
1188 && first_error.is_none()
1189 {
1190 first_error = Some(error);
1191 }
1192 }
1193 if let Some(writer) = &self.metadata_writer
1194 && let Err(error) = writer.flush()
1195 && first_error.is_none()
1196 {
1197 first_error = Some(AdapterError::Transport(error.to_string()));
1198 }
1199 first_error.map_or(Ok(()), Err)
1200 }
1201
1202 pub async fn answer_permission(
1205 &mut self,
1206 slot: RosterSlot,
1207 request_id: String,
1208 answer: PermissionAnswer,
1209 ) -> AdapterResult<()> {
1210 let host = self
1211 .hosts
1212 .get_mut(slot)
1213 .ok_or_else(|| AdapterError::Transport("permission target is missing".into()))?;
1214 host.answer_permission(request_id, answer).await
1215 }
1216
1217 pub fn pause(&mut self) {
1218 self.relay.pause();
1219 }
1220
1221 pub fn resume(&mut self) {
1222 self.relay.resume();
1223 }
1224
1225 pub fn set_strategy(&mut self, strategy: CollaborationStrategy) {
1228 self.relay.set_strategy(strategy);
1229 if strategy != CollaborationStrategy::Pair {
1230 self.pair_implementer = None;
1231 }
1232 let _ = self.queue_session_metadata();
1233 }
1234
1235 pub fn strategy(&self) -> CollaborationStrategy {
1236 self.relay.strategy()
1237 }
1238
1239 pub fn roster_identity(&self, slot: RosterSlot) -> Option<&str> {
1240 self.roster_identities.get(slot).map(String::as_str)
1241 }
1242
1243 pub fn active_slot_for_identity(&self, identity: &str) -> Option<RosterSlot> {
1244 self.relay.active_slots().find(|slot| {
1245 self.roster_identity(*slot)
1246 .is_some_and(|candidate| candidate.eq_ignore_ascii_case(identity))
1247 })
1248 }
1249
1250 pub async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
1253 let active = self.relay.active_slots().collect::<Vec<_>>();
1254 for slot in active {
1255 let Some(host) = self.hosts.get_mut(slot) else {
1256 continue;
1257 };
1258 if host.adapter().capabilities().supports_modes {
1259 host.set_mode(mode.clone()).await?;
1260 }
1261 }
1262 Ok(())
1263 }
1264
1265 pub async fn set_model(&mut self, slot: RosterSlot, model: String) -> AdapterResult<()> {
1266 let host = self
1267 .hosts
1268 .get_mut(slot)
1269 .ok_or_else(|| AdapterError::Transport("model target is missing".into()))?;
1270 if !host.adapter().capabilities().supports_models {
1271 return Err(AdapterError::Unsupported("set_model"));
1272 }
1273 host.set_model(model.clone()).await?;
1274 if let Some(sink) = &self.event_sink {
1275 sink(AgentEvent::ModelUpdated {
1276 slot,
1277 current_model: model,
1278 });
1279 }
1280 Ok(())
1281 }
1282
1283 pub async fn set_policy(&mut self, policy: String) -> AdapterResult<()> {
1287 let desired_policy = canonical_policy_id(&policy).to_owned();
1288 let active = self.relay.active_slots().collect::<Vec<_>>();
1292 for active_slot in &active {
1293 let Some(host) = self.hosts.get(*active_slot) else {
1294 continue;
1295 };
1296 if !host.adapter().capabilities().supports_modes {
1297 continue;
1298 }
1299 let advertised = host
1300 .state
1301 .slots
1302 .get(*active_slot)
1303 .map(|agent| agent.modes.as_slice())
1304 .unwrap_or_default();
1305 if !advertised.is_empty()
1306 && crate::policy::resolve(&desired_policy, advertised).is_none()
1307 {
1308 return Err(AdapterError::Unsupported(
1309 "desired policy is unavailable for an active adapter",
1310 ));
1311 }
1312 }
1313 self.desired_policy = desired_policy.clone();
1314 let mut first_error = None;
1315 for active_slot in active {
1316 let Some(host) = self.hosts.get_mut(active_slot) else {
1317 continue;
1318 };
1319 if let Err(error) = apply_policy_to_host(host, &desired_policy).await
1320 && first_error.is_none()
1321 {
1322 first_error = Some(error);
1323 }
1324 }
1325 first_error.map_or(Ok(()), Err)
1326 }
1327
1328 pub async fn reload(&mut self, slot: RosterSlot) -> AdapterResult<()> {
1329 let desired_policy = self.desired_policy.clone();
1330 let _ = self.queue_session_metadata();
1331 let event_sink = self.event_sink.clone();
1332 let host = self
1333 .hosts
1334 .get_mut(slot)
1335 .ok_or_else(|| AdapterError::Transport("reload target is missing".into()))?;
1336 host.reload().await?;
1337 refresh_adapter_startup(host, &event_sink).await?;
1338 apply_policy_to_host(host, &desired_policy).await?;
1339 if let Some(introduced) = self.introduced.get_mut(slot) {
1340 *introduced = false;
1341 }
1342 self.relay
1343 .reactivate(slot)
1344 .map_err(|error| AdapterError::Transport(error.into()))?;
1345 let _ = self.queue_session_metadata();
1346 if let Some(sink) = &self.event_sink {
1347 sink(AgentEvent::RosterUpdated {
1348 update: RosterUpdate::Reloaded { slot },
1349 });
1350 }
1351 let _ = self.relay.clear_limited(slot);
1352 Ok(())
1353 }
1354
1355 pub async fn drop_agent(&mut self, slot: RosterSlot) -> AdapterResult<()> {
1357 self.relay
1358 .drop_agent(slot)
1359 .map_err(|error| AdapterError::Transport(error.into()))?;
1360 let _stop_result = if let Some(host) = self.hosts.get_mut(slot) {
1361 host.stop().await
1362 } else {
1363 Ok(())
1364 };
1365 let _ = self.queue_session_metadata();
1369 if let Some(sink) = &self.event_sink {
1370 sink(AgentEvent::RosterUpdated {
1371 update: RosterUpdate::Dropped { slot },
1372 });
1373 }
1374 Ok(())
1375 }
1376
1377 pub async fn add_agent(
1381 &mut self,
1382 mut host: AdapterHost,
1383 name: impl Into<String>,
1384 identity: impl Into<String>,
1385 command: impl Into<String>,
1386 ) -> AdapterResult<RosterSlot> {
1387 let slot = self.hosts.len();
1388 if host.adapter().slot() != slot {
1389 return Err(AdapterError::Transport(
1390 "new adapter slot must append after the existing roster".into(),
1391 ));
1392 }
1393 if let Err(error) = host.start().await {
1394 let _ = host.stop().await;
1395 return Err(error);
1396 }
1397 if let Err(error) = refresh_adapter_startup(&mut host, &self.event_sink).await {
1398 let _ = host.stop().await;
1399 return Err(error);
1400 }
1401 if let Err(error) = apply_policy_to_host(&mut host, &self.desired_policy).await {
1402 let _ = host.stop().await;
1403 return Err(error);
1404 }
1405 let capabilities = host.adapter().capabilities();
1406 self.hosts.push(host);
1407 self.relay.add_agent();
1408 self.introduced.push(false);
1409 let name = name.into();
1410 let identity = identity.into();
1411 self.roster_names.push(name.clone());
1412 self.roster_identities.push(identity.clone());
1413 self.roster_launch_specs
1414 .push((self.hosts[slot].adapter().protocol().into(), command.into()));
1415 let _ = self.queue_session_metadata();
1416 if let Some(sink) = &self.event_sink {
1417 sink(AgentEvent::RosterUpdated {
1418 update: RosterUpdate::Added {
1419 slot,
1420 name,
1421 identity,
1422 },
1423 });
1424 sink(AgentEvent::Ready { slot, capabilities });
1425 }
1426 Ok(slot)
1427 }
1428
1429 pub fn swap_agents(&mut self, first: RosterSlot, second: RosterSlot) -> AdapterResult<()> {
1433 if first == second {
1434 return Ok(());
1435 }
1436 if first >= self.hosts.len() || second >= self.hosts.len() {
1437 return Err(AdapterError::Transport("roster slot out of range".into()));
1438 }
1439 self.relay
1440 .swap_agents(first, second)
1441 .map_err(|error| AdapterError::Transport(error.into()))?;
1442 let low = first.min(second);
1443 let high = first.max(second);
1444 let high_host = self.hosts.remove(high);
1445 let low_host = self.hosts.remove(low);
1446 self.hosts.insert(low, high_host.remap(low));
1447 self.hosts.insert(high, low_host.remap(high));
1448 self.roster_names.swap(first, second);
1449 self.roster_identities.swap(first, second);
1450 if self.roster_launch_specs.len() == self.hosts.len() {
1451 self.roster_launch_specs.swap(first, second);
1452 }
1453 self.introduced.swap(first, second);
1454 if self.pair_implementer == Some(first) {
1455 self.pair_implementer = Some(second);
1456 } else if self.pair_implementer == Some(second) {
1457 self.pair_implementer = Some(first);
1458 }
1459 if self.last_public_dispatch == Some(first) {
1460 self.last_public_dispatch = Some(second);
1461 } else if self.last_public_dispatch == Some(second) {
1462 self.last_public_dispatch = Some(first);
1463 }
1464 if let Some(sink) = &self.event_sink {
1465 sink(AgentEvent::RosterUpdated {
1466 update: RosterUpdate::Swapped { first, second },
1467 });
1468 sink(AgentEvent::Ready {
1469 slot: first,
1470 capabilities: self.hosts[first].adapter().capabilities(),
1471 });
1472 sink(AgentEvent::Ready {
1473 slot: second,
1474 capabilities: self.hosts[second].adapter().capabilities(),
1475 });
1476 }
1477 let _ = self.queue_session_metadata();
1478 Ok(())
1479 }
1480
1481 pub fn relay(&self) -> &Relay {
1482 &self.relay
1483 }
1484
1485 pub fn next_slot(&self) -> RosterSlot {
1486 self.hosts.len()
1487 }
1488
1489 pub fn relay_mut(&mut self) -> &mut Relay {
1490 &mut self.relay
1491 }
1492
1493 pub fn cancellation(&self) -> RelayCancellation {
1494 RelayCancellation {
1495 requested: Arc::clone(&self.cancel_requested),
1496 active_turn_slot: Arc::clone(&self.active_turn_slot),
1497 notify: Arc::clone(&self.cancel_notify),
1498 }
1499 }
1500
1501 pub fn dispatches(&self) -> &[(RosterSlot, String)] {
1505 &self.dispatches
1506 }
1507
1508 pub async fn run_turn(
1509 &mut self,
1510 task: impl Into<String>,
1511 first_slot: RosterSlot,
1512 ) -> AdapterResult<RelayDecision> {
1513 self.run_turn_inner(task.into(), first_slot, None).await
1514 }
1515
1516 pub async fn run_turn_with_permissions(
1517 &mut self,
1518 task: impl Into<String>,
1519 first_slot: RosterSlot,
1520 permissions: &mut tokio::sync::mpsc::UnboundedReceiver<RelayPermissionAnswer>,
1521 ) -> AdapterResult<RelayDecision> {
1522 self.run_turn_inner(task.into(), first_slot, Some(permissions))
1523 .await
1524 }
1525
1526 async fn run_turn_inner(
1527 &mut self,
1528 task: String,
1529 first_slot: RosterSlot,
1530 mut permissions: Option<&mut tokio::sync::mpsc::UnboundedReceiver<RelayPermissionAnswer>>,
1531 ) -> AdapterResult<RelayDecision> {
1532 let decision = self.relay.begin(task, first_slot);
1533 let RelayDecision::Dispatch {
1534 slot,
1535 prompt,
1536 direct,
1537 can_stop,
1538 } = &decision
1539 else {
1540 return Ok(decision);
1541 };
1542 self.active_turn_slot.store(*slot, Ordering::Release);
1543 let _active_turn = ActiveTurnGuard(Arc::clone(&self.active_turn_slot));
1544 let event_sink = self.event_sink.clone();
1545 if self
1549 .hosts
1550 .get(*slot)
1551 .is_some_and(|host| host.adapter().needs_restart())
1552 && let Err(error) = self.reload(*slot).await
1553 {
1554 let limited =
1555 report_relay_failure(&mut self.relay, &event_sink, *slot, true, error.to_string());
1556 if limited {
1557 self.relay.finish(*slot, *direct, false);
1558 }
1559 let _ = self.queue_session_metadata();
1560 if limited {
1561 return Ok(decision);
1562 }
1563 return Err(error);
1564 }
1565 let speaker_name = self
1566 .roster_names
1567 .get(*slot)
1568 .cloned()
1569 .unwrap_or_else(|| self.hosts[*slot].adapter().display_name());
1570 let unseen = self.relay.unseen_context(*slot);
1571 if !*direct && !prompt.trim().is_empty() {
1578 if self.relay.shared_task().is_none() {
1579 self.relay.set_shared_task(prompt.clone());
1580 }
1581 self.relay
1582 .record_public(public_context_speaker("User"), prompt.clone());
1583 }
1584 let prompt = if unseen.is_empty() {
1585 prompt.clone()
1586 } else {
1587 format!("{prompt}\n\nPublic updates:\n{unseen}")
1588 };
1589 let introduction = if !self.introduced.get(*slot).copied().unwrap_or(false) {
1590 let self_name = speaker_name.clone();
1591 let roster = self
1592 .relay
1593 .active_slots()
1594 .map(|candidate| {
1595 let name = self
1596 .roster_names
1597 .get(candidate)
1598 .cloned()
1599 .unwrap_or_else(|| self.hosts[candidate].adapter().display_name());
1600 let role = if candidate == *slot { " — you" } else { "" };
1601 format!("{}. {name}{role}", candidate.saturating_add(1))
1602 })
1603 .collect::<Vec<_>>();
1604 let shared_task = self
1605 .relay
1606 .shared_task()
1607 .filter(|task| *task != prompt)
1608 .map(|task| format!("\n\nShared task:\n{task}"))
1609 .unwrap_or_default();
1610 format!(
1611 "You are {self_name}.\nCodeSwarm roster (ordered):\n{}\n\
1612 Turns relay sequentially through this roster. Treat the user request as the shared task; use timestamped public updates as conversation context.{shared_task}",
1613 roster.join("\n")
1614 )
1615 } else {
1616 String::new()
1617 };
1618 if !*direct && !*can_stop && self.relay.strategy() == CollaborationStrategy::Pair {
1622 self.pair_implementer = Some(*slot);
1623 }
1624 let role_block = if !*direct
1625 && self.relay.strategy() == CollaborationStrategy::Pair
1626 && self.relay.active_slots().count() >= 2
1627 {
1628 crate::workflow::pair_role(self.pair_implementer, *slot)
1629 .map(|role| {
1630 let peer = match role {
1631 crate::workflow::PairRole::Reviewer => self
1632 .last_public_dispatch
1633 .and_then(|previous| self.roster_names.get(previous).cloned()),
1634 crate::workflow::PairRole::Implementer => None,
1635 };
1636 crate::workflow::role_fragment(role, peer.as_deref())
1637 })
1638 .unwrap_or_default()
1639 } else {
1640 String::new()
1641 };
1642 let effective_can_stop = *can_stop
1643 && !(self.relay.strategy() == CollaborationStrategy::Pair
1644 && self.relay.routable_slots().count() >= 2
1645 && self.pair_implementer == Some(*slot));
1646 let role_separator =
1647 if role_block.is_empty() || (introduction.is_empty() && prompt.is_empty()) {
1648 ""
1649 } else {
1650 "\n\n"
1651 };
1652 let handoff_block = if self.relay.strategy() == CollaborationStrategy::Roster && !*direct {
1655 let targets = self
1656 .relay
1657 .routable_slots()
1658 .filter(|candidate| candidate != slot)
1659 .map(|candidate| {
1660 format!(
1661 "[CODESWARM:NEXT:{}] → {}",
1662 candidate + 1,
1663 self.roster_names
1664 .get(candidate)
1665 .cloned()
1666 .unwrap_or_else(|| self.hosts[candidate].adapter().display_name())
1667 )
1668 })
1669 .collect::<Vec<_>>()
1670 .join("\n");
1671 format!(
1672 "\n\nRoster handoff: to choose the next agent instead of normal roster order, end your final message with exactly one of the following markers:\n{targets}\nUse only a listed target, never yourself. The marker must follow all text, reasoning, and tool activity; only trailing whitespace is allowed. CodeSwarm hides it and routes at turn completion. Without a valid marker, normal roster order applies. Unavailable targets are ignored. Queued user input takes priority; turn limits and review-stop rules still apply. Choose either a handoff marker or the stop marker, not both."
1673 )
1674 } else {
1675 "\n\nAgent-directed handoff markers are disabled on this turn; they only route public turns in Roster mode.".to_owned()
1676 };
1677 let prompt = format!(
1678 "{introduction}{separator}{prompt}{role_separator}{role_block}\n\n{}{handoff_block}",
1679 if effective_can_stop {
1680 format!(
1681 "You are reviewing another agent. {STOP_TOKEN} is a global batch stop: it stops all other agents and ends the entire automated relay, not just your turn. Use it with extreme care.\nUse it only when the shared task is fully complete, no meaningful correction is needed, and no other agent should continue working. If there is any uncertainty, do not use it; state what remains and let the relay continue.\nWhen—and only when—those conditions are met, end your final response with {STOP_TOKEN}, optionally preceded by an emoji.\nOnly a terminal marker after all reasoning and tool activity requests a stop. A marker followed by more output or activity is non-stopping reasoning. Trailing whitespace is allowed.\nCodeSwarm hides the token and evaluates it only when your turn is complete."
1682 )
1683 } else {
1684 format!(
1685 "Do not use {STOP_TOKEN} on this turn. Your response must be offered to another agent for review."
1686 )
1687 },
1688 separator = if introduction.is_empty() { "" } else { "\n\n" },
1689 );
1690 let host = self
1691 .hosts
1692 .get_mut(*slot)
1693 .ok_or_else(|| AdapterError::Transport("relay selected missing adapter".into()))?;
1694 let prompt = crate::goal::prompt(self.goal.as_ref(), &prompt);
1695 if let Err(error) = host.send_prompt(prompt.clone()).await {
1696 let limited =
1697 report_relay_failure(&mut self.relay, &event_sink, *slot, true, error.to_string());
1698 if limited {
1699 self.relay.finish(*slot, *direct, false);
1700 }
1701 let _ = self.queue_session_metadata();
1702 if limited {
1703 return Ok(decision);
1704 }
1705 return Err(error);
1706 }
1707 if let Some(sink) = &self.event_sink {
1708 sink(AgentEvent::TurnStarted { slot: *slot });
1709 }
1710 if let Some(introduced) = self.introduced.get_mut(*slot) {
1711 *introduced = true;
1712 }
1713 self.dispatches.push((*slot, prompt));
1714 if !*direct {
1715 self.last_public_dispatch = Some(*slot);
1716 }
1717 let mut response = String::new();
1718 let mut stop_segment_start = 0;
1721 let mut emitted_text = 0usize;
1722 let completion_event = loop {
1723 if self.cancel_requested.swap(false, Ordering::AcqRel) {
1724 if let Err(error) = cancel_with_timeout(host).await {
1725 report_relay_failure(
1726 &mut self.relay,
1727 &event_sink,
1728 *slot,
1729 true,
1730 error.to_string(),
1731 );
1732 let _ = self.queue_session_metadata();
1733 return Err(error);
1734 }
1735 return Err(AdapterError::Transport("relay turn cancelled".into()));
1736 }
1737 let update = tokio::select! {
1738 update = host.next_update() => match update {
1739 Some(Ok(update)) => update,
1740 Some(Err(error)) => {
1741 let limited = report_relay_failure(
1742 &mut self.relay,
1743 &event_sink,
1744 *slot,
1745 true,
1746 error.to_string(),
1747 );
1748 if limited {
1749 self.relay.finish(*slot, *direct, false);
1750 }
1751 let _ = self.queue_session_metadata();
1752 if limited {
1753 return Ok(decision);
1754 }
1755 return Err(error);
1756 }
1757 None => {
1758 let error = AdapterError::Transport("adapter ended during turn".into());
1759 let limited = report_relay_failure(
1760 &mut self.relay,
1761 &event_sink,
1762 *slot,
1763 true,
1764 error.to_string(),
1765 );
1766 if limited {
1767 self.relay.finish(*slot, *direct, false);
1768 }
1769 let _ = self.queue_session_metadata();
1770 if limited {
1771 return Ok(decision);
1772 }
1773 return Err(error);
1774 }
1775 },
1776 _ = self.cancel_notify.notified() => {
1777 if !self.cancel_requested.swap(false, Ordering::AcqRel) {
1778 continue;
1779 }
1780 if let Err(error) = cancel_with_timeout(host).await {
1781 report_relay_failure(
1782 &mut self.relay,
1783 &event_sink,
1784 *slot,
1785 true,
1786 error.to_string(),
1787 );
1788 let _ = self.queue_session_metadata();
1789 return Err(error);
1790 }
1791 return Err(AdapterError::Transport("relay turn cancelled".into()));
1792 },
1793 permission = async {
1794 match permissions.as_mut() {
1795 Some(receiver) => receiver.recv().await,
1796 None => std::future::pending().await,
1797 }
1798 } => {
1799 let Some(permission) = permission else {
1800 permissions = None;
1801 continue;
1802 };
1803 if permission.slot != *slot {
1804 return Err(AdapterError::Transport(
1805 "permission response targets an inactive relay slot".into(),
1806 ));
1807 }
1808 host.answer_permission(permission.request_id, permission.answer).await?;
1809 continue;
1810 },
1811 };
1812 match &update.event {
1813 AgentEvent::Text { text, .. } => response.push_str(text),
1814 AgentEvent::Thought { text, .. } | AgentEvent::UserText { text, .. }
1815 if !text.trim().is_empty() =>
1816 {
1817 stop_segment_start = response.len();
1818 }
1819 AgentEvent::Tool { .. }
1820 | AgentEvent::Permission { .. }
1821 | AgentEvent::Terminal { .. } => {
1822 stop_segment_start = response.len();
1823 }
1824 AgentEvent::TurnComplete { .. } => {
1825 let visible_response = strip_control_tokens(&response);
1826 let visible_start = emitted_text.min(visible_response.len());
1827 let visible_start = floor_char_boundary(&visible_response, visible_start);
1828 if visible_start < visible_response.len()
1829 && let Some(sink) = &self.event_sink
1830 {
1831 sink(AgentEvent::Text {
1832 slot: *slot,
1833 text: visible_response[visible_start..].to_owned(),
1834 });
1835 }
1836 self.cancel_requested.store(false, Ordering::Release);
1837 break update.event.clone();
1838 }
1839 AgentEvent::Failed {
1840 started, detail, ..
1841 } => {
1842 let limited = report_relay_failure(
1843 &mut self.relay,
1844 &event_sink,
1845 *slot,
1846 *started,
1847 detail.clone(),
1848 );
1849 if limited {
1850 self.relay.finish(*slot, *direct, false);
1851 }
1852 let _ = self.queue_session_metadata();
1853 if limited {
1854 return Ok(decision);
1855 }
1856 return Err(AdapterError::Transport(detail.clone()));
1857 }
1858 _ => {}
1859 }
1860 if let AgentEvent::Text { .. } = &update.event {
1861 let visible_response = strip_control_tokens(&response);
1863 let visible_end = control_token_visible_end(&visible_response);
1864 if emitted_text < visible_end {
1865 if let Some(sink) = &self.event_sink {
1866 sink(AgentEvent::Text {
1867 slot: *slot,
1868 text: visible_response[emitted_text..visible_end].to_owned(),
1869 });
1870 }
1871 emitted_text = visible_end;
1872 }
1873 } else if let Some(sink) = &self.event_sink {
1874 sink(update.event.clone());
1875 }
1876 };
1877 let requested_stop = response[stop_segment_start..]
1878 .trim_end()
1879 .ends_with(STOP_TOKEN);
1880 let next_slot = requested_next_slot(&response[stop_segment_start..]);
1881 let (response, _) = strip_stop_token(&response);
1882 let response = strip_control_tokens(&response);
1883 let accepted_stop = requested_stop && effective_can_stop;
1884 let needs_stop_acknowledgment = accepted_stop && response.is_empty();
1885 let response = if needs_stop_acknowledgment {
1886 DEFAULT_STOP_ACKNOWLEDGMENT.to_owned()
1887 } else {
1888 response
1889 };
1890 if needs_stop_acknowledgment && let Some(sink) = &self.event_sink {
1895 sink(AgentEvent::Text {
1896 slot: *slot,
1897 text: response.clone(),
1898 });
1899 }
1900 if let Some(sink) = &self.event_sink {
1901 sink(completion_event);
1902 }
1903 if is_usage_limit_response(&response) {
1906 let detail = response.clone();
1907 let _ = self.relay.mark_limited(*slot);
1908 self.relay.finish(*slot, *direct, false);
1911 self.queue_session_metadata()?;
1912 if let Some(sink) = &self.event_sink {
1913 sink(AgentEvent::UsageLimitReached {
1914 slot: *slot,
1915 detail,
1916 });
1917 }
1918 return Ok(decision);
1919 }
1920 if !*direct && !response.is_empty() {
1921 self.relay
1922 .record_public(public_context_speaker(&speaker_name), response);
1923 }
1924 self.relay.mark_context_seen(*slot);
1925 self.relay
1926 .finish_with_handoff(*slot, *direct, accepted_stop, next_slot);
1927 self.queue_session_metadata()?;
1928 Ok(decision)
1929 }
1930}
1931
1932fn report_relay_failure(
1933 relay: &mut Relay,
1934 event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
1935 slot: RosterSlot,
1936 started: bool,
1937 detail: String,
1938) -> bool {
1939 if is_usage_limit_response(&detail) {
1943 let _ = relay.mark_limited(slot);
1944 if let Some(sink) = event_sink {
1945 sink(AgentEvent::UsageLimitReached { slot, detail });
1946 }
1947 return true;
1948 }
1949 if started {
1950 let _ = relay.mark_limited(slot);
1951 } else {
1952 let _ = relay.tombstone(slot);
1953 }
1954 if let Some(sink) = event_sink {
1955 sink(AgentEvent::Failed {
1956 slot,
1957 started,
1958 detail,
1959 });
1960 }
1961 started
1962}
1963
1964#[derive(Debug)]
1966pub struct ScriptedAdapter {
1967 slot: RosterSlot,
1968 capabilities: AgentCapabilities,
1969 events: VecDeque<AdapterResult<AgentEvent>>,
1970 prompts: Vec<String>,
1971}
1972
1973impl ScriptedAdapter {
1974 pub fn new(
1975 slot: RosterSlot,
1976 capabilities: AgentCapabilities,
1977 events: impl IntoIterator<Item = AgentEvent>,
1978 ) -> Self {
1979 Self {
1980 slot,
1981 capabilities,
1982 events: events.into_iter().map(Ok).collect(),
1983 prompts: Vec::new(),
1984 }
1985 }
1986
1987 pub fn prompts(&self) -> &[String] {
1988 &self.prompts
1989 }
1990}
1991
1992#[async_trait]
1993impl AgentAdapter for ScriptedAdapter {
1994 fn slot(&self) -> RosterSlot {
1995 self.slot
1996 }
1997
1998 fn capabilities(&self) -> AgentCapabilities {
1999 self.capabilities.clone()
2000 }
2001
2002 async fn start(&mut self) -> AdapterResult<()> {
2003 Ok(())
2004 }
2005
2006 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
2007 self.prompts.push(prompt);
2008 Ok(())
2009 }
2010
2011 async fn cancel(&mut self) -> AdapterResult<bool> {
2012 Ok(self.capabilities.supports_cancel)
2013 }
2014
2015 async fn answer_permission(
2016 &mut self,
2017 _request_id: String,
2018 _answer: PermissionAnswer,
2019 ) -> AdapterResult<()> {
2020 if self.capabilities.supports_permissions {
2021 Ok(())
2022 } else {
2023 Err(AdapterError::Unsupported("permission answer"))
2024 }
2025 }
2026
2027 async fn set_mode(&mut self, _mode: String) -> AdapterResult<()> {
2028 if self.capabilities.supports_modes {
2029 Ok(())
2030 } else {
2031 Err(AdapterError::Unsupported("set_mode"))
2032 }
2033 }
2034
2035 async fn reload(&mut self) -> AdapterResult<()> {
2036 Ok(())
2037 }
2038
2039 async fn stop(&mut self) -> AdapterResult<()> {
2040 Ok(())
2041 }
2042
2043 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
2044 self.events.pop_front()
2045 }
2046}
2047
2048#[derive(Debug)]
2051pub struct AgyAdapter {
2052 slot: RosterSlot,
2053 cwd: PathBuf,
2054 command: String,
2055 mode: String,
2056 mode_policy: String,
2057 session_id: Option<String>,
2058 child: Option<Child>,
2059 sender: mpsc::Sender<AdapterResult<AgentEvent>>,
2060 receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
2061 announced_session: Arc<Mutex<Option<String>>>,
2065 cancel_requested: Arc<AtomicBool>,
2066}
2067
2068impl AgyAdapter {
2069 pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
2070 let (sender, receiver) = mpsc::channel(256);
2071 Self {
2072 slot,
2073 cwd,
2074 command: command.into(),
2075 mode: "default".into(),
2076 mode_policy: "agy:full-access".into(),
2077 session_id: None,
2078 child: None,
2079 sender,
2080 receiver,
2081 announced_session: Arc::new(Mutex::new(None)),
2082 cancel_requested: Arc::new(AtomicBool::new(false)),
2083 }
2084 }
2085
2086 pub fn with_session_id(
2087 slot: RosterSlot,
2088 cwd: PathBuf,
2089 command: impl Into<String>,
2090 session_id: impl Into<String>,
2091 ) -> Self {
2092 let mut adapter = Self::new(slot, cwd, command);
2093 adapter.session_id = Some(session_id.into());
2094 adapter
2095 }
2096
2097 fn modes() -> Vec<Mode> {
2098 vec![
2099 Mode {
2100 id: "agy:full-access".into(),
2101 label: "Auto pilot".into(),
2102 },
2103 Mode {
2104 id: "agy:manual".into(),
2105 label: "Manual".into(),
2106 },
2107 Mode {
2108 id: "accept-edits".into(),
2109 label: "Accept Edits".into(),
2110 },
2111 Mode {
2112 id: "plan".into(),
2113 label: "Plan".into(),
2114 },
2115 ]
2116 }
2117
2118 async fn emit(&self, event: AdapterResult<AgentEvent>) {
2119 let _ = self.sender.send(event).await;
2120 }
2121}
2122
2123#[async_trait]
2124impl AgentAdapter for AgyAdapter {
2125 fn slot(&self) -> RosterSlot {
2126 self.slot
2127 }
2128
2129 fn session_id(&self) -> Option<String> {
2130 self.session_id.clone()
2131 }
2132
2133 fn protocol(&self) -> &'static str {
2134 "native"
2135 }
2136
2137 fn capabilities(&self) -> AgentCapabilities {
2138 AgentCapabilities {
2139 supports_cancel: true,
2140 supports_modes: true,
2141 supports_permissions: false,
2142 supports_terminals: true,
2143 supports_session_load: true,
2144 supports_models: false,
2145 }
2146 }
2147
2148 async fn start(&mut self) -> AdapterResult<()> {
2149 self.cancel_requested.store(false, Ordering::Release);
2150 self.emit(Ok(AgentEvent::Ready {
2151 slot: self.slot,
2152 capabilities: self.capabilities(),
2153 }))
2154 .await;
2155 self.emit(Ok(AgentEvent::ModesReplaced {
2156 slot: self.slot,
2157 modes: Self::modes(),
2158 current_mode: Some(self.mode_policy.clone()),
2159 }))
2160 .await;
2161 Ok(())
2162 }
2163
2164 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
2165 if self.child.is_some() {
2166 return Err(AdapterError::Transport(
2167 "agent is already handling a turn".into(),
2168 ));
2169 }
2170 if self.session_id.is_none()
2171 && let Ok(session) = self.announced_session.lock()
2172 {
2173 self.session_id = session.clone();
2174 }
2175 self.cancel_requested.store(false, Ordering::Release);
2176 let (program, args) = parse_command_line(&self.command)
2177 .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
2178 let mut command = Command::new(program);
2179 isolate_process_group(&mut command);
2180 command
2181 .args(args)
2182 .arg("--print")
2183 .arg(prompt)
2184 .arg("--print-timeout")
2185 .arg("1440m")
2187 .arg("--output-format")
2188 .arg("stream-json")
2189 .current_dir(&self.cwd)
2190 .env("CODESWARM_CWD", &self.cwd)
2191 .stdout(Stdio::piped())
2192 .stderr(Stdio::piped());
2193 if let Some(session_id) = &self.session_id {
2194 command.arg("--conversation").arg(session_id);
2195 }
2196 if self.mode != "default" {
2197 command.arg("--mode").arg(&self.mode);
2198 }
2199 let mut child = command
2200 .spawn()
2201 .map_err(|error| AdapterError::Spawn(error.to_string()))?;
2202 let stdout = match child.stdout.take() {
2203 Some(stdout) => stdout,
2204 None => {
2205 let _ = terminate_child(&mut child).await;
2206 return Err(AdapterError::Transport("agent has no stdout".into()));
2207 }
2208 };
2209 let stderr = match child.stderr.take() {
2210 Some(stderr) => stderr,
2211 None => {
2212 let _ = terminate_child(&mut child).await;
2213 return Err(AdapterError::Transport("agent has no stderr".into()));
2214 }
2215 };
2216 let sender = self.sender.clone();
2217 let slot = self.slot;
2218 let announced_session = Arc::clone(&self.announced_session);
2219 let cancel_requested = Arc::clone(&self.cancel_requested);
2220 tokio::spawn(async move {
2221 let stderr_task = tokio::spawn(async move {
2222 const MAX_STDERR: usize = 32 * 1024;
2223 let mut stderr = BufReader::new(stderr);
2224 let mut bytes = Vec::new();
2225 let mut chunk = [0_u8; 4096];
2226 while let Ok(count) = stderr.read(&mut chunk).await {
2227 if count == 0 {
2228 break;
2229 }
2230 bytes.extend_from_slice(&chunk[..count]);
2231 if bytes.len() > MAX_STDERR {
2232 let keep_from = bytes.len() - MAX_STDERR;
2233 bytes.drain(..keep_from);
2234 }
2235 }
2236 String::from_utf8_lossy(&bytes).trim().to_owned()
2237 });
2238 let mut lines = BufReader::new(stdout).lines();
2239 let mut result: Option<Value> = None;
2240 let mut streamed_response = false;
2241 while let Ok(Some(line)) = lines.next_line().await {
2242 let value = match serde_json::from_str::<Value>(&line) {
2243 Ok(value) => value,
2244 Err(_) => {
2245 continue;
2250 }
2251 };
2252 if value.get("event").and_then(Value::as_str) == Some("init")
2253 && let Some(session_id) = value
2254 .get("conversation_id")
2255 .or_else(|| value.get("conversationId"))
2256 .and_then(Value::as_str)
2257 .filter(|id| !id.is_empty())
2258 && let Ok(mut announced) = announced_session.lock()
2259 {
2260 *announced = Some(session_id.to_owned());
2261 }
2262 if value.get("event").and_then(Value::as_str) == Some("result") {
2263 result = value.get("result").cloned();
2264 }
2265 match parse_agy_value(slot, &value) {
2266 Ok(Some(event)) => {
2267 if matches!(event, AgentEvent::Text { .. }) {
2268 streamed_response = true;
2269 }
2270 if sender.send(Ok(event)).await.is_err() {
2271 break;
2272 }
2273 }
2274 Ok(None) => {}
2275 Err(error) => {
2276 let _ = sender.send(Err(error)).await;
2277 }
2278 }
2279 }
2280 let stderr = stderr_task.await.ok().unwrap_or_default();
2281 let succeeded = cancel_requested.load(Ordering::Acquire)
2282 || result
2283 .as_ref()
2284 .and_then(|result| result.get("status"))
2285 .and_then(Value::as_str)
2286 == Some("SUCCESS");
2287 if succeeded {
2288 if !streamed_response
2293 && let Some(response) = result
2294 .as_ref()
2295 .and_then(|result| result.get("response"))
2296 .and_then(Value::as_str)
2297 .filter(|response| !response.is_empty())
2298 {
2299 let _ = sender
2300 .send(Ok(AgentEvent::Text {
2301 slot,
2302 text: response.to_owned(),
2303 }))
2304 .await;
2305 }
2306 let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
2307 } else {
2308 let detail = result
2309 .as_ref()
2310 .and_then(|result| result.get("error"))
2311 .and_then(Value::as_str)
2312 .filter(|detail| !detail.is_empty())
2313 .map(str::to_owned)
2314 .or_else(|| (!stderr.is_empty()).then_some(stderr))
2315 .unwrap_or_else(|| "native stream ended before a successful result".into());
2316 let _ = sender
2317 .send(Ok(AgentEvent::Failed {
2318 slot,
2319 started: true,
2320 detail,
2321 }))
2322 .await;
2323 }
2324 });
2325 self.child = Some(child);
2326 Ok(())
2327 }
2328
2329 async fn cancel(&mut self) -> AdapterResult<bool> {
2330 self.cancel_requested.store(true, Ordering::Release);
2331 let Some(mut child) = self.child.take() else {
2332 return Ok(false);
2333 };
2334 terminate_child(&mut child).await?;
2338 let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
2339 while let Some(event) = self.receiver.recv().await {
2340 if matches!(
2341 event,
2342 Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
2343 ) {
2344 break;
2345 }
2346 }
2347 })
2348 .await;
2349 Ok(true)
2350 }
2351
2352 async fn answer_permission(
2353 &mut self,
2354 _request_id: String,
2355 _answer: PermissionAnswer,
2356 ) -> AdapterResult<()> {
2357 Err(AdapterError::Unsupported("permission answer"))
2358 }
2359
2360 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
2361 let (mode, mode_policy) = match mode.as_str() {
2362 "full-access" | "codeswarm:mode:full-access" | "auto" | "autopilot" => {
2363 ("default".to_owned(), "agy:full-access".to_owned())
2364 }
2365 "codeswarm:mode:plan" | "readonly" | "plan" => ("plan".to_owned(), "plan".to_owned()),
2366 "codeswarm:mode:accept-edits" | "acceptedits" | "accept-edits" => {
2367 ("accept-edits".to_owned(), "accept-edits".to_owned())
2368 }
2369 "codeswarm:mode:manual" | "manual" | "ask" | "default" => {
2370 ("default".to_owned(), "agy:manual".to_owned())
2371 }
2372 "agy:full-access" => ("default".to_owned(), "agy:full-access".to_owned()),
2373 "agy:manual" => ("default".to_owned(), "agy:manual".to_owned()),
2374 _ => return Err(AdapterError::Unsupported("requested Agy mode")),
2375 };
2376 self.mode = mode;
2377 self.mode_policy = mode_policy.clone();
2378 self.emit(Ok(AgentEvent::ModesReplaced {
2379 slot: self.slot,
2380 modes: Self::modes(),
2381 current_mode: Some(mode_policy),
2382 }))
2383 .await;
2384 Ok(())
2385 }
2386
2387 async fn reload(&mut self) -> AdapterResult<()> {
2388 self.stop().await?;
2389 self.start().await
2390 }
2391
2392 async fn stop(&mut self) -> AdapterResult<()> {
2393 let _ = self.cancel().await?;
2394 Ok(())
2395 }
2396
2397 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
2398 let event = self.receiver.recv().await;
2399 if matches!(
2400 event.as_ref(),
2401 Some(Ok(
2402 AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. }
2403 ))
2404 ) {
2405 if self.session_id.is_none()
2406 && let Ok(session) = self.announced_session.lock()
2407 {
2408 self.session_id = session.clone();
2409 }
2410 if let Some(mut child) = self.child.take() {
2414 let _ = child.wait().await;
2415 }
2416 }
2417 event
2418 }
2419}
2420
2421#[cfg(test)]
2422fn parse_agy_line(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
2423 let value: Value =
2424 serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
2425 parse_agy_value(slot, &value)
2426}
2427
2428fn parse_agy_value(slot: RosterSlot, value: &Value) -> AdapterResult<Option<AgentEvent>> {
2429 let event = value.get("event").and_then(Value::as_str);
2430 if let Some(terminal) = parse_terminal_event(value, event) {
2431 return Ok(Some(AgentEvent::Terminal {
2432 slot,
2433 event: terminal,
2434 }));
2435 }
2436 match event {
2437 Some("step_update") => {
2438 let Some(update) = value.get("step_update") else {
2439 return Ok(None);
2440 };
2441 let is_response = update
2442 .get("step_type")
2443 .and_then(Value::as_str)
2444 .is_some_and(|kind| kind == "agent_response");
2445 let text = update.get("text_delta").and_then(Value::as_str);
2446 let response = is_response
2447 .then(|| text.map(str::to_owned))
2448 .flatten()
2449 .filter(|text| !text.is_empty())
2450 .map(|text| AgentEvent::Text { slot, text });
2451 Ok(response.or_else(|| parse_agy_tool(slot, value)))
2452 }
2453 _ => Ok(None),
2454 }
2455}
2456
2457fn parse_agy_tool(slot: RosterSlot, value: &Value) -> Option<AgentEvent> {
2458 let update = value.get("step_update")?;
2459 if update.get("step_type")?.as_str()? != "tool" {
2460 return None;
2461 }
2462 let step_index = update.get("step_index")?.as_i64()?;
2463 let title = update
2464 .get("tool_name")
2465 .and_then(Value::as_str)
2466 .unwrap_or("Tool call")
2467 .replace('_', " ");
2468 let status = match update.get("state").and_then(Value::as_str) {
2469 Some("DONE") => ToolStatus::Completed,
2470 Some("FAILED") => ToolStatus::Failed,
2471 Some("ACTIVE") => ToolStatus::Running,
2472 _ => ToolStatus::Pending,
2473 };
2474 let detail = update
2475 .get("tool_info")
2476 .and_then(|info| info.get("output"))
2477 .and_then(Value::as_str)
2478 .map(str::to_owned);
2479 Some(AgentEvent::Tool {
2480 slot,
2481 update: ToolUpdate {
2482 id: format!("agy-tool-{step_index}"),
2483 title,
2484 status,
2485 detail,
2486 },
2487 })
2488}
2489
2490#[derive(Debug)]
2493pub struct AcpAdapter {
2494 slot: RosterSlot,
2495 program: String,
2496 args: Vec<String>,
2497 cwd: PathBuf,
2498 child: Option<Child>,
2499 reader: Option<BufReader<ChildStdout>>,
2500 capabilities: AgentCapabilities,
2501 modes: Vec<Mode>,
2502 models: Vec<Mode>,
2503 model_config_id: Option<String>,
2504 session_id: Option<String>,
2505 next_request_id: u64,
2506 prompt_request_id: Option<u64>,
2507 prompt_had_output: bool,
2508 queued_events: VecDeque<AdapterResult<AgentEvent>>,
2509 tool_updates: BTreeMap<String, ToolUpdate>,
2510 stderr_task: Option<tokio::task::JoinHandle<String>>,
2511 terminals: BTreeMap<String, TerminalProcess>,
2512 next_terminal_id: u64,
2513}
2514
2515impl AcpAdapter {
2516 pub fn new(
2517 slot: RosterSlot,
2518 cwd: PathBuf,
2519 program: impl Into<String>,
2520 args: Vec<String>,
2521 ) -> Self {
2522 Self {
2523 slot,
2524 program: program.into(),
2525 args,
2526 cwd,
2527 child: None,
2528 reader: None,
2529 capabilities: AgentCapabilities::default(),
2530 modes: Vec::new(),
2531 models: Vec::new(),
2532 model_config_id: None,
2533 session_id: None,
2534 next_request_id: 1,
2535 prompt_request_id: None,
2536 prompt_had_output: false,
2537 queued_events: VecDeque::new(),
2538 tool_updates: BTreeMap::new(),
2539 stderr_task: None,
2540 terminals: BTreeMap::new(),
2541 next_terminal_id: 1,
2542 }
2543 }
2544
2545 pub fn with_session_id(
2546 slot: RosterSlot,
2547 cwd: PathBuf,
2548 program: impl Into<String>,
2549 args: Vec<String>,
2550 session_id: impl Into<String>,
2551 ) -> Self {
2552 let mut adapter = Self::new(slot, cwd, program, args);
2553 adapter.session_id = Some(session_id.into());
2554 adapter
2555 }
2556
2557 async fn request(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2558 self.request_with_timeout(method, params, std::time::Duration::from_secs(30))
2559 .await
2560 }
2561
2562 async fn request_with_timeout(
2563 &mut self,
2564 method: &str,
2565 params: Value,
2566 deadline: std::time::Duration,
2567 ) -> AdapterResult<Value> {
2568 tokio::time::timeout(deadline, self.request_inner(method, params))
2569 .await
2570 .map_err(|_| {
2571 AdapterError::Transport(format!(
2572 "ACP {method} timed out; reload the agent to retry"
2573 ))
2574 })?
2575 }
2576
2577 async fn request_inner(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2578 let request_id = self.next_request_id;
2579 self.next_request_id += 1;
2580 self.write_json(serde_json::json!({
2581 "jsonrpc": "2.0",
2582 "id": request_id,
2583 "method": method,
2584 "params": params,
2585 }))
2586 .await?;
2587 loop {
2588 let line = self.read_line().await?;
2589 let value: Value = match serde_json::from_str(&line) {
2590 Ok(value) => value,
2591 Err(_) => {
2592 continue;
2596 }
2597 };
2598 if self.reject_empty_permission_request(&value).await? {
2599 continue;
2600 }
2601 if self.handle_client_request(&value).await? {
2602 continue;
2603 }
2604 if value
2605 .get("id")
2606 .is_some_and(|id| rpc_id_to_string(id) == request_id.to_string())
2607 {
2608 if let Some(error) = value.get("error") {
2609 return Err(AdapterError::Protocol(error.to_string()));
2610 }
2611 return value
2612 .get("result")
2613 .cloned()
2614 .ok_or_else(|| AdapterError::Protocol("response has no result".into()));
2615 }
2616 if let Some(event) = parse_acp_value(self.slot, &value, &mut self.tool_updates)? {
2617 let event = if method == "session/load" {
2618 restored_history_event(event)
2619 } else {
2620 event
2621 };
2622 self.queued_events.push_back(Ok(event));
2623 }
2624 }
2625 }
2626
2627 async fn write_json(&mut self, value: Value) -> AdapterResult<()> {
2628 let child = self
2629 .child
2630 .as_mut()
2631 .ok_or_else(|| AdapterError::Transport("ACP agent is not running".into()))?;
2632 let stdin = child
2633 .stdin
2634 .as_mut()
2635 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdin".into()))?;
2636 stdin
2637 .write_all(value.to_string().as_bytes())
2638 .await
2639 .map_err(|error| AdapterError::Transport(error.to_string()))?;
2640 stdin
2641 .write_all(b"\n")
2642 .await
2643 .map_err(|error| AdapterError::Transport(error.to_string()))
2644 }
2645
2646 async fn reset_transport(&mut self) {
2649 let terminals = std::mem::take(&mut self.terminals);
2650 for terminal in terminals.values() {
2651 terminal.stop().await;
2652 }
2653 self.queued_events.clear();
2654 self.tool_updates.clear();
2655 if let Some(mut child) = self.child.take() {
2656 let _ = terminate_child(&mut child).await;
2657 }
2658 self.reader = None;
2659 if let Some(task) = self.stderr_task.take() {
2660 task.abort();
2661 }
2662 self.prompt_request_id = None;
2663 }
2664
2665 async fn reject_empty_permission_request(&mut self, value: &Value) -> AdapterResult<bool> {
2670 if value.get("method").and_then(Value::as_str) != Some("session/request_permission")
2671 || value.get("id").is_none()
2672 {
2673 return Ok(false);
2674 }
2675 let valid = value
2676 .get("params")
2677 .and_then(|params| params.get("options"))
2678 .and_then(Value::as_array)
2679 .is_some_and(|options| !options.is_empty());
2680 if valid {
2681 return Ok(false);
2682 }
2683 self.write_json(serde_json::json!({
2684 "jsonrpc": "2.0",
2685 "id": value.get("id").cloned().unwrap_or(Value::Null),
2686 "error": {
2687 "code": -32602,
2688 "message": "Permission request requires at least one option",
2689 },
2690 }))
2691 .await?;
2692 Ok(true)
2693 }
2694
2695 fn workspace_path(&self, path: &str) -> Result<PathBuf, String> {
2696 let root = self
2697 .cwd
2698 .canonicalize()
2699 .map_err(|error| format!("unable to resolve workspace: {error}"))?;
2700 let requested = Path::new(path);
2701 let candidate = if requested.is_absolute() {
2702 requested.to_path_buf()
2703 } else {
2704 root.join(requested)
2705 };
2706 let resolved = if !candidate.exists() {
2707 let parent = candidate
2708 .parent()
2709 .ok_or_else(|| "file path has no parent".to_owned())?
2710 .canonicalize()
2711 .map_err(|error| format!("unable to resolve parent directory: {error}"))?;
2712 parent.join(
2713 candidate
2714 .file_name()
2715 .ok_or_else(|| "file path has no filename".to_owned())?,
2716 )
2717 } else {
2718 candidate
2719 .canonicalize()
2720 .map_err(|error| format!("unable to resolve file path: {error}"))?
2721 };
2722 if !resolved.starts_with(&root) {
2723 return Err("file path is outside the project".into());
2724 }
2725 Ok(resolved)
2726 }
2727
2728 fn read_workspace_text(
2729 &self,
2730 path: &str,
2731 line: Option<i64>,
2732 limit: Option<i64>,
2733 ) -> Result<String, String> {
2734 if line.is_some_and(|line| line < 1) {
2735 return Err("line must be positive".into());
2736 }
2737 if limit.is_some_and(|limit| limit < 0) {
2738 return Err("limit must not be negative".into());
2739 }
2740 let path = self.workspace_path(path)?;
2741 let mut bytes = Vec::new();
2742 let mut source = match std::fs::File::open(path) {
2743 Ok(source) => source,
2744 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(String::new()),
2745 Err(error) => return Err(error.to_string()),
2746 };
2747 source
2748 .by_ref()
2749 .take((MAX_FILE_READ_BYTES as u64).saturating_add(1))
2750 .read_to_end(&mut bytes)
2751 .map_err(|error| error.to_string())?;
2752 bytes.truncate(MAX_FILE_READ_BYTES);
2753 let text = String::from_utf8_lossy(&bytes);
2754 if line.is_none() && limit.is_none() {
2755 return Ok(text.into_owned());
2756 }
2757 let start = line.map_or(0, |line| line as usize - 1);
2758 let limit = limit.unwrap_or(i64::MAX) as usize;
2759 let selected = text
2760 .split_inclusive('\n')
2761 .skip(start)
2762 .take(limit)
2763 .collect::<String>();
2764 if line.is_some() {
2765 Ok(selected.trim_end_matches('\n').to_owned())
2766 } else {
2767 Ok(selected)
2768 }
2769 }
2770
2771 fn write_workspace_text(&self, params: &Value) -> Result<(), String> {
2772 let path = params
2773 .get("path")
2774 .and_then(Value::as_str)
2775 .filter(|path| !path.is_empty())
2776 .ok_or("path must be a non-empty string")?;
2777 let content = params
2778 .get("content")
2779 .and_then(Value::as_str)
2780 .ok_or("content must be a string")?;
2781 let path = self.workspace_path(path)?;
2782 std::fs::write(path, content).map_err(|error| error.to_string())
2783 }
2784
2785 async fn terminal_create(&mut self, params: &Value) -> Result<Value, String> {
2786 let command = params
2787 .get("command")
2788 .and_then(Value::as_str)
2789 .filter(|command| !command.trim().is_empty())
2790 .ok_or_else(|| "terminal command is required".to_owned())?;
2791 let cwd = params.get("cwd").and_then(Value::as_str).unwrap_or(".");
2792 let cwd = self.workspace_path(cwd)?;
2793 if !cwd.is_dir() {
2794 return Err("terminal cwd is not a directory".into());
2795 }
2796 let mut process = Command::new(command);
2797 isolate_process_group(&mut process);
2798 if let Some(args) = params.get("args").and_then(Value::as_array) {
2799 process.args(args.iter().filter_map(Value::as_str));
2800 }
2801 process
2802 .current_dir(&cwd)
2803 .stdin(Stdio::null())
2804 .stdout(Stdio::piped())
2805 .stderr(Stdio::piped());
2806 if let Some(env) = params.get("env") {
2807 if let Some(entries) = env.as_array() {
2808 for entry in entries {
2809 if let (Some(name), Some(value)) = (
2810 entry.get("name").and_then(Value::as_str),
2811 entry.get("value").and_then(Value::as_str),
2812 ) {
2813 process.env(name, value);
2814 }
2815 }
2816 } else if let Some(entries) = env.as_object() {
2817 for (name, value) in entries {
2818 if let Some(value) = value.as_str() {
2819 process.env(name, value);
2820 }
2821 }
2822 }
2823 }
2824 let mut child = process.spawn().map_err(|error| error.to_string())?;
2825 let stdout = child.stdout.take();
2826 let stderr = child.stderr.take();
2827 let (Some(stdout), Some(stderr)) = (stdout, stderr) else {
2828 let _ = terminate_child(&mut child).await;
2829 return Err("terminal has no output pipes".into());
2830 };
2831 let output = Arc::new(Mutex::new(Vec::new()));
2832 let truncated = Arc::new(AtomicBool::new(false));
2833 let output_readers = Arc::new(AtomicUsize::new(2));
2834 let output_limit = params
2835 .get("outputByteLimit")
2836 .and_then(Value::as_u64)
2837 .map_or(MAX_TERMINAL_OUTPUT_BYTES, |limit| {
2838 usize::try_from(limit)
2839 .unwrap_or(MAX_TERMINAL_OUTPUT_BYTES)
2840 .min(MAX_TERMINAL_OUTPUT_BYTES)
2841 });
2842 tokio::spawn(drain_terminal_output(
2843 stdout,
2844 Arc::clone(&output),
2845 Arc::clone(&truncated),
2846 Arc::clone(&output_readers),
2847 output_limit,
2848 ));
2849 tokio::spawn(drain_terminal_output(
2850 stderr,
2851 Arc::clone(&output),
2852 Arc::clone(&truncated),
2853 Arc::clone(&output_readers),
2854 output_limit,
2855 ));
2856 let id = format!("terminal-{}", self.next_terminal_id);
2857 self.next_terminal_id = self.next_terminal_id.saturating_add(1);
2858 let state = TerminalProcess {
2859 child: Arc::new(AsyncMutex::new(Some(child))),
2860 output,
2861 truncated,
2862 output_readers,
2863 };
2864 self.terminals.insert(id.clone(), state);
2865 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2866 slot: self.slot,
2867 event: TerminalEvent::Created {
2868 id: id.clone(),
2869 command: std::iter::once(command)
2870 .chain(
2871 params
2872 .get("args")
2873 .and_then(Value::as_array)
2874 .into_iter()
2875 .flatten()
2876 .filter_map(Value::as_str),
2877 )
2878 .collect::<Vec<_>>()
2879 .join(" "),
2880 },
2881 }));
2882 Ok(serde_json::json!({"terminalId": id}))
2883 }
2884
2885 async fn terminal_output(&mut self, id: &str) -> Result<Value, String> {
2886 let terminal = self
2887 .terminals
2888 .get(id)
2889 .ok_or_else(|| "terminal not found".to_owned())?
2890 .clone();
2891 let output = terminal
2892 .output
2893 .lock()
2894 .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
2895 .unwrap_or_default();
2896 let exit_code = terminal.exit_code().await;
2897 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2898 slot: self.slot,
2899 event: TerminalEvent::Output {
2900 id: id.to_owned(),
2901 text: output.clone(),
2902 },
2903 }));
2904 let mut response = serde_json::json!({
2905 "output": output,
2906 "truncated": terminal.truncated.load(Ordering::Acquire),
2907 });
2908 if let Some(code) = exit_code {
2909 response["exitStatus"] = serde_json::json!({"exitCode": code});
2910 }
2911 Ok(response)
2912 }
2913
2914 async fn terminal_wait(&mut self, id: &str) -> Result<Value, String> {
2915 let terminal = self
2916 .terminals
2917 .get(id)
2918 .ok_or_else(|| "terminal not found".to_owned())?
2919 .clone();
2920 let exit_code = terminal.wait().await;
2921 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2922 slot: self.slot,
2923 event: TerminalEvent::Exited {
2924 id: id.to_owned(),
2925 code: exit_code.unwrap_or(-1),
2926 },
2927 }));
2928 Ok(serde_json::json!({"exitCode": exit_code, "signal": Value::Null}))
2929 }
2930
2931 async fn handle_client_request(&mut self, value: &Value) -> AdapterResult<bool> {
2935 let Some(method) = value.get("method").and_then(Value::as_str) else {
2936 return Ok(false);
2937 };
2938 let Some(id) = value.get("id").cloned() else {
2939 return Ok(false);
2940 };
2941 if method == "session/request_permission" {
2944 return Ok(false);
2945 }
2946 let params = value.get("params").cloned().unwrap_or(Value::Null);
2947 let response = match method {
2948 "fs/read_text_file" => {
2949 let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2950 let line = params.get("line").and_then(Value::as_i64);
2951 let limit = params.get("limit").and_then(Value::as_i64);
2952 match self.read_workspace_text(path, line, limit) {
2953 Ok(content) => serde_json::json!({
2954 "jsonrpc": "2.0",
2955 "id": id,
2956 "result": {"content": content},
2957 }),
2958 Err(message) => serde_json::json!({
2959 "jsonrpc": "2.0",
2960 "id": id,
2961 "error": {"code": -32602, "message": message},
2962 }),
2963 }
2964 }
2965 "fs/write_text_file" => {
2966 let result = self.write_workspace_text(¶ms);
2967 match result {
2968 Ok(()) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}}),
2969 Err(message) => serde_json::json!({
2970 "jsonrpc": "2.0",
2971 "id": id,
2972 "error": {"code": -32602, "message": message},
2973 }),
2974 }
2975 }
2976 "terminal/create" => match self.terminal_create(¶ms).await {
2977 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2978 Err(message) => serde_json::json!({
2979 "jsonrpc": "2.0",
2980 "id": id,
2981 "error": {"code": -32602, "message": message},
2982 }),
2983 },
2984 "terminal/output" => {
2985 let terminal_id = params
2986 .get("terminalId")
2987 .and_then(Value::as_str)
2988 .unwrap_or("");
2989 match self.terminal_output(terminal_id).await {
2990 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2991 Err(message) => serde_json::json!({
2992 "jsonrpc": "2.0",
2993 "id": id,
2994 "error": {"code": -32602, "message": message},
2995 }),
2996 }
2997 }
2998 "terminal/wait_for_exit" => {
2999 let terminal_id = params
3000 .get("terminalId")
3001 .and_then(Value::as_str)
3002 .unwrap_or("");
3003 match self.terminal_wait(terminal_id).await {
3004 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
3005 Err(message) => serde_json::json!({
3006 "jsonrpc": "2.0",
3007 "id": id,
3008 "error": {"code": -32602, "message": message},
3009 }),
3010 }
3011 }
3012 "terminal/kill" => {
3013 let terminal_id = params
3014 .get("terminalId")
3015 .and_then(Value::as_str)
3016 .unwrap_or("");
3017 if let Some(terminal) = self.terminals.get(terminal_id) {
3018 terminal.kill().await;
3019 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
3020 } else {
3021 serde_json::json!({
3022 "jsonrpc": "2.0",
3023 "id": id,
3024 "error": {"code": -32602, "message": "terminal not found"},
3025 })
3026 }
3027 }
3028 "terminal/release" => {
3029 let terminal_id = params
3030 .get("terminalId")
3031 .and_then(Value::as_str)
3032 .unwrap_or("");
3033 if let Some(terminal) = self.terminals.remove(terminal_id) {
3034 terminal.stop().await;
3035 self.queued_events.push_back(Ok(AgentEvent::Terminal {
3036 slot: self.slot,
3037 event: TerminalEvent::Released {
3038 id: terminal_id.to_owned(),
3039 },
3040 }));
3041 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
3042 } else {
3043 serde_json::json!({
3044 "jsonrpc": "2.0",
3045 "id": id,
3046 "error": {"code": -32602, "message": "terminal not found"},
3047 })
3048 }
3049 }
3050 _ => serde_json::json!({
3051 "jsonrpc": "2.0",
3052 "id": id,
3053 "error": {"code": -32601, "message": format!("unsupported client method: {method}")},
3054 }),
3055 };
3056 self.write_json(response).await?;
3057 Ok(true)
3058 }
3059
3060 async fn read_line(&mut self) -> AdapterResult<String> {
3061 let reader = self
3062 .reader
3063 .as_mut()
3064 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdout".into()))?;
3065 read_bounded_line(reader).await
3066 }
3067
3068 async fn start(&mut self) -> AdapterResult<()> {
3069 self.modes.clear();
3070 if self.child.is_some() {
3072 self.stop().await?;
3073 }
3074 let mut command = Command::new(&self.program);
3075 isolate_process_group(&mut command);
3076 command
3077 .args(&self.args)
3078 .current_dir(&self.cwd)
3079 .stdin(Stdio::piped())
3080 .stdout(Stdio::piped())
3081 .stderr(Stdio::piped())
3082 .env("CODESWARM_CWD", &self.cwd);
3083 if self.program.to_ascii_lowercase().contains("gemini")
3084 || self
3085 .args
3086 .iter()
3087 .any(|arg| arg.to_ascii_lowercase().contains("gemini"))
3088 {
3089 command.env("GEMINI_TELEMETRY_ENABLED", "false");
3090 }
3091 let mut child = command
3092 .spawn()
3093 .map_err(|error| AdapterError::Spawn(error.to_string()))?;
3094 let stdout = match child.stdout.take() {
3095 Some(stdout) => stdout,
3096 None => {
3097 let _ = terminate_child(&mut child).await;
3098 return Err(AdapterError::Transport("ACP agent has no stdout".into()));
3099 }
3100 };
3101 let stderr = match child.stderr.take() {
3102 Some(stderr) => stderr,
3103 None => {
3104 let _ = terminate_child(&mut child).await;
3105 return Err(AdapterError::Transport("ACP agent has no stderr".into()));
3106 }
3107 };
3108 self.child = Some(child);
3109 self.reader = Some(BufReader::new(stdout));
3110 self.stderr_task = Some(tokio::spawn(drain_bounded(stderr, 32 * 1024)));
3111
3112 let initialize = match self
3113 .request(
3114 "initialize",
3115 serde_json::json!({
3116 "protocolVersion": 1,
3117 "clientCapabilities": {
3118 "fs": {"readTextFile": true, "writeTextFile": true},
3119 "terminal": true,
3120 },
3121 "clientInfo": {
3122 "name": "CodeSwarm",
3123 "title": "CodeSwarm",
3124 "version": env!("CARGO_PKG_VERSION"),
3125 },
3126 }),
3127 )
3128 .await
3129 {
3130 Ok(value) => value,
3131 Err(error) => {
3132 let _ = self.stop().await;
3133 return Err(error);
3134 }
3135 };
3136 let agent_capabilities = initialize
3137 .get("agentCapabilities")
3138 .cloned()
3139 .unwrap_or(Value::Null);
3140 self.capabilities = AgentCapabilities {
3141 supports_cancel: true,
3142 supports_modes: true,
3143 supports_permissions: true,
3144 supports_terminals: true,
3145 supports_session_load: agent_capabilities
3146 .get("loadSession")
3147 .and_then(Value::as_bool)
3148 .unwrap_or(false),
3149 supports_models: false,
3150 };
3151 let session = if let Some(session_id) = self.session_id.clone() {
3152 if !self.capabilities.supports_session_load {
3153 let _ = self.stop().await;
3154 return Err(AdapterError::Unsupported("session/load"));
3155 }
3156 match self
3157 .request(
3158 "session/load",
3159 serde_json::json!({
3160 "cwd": self.cwd,
3161 "mcpServers": [],
3162 "sessionId": session_id,
3163 }),
3164 )
3165 .await
3166 {
3167 Ok(value) => value,
3168 Err(error) => {
3169 let _ = self.stop().await;
3170 return Err(error);
3171 }
3172 }
3173 } else {
3174 let session = match self
3175 .request(
3176 "session/new",
3177 serde_json::json!({"cwd": self.cwd, "mcpServers": []}),
3178 )
3179 .await
3180 {
3181 Ok(value) => value,
3182 Err(error) => {
3183 let _ = self.stop().await;
3184 return Err(error);
3185 }
3186 };
3187 self.session_id = session
3188 .get("sessionId")
3189 .and_then(Value::as_str)
3190 .map(str::to_owned);
3191 if self.session_id.is_none() {
3192 let _ = self.stop().await;
3193 return Err(AdapterError::Protocol(
3194 "session/new returned no sessionId".into(),
3195 ));
3196 }
3197 session
3198 };
3199 self.capabilities.supports_modes = false;
3200 if let Some(modes) = session.get("modes") {
3201 let available = modes
3202 .get("availableModes")
3203 .and_then(Value::as_array)
3204 .map(|modes| {
3205 modes
3206 .iter()
3207 .filter_map(|mode| {
3208 Some(Mode {
3209 id: mode.get("id")?.as_str()?.to_owned(),
3210 label: mode.get("name")?.as_str()?.to_owned(),
3211 })
3212 })
3213 .collect::<Vec<_>>()
3214 })
3215 .unwrap_or_default();
3216 self.modes = available.clone();
3217 self.capabilities.supports_modes = !available.is_empty();
3218 self.queued_events.push_back(Ok(AgentEvent::ModesReplaced {
3219 slot: self.slot,
3220 modes: available,
3221 current_mode: modes
3222 .get("currentModeId")
3223 .and_then(Value::as_str)
3224 .map(str::to_owned),
3225 }));
3226 }
3227 self.models.clear();
3228 self.model_config_id = None;
3229 let current_model =
3230 parse_model_config(&session).and_then(|(config_id, models, current)| {
3231 self.model_config_id = Some(config_id);
3232 self.models = models;
3233 current
3234 });
3235 self.capabilities.supports_models =
3236 self.model_config_id.is_some() && !self.models.is_empty();
3237 if let Some(config_id) = self.model_config_id.clone()
3238 && !self.models.is_empty()
3239 {
3240 self.queued_events.push_back(Ok(AgentEvent::ModelsReplaced {
3241 slot: self.slot,
3242 config_id,
3243 models: self.models.clone(),
3244 current_model,
3245 }));
3246 }
3247 self.queued_events.push_back(Ok(AgentEvent::Ready {
3248 slot: self.slot,
3249 capabilities: self.capabilities(),
3250 }));
3251 Ok(())
3252 }
3253}
3254
3255fn prompt_resource_paths(prompt: &str) -> Vec<String> {
3256 let characters = prompt.chars().collect::<Vec<_>>();
3257 let mut paths = Vec::new();
3258 let mut index = 0;
3259 while index < characters.len() {
3260 if characters[index] != '@' {
3261 index += 1;
3262 continue;
3263 }
3264 index += 1;
3265 let quoted = characters.get(index) == Some(&'"');
3266 if quoted {
3267 index += 1;
3268 }
3269 let start = index;
3270 while index < characters.len()
3271 && if quoted {
3272 characters[index] != '"'
3273 } else {
3274 !characters[index].is_whitespace()
3275 }
3276 {
3277 index += 1;
3278 }
3279 if index > start {
3280 paths.push(characters[start..index].iter().collect());
3281 }
3282 if quoted && index < characters.len() {
3283 index += 1;
3284 }
3285 }
3286 paths
3287}
3288
3289fn prompt_content_blocks(cwd: &Path, prompt: &str) -> Vec<Value> {
3290 let mut blocks = vec![serde_json::json!({"type": "text", "text": prompt})];
3291 for path in prompt_resource_paths(prompt) {
3292 if path.ends_with('/') {
3293 continue;
3294 }
3295 let Ok(resource) = resources::load(cwd, &path) else {
3296 continue;
3297 };
3298 let uri = format!("file://{}", resource.path.display());
3299 let resource_value = if let Some(text) = resource.text {
3300 serde_json::json!({
3301 "uri": uri,
3302 "text": text,
3303 "mimeType": resource.mime_type,
3304 })
3305 } else if let Some(data) = resource.data {
3306 serde_json::json!({
3307 "uri": uri,
3308 "blob": BASE64.encode(data),
3309 "mimeType": resource.mime_type,
3310 })
3311 } else {
3312 continue;
3313 };
3314 blocks.push(serde_json::json!({
3315 "type": "resource",
3316 "resource": resource_value,
3317 }));
3318 }
3319 blocks
3320}
3321
3322#[async_trait]
3323impl AgentAdapter for AcpAdapter {
3324 fn slot(&self) -> RosterSlot {
3325 self.slot
3326 }
3327
3328 fn session_id(&self) -> Option<String> {
3329 self.session_id.clone()
3330 }
3331
3332 fn protocol(&self) -> &'static str {
3333 "acp"
3334 }
3335
3336 fn needs_restart(&self) -> bool {
3337 self.reader.is_none() || self.child.is_none()
3338 }
3339
3340 fn capabilities(&self) -> AgentCapabilities {
3341 self.capabilities.clone()
3342 }
3343
3344 async fn start(&mut self) -> AdapterResult<()> {
3345 AcpAdapter::start(self).await
3348 }
3349
3350 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
3351 if self.needs_restart() {
3352 return Err(AdapterError::Transport(
3353 "ACP agent transport is not running; reload the agent before retrying".into(),
3354 ));
3355 }
3356 let session_id = self
3357 .session_id
3358 .as_ref()
3359 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3360 self.tool_updates.clear();
3361 let request_id = self.next_request_id;
3362 self.next_request_id += 1;
3363 let prompt_blocks = prompt_content_blocks(&self.cwd, &prompt);
3364 let write_result = self
3365 .write_json(serde_json::json!({
3366 "jsonrpc": "2.0",
3367 "id": request_id,
3368 "method": "session/prompt",
3369 "params": {
3370 "sessionId": session_id,
3371 "prompt": prompt_blocks,
3372 },
3373 }))
3374 .await;
3375 if let Err(error) = write_result {
3376 self.reset_transport().await;
3377 return Err(error);
3378 }
3379 self.prompt_request_id = Some(request_id);
3380 self.prompt_had_output = false;
3381 Ok(())
3382 }
3383
3384 async fn cancel(&mut self) -> AdapterResult<bool> {
3385 let Some(session_id) = &self.session_id else {
3386 return Ok(false);
3387 };
3388 self.write_json(serde_json::json!({
3389 "jsonrpc": "2.0",
3390 "method": "session/cancel",
3391 "params": {"sessionId": session_id, "_meta": {}},
3392 }))
3393 .await?;
3394 let settled = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
3395 loop {
3396 match <Self as AgentAdapter>::next_event(self).await {
3397 Some(Ok(AgentEvent::TurnComplete { .. })) | None => break,
3398 Some(Ok(_)) => {}
3399 Some(Err(_)) => break,
3400 }
3401 }
3402 })
3403 .await
3404 .is_ok();
3405 if !settled {
3406 self.reload().await?;
3410 }
3411 Ok(true)
3412 }
3413
3414 async fn answer_permission(
3415 &mut self,
3416 request_id: String,
3417 answer: PermissionAnswer,
3418 ) -> AdapterResult<()> {
3419 let id = request_id
3420 .parse::<u64>()
3421 .map(Value::from)
3422 .unwrap_or_else(|_| Value::String(request_id));
3423 let outcome = match answer {
3424 PermissionAnswer::Selected { option_id } => {
3425 serde_json::json!({"outcome": "selected", "optionId": option_id})
3426 }
3427 PermissionAnswer::Cancelled => serde_json::json!({"outcome": "cancelled"}),
3428 };
3429 self.write_json(serde_json::json!({
3430 "jsonrpc": "2.0",
3431 "id": id,
3432 "result": {"outcome": outcome},
3436 }))
3437 .await
3438 }
3439
3440 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
3441 let session_id = self
3442 .session_id
3443 .as_ref()
3444 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3445 let policy = match mode.as_str() {
3446 "plan" => "codeswarm:mode:plan",
3447 "default" | "manual" => "codeswarm:mode:manual",
3448 "accept-edits" => "codeswarm:mode:accept-edits",
3449 "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
3450 other => other,
3451 };
3452 let native_mode = crate::policy::resolve(policy, &self.modes)
3453 .map(|mode| mode.id)
3454 .unwrap_or(mode);
3455 let _ = self
3456 .request(
3457 "session/set_mode",
3458 serde_json::json!({"sessionId": session_id, "modeId": native_mode.clone()}),
3459 )
3460 .await?;
3461 self.queued_events.push_back(Ok(AgentEvent::ModeUpdated {
3462 slot: self.slot,
3463 current_mode: native_mode,
3464 }));
3465 Ok(())
3466 }
3467
3468 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
3469 let session_id = self
3470 .session_id
3471 .clone()
3472 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3473 let config_id = self
3474 .model_config_id
3475 .clone()
3476 .ok_or(AdapterError::Unsupported("set_model"))?;
3477 if !self.models.iter().any(|candidate| candidate.id == model) {
3478 return Err(AdapterError::Protocol(
3479 "model is not advertised by the agent".into(),
3480 ));
3481 }
3482 let _ = self
3483 .request(
3484 "session/set_config_option",
3485 serde_json::json!({
3486 "sessionId": session_id,
3487 "configId": config_id,
3488 "value": model,
3489 }),
3490 )
3491 .await?;
3492 Ok(())
3493 }
3494
3495 async fn reload(&mut self) -> AdapterResult<()> {
3496 let session_id = self
3501 .capabilities
3502 .supports_session_load
3503 .then(|| self.session_id.clone())
3504 .flatten();
3505 self.stop().await?;
3506 self.session_id = session_id.clone();
3507 let result = self.start().await;
3508 if result.is_err() {
3509 self.session_id = session_id;
3513 }
3514 result
3515 }
3516
3517 async fn stop(&mut self) -> AdapterResult<()> {
3518 let terminals = std::mem::take(&mut self.terminals);
3519 self.queued_events.clear();
3520 self.tool_updates.clear();
3521 for terminal in terminals.values() {
3522 terminal.stop().await;
3523 }
3524 if let Some(mut child) = self.child.take() {
3525 terminate_child(&mut child).await?;
3526 }
3527 self.reader = None;
3528 self.session_id = None;
3529 self.prompt_request_id = None;
3530 if let Some(task) = self.stderr_task.take() {
3531 task.abort();
3532 let _ = task.await;
3533 }
3534 Ok(())
3535 }
3536
3537 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
3538 if let Some(event) = self.queued_events.pop_front() {
3539 return Some(event);
3540 }
3541 loop {
3542 let line = match self.read_line().await {
3543 Ok(line) => line,
3544 Err(error) => {
3545 self.reset_transport().await;
3549 return Some(Err(error));
3550 }
3551 };
3552 let value: Value = match serde_json::from_str(&line) {
3553 Ok(value) => value,
3554 Err(_) => continue,
3555 };
3556 match self.reject_empty_permission_request(&value).await {
3557 Ok(true) => continue,
3558 Ok(false) => {}
3559 Err(error) => return Some(Err(error)),
3560 }
3561 match self.handle_client_request(&value).await {
3562 Ok(true) => {
3563 if self.prompt_request_id.is_some() {
3566 self.prompt_had_output = true;
3567 }
3568 continue;
3569 }
3570 Ok(false) => {}
3571 Err(error) => return Some(Err(error)),
3572 }
3573 match parse_acp_value(self.slot, &value, &mut self.tool_updates) {
3574 Ok(Some(event)) => {
3575 if let AgentEvent::ModelsReplaced {
3576 config_id, models, ..
3577 } = &event
3578 {
3579 self.model_config_id = Some(config_id.clone());
3580 self.models = models.clone();
3581 self.capabilities.supports_models = !models.is_empty();
3582 }
3583 if self.prompt_request_id.is_some()
3584 && matches!(
3585 event,
3586 AgentEvent::Text { .. }
3587 | AgentEvent::Thought { .. }
3588 | AgentEvent::Tool { .. }
3589 | AgentEvent::Permission { .. }
3590 | AgentEvent::Terminal { .. }
3591 )
3592 {
3593 self.prompt_had_output = true;
3594 }
3595 return Some(Ok(event));
3596 }
3597 Ok(None) => {}
3598 Err(error) => return Some(Err(error)),
3599 }
3600 if value.get("id").is_some_and(|id| {
3601 self.prompt_request_id
3602 .is_some_and(|expected| rpc_id_to_string(id) == expected.to_string())
3603 }) {
3604 if let Some(error) = value.get("error") {
3605 self.prompt_request_id = None;
3606 return Some(Err(AdapterError::Protocol(error.to_string())));
3607 }
3608 self.prompt_request_id = None;
3609 let stop_reason = value
3610 .get("result")
3611 .and_then(|result| result.get("stopReason"))
3612 .and_then(Value::as_str);
3613 if let Some(reason) =
3614 stop_reason.filter(|reason| !matches!(*reason, "end_turn" | "cancelled"))
3615 {
3616 let detail = match reason {
3617 "max_tokens" => {
3618 "ACP turn stopped because the output token limit was reached".into()
3619 }
3620 "max_turn_requests" => {
3621 "ACP turn stopped because the tool-turn limit was reached".into()
3622 }
3623 other => format!("ACP turn stopped before responding: {other}"),
3624 };
3625 return Some(Ok(AgentEvent::Failed {
3626 slot: self.slot,
3627 started: true,
3628 detail,
3629 }));
3630 }
3631 if !self.prompt_had_output && !matches!(stop_reason, Some("cancelled")) {
3636 return Some(Ok(AgentEvent::Failed {
3637 slot: self.slot,
3638 started: true,
3639 detail: "ACP turn ended with no agent output; the configured model may be unauthorized or out of quota — pick another model with /model or verify with `opencode run`".into(),
3640 }));
3641 }
3642 return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3643 }
3644 }
3645 }
3646}
3647
3648#[cfg(test)]
3649fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3650 let value: Value =
3651 serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3652 parse_acp_value(slot, &value, &mut BTreeMap::new())
3653}
3654
3655fn parse_acp_value(
3656 slot: RosterSlot,
3657 value: &Value,
3658 tools: &mut BTreeMap<String, ToolUpdate>,
3659) -> AdapterResult<Option<AgentEvent>> {
3660 let method = value.get("method").and_then(Value::as_str);
3661 if method == Some("session/request_permission") {
3662 let params = value.get("params").cloned().unwrap_or(Value::Null);
3663 let request_id = value
3664 .get("id")
3665 .map(rpc_id_to_string)
3666 .unwrap_or_else(|| "permission".into());
3667 return Ok(parse_permission_event(
3668 slot,
3669 ¶ms,
3670 &request_id,
3671 params.get("options"),
3672 ));
3673 }
3674 if method != Some("session/update") {
3675 return Ok(None);
3676 }
3677 let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3678 return Ok(None);
3679 };
3680 let kind = update.get("sessionUpdate").and_then(Value::as_str);
3681 if kind == Some("config_option_update")
3682 && let Some((config_id, models, current_model)) = parse_model_config(update)
3683 {
3684 return Ok(Some(AgentEvent::ModelsReplaced {
3685 slot,
3686 config_id,
3687 models,
3688 current_model,
3689 }));
3690 }
3691 if kind == Some("request_permission") {
3692 let request_id = update
3693 .get("toolCall")
3694 .and_then(|tool| tool.get("toolCallId"))
3695 .and_then(Value::as_str)
3696 .unwrap_or("permission");
3697 return Ok(parse_permission_event(
3698 slot,
3699 update,
3700 request_id,
3701 update.get("options"),
3702 ));
3703 }
3704 if kind == Some("available_commands_update") {
3705 let commands = update
3706 .get("availableCommands")
3707 .and_then(Value::as_array)
3708 .map(|commands| {
3709 commands
3710 .iter()
3711 .filter_map(|command| {
3712 let name = command.get("name").and_then(Value::as_str)?.trim();
3713 (!name.is_empty()).then(|| AgentCommand {
3714 name: name.to_owned(),
3715 })
3716 })
3717 .collect::<Vec<_>>()
3718 })
3719 .unwrap_or_default();
3720 return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3721 }
3722 if kind == Some("current_mode_update") {
3723 if let Some(mode) = update
3724 .get("currentModeId")
3725 .and_then(Value::as_str)
3726 .filter(|mode| !mode.trim().is_empty())
3727 {
3728 return Ok(Some(AgentEvent::ModeUpdated {
3729 slot,
3730 current_mode: mode.to_owned(),
3731 }));
3732 }
3733 return Ok(None);
3734 }
3735 if kind == Some("usage_update") {
3736 let Some(used) = update.get("used").and_then(Value::as_u64) else {
3737 return Ok(None);
3738 };
3739 let Some(size) = update.get("size").and_then(Value::as_u64) else {
3740 return Ok(None);
3741 };
3742 return Ok(Some(AgentEvent::UsageUpdated {
3743 slot,
3744 usage: UsageUpdate { used, size },
3745 }));
3746 }
3747 if let Some(terminal) = parse_terminal_event(update, kind) {
3748 return Ok(Some(AgentEvent::Terminal {
3749 slot,
3750 event: terminal,
3751 }));
3752 }
3753 let text = update
3754 .get("content")
3755 .and_then(|content| content.get("text"))
3756 .and_then(Value::as_str)
3757 .map(str::to_owned);
3758 if kind == Some("user_message_chunk") {
3759 return Ok(text
3760 .filter(|text| !text.is_empty())
3761 .map(|text| AgentEvent::UserText { slot, text }));
3762 }
3763 if kind == Some("agent_message_chunk")
3764 && let Some(mode) = text
3765 .as_deref()
3766 .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3767 .map(str::trim)
3768 .filter(|mode| !mode.is_empty())
3769 {
3770 return Ok(Some(AgentEvent::ModesReplaced {
3775 slot,
3776 modes: vec![Mode {
3777 id: mode.to_owned(),
3778 label: mode.to_owned(),
3779 }],
3780 current_mode: Some(mode.to_owned()),
3781 }));
3782 }
3783 match (kind, text) {
3784 (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3785 Ok(Some(AgentEvent::Text { slot, text }))
3786 }
3787 (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3788 Ok(Some(AgentEvent::Thought { slot, text }))
3789 }
3790 (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3791 Ok(normalize_acp_tool(update, tools).map(|update| AgentEvent::Tool { slot, update }))
3792 }
3793 _ => Ok(None),
3794 }
3795}
3796
3797fn normalize_acp_tool(
3800 value: &Value,
3801 tools: &mut BTreeMap<String, ToolUpdate>,
3802) -> Option<ToolUpdate> {
3803 let id = value.get("toolCallId")?.as_str()?;
3804 if id.trim().is_empty() {
3805 return None;
3806 }
3807 if value.get("sessionUpdate").and_then(Value::as_str) == Some("tool_call") {
3808 tools.remove(id);
3809 }
3810 let tool = tools.entry(id.to_owned()).or_insert_with(|| ToolUpdate {
3811 id: id.to_owned(),
3812 title: "Tool call".into(),
3813 status: ToolStatus::Pending,
3814 detail: None,
3815 });
3816 if let Some(title) = value.get("title").and_then(Value::as_str) {
3817 tool.title = title.to_owned();
3818 }
3819 if let Some(status) =
3820 value
3821 .get("status")
3822 .and_then(Value::as_str)
3823 .and_then(|status| match status {
3824 "pending" => Some(ToolStatus::Pending),
3825 "in_progress" => Some(ToolStatus::Running),
3826 "completed" => Some(ToolStatus::Completed),
3827 "failed" => Some(ToolStatus::Failed),
3828 _ => None,
3829 })
3830 {
3831 tool.status = status;
3832 }
3833 if let Some(content) = value.get("content").and_then(Value::as_array) {
3834 let text = content
3835 .iter()
3836 .filter_map(|entry| match entry.get("type").and_then(Value::as_str) {
3837 Some("content") => entry.get("content")?.get("text")?.as_str(),
3838 Some("diff") => entry.get("newText")?.as_str(),
3839 _ => None,
3840 })
3841 .collect::<Vec<_>>()
3842 .join("\n");
3843 tool.detail = (!text.is_empty()).then_some(text);
3844 } else if let Some(output) = value.get("rawOutput").filter(|output| !output.is_null()) {
3845 tool.detail = Some(
3846 output
3847 .as_str()
3848 .map(str::to_owned)
3849 .unwrap_or_else(|| output.to_string()),
3850 );
3851 }
3852 Some(tool.clone())
3853}
3854
3855fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3856 if let Some(object) = value.get("models").and_then(Value::as_object)
3860 && value.get("configOptions").is_none()
3861 {
3862 let available = object.get("availableModels")?.as_array()?;
3863 let modes = available
3864 .iter()
3865 .filter_map(|option| {
3866 let id = option.get("modelId")?.as_str()?.to_owned();
3867 let label = option
3868 .get("name")
3869 .or_else(|| option.get("label"))
3870 .and_then(Value::as_str)
3871 .unwrap_or(&id)
3872 .to_owned();
3873 Some(Mode { id, label })
3874 })
3875 .collect::<Vec<_>>();
3876 return (!modes.is_empty()).then(|| {
3877 let current = object
3878 .get("currentModelId")
3879 .and_then(Value::as_str)
3880 .map(str::to_owned);
3881 ("model".to_owned(), modes, current)
3882 });
3883 }
3884 let config = value
3885 .get("configOptions")?
3886 .as_array()?
3887 .iter()
3888 .find(|option| {
3889 option.get("category").and_then(Value::as_str) == Some("model")
3890 && matches!(
3891 option.get("type").and_then(Value::as_str),
3892 Some("select" | "enum")
3893 )
3894 })?;
3895 let config_id = config.get("id")?.as_str()?.to_owned();
3896 let models = config
3897 .get("options")?
3898 .as_array()?
3899 .iter()
3900 .filter_map(|option| {
3901 let id = option.get("value")?.as_str()?.to_owned();
3902 let label = option
3903 .get("name")
3904 .or_else(|| option.get("label"))
3905 .and_then(Value::as_str)
3906 .unwrap_or(&id)
3907 .to_owned();
3908 Some(Mode { id, label })
3909 })
3910 .collect::<Vec<_>>();
3911 (!models.is_empty()).then(|| {
3912 let current = config
3913 .get("currentValue")
3914 .and_then(Value::as_str)
3915 .map(str::to_owned);
3916 (config_id, models, current)
3917 })
3918}
3919
3920fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3925 let nested = value.get("terminal").unwrap_or(value);
3926 let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3927 let id = nested
3928 .get("terminalId")
3929 .or_else(|| nested.get("terminal_id"))
3930 .or_else(|| nested.get("id"))
3931 .and_then(Value::as_str)
3932 .unwrap_or("terminal")
3933 .to_owned();
3934 match kind {
3935 "terminal_created" | "terminal_create" | "terminal_started" => {
3936 let command = nested
3937 .get("command")
3938 .and_then(Value::as_str)
3939 .unwrap_or("")
3940 .to_owned();
3941 Some(TerminalEvent::Created { id, command })
3942 }
3943 "terminal_output" | "terminal_output_chunk" => {
3944 let text = nested
3945 .get("output")
3946 .or_else(|| nested.get("text"))
3947 .and_then(Value::as_str)
3948 .unwrap_or("")
3949 .to_owned();
3950 Some(TerminalEvent::Output { id, text })
3951 }
3952 "terminal_exited" | "terminal_exit" => {
3953 let code = nested
3954 .get("exitCode")
3955 .or_else(|| nested.get("exit_code"))
3956 .or_else(|| nested.get("code"))
3957 .and_then(Value::as_i64)
3958 .unwrap_or(0) as i32;
3959 Some(TerminalEvent::Exited { id, code })
3960 }
3961 "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3962 _ => None,
3963 }
3964}
3965
3966fn parse_permission_event(
3967 slot: RosterSlot,
3968 value: &Value,
3969 request_id: &str,
3970 options: Option<&Value>,
3971) -> Option<AgentEvent> {
3972 let tool = value.get("toolCall").unwrap_or(value);
3973 let title = tool
3974 .get("title")
3975 .and_then(Value::as_str)
3976 .unwrap_or("Agent requests permission")
3977 .to_owned();
3978 let (options, option_ids, option_kinds): (Vec<String>, Vec<String>, Vec<String>) = options
3979 .and_then(Value::as_array)
3980 .map(|options| {
3981 let parsed = options
3982 .iter()
3983 .filter_map(|option| {
3984 let label = option
3985 .get("name")
3986 .or_else(|| option.get("optionId"))
3987 .and_then(Value::as_str)?
3988 .to_owned();
3989 let option_id = option
3990 .get("optionId")
3991 .or_else(|| option.get("id"))
3992 .and_then(Value::as_str)
3993 .map(str::to_owned)
3994 .unwrap_or_else(|| label.clone());
3995 let kind = option
3996 .get("kind")
3997 .and_then(Value::as_str)
3998 .unwrap_or_default()
3999 .to_owned();
4000 Some((label, option_id, kind))
4001 })
4002 .collect::<Vec<_>>();
4003 let mut labels = Vec::with_capacity(parsed.len());
4004 let mut ids = Vec::with_capacity(parsed.len());
4005 let mut kinds = Vec::with_capacity(parsed.len());
4006 for (label, id, kind) in parsed {
4007 labels.push(label);
4008 ids.push(id);
4009 kinds.push(kind);
4010 }
4011 (labels, ids, kinds)
4012 })
4013 .unwrap_or_default();
4014 if options.is_empty() {
4015 return None;
4016 }
4017 Some(AgentEvent::Permission {
4018 slot,
4019 request: PermissionRequest {
4020 id: request_id.to_owned(),
4021 title,
4022 options,
4023 option_ids,
4024 option_kinds,
4025 },
4026 })
4027}
4028
4029fn rpc_id_to_string(value: &Value) -> String {
4030 value
4031 .as_str()
4032 .map(str::to_owned)
4033 .or_else(|| value.as_u64().map(|id| id.to_string()))
4034 .unwrap_or_else(|| value.to_string())
4035}
4036
4037#[cfg(test)]
4038mod tests {
4039 use super::{
4040 AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
4041 RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
4042 parse_model_config, prompt_content_blocks, read_bounded_line,
4043 };
4044 #[cfg(target_os = "linux")]
4045 use super::{isolate_process_group, terminate_child};
4046 use crate::TerminalEvent;
4047 use crate::{
4048 AdapterError, AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
4049 persistence::SessionMetadataStore,
4050 relay::{CollaborationStrategy, DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
4051 };
4052 use async_trait::async_trait;
4053 use serde_json::Value;
4054 use std::collections::VecDeque;
4055 use std::sync::{
4056 Arc, Mutex,
4057 atomic::{AtomicUsize, Ordering},
4058 };
4059
4060 fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
4061 let nonce = std::time::SystemTime::now()
4062 .duration_since(std::time::UNIX_EPOCH)
4063 .expect("clock")
4064 .as_nanos();
4065 std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
4066 }
4067
4068 #[test]
4069 fn malformed_file_writes_preserve_existing_content() {
4070 let root = unique_test_path("codeswarm-write-validation", "dir");
4071 std::fs::create_dir_all(&root).unwrap();
4072 let file = root.join("keep.txt");
4073 std::fs::write(&file, "valuable content").unwrap();
4074 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
4075 for content in [
4076 Value::Null,
4077 serde_json::json!(false),
4078 serde_json::json!(42),
4079 serde_json::json!([]),
4080 ] {
4081 assert!(
4082 adapter
4083 .write_workspace_text(
4084 &serde_json::json!({"path":"keep.txt", "content": content})
4085 )
4086 .is_err()
4087 );
4088 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4089 }
4090 assert!(
4091 adapter
4092 .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
4093 .is_err()
4094 );
4095 assert!(
4096 adapter
4097 .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
4098 .is_err()
4099 );
4100 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4101 adapter
4102 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
4103 .unwrap();
4104 assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
4105 adapter
4106 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
4107 .unwrap();
4108 assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
4109 std::fs::remove_dir_all(root).unwrap();
4110 }
4111
4112 #[tokio::test]
4113 async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
4114 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
4115 let mut adapter = AcpAdapter::new(
4116 0,
4117 std::env::current_dir().unwrap(),
4118 "sh",
4119 vec!["-c".into(), script.into()],
4120 );
4121 adapter.start().await.unwrap();
4122 let error = adapter
4123 .request_with_timeout(
4124 "session/set_mode",
4125 serde_json::json!({}),
4126 std::time::Duration::from_millis(10),
4127 )
4128 .await
4129 .unwrap_err();
4130 assert!(error.to_string().contains("session/set_mode timed out"));
4131 adapter.stop().await.unwrap();
4132 assert!(adapter.child.is_none());
4133 }
4134
4135 #[tokio::test]
4136 async fn goals_reach_every_roster_slot_without_native_goal_support() {
4137 use crate::goal::GoalCommand;
4138 let hosts = (0..3)
4139 .map(|slot| {
4140 AdapterHost::new(
4141 Box::new(ScriptedAdapter::new(
4142 slot,
4143 AgentCapabilities::default(),
4144 [
4145 AgentEvent::TurnComplete { slot },
4146 AgentEvent::TurnComplete { slot },
4147 ],
4148 )),
4149 None,
4150 )
4151 })
4152 .collect();
4153 let mut relay = RelayHost::new(hosts, 10).unwrap();
4154 relay.start().await.unwrap();
4155 let task = relay
4156 .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
4157 .unwrap()
4158 .unwrap();
4159 relay.relay_mut().enqueue_human(task, Some(0));
4160 for slot in 0..3 {
4161 relay.run_turn("", 0).await.unwrap();
4162 let (actual, prompt) = relay.dispatches().last().unwrap();
4163 assert_eq!(*actual, slot);
4164 assert!(prompt.contains("Active shared goal: Ship the settings screen"));
4165 }
4166 let snapshot = relay.session_metadata();
4167 let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
4168 assert!(restored.is_some());
4169 relay.restore_goal(restored);
4170 relay.reload(0).await.unwrap();
4171 relay.run_turn("", 0).await.unwrap();
4172 assert!(
4173 relay
4174 .dispatches()
4175 .last()
4176 .unwrap()
4177 .1
4178 .contains("Active shared goal: Ship the settings screen")
4179 );
4180 relay.apply_goal(GoalCommand::Done).unwrap();
4181 relay.run_turn("", 0).await.unwrap();
4182 assert!(
4183 relay
4184 .dispatches()
4185 .last()
4186 .unwrap()
4187 .1
4188 .contains("No active shared goal")
4189 );
4190 relay.apply_goal(GoalCommand::Clear).unwrap();
4191 assert!(relay.session_metadata().get("goal").unwrap().is_null());
4192 }
4193
4194 #[tokio::test]
4195 async fn replacement_agent_receives_task_after_public_journal_pruning() {
4196 let hosts = (0..2)
4197 .map(|slot| {
4198 AdapterHost::new(
4199 Box::new(ScriptedAdapter::new(
4200 slot,
4201 AgentCapabilities::default(),
4202 [
4203 AgentEvent::Text {
4204 slot,
4205 text: "progress".into(),
4206 },
4207 AgentEvent::TurnComplete { slot },
4208 AgentEvent::Text {
4209 slot,
4210 text: "more progress".into(),
4211 },
4212 AgentEvent::TurnComplete { slot },
4213 ],
4214 )),
4215 None,
4216 )
4217 })
4218 .collect();
4219 let mut relay = RelayHost::new(hosts, 10).unwrap();
4220 relay.start().await.unwrap();
4221 relay
4222 .relay_mut()
4223 .enqueue_human("Fix the login bug", Some(0));
4224 relay.run_turn("", 0).await.unwrap();
4225 relay.run_turn("", 0).await.unwrap();
4226 relay.run_turn("", 0).await.unwrap();
4227 relay.reload(1).await.unwrap();
4228 assert!(
4229 !relay
4230 .relay_mut()
4231 .unseen_context(1)
4232 .contains("Fix the login bug")
4233 );
4234 relay.run_turn("", 0).await.unwrap();
4235 assert!(
4236 relay
4237 .dispatches()
4238 .last()
4239 .unwrap()
4240 .1
4241 .contains("Shared task:\nFix the login bug")
4242 );
4243 }
4244
4245 #[cfg(target_os = "linux")]
4246 #[tokio::test]
4247 async fn termination_kills_only_the_verified_isolated_child_group() {
4248 use nix::unistd::{Pid, getpgid, getpgrp};
4249 use tokio::io::{AsyncBufReadExt, BufReader};
4250
4251 let own_group = getpgrp();
4252 let mut command = tokio::process::Command::new("sh");
4253 isolate_process_group(&mut command);
4254 command
4255 .arg("-c")
4256 .arg("sleep 60 & echo $!; wait")
4257 .stdout(std::process::Stdio::piped());
4258 let mut child = command.spawn().expect("spawn isolated shell");
4259 let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
4260 assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
4261 assert_ne!(leader, own_group);
4262
4263 let stdout = child.stdout.take().expect("child stdout");
4264 let mut lines = BufReader::new(stdout).lines();
4265 let descendant = lines
4266 .next_line()
4267 .await
4268 .expect("read descendant pid")
4269 .expect("descendant pid")
4270 .parse::<i32>()
4271 .expect("numeric descendant pid");
4272 let descendant = Pid::from_raw(descendant);
4273 assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
4274
4275 terminate_child(&mut child).await.expect("terminate group");
4276 for _ in 0..100 {
4277 if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
4278 return;
4279 }
4280 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
4281 }
4282 panic!("descendant {descendant} survived isolated group termination");
4283 }
4284
4285 #[test]
4286 fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
4287 assert_eq!(
4288 parse_command_line(r#"agent --name "local bridge" --flag 'two words'"#),
4289 Ok((
4290 "agent".into(),
4291 vec![
4292 "--name".into(),
4293 "local bridge".into(),
4294 "--flag".into(),
4295 "two words".into(),
4296 ]
4297 ),)
4298 );
4299 assert_eq!(
4300 parse_command_line(r#"agent "" escaped\ argument"#),
4301 Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
4302 );
4303 }
4304
4305 #[test]
4306 fn acp_prompt_expands_safe_at_path_resources() {
4307 let root = unique_test_path("codeswarm-prompt-resource", "dir");
4308 std::fs::create_dir_all(&root).expect("workspace");
4309 std::fs::write(root.join("note.md"), "resource text").expect("resource");
4310 let blocks = prompt_content_blocks(&root, "inspect @note.md");
4311 assert_eq!(blocks[0]["type"], "text");
4312 assert_eq!(blocks[0]["text"], "inspect @note.md");
4313 assert_eq!(blocks[1]["type"], "resource");
4314 assert_eq!(blocks[1]["resource"]["text"], "resource text");
4315 assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
4316 std::fs::remove_dir_all(root).expect("cleanup workspace");
4317 }
4318
4319 #[tokio::test]
4320 async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
4321 let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
4322 bytes.push(b'\n');
4323 let mut reader = tokio::io::BufReader::new(bytes.as_slice());
4324 assert!(matches!(
4325 read_bounded_line(&mut reader).await,
4326 Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
4327 ));
4328 }
4329
4330 #[test]
4331 fn rejects_malformed_configured_commands_before_spawn() {
4332 assert_eq!(
4333 parse_command_line("agent 'unfinished"),
4334 Err(super::CommandParseError::UnterminatedQuote)
4335 );
4336 assert_eq!(
4337 parse_command_line("agent\\"),
4338 Err(super::CommandParseError::TrailingEscape)
4339 );
4340 assert_eq!(
4341 parse_command_line(" \t"),
4342 Err(super::CommandParseError::Empty)
4343 );
4344 }
4345
4346 #[derive(Debug)]
4347 struct PendingAdapter {
4348 slot: usize,
4349 hang_on_cancel: bool,
4350 }
4351
4352 #[derive(Debug)]
4353 struct ConcurrentStartAdapter {
4354 slot: usize,
4355 barrier: Arc<tokio::sync::Barrier>,
4356 }
4357
4358 #[derive(Debug)]
4359 struct ReloadProbeAdapter {
4360 slot: usize,
4361 crashed: bool,
4362 reloaded: bool,
4363 events: VecDeque<AgentEvent>,
4364 prompts: Arc<Mutex<Vec<String>>>,
4365 }
4366
4367 #[async_trait]
4368 impl AgentAdapter for ReloadProbeAdapter {
4369 fn slot(&self) -> usize {
4370 self.slot
4371 }
4372
4373 fn display_name(&self) -> String {
4374 "Reload probe".into()
4375 }
4376
4377 fn protocol(&self) -> &'static str {
4378 "native"
4379 }
4380
4381 fn capabilities(&self) -> AgentCapabilities {
4382 AgentCapabilities {
4383 supports_modes: true,
4384 ..AgentCapabilities::default()
4385 }
4386 }
4387
4388 fn needs_restart(&self) -> bool {
4389 self.crashed && !self.reloaded
4390 }
4391
4392 async fn start(&mut self) -> super::AdapterResult<()> {
4393 self.events.push_back(AgentEvent::ModesReplaced {
4394 slot: self.slot,
4395 modes: vec![
4396 Mode {
4397 id: "codeswarm:mode:full-access".into(),
4398 label: "Auto pilot".into(),
4399 },
4400 Mode {
4401 id: "codeswarm:mode:plan".into(),
4402 label: "Plan".into(),
4403 },
4404 ],
4405 current_mode: Some("codeswarm:mode:full-access".into()),
4406 });
4407 self.events.push_back(AgentEvent::Ready {
4408 slot: self.slot,
4409 capabilities: self.capabilities(),
4410 });
4411 Ok(())
4412 }
4413
4414 async fn send_prompt(&mut self, prompt: String) -> super::AdapterResult<()> {
4415 self.prompts.lock().expect("prompts").push(prompt);
4416 if !self.crashed {
4417 self.crashed = true;
4418 self.events.push_back(AgentEvent::Failed {
4419 slot: self.slot,
4420 started: true,
4421 detail: "probe crashed".into(),
4422 });
4423 } else {
4424 self.events.push_back(AgentEvent::Text {
4425 slot: self.slot,
4426 text: "recovered".into(),
4427 });
4428 self.events
4429 .push_back(AgentEvent::TurnComplete { slot: self.slot });
4430 }
4431 Ok(())
4432 }
4433
4434 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4435 Ok(false)
4436 }
4437
4438 async fn answer_permission(
4439 &mut self,
4440 _request_id: String,
4441 _answer: PermissionAnswer,
4442 ) -> super::AdapterResult<()> {
4443 Err(super::AdapterError::Unsupported("permission answer"))
4444 }
4445
4446 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4447 Ok(())
4448 }
4449
4450 async fn reload(&mut self) -> super::AdapterResult<()> {
4451 self.reloaded = true;
4452 self.start().await
4453 }
4454
4455 async fn stop(&mut self) -> super::AdapterResult<()> {
4456 Ok(())
4457 }
4458
4459 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4460 self.events.pop_front().map(Ok)
4461 }
4462 }
4463
4464 #[async_trait]
4465 impl AgentAdapter for ConcurrentStartAdapter {
4466 fn slot(&self) -> usize {
4467 self.slot
4468 }
4469
4470 fn capabilities(&self) -> AgentCapabilities {
4471 AgentCapabilities::default()
4472 }
4473
4474 async fn start(&mut self) -> super::AdapterResult<()> {
4475 self.barrier.wait().await;
4476 Ok(())
4477 }
4478
4479 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4480 Ok(())
4481 }
4482
4483 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4484 Ok(true)
4485 }
4486
4487 async fn answer_permission(
4488 &mut self,
4489 _request_id: String,
4490 _answer: PermissionAnswer,
4491 ) -> super::AdapterResult<()> {
4492 Ok(())
4493 }
4494
4495 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4496 Ok(())
4497 }
4498
4499 async fn reload(&mut self) -> super::AdapterResult<()> {
4500 Ok(())
4501 }
4502
4503 async fn stop(&mut self) -> super::AdapterResult<()> {
4504 Ok(())
4505 }
4506
4507 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4508 std::future::pending().await
4509 }
4510 }
4511
4512 #[derive(Debug)]
4513 struct PermissionBlockingAdapter {
4514 slot: usize,
4515 phase: u8,
4516 }
4517
4518 #[async_trait]
4519 impl AgentAdapter for PermissionBlockingAdapter {
4520 fn slot(&self) -> usize {
4521 self.slot
4522 }
4523
4524 fn capabilities(&self) -> AgentCapabilities {
4525 AgentCapabilities {
4526 supports_permissions: true,
4527 ..AgentCapabilities::default()
4528 }
4529 }
4530
4531 async fn start(&mut self) -> super::AdapterResult<()> {
4532 Ok(())
4533 }
4534
4535 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4536 Ok(())
4537 }
4538
4539 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4540 Ok(true)
4541 }
4542
4543 async fn answer_permission(
4544 &mut self,
4545 request_id: String,
4546 answer: PermissionAnswer,
4547 ) -> super::AdapterResult<()> {
4548 if self.phase != 1 || request_id != "permission-1" {
4549 return Err(super::AdapterError::Protocol(
4550 "unexpected permission response".into(),
4551 ));
4552 }
4553 assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4554 self.phase = 2;
4555 Ok(())
4556 }
4557
4558 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4559 Ok(())
4560 }
4561
4562 async fn reload(&mut self) -> super::AdapterResult<()> {
4563 Ok(())
4564 }
4565
4566 async fn stop(&mut self) -> super::AdapterResult<()> {
4567 Ok(())
4568 }
4569
4570 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4571 match self.phase {
4572 0 => {
4573 self.phase = 1;
4574 Some(Ok(AgentEvent::Permission {
4575 slot: self.slot,
4576 request: crate::PermissionRequest {
4577 id: "permission-1".into(),
4578 title: "Allow?".into(),
4579 options: vec!["Allow".into()],
4580 option_ids: vec!["allow".into()],
4581 option_kinds: Vec::new(),
4582 },
4583 }))
4584 }
4585 1 => std::future::pending().await,
4586 _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4587 }
4588 }
4589 }
4590
4591 #[async_trait]
4592 impl AgentAdapter for PendingAdapter {
4593 fn slot(&self) -> usize {
4594 self.slot
4595 }
4596
4597 fn capabilities(&self) -> AgentCapabilities {
4598 AgentCapabilities {
4599 supports_cancel: true,
4600 ..AgentCapabilities::default()
4601 }
4602 }
4603
4604 async fn start(&mut self) -> super::AdapterResult<()> {
4605 Ok(())
4606 }
4607
4608 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4609 Ok(())
4610 }
4611
4612 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4613 if self.hang_on_cancel {
4614 return std::future::pending().await;
4615 }
4616 Ok(true)
4617 }
4618
4619 async fn answer_permission(
4620 &mut self,
4621 _request_id: String,
4622 _answer: PermissionAnswer,
4623 ) -> super::AdapterResult<()> {
4624 Ok(())
4625 }
4626
4627 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4628 Ok(())
4629 }
4630
4631 async fn reload(&mut self) -> super::AdapterResult<()> {
4632 Ok(())
4633 }
4634
4635 async fn stop(&mut self) -> super::AdapterResult<()> {
4636 Ok(())
4637 }
4638
4639 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4640 std::future::pending().await
4641 }
4642 }
4643
4644 #[derive(Debug)]
4645 struct StopTrackingAdapter {
4646 slot: usize,
4647 stopped: Arc<AtomicUsize>,
4648 fail_stop: bool,
4649 }
4650
4651 #[derive(Debug)]
4652 struct ModeOrderAdapter {
4653 slot: usize,
4654 log: Arc<Mutex<Vec<String>>>,
4655 phase: u8,
4656 }
4657
4658 #[derive(Debug)]
4659 struct StartupAcpAdapter {
4660 slot: usize,
4661 events: std::collections::VecDeque<AgentEvent>,
4662 }
4663
4664 impl StartupAcpAdapter {
4665 fn new(slot: usize) -> Self {
4666 Self {
4667 slot,
4668 events: [
4669 AgentEvent::ModesReplaced {
4670 slot,
4671 modes: vec![Mode {
4672 id: "full-access".into(),
4673 label: "Auto pilot".into(),
4674 }],
4675 current_mode: Some("full-access".into()),
4676 },
4677 AgentEvent::Ready {
4678 slot,
4679 capabilities: AgentCapabilities {
4680 supports_modes: true,
4681 ..AgentCapabilities::default()
4682 },
4683 },
4684 ]
4685 .into(),
4686 }
4687 }
4688 }
4689
4690 #[async_trait]
4691 impl AgentAdapter for StartupAcpAdapter {
4692 fn slot(&self) -> usize {
4693 self.slot
4694 }
4695
4696 fn protocol(&self) -> &'static str {
4697 "acp"
4698 }
4699
4700 fn capabilities(&self) -> AgentCapabilities {
4701 AgentCapabilities {
4702 supports_modes: true,
4703 ..AgentCapabilities::default()
4704 }
4705 }
4706
4707 async fn start(&mut self) -> super::AdapterResult<()> {
4708 Ok(())
4709 }
4710
4711 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4712 Ok(())
4713 }
4714
4715 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4716 Ok(true)
4717 }
4718
4719 async fn answer_permission(
4720 &mut self,
4721 _request_id: String,
4722 _answer: PermissionAnswer,
4723 ) -> super::AdapterResult<()> {
4724 Ok(())
4725 }
4726
4727 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4728 Ok(())
4729 }
4730
4731 async fn reload(&mut self) -> super::AdapterResult<()> {
4732 self.events = Self::new(self.slot).events;
4733 Ok(())
4734 }
4735
4736 async fn stop(&mut self) -> super::AdapterResult<()> {
4737 Ok(())
4738 }
4739
4740 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4741 self.events.pop_front().map(Ok)
4742 }
4743 }
4744
4745 #[async_trait]
4746 impl AgentAdapter for ModeOrderAdapter {
4747 fn slot(&self) -> usize {
4748 self.slot
4749 }
4750
4751 fn capabilities(&self) -> AgentCapabilities {
4752 AgentCapabilities {
4753 supports_modes: true,
4754 ..AgentCapabilities::default()
4755 }
4756 }
4757
4758 async fn start(&mut self) -> super::AdapterResult<()> {
4759 self.log.lock().expect("log").push("start".into());
4760 Ok(())
4761 }
4762
4763 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4764 self.log.lock().expect("log").push("prompt".into());
4765 Ok(())
4766 }
4767
4768 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4769 Ok(true)
4770 }
4771
4772 async fn answer_permission(
4773 &mut self,
4774 _request_id: String,
4775 _answer: PermissionAnswer,
4776 ) -> super::AdapterResult<()> {
4777 Ok(())
4778 }
4779
4780 async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4781 self.log.lock().expect("log").push(format!("mode:{mode}"));
4782 Ok(())
4783 }
4784
4785 async fn reload(&mut self) -> super::AdapterResult<()> {
4786 self.log.lock().expect("log").push("reload".into());
4787 self.phase = 0;
4788 Ok(())
4789 }
4790
4791 async fn stop(&mut self) -> super::AdapterResult<()> {
4792 Ok(())
4793 }
4794
4795 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4796 match self.phase {
4797 0 => {
4798 self.phase = 1;
4799 Some(Ok(AgentEvent::ModesReplaced {
4800 slot: self.slot,
4801 modes: vec![Mode {
4802 id: "yolo".into(),
4803 label: "YOLO".into(),
4804 }],
4805 current_mode: None,
4806 }))
4807 }
4808 1 => {
4809 self.phase = 2;
4810 Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4811 }
4812 _ => std::future::pending().await,
4813 }
4814 }
4815 }
4816
4817 #[async_trait]
4818 impl AgentAdapter for StopTrackingAdapter {
4819 fn slot(&self) -> usize {
4820 self.slot
4821 }
4822
4823 fn capabilities(&self) -> AgentCapabilities {
4824 AgentCapabilities::default()
4825 }
4826
4827 async fn start(&mut self) -> super::AdapterResult<()> {
4828 Ok(())
4829 }
4830
4831 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4832 Ok(())
4833 }
4834
4835 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4836 Ok(false)
4837 }
4838
4839 async fn answer_permission(
4840 &mut self,
4841 _request_id: String,
4842 _answer: PermissionAnswer,
4843 ) -> super::AdapterResult<()> {
4844 Ok(())
4845 }
4846
4847 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4848 Ok(())
4849 }
4850
4851 async fn reload(&mut self) -> super::AdapterResult<()> {
4852 Ok(())
4853 }
4854
4855 async fn stop(&mut self) -> super::AdapterResult<()> {
4856 self.stopped.fetch_add(1, Ordering::Relaxed);
4857 if self.fail_stop {
4858 Err(super::AdapterError::Transport("stop failed".into()))
4859 } else {
4860 Ok(())
4861 }
4862 }
4863
4864 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4865 None
4866 }
4867 }
4868
4869 #[derive(Debug)]
4873 struct FailingStartAdapter {
4874 slot: usize,
4875 stopped: Arc<AtomicUsize>,
4876 }
4877
4878 #[async_trait]
4879 impl AgentAdapter for FailingStartAdapter {
4880 fn slot(&self) -> usize {
4881 self.slot
4882 }
4883
4884 fn capabilities(&self) -> AgentCapabilities {
4885 AgentCapabilities::default()
4886 }
4887
4888 async fn start(&mut self) -> super::AdapterResult<()> {
4889 Err(super::AdapterError::Spawn("startup failed".into()))
4890 }
4891
4892 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4893 Ok(())
4894 }
4895
4896 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4897 Ok(false)
4898 }
4899
4900 async fn answer_permission(
4901 &mut self,
4902 _request_id: String,
4903 _answer: PermissionAnswer,
4904 ) -> super::AdapterResult<()> {
4905 Ok(())
4906 }
4907
4908 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4909 Ok(())
4910 }
4911
4912 async fn reload(&mut self) -> super::AdapterResult<()> {
4913 Ok(())
4914 }
4915
4916 async fn stop(&mut self) -> super::AdapterResult<()> {
4917 self.stopped.fetch_add(1, Ordering::Relaxed);
4918 Ok(())
4919 }
4920
4921 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4922 None
4923 }
4924 }
4925
4926 #[tokio::test]
4927 async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4928 let stopped = Arc::new(AtomicUsize::new(0));
4929 let relay = RelayHost::new(
4930 vec![
4931 AdapterHost::new(
4932 Box::new(StopTrackingAdapter {
4933 slot: 0,
4934 stopped: Arc::clone(&stopped),
4935 fail_stop: true,
4936 }),
4937 None,
4938 ),
4939 AdapterHost::new(
4940 Box::new(StopTrackingAdapter {
4941 slot: 1,
4942 stopped: Arc::clone(&stopped),
4943 fail_stop: false,
4944 }),
4945 None,
4946 ),
4947 ],
4948 4,
4949 )
4950 .expect("relay");
4951 let mut relay = relay;
4952
4953 let error = relay.stop().await.expect_err("first stop failure");
4954 assert!(error.to_string().contains("stop failed"));
4955 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4956 }
4957
4958 #[tokio::test]
4959 async fn relay_start_isolates_a_failed_adapter_and_keeps_healthy_peers() {
4960 let stopped = Arc::new(AtomicUsize::new(0));
4961 let mut relay = RelayHost::new(
4962 vec![
4963 AdapterHost::new(
4964 Box::new(StopTrackingAdapter {
4965 slot: 0,
4966 stopped: Arc::clone(&stopped),
4967 fail_stop: false,
4968 }),
4969 None,
4970 ),
4971 AdapterHost::new(
4972 Box::new(FailingStartAdapter {
4973 slot: 1,
4974 stopped: Arc::clone(&stopped),
4975 }),
4976 None,
4977 ),
4978 ],
4979 4,
4980 )
4981 .expect("relay");
4982
4983 relay.start().await.expect("healthy peer remains available");
4984 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0]);
4985 assert_eq!(stopped.load(Ordering::Relaxed), 1);
4986 relay.stop().await.unwrap();
4987 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4988 }
4989
4990 #[test]
4991 fn parses_acp_text_without_ui_dependency() {
4992 let event = parse_acp_notification(
4993 2,
4994 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4995 )
4996 .expect("valid ACP")
4997 .expect("text event");
4998 assert_eq!(
4999 event,
5000 AgentEvent::Text {
5001 slot: 2,
5002 text: "hello".into(),
5003 }
5004 );
5005 }
5006
5007 #[test]
5008 fn parses_acp_state_notifications_at_the_adapter_boundary() {
5009 let commands = parse_acp_notification(
5010 3,
5011 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
5012 )
5013 .expect("valid ACP")
5014 .expect("commands event");
5015 assert_eq!(
5016 commands,
5017 AgentEvent::CommandsReplaced {
5018 slot: 3,
5019 commands: vec![crate::AgentCommand {
5020 name: "review".into()
5021 }]
5022 }
5023 );
5024
5025 let mode = parse_acp_notification(
5026 3,
5027 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
5028 )
5029 .expect("valid ACP")
5030 .expect("mode event");
5031 assert_eq!(
5032 mode,
5033 AgentEvent::ModeUpdated {
5034 slot: 3,
5035 current_mode: "review".into()
5036 }
5037 );
5038
5039 let usage = parse_acp_notification(
5040 3,
5041 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
5042 )
5043 .expect("valid ACP")
5044 .expect("usage event");
5045 assert_eq!(
5046 usage,
5047 AgentEvent::UsageUpdated {
5048 slot: 3,
5049 usage: crate::UsageUpdate {
5050 used: 4200,
5051 size: 128000
5052 }
5053 }
5054 );
5055
5056 let models = parse_acp_notification(
5057 3,
5058 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"}]}]}}}"#,
5059 )
5060 .expect("valid ACP")
5061 .expect("models event");
5062 assert!(matches!(
5063 models,
5064 AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
5065 if models.len() == 2 && current_model.as_deref() == Some("smart")
5066 ));
5067 assert_eq!(
5068 parse_model_config(&serde_json::json!({
5069 "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
5070 })),
5071 None
5072 );
5073
5074 let user = parse_acp_notification(
5075 3,
5076 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
5077 )
5078 .expect("valid ACP")
5079 .expect("user event");
5080 assert_eq!(
5081 user,
5082 AgentEvent::UserText {
5083 slot: 3,
5084 text: "context".into()
5085 }
5086 );
5087 }
5088
5089 #[test]
5090 fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
5091 let event = parse_acp_notification(
5092 0,
5093 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
5094 )
5095 .expect("valid ACP")
5096 .expect("mode event");
5097 assert!(matches!(
5098 event,
5099 AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
5100 if mode == "yolo" && modes[0].id == "yolo"
5101 ));
5102 }
5103
5104 #[test]
5105 fn parses_native_agy_text_without_acp_bridge() {
5106 let event = parse_agy_line(
5107 1,
5108 r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
5109 )
5110 .expect("valid stream-json")
5111 .expect("text event");
5112 assert_eq!(
5113 event,
5114 AgentEvent::Text {
5115 slot: 1,
5116 text: "hello".into(),
5117 }
5118 );
5119 }
5120
5121 #[test]
5122 fn parses_tool_lifecycle_from_each_protocol() {
5123 let agy = parse_agy_line(
5124 1,
5125 r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
5126 )
5127 .expect("valid native tool")
5128 .expect("tool event");
5129 assert!(matches!(
5130 agy,
5131 AgentEvent::Tool {
5132 update: crate::ToolUpdate {
5133 status: ToolStatus::Completed,
5134 ..
5135 },
5136 ..
5137 }
5138 ));
5139
5140 let acp = parse_acp_notification(
5141 1,
5142 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
5143 )
5144 .expect("valid ACP tool")
5145 .expect("tool event");
5146 assert!(matches!(
5147 acp,
5148 AgentEvent::Tool {
5149 update: crate::ToolUpdate {
5150 status: ToolStatus::Failed,
5151 ..
5152 },
5153 ..
5154 }
5155 ));
5156 }
5157
5158 #[test]
5159 fn parses_terminal_lifecycle_from_acp_and_native_events() {
5160 let created = parse_acp_notification(
5161 0,
5162 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
5163 )
5164 .expect("valid ACP terminal")
5165 .expect("terminal event");
5166 assert_eq!(
5167 created,
5168 AgentEvent::Terminal {
5169 slot: 0,
5170 event: TerminalEvent::Created {
5171 id: "term-1".into(),
5172 command: "cargo test".into(),
5173 },
5174 }
5175 );
5176 let output = parse_agy_line(
5177 1,
5178 r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
5179 )
5180 .expect("valid native terminal")
5181 .expect("terminal event");
5182 assert_eq!(
5183 output,
5184 AgentEvent::Terminal {
5185 slot: 1,
5186 event: TerminalEvent::Output {
5187 id: "term-1".into(),
5188 text: "ok\n".into(),
5189 },
5190 }
5191 );
5192 let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
5193 .expect("valid native release")
5194 .expect("terminal event");
5195 assert!(matches!(
5196 released,
5197 AgentEvent::Terminal {
5198 event: TerminalEvent::Released { id },
5199 ..
5200 } if id == "term-1"
5201 ));
5202 }
5203
5204 #[test]
5205 fn parses_acp_permission_requests() {
5206 let event = parse_acp_notification(
5207 0,
5208 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"request_permission","toolCall":{"toolCallId":"t1","title":"Write file"},"options":[{"name":"Yes","optionId":"opaque-approval","kind":"allow_once"},{"name":"No","optionId":"reject","kind":"reject_once"}]}}}"#,
5209 )
5210 .expect("valid permission")
5211 .expect("permission event");
5212 assert!(matches!(
5213 event,
5214 AgentEvent::Permission { request, .. }
5215 if request.id == "t1"
5216 && request.title == "Write file"
5217 && request.options == ["Yes", "No"]
5218 && request.option_ids == ["opaque-approval", "reject"]
5219 && request.option_kinds == ["allow_once", "reject_once"]
5220 ));
5221 }
5222
5223 #[test]
5224 fn parses_acp_permission_request_as_json_rpc_request() {
5225 let event = parse_acp_notification(
5226 2,
5227 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
5228 )
5229 .expect("valid permission request")
5230 .expect("permission event");
5231 assert!(matches!(
5232 event,
5233 AgentEvent::Permission { request, .. }
5234 if request.id == "17"
5235 && request.title == "Write file"
5236 && request.options == ["allow-once", "reject"]
5237 && request.option_ids == ["allow-once", "reject"]
5238 && request.option_kinds == ["", ""]
5239 ));
5240 }
5241
5242 #[tokio::test]
5243 async fn native_adapter_explicitly_rejects_permission_answers() {
5244 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5245 assert_eq!(
5246 adapter
5247 .answer_permission(
5248 "request".into(),
5249 PermissionAnswer::Selected {
5250 option_id: "allow".into()
5251 },
5252 )
5253 .await,
5254 Err(super::AdapterError::Unsupported("permission answer"))
5255 );
5256 }
5257
5258 #[tokio::test]
5259 async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
5260 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5261 adapter
5262 .set_mode("full-access".into())
5263 .await
5264 .expect("auto-pilot alias");
5265 assert!(matches!(
5266 adapter.next_event().await,
5267 Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
5268 ));
5269 }
5270
5271 #[tokio::test]
5272 async fn native_turns_receive_a_twenty_four_hour_timeout() {
5273 let script_path = unique_test_path("codeswarm-native-timeout", "sh");
5274 std::fs::write(
5275 &script_path,
5276 r#"#!/bin/sh
5277seen=0
5278while [ "$#" -gt 0 ]; do
5279 case "$1" in
5280 --print-timeout)
5281 shift
5282 [ "$1" = "1440m" ] || exit 2
5283 seen=$((seen + 1))
5284 ;;
5285 esac
5286 shift
5287done
5288[ "$seen" = 1 ] || exit 3
5289printf '%s\n' '{"event":"result","result":{"status":"SUCCESS","response":"timeout accepted"}}'
5290"#,
5291 )
5292 .unwrap();
5293 let mut adapter = AgyAdapter::with_session_id(
5294 0,
5295 std::env::current_dir().unwrap(),
5296 format!("sh {}", script_path.display()),
5297 "saved-session",
5298 );
5299 adapter.start().await.unwrap();
5300 adapter.next_event().await.unwrap().unwrap();
5301 adapter.next_event().await.unwrap().unwrap();
5302 for prompt in ["first task", "follow-up task"] {
5303 adapter.send_prompt(prompt.into()).await.unwrap();
5304 assert!(
5305 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "timeout accepted")
5306 );
5307 assert!(matches!(
5308 adapter.next_event().await,
5309 Some(Ok(AgentEvent::TurnComplete { .. }))
5310 ));
5311 }
5312 adapter.stop().await.unwrap();
5313 std::fs::remove_file(script_path).unwrap();
5314 }
5315
5316 #[tokio::test]
5317 async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
5318 let script_path = unique_test_path("codeswarm-agy-session", "sh");
5319 std::fs::write(
5320 &script_path,
5321 "#!/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",
5322 )
5323 .expect("write native test script");
5324 let mut adapter = AgyAdapter::new(
5325 0,
5326 std::env::current_dir().expect("cwd"),
5327 format!("sh {}", script_path.display()),
5328 );
5329 adapter.start().await.expect("start native adapter");
5330 assert!(adapter.next_event().await.is_some());
5332 assert!(adapter.next_event().await.is_some());
5333 adapter
5334 .send_prompt("first".into())
5335 .await
5336 .expect("first prompt");
5337 while !matches!(
5338 adapter.next_event().await,
5339 Some(Ok(AgentEvent::TurnComplete { .. }))
5340 ) {}
5341 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5342 adapter
5343 .send_prompt("follow up".into())
5344 .await
5345 .expect("follow-up prompt");
5346 while !matches!(
5347 adapter.next_event().await,
5348 Some(Ok(AgentEvent::TurnComplete { .. }))
5349 ) {}
5350 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5351 adapter.stop().await.expect("stop native adapter");
5352 std::fs::remove_file(script_path).expect("cleanup native script");
5353 }
5354
5355 #[tokio::test]
5356 async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
5357 let script_path = unique_test_path("codeswarm-agy-failure", "sh");
5358 std::fs::write(
5359 &script_path,
5360 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
5361 )
5362 .expect("write native test script");
5363 let mut adapter = AgyAdapter::new(
5364 0,
5365 std::env::current_dir().expect("cwd"),
5366 format!("sh {}", script_path.display()),
5367 );
5368 adapter.start().await.expect("start native adapter");
5369 assert!(adapter.next_event().await.is_some());
5370 assert!(adapter.next_event().await.is_some());
5371 adapter.send_prompt("fail".into()).await.expect("prompt");
5372 assert!(matches!(
5373 adapter.next_event().await,
5374 Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
5375 if detail == "agent failed"
5376 ));
5377 adapter.stop().await.expect("stop native adapter");
5378 std::fs::remove_file(script_path).expect("cleanup native script");
5379 }
5380
5381 #[tokio::test]
5382 async fn native_crash_reaps_process_and_retries_on_next_prompt() {
5383 let script_path = unique_test_path("codeswarm-agy-retry", "sh");
5384 let marker_path = unique_test_path("codeswarm-agy-retry-marker", "txt");
5385 std::fs::write(
5386 &script_path,
5387 format!(
5388 "#!/bin/sh\ncount=0\nif [ -f '{}' ]; then count=$(cat '{}'); fi\ncount=$((count + 1))\nprintf '%s' \"$count\" > '{}'\nif [ \"$count\" = 1 ]; then\n printf '%s\\n' '{{\"event\":\"result\",\"result\":{{\"status\":\"FAILURE\",\"error\":\"first crash\"}}}}'\nelse\n printf '%s\\n' '{{\"event\":\"result\",\"result\":{{\"status\":\"SUCCESS\",\"response\":\"recovered\"}}}}'\nfi\n",
5389 marker_path.display(),
5390 marker_path.display(),
5391 marker_path.display(),
5392 ),
5393 )
5394 .expect("write retry script");
5395 let mut adapter = AgyAdapter::new(
5396 0,
5397 std::env::current_dir().expect("cwd"),
5398 format!("sh {}", script_path.display()),
5399 );
5400 adapter.start().await.expect("start native adapter");
5401 assert!(adapter.next_event().await.is_some());
5402 assert!(adapter.next_event().await.is_some());
5403 adapter
5404 .send_prompt("first".into())
5405 .await
5406 .expect("first prompt");
5407 assert!(matches!(
5408 adapter.next_event().await,
5409 Some(Ok(AgentEvent::Failed { detail, .. })) if detail == "first crash"
5410 ));
5411 adapter
5412 .send_prompt("retry".into())
5413 .await
5414 .expect("retry prompt starts a fresh process");
5415 assert!(matches!(
5416 adapter.next_event().await,
5417 Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5418 ));
5419 assert!(matches!(
5420 adapter.next_event().await,
5421 Some(Ok(AgentEvent::TurnComplete { .. }))
5422 ));
5423 adapter.stop().await.expect("stop native adapter");
5424 std::fs::remove_file(script_path).expect("cleanup retry script");
5425 std::fs::remove_file(marker_path).expect("cleanup retry marker");
5426 }
5427
5428 #[tokio::test]
5429 async fn acp_adapter_initializes_session_and_completes_a_prompt() {
5430 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"}}'"#;
5431 let cwd = std::env::current_dir().expect("cwd");
5432 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5433 adapter.start().await.expect("initialize");
5434 assert!(matches!(
5435 adapter.next_event().await,
5436 Some(Ok(AgentEvent::ModesReplaced { .. }))
5437 ));
5438 assert!(matches!(
5439 adapter.next_event().await,
5440 Some(Ok(AgentEvent::Ready { .. }))
5441 ));
5442 adapter.send_prompt("hello".into()).await.expect("prompt");
5443 assert!(matches!(
5444 adapter.next_event().await,
5445 Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
5446 ));
5447 assert!(matches!(
5448 adapter.next_event().await,
5449 Some(Ok(AgentEvent::TurnComplete { .. }))
5450 ));
5451 }
5452
5453 #[tokio::test]
5454 async fn acp_empty_end_turn_is_a_failed_turn_not_an_echo() {
5455 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1"}}'; read _; echo '{"jsonrpc":"2.0","method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"say hello"}}}}'; echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'"#;
5458 let cwd = std::env::current_dir().expect("cwd");
5459 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5460 adapter.start().await.expect("initialize");
5461 assert!(matches!(
5462 adapter.next_event().await,
5463 Some(Ok(AgentEvent::Ready { .. }))
5464 ));
5465 adapter
5466 .send_prompt("say hello".into())
5467 .await
5468 .expect("prompt");
5469 assert!(matches!(
5470 adapter.next_event().await,
5471 Some(Ok(AgentEvent::UserText { .. }))
5472 ));
5473 assert!(matches!(
5474 adapter.next_event().await,
5475 Some(Ok(AgentEvent::Failed {
5476 slot: 1,
5477 started: true,
5478 detail,
5479 })) if detail.contains("no agent output")
5480 ));
5481 }
5482
5483 #[test]
5484 fn parses_opencode_models_shape_without_config_options() {
5485 let session = serde_json::json!({
5486 "sessionId": "s",
5487 "models": {
5488 "currentModelId": "opencode-go/m",
5489 "availableModels": [
5490 {"modelId": "opencode-go/m", "name": "M"},
5491 {"modelId": "other/m2", "name": "M2"},
5492 ],
5493 },
5494 "modes": {"currentModeId": "build", "availableModes": []},
5495 });
5496 let (config_id, models, current) =
5497 super::parse_model_config(&session).expect("models shape");
5498 assert_eq!(config_id, "model");
5499 assert_eq!(models.len(), 2);
5500 assert_eq!(current.as_deref(), Some("opencode-go/m"));
5501 }
5502
5503 #[tokio::test]
5504 async fn acp_output_token_limit_is_reported_as_a_failed_turn() {
5505 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1"}}'; read _; echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"max_tokens"}}'"#;
5506 let cwd = std::env::current_dir().expect("cwd");
5507 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5508 adapter.start().await.expect("initialize");
5509 assert!(matches!(
5510 adapter.next_event().await,
5511 Some(Ok(AgentEvent::Ready { .. }))
5512 ));
5513 adapter.send_prompt("hello".into()).await.expect("prompt");
5514 assert!(matches!(
5515 adapter.next_event().await,
5516 Some(Ok(AgentEvent::Failed {
5517 slot: 1,
5518 started: true,
5519 detail,
5520 })) if detail.contains("output token limit")
5521 ));
5522 }
5523
5524 #[tokio::test]
5525 async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
5526 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"}}'"#;
5527 let cwd = std::env::current_dir().expect("cwd");
5528 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5529 adapter.start().await.expect("initialize");
5530 assert!(adapter.next_event().await.is_some());
5531 assert!(adapter.next_event().await.is_some());
5532
5533 for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
5534 adapter.send_prompt(prompt.into()).await.expect("prompt");
5535 assert!(matches!(
5536 adapter.next_event().await,
5537 Some(Ok(AgentEvent::Text { text, .. })) if text == expected
5538 ));
5539 assert!(matches!(
5540 adapter.next_event().await,
5541 Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
5542 ));
5543 }
5544 adapter.stop().await.expect("stop");
5545 }
5546
5547 #[tokio::test]
5548 async fn empty_acp_mode_catalog_disables_mode_control() {
5549 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":[]}}}'"#;
5550 let cwd = std::env::current_dir().expect("cwd");
5551 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5552 adapter.start().await.expect("initialize");
5553 assert!(!adapter.capabilities().supports_modes);
5554 assert!(matches!(
5555 adapter.next_event().await,
5556 Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
5557 ));
5558 assert!(matches!(
5559 adapter.next_event().await,
5560 Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
5561 ));
5562 adapter.stop().await.expect("stop");
5563 }
5564
5565 #[tokio::test]
5566 async fn acp_models_are_discovered_live_and_changed_through_session_config() {
5567 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"#;
5568 let cwd = std::env::current_dir().expect("cwd");
5569 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5570 adapter.start().await.expect("initialize");
5571 assert!(adapter.capabilities().supports_models);
5572 assert!(matches!(
5573 adapter.next_event().await,
5574 Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
5575 if config_id == "model"
5576 && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
5577 && current_model.as_deref() == Some("fast")
5578 ));
5579 assert!(matches!(
5580 adapter.next_event().await,
5581 Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
5582 ));
5583 adapter.set_model("smart".into()).await.expect("set model");
5584 assert!(adapter.set_model("invented".into()).await.is_err());
5585 adapter.stop().await.expect("stop");
5586 }
5587
5588 #[tokio::test]
5589 async fn acp_mode_change_is_acknowledged_without_provider_notification() {
5590 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":{}}'"#;
5591 let cwd = std::env::current_dir().expect("cwd");
5592 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5593 adapter.start().await.expect("initialize");
5594 adapter
5595 .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
5596 .await
5597 .expect("set mode");
5598 assert!(matches!(
5599 adapter.next_event().await,
5600 Some(Ok(AgentEvent::ModesReplaced { .. }))
5601 ));
5602 assert!(matches!(
5603 adapter.next_event().await,
5604 Some(Ok(AgentEvent::Ready { .. }))
5605 ));
5606 assert!(matches!(
5607 adapter.next_event().await,
5608 Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
5609 ));
5610 adapter.stop().await.expect("stop");
5611 }
5612
5613 #[tokio::test]
5614 async fn acp_reload_preserves_a_loadable_session_id() {
5615 let cwd = std::env::current_dir().expect("cwd");
5616 let mut adapter = AcpAdapter::with_session_id(
5617 0,
5618 cwd,
5619 "__codeswarm_missing_acp_for_reload_test__",
5620 Vec::new(),
5621 "saved-session",
5622 );
5623 adapter.capabilities.supports_session_load = true;
5624 assert!(adapter.reload().await.is_err());
5628 assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
5629 }
5630
5631 #[tokio::test]
5632 async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
5633 let cwd = std::env::current_dir().expect("cwd");
5634 let mut adapter = AcpAdapter::with_session_id(
5635 0,
5636 cwd,
5637 "__codeswarm_missing_nonloadable_acp__",
5638 Vec::new(),
5639 "stale-session",
5640 );
5641 adapter.capabilities.supports_session_load = false;
5642 assert!(adapter.reload().await.is_err());
5643 assert_eq!(adapter.session_id, None);
5644 }
5645
5646 #[tokio::test]
5647 async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
5648 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"}}'"#;
5649 let cwd = std::env::current_dir().expect("cwd");
5650 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5651 adapter.start().await.expect("initialize");
5652 assert!(matches!(
5653 adapter.next_event().await,
5654 Some(Ok(AgentEvent::Ready { .. }))
5655 ));
5656 adapter.send_prompt("hello".into()).await.expect("prompt");
5657 assert!(matches!(
5658 adapter.next_event().await,
5659 Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
5660 ));
5661 assert!(matches!(
5662 adapter.next_event().await,
5663 Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
5664 ));
5665 }
5666
5667 #[test]
5668 fn acp_tool_patches_preserve_fields_and_honor_explicit_replacements() {
5669 let mut tools = std::collections::BTreeMap::new();
5670 let first = serde_json::json!({"sessionUpdate":"tool_call", "toolCallId":"read", "title":"Read config", "status":"in_progress",
5671 "content":[{"type":"content", "content":{"type":"text", "text":"old output"}}]});
5672 let initial = super::normalize_acp_tool(&first, &mut tools).unwrap();
5673 assert_eq!(initial.detail.as_deref(), Some("old output"));
5674 let completed = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call_update","toolCallId":"read","status":"completed"}), &mut tools).unwrap();
5675 assert_eq!(completed.title, "Read config");
5676 assert_eq!(completed.detail.as_deref(), Some("old output"));
5677 assert_eq!(completed.status, ToolStatus::Completed);
5678 let malformed = super::normalize_acp_tool(
5679 &serde_json::json!({"toolCallId":"read","title":3,"status":"unknown","content":null}),
5680 &mut tools,
5681 )
5682 .unwrap();
5683 assert_eq!(malformed, completed);
5684 let replaced = super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read","content":[false,{"type":"content","content":{"type":"text","text":"new output"}}]}), &mut tools).unwrap();
5685 assert_eq!(replaced.detail.as_deref(), Some("new output"));
5686 let cleared = super::normalize_acp_tool(
5687 &serde_json::json!({"toolCallId":"read","content":[]}),
5688 &mut tools,
5689 )
5690 .unwrap();
5691 assert_eq!(cleared.detail, None);
5692 let raw = super::normalize_acp_tool(
5693 &serde_json::json!({"toolCallId":"read","rawOutput":{"ok":true}}),
5694 &mut tools,
5695 )
5696 .unwrap();
5697 assert_eq!(raw.detail.as_deref(), Some("{\"ok\":true}"));
5698 let fresh = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call","toolCallId":"read","title":"New call"}), &mut tools).unwrap();
5699 assert_eq!(fresh.status, ToolStatus::Pending);
5700 assert_eq!(fresh.detail, None);
5701 for invalid in [
5702 serde_json::json!({}),
5703 serde_json::json!({"toolCallId":7}),
5704 serde_json::json!({"toolCallId":" "}),
5705 ] {
5706 assert!(super::normalize_acp_tool(&invalid, &mut tools).is_none());
5707 }
5708 assert_eq!(tools.len(), 1);
5709 super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read "}), &mut tools).unwrap();
5711 assert_eq!(tools.len(), 2);
5712 }
5713
5714 #[tokio::test]
5715 async fn acp_tool_status_only_notifications_retain_name_and_output() {
5716 let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'
5717read _; echo '{"id":2,"result":{"sessionId":"s"}}'
5718read _
5719echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call","toolCallId":"r","title":"Read config","status":"in_progress","content":[{"type":"content","content":{"type":"text","text":"file content"}}]}}}'
5720echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"r","status":"completed"}}}'
5721echo '{"id":3,"result":{"stopReason":"end_turn"}}'"#;
5722 let mut adapter = AcpAdapter::new(
5723 0,
5724 std::env::current_dir().unwrap(),
5725 "sh",
5726 vec!["-c".into(), script.into()],
5727 );
5728 adapter.start().await.unwrap();
5729 adapter.next_event().await.unwrap().unwrap();
5730 adapter.send_prompt("read".into()).await.unwrap();
5731 for status in [ToolStatus::Running, ToolStatus::Completed] {
5732 let Some(Ok(AgentEvent::Tool { update, .. })) = adapter.next_event().await else {
5733 panic!("tool event");
5734 };
5735 assert_eq!(update.status, status);
5736 assert_eq!(update.title, "Read config");
5737 assert_eq!(update.detail.as_deref(), Some("file content"));
5738 }
5739 assert!(matches!(
5740 adapter.next_event().await,
5741 Some(Ok(AgentEvent::TurnComplete { .. }))
5742 ));
5743 adapter.stop().await.unwrap();
5744 }
5745
5746 #[tokio::test]
5747 async fn acp_reload_discards_old_queued_events_and_catalogs() {
5748 let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"id":2,"result":{"sessionId":"new"}}'"#;
5749 let mut adapter = AcpAdapter::new(
5750 0,
5751 std::env::current_dir().unwrap(),
5752 "sh",
5753 vec!["-c".into(), script.into()],
5754 );
5755 adapter.start().await.unwrap();
5756 adapter.queued_events.push_back(Ok(AgentEvent::Text {
5757 slot: 0,
5758 text: "stale".into(),
5759 }));
5760 adapter.modes = vec![Mode {
5761 id: "stale".into(),
5762 label: "Stale".into(),
5763 }];
5764 adapter.next_request_id = 1;
5766 adapter.reload().await.unwrap();
5767 assert!(adapter.modes.is_empty());
5768 assert_eq!(adapter.queued_events.len(), 1);
5769 assert!(matches!(
5770 adapter.next_event().await,
5771 Some(Ok(AgentEvent::Ready { .. }))
5772 ));
5773 adapter.stop().await.unwrap();
5774 assert!(adapter.queued_events.is_empty());
5775 }
5776
5777 #[tokio::test]
5778 async fn acp_load_replays_history_without_starting_a_turn() {
5779 let script = r#"
5780read _
5781echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'
5782read request
5783case "$request" in *session/load*) ;; *) exit 2;; esac
5784echo '{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"text":"old question"}}}}'
5785echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"old answer"}}}}'
5786echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_thought_chunk","content":{"text":"old reasoning"}}}}'
5787echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call","toolCallId":"old-tool","title":"Read","status":"in_progress"}}}'
5788echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"old-tool","title":"Read","status":"completed"}}}'
5789echo '{"jsonrpc":"2.0","id":2,"result":{}}'
5790read request
5791case "$request" in *session/prompt*) ;; *) exit 3;; esac
5792echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"new answer"}}}}'
5793echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
5794"#;
5795 let mut adapter = AcpAdapter::with_session_id(
5796 2,
5797 std::env::current_dir().unwrap(),
5798 "sh",
5799 vec!["-c".into(), script.into()],
5800 "saved",
5801 );
5802 adapter.start().await.unwrap();
5803 let mut state = crate::SessionState::new(3);
5804 for _ in 0..5 {
5805 let event = adapter.next_event().await.unwrap().unwrap();
5806 assert!(matches!(&event, AgentEvent::History { slot: 2, .. }));
5807 crate::reduce(&mut state, event);
5808 assert_eq!(state.active_slot, None);
5809 }
5810 assert!(matches!(
5811 adapter.next_event().await,
5812 Some(Ok(AgentEvent::Ready { slot: 2, .. }))
5813 ));
5814 adapter.send_prompt("new question".into()).await.unwrap();
5815 assert!(
5816 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "new answer")
5817 );
5818 assert!(matches!(
5819 adapter.next_event().await,
5820 Some(Ok(AgentEvent::TurnComplete { slot: 2 }))
5821 ));
5822 adapter.stop().await.unwrap();
5823 }
5824
5825 #[tokio::test]
5826 async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
5827 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
5828 let cwd = std::env::current_dir().expect("cwd");
5829 let mut adapter = AcpAdapter::with_session_id(
5830 0,
5831 cwd,
5832 "sh",
5833 vec!["-c".into(), script.into()],
5834 "existing-session",
5835 );
5836 adapter.start().await.expect("load existing session");
5837 assert!(matches!(
5838 adapter.next_event().await,
5839 Some(Ok(AgentEvent::Ready { .. }))
5840 ));
5841 }
5842
5843 #[tokio::test]
5844 async fn acp_start_failure_reaps_transport_process() {
5845 let mut adapter = AcpAdapter::new(
5850 0,
5851 std::env::current_dir().expect("cwd"),
5852 "sh",
5853 vec!["-c".into(), "printf 'not-json\\n'".into()],
5854 );
5855 assert!(adapter.start().await.is_err());
5856 assert!(adapter.child.is_none());
5857 assert!(adapter.reader.is_none());
5858 }
5859
5860 #[tokio::test]
5861 async fn acp_transport_crash_is_reloaded_before_the_next_prompt() {
5862 let marker = unique_test_path("codeswarm-acp-retry", "count");
5863 let script = format!(
5864 r#"count=0
5865if [ -f '{0}' ]; then count=$(cat '{0}'); fi
5866count=$((count + 1))
5867printf '%s' "$count" > '{0}'
5868while IFS= read -r request; do
5869 id=$(printf '%s' "$request" | sed -n 's/.*"id":\([0-9][0-9]*\).*/\1/p')
5870 case "$request" in
5871 *initialize*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"agentCapabilities":{{"loadSession":true}}}}}}' ;;
5872 *session/new*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"sessionId":"saved-session"}}}}' ;;
5873 *session/load*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{}}}}' ;;
5874 *session/prompt*)
5875 if [ "$count" = 1 ]; then exit 0; fi
5876 printf '%s\n' '{{"jsonrpc":"2.0","method":"session/update","params":{{"update":{{"sessionUpdate":"agent_message_chunk","content":{{"text":"recovered"}}}}}}}}'
5877 printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"stopReason":"end_turn"}}}}'
5878 ;;
5879 esac
5880done
5881"#,
5882 marker.display()
5883 );
5884 let mut adapter = AcpAdapter::new(
5885 0,
5886 std::env::current_dir().expect("cwd"),
5887 "sh",
5888 vec!["-c".into(), script],
5889 );
5890 adapter.start().await.expect("initial ACP startup");
5891 assert!(matches!(
5892 adapter.next_event().await,
5893 Some(Ok(AgentEvent::Ready { .. }))
5894 ));
5895 adapter
5896 .send_prompt("first".into())
5897 .await
5898 .expect("first prompt");
5899 assert!(matches!(
5900 adapter.next_event().await,
5901 Some(Err(AdapterError::Transport(_)))
5902 ));
5903 assert!(adapter.child.is_none());
5904 assert!(adapter.reader.is_none());
5905 assert_eq!(adapter.session_id(), Some("saved-session".into()));
5906
5907 adapter.reload().await.expect("reload ACP transport");
5908 assert!(matches!(
5909 adapter.next_event().await,
5910 Some(Ok(AgentEvent::Ready { .. }))
5911 ));
5912 adapter
5913 .send_prompt("retry".into())
5914 .await
5915 .expect("retry prompt");
5916 assert!(matches!(
5917 adapter.next_event().await,
5918 Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5919 ));
5920 assert!(matches!(
5921 adapter.next_event().await,
5922 Some(Ok(AgentEvent::TurnComplete { .. }))
5923 ));
5924 adapter.stop().await.expect("stop ACP");
5925 std::fs::remove_file(marker).expect("cleanup marker");
5926 }
5927
5928 #[tokio::test]
5929 async fn coordinator_reload_replays_context_and_reintroduces_a_crashed_slot() {
5930 let prompts = Arc::new(Mutex::new(Vec::new()));
5931 let healthy = ScriptedAdapter::new(
5932 0,
5933 AgentCapabilities::default(),
5934 [
5935 AgentEvent::Text {
5936 slot: 0,
5937 text: "peer context".into(),
5938 },
5939 AgentEvent::TurnComplete { slot: 0 },
5940 ],
5941 );
5942 let probe = ReloadProbeAdapter {
5943 slot: 1,
5944 crashed: false,
5945 reloaded: false,
5946 events: VecDeque::new(),
5947 prompts: Arc::clone(&prompts),
5948 };
5949 let mut relay = RelayHost::new(
5950 vec![
5951 AdapterHost::new(Box::new(healthy), None),
5952 AdapterHost::new(Box::new(probe), None),
5953 ],
5954 8,
5955 )
5956 .expect("relay");
5957 relay.start().await.expect("start");
5958 relay
5959 .run_turn("original task", 0)
5960 .await
5961 .expect("first turn");
5962 relay.run_turn("", 0).await.expect("crashed turn");
5963 relay.relay_mut().enqueue_human("retry", Some(1));
5964 relay.run_turn("", 0).await.expect("reloaded turn");
5965
5966 let prompts = prompts.lock().expect("prompts");
5967 let retry = prompts.last().expect("retry prompt");
5968 assert!(retry.contains("You are Reload probe"), "{retry}");
5969 assert!(retry.contains("original task"), "{retry}");
5970 assert!(retry.contains("peer context"), "{retry}");
5971 assert!(retry.contains("retry"), "{retry}");
5972 }
5973
5974 #[tokio::test]
5975 async fn acp_adapter_answers_permission_json_rpc_requests() {
5976 let path = std::env::temp_dir().join(format!(
5977 "codeswarm-permission-answer-{}",
5978 std::process::id()
5979 ));
5980 let script = format!(
5981 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"}}}}'"#,
5982 path.display()
5983 );
5984 let mut adapter = AcpAdapter::new(
5985 0,
5986 std::env::current_dir().expect("cwd"),
5987 "sh",
5988 vec!["-c".into(), script],
5989 );
5990 adapter.start().await.expect("start ACP");
5991 assert!(matches!(
5992 adapter.next_event().await,
5993 Some(Ok(AgentEvent::Ready { .. }))
5994 ));
5995 adapter.send_prompt("do it".into()).await.expect("prompt");
5996 assert!(matches!(
5997 adapter.next_event().await,
5998 Some(Ok(AgentEvent::Permission { request, .. }))
5999 if request.id == "9"
6000 && request.options == ["Allow once"]
6001 && request.option_ids == ["allow-once"]
6002 ));
6003 adapter
6004 .answer_permission(
6005 "9".into(),
6006 PermissionAnswer::Selected {
6007 option_id: "allow-once".into(),
6008 },
6009 )
6010 .await
6011 .expect("permission answer");
6012 assert!(matches!(
6013 adapter.next_event().await,
6014 Some(Ok(AgentEvent::TurnComplete { .. }))
6015 ));
6016 let answer: Value = serde_json::from_str(
6017 &std::fs::read_to_string(&path).expect("captured permission answer"),
6018 )
6019 .expect("valid JSON-RPC answer");
6020 assert_eq!(answer["id"], 9);
6021 assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
6022 assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
6023 std::fs::remove_file(path).expect("cleanup");
6024 }
6025
6026 #[test]
6027 fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
6028 let event = parse_acp_notification(
6029 0,
6030 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
6031 )
6032 .expect("valid JSON-RPC request");
6033 assert!(event.is_none());
6034 }
6035
6036 #[tokio::test]
6037 async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
6038 let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
6039 std::fs::write(
6040 &script_path,
6041 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
6042 )
6043 .expect("write native test script");
6044 let mut adapter = AgyAdapter::new(
6045 0,
6046 std::env::current_dir().expect("cwd"),
6047 format!("sh {}", script_path.display()),
6048 );
6049 adapter.start().await.expect("start native adapter");
6050 assert!(adapter.next_event().await.is_some());
6051 assert!(adapter.next_event().await.is_some());
6052 adapter
6053 .send_prompt("continue".into())
6054 .await
6055 .expect("prompt");
6056 assert!(matches!(
6057 adapter.next_event().await,
6058 Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
6059 ));
6060 assert!(matches!(
6061 adapter.next_event().await,
6062 Some(Ok(AgentEvent::TurnComplete { .. }))
6063 ));
6064 adapter.stop().await.expect("stop native adapter");
6065 std::fs::remove_file(script_path).expect("cleanup native script");
6066 }
6067
6068 #[test]
6069 fn acp_workspace_file_access_is_root_bound_and_size_limited() {
6070 let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
6071 let _ = std::fs::remove_dir_all(&root);
6072 std::fs::create_dir_all(&root).expect("workspace");
6073 std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
6074 let outside =
6075 std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
6076 std::fs::write(&outside, "secret").expect("outside file");
6077 let link = root.join("outside-link");
6078 #[cfg(unix)]
6079 std::os::unix::fs::symlink(&outside, &link).expect("symlink");
6080 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
6081
6082 assert_eq!(
6083 adapter
6084 .read_workspace_text("inside.txt", Some(2), Some(1))
6085 .expect("read inside"),
6086 "two"
6087 );
6088 std::fs::write(
6089 root.join("large.txt"),
6090 vec![b'x'; MAX_FILE_READ_BYTES + 1024],
6091 )
6092 .expect("large file");
6093 let bounded = adapter
6094 .read_workspace_text("large.txt", None, None)
6095 .expect("bounded read");
6096 assert!(bounded.len() <= MAX_FILE_READ_BYTES);
6097 #[cfg(unix)]
6098 {
6099 std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
6100 .expect("internal symlink");
6101 assert_eq!(
6102 adapter
6103 .read_workspace_text("inside-link", None, None)
6104 .expect("read internal symlink"),
6105 "one\ntwo\nthree\n"
6106 );
6107 }
6108 assert!(adapter.workspace_path("../codeswarm-outside").is_err());
6109 assert!(
6110 adapter
6111 .workspace_path(&outside.display().to_string())
6112 .is_err()
6113 );
6114 #[cfg(unix)]
6115 assert!(adapter.workspace_path("outside-link").is_err());
6116 #[cfg(unix)]
6117 std::fs::remove_file(link).expect("cleanup symlink");
6118 #[cfg(unix)]
6119 std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
6120 std::fs::remove_file(outside).expect("cleanup outside");
6121 std::fs::remove_dir_all(root).expect("cleanup workspace");
6122 }
6123
6124 #[tokio::test]
6125 async fn running_terminal_output_omits_exit_status_until_completion() {
6126 let root = unique_test_path("codeswarm-terminal-output", "dir");
6127 std::fs::create_dir_all(&root).expect("workspace");
6128 let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
6129 let result = adapter
6130 .terminal_create(&serde_json::json!({
6131 "command": "sh",
6132 "args": ["-c", "sleep 0.2; printf done"],
6133 "cwd": ".",
6134 }))
6135 .await
6136 .expect("terminal create");
6137 let id = result["terminalId"].as_str().expect("terminal id");
6138 let output = adapter.terminal_output(id).await.expect("terminal output");
6139 assert!(output.get("exitStatus").is_none());
6140 if let Some(terminal) = adapter.terminals.remove(id) {
6141 terminal.stop().await;
6142 }
6143 std::fs::remove_dir_all(root).expect("cleanup workspace");
6144 }
6145
6146 #[tokio::test]
6147 async fn acp_adapter_answers_workspace_read_requests() {
6148 let root =
6149 std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
6150 let _ = std::fs::remove_dir_all(&root);
6151 std::fs::create_dir_all(&root).expect("workspace");
6152 let source = root.join("inside.txt");
6153 let answer = root.join("answer.json");
6154 std::fs::write(&source, "workspace content").expect("source");
6155 let script = format!(
6156 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"}}}}'"#,
6157 source.display(),
6158 answer.display(),
6159 );
6160 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6161 adapter.start().await.expect("start ACP");
6162 assert!(matches!(
6163 adapter.next_event().await,
6164 Some(Ok(AgentEvent::Ready { .. }))
6165 ));
6166 adapter.send_prompt("read it".into()).await.expect("prompt");
6167 assert!(matches!(
6168 adapter.next_event().await,
6169 Some(Ok(AgentEvent::TurnComplete { .. }))
6170 ));
6171 let response: Value =
6172 serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
6173 .expect("response JSON");
6174 assert_eq!(response["id"], 9);
6175 assert_eq!(response["result"]["content"], "workspace content");
6176 adapter.stop().await.expect("stop ACP");
6177 std::fs::remove_dir_all(root).expect("cleanup workspace");
6178 }
6179
6180 #[tokio::test]
6181 async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
6182 let root =
6183 std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
6184 let _ = std::fs::remove_dir_all(&root);
6185 std::fs::create_dir_all(&root).expect("workspace");
6186 let create_request = serde_json::json!({
6187 "jsonrpc": "2.0",
6188 "id": 9,
6189 "method": "terminal/create",
6190 "params": {
6191 "sessionId": "s1",
6192 "command": "sh",
6193 "args": ["-c", "sleep 0.1; printf terminal-ok"],
6194 "cwd": ".",
6195 },
6196 });
6197 let wait_request = serde_json::json!({
6198 "jsonrpc": "2.0",
6199 "id": 10,
6200 "method": "terminal/wait_for_exit",
6201 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6202 });
6203 let output_request = serde_json::json!({
6204 "jsonrpc": "2.0",
6205 "id": 11,
6206 "method": "terminal/output",
6207 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6208 });
6209 let create_answer = root.join("create-answer.json");
6210 let wait_answer = root.join("wait-answer.json");
6211 let output_answer = root.join("output-answer.json");
6212 let script = format!(
6213 "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\"}}}}'",
6214 create_request,
6215 create_answer.display(),
6216 wait_request,
6217 wait_answer.display(),
6218 output_request,
6219 output_answer.display(),
6220 );
6221 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6222 adapter.start().await.expect("start ACP");
6223 assert!(matches!(
6224 adapter.next_event().await,
6225 Some(Ok(AgentEvent::Ready { .. }))
6226 ));
6227 adapter
6228 .send_prompt("run terminal".into())
6229 .await
6230 .expect("prompt");
6231 let mut saw_complete = false;
6232 for _ in 0..6 {
6233 match adapter.next_event().await {
6234 Some(Ok(AgentEvent::TurnComplete { .. })) => {
6235 saw_complete = true;
6236 break;
6237 }
6238 Some(_) => {}
6239 None => break,
6240 }
6241 }
6242 assert!(saw_complete, "terminal requests should not stall ACP");
6243 let create: Value = serde_json::from_str(
6244 &std::fs::read_to_string(&create_answer).expect("captured create response"),
6245 )
6246 .expect("create JSON");
6247 assert_eq!(create["result"]["terminalId"], "terminal-1");
6248 let output: Value = serde_json::from_str(
6249 &std::fs::read_to_string(&output_answer).expect("captured output response"),
6250 )
6251 .expect("output JSON");
6252 assert!(
6253 output["result"]["output"]
6254 .as_str()
6255 .unwrap_or_default()
6256 .contains("terminal-ok"),
6257 "output response: {output}"
6258 );
6259 adapter.stop().await.expect("stop ACP");
6260 std::fs::remove_dir_all(root).expect("cleanup workspace");
6261 }
6262
6263 #[tokio::test]
6264 async fn host_reduces_and_persists_adapter_events() {
6265 let path =
6266 std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
6267 let adapter = ScriptedAdapter::new(
6268 0,
6269 AgentCapabilities::default(),
6270 [AgentEvent::Text {
6271 slot: 0,
6272 text: "hello".into(),
6273 }],
6274 );
6275 let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
6276 host.start().await.expect("start");
6277 host.next_effects()
6278 .await
6279 .expect("event")
6280 .expect("valid event");
6281 assert_eq!(host.state.public_text[0].1, "hello");
6282 assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
6283 std::fs::remove_file(path).expect("cleanup");
6284 }
6285
6286 #[tokio::test]
6287 async fn relay_applies_default_policy_before_the_first_prompt() {
6288 let first_log = Arc::new(Mutex::new(Vec::new()));
6289 let second_log = Arc::new(Mutex::new(Vec::new()));
6290 let first = AdapterHost::new(
6291 Box::new(ModeOrderAdapter {
6292 slot: 0,
6293 log: Arc::clone(&first_log),
6294 phase: 0,
6295 }),
6296 None,
6297 );
6298 let second = AdapterHost::new(
6299 Box::new(ModeOrderAdapter {
6300 slot: 1,
6301 log: Arc::clone(&second_log),
6302 phase: 0,
6303 }),
6304 None,
6305 );
6306 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6307 relay.start().await.expect("start and synchronize policy");
6308 relay.run_turn("task", 0).await.expect("first turn");
6309 {
6310 let log = first_log.lock().expect("log");
6311 assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
6312 }
6313 assert_eq!(
6314 second_log.lock().expect("log").as_slice(),
6315 ["start", "mode:yolo"]
6316 );
6317 let added_log = Arc::new(Mutex::new(Vec::new()));
6318 relay
6319 .add_agent(
6320 AdapterHost::new(
6321 Box::new(ModeOrderAdapter {
6322 slot: 2,
6323 log: Arc::clone(&added_log),
6324 phase: 0,
6325 }),
6326 None,
6327 ),
6328 "Added",
6329 "added.example",
6330 "added-agent",
6331 )
6332 .await
6333 .expect("add with synchronized policy");
6334 assert_eq!(
6335 added_log.lock().expect("log").as_slice(),
6336 ["start", "mode:yolo"]
6337 );
6338 relay.drop_agent(2).await.expect("drop added agent");
6339 added_log.lock().expect("log").clear();
6340 relay
6341 .reload(2)
6342 .await
6343 .expect("reload with synchronized policy");
6344 assert_eq!(
6345 added_log.lock().expect("log").as_slice(),
6346 ["reload", "mode:yolo"]
6347 );
6348 }
6349
6350 #[tokio::test]
6351 async fn acp_roster_is_ready_before_any_prompt_is_sent() {
6352 let hosts = (0..2)
6353 .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
6354 .collect::<Vec<_>>();
6355 let startup_events = Arc::new(Mutex::new(Vec::new()));
6356 let captured = Arc::clone(&startup_events);
6357 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6358 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6359
6360 relay.start().await.expect("complete startup handshake");
6361
6362 assert!(relay.dispatches().is_empty());
6363 let ready_slots = startup_events
6364 .lock()
6365 .expect("events")
6366 .iter()
6367 .filter_map(|event| match event {
6368 AgentEvent::Ready { slot, .. } => Some(*slot),
6369 _ => None,
6370 })
6371 .collect::<Vec<_>>();
6372 assert_eq!(ready_slots, vec![0, 1]);
6373 }
6374
6375 #[tokio::test]
6376 async fn independent_roster_adapters_start_concurrently() {
6377 let barrier = Arc::new(tokio::sync::Barrier::new(2));
6378 let hosts = (0..2)
6379 .map(|slot| {
6380 AdapterHost::new(
6381 Box::new(ConcurrentStartAdapter {
6382 slot,
6383 barrier: Arc::clone(&barrier),
6384 }),
6385 None,
6386 )
6387 })
6388 .collect::<Vec<_>>();
6389 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6390 tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
6391 .await
6392 .expect("startup should not serialize barrier participants")
6393 .expect("startup succeeds");
6394 }
6395
6396 #[tokio::test]
6397 async fn relay_host_dispatches_turns_sequentially() {
6398 let capabilities = AgentCapabilities {
6399 supports_cancel: true,
6400 ..AgentCapabilities::default()
6401 };
6402 let first = ScriptedAdapter::new(
6403 0,
6404 capabilities.clone(),
6405 [
6406 AgentEvent::Text {
6407 slot: 0,
6408 text: "first".into(),
6409 },
6410 AgentEvent::TurnComplete { slot: 0 },
6411 ],
6412 );
6413 let second = ScriptedAdapter::new(
6414 1,
6415 capabilities,
6416 [
6417 AgentEvent::Text {
6418 slot: 1,
6419 text: "review".into(),
6420 },
6421 AgentEvent::TurnComplete { slot: 1 },
6422 ],
6423 );
6424 let hosts = vec![
6425 AdapterHost::new(Box::new(first), None),
6426 AdapterHost::new(Box::new(second), None),
6427 ];
6428 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6429 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6430 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6431 let captured = std::sync::Arc::clone(&events);
6432 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6433 relay.start().await.expect("start");
6434 events.lock().expect("events").clear();
6435 assert!(matches!(
6436 relay.run_turn("task", 0).await.expect("first turn"),
6437 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6438 ));
6439 assert!(matches!(
6440 relay.run_turn("first", 0).await.expect("second turn"),
6441 crate::relay::RelayDecision::Dispatch {
6442 slot: 1,
6443 can_stop: true,
6444 ..
6445 }
6446 ));
6447 assert_eq!(
6448 relay
6449 .dispatches()
6450 .iter()
6451 .map(|(slot, _)| *slot)
6452 .collect::<Vec<_>>(),
6453 [0, 1]
6454 );
6455 assert!(relay.dispatches()[0].1.contains("You are Claude"));
6456 assert!(
6457 relay.dispatches()[0]
6458 .1
6459 .contains("CodeSwarm roster (ordered)")
6460 );
6461 assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
6462 assert!(relay.dispatches()[0].1.contains("2. Codex"));
6463 assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
6464 assert!(
6465 relay.dispatches()[1]
6466 .1
6467 .contains("stops all other agents and ends the entire automated relay")
6468 );
6469 assert!(relay.dispatches()[1].1.contains("Use it with extreme care"));
6470 assert!(
6471 relay.dispatches()[1]
6472 .1
6473 .contains("If there is any uncertainty, do not use it")
6474 );
6475 assert!(relay.dispatches()[0].1.contains("Do not use"));
6476 let lifecycle = events.lock().expect("events");
6477 let positions = lifecycle
6478 .iter()
6479 .filter_map(|event| match event {
6480 AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
6481 AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
6482 _ => None,
6483 })
6484 .collect::<Vec<_>>();
6485 assert_eq!(
6486 positions,
6487 [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
6488 );
6489 }
6490
6491 #[tokio::test]
6492 async fn failed_resume_does_not_stop_or_dispatch_to_healthy_peer() {
6493 let stops = Arc::new(AtomicUsize::new(0));
6494 let failed = FailingStartAdapter {
6495 slot: 0,
6496 stopped: stops.clone(),
6497 };
6498 let healthy = ScriptedAdapter::new(
6499 1,
6500 AgentCapabilities::default(),
6501 [
6502 AgentEvent::Text {
6503 slot: 1,
6504 text: "healthy response".into(),
6505 },
6506 AgentEvent::TurnComplete { slot: 1 },
6507 ],
6508 );
6509 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6510 let captured = events.clone();
6511 let mut relay = RelayHost::new(
6512 vec![
6513 AdapterHost::new(Box::new(failed), None),
6514 AdapterHost::new(Box::new(healthy), None),
6515 ],
6516 4,
6517 )
6518 .unwrap();
6519 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6520 relay.start_resuming().await.unwrap();
6521 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
6522 assert!(relay.dispatches().is_empty());
6523 assert_eq!(stops.load(Ordering::Relaxed), 1);
6524 assert!(
6525 events
6526 .lock()
6527 .unwrap()
6528 .iter()
6529 .any(|event| matches!(event, AgentEvent::Failed { slot: 0, .. }))
6530 );
6531 assert!(
6532 !events
6533 .lock()
6534 .unwrap()
6535 .iter()
6536 .any(|event| matches!(event, AgentEvent::Failed { slot: 1, .. }))
6537 );
6538 assert!(!relay.relay_mut().enqueue_human("do not retarget", Some(0)));
6539 assert!(
6540 relay
6541 .relay_mut()
6542 .enqueue_human("explicit healthy target", Some(1))
6543 );
6544 assert!(matches!(
6545 relay.run_turn("", 1).await.unwrap(),
6546 RelayDecision::Dispatch { slot: 1, .. }
6547 ));
6548 relay.stop().await.unwrap();
6549 }
6550
6551 #[tokio::test]
6552 async fn pair_strategy_wires_roles_into_non_direct_prompts() {
6553 let capabilities = AgentCapabilities::default();
6554 let first = ScriptedAdapter::new(
6555 0,
6556 capabilities.clone(),
6557 [
6558 AgentEvent::Text {
6559 slot: 0,
6560 text: "implemented".into(),
6561 },
6562 AgentEvent::TurnComplete { slot: 0 },
6563 AgentEvent::Text {
6564 slot: 0,
6565 text: format!("fixed review findings {STOP_TOKEN}"),
6566 },
6567 AgentEvent::TurnComplete { slot: 0 },
6568 ],
6569 );
6570 let second = ScriptedAdapter::new(
6571 1,
6572 capabilities,
6573 [
6574 AgentEvent::Text {
6575 slot: 1,
6576 text: "reviewed".into(),
6577 },
6578 AgentEvent::TurnComplete { slot: 1 },
6579 AgentEvent::Text {
6580 slot: 1,
6581 text: format!("approved {STOP_TOKEN}"),
6582 },
6583 AgentEvent::TurnComplete { slot: 1 },
6584 AgentEvent::Text {
6585 slot: 1,
6586 text: "new task".into(),
6587 },
6588 AgentEvent::TurnComplete { slot: 1 },
6589 ],
6590 );
6591 let hosts = vec![
6592 AdapterHost::new(Box::new(first), None),
6593 AdapterHost::new(Box::new(second), None),
6594 ];
6595 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6596 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6597 relay.relay_mut().set_strategy(CollaborationStrategy::Pair);
6598 relay.start().await.expect("start");
6599 assert!(matches!(
6600 relay.run_turn("task", 0).await.expect("implementer turn"),
6601 RelayDecision::Dispatch {
6602 slot: 0,
6603 can_stop: false,
6604 ..
6605 }
6606 ));
6607 let implementer_prompt = &relay.dispatches()[0].1;
6608 assert!(implementer_prompt.contains("you are the implementer"));
6609 assert!(implementer_prompt.contains("pair reviewer will review the result next"));
6610 assert!(!implementer_prompt.contains("you are the reviewer"));
6611 assert!(implementer_prompt.contains("Do not use"));
6612 assert!(matches!(
6613 relay.run_turn("", 0).await.expect("reviewer turn"),
6614 RelayDecision::Dispatch {
6615 slot: 1,
6616 can_stop: true,
6617 ..
6618 }
6619 ));
6620 let reviewer_prompt = &relay.dispatches()[1].1;
6621 assert!(reviewer_prompt.contains("you are the reviewer"));
6622 assert!(reviewer_prompt.contains("Claude handed off"));
6623 assert!(reviewer_prompt.contains("concrete defects"));
6624 assert!(reviewer_prompt.contains("concise approval"));
6625 assert!(reviewer_prompt.contains(STOP_TOKEN));
6626 assert!(!reviewer_prompt.contains("you are the implementer"));
6627 relay.run_turn("", 0).await.unwrap();
6628 assert!(relay.dispatches()[2].1.contains("you are the implementer"));
6629 assert!(relay.dispatches()[2].1.contains("Do not use"));
6630 assert!(matches!(
6631 relay.run_turn("", 0).await.unwrap(),
6632 RelayDecision::Dispatch { slot: 1, .. }
6633 ));
6634 assert!(relay.dispatches()[3].1.contains("you are the reviewer"));
6635 assert!(relay.relay_mut().enqueue_human("new task", Some(1)));
6636 relay.run_turn("", 1).await.unwrap();
6637 assert!(relay.dispatches()[4].1.contains("you are the implementer"));
6638 }
6639
6640 #[tokio::test]
6641 async fn solo_roster_and_direct_prompts_omit_pair_roles() {
6642 let solo = ScriptedAdapter::new(
6643 0,
6644 AgentCapabilities::default(),
6645 [
6646 AgentEvent::Text {
6647 slot: 0,
6648 text: "solo".into(),
6649 },
6650 AgentEvent::TurnComplete { slot: 0 },
6651 ],
6652 );
6653 let mut solo_relay =
6654 RelayHost::new(vec![AdapterHost::new(Box::new(solo), None)], 4).expect("relay");
6655 solo_relay
6656 .relay_mut()
6657 .set_strategy(CollaborationStrategy::Pair);
6658 solo_relay.start().await.expect("start");
6659 solo_relay.run_turn("task", 0).await.expect("solo turn");
6660 assert!(!solo_relay.dispatches()[0].1.contains("Pair role"));
6661
6662 let roster_first = ScriptedAdapter::new(
6663 0,
6664 AgentCapabilities::default(),
6665 [AgentEvent::TurnComplete { slot: 0 }],
6666 );
6667 let roster_second = ScriptedAdapter::new(
6668 1,
6669 AgentCapabilities::default(),
6670 [AgentEvent::TurnComplete { slot: 1 }],
6671 );
6672 let mut roster = RelayHost::new(
6673 vec![
6674 AdapterHost::new(Box::new(roster_first), None),
6675 AdapterHost::new(Box::new(roster_second), None),
6676 ],
6677 4,
6678 )
6679 .expect("relay");
6680 roster.start().await.expect("start");
6681 roster.run_turn("task", 0).await.expect("first turn");
6682 roster.run_turn("", 0).await.expect("second turn");
6683 assert!(!roster.dispatches()[0].1.contains("Pair role"));
6684 assert!(!roster.dispatches()[1].1.contains("Pair role"));
6685
6686 let pair_first = ScriptedAdapter::new(
6687 0,
6688 AgentCapabilities::default(),
6689 [AgentEvent::TurnComplete { slot: 0 }],
6690 );
6691 let pair_second = ScriptedAdapter::new(
6692 1,
6693 AgentCapabilities::default(),
6694 [AgentEvent::TurnComplete { slot: 1 }],
6695 );
6696 let mut pair = RelayHost::new(
6697 vec![
6698 AdapterHost::new(Box::new(pair_first), None),
6699 AdapterHost::new(Box::new(pair_second), None),
6700 ],
6701 4,
6702 )
6703 .expect("relay");
6704 pair.relay_mut().set_strategy(CollaborationStrategy::Pair);
6705 assert_eq!(pair.relay_mut().enqueue_direct(1, "private"), Ok(true));
6706 pair.start().await.expect("start");
6707 assert!(matches!(
6708 pair.run_turn("ignored", 0).await.expect("direct turn"),
6709 RelayDecision::Dispatch {
6710 slot: 1,
6711 direct: true,
6712 ..
6713 }
6714 ));
6715 let direct_prompt = &pair.dispatches()[0].1;
6716 assert!(direct_prompt.contains("private"));
6717 assert!(!direct_prompt.contains("Pair role"));
6718 }
6719
6720 #[tokio::test]
6721 async fn relay_host_routes_around_a_usage_limited_agent() {
6722 let capabilities = AgentCapabilities::default();
6723 let first = ScriptedAdapter::new(
6724 0,
6725 capabilities.clone(),
6726 [
6727 AgentEvent::Text {
6728 slot: 0,
6729 text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
6730 credits or try again later."
6731 .into(),
6732 },
6733 AgentEvent::TurnComplete { slot: 0 },
6734 ],
6735 );
6736 let second = ScriptedAdapter::new(
6737 1,
6738 capabilities,
6739 [
6740 AgentEvent::Text {
6741 slot: 1,
6742 text: "review done".into(),
6743 },
6744 AgentEvent::TurnComplete { slot: 1 },
6745 ],
6746 );
6747 let hosts = vec![
6748 AdapterHost::new(Box::new(first), None),
6749 AdapterHost::new(Box::new(second), None),
6750 ];
6751 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6752 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6753 let captured = std::sync::Arc::clone(&events);
6754 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6755 relay.start().await.expect("start");
6756 events.lock().expect("events").clear();
6757 assert!(matches!(
6758 relay.run_turn("task", 0).await.expect("limited turn"),
6759 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6760 ));
6761 assert!(
6762 events
6763 .lock()
6764 .expect("events")
6765 .iter()
6766 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6767 );
6768 assert!(matches!(
6770 relay.run_turn("", 0).await.expect("next turn"),
6771 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6772 ));
6773 assert!(relay.relay().is_limited(0));
6774 relay.reload(0).await.expect("reload");
6778 assert!(!relay.relay().is_limited(0));
6779 }
6780
6781 #[tokio::test]
6782 async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
6783 let limited = ScriptedAdapter::new(
6784 0,
6785 AgentCapabilities::default(),
6786 [AgentEvent::Failed {
6787 slot: 0,
6788 started: true,
6789 detail: "request failed: insufficient_quota".into(),
6790 }],
6791 );
6792 let healthy = ScriptedAdapter::new(
6793 1,
6794 AgentCapabilities::default(),
6795 [AgentEvent::TurnComplete { slot: 1 }],
6796 );
6797 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6798 let captured = std::sync::Arc::clone(&events);
6799 let mut relay = RelayHost::new(
6800 vec![
6801 AdapterHost::new(Box::new(limited), None),
6802 AdapterHost::new(Box::new(healthy), None),
6803 ],
6804 4,
6805 )
6806 .expect("relay");
6807 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6808 relay.start().await.expect("start");
6809 events.lock().expect("events").clear();
6810
6811 assert!(matches!(
6812 relay.run_turn("task", 0).await.expect("limited failure"),
6813 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6814 ));
6815 assert!(relay.relay().is_limited(0));
6816 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
6817 {
6818 let events = events.lock().expect("events");
6819 assert!(
6820 events
6821 .iter()
6822 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6823 );
6824 assert!(
6825 !events
6826 .iter()
6827 .any(|event| matches!(event, AgentEvent::Failed { .. }))
6828 );
6829 }
6830
6831 assert!(matches!(
6832 relay.run_turn("", 0).await.expect("healthy peer"),
6833 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6834 ));
6835 }
6836
6837 #[tokio::test]
6838 async fn relay_failure_is_skipped_for_one_batch_without_changing_the_roster() {
6839 let failed = ScriptedAdapter::new(
6840 0,
6841 AgentCapabilities::default(),
6842 [AgentEvent::Failed {
6843 slot: 0,
6844 started: true,
6845 detail: "connection lost".into(),
6846 }],
6847 );
6848 let healthy = ScriptedAdapter::new(
6849 1,
6850 AgentCapabilities::default(),
6851 [AgentEvent::TurnComplete { slot: 1 }],
6852 );
6853 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6854 let captured = std::sync::Arc::clone(&events);
6855 let mut relay = RelayHost::new(
6856 vec![
6857 AdapterHost::new(Box::new(failed), None),
6858 AdapterHost::new(Box::new(healthy), None),
6859 ],
6860 4,
6861 )
6862 .expect("relay");
6863 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6864 relay.start().await.expect("start");
6865
6866 assert!(matches!(
6867 relay.run_turn("task", 0).await.expect("handled failure"),
6868 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6869 ));
6870 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![0, 1]);
6871 assert!(relay.relay().is_limited(0));
6872 assert!(events.lock().expect("lock").iter().any(|event| {
6873 matches!(
6874 event,
6875 AgentEvent::Failed {
6876 slot: 0,
6877 started: true,
6878 ..
6879 }
6880 )
6881 }));
6882 assert!(matches!(
6883 relay.run_turn("", 0).await.expect("healthy peer"),
6884 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6885 ));
6886 }
6887
6888 #[tokio::test]
6889 async fn codex_stop_does_not_skip_later_roster_reviewers() {
6890 let hosts = (0..3)
6891 .map(|slot| {
6892 AdapterHost::new(
6893 Box::new(ScriptedAdapter::new(
6894 slot,
6895 AgentCapabilities::default(),
6896 [
6897 AgentEvent::Text {
6898 slot,
6899 text: STOP_TOKEN.into(),
6900 },
6901 AgentEvent::TurnComplete { slot },
6902 ],
6903 )),
6904 None,
6905 )
6906 })
6907 .collect();
6908 let mut relay = RelayHost::new(hosts, 10).expect("relay");
6909 relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
6910 relay.start().await.expect("start");
6911 for expected in 0..3 {
6912 assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
6913 RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
6914 }
6915 assert_eq!(
6916 relay.run_turn("", 0).await.expect("complete"),
6917 RelayDecision::Complete
6918 );
6919 }
6920
6921 #[tokio::test]
6922 async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
6923 let first = ScriptedAdapter::new(
6924 0,
6925 AgentCapabilities::default(),
6926 [
6927 AgentEvent::Text {
6928 slot: 0,
6929 text: "done".into(),
6930 },
6931 AgentEvent::TurnComplete { slot: 0 },
6932 ],
6933 );
6934 let reviewer = ScriptedAdapter::new(
6935 1,
6936 AgentCapabilities::default(),
6937 [
6938 AgentEvent::Text {
6939 slot: 1,
6940 text: STOP_TOKEN.into(),
6941 },
6942 AgentEvent::TurnComplete { slot: 1 },
6943 ],
6944 );
6945 let mut relay = RelayHost::new(
6946 vec![
6947 AdapterHost::new(Box::new(first), None),
6948 AdapterHost::new(Box::new(reviewer), None),
6949 ],
6950 10,
6951 )
6952 .expect("relay");
6953 relay.start().await.expect("start");
6954 let first_decision = relay.run_turn("task", 0).await.expect("first");
6955 assert!(matches!(
6956 first_decision,
6957 RelayDecision::Dispatch { slot: 0, .. }
6958 ));
6959 let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
6960 assert!(matches!(
6961 reviewer_decision,
6962 RelayDecision::Dispatch {
6963 slot: 1,
6964 can_stop: true,
6965 ..
6966 }
6967 ));
6968 assert_eq!(
6969 relay.run_turn("", 0).await.expect("complete"),
6970 RelayDecision::Complete
6971 );
6972 }
6973
6974 #[tokio::test]
6975 async fn relay_stream_emits_text_and_thought_endings_before_tools() {
6976 let tool = AgentEvent::Tool {
6977 slot: 0,
6978 update: crate::ToolUpdate {
6979 id: "read".into(),
6980 title: "Read file".into(),
6981 status: ToolStatus::Running,
6982 detail: None,
6983 },
6984 };
6985 let updates = vec![
6986 AgentEvent::Thought {
6987 slot: 0,
6988 text: "Check the buffer. ✈".into(),
6989 },
6990 AgentEvent::Text {
6991 slot: 0,
6992 text: "Let me check.".into(),
6993 },
6994 tool.clone(),
6995 AgentEvent::Text {
6996 slot: 0,
6997 text: "[CODE".into(),
6998 },
6999 AgentEvent::Text {
7000 slot: 0,
7001 text: " is ordinary.".into(),
7002 },
7003 AgentEvent::Text {
7004 slot: 0,
7005 text: "[CODESWARM:".into(),
7006 },
7007 AgentEvent::Text {
7008 slot: 0,
7009 text: "STOP] Done.".into(),
7010 },
7011 tool,
7012 AgentEvent::TurnComplete { slot: 0 },
7013 ];
7014 let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
7015 let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
7016 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7017 let captured = std::sync::Arc::clone(&events);
7018 let mut relay = RelayHost::new(
7019 vec![
7020 AdapterHost::new(Box::new(first), None),
7021 AdapterHost::new(Box::new(reviewer), None),
7022 ],
7023 2,
7024 )
7025 .expect("relay");
7026 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
7027 relay.start().await.unwrap();
7028 relay.run_turn("task", 0).await.unwrap();
7029 let captured = events.lock().unwrap();
7030 let visible: Vec<_> = captured
7031 .iter()
7032 .filter(|event| {
7033 matches!(
7034 event,
7035 AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
7036 )
7037 })
7038 .cloned()
7039 .collect();
7040 assert_eq!(
7041 visible,
7042 vec![
7043 updates[0].clone(),
7044 updates[1].clone(),
7045 updates[2].clone(),
7046 AgentEvent::Text {
7047 slot: 0,
7048 text: "[CODE is ordinary.".into()
7049 },
7050 AgentEvent::Text {
7051 slot: 0,
7052 text: " Done.".into()
7053 },
7054 updates[7].clone(),
7055 ]
7056 );
7057 }
7058
7059 #[tokio::test]
7060 async fn roster_handoff_routes_only_terminal_message_markers_and_refreshes_targets() {
7061 let text = |value: &str| AgentEvent::Text {
7062 slot: 0,
7063 text: value.into(),
7064 };
7065 let thought = || AgentEvent::Thought {
7066 slot: 0,
7067 text: "still checking".into(),
7068 };
7069 let tool = || AgentEvent::Tool {
7070 slot: 0,
7071 update: crate::ToolUpdate {
7072 id: "read".into(),
7073 title: "Read file".into(),
7074 status: ToolStatus::Running,
7075 detail: None,
7076 },
7077 };
7078 let cases = vec![
7079 (vec![text("result [CODESWARM:NEXT:3]\n ")], 2),
7080 (
7081 vec![text("result [CODESWARM:"), text("NEXT:"), text("3]")],
7082 2,
7083 ),
7084 (vec![text("result [CODESWARM:NEXT:3]"), text(" more")], 1),
7085 (vec![text("result [CODESWARM:NEXT:3]"), thought()], 1),
7086 (
7087 vec![text("result [CODESWARM:NEXT:3]"), tool(), text(" ")],
7088 1,
7089 ),
7090 (
7091 vec![text("result [CODESWARM:NEXT:"), thought(), text("3]")],
7092 1,
7093 ),
7094 (
7095 vec![text("result"), thought(), text("[CODESWARM:NEXT:3]")],
7096 2,
7097 ),
7098 (vec![text("result [CODESWARM:NEXT:1]")], 1),
7099 (vec![text("result [CODESWARM:NEXT:0]")], 1),
7100 (vec![text("result [CODESWARM:NEXT:99]")], 1),
7101 (
7102 vec![
7103 text("result [CODESWARM:NEXT:3]"),
7104 AgentEvent::UsageUpdated {
7105 slot: 0,
7106 usage: crate::UsageUpdate { used: 1, size: 100 },
7107 },
7108 ],
7109 2,
7110 ),
7111 ];
7112 for (mut updates, expected) in cases {
7113 updates.push(AgentEvent::TurnComplete { slot: 0 });
7114 let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates);
7115 let hosts = std::iter::once(AdapterHost::new(Box::new(first), None))
7116 .chain((1..3).map(|slot| {
7117 AdapterHost::new(
7118 Box::new(ScriptedAdapter::new(
7119 slot,
7120 AgentCapabilities::default(),
7121 [AgentEvent::TurnComplete { slot }],
7122 )),
7123 None,
7124 )
7125 }))
7126 .collect();
7127 let mut relay = RelayHost::new(hosts, 10).unwrap();
7128 relay.set_roster_names(vec!["Worker".into(), "Codex".into(), "Codex".into()]);
7129 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
7130 let captured = events.clone();
7131 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
7132 relay.start().await.unwrap();
7133 relay.run_turn("task", 0).await.unwrap();
7134 let prompt = &relay.dispatches()[0].1;
7135 assert!(prompt.contains("[CODESWARM:NEXT:2] → Codex"));
7136 assert!(prompt.contains("[CODESWARM:NEXT:3] → Codex"));
7137 assert!(!prompt.contains("[CODESWARM:NEXT:1]"));
7138 relay.introduced.fill(true);
7140 relay.set_roster_names(vec!["Replacement".into(), "Codex".into(), "Codex".into()]);
7141 let next = relay.run_turn("", 0).await.unwrap();
7142 assert!(
7143 matches!(next, RelayDecision::Dispatch { slot, can_stop: false, .. } if slot == expected),
7144 "{next:?}"
7145 );
7146 let prompt = &relay.dispatches()[1].1;
7147 assert!(prompt.contains("[CODESWARM:NEXT:1] → Replacement"));
7148 assert!(prompt.contains("result"));
7149 let public = prompt
7150 .split("Public updates:\n")
7151 .nth(1)
7152 .unwrap()
7153 .split("\n\nDo not use")
7154 .next()
7155 .unwrap();
7156 assert!(!public.contains("[CODESWARM:NEXT:"));
7157 let visible = events
7158 .lock()
7159 .unwrap()
7160 .iter()
7161 .filter_map(|event| match event {
7162 AgentEvent::Text { text, .. } => Some(text.clone()),
7163 _ => None,
7164 })
7165 .collect::<String>();
7166 assert!(visible.contains("result"));
7167 assert!(!visible.contains("[CODESWARM:"), "{visible}");
7168 }
7169 }
7170
7171 #[tokio::test]
7172 async fn reviewer_stop_requires_a_terminal_marker_after_all_activity() {
7173 let text = |value: &str| AgentEvent::Text {
7174 slot: 1,
7175 text: value.into(),
7176 };
7177 let thought = || AgentEvent::Thought {
7178 slot: 1,
7179 text: "still checking".into(),
7180 };
7181 let tool = || AgentEvent::Tool {
7182 slot: 1,
7183 update: crate::ToolUpdate {
7184 id: "read".into(),
7185 title: "Read file".into(),
7186 status: ToolStatus::Running,
7187 detail: None,
7188 },
7189 };
7190 let cases = vec![
7191 (vec![text(&format!("done {STOP_TOKEN}"))], true),
7192 (vec![text(STOP_TOKEN), text("\n ")], true),
7193 (vec![text(STOP_TOKEN), text(" actually keep going")], false),
7194 (vec![text(STOP_TOKEN), thought()], false),
7195 (vec![text(STOP_TOKEN), tool()], false),
7196 (vec![text(STOP_TOKEN), tool(), text(" ")], false),
7197 (vec![text(STOP_TOKEN), tool(), text(STOP_TOKEN)], true),
7198 (vec![text("[CODESWARM:"), text("STOP]")], true),
7199 (vec![text("[CODESWARM:"), thought(), text("STOP]")], false),
7200 (
7201 vec![AgentEvent::Thought {
7202 slot: 1,
7203 text: STOP_TOKEN.into(),
7204 }],
7205 false,
7206 ),
7207 (
7208 vec![
7209 text(STOP_TOKEN),
7210 AgentEvent::UsageUpdated {
7211 slot: 1,
7212 usage: crate::UsageUpdate { used: 1, size: 100 },
7213 },
7214 ],
7215 true,
7216 ),
7217 ];
7218 for (mut events, stop) in cases {
7219 let first = ScriptedAdapter::new(
7220 0,
7221 AgentCapabilities::default(),
7222 [
7223 AgentEvent::Text {
7224 slot: 0,
7225 text: "initial response".into(),
7226 },
7227 AgentEvent::TurnComplete { slot: 0 },
7228 AgentEvent::TurnComplete { slot: 0 },
7229 ],
7230 );
7231 events.push(AgentEvent::TurnComplete { slot: 1 });
7232 let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), events.clone());
7233 let mut relay = RelayHost::new(
7234 vec![
7235 AdapterHost::new(Box::new(first), None),
7236 AdapterHost::new(Box::new(reviewer), None),
7237 ],
7238 4,
7239 )
7240 .unwrap();
7241 relay.start().await.unwrap();
7242 relay.run_turn("task", 0).await.unwrap();
7243 relay.run_turn("", 0).await.unwrap();
7244 let next = relay.run_turn("", 0).await.unwrap();
7245 assert_eq!(
7246 matches!(next, RelayDecision::Complete),
7247 stop,
7248 "events={events:?}"
7249 );
7250 relay.stop().await.unwrap();
7251 }
7252 }
7253
7254 #[tokio::test]
7255 async fn stop_token_is_filtered_from_streamed_ui_events() {
7256 let first = ScriptedAdapter::new(
7257 0,
7258 AgentCapabilities::default(),
7259 [
7260 AgentEvent::Text {
7261 slot: 0,
7262 text: format!("visible {STOP_TOKEN} trailing"),
7263 },
7264 AgentEvent::TurnComplete { slot: 0 },
7265 ],
7266 );
7267 let reviewer = ScriptedAdapter::new(
7268 1,
7269 AgentCapabilities::default(),
7270 [AgentEvent::TurnComplete { slot: 1 }],
7271 );
7272 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7273 let captured = std::sync::Arc::clone(&events);
7274 let mut relay = RelayHost::new(
7275 vec![
7276 AdapterHost::new(Box::new(first), None),
7277 AdapterHost::new(Box::new(reviewer), None),
7278 ],
7279 2,
7280 )
7281 .expect("relay");
7282 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7283 relay.start().await.expect("start");
7284 relay.run_turn("task", 0).await.expect("turn");
7285 let captured = events.lock().expect("lock");
7286 assert!(captured.iter().all(|event| match event {
7287 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7288 _ => true,
7289 }));
7290 let visible = captured
7291 .iter()
7292 .filter_map(|event| match event {
7293 AgentEvent::Text { text, .. } => Some(text.as_str()),
7294 _ => None,
7295 })
7296 .collect::<String>();
7297 assert_eq!(visible, "visible trailing");
7298 }
7299
7300 #[tokio::test]
7301 async fn token_only_reviewer_response_emits_visible_acknowledgment() {
7302 let first = ScriptedAdapter::new(
7303 0,
7304 AgentCapabilities::default(),
7305 [
7306 AgentEvent::Text {
7307 slot: 0,
7308 text: "done".into(),
7309 },
7310 AgentEvent::TurnComplete { slot: 0 },
7311 ],
7312 );
7313 let reviewer = ScriptedAdapter::new(
7314 1,
7315 AgentCapabilities::default(),
7316 [
7317 AgentEvent::Text {
7318 slot: 1,
7319 text: STOP_TOKEN.into(),
7320 },
7321 AgentEvent::TurnComplete { slot: 1 },
7322 ],
7323 );
7324 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7325 let captured = std::sync::Arc::clone(&events);
7326 let mut relay = RelayHost::new(
7327 vec![
7328 AdapterHost::new(Box::new(first), None),
7329 AdapterHost::new(Box::new(reviewer), None),
7330 ],
7331 4,
7332 )
7333 .expect("relay");
7334 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7335 relay.start().await.expect("start");
7336 relay.run_turn("task", 0).await.expect("first turn");
7337 relay.run_turn("", 0).await.expect("review turn");
7338 let captured = events.lock().expect("lock");
7339 assert!(captured.iter().any(|event| {
7340 matches!(
7341 event,
7342 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7343 )
7344 }));
7345 assert!(captured.iter().all(|event| match event {
7346 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7347 _ => true,
7348 }));
7349 let acknowledgment = captured
7350 .iter()
7351 .position(|event| {
7352 matches!(
7353 event,
7354 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7355 )
7356 })
7357 .expect("visible acknowledgment");
7358 let completion = captured
7359 .iter()
7360 .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
7361 .expect("reviewer completion");
7362 assert!(acknowledgment < completion);
7363 }
7364
7365 #[tokio::test]
7366 async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
7367 let first = ScriptedAdapter::new(
7368 0,
7369 AgentCapabilities::default(),
7370 [
7371 AgentEvent::Text {
7372 slot: 0,
7373 text: "done".into(),
7374 },
7375 AgentEvent::TurnComplete { slot: 0 },
7376 ],
7377 );
7378 let reviewer = ScriptedAdapter::new(
7379 1,
7380 AgentCapabilities::default(),
7381 [
7382 AgentEvent::Text {
7383 slot: 1,
7384 text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
7385 },
7386 AgentEvent::TurnComplete { slot: 1 },
7387 ],
7388 );
7389 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7390 let captured = std::sync::Arc::clone(&events);
7391 let mut relay = RelayHost::new(
7392 vec![
7393 AdapterHost::new(Box::new(first), None),
7394 AdapterHost::new(Box::new(reviewer), None),
7395 ],
7396 4,
7397 )
7398 .expect("relay");
7399 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7400 relay.start().await.expect("start");
7401 relay.run_turn("task", 0).await.expect("first turn");
7402 relay.run_turn("", 0).await.expect("review turn");
7403
7404 let visible = events
7405 .lock()
7406 .expect("lock")
7407 .iter()
7408 .filter_map(|event| match event {
7409 AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
7410 _ => None,
7411 })
7412 .collect::<String>();
7413 assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
7414 assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
7415 }
7416
7417 #[tokio::test]
7418 async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
7419 let first = AdapterHost::new(
7420 Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
7421 None,
7422 );
7423 let second = AdapterHost::new(
7424 Box::new(ScriptedAdapter::new(
7425 1,
7426 AgentCapabilities::default(),
7427 [AgentEvent::TurnComplete { slot: 1 }],
7428 )),
7429 None,
7430 );
7431 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7432 let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
7433 relay.set_event_sink(move |event| {
7434 if matches!(event, AgentEvent::Permission { .. }) {
7435 let _ = seen_sender.send(());
7436 }
7437 });
7438 relay.start().await.expect("start");
7439 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
7440 let answer = async move {
7441 seen_receiver.recv().await.expect("permission request");
7442 sender
7443 .send(super::RelayPermissionAnswer {
7444 slot: 0,
7445 request_id: "permission-1".into(),
7446 answer: PermissionAnswer::Selected {
7447 option_id: "allow".into(),
7448 },
7449 })
7450 .expect("queue permission answer");
7451 };
7452 tokio::time::timeout(std::time::Duration::from_millis(100), async {
7453 let ((), result) = tokio::join!(
7454 answer,
7455 relay.run_turn_with_permissions("task", 0, &mut receiver)
7456 );
7457 result
7458 })
7459 .await
7460 .expect("permission-gated turn should not deadlock")
7461 .expect("turn completes");
7462 }
7463
7464 #[tokio::test]
7465 async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
7466 let first = AdapterHost::new(
7467 Box::new(PendingAdapter {
7468 slot: 0,
7469 hang_on_cancel: false,
7470 }),
7471 None,
7472 );
7473 let second = AdapterHost::new(
7474 Box::new(ScriptedAdapter::new(
7475 1,
7476 AgentCapabilities::default(),
7477 [AgentEvent::TurnComplete { slot: 1 }],
7478 )),
7479 None,
7480 );
7481 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7482 relay.start().await.expect("start");
7483 let cancellation = relay.cancellation();
7484 let error = {
7485 let turn = relay.run_turn("task", 0);
7486 tokio::pin!(turn);
7487 cancellation.request();
7488 turn.await.expect_err("cancellation should stop turn")
7489 };
7490 assert!(error.to_string().contains("relay turn cancelled"));
7491
7492 assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
7493 relay
7494 .run_turn("", 1)
7495 .await
7496 .expect("replacement job reaches the selected peer");
7497 let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
7498 assert!(replacement.contains("replacement job"));
7499 assert!(replacement.contains("User "));
7500 assert!(replacement.contains(":\ntask"));
7501 let owner_updates = relay.relay_mut().unseen_context(0);
7502 assert!(owner_updates.contains("User "));
7503 assert!(owner_updates.contains(":\ntask"));
7504 assert!(owner_updates.contains(":\nreplacement job"));
7505 }
7506
7507 #[tokio::test]
7508 async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
7509 let first = AdapterHost::new(
7510 Box::new(PendingAdapter {
7511 slot: 0,
7512 hang_on_cancel: true,
7513 }),
7514 None,
7515 );
7516 let second = AdapterHost::new(
7517 Box::new(PendingAdapter {
7518 slot: 1,
7519 hang_on_cancel: false,
7520 }),
7521 None,
7522 );
7523 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7524 relay.start().await.expect("start");
7525 let cancellation = relay.cancellation();
7526 let turn = relay.run_turn("task", 0);
7527 tokio::pin!(turn);
7528 cancellation.request();
7529 let error = turn.await.expect_err("cancellation should stop turn");
7530 assert!(error.to_string().contains("timed out"));
7531 }
7532
7533 #[tokio::test]
7534 async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
7535 let event = [AgentEvent::TurnComplete { slot: 0 }];
7536 let first = AdapterHost::new(
7537 Box::new(ScriptedAdapter::new(
7538 0,
7539 AgentCapabilities::default(),
7540 event.clone(),
7541 )),
7542 None,
7543 );
7544 let second = AdapterHost::new(
7545 Box::new(ScriptedAdapter::new(
7546 1,
7547 AgentCapabilities::default(),
7548 [AgentEvent::TurnComplete { slot: 1 }],
7549 )),
7550 None,
7551 );
7552 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7553 relay.start().await.expect("start");
7554
7555 relay.pause();
7556 assert_eq!(
7557 relay.run_turn("paused", 0).await.expect("paused turn"),
7558 crate::relay::RelayDecision::Paused
7559 );
7560 assert!(relay.dispatches().is_empty());
7561
7562 relay.resume();
7563 relay.relay_mut().drop_agent(1).expect("drop reviewer");
7564 assert!(matches!(
7565 relay
7566 .run_turn("solo follow-up", 0)
7567 .await
7568 .expect("solo turn"),
7569 crate::relay::RelayDecision::Dispatch {
7570 slot: 0,
7571 can_stop: false,
7572 ..
7573 }
7574 ));
7575 assert_eq!(relay.dispatches().len(), 1);
7576 }
7577
7578 #[tokio::test]
7579 async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
7580 let first = AdapterHost::new(
7581 Box::new(ScriptedAdapter::new(
7582 0,
7583 AgentCapabilities::default(),
7584 [AgentEvent::TurnComplete { slot: 0 }],
7585 )),
7586 None,
7587 );
7588 let second = AdapterHost::new(
7589 Box::new(ScriptedAdapter::new(
7590 1,
7591 AgentCapabilities::default(),
7592 [AgentEvent::TurnComplete { slot: 1 }],
7593 )),
7594 None,
7595 );
7596 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7597 relay.set_roster_names(vec!["First".into(), "Second".into()]);
7598 relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
7599 relay.set_roster_launch_specs(vec![
7600 ("custom".into(), "owner".into()),
7601 ("custom".into(), "peer".into()),
7602 ]);
7603 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7604 let captured = std::sync::Arc::clone(&events);
7605 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7606 relay.start().await.expect("start");
7607 let slot = relay
7608 .add_agent(
7609 AdapterHost::new(
7610 Box::new(ScriptedAdapter::new(
7611 2,
7612 AgentCapabilities::default(),
7613 [AgentEvent::TurnComplete { slot: 2 }],
7614 )),
7615 None,
7616 ),
7617 "Reviewer",
7618 "reviewer.example",
7619 "reviewer --acp",
7620 )
7621 .await
7622 .expect("append agent");
7623 assert_eq!(slot, 2);
7624 assert_eq!(
7625 relay.relay().active_slots().collect::<Vec<_>>(),
7626 vec![0, 1, 2]
7627 );
7628 assert_eq!(
7629 relay
7630 .session_metadata()
7631 .get("agents")
7632 .and_then(|value| value.as_array())
7633 .map(Vec::len),
7634 Some(3)
7635 );
7636 relay.drop_agent(1).await.expect("drop middle peer");
7637 let metadata = relay.session_metadata();
7638 assert_eq!(
7639 metadata.get("agents"),
7640 Some(&serde_json::json!([
7641 {"slot": 0, "name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
7642 {"slot": 2, "name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
7643 ]))
7644 );
7645 assert!(
7646 events
7647 .lock()
7648 .expect("lock")
7649 .iter()
7650 .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
7651 );
7652 }
7653
7654 #[tokio::test]
7655 async fn relay_host_persists_coordinator_owned_runtime_metadata() {
7656 let path = unique_test_path("codeswarm-session-metadata", "json");
7657 let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
7658 let writer = metadata_store.buffered().expect("metadata writer");
7659 let first = AdapterHost::new(
7660 Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
7661 None,
7662 );
7663 let second = AdapterHost::new(
7664 Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
7665 None,
7666 );
7667 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7668 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7669 relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
7670 relay.set_roster_launch_specs(vec![
7671 ("custom".into(), "claude".into()),
7672 ("custom".into(), "codex".into()),
7673 ]);
7674 relay.set_session_metadata_writer(writer);
7675 relay.start().await.expect("start");
7676 relay.drop_agent(0).await.expect("drop first agent");
7677 relay.stop().await.expect("stop");
7678
7679 let loaded = metadata_store
7680 .read()
7681 .expect("read metadata")
7682 .expect("metadata snapshot");
7683 assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
7684 assert_eq!(
7685 loaded.get("agents"),
7686 Some(&serde_json::json!([{
7687 "slot": 1, "name": "Codex", "identity": "openai.com", "protocol": "custom",
7688 "command": "codex", "supports_load_session": false
7689 }]))
7690 );
7691 assert!(loaded.get("owner").is_none());
7692 let _ = std::fs::remove_file(path);
7693 }
7694
7695 #[tokio::test]
7696 async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
7697 let first = AdapterHost::new(
7698 Box::new(ScriptedAdapter::new(
7699 0,
7700 AgentCapabilities::default(),
7701 [
7702 AgentEvent::Text {
7703 slot: 0,
7704 text: "owner stream".into(),
7705 },
7706 AgentEvent::TurnComplete { slot: 0 },
7707 ],
7708 )),
7709 None,
7710 );
7711 let second = AdapterHost::new(
7712 Box::new(ScriptedAdapter::new(
7713 1,
7714 AgentCapabilities::default(),
7715 [
7716 AgentEvent::Text {
7717 slot: 1,
7718 text: "peer stream".into(),
7719 },
7720 AgentEvent::TurnComplete { slot: 1 },
7721 ],
7722 )),
7723 None,
7724 );
7725 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7726 relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
7727 relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
7728 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7729 let captured = std::sync::Arc::clone(&events);
7730 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7731 relay.start().await.expect("start");
7732
7733 relay.swap_agents(0, 1).expect("swap peers");
7734 assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
7735 assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
7736 relay.run_turn("task", 0).await.expect("swapped turn");
7737 let events = events.lock().expect("events");
7738 assert!(events.iter().any(|event| {
7739 matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
7740 }));
7741 assert!(relay.dispatches()[0].1.contains("You are Peer"));
7742 }
7743
7744 #[tokio::test]
7745 async fn relay_host_persists_all_active_agent_metadata_off_thread() {
7746 let path = unique_test_path("codeswarm-session-metadata", "json");
7747 let first = AdapterHost::new(
7748 Box::new(ScriptedAdapter::new(
7749 0,
7750 AgentCapabilities::default(),
7751 [AgentEvent::TurnComplete { slot: 0 }],
7752 )),
7753 None,
7754 );
7755 let second = AdapterHost::new(
7756 Box::new(ScriptedAdapter::new(
7757 1,
7758 AgentCapabilities::default(),
7759 [AgentEvent::TurnComplete { slot: 1 }],
7760 )),
7761 None,
7762 );
7763 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7764 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7765 relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
7766 relay.set_roster_launch_specs(vec![
7767 ("custom".into(), "claude".into()),
7768 ("custom".into(), "codex".into()),
7769 ]);
7770 let writer = SessionMetadataStore::open(&path)
7771 .buffered()
7772 .expect("metadata writer");
7773 relay.set_session_metadata_writer(writer);
7774 relay.start().await.expect("start");
7775 relay.stop().await.expect("stop");
7776 let loaded = SessionMetadataStore::open(&path)
7777 .read()
7778 .expect("read metadata")
7779 .expect("metadata snapshot");
7780 let agents = loaded
7781 .get("agents")
7782 .and_then(|value| value.as_array())
7783 .expect("agents");
7784 assert_eq!(agents.len(), 2);
7785 assert_eq!(agents[0]["identity"], "claude.com");
7786 assert_eq!(agents[1]["identity"], "openai.com");
7787 let _ = std::fs::remove_file(path);
7788 }
7789
7790 #[tokio::test]
7791 async fn relay_host_routes_unseen_public_context_to_next_agent() {
7792 let first = AdapterHost::new(
7793 Box::new(ScriptedAdapter::new(
7794 0,
7795 AgentCapabilities::default(),
7796 [
7797 AgentEvent::Text {
7798 slot: 0,
7799 text: "implemented the fix".into(),
7800 },
7801 AgentEvent::TurnComplete { slot: 0 },
7802 ],
7803 )),
7804 None,
7805 );
7806 let second = AdapterHost::new(
7807 Box::new(ScriptedAdapter::new(
7808 1,
7809 AgentCapabilities::default(),
7810 [AgentEvent::TurnComplete { slot: 1 }],
7811 )),
7812 None,
7813 );
7814 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7815 relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
7816 relay.start().await.expect("start");
7817 relay.run_turn("task", 0).await.expect("first turn");
7818 relay.run_turn("review this", 0).await.expect("review turn");
7819
7820 assert_eq!(relay.dispatches().len(), 2);
7821 assert_eq!(relay.dispatches()[0].0, 0);
7822 assert!(relay.dispatches()[0].1.contains("task"));
7823 assert!(relay.dispatches()[0].1.contains("You are Codex"));
7824 assert!(relay.dispatches()[0].1.contains("2. Qwen"));
7825 assert_eq!(relay.dispatches()[1].0, 1);
7826 assert!(relay.dispatches()[1].1.contains("review this"));
7827 let public = relay.dispatches()[1]
7828 .1
7829 .split_once("Public updates:\n")
7830 .map(|(_, updates)| updates)
7831 .expect("review receives public context");
7832 let header = public
7833 .lines()
7834 .find(|line| line.starts_with("Codex "))
7835 .expect("named previous agent");
7836 let timestamp = header
7837 .strip_prefix("Codex ")
7838 .and_then(|value| value.strip_suffix(':'))
7839 .expect("timestamped header");
7840 assert_eq!(timestamp.len(), 5);
7841 assert_eq!(timestamp.as_bytes()[2], b':');
7842 assert!(
7843 timestamp
7844 .bytes()
7845 .enumerate()
7846 .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
7847 );
7848 assert!(public.contains("implemented the fix"));
7849 assert!(!public.contains("Agent 0"));
7850 }
7851}