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): (Vec<String>, Vec<String>) = options
3979 .and_then(Value::as_array)
3980 .map(|options| {
3981 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 Some((label, option_id))
3996 })
3997 .unzip()
3998 })
3999 .unwrap_or_default();
4000 if options.is_empty() {
4001 return None;
4002 }
4003 Some(AgentEvent::Permission {
4004 slot,
4005 request: PermissionRequest {
4006 id: request_id.to_owned(),
4007 title,
4008 options,
4009 option_ids,
4010 },
4011 })
4012}
4013
4014fn rpc_id_to_string(value: &Value) -> String {
4015 value
4016 .as_str()
4017 .map(str::to_owned)
4018 .or_else(|| value.as_u64().map(|id| id.to_string()))
4019 .unwrap_or_else(|| value.to_string())
4020}
4021
4022#[cfg(test)]
4023mod tests {
4024 use super::{
4025 AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
4026 RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
4027 parse_model_config, prompt_content_blocks, read_bounded_line,
4028 };
4029 #[cfg(target_os = "linux")]
4030 use super::{isolate_process_group, terminate_child};
4031 use crate::TerminalEvent;
4032 use crate::{
4033 AdapterError, AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
4034 persistence::SessionMetadataStore,
4035 relay::{CollaborationStrategy, DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
4036 };
4037 use async_trait::async_trait;
4038 use serde_json::Value;
4039 use std::collections::VecDeque;
4040 use std::sync::{
4041 Arc, Mutex,
4042 atomic::{AtomicUsize, Ordering},
4043 };
4044
4045 fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
4046 let nonce = std::time::SystemTime::now()
4047 .duration_since(std::time::UNIX_EPOCH)
4048 .expect("clock")
4049 .as_nanos();
4050 std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
4051 }
4052
4053 #[test]
4054 fn malformed_file_writes_preserve_existing_content() {
4055 let root = unique_test_path("codeswarm-write-validation", "dir");
4056 std::fs::create_dir_all(&root).unwrap();
4057 let file = root.join("keep.txt");
4058 std::fs::write(&file, "valuable content").unwrap();
4059 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
4060 for content in [
4061 Value::Null,
4062 serde_json::json!(false),
4063 serde_json::json!(42),
4064 serde_json::json!([]),
4065 ] {
4066 assert!(
4067 adapter
4068 .write_workspace_text(
4069 &serde_json::json!({"path":"keep.txt", "content": content})
4070 )
4071 .is_err()
4072 );
4073 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4074 }
4075 assert!(
4076 adapter
4077 .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
4078 .is_err()
4079 );
4080 assert!(
4081 adapter
4082 .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
4083 .is_err()
4084 );
4085 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4086 adapter
4087 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
4088 .unwrap();
4089 assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
4090 adapter
4091 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
4092 .unwrap();
4093 assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
4094 std::fs::remove_dir_all(root).unwrap();
4095 }
4096
4097 #[tokio::test]
4098 async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
4099 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
4100 let mut adapter = AcpAdapter::new(
4101 0,
4102 std::env::current_dir().unwrap(),
4103 "sh",
4104 vec!["-c".into(), script.into()],
4105 );
4106 adapter.start().await.unwrap();
4107 let error = adapter
4108 .request_with_timeout(
4109 "session/set_mode",
4110 serde_json::json!({}),
4111 std::time::Duration::from_millis(10),
4112 )
4113 .await
4114 .unwrap_err();
4115 assert!(error.to_string().contains("session/set_mode timed out"));
4116 adapter.stop().await.unwrap();
4117 assert!(adapter.child.is_none());
4118 }
4119
4120 #[tokio::test]
4121 async fn goals_reach_every_roster_slot_without_native_goal_support() {
4122 use crate::goal::GoalCommand;
4123 let hosts = (0..3)
4124 .map(|slot| {
4125 AdapterHost::new(
4126 Box::new(ScriptedAdapter::new(
4127 slot,
4128 AgentCapabilities::default(),
4129 [
4130 AgentEvent::TurnComplete { slot },
4131 AgentEvent::TurnComplete { slot },
4132 ],
4133 )),
4134 None,
4135 )
4136 })
4137 .collect();
4138 let mut relay = RelayHost::new(hosts, 10).unwrap();
4139 relay.start().await.unwrap();
4140 let task = relay
4141 .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
4142 .unwrap()
4143 .unwrap();
4144 relay.relay_mut().enqueue_human(task, Some(0));
4145 for slot in 0..3 {
4146 relay.run_turn("", 0).await.unwrap();
4147 let (actual, prompt) = relay.dispatches().last().unwrap();
4148 assert_eq!(*actual, slot);
4149 assert!(prompt.contains("Active shared goal: Ship the settings screen"));
4150 }
4151 let snapshot = relay.session_metadata();
4152 let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
4153 assert!(restored.is_some());
4154 relay.restore_goal(restored);
4155 relay.reload(0).await.unwrap();
4156 relay.run_turn("", 0).await.unwrap();
4157 assert!(
4158 relay
4159 .dispatches()
4160 .last()
4161 .unwrap()
4162 .1
4163 .contains("Active shared goal: Ship the settings screen")
4164 );
4165 relay.apply_goal(GoalCommand::Done).unwrap();
4166 relay.run_turn("", 0).await.unwrap();
4167 assert!(
4168 relay
4169 .dispatches()
4170 .last()
4171 .unwrap()
4172 .1
4173 .contains("No active shared goal")
4174 );
4175 relay.apply_goal(GoalCommand::Clear).unwrap();
4176 assert!(relay.session_metadata().get("goal").unwrap().is_null());
4177 }
4178
4179 #[tokio::test]
4180 async fn replacement_agent_receives_task_after_public_journal_pruning() {
4181 let hosts = (0..2)
4182 .map(|slot| {
4183 AdapterHost::new(
4184 Box::new(ScriptedAdapter::new(
4185 slot,
4186 AgentCapabilities::default(),
4187 [
4188 AgentEvent::Text {
4189 slot,
4190 text: "progress".into(),
4191 },
4192 AgentEvent::TurnComplete { slot },
4193 AgentEvent::Text {
4194 slot,
4195 text: "more progress".into(),
4196 },
4197 AgentEvent::TurnComplete { slot },
4198 ],
4199 )),
4200 None,
4201 )
4202 })
4203 .collect();
4204 let mut relay = RelayHost::new(hosts, 10).unwrap();
4205 relay.start().await.unwrap();
4206 relay
4207 .relay_mut()
4208 .enqueue_human("Fix the login bug", Some(0));
4209 relay.run_turn("", 0).await.unwrap();
4210 relay.run_turn("", 0).await.unwrap();
4211 relay.run_turn("", 0).await.unwrap();
4212 relay.reload(1).await.unwrap();
4213 assert!(
4214 !relay
4215 .relay_mut()
4216 .unseen_context(1)
4217 .contains("Fix the login bug")
4218 );
4219 relay.run_turn("", 0).await.unwrap();
4220 assert!(
4221 relay
4222 .dispatches()
4223 .last()
4224 .unwrap()
4225 .1
4226 .contains("Shared task:\nFix the login bug")
4227 );
4228 }
4229
4230 #[cfg(target_os = "linux")]
4231 #[tokio::test]
4232 async fn termination_kills_only_the_verified_isolated_child_group() {
4233 use nix::unistd::{Pid, getpgid, getpgrp};
4234 use tokio::io::{AsyncBufReadExt, BufReader};
4235
4236 let own_group = getpgrp();
4237 let mut command = tokio::process::Command::new("sh");
4238 isolate_process_group(&mut command);
4239 command
4240 .arg("-c")
4241 .arg("sleep 60 & echo $!; wait")
4242 .stdout(std::process::Stdio::piped());
4243 let mut child = command.spawn().expect("spawn isolated shell");
4244 let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
4245 assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
4246 assert_ne!(leader, own_group);
4247
4248 let stdout = child.stdout.take().expect("child stdout");
4249 let mut lines = BufReader::new(stdout).lines();
4250 let descendant = lines
4251 .next_line()
4252 .await
4253 .expect("read descendant pid")
4254 .expect("descendant pid")
4255 .parse::<i32>()
4256 .expect("numeric descendant pid");
4257 let descendant = Pid::from_raw(descendant);
4258 assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
4259
4260 terminate_child(&mut child).await.expect("terminate group");
4261 for _ in 0..100 {
4262 if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
4263 return;
4264 }
4265 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
4266 }
4267 panic!("descendant {descendant} survived isolated group termination");
4268 }
4269
4270 #[test]
4271 fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
4272 assert_eq!(
4273 parse_command_line(r#"agent --name "local bridge" --flag 'two words'"#),
4274 Ok((
4275 "agent".into(),
4276 vec![
4277 "--name".into(),
4278 "local bridge".into(),
4279 "--flag".into(),
4280 "two words".into(),
4281 ]
4282 ),)
4283 );
4284 assert_eq!(
4285 parse_command_line(r#"agent "" escaped\ argument"#),
4286 Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
4287 );
4288 }
4289
4290 #[test]
4291 fn acp_prompt_expands_safe_at_path_resources() {
4292 let root = unique_test_path("codeswarm-prompt-resource", "dir");
4293 std::fs::create_dir_all(&root).expect("workspace");
4294 std::fs::write(root.join("note.md"), "resource text").expect("resource");
4295 let blocks = prompt_content_blocks(&root, "inspect @note.md");
4296 assert_eq!(blocks[0]["type"], "text");
4297 assert_eq!(blocks[0]["text"], "inspect @note.md");
4298 assert_eq!(blocks[1]["type"], "resource");
4299 assert_eq!(blocks[1]["resource"]["text"], "resource text");
4300 assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
4301 std::fs::remove_dir_all(root).expect("cleanup workspace");
4302 }
4303
4304 #[tokio::test]
4305 async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
4306 let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
4307 bytes.push(b'\n');
4308 let mut reader = tokio::io::BufReader::new(bytes.as_slice());
4309 assert!(matches!(
4310 read_bounded_line(&mut reader).await,
4311 Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
4312 ));
4313 }
4314
4315 #[test]
4316 fn rejects_malformed_configured_commands_before_spawn() {
4317 assert_eq!(
4318 parse_command_line("agent 'unfinished"),
4319 Err(super::CommandParseError::UnterminatedQuote)
4320 );
4321 assert_eq!(
4322 parse_command_line("agent\\"),
4323 Err(super::CommandParseError::TrailingEscape)
4324 );
4325 assert_eq!(
4326 parse_command_line(" \t"),
4327 Err(super::CommandParseError::Empty)
4328 );
4329 }
4330
4331 #[derive(Debug)]
4332 struct PendingAdapter {
4333 slot: usize,
4334 hang_on_cancel: bool,
4335 }
4336
4337 #[derive(Debug)]
4338 struct ConcurrentStartAdapter {
4339 slot: usize,
4340 barrier: Arc<tokio::sync::Barrier>,
4341 }
4342
4343 #[derive(Debug)]
4344 struct ReloadProbeAdapter {
4345 slot: usize,
4346 crashed: bool,
4347 reloaded: bool,
4348 events: VecDeque<AgentEvent>,
4349 prompts: Arc<Mutex<Vec<String>>>,
4350 }
4351
4352 #[async_trait]
4353 impl AgentAdapter for ReloadProbeAdapter {
4354 fn slot(&self) -> usize {
4355 self.slot
4356 }
4357
4358 fn display_name(&self) -> String {
4359 "Reload probe".into()
4360 }
4361
4362 fn protocol(&self) -> &'static str {
4363 "native"
4364 }
4365
4366 fn capabilities(&self) -> AgentCapabilities {
4367 AgentCapabilities {
4368 supports_modes: true,
4369 ..AgentCapabilities::default()
4370 }
4371 }
4372
4373 fn needs_restart(&self) -> bool {
4374 self.crashed && !self.reloaded
4375 }
4376
4377 async fn start(&mut self) -> super::AdapterResult<()> {
4378 self.events.push_back(AgentEvent::ModesReplaced {
4379 slot: self.slot,
4380 modes: vec![
4381 Mode {
4382 id: "codeswarm:mode:full-access".into(),
4383 label: "Auto pilot".into(),
4384 },
4385 Mode {
4386 id: "codeswarm:mode:plan".into(),
4387 label: "Plan".into(),
4388 },
4389 ],
4390 current_mode: Some("codeswarm:mode:full-access".into()),
4391 });
4392 self.events.push_back(AgentEvent::Ready {
4393 slot: self.slot,
4394 capabilities: self.capabilities(),
4395 });
4396 Ok(())
4397 }
4398
4399 async fn send_prompt(&mut self, prompt: String) -> super::AdapterResult<()> {
4400 self.prompts.lock().expect("prompts").push(prompt);
4401 if !self.crashed {
4402 self.crashed = true;
4403 self.events.push_back(AgentEvent::Failed {
4404 slot: self.slot,
4405 started: true,
4406 detail: "probe crashed".into(),
4407 });
4408 } else {
4409 self.events.push_back(AgentEvent::Text {
4410 slot: self.slot,
4411 text: "recovered".into(),
4412 });
4413 self.events
4414 .push_back(AgentEvent::TurnComplete { slot: self.slot });
4415 }
4416 Ok(())
4417 }
4418
4419 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4420 Ok(false)
4421 }
4422
4423 async fn answer_permission(
4424 &mut self,
4425 _request_id: String,
4426 _answer: PermissionAnswer,
4427 ) -> super::AdapterResult<()> {
4428 Err(super::AdapterError::Unsupported("permission answer"))
4429 }
4430
4431 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4432 Ok(())
4433 }
4434
4435 async fn reload(&mut self) -> super::AdapterResult<()> {
4436 self.reloaded = true;
4437 self.start().await
4438 }
4439
4440 async fn stop(&mut self) -> super::AdapterResult<()> {
4441 Ok(())
4442 }
4443
4444 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4445 self.events.pop_front().map(Ok)
4446 }
4447 }
4448
4449 #[async_trait]
4450 impl AgentAdapter for ConcurrentStartAdapter {
4451 fn slot(&self) -> usize {
4452 self.slot
4453 }
4454
4455 fn capabilities(&self) -> AgentCapabilities {
4456 AgentCapabilities::default()
4457 }
4458
4459 async fn start(&mut self) -> super::AdapterResult<()> {
4460 self.barrier.wait().await;
4461 Ok(())
4462 }
4463
4464 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4465 Ok(())
4466 }
4467
4468 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4469 Ok(true)
4470 }
4471
4472 async fn answer_permission(
4473 &mut self,
4474 _request_id: String,
4475 _answer: PermissionAnswer,
4476 ) -> super::AdapterResult<()> {
4477 Ok(())
4478 }
4479
4480 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4481 Ok(())
4482 }
4483
4484 async fn reload(&mut self) -> super::AdapterResult<()> {
4485 Ok(())
4486 }
4487
4488 async fn stop(&mut self) -> super::AdapterResult<()> {
4489 Ok(())
4490 }
4491
4492 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4493 std::future::pending().await
4494 }
4495 }
4496
4497 #[derive(Debug)]
4498 struct PermissionBlockingAdapter {
4499 slot: usize,
4500 phase: u8,
4501 }
4502
4503 #[async_trait]
4504 impl AgentAdapter for PermissionBlockingAdapter {
4505 fn slot(&self) -> usize {
4506 self.slot
4507 }
4508
4509 fn capabilities(&self) -> AgentCapabilities {
4510 AgentCapabilities {
4511 supports_permissions: true,
4512 ..AgentCapabilities::default()
4513 }
4514 }
4515
4516 async fn start(&mut self) -> super::AdapterResult<()> {
4517 Ok(())
4518 }
4519
4520 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4521 Ok(())
4522 }
4523
4524 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4525 Ok(true)
4526 }
4527
4528 async fn answer_permission(
4529 &mut self,
4530 request_id: String,
4531 answer: PermissionAnswer,
4532 ) -> super::AdapterResult<()> {
4533 if self.phase != 1 || request_id != "permission-1" {
4534 return Err(super::AdapterError::Protocol(
4535 "unexpected permission response".into(),
4536 ));
4537 }
4538 assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4539 self.phase = 2;
4540 Ok(())
4541 }
4542
4543 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4544 Ok(())
4545 }
4546
4547 async fn reload(&mut self) -> super::AdapterResult<()> {
4548 Ok(())
4549 }
4550
4551 async fn stop(&mut self) -> super::AdapterResult<()> {
4552 Ok(())
4553 }
4554
4555 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4556 match self.phase {
4557 0 => {
4558 self.phase = 1;
4559 Some(Ok(AgentEvent::Permission {
4560 slot: self.slot,
4561 request: crate::PermissionRequest {
4562 id: "permission-1".into(),
4563 title: "Allow?".into(),
4564 options: vec!["Allow".into()],
4565 option_ids: vec!["allow".into()],
4566 },
4567 }))
4568 }
4569 1 => std::future::pending().await,
4570 _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4571 }
4572 }
4573 }
4574
4575 #[async_trait]
4576 impl AgentAdapter for PendingAdapter {
4577 fn slot(&self) -> usize {
4578 self.slot
4579 }
4580
4581 fn capabilities(&self) -> AgentCapabilities {
4582 AgentCapabilities {
4583 supports_cancel: true,
4584 ..AgentCapabilities::default()
4585 }
4586 }
4587
4588 async fn start(&mut self) -> super::AdapterResult<()> {
4589 Ok(())
4590 }
4591
4592 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4593 Ok(())
4594 }
4595
4596 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4597 if self.hang_on_cancel {
4598 return std::future::pending().await;
4599 }
4600 Ok(true)
4601 }
4602
4603 async fn answer_permission(
4604 &mut self,
4605 _request_id: String,
4606 _answer: PermissionAnswer,
4607 ) -> super::AdapterResult<()> {
4608 Ok(())
4609 }
4610
4611 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4612 Ok(())
4613 }
4614
4615 async fn reload(&mut self) -> super::AdapterResult<()> {
4616 Ok(())
4617 }
4618
4619 async fn stop(&mut self) -> super::AdapterResult<()> {
4620 Ok(())
4621 }
4622
4623 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4624 std::future::pending().await
4625 }
4626 }
4627
4628 #[derive(Debug)]
4629 struct StopTrackingAdapter {
4630 slot: usize,
4631 stopped: Arc<AtomicUsize>,
4632 fail_stop: bool,
4633 }
4634
4635 #[derive(Debug)]
4636 struct ModeOrderAdapter {
4637 slot: usize,
4638 log: Arc<Mutex<Vec<String>>>,
4639 phase: u8,
4640 }
4641
4642 #[derive(Debug)]
4643 struct StartupAcpAdapter {
4644 slot: usize,
4645 events: std::collections::VecDeque<AgentEvent>,
4646 }
4647
4648 impl StartupAcpAdapter {
4649 fn new(slot: usize) -> Self {
4650 Self {
4651 slot,
4652 events: [
4653 AgentEvent::ModesReplaced {
4654 slot,
4655 modes: vec![Mode {
4656 id: "full-access".into(),
4657 label: "Auto pilot".into(),
4658 }],
4659 current_mode: Some("full-access".into()),
4660 },
4661 AgentEvent::Ready {
4662 slot,
4663 capabilities: AgentCapabilities {
4664 supports_modes: true,
4665 ..AgentCapabilities::default()
4666 },
4667 },
4668 ]
4669 .into(),
4670 }
4671 }
4672 }
4673
4674 #[async_trait]
4675 impl AgentAdapter for StartupAcpAdapter {
4676 fn slot(&self) -> usize {
4677 self.slot
4678 }
4679
4680 fn protocol(&self) -> &'static str {
4681 "acp"
4682 }
4683
4684 fn capabilities(&self) -> AgentCapabilities {
4685 AgentCapabilities {
4686 supports_modes: true,
4687 ..AgentCapabilities::default()
4688 }
4689 }
4690
4691 async fn start(&mut self) -> super::AdapterResult<()> {
4692 Ok(())
4693 }
4694
4695 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4696 Ok(())
4697 }
4698
4699 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4700 Ok(true)
4701 }
4702
4703 async fn answer_permission(
4704 &mut self,
4705 _request_id: String,
4706 _answer: PermissionAnswer,
4707 ) -> super::AdapterResult<()> {
4708 Ok(())
4709 }
4710
4711 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4712 Ok(())
4713 }
4714
4715 async fn reload(&mut self) -> super::AdapterResult<()> {
4716 self.events = Self::new(self.slot).events;
4717 Ok(())
4718 }
4719
4720 async fn stop(&mut self) -> super::AdapterResult<()> {
4721 Ok(())
4722 }
4723
4724 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4725 self.events.pop_front().map(Ok)
4726 }
4727 }
4728
4729 #[async_trait]
4730 impl AgentAdapter for ModeOrderAdapter {
4731 fn slot(&self) -> usize {
4732 self.slot
4733 }
4734
4735 fn capabilities(&self) -> AgentCapabilities {
4736 AgentCapabilities {
4737 supports_modes: true,
4738 ..AgentCapabilities::default()
4739 }
4740 }
4741
4742 async fn start(&mut self) -> super::AdapterResult<()> {
4743 self.log.lock().expect("log").push("start".into());
4744 Ok(())
4745 }
4746
4747 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4748 self.log.lock().expect("log").push("prompt".into());
4749 Ok(())
4750 }
4751
4752 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4753 Ok(true)
4754 }
4755
4756 async fn answer_permission(
4757 &mut self,
4758 _request_id: String,
4759 _answer: PermissionAnswer,
4760 ) -> super::AdapterResult<()> {
4761 Ok(())
4762 }
4763
4764 async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4765 self.log.lock().expect("log").push(format!("mode:{mode}"));
4766 Ok(())
4767 }
4768
4769 async fn reload(&mut self) -> super::AdapterResult<()> {
4770 self.log.lock().expect("log").push("reload".into());
4771 self.phase = 0;
4772 Ok(())
4773 }
4774
4775 async fn stop(&mut self) -> super::AdapterResult<()> {
4776 Ok(())
4777 }
4778
4779 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4780 match self.phase {
4781 0 => {
4782 self.phase = 1;
4783 Some(Ok(AgentEvent::ModesReplaced {
4784 slot: self.slot,
4785 modes: vec![Mode {
4786 id: "yolo".into(),
4787 label: "YOLO".into(),
4788 }],
4789 current_mode: None,
4790 }))
4791 }
4792 1 => {
4793 self.phase = 2;
4794 Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4795 }
4796 _ => std::future::pending().await,
4797 }
4798 }
4799 }
4800
4801 #[async_trait]
4802 impl AgentAdapter for StopTrackingAdapter {
4803 fn slot(&self) -> usize {
4804 self.slot
4805 }
4806
4807 fn capabilities(&self) -> AgentCapabilities {
4808 AgentCapabilities::default()
4809 }
4810
4811 async fn start(&mut self) -> super::AdapterResult<()> {
4812 Ok(())
4813 }
4814
4815 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4816 Ok(())
4817 }
4818
4819 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4820 Ok(false)
4821 }
4822
4823 async fn answer_permission(
4824 &mut self,
4825 _request_id: String,
4826 _answer: PermissionAnswer,
4827 ) -> super::AdapterResult<()> {
4828 Ok(())
4829 }
4830
4831 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4832 Ok(())
4833 }
4834
4835 async fn reload(&mut self) -> super::AdapterResult<()> {
4836 Ok(())
4837 }
4838
4839 async fn stop(&mut self) -> super::AdapterResult<()> {
4840 self.stopped.fetch_add(1, Ordering::Relaxed);
4841 if self.fail_stop {
4842 Err(super::AdapterError::Transport("stop failed".into()))
4843 } else {
4844 Ok(())
4845 }
4846 }
4847
4848 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4849 None
4850 }
4851 }
4852
4853 #[derive(Debug)]
4857 struct FailingStartAdapter {
4858 slot: usize,
4859 stopped: Arc<AtomicUsize>,
4860 }
4861
4862 #[async_trait]
4863 impl AgentAdapter for FailingStartAdapter {
4864 fn slot(&self) -> usize {
4865 self.slot
4866 }
4867
4868 fn capabilities(&self) -> AgentCapabilities {
4869 AgentCapabilities::default()
4870 }
4871
4872 async fn start(&mut self) -> super::AdapterResult<()> {
4873 Err(super::AdapterError::Spawn("startup failed".into()))
4874 }
4875
4876 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4877 Ok(())
4878 }
4879
4880 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4881 Ok(false)
4882 }
4883
4884 async fn answer_permission(
4885 &mut self,
4886 _request_id: String,
4887 _answer: PermissionAnswer,
4888 ) -> super::AdapterResult<()> {
4889 Ok(())
4890 }
4891
4892 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4893 Ok(())
4894 }
4895
4896 async fn reload(&mut self) -> super::AdapterResult<()> {
4897 Ok(())
4898 }
4899
4900 async fn stop(&mut self) -> super::AdapterResult<()> {
4901 self.stopped.fetch_add(1, Ordering::Relaxed);
4902 Ok(())
4903 }
4904
4905 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4906 None
4907 }
4908 }
4909
4910 #[tokio::test]
4911 async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4912 let stopped = Arc::new(AtomicUsize::new(0));
4913 let relay = RelayHost::new(
4914 vec![
4915 AdapterHost::new(
4916 Box::new(StopTrackingAdapter {
4917 slot: 0,
4918 stopped: Arc::clone(&stopped),
4919 fail_stop: true,
4920 }),
4921 None,
4922 ),
4923 AdapterHost::new(
4924 Box::new(StopTrackingAdapter {
4925 slot: 1,
4926 stopped: Arc::clone(&stopped),
4927 fail_stop: false,
4928 }),
4929 None,
4930 ),
4931 ],
4932 4,
4933 )
4934 .expect("relay");
4935 let mut relay = relay;
4936
4937 let error = relay.stop().await.expect_err("first stop failure");
4938 assert!(error.to_string().contains("stop failed"));
4939 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4940 }
4941
4942 #[tokio::test]
4943 async fn relay_start_isolates_a_failed_adapter_and_keeps_healthy_peers() {
4944 let stopped = Arc::new(AtomicUsize::new(0));
4945 let mut relay = RelayHost::new(
4946 vec![
4947 AdapterHost::new(
4948 Box::new(StopTrackingAdapter {
4949 slot: 0,
4950 stopped: Arc::clone(&stopped),
4951 fail_stop: false,
4952 }),
4953 None,
4954 ),
4955 AdapterHost::new(
4956 Box::new(FailingStartAdapter {
4957 slot: 1,
4958 stopped: Arc::clone(&stopped),
4959 }),
4960 None,
4961 ),
4962 ],
4963 4,
4964 )
4965 .expect("relay");
4966
4967 relay.start().await.expect("healthy peer remains available");
4968 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0]);
4969 assert_eq!(stopped.load(Ordering::Relaxed), 1);
4970 relay.stop().await.unwrap();
4971 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4972 }
4973
4974 #[test]
4975 fn parses_acp_text_without_ui_dependency() {
4976 let event = parse_acp_notification(
4977 2,
4978 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4979 )
4980 .expect("valid ACP")
4981 .expect("text event");
4982 assert_eq!(
4983 event,
4984 AgentEvent::Text {
4985 slot: 2,
4986 text: "hello".into(),
4987 }
4988 );
4989 }
4990
4991 #[test]
4992 fn parses_acp_state_notifications_at_the_adapter_boundary() {
4993 let commands = parse_acp_notification(
4994 3,
4995 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
4996 )
4997 .expect("valid ACP")
4998 .expect("commands event");
4999 assert_eq!(
5000 commands,
5001 AgentEvent::CommandsReplaced {
5002 slot: 3,
5003 commands: vec![crate::AgentCommand {
5004 name: "review".into()
5005 }]
5006 }
5007 );
5008
5009 let mode = parse_acp_notification(
5010 3,
5011 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
5012 )
5013 .expect("valid ACP")
5014 .expect("mode event");
5015 assert_eq!(
5016 mode,
5017 AgentEvent::ModeUpdated {
5018 slot: 3,
5019 current_mode: "review".into()
5020 }
5021 );
5022
5023 let usage = parse_acp_notification(
5024 3,
5025 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
5026 )
5027 .expect("valid ACP")
5028 .expect("usage event");
5029 assert_eq!(
5030 usage,
5031 AgentEvent::UsageUpdated {
5032 slot: 3,
5033 usage: crate::UsageUpdate {
5034 used: 4200,
5035 size: 128000
5036 }
5037 }
5038 );
5039
5040 let models = parse_acp_notification(
5041 3,
5042 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"}]}]}}}"#,
5043 )
5044 .expect("valid ACP")
5045 .expect("models event");
5046 assert!(matches!(
5047 models,
5048 AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
5049 if models.len() == 2 && current_model.as_deref() == Some("smart")
5050 ));
5051 assert_eq!(
5052 parse_model_config(&serde_json::json!({
5053 "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
5054 })),
5055 None
5056 );
5057
5058 let user = parse_acp_notification(
5059 3,
5060 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
5061 )
5062 .expect("valid ACP")
5063 .expect("user event");
5064 assert_eq!(
5065 user,
5066 AgentEvent::UserText {
5067 slot: 3,
5068 text: "context".into()
5069 }
5070 );
5071 }
5072
5073 #[test]
5074 fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
5075 let event = parse_acp_notification(
5076 0,
5077 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
5078 )
5079 .expect("valid ACP")
5080 .expect("mode event");
5081 assert!(matches!(
5082 event,
5083 AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
5084 if mode == "yolo" && modes[0].id == "yolo"
5085 ));
5086 }
5087
5088 #[test]
5089 fn parses_native_agy_text_without_acp_bridge() {
5090 let event = parse_agy_line(
5091 1,
5092 r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
5093 )
5094 .expect("valid stream-json")
5095 .expect("text event");
5096 assert_eq!(
5097 event,
5098 AgentEvent::Text {
5099 slot: 1,
5100 text: "hello".into(),
5101 }
5102 );
5103 }
5104
5105 #[test]
5106 fn parses_tool_lifecycle_from_each_protocol() {
5107 let agy = parse_agy_line(
5108 1,
5109 r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
5110 )
5111 .expect("valid native tool")
5112 .expect("tool event");
5113 assert!(matches!(
5114 agy,
5115 AgentEvent::Tool {
5116 update: crate::ToolUpdate {
5117 status: ToolStatus::Completed,
5118 ..
5119 },
5120 ..
5121 }
5122 ));
5123
5124 let acp = parse_acp_notification(
5125 1,
5126 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
5127 )
5128 .expect("valid ACP tool")
5129 .expect("tool event");
5130 assert!(matches!(
5131 acp,
5132 AgentEvent::Tool {
5133 update: crate::ToolUpdate {
5134 status: ToolStatus::Failed,
5135 ..
5136 },
5137 ..
5138 }
5139 ));
5140 }
5141
5142 #[test]
5143 fn parses_terminal_lifecycle_from_acp_and_native_events() {
5144 let created = parse_acp_notification(
5145 0,
5146 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
5147 )
5148 .expect("valid ACP terminal")
5149 .expect("terminal event");
5150 assert_eq!(
5151 created,
5152 AgentEvent::Terminal {
5153 slot: 0,
5154 event: TerminalEvent::Created {
5155 id: "term-1".into(),
5156 command: "cargo test".into(),
5157 },
5158 }
5159 );
5160 let output = parse_agy_line(
5161 1,
5162 r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
5163 )
5164 .expect("valid native terminal")
5165 .expect("terminal event");
5166 assert_eq!(
5167 output,
5168 AgentEvent::Terminal {
5169 slot: 1,
5170 event: TerminalEvent::Output {
5171 id: "term-1".into(),
5172 text: "ok\n".into(),
5173 },
5174 }
5175 );
5176 let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
5177 .expect("valid native release")
5178 .expect("terminal event");
5179 assert!(matches!(
5180 released,
5181 AgentEvent::Terminal {
5182 event: TerminalEvent::Released { id },
5183 ..
5184 } if id == "term-1"
5185 ));
5186 }
5187
5188 #[test]
5189 fn parses_acp_permission_requests() {
5190 let event = parse_acp_notification(
5191 0,
5192 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"request_permission","toolCall":{"toolCallId":"t1","title":"Write file"},"options":[{"name":"Allow once","optionId":"allow-once"},{"name":"Reject","optionId":"reject"}]}}}"#,
5193 )
5194 .expect("valid permission")
5195 .expect("permission event");
5196 assert!(matches!(
5197 event,
5198 AgentEvent::Permission { request, .. }
5199 if request.id == "t1"
5200 && request.title == "Write file"
5201 && request.options == ["Allow once", "Reject"]
5202 && request.option_ids == ["allow-once", "reject"]
5203 ));
5204 }
5205
5206 #[test]
5207 fn parses_acp_permission_request_as_json_rpc_request() {
5208 let event = parse_acp_notification(
5209 2,
5210 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
5211 )
5212 .expect("valid permission request")
5213 .expect("permission event");
5214 assert!(matches!(
5215 event,
5216 AgentEvent::Permission { request, .. }
5217 if request.id == "17"
5218 && request.title == "Write file"
5219 && request.options == ["allow-once", "reject"]
5220 && request.option_ids == ["allow-once", "reject"]
5221 ));
5222 }
5223
5224 #[tokio::test]
5225 async fn native_adapter_explicitly_rejects_permission_answers() {
5226 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5227 assert_eq!(
5228 adapter
5229 .answer_permission(
5230 "request".into(),
5231 PermissionAnswer::Selected {
5232 option_id: "allow".into()
5233 },
5234 )
5235 .await,
5236 Err(super::AdapterError::Unsupported("permission answer"))
5237 );
5238 }
5239
5240 #[tokio::test]
5241 async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
5242 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5243 adapter
5244 .set_mode("full-access".into())
5245 .await
5246 .expect("auto-pilot alias");
5247 assert!(matches!(
5248 adapter.next_event().await,
5249 Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
5250 ));
5251 }
5252
5253 #[tokio::test]
5254 async fn native_turns_receive_a_twenty_four_hour_timeout() {
5255 let script_path = unique_test_path("codeswarm-native-timeout", "sh");
5256 std::fs::write(
5257 &script_path,
5258 r#"#!/bin/sh
5259seen=0
5260while [ "$#" -gt 0 ]; do
5261 case "$1" in
5262 --print-timeout)
5263 shift
5264 [ "$1" = "1440m" ] || exit 2
5265 seen=$((seen + 1))
5266 ;;
5267 esac
5268 shift
5269done
5270[ "$seen" = 1 ] || exit 3
5271printf '%s\n' '{"event":"result","result":{"status":"SUCCESS","response":"timeout accepted"}}'
5272"#,
5273 )
5274 .unwrap();
5275 let mut adapter = AgyAdapter::with_session_id(
5276 0,
5277 std::env::current_dir().unwrap(),
5278 format!("sh {}", script_path.display()),
5279 "saved-session",
5280 );
5281 adapter.start().await.unwrap();
5282 adapter.next_event().await.unwrap().unwrap();
5283 adapter.next_event().await.unwrap().unwrap();
5284 for prompt in ["first task", "follow-up task"] {
5285 adapter.send_prompt(prompt.into()).await.unwrap();
5286 assert!(
5287 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "timeout accepted")
5288 );
5289 assert!(matches!(
5290 adapter.next_event().await,
5291 Some(Ok(AgentEvent::TurnComplete { .. }))
5292 ));
5293 }
5294 adapter.stop().await.unwrap();
5295 std::fs::remove_file(script_path).unwrap();
5296 }
5297
5298 #[tokio::test]
5299 async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
5300 let script_path = unique_test_path("codeswarm-agy-session", "sh");
5301 std::fs::write(
5302 &script_path,
5303 "#!/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",
5304 )
5305 .expect("write native test script");
5306 let mut adapter = AgyAdapter::new(
5307 0,
5308 std::env::current_dir().expect("cwd"),
5309 format!("sh {}", script_path.display()),
5310 );
5311 adapter.start().await.expect("start native adapter");
5312 assert!(adapter.next_event().await.is_some());
5314 assert!(adapter.next_event().await.is_some());
5315 adapter
5316 .send_prompt("first".into())
5317 .await
5318 .expect("first prompt");
5319 while !matches!(
5320 adapter.next_event().await,
5321 Some(Ok(AgentEvent::TurnComplete { .. }))
5322 ) {}
5323 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5324 adapter
5325 .send_prompt("follow up".into())
5326 .await
5327 .expect("follow-up prompt");
5328 while !matches!(
5329 adapter.next_event().await,
5330 Some(Ok(AgentEvent::TurnComplete { .. }))
5331 ) {}
5332 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5333 adapter.stop().await.expect("stop native adapter");
5334 std::fs::remove_file(script_path).expect("cleanup native script");
5335 }
5336
5337 #[tokio::test]
5338 async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
5339 let script_path = unique_test_path("codeswarm-agy-failure", "sh");
5340 std::fs::write(
5341 &script_path,
5342 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
5343 )
5344 .expect("write native test script");
5345 let mut adapter = AgyAdapter::new(
5346 0,
5347 std::env::current_dir().expect("cwd"),
5348 format!("sh {}", script_path.display()),
5349 );
5350 adapter.start().await.expect("start native adapter");
5351 assert!(adapter.next_event().await.is_some());
5352 assert!(adapter.next_event().await.is_some());
5353 adapter.send_prompt("fail".into()).await.expect("prompt");
5354 assert!(matches!(
5355 adapter.next_event().await,
5356 Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
5357 if detail == "agent failed"
5358 ));
5359 adapter.stop().await.expect("stop native adapter");
5360 std::fs::remove_file(script_path).expect("cleanup native script");
5361 }
5362
5363 #[tokio::test]
5364 async fn native_crash_reaps_process_and_retries_on_next_prompt() {
5365 let script_path = unique_test_path("codeswarm-agy-retry", "sh");
5366 let marker_path = unique_test_path("codeswarm-agy-retry-marker", "txt");
5367 std::fs::write(
5368 &script_path,
5369 format!(
5370 "#!/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",
5371 marker_path.display(),
5372 marker_path.display(),
5373 marker_path.display(),
5374 ),
5375 )
5376 .expect("write retry script");
5377 let mut adapter = AgyAdapter::new(
5378 0,
5379 std::env::current_dir().expect("cwd"),
5380 format!("sh {}", script_path.display()),
5381 );
5382 adapter.start().await.expect("start native adapter");
5383 assert!(adapter.next_event().await.is_some());
5384 assert!(adapter.next_event().await.is_some());
5385 adapter
5386 .send_prompt("first".into())
5387 .await
5388 .expect("first prompt");
5389 assert!(matches!(
5390 adapter.next_event().await,
5391 Some(Ok(AgentEvent::Failed { detail, .. })) if detail == "first crash"
5392 ));
5393 adapter
5394 .send_prompt("retry".into())
5395 .await
5396 .expect("retry prompt starts a fresh process");
5397 assert!(matches!(
5398 adapter.next_event().await,
5399 Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5400 ));
5401 assert!(matches!(
5402 adapter.next_event().await,
5403 Some(Ok(AgentEvent::TurnComplete { .. }))
5404 ));
5405 adapter.stop().await.expect("stop native adapter");
5406 std::fs::remove_file(script_path).expect("cleanup retry script");
5407 std::fs::remove_file(marker_path).expect("cleanup retry marker");
5408 }
5409
5410 #[tokio::test]
5411 async fn acp_adapter_initializes_session_and_completes_a_prompt() {
5412 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"}}'"#;
5413 let cwd = std::env::current_dir().expect("cwd");
5414 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5415 adapter.start().await.expect("initialize");
5416 assert!(matches!(
5417 adapter.next_event().await,
5418 Some(Ok(AgentEvent::ModesReplaced { .. }))
5419 ));
5420 assert!(matches!(
5421 adapter.next_event().await,
5422 Some(Ok(AgentEvent::Ready { .. }))
5423 ));
5424 adapter.send_prompt("hello".into()).await.expect("prompt");
5425 assert!(matches!(
5426 adapter.next_event().await,
5427 Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
5428 ));
5429 assert!(matches!(
5430 adapter.next_event().await,
5431 Some(Ok(AgentEvent::TurnComplete { .. }))
5432 ));
5433 }
5434
5435 #[tokio::test]
5436 async fn acp_empty_end_turn_is_a_failed_turn_not_an_echo() {
5437 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"}}'"#;
5440 let cwd = std::env::current_dir().expect("cwd");
5441 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5442 adapter.start().await.expect("initialize");
5443 assert!(matches!(
5444 adapter.next_event().await,
5445 Some(Ok(AgentEvent::Ready { .. }))
5446 ));
5447 adapter
5448 .send_prompt("say hello".into())
5449 .await
5450 .expect("prompt");
5451 assert!(matches!(
5452 adapter.next_event().await,
5453 Some(Ok(AgentEvent::UserText { .. }))
5454 ));
5455 assert!(matches!(
5456 adapter.next_event().await,
5457 Some(Ok(AgentEvent::Failed {
5458 slot: 1,
5459 started: true,
5460 detail,
5461 })) if detail.contains("no agent output")
5462 ));
5463 }
5464
5465 #[test]
5466 fn parses_opencode_models_shape_without_config_options() {
5467 let session = serde_json::json!({
5468 "sessionId": "s",
5469 "models": {
5470 "currentModelId": "opencode-go/m",
5471 "availableModels": [
5472 {"modelId": "opencode-go/m", "name": "M"},
5473 {"modelId": "other/m2", "name": "M2"},
5474 ],
5475 },
5476 "modes": {"currentModeId": "build", "availableModes": []},
5477 });
5478 let (config_id, models, current) =
5479 super::parse_model_config(&session).expect("models shape");
5480 assert_eq!(config_id, "model");
5481 assert_eq!(models.len(), 2);
5482 assert_eq!(current.as_deref(), Some("opencode-go/m"));
5483 }
5484
5485 #[tokio::test]
5486 async fn acp_output_token_limit_is_reported_as_a_failed_turn() {
5487 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"}}'"#;
5488 let cwd = std::env::current_dir().expect("cwd");
5489 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5490 adapter.start().await.expect("initialize");
5491 assert!(matches!(
5492 adapter.next_event().await,
5493 Some(Ok(AgentEvent::Ready { .. }))
5494 ));
5495 adapter.send_prompt("hello".into()).await.expect("prompt");
5496 assert!(matches!(
5497 adapter.next_event().await,
5498 Some(Ok(AgentEvent::Failed {
5499 slot: 1,
5500 started: true,
5501 detail,
5502 })) if detail.contains("output token limit")
5503 ));
5504 }
5505
5506 #[tokio::test]
5507 async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
5508 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"}}'"#;
5509 let cwd = std::env::current_dir().expect("cwd");
5510 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5511 adapter.start().await.expect("initialize");
5512 assert!(adapter.next_event().await.is_some());
5513 assert!(adapter.next_event().await.is_some());
5514
5515 for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
5516 adapter.send_prompt(prompt.into()).await.expect("prompt");
5517 assert!(matches!(
5518 adapter.next_event().await,
5519 Some(Ok(AgentEvent::Text { text, .. })) if text == expected
5520 ));
5521 assert!(matches!(
5522 adapter.next_event().await,
5523 Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
5524 ));
5525 }
5526 adapter.stop().await.expect("stop");
5527 }
5528
5529 #[tokio::test]
5530 async fn empty_acp_mode_catalog_disables_mode_control() {
5531 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":[]}}}'"#;
5532 let cwd = std::env::current_dir().expect("cwd");
5533 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5534 adapter.start().await.expect("initialize");
5535 assert!(!adapter.capabilities().supports_modes);
5536 assert!(matches!(
5537 adapter.next_event().await,
5538 Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
5539 ));
5540 assert!(matches!(
5541 adapter.next_event().await,
5542 Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
5543 ));
5544 adapter.stop().await.expect("stop");
5545 }
5546
5547 #[tokio::test]
5548 async fn acp_models_are_discovered_live_and_changed_through_session_config() {
5549 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"#;
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_models);
5554 assert!(matches!(
5555 adapter.next_event().await,
5556 Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
5557 if config_id == "model"
5558 && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
5559 && current_model.as_deref() == Some("fast")
5560 ));
5561 assert!(matches!(
5562 adapter.next_event().await,
5563 Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
5564 ));
5565 adapter.set_model("smart".into()).await.expect("set model");
5566 assert!(adapter.set_model("invented".into()).await.is_err());
5567 adapter.stop().await.expect("stop");
5568 }
5569
5570 #[tokio::test]
5571 async fn acp_mode_change_is_acknowledged_without_provider_notification() {
5572 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":{}}'"#;
5573 let cwd = std::env::current_dir().expect("cwd");
5574 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5575 adapter.start().await.expect("initialize");
5576 adapter
5577 .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
5578 .await
5579 .expect("set mode");
5580 assert!(matches!(
5581 adapter.next_event().await,
5582 Some(Ok(AgentEvent::ModesReplaced { .. }))
5583 ));
5584 assert!(matches!(
5585 adapter.next_event().await,
5586 Some(Ok(AgentEvent::Ready { .. }))
5587 ));
5588 assert!(matches!(
5589 adapter.next_event().await,
5590 Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
5591 ));
5592 adapter.stop().await.expect("stop");
5593 }
5594
5595 #[tokio::test]
5596 async fn acp_reload_preserves_a_loadable_session_id() {
5597 let cwd = std::env::current_dir().expect("cwd");
5598 let mut adapter = AcpAdapter::with_session_id(
5599 0,
5600 cwd,
5601 "__codeswarm_missing_acp_for_reload_test__",
5602 Vec::new(),
5603 "saved-session",
5604 );
5605 adapter.capabilities.supports_session_load = true;
5606 assert!(adapter.reload().await.is_err());
5610 assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
5611 }
5612
5613 #[tokio::test]
5614 async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
5615 let cwd = std::env::current_dir().expect("cwd");
5616 let mut adapter = AcpAdapter::with_session_id(
5617 0,
5618 cwd,
5619 "__codeswarm_missing_nonloadable_acp__",
5620 Vec::new(),
5621 "stale-session",
5622 );
5623 adapter.capabilities.supports_session_load = false;
5624 assert!(adapter.reload().await.is_err());
5625 assert_eq!(adapter.session_id, None);
5626 }
5627
5628 #[tokio::test]
5629 async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
5630 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"}}'"#;
5631 let cwd = std::env::current_dir().expect("cwd");
5632 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5633 adapter.start().await.expect("initialize");
5634 assert!(matches!(
5635 adapter.next_event().await,
5636 Some(Ok(AgentEvent::Ready { .. }))
5637 ));
5638 adapter.send_prompt("hello".into()).await.expect("prompt");
5639 assert!(matches!(
5640 adapter.next_event().await,
5641 Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
5642 ));
5643 assert!(matches!(
5644 adapter.next_event().await,
5645 Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
5646 ));
5647 }
5648
5649 #[test]
5650 fn acp_tool_patches_preserve_fields_and_honor_explicit_replacements() {
5651 let mut tools = std::collections::BTreeMap::new();
5652 let first = serde_json::json!({"sessionUpdate":"tool_call", "toolCallId":"read", "title":"Read config", "status":"in_progress",
5653 "content":[{"type":"content", "content":{"type":"text", "text":"old output"}}]});
5654 let initial = super::normalize_acp_tool(&first, &mut tools).unwrap();
5655 assert_eq!(initial.detail.as_deref(), Some("old output"));
5656 let completed = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call_update","toolCallId":"read","status":"completed"}), &mut tools).unwrap();
5657 assert_eq!(completed.title, "Read config");
5658 assert_eq!(completed.detail.as_deref(), Some("old output"));
5659 assert_eq!(completed.status, ToolStatus::Completed);
5660 let malformed = super::normalize_acp_tool(
5661 &serde_json::json!({"toolCallId":"read","title":3,"status":"unknown","content":null}),
5662 &mut tools,
5663 )
5664 .unwrap();
5665 assert_eq!(malformed, completed);
5666 let replaced = super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read","content":[false,{"type":"content","content":{"type":"text","text":"new output"}}]}), &mut tools).unwrap();
5667 assert_eq!(replaced.detail.as_deref(), Some("new output"));
5668 let cleared = super::normalize_acp_tool(
5669 &serde_json::json!({"toolCallId":"read","content":[]}),
5670 &mut tools,
5671 )
5672 .unwrap();
5673 assert_eq!(cleared.detail, None);
5674 let raw = super::normalize_acp_tool(
5675 &serde_json::json!({"toolCallId":"read","rawOutput":{"ok":true}}),
5676 &mut tools,
5677 )
5678 .unwrap();
5679 assert_eq!(raw.detail.as_deref(), Some("{\"ok\":true}"));
5680 let fresh = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call","toolCallId":"read","title":"New call"}), &mut tools).unwrap();
5681 assert_eq!(fresh.status, ToolStatus::Pending);
5682 assert_eq!(fresh.detail, None);
5683 for invalid in [
5684 serde_json::json!({}),
5685 serde_json::json!({"toolCallId":7}),
5686 serde_json::json!({"toolCallId":" "}),
5687 ] {
5688 assert!(super::normalize_acp_tool(&invalid, &mut tools).is_none());
5689 }
5690 assert_eq!(tools.len(), 1);
5691 super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read "}), &mut tools).unwrap();
5693 assert_eq!(tools.len(), 2);
5694 }
5695
5696 #[tokio::test]
5697 async fn acp_tool_status_only_notifications_retain_name_and_output() {
5698 let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'
5699read _; echo '{"id":2,"result":{"sessionId":"s"}}'
5700read _
5701echo '{"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"}}]}}}'
5702echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"r","status":"completed"}}}'
5703echo '{"id":3,"result":{"stopReason":"end_turn"}}'"#;
5704 let mut adapter = AcpAdapter::new(
5705 0,
5706 std::env::current_dir().unwrap(),
5707 "sh",
5708 vec!["-c".into(), script.into()],
5709 );
5710 adapter.start().await.unwrap();
5711 adapter.next_event().await.unwrap().unwrap();
5712 adapter.send_prompt("read".into()).await.unwrap();
5713 for status in [ToolStatus::Running, ToolStatus::Completed] {
5714 let Some(Ok(AgentEvent::Tool { update, .. })) = adapter.next_event().await else {
5715 panic!("tool event");
5716 };
5717 assert_eq!(update.status, status);
5718 assert_eq!(update.title, "Read config");
5719 assert_eq!(update.detail.as_deref(), Some("file content"));
5720 }
5721 assert!(matches!(
5722 adapter.next_event().await,
5723 Some(Ok(AgentEvent::TurnComplete { .. }))
5724 ));
5725 adapter.stop().await.unwrap();
5726 }
5727
5728 #[tokio::test]
5729 async fn acp_reload_discards_old_queued_events_and_catalogs() {
5730 let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"id":2,"result":{"sessionId":"new"}}'"#;
5731 let mut adapter = AcpAdapter::new(
5732 0,
5733 std::env::current_dir().unwrap(),
5734 "sh",
5735 vec!["-c".into(), script.into()],
5736 );
5737 adapter.start().await.unwrap();
5738 adapter.queued_events.push_back(Ok(AgentEvent::Text {
5739 slot: 0,
5740 text: "stale".into(),
5741 }));
5742 adapter.modes = vec![Mode {
5743 id: "stale".into(),
5744 label: "Stale".into(),
5745 }];
5746 adapter.next_request_id = 1;
5748 adapter.reload().await.unwrap();
5749 assert!(adapter.modes.is_empty());
5750 assert_eq!(adapter.queued_events.len(), 1);
5751 assert!(matches!(
5752 adapter.next_event().await,
5753 Some(Ok(AgentEvent::Ready { .. }))
5754 ));
5755 adapter.stop().await.unwrap();
5756 assert!(adapter.queued_events.is_empty());
5757 }
5758
5759 #[tokio::test]
5760 async fn acp_load_replays_history_without_starting_a_turn() {
5761 let script = r#"
5762read _
5763echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'
5764read request
5765case "$request" in *session/load*) ;; *) exit 2;; esac
5766echo '{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"text":"old question"}}}}'
5767echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"old answer"}}}}'
5768echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_thought_chunk","content":{"text":"old reasoning"}}}}'
5769echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call","toolCallId":"old-tool","title":"Read","status":"in_progress"}}}'
5770echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"old-tool","title":"Read","status":"completed"}}}'
5771echo '{"jsonrpc":"2.0","id":2,"result":{}}'
5772read request
5773case "$request" in *session/prompt*) ;; *) exit 3;; esac
5774echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"new answer"}}}}'
5775echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
5776"#;
5777 let mut adapter = AcpAdapter::with_session_id(
5778 2,
5779 std::env::current_dir().unwrap(),
5780 "sh",
5781 vec!["-c".into(), script.into()],
5782 "saved",
5783 );
5784 adapter.start().await.unwrap();
5785 let mut state = crate::SessionState::new(3);
5786 for _ in 0..5 {
5787 let event = adapter.next_event().await.unwrap().unwrap();
5788 assert!(matches!(&event, AgentEvent::History { slot: 2, .. }));
5789 crate::reduce(&mut state, event);
5790 assert_eq!(state.active_slot, None);
5791 }
5792 assert!(matches!(
5793 adapter.next_event().await,
5794 Some(Ok(AgentEvent::Ready { slot: 2, .. }))
5795 ));
5796 adapter.send_prompt("new question".into()).await.unwrap();
5797 assert!(
5798 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "new answer")
5799 );
5800 assert!(matches!(
5801 adapter.next_event().await,
5802 Some(Ok(AgentEvent::TurnComplete { slot: 2 }))
5803 ));
5804 adapter.stop().await.unwrap();
5805 }
5806
5807 #[tokio::test]
5808 async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
5809 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
5810 let cwd = std::env::current_dir().expect("cwd");
5811 let mut adapter = AcpAdapter::with_session_id(
5812 0,
5813 cwd,
5814 "sh",
5815 vec!["-c".into(), script.into()],
5816 "existing-session",
5817 );
5818 adapter.start().await.expect("load existing session");
5819 assert!(matches!(
5820 adapter.next_event().await,
5821 Some(Ok(AgentEvent::Ready { .. }))
5822 ));
5823 }
5824
5825 #[tokio::test]
5826 async fn acp_start_failure_reaps_transport_process() {
5827 let mut adapter = AcpAdapter::new(
5832 0,
5833 std::env::current_dir().expect("cwd"),
5834 "sh",
5835 vec!["-c".into(), "printf 'not-json\\n'".into()],
5836 );
5837 assert!(adapter.start().await.is_err());
5838 assert!(adapter.child.is_none());
5839 assert!(adapter.reader.is_none());
5840 }
5841
5842 #[tokio::test]
5843 async fn acp_transport_crash_is_reloaded_before_the_next_prompt() {
5844 let marker = unique_test_path("codeswarm-acp-retry", "count");
5845 let script = format!(
5846 r#"count=0
5847if [ -f '{0}' ]; then count=$(cat '{0}'); fi
5848count=$((count + 1))
5849printf '%s' "$count" > '{0}'
5850while IFS= read -r request; do
5851 id=$(printf '%s' "$request" | sed -n 's/.*"id":\([0-9][0-9]*\).*/\1/p')
5852 case "$request" in
5853 *initialize*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"agentCapabilities":{{"loadSession":true}}}}}}' ;;
5854 *session/new*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"sessionId":"saved-session"}}}}' ;;
5855 *session/load*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{}}}}' ;;
5856 *session/prompt*)
5857 if [ "$count" = 1 ]; then exit 0; fi
5858 printf '%s\n' '{{"jsonrpc":"2.0","method":"session/update","params":{{"update":{{"sessionUpdate":"agent_message_chunk","content":{{"text":"recovered"}}}}}}}}'
5859 printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"stopReason":"end_turn"}}}}'
5860 ;;
5861 esac
5862done
5863"#,
5864 marker.display()
5865 );
5866 let mut adapter = AcpAdapter::new(
5867 0,
5868 std::env::current_dir().expect("cwd"),
5869 "sh",
5870 vec!["-c".into(), script],
5871 );
5872 adapter.start().await.expect("initial ACP startup");
5873 assert!(matches!(
5874 adapter.next_event().await,
5875 Some(Ok(AgentEvent::Ready { .. }))
5876 ));
5877 adapter
5878 .send_prompt("first".into())
5879 .await
5880 .expect("first prompt");
5881 assert!(matches!(
5882 adapter.next_event().await,
5883 Some(Err(AdapterError::Transport(_)))
5884 ));
5885 assert!(adapter.child.is_none());
5886 assert!(adapter.reader.is_none());
5887 assert_eq!(adapter.session_id(), Some("saved-session".into()));
5888
5889 adapter.reload().await.expect("reload ACP transport");
5890 assert!(matches!(
5891 adapter.next_event().await,
5892 Some(Ok(AgentEvent::Ready { .. }))
5893 ));
5894 adapter
5895 .send_prompt("retry".into())
5896 .await
5897 .expect("retry prompt");
5898 assert!(matches!(
5899 adapter.next_event().await,
5900 Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5901 ));
5902 assert!(matches!(
5903 adapter.next_event().await,
5904 Some(Ok(AgentEvent::TurnComplete { .. }))
5905 ));
5906 adapter.stop().await.expect("stop ACP");
5907 std::fs::remove_file(marker).expect("cleanup marker");
5908 }
5909
5910 #[tokio::test]
5911 async fn coordinator_reload_replays_context_and_reintroduces_a_crashed_slot() {
5912 let prompts = Arc::new(Mutex::new(Vec::new()));
5913 let healthy = ScriptedAdapter::new(
5914 0,
5915 AgentCapabilities::default(),
5916 [
5917 AgentEvent::Text {
5918 slot: 0,
5919 text: "peer context".into(),
5920 },
5921 AgentEvent::TurnComplete { slot: 0 },
5922 ],
5923 );
5924 let probe = ReloadProbeAdapter {
5925 slot: 1,
5926 crashed: false,
5927 reloaded: false,
5928 events: VecDeque::new(),
5929 prompts: Arc::clone(&prompts),
5930 };
5931 let mut relay = RelayHost::new(
5932 vec![
5933 AdapterHost::new(Box::new(healthy), None),
5934 AdapterHost::new(Box::new(probe), None),
5935 ],
5936 8,
5937 )
5938 .expect("relay");
5939 relay.start().await.expect("start");
5940 relay
5941 .run_turn("original task", 0)
5942 .await
5943 .expect("first turn");
5944 relay.run_turn("", 0).await.expect("crashed turn");
5945 relay.relay_mut().enqueue_human("retry", Some(1));
5946 relay.run_turn("", 0).await.expect("reloaded turn");
5947
5948 let prompts = prompts.lock().expect("prompts");
5949 let retry = prompts.last().expect("retry prompt");
5950 assert!(retry.contains("You are Reload probe"), "{retry}");
5951 assert!(retry.contains("original task"), "{retry}");
5952 assert!(retry.contains("peer context"), "{retry}");
5953 assert!(retry.contains("retry"), "{retry}");
5954 }
5955
5956 #[tokio::test]
5957 async fn acp_adapter_answers_permission_json_rpc_requests() {
5958 let path = std::env::temp_dir().join(format!(
5959 "codeswarm-permission-answer-{}",
5960 std::process::id()
5961 ));
5962 let script = format!(
5963 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"}}}}'"#,
5964 path.display()
5965 );
5966 let mut adapter = AcpAdapter::new(
5967 0,
5968 std::env::current_dir().expect("cwd"),
5969 "sh",
5970 vec!["-c".into(), script],
5971 );
5972 adapter.start().await.expect("start ACP");
5973 assert!(matches!(
5974 adapter.next_event().await,
5975 Some(Ok(AgentEvent::Ready { .. }))
5976 ));
5977 adapter.send_prompt("do it".into()).await.expect("prompt");
5978 assert!(matches!(
5979 adapter.next_event().await,
5980 Some(Ok(AgentEvent::Permission { request, .. }))
5981 if request.id == "9"
5982 && request.options == ["Allow once"]
5983 && request.option_ids == ["allow-once"]
5984 ));
5985 adapter
5986 .answer_permission(
5987 "9".into(),
5988 PermissionAnswer::Selected {
5989 option_id: "allow-once".into(),
5990 },
5991 )
5992 .await
5993 .expect("permission answer");
5994 assert!(matches!(
5995 adapter.next_event().await,
5996 Some(Ok(AgentEvent::TurnComplete { .. }))
5997 ));
5998 let answer: Value = serde_json::from_str(
5999 &std::fs::read_to_string(&path).expect("captured permission answer"),
6000 )
6001 .expect("valid JSON-RPC answer");
6002 assert_eq!(answer["id"], 9);
6003 assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
6004 assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
6005 std::fs::remove_file(path).expect("cleanup");
6006 }
6007
6008 #[test]
6009 fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
6010 let event = parse_acp_notification(
6011 0,
6012 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
6013 )
6014 .expect("valid JSON-RPC request");
6015 assert!(event.is_none());
6016 }
6017
6018 #[tokio::test]
6019 async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
6020 let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
6021 std::fs::write(
6022 &script_path,
6023 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
6024 )
6025 .expect("write native test script");
6026 let mut adapter = AgyAdapter::new(
6027 0,
6028 std::env::current_dir().expect("cwd"),
6029 format!("sh {}", script_path.display()),
6030 );
6031 adapter.start().await.expect("start native adapter");
6032 assert!(adapter.next_event().await.is_some());
6033 assert!(adapter.next_event().await.is_some());
6034 adapter
6035 .send_prompt("continue".into())
6036 .await
6037 .expect("prompt");
6038 assert!(matches!(
6039 adapter.next_event().await,
6040 Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
6041 ));
6042 assert!(matches!(
6043 adapter.next_event().await,
6044 Some(Ok(AgentEvent::TurnComplete { .. }))
6045 ));
6046 adapter.stop().await.expect("stop native adapter");
6047 std::fs::remove_file(script_path).expect("cleanup native script");
6048 }
6049
6050 #[test]
6051 fn acp_workspace_file_access_is_root_bound_and_size_limited() {
6052 let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
6053 let _ = std::fs::remove_dir_all(&root);
6054 std::fs::create_dir_all(&root).expect("workspace");
6055 std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
6056 let outside =
6057 std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
6058 std::fs::write(&outside, "secret").expect("outside file");
6059 let link = root.join("outside-link");
6060 #[cfg(unix)]
6061 std::os::unix::fs::symlink(&outside, &link).expect("symlink");
6062 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
6063
6064 assert_eq!(
6065 adapter
6066 .read_workspace_text("inside.txt", Some(2), Some(1))
6067 .expect("read inside"),
6068 "two"
6069 );
6070 std::fs::write(
6071 root.join("large.txt"),
6072 vec![b'x'; MAX_FILE_READ_BYTES + 1024],
6073 )
6074 .expect("large file");
6075 let bounded = adapter
6076 .read_workspace_text("large.txt", None, None)
6077 .expect("bounded read");
6078 assert!(bounded.len() <= MAX_FILE_READ_BYTES);
6079 #[cfg(unix)]
6080 {
6081 std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
6082 .expect("internal symlink");
6083 assert_eq!(
6084 adapter
6085 .read_workspace_text("inside-link", None, None)
6086 .expect("read internal symlink"),
6087 "one\ntwo\nthree\n"
6088 );
6089 }
6090 assert!(adapter.workspace_path("../codeswarm-outside").is_err());
6091 assert!(
6092 adapter
6093 .workspace_path(&outside.display().to_string())
6094 .is_err()
6095 );
6096 #[cfg(unix)]
6097 assert!(adapter.workspace_path("outside-link").is_err());
6098 #[cfg(unix)]
6099 std::fs::remove_file(link).expect("cleanup symlink");
6100 #[cfg(unix)]
6101 std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
6102 std::fs::remove_file(outside).expect("cleanup outside");
6103 std::fs::remove_dir_all(root).expect("cleanup workspace");
6104 }
6105
6106 #[tokio::test]
6107 async fn running_terminal_output_omits_exit_status_until_completion() {
6108 let root = unique_test_path("codeswarm-terminal-output", "dir");
6109 std::fs::create_dir_all(&root).expect("workspace");
6110 let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
6111 let result = adapter
6112 .terminal_create(&serde_json::json!({
6113 "command": "sh",
6114 "args": ["-c", "sleep 0.2; printf done"],
6115 "cwd": ".",
6116 }))
6117 .await
6118 .expect("terminal create");
6119 let id = result["terminalId"].as_str().expect("terminal id");
6120 let output = adapter.terminal_output(id).await.expect("terminal output");
6121 assert!(output.get("exitStatus").is_none());
6122 if let Some(terminal) = adapter.terminals.remove(id) {
6123 terminal.stop().await;
6124 }
6125 std::fs::remove_dir_all(root).expect("cleanup workspace");
6126 }
6127
6128 #[tokio::test]
6129 async fn acp_adapter_answers_workspace_read_requests() {
6130 let root =
6131 std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
6132 let _ = std::fs::remove_dir_all(&root);
6133 std::fs::create_dir_all(&root).expect("workspace");
6134 let source = root.join("inside.txt");
6135 let answer = root.join("answer.json");
6136 std::fs::write(&source, "workspace content").expect("source");
6137 let script = format!(
6138 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"}}}}'"#,
6139 source.display(),
6140 answer.display(),
6141 );
6142 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6143 adapter.start().await.expect("start ACP");
6144 assert!(matches!(
6145 adapter.next_event().await,
6146 Some(Ok(AgentEvent::Ready { .. }))
6147 ));
6148 adapter.send_prompt("read it".into()).await.expect("prompt");
6149 assert!(matches!(
6150 adapter.next_event().await,
6151 Some(Ok(AgentEvent::TurnComplete { .. }))
6152 ));
6153 let response: Value =
6154 serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
6155 .expect("response JSON");
6156 assert_eq!(response["id"], 9);
6157 assert_eq!(response["result"]["content"], "workspace content");
6158 adapter.stop().await.expect("stop ACP");
6159 std::fs::remove_dir_all(root).expect("cleanup workspace");
6160 }
6161
6162 #[tokio::test]
6163 async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
6164 let root =
6165 std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
6166 let _ = std::fs::remove_dir_all(&root);
6167 std::fs::create_dir_all(&root).expect("workspace");
6168 let create_request = serde_json::json!({
6169 "jsonrpc": "2.0",
6170 "id": 9,
6171 "method": "terminal/create",
6172 "params": {
6173 "sessionId": "s1",
6174 "command": "sh",
6175 "args": ["-c", "sleep 0.1; printf terminal-ok"],
6176 "cwd": ".",
6177 },
6178 });
6179 let wait_request = serde_json::json!({
6180 "jsonrpc": "2.0",
6181 "id": 10,
6182 "method": "terminal/wait_for_exit",
6183 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6184 });
6185 let output_request = serde_json::json!({
6186 "jsonrpc": "2.0",
6187 "id": 11,
6188 "method": "terminal/output",
6189 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6190 });
6191 let create_answer = root.join("create-answer.json");
6192 let wait_answer = root.join("wait-answer.json");
6193 let output_answer = root.join("output-answer.json");
6194 let script = format!(
6195 "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\"}}}}'",
6196 create_request,
6197 create_answer.display(),
6198 wait_request,
6199 wait_answer.display(),
6200 output_request,
6201 output_answer.display(),
6202 );
6203 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6204 adapter.start().await.expect("start ACP");
6205 assert!(matches!(
6206 adapter.next_event().await,
6207 Some(Ok(AgentEvent::Ready { .. }))
6208 ));
6209 adapter
6210 .send_prompt("run terminal".into())
6211 .await
6212 .expect("prompt");
6213 let mut saw_complete = false;
6214 for _ in 0..6 {
6215 match adapter.next_event().await {
6216 Some(Ok(AgentEvent::TurnComplete { .. })) => {
6217 saw_complete = true;
6218 break;
6219 }
6220 Some(_) => {}
6221 None => break,
6222 }
6223 }
6224 assert!(saw_complete, "terminal requests should not stall ACP");
6225 let create: Value = serde_json::from_str(
6226 &std::fs::read_to_string(&create_answer).expect("captured create response"),
6227 )
6228 .expect("create JSON");
6229 assert_eq!(create["result"]["terminalId"], "terminal-1");
6230 let output: Value = serde_json::from_str(
6231 &std::fs::read_to_string(&output_answer).expect("captured output response"),
6232 )
6233 .expect("output JSON");
6234 assert!(
6235 output["result"]["output"]
6236 .as_str()
6237 .unwrap_or_default()
6238 .contains("terminal-ok"),
6239 "output response: {output}"
6240 );
6241 adapter.stop().await.expect("stop ACP");
6242 std::fs::remove_dir_all(root).expect("cleanup workspace");
6243 }
6244
6245 #[tokio::test]
6246 async fn host_reduces_and_persists_adapter_events() {
6247 let path =
6248 std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
6249 let adapter = ScriptedAdapter::new(
6250 0,
6251 AgentCapabilities::default(),
6252 [AgentEvent::Text {
6253 slot: 0,
6254 text: "hello".into(),
6255 }],
6256 );
6257 let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
6258 host.start().await.expect("start");
6259 host.next_effects()
6260 .await
6261 .expect("event")
6262 .expect("valid event");
6263 assert_eq!(host.state.public_text[0].1, "hello");
6264 assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
6265 std::fs::remove_file(path).expect("cleanup");
6266 }
6267
6268 #[tokio::test]
6269 async fn relay_applies_default_policy_before_the_first_prompt() {
6270 let first_log = Arc::new(Mutex::new(Vec::new()));
6271 let second_log = Arc::new(Mutex::new(Vec::new()));
6272 let first = AdapterHost::new(
6273 Box::new(ModeOrderAdapter {
6274 slot: 0,
6275 log: Arc::clone(&first_log),
6276 phase: 0,
6277 }),
6278 None,
6279 );
6280 let second = AdapterHost::new(
6281 Box::new(ModeOrderAdapter {
6282 slot: 1,
6283 log: Arc::clone(&second_log),
6284 phase: 0,
6285 }),
6286 None,
6287 );
6288 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6289 relay.start().await.expect("start and synchronize policy");
6290 relay.run_turn("task", 0).await.expect("first turn");
6291 {
6292 let log = first_log.lock().expect("log");
6293 assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
6294 }
6295 assert_eq!(
6296 second_log.lock().expect("log").as_slice(),
6297 ["start", "mode:yolo"]
6298 );
6299 let added_log = Arc::new(Mutex::new(Vec::new()));
6300 relay
6301 .add_agent(
6302 AdapterHost::new(
6303 Box::new(ModeOrderAdapter {
6304 slot: 2,
6305 log: Arc::clone(&added_log),
6306 phase: 0,
6307 }),
6308 None,
6309 ),
6310 "Added",
6311 "added.example",
6312 "added-agent",
6313 )
6314 .await
6315 .expect("add with synchronized policy");
6316 assert_eq!(
6317 added_log.lock().expect("log").as_slice(),
6318 ["start", "mode:yolo"]
6319 );
6320 relay.drop_agent(2).await.expect("drop added agent");
6321 added_log.lock().expect("log").clear();
6322 relay
6323 .reload(2)
6324 .await
6325 .expect("reload with synchronized policy");
6326 assert_eq!(
6327 added_log.lock().expect("log").as_slice(),
6328 ["reload", "mode:yolo"]
6329 );
6330 }
6331
6332 #[tokio::test]
6333 async fn acp_roster_is_ready_before_any_prompt_is_sent() {
6334 let hosts = (0..2)
6335 .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
6336 .collect::<Vec<_>>();
6337 let startup_events = Arc::new(Mutex::new(Vec::new()));
6338 let captured = Arc::clone(&startup_events);
6339 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6340 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6341
6342 relay.start().await.expect("complete startup handshake");
6343
6344 assert!(relay.dispatches().is_empty());
6345 let ready_slots = startup_events
6346 .lock()
6347 .expect("events")
6348 .iter()
6349 .filter_map(|event| match event {
6350 AgentEvent::Ready { slot, .. } => Some(*slot),
6351 _ => None,
6352 })
6353 .collect::<Vec<_>>();
6354 assert_eq!(ready_slots, vec![0, 1]);
6355 }
6356
6357 #[tokio::test]
6358 async fn independent_roster_adapters_start_concurrently() {
6359 let barrier = Arc::new(tokio::sync::Barrier::new(2));
6360 let hosts = (0..2)
6361 .map(|slot| {
6362 AdapterHost::new(
6363 Box::new(ConcurrentStartAdapter {
6364 slot,
6365 barrier: Arc::clone(&barrier),
6366 }),
6367 None,
6368 )
6369 })
6370 .collect::<Vec<_>>();
6371 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6372 tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
6373 .await
6374 .expect("startup should not serialize barrier participants")
6375 .expect("startup succeeds");
6376 }
6377
6378 #[tokio::test]
6379 async fn relay_host_dispatches_turns_sequentially() {
6380 let capabilities = AgentCapabilities {
6381 supports_cancel: true,
6382 ..AgentCapabilities::default()
6383 };
6384 let first = ScriptedAdapter::new(
6385 0,
6386 capabilities.clone(),
6387 [
6388 AgentEvent::Text {
6389 slot: 0,
6390 text: "first".into(),
6391 },
6392 AgentEvent::TurnComplete { slot: 0 },
6393 ],
6394 );
6395 let second = ScriptedAdapter::new(
6396 1,
6397 capabilities,
6398 [
6399 AgentEvent::Text {
6400 slot: 1,
6401 text: "review".into(),
6402 },
6403 AgentEvent::TurnComplete { slot: 1 },
6404 ],
6405 );
6406 let hosts = vec![
6407 AdapterHost::new(Box::new(first), None),
6408 AdapterHost::new(Box::new(second), None),
6409 ];
6410 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6411 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6412 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6413 let captured = std::sync::Arc::clone(&events);
6414 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6415 relay.start().await.expect("start");
6416 events.lock().expect("events").clear();
6417 assert!(matches!(
6418 relay.run_turn("task", 0).await.expect("first turn"),
6419 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6420 ));
6421 assert!(matches!(
6422 relay.run_turn("first", 0).await.expect("second turn"),
6423 crate::relay::RelayDecision::Dispatch {
6424 slot: 1,
6425 can_stop: true,
6426 ..
6427 }
6428 ));
6429 assert_eq!(
6430 relay
6431 .dispatches()
6432 .iter()
6433 .map(|(slot, _)| *slot)
6434 .collect::<Vec<_>>(),
6435 [0, 1]
6436 );
6437 assert!(relay.dispatches()[0].1.contains("You are Claude"));
6438 assert!(
6439 relay.dispatches()[0]
6440 .1
6441 .contains("CodeSwarm roster (ordered)")
6442 );
6443 assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
6444 assert!(relay.dispatches()[0].1.contains("2. Codex"));
6445 assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
6446 assert!(
6447 relay.dispatches()[1]
6448 .1
6449 .contains("stops all other agents and ends the entire automated relay")
6450 );
6451 assert!(relay.dispatches()[1].1.contains("Use it with extreme care"));
6452 assert!(
6453 relay.dispatches()[1]
6454 .1
6455 .contains("If there is any uncertainty, do not use it")
6456 );
6457 assert!(relay.dispatches()[0].1.contains("Do not use"));
6458 let lifecycle = events.lock().expect("events");
6459 let positions = lifecycle
6460 .iter()
6461 .filter_map(|event| match event {
6462 AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
6463 AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
6464 _ => None,
6465 })
6466 .collect::<Vec<_>>();
6467 assert_eq!(
6468 positions,
6469 [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
6470 );
6471 }
6472
6473 #[tokio::test]
6474 async fn failed_resume_does_not_stop_or_dispatch_to_healthy_peer() {
6475 let stops = Arc::new(AtomicUsize::new(0));
6476 let failed = FailingStartAdapter {
6477 slot: 0,
6478 stopped: stops.clone(),
6479 };
6480 let healthy = ScriptedAdapter::new(
6481 1,
6482 AgentCapabilities::default(),
6483 [
6484 AgentEvent::Text {
6485 slot: 1,
6486 text: "healthy response".into(),
6487 },
6488 AgentEvent::TurnComplete { slot: 1 },
6489 ],
6490 );
6491 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6492 let captured = events.clone();
6493 let mut relay = RelayHost::new(
6494 vec![
6495 AdapterHost::new(Box::new(failed), None),
6496 AdapterHost::new(Box::new(healthy), None),
6497 ],
6498 4,
6499 )
6500 .unwrap();
6501 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6502 relay.start_resuming().await.unwrap();
6503 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
6504 assert!(relay.dispatches().is_empty());
6505 assert_eq!(stops.load(Ordering::Relaxed), 1);
6506 assert!(
6507 events
6508 .lock()
6509 .unwrap()
6510 .iter()
6511 .any(|event| matches!(event, AgentEvent::Failed { slot: 0, .. }))
6512 );
6513 assert!(
6514 !events
6515 .lock()
6516 .unwrap()
6517 .iter()
6518 .any(|event| matches!(event, AgentEvent::Failed { slot: 1, .. }))
6519 );
6520 assert!(!relay.relay_mut().enqueue_human("do not retarget", Some(0)));
6521 assert!(
6522 relay
6523 .relay_mut()
6524 .enqueue_human("explicit healthy target", Some(1))
6525 );
6526 assert!(matches!(
6527 relay.run_turn("", 1).await.unwrap(),
6528 RelayDecision::Dispatch { slot: 1, .. }
6529 ));
6530 relay.stop().await.unwrap();
6531 }
6532
6533 #[tokio::test]
6534 async fn pair_strategy_wires_roles_into_non_direct_prompts() {
6535 let capabilities = AgentCapabilities::default();
6536 let first = ScriptedAdapter::new(
6537 0,
6538 capabilities.clone(),
6539 [
6540 AgentEvent::Text {
6541 slot: 0,
6542 text: "implemented".into(),
6543 },
6544 AgentEvent::TurnComplete { slot: 0 },
6545 AgentEvent::Text {
6546 slot: 0,
6547 text: format!("fixed review findings {STOP_TOKEN}"),
6548 },
6549 AgentEvent::TurnComplete { slot: 0 },
6550 ],
6551 );
6552 let second = ScriptedAdapter::new(
6553 1,
6554 capabilities,
6555 [
6556 AgentEvent::Text {
6557 slot: 1,
6558 text: "reviewed".into(),
6559 },
6560 AgentEvent::TurnComplete { slot: 1 },
6561 AgentEvent::Text {
6562 slot: 1,
6563 text: format!("approved {STOP_TOKEN}"),
6564 },
6565 AgentEvent::TurnComplete { slot: 1 },
6566 AgentEvent::Text {
6567 slot: 1,
6568 text: "new task".into(),
6569 },
6570 AgentEvent::TurnComplete { slot: 1 },
6571 ],
6572 );
6573 let hosts = vec![
6574 AdapterHost::new(Box::new(first), None),
6575 AdapterHost::new(Box::new(second), None),
6576 ];
6577 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6578 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6579 relay.relay_mut().set_strategy(CollaborationStrategy::Pair);
6580 relay.start().await.expect("start");
6581 assert!(matches!(
6582 relay.run_turn("task", 0).await.expect("implementer turn"),
6583 RelayDecision::Dispatch {
6584 slot: 0,
6585 can_stop: false,
6586 ..
6587 }
6588 ));
6589 let implementer_prompt = &relay.dispatches()[0].1;
6590 assert!(implementer_prompt.contains("you are the implementer"));
6591 assert!(implementer_prompt.contains("pair reviewer will review the result next"));
6592 assert!(!implementer_prompt.contains("you are the reviewer"));
6593 assert!(implementer_prompt.contains("Do not use"));
6594 assert!(matches!(
6595 relay.run_turn("", 0).await.expect("reviewer turn"),
6596 RelayDecision::Dispatch {
6597 slot: 1,
6598 can_stop: true,
6599 ..
6600 }
6601 ));
6602 let reviewer_prompt = &relay.dispatches()[1].1;
6603 assert!(reviewer_prompt.contains("you are the reviewer"));
6604 assert!(reviewer_prompt.contains("Claude handed off"));
6605 assert!(reviewer_prompt.contains("concrete defects"));
6606 assert!(reviewer_prompt.contains("concise approval"));
6607 assert!(reviewer_prompt.contains(STOP_TOKEN));
6608 assert!(!reviewer_prompt.contains("you are the implementer"));
6609 relay.run_turn("", 0).await.unwrap();
6610 assert!(relay.dispatches()[2].1.contains("you are the implementer"));
6611 assert!(relay.dispatches()[2].1.contains("Do not use"));
6612 assert!(matches!(
6613 relay.run_turn("", 0).await.unwrap(),
6614 RelayDecision::Dispatch { slot: 1, .. }
6615 ));
6616 assert!(relay.dispatches()[3].1.contains("you are the reviewer"));
6617 assert!(relay.relay_mut().enqueue_human("new task", Some(1)));
6618 relay.run_turn("", 1).await.unwrap();
6619 assert!(relay.dispatches()[4].1.contains("you are the implementer"));
6620 }
6621
6622 #[tokio::test]
6623 async fn solo_roster_and_direct_prompts_omit_pair_roles() {
6624 let solo = ScriptedAdapter::new(
6625 0,
6626 AgentCapabilities::default(),
6627 [
6628 AgentEvent::Text {
6629 slot: 0,
6630 text: "solo".into(),
6631 },
6632 AgentEvent::TurnComplete { slot: 0 },
6633 ],
6634 );
6635 let mut solo_relay =
6636 RelayHost::new(vec![AdapterHost::new(Box::new(solo), None)], 4).expect("relay");
6637 solo_relay
6638 .relay_mut()
6639 .set_strategy(CollaborationStrategy::Pair);
6640 solo_relay.start().await.expect("start");
6641 solo_relay.run_turn("task", 0).await.expect("solo turn");
6642 assert!(!solo_relay.dispatches()[0].1.contains("Pair role"));
6643
6644 let roster_first = ScriptedAdapter::new(
6645 0,
6646 AgentCapabilities::default(),
6647 [AgentEvent::TurnComplete { slot: 0 }],
6648 );
6649 let roster_second = ScriptedAdapter::new(
6650 1,
6651 AgentCapabilities::default(),
6652 [AgentEvent::TurnComplete { slot: 1 }],
6653 );
6654 let mut roster = RelayHost::new(
6655 vec![
6656 AdapterHost::new(Box::new(roster_first), None),
6657 AdapterHost::new(Box::new(roster_second), None),
6658 ],
6659 4,
6660 )
6661 .expect("relay");
6662 roster.start().await.expect("start");
6663 roster.run_turn("task", 0).await.expect("first turn");
6664 roster.run_turn("", 0).await.expect("second turn");
6665 assert!(!roster.dispatches()[0].1.contains("Pair role"));
6666 assert!(!roster.dispatches()[1].1.contains("Pair role"));
6667
6668 let pair_first = ScriptedAdapter::new(
6669 0,
6670 AgentCapabilities::default(),
6671 [AgentEvent::TurnComplete { slot: 0 }],
6672 );
6673 let pair_second = ScriptedAdapter::new(
6674 1,
6675 AgentCapabilities::default(),
6676 [AgentEvent::TurnComplete { slot: 1 }],
6677 );
6678 let mut pair = RelayHost::new(
6679 vec![
6680 AdapterHost::new(Box::new(pair_first), None),
6681 AdapterHost::new(Box::new(pair_second), None),
6682 ],
6683 4,
6684 )
6685 .expect("relay");
6686 pair.relay_mut().set_strategy(CollaborationStrategy::Pair);
6687 assert_eq!(pair.relay_mut().enqueue_direct(1, "private"), Ok(true));
6688 pair.start().await.expect("start");
6689 assert!(matches!(
6690 pair.run_turn("ignored", 0).await.expect("direct turn"),
6691 RelayDecision::Dispatch {
6692 slot: 1,
6693 direct: true,
6694 ..
6695 }
6696 ));
6697 let direct_prompt = &pair.dispatches()[0].1;
6698 assert!(direct_prompt.contains("private"));
6699 assert!(!direct_prompt.contains("Pair role"));
6700 }
6701
6702 #[tokio::test]
6703 async fn relay_host_routes_around_a_usage_limited_agent() {
6704 let capabilities = AgentCapabilities::default();
6705 let first = ScriptedAdapter::new(
6706 0,
6707 capabilities.clone(),
6708 [
6709 AgentEvent::Text {
6710 slot: 0,
6711 text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
6712 credits or try again later."
6713 .into(),
6714 },
6715 AgentEvent::TurnComplete { slot: 0 },
6716 ],
6717 );
6718 let second = ScriptedAdapter::new(
6719 1,
6720 capabilities,
6721 [
6722 AgentEvent::Text {
6723 slot: 1,
6724 text: "review done".into(),
6725 },
6726 AgentEvent::TurnComplete { slot: 1 },
6727 ],
6728 );
6729 let hosts = vec![
6730 AdapterHost::new(Box::new(first), None),
6731 AdapterHost::new(Box::new(second), None),
6732 ];
6733 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6734 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6735 let captured = std::sync::Arc::clone(&events);
6736 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6737 relay.start().await.expect("start");
6738 events.lock().expect("events").clear();
6739 assert!(matches!(
6740 relay.run_turn("task", 0).await.expect("limited turn"),
6741 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6742 ));
6743 assert!(
6744 events
6745 .lock()
6746 .expect("events")
6747 .iter()
6748 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6749 );
6750 assert!(matches!(
6752 relay.run_turn("", 0).await.expect("next turn"),
6753 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6754 ));
6755 assert!(relay.relay().is_limited(0));
6756 relay.reload(0).await.expect("reload");
6760 assert!(!relay.relay().is_limited(0));
6761 }
6762
6763 #[tokio::test]
6764 async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
6765 let limited = ScriptedAdapter::new(
6766 0,
6767 AgentCapabilities::default(),
6768 [AgentEvent::Failed {
6769 slot: 0,
6770 started: true,
6771 detail: "request failed: insufficient_quota".into(),
6772 }],
6773 );
6774 let healthy = ScriptedAdapter::new(
6775 1,
6776 AgentCapabilities::default(),
6777 [AgentEvent::TurnComplete { slot: 1 }],
6778 );
6779 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6780 let captured = std::sync::Arc::clone(&events);
6781 let mut relay = RelayHost::new(
6782 vec![
6783 AdapterHost::new(Box::new(limited), None),
6784 AdapterHost::new(Box::new(healthy), None),
6785 ],
6786 4,
6787 )
6788 .expect("relay");
6789 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6790 relay.start().await.expect("start");
6791 events.lock().expect("events").clear();
6792
6793 assert!(matches!(
6794 relay.run_turn("task", 0).await.expect("limited failure"),
6795 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6796 ));
6797 assert!(relay.relay().is_limited(0));
6798 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
6799 {
6800 let events = events.lock().expect("events");
6801 assert!(
6802 events
6803 .iter()
6804 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6805 );
6806 assert!(
6807 !events
6808 .iter()
6809 .any(|event| matches!(event, AgentEvent::Failed { .. }))
6810 );
6811 }
6812
6813 assert!(matches!(
6814 relay.run_turn("", 0).await.expect("healthy peer"),
6815 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6816 ));
6817 }
6818
6819 #[tokio::test]
6820 async fn relay_failure_is_skipped_for_one_batch_without_changing_the_roster() {
6821 let failed = ScriptedAdapter::new(
6822 0,
6823 AgentCapabilities::default(),
6824 [AgentEvent::Failed {
6825 slot: 0,
6826 started: true,
6827 detail: "connection lost".into(),
6828 }],
6829 );
6830 let healthy = ScriptedAdapter::new(
6831 1,
6832 AgentCapabilities::default(),
6833 [AgentEvent::TurnComplete { slot: 1 }],
6834 );
6835 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6836 let captured = std::sync::Arc::clone(&events);
6837 let mut relay = RelayHost::new(
6838 vec![
6839 AdapterHost::new(Box::new(failed), None),
6840 AdapterHost::new(Box::new(healthy), None),
6841 ],
6842 4,
6843 )
6844 .expect("relay");
6845 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6846 relay.start().await.expect("start");
6847
6848 assert!(matches!(
6849 relay.run_turn("task", 0).await.expect("handled failure"),
6850 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6851 ));
6852 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![0, 1]);
6853 assert!(relay.relay().is_limited(0));
6854 assert!(events.lock().expect("lock").iter().any(|event| {
6855 matches!(
6856 event,
6857 AgentEvent::Failed {
6858 slot: 0,
6859 started: true,
6860 ..
6861 }
6862 )
6863 }));
6864 assert!(matches!(
6865 relay.run_turn("", 0).await.expect("healthy peer"),
6866 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6867 ));
6868 }
6869
6870 #[tokio::test]
6871 async fn codex_stop_does_not_skip_later_roster_reviewers() {
6872 let hosts = (0..3)
6873 .map(|slot| {
6874 AdapterHost::new(
6875 Box::new(ScriptedAdapter::new(
6876 slot,
6877 AgentCapabilities::default(),
6878 [
6879 AgentEvent::Text {
6880 slot,
6881 text: STOP_TOKEN.into(),
6882 },
6883 AgentEvent::TurnComplete { slot },
6884 ],
6885 )),
6886 None,
6887 )
6888 })
6889 .collect();
6890 let mut relay = RelayHost::new(hosts, 10).expect("relay");
6891 relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
6892 relay.start().await.expect("start");
6893 for expected in 0..3 {
6894 assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
6895 RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
6896 }
6897 assert_eq!(
6898 relay.run_turn("", 0).await.expect("complete"),
6899 RelayDecision::Complete
6900 );
6901 }
6902
6903 #[tokio::test]
6904 async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
6905 let first = ScriptedAdapter::new(
6906 0,
6907 AgentCapabilities::default(),
6908 [
6909 AgentEvent::Text {
6910 slot: 0,
6911 text: "done".into(),
6912 },
6913 AgentEvent::TurnComplete { slot: 0 },
6914 ],
6915 );
6916 let reviewer = ScriptedAdapter::new(
6917 1,
6918 AgentCapabilities::default(),
6919 [
6920 AgentEvent::Text {
6921 slot: 1,
6922 text: STOP_TOKEN.into(),
6923 },
6924 AgentEvent::TurnComplete { slot: 1 },
6925 ],
6926 );
6927 let mut relay = RelayHost::new(
6928 vec![
6929 AdapterHost::new(Box::new(first), None),
6930 AdapterHost::new(Box::new(reviewer), None),
6931 ],
6932 10,
6933 )
6934 .expect("relay");
6935 relay.start().await.expect("start");
6936 let first_decision = relay.run_turn("task", 0).await.expect("first");
6937 assert!(matches!(
6938 first_decision,
6939 RelayDecision::Dispatch { slot: 0, .. }
6940 ));
6941 let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
6942 assert!(matches!(
6943 reviewer_decision,
6944 RelayDecision::Dispatch {
6945 slot: 1,
6946 can_stop: true,
6947 ..
6948 }
6949 ));
6950 assert_eq!(
6951 relay.run_turn("", 0).await.expect("complete"),
6952 RelayDecision::Complete
6953 );
6954 }
6955
6956 #[tokio::test]
6957 async fn relay_stream_emits_text_and_thought_endings_before_tools() {
6958 let tool = AgentEvent::Tool {
6959 slot: 0,
6960 update: crate::ToolUpdate {
6961 id: "read".into(),
6962 title: "Read file".into(),
6963 status: ToolStatus::Running,
6964 detail: None,
6965 },
6966 };
6967 let updates = vec![
6968 AgentEvent::Thought {
6969 slot: 0,
6970 text: "Check the buffer. ✈".into(),
6971 },
6972 AgentEvent::Text {
6973 slot: 0,
6974 text: "Let me check.".into(),
6975 },
6976 tool.clone(),
6977 AgentEvent::Text {
6978 slot: 0,
6979 text: "[CODE".into(),
6980 },
6981 AgentEvent::Text {
6982 slot: 0,
6983 text: " is ordinary.".into(),
6984 },
6985 AgentEvent::Text {
6986 slot: 0,
6987 text: "[CODESWARM:".into(),
6988 },
6989 AgentEvent::Text {
6990 slot: 0,
6991 text: "STOP] Done.".into(),
6992 },
6993 tool,
6994 AgentEvent::TurnComplete { slot: 0 },
6995 ];
6996 let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
6997 let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
6998 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6999 let captured = std::sync::Arc::clone(&events);
7000 let mut relay = RelayHost::new(
7001 vec![
7002 AdapterHost::new(Box::new(first), None),
7003 AdapterHost::new(Box::new(reviewer), None),
7004 ],
7005 2,
7006 )
7007 .expect("relay");
7008 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
7009 relay.start().await.unwrap();
7010 relay.run_turn("task", 0).await.unwrap();
7011 let captured = events.lock().unwrap();
7012 let visible: Vec<_> = captured
7013 .iter()
7014 .filter(|event| {
7015 matches!(
7016 event,
7017 AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
7018 )
7019 })
7020 .cloned()
7021 .collect();
7022 assert_eq!(
7023 visible,
7024 vec![
7025 updates[0].clone(),
7026 updates[1].clone(),
7027 updates[2].clone(),
7028 AgentEvent::Text {
7029 slot: 0,
7030 text: "[CODE is ordinary.".into()
7031 },
7032 AgentEvent::Text {
7033 slot: 0,
7034 text: " Done.".into()
7035 },
7036 updates[7].clone(),
7037 ]
7038 );
7039 }
7040
7041 #[tokio::test]
7042 async fn roster_handoff_routes_only_terminal_message_markers_and_refreshes_targets() {
7043 let text = |value: &str| AgentEvent::Text {
7044 slot: 0,
7045 text: value.into(),
7046 };
7047 let thought = || AgentEvent::Thought {
7048 slot: 0,
7049 text: "still checking".into(),
7050 };
7051 let tool = || AgentEvent::Tool {
7052 slot: 0,
7053 update: crate::ToolUpdate {
7054 id: "read".into(),
7055 title: "Read file".into(),
7056 status: ToolStatus::Running,
7057 detail: None,
7058 },
7059 };
7060 let cases = vec![
7061 (vec![text("result [CODESWARM:NEXT:3]\n ")], 2),
7062 (
7063 vec![text("result [CODESWARM:"), text("NEXT:"), text("3]")],
7064 2,
7065 ),
7066 (vec![text("result [CODESWARM:NEXT:3]"), text(" more")], 1),
7067 (vec![text("result [CODESWARM:NEXT:3]"), thought()], 1),
7068 (
7069 vec![text("result [CODESWARM:NEXT:3]"), tool(), text(" ")],
7070 1,
7071 ),
7072 (
7073 vec![text("result [CODESWARM:NEXT:"), thought(), text("3]")],
7074 1,
7075 ),
7076 (
7077 vec![text("result"), thought(), text("[CODESWARM:NEXT:3]")],
7078 2,
7079 ),
7080 (vec![text("result [CODESWARM:NEXT:1]")], 1),
7081 (vec![text("result [CODESWARM:NEXT:0]")], 1),
7082 (vec![text("result [CODESWARM:NEXT:99]")], 1),
7083 (
7084 vec![
7085 text("result [CODESWARM:NEXT:3]"),
7086 AgentEvent::UsageUpdated {
7087 slot: 0,
7088 usage: crate::UsageUpdate { used: 1, size: 100 },
7089 },
7090 ],
7091 2,
7092 ),
7093 ];
7094 for (mut updates, expected) in cases {
7095 updates.push(AgentEvent::TurnComplete { slot: 0 });
7096 let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates);
7097 let hosts = std::iter::once(AdapterHost::new(Box::new(first), None))
7098 .chain((1..3).map(|slot| {
7099 AdapterHost::new(
7100 Box::new(ScriptedAdapter::new(
7101 slot,
7102 AgentCapabilities::default(),
7103 [AgentEvent::TurnComplete { slot }],
7104 )),
7105 None,
7106 )
7107 }))
7108 .collect();
7109 let mut relay = RelayHost::new(hosts, 10).unwrap();
7110 relay.set_roster_names(vec!["Worker".into(), "Codex".into(), "Codex".into()]);
7111 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
7112 let captured = events.clone();
7113 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
7114 relay.start().await.unwrap();
7115 relay.run_turn("task", 0).await.unwrap();
7116 let prompt = &relay.dispatches()[0].1;
7117 assert!(prompt.contains("[CODESWARM:NEXT:2] → Codex"));
7118 assert!(prompt.contains("[CODESWARM:NEXT:3] → Codex"));
7119 assert!(!prompt.contains("[CODESWARM:NEXT:1]"));
7120 relay.introduced.fill(true);
7122 relay.set_roster_names(vec!["Replacement".into(), "Codex".into(), "Codex".into()]);
7123 let next = relay.run_turn("", 0).await.unwrap();
7124 assert!(
7125 matches!(next, RelayDecision::Dispatch { slot, can_stop: false, .. } if slot == expected),
7126 "{next:?}"
7127 );
7128 let prompt = &relay.dispatches()[1].1;
7129 assert!(prompt.contains("[CODESWARM:NEXT:1] → Replacement"));
7130 assert!(prompt.contains("result"));
7131 let public = prompt
7132 .split("Public updates:\n")
7133 .nth(1)
7134 .unwrap()
7135 .split("\n\nDo not use")
7136 .next()
7137 .unwrap();
7138 assert!(!public.contains("[CODESWARM:NEXT:"));
7139 let visible = events
7140 .lock()
7141 .unwrap()
7142 .iter()
7143 .filter_map(|event| match event {
7144 AgentEvent::Text { text, .. } => Some(text.clone()),
7145 _ => None,
7146 })
7147 .collect::<String>();
7148 assert!(visible.contains("result"));
7149 assert!(!visible.contains("[CODESWARM:"), "{visible}");
7150 }
7151 }
7152
7153 #[tokio::test]
7154 async fn reviewer_stop_requires_a_terminal_marker_after_all_activity() {
7155 let text = |value: &str| AgentEvent::Text {
7156 slot: 1,
7157 text: value.into(),
7158 };
7159 let thought = || AgentEvent::Thought {
7160 slot: 1,
7161 text: "still checking".into(),
7162 };
7163 let tool = || AgentEvent::Tool {
7164 slot: 1,
7165 update: crate::ToolUpdate {
7166 id: "read".into(),
7167 title: "Read file".into(),
7168 status: ToolStatus::Running,
7169 detail: None,
7170 },
7171 };
7172 let cases = vec![
7173 (vec![text(&format!("done {STOP_TOKEN}"))], true),
7174 (vec![text(STOP_TOKEN), text("\n ")], true),
7175 (vec![text(STOP_TOKEN), text(" actually keep going")], false),
7176 (vec![text(STOP_TOKEN), thought()], false),
7177 (vec![text(STOP_TOKEN), tool()], false),
7178 (vec![text(STOP_TOKEN), tool(), text(" ")], false),
7179 (vec![text(STOP_TOKEN), tool(), text(STOP_TOKEN)], true),
7180 (vec![text("[CODESWARM:"), text("STOP]")], true),
7181 (vec![text("[CODESWARM:"), thought(), text("STOP]")], false),
7182 (
7183 vec![AgentEvent::Thought {
7184 slot: 1,
7185 text: STOP_TOKEN.into(),
7186 }],
7187 false,
7188 ),
7189 (
7190 vec![
7191 text(STOP_TOKEN),
7192 AgentEvent::UsageUpdated {
7193 slot: 1,
7194 usage: crate::UsageUpdate { used: 1, size: 100 },
7195 },
7196 ],
7197 true,
7198 ),
7199 ];
7200 for (mut events, stop) in cases {
7201 let first = ScriptedAdapter::new(
7202 0,
7203 AgentCapabilities::default(),
7204 [
7205 AgentEvent::Text {
7206 slot: 0,
7207 text: "initial response".into(),
7208 },
7209 AgentEvent::TurnComplete { slot: 0 },
7210 AgentEvent::TurnComplete { slot: 0 },
7211 ],
7212 );
7213 events.push(AgentEvent::TurnComplete { slot: 1 });
7214 let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), events.clone());
7215 let mut relay = RelayHost::new(
7216 vec![
7217 AdapterHost::new(Box::new(first), None),
7218 AdapterHost::new(Box::new(reviewer), None),
7219 ],
7220 4,
7221 )
7222 .unwrap();
7223 relay.start().await.unwrap();
7224 relay.run_turn("task", 0).await.unwrap();
7225 relay.run_turn("", 0).await.unwrap();
7226 let next = relay.run_turn("", 0).await.unwrap();
7227 assert_eq!(
7228 matches!(next, RelayDecision::Complete),
7229 stop,
7230 "events={events:?}"
7231 );
7232 relay.stop().await.unwrap();
7233 }
7234 }
7235
7236 #[tokio::test]
7237 async fn stop_token_is_filtered_from_streamed_ui_events() {
7238 let first = ScriptedAdapter::new(
7239 0,
7240 AgentCapabilities::default(),
7241 [
7242 AgentEvent::Text {
7243 slot: 0,
7244 text: format!("visible {STOP_TOKEN} trailing"),
7245 },
7246 AgentEvent::TurnComplete { slot: 0 },
7247 ],
7248 );
7249 let reviewer = ScriptedAdapter::new(
7250 1,
7251 AgentCapabilities::default(),
7252 [AgentEvent::TurnComplete { slot: 1 }],
7253 );
7254 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7255 let captured = std::sync::Arc::clone(&events);
7256 let mut relay = RelayHost::new(
7257 vec![
7258 AdapterHost::new(Box::new(first), None),
7259 AdapterHost::new(Box::new(reviewer), None),
7260 ],
7261 2,
7262 )
7263 .expect("relay");
7264 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7265 relay.start().await.expect("start");
7266 relay.run_turn("task", 0).await.expect("turn");
7267 let captured = events.lock().expect("lock");
7268 assert!(captured.iter().all(|event| match event {
7269 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7270 _ => true,
7271 }));
7272 let visible = captured
7273 .iter()
7274 .filter_map(|event| match event {
7275 AgentEvent::Text { text, .. } => Some(text.as_str()),
7276 _ => None,
7277 })
7278 .collect::<String>();
7279 assert_eq!(visible, "visible trailing");
7280 }
7281
7282 #[tokio::test]
7283 async fn token_only_reviewer_response_emits_visible_acknowledgment() {
7284 let first = ScriptedAdapter::new(
7285 0,
7286 AgentCapabilities::default(),
7287 [
7288 AgentEvent::Text {
7289 slot: 0,
7290 text: "done".into(),
7291 },
7292 AgentEvent::TurnComplete { slot: 0 },
7293 ],
7294 );
7295 let reviewer = ScriptedAdapter::new(
7296 1,
7297 AgentCapabilities::default(),
7298 [
7299 AgentEvent::Text {
7300 slot: 1,
7301 text: STOP_TOKEN.into(),
7302 },
7303 AgentEvent::TurnComplete { slot: 1 },
7304 ],
7305 );
7306 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7307 let captured = std::sync::Arc::clone(&events);
7308 let mut relay = RelayHost::new(
7309 vec![
7310 AdapterHost::new(Box::new(first), None),
7311 AdapterHost::new(Box::new(reviewer), None),
7312 ],
7313 4,
7314 )
7315 .expect("relay");
7316 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7317 relay.start().await.expect("start");
7318 relay.run_turn("task", 0).await.expect("first turn");
7319 relay.run_turn("", 0).await.expect("review turn");
7320 let captured = events.lock().expect("lock");
7321 assert!(captured.iter().any(|event| {
7322 matches!(
7323 event,
7324 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7325 )
7326 }));
7327 assert!(captured.iter().all(|event| match event {
7328 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7329 _ => true,
7330 }));
7331 let acknowledgment = captured
7332 .iter()
7333 .position(|event| {
7334 matches!(
7335 event,
7336 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7337 )
7338 })
7339 .expect("visible acknowledgment");
7340 let completion = captured
7341 .iter()
7342 .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
7343 .expect("reviewer completion");
7344 assert!(acknowledgment < completion);
7345 }
7346
7347 #[tokio::test]
7348 async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
7349 let first = ScriptedAdapter::new(
7350 0,
7351 AgentCapabilities::default(),
7352 [
7353 AgentEvent::Text {
7354 slot: 0,
7355 text: "done".into(),
7356 },
7357 AgentEvent::TurnComplete { slot: 0 },
7358 ],
7359 );
7360 let reviewer = ScriptedAdapter::new(
7361 1,
7362 AgentCapabilities::default(),
7363 [
7364 AgentEvent::Text {
7365 slot: 1,
7366 text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
7367 },
7368 AgentEvent::TurnComplete { slot: 1 },
7369 ],
7370 );
7371 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7372 let captured = std::sync::Arc::clone(&events);
7373 let mut relay = RelayHost::new(
7374 vec![
7375 AdapterHost::new(Box::new(first), None),
7376 AdapterHost::new(Box::new(reviewer), None),
7377 ],
7378 4,
7379 )
7380 .expect("relay");
7381 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7382 relay.start().await.expect("start");
7383 relay.run_turn("task", 0).await.expect("first turn");
7384 relay.run_turn("", 0).await.expect("review turn");
7385
7386 let visible = events
7387 .lock()
7388 .expect("lock")
7389 .iter()
7390 .filter_map(|event| match event {
7391 AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
7392 _ => None,
7393 })
7394 .collect::<String>();
7395 assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
7396 assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
7397 }
7398
7399 #[tokio::test]
7400 async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
7401 let first = AdapterHost::new(
7402 Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
7403 None,
7404 );
7405 let second = AdapterHost::new(
7406 Box::new(ScriptedAdapter::new(
7407 1,
7408 AgentCapabilities::default(),
7409 [AgentEvent::TurnComplete { slot: 1 }],
7410 )),
7411 None,
7412 );
7413 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7414 let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
7415 relay.set_event_sink(move |event| {
7416 if matches!(event, AgentEvent::Permission { .. }) {
7417 let _ = seen_sender.send(());
7418 }
7419 });
7420 relay.start().await.expect("start");
7421 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
7422 let answer = async move {
7423 seen_receiver.recv().await.expect("permission request");
7424 sender
7425 .send(super::RelayPermissionAnswer {
7426 slot: 0,
7427 request_id: "permission-1".into(),
7428 answer: PermissionAnswer::Selected {
7429 option_id: "allow".into(),
7430 },
7431 })
7432 .expect("queue permission answer");
7433 };
7434 tokio::time::timeout(std::time::Duration::from_millis(100), async {
7435 let ((), result) = tokio::join!(
7436 answer,
7437 relay.run_turn_with_permissions("task", 0, &mut receiver)
7438 );
7439 result
7440 })
7441 .await
7442 .expect("permission-gated turn should not deadlock")
7443 .expect("turn completes");
7444 }
7445
7446 #[tokio::test]
7447 async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
7448 let first = AdapterHost::new(
7449 Box::new(PendingAdapter {
7450 slot: 0,
7451 hang_on_cancel: false,
7452 }),
7453 None,
7454 );
7455 let second = AdapterHost::new(
7456 Box::new(ScriptedAdapter::new(
7457 1,
7458 AgentCapabilities::default(),
7459 [AgentEvent::TurnComplete { slot: 1 }],
7460 )),
7461 None,
7462 );
7463 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7464 relay.start().await.expect("start");
7465 let cancellation = relay.cancellation();
7466 let error = {
7467 let turn = relay.run_turn("task", 0);
7468 tokio::pin!(turn);
7469 cancellation.request();
7470 turn.await.expect_err("cancellation should stop turn")
7471 };
7472 assert!(error.to_string().contains("relay turn cancelled"));
7473
7474 assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
7475 relay
7476 .run_turn("", 1)
7477 .await
7478 .expect("replacement job reaches the selected peer");
7479 let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
7480 assert!(replacement.contains("replacement job"));
7481 assert!(replacement.contains("User "));
7482 assert!(replacement.contains(":\ntask"));
7483 let owner_updates = relay.relay_mut().unseen_context(0);
7484 assert!(owner_updates.contains("User "));
7485 assert!(owner_updates.contains(":\ntask"));
7486 assert!(owner_updates.contains(":\nreplacement job"));
7487 }
7488
7489 #[tokio::test]
7490 async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
7491 let first = AdapterHost::new(
7492 Box::new(PendingAdapter {
7493 slot: 0,
7494 hang_on_cancel: true,
7495 }),
7496 None,
7497 );
7498 let second = AdapterHost::new(
7499 Box::new(PendingAdapter {
7500 slot: 1,
7501 hang_on_cancel: false,
7502 }),
7503 None,
7504 );
7505 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7506 relay.start().await.expect("start");
7507 let cancellation = relay.cancellation();
7508 let turn = relay.run_turn("task", 0);
7509 tokio::pin!(turn);
7510 cancellation.request();
7511 let error = turn.await.expect_err("cancellation should stop turn");
7512 assert!(error.to_string().contains("timed out"));
7513 }
7514
7515 #[tokio::test]
7516 async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
7517 let event = [AgentEvent::TurnComplete { slot: 0 }];
7518 let first = AdapterHost::new(
7519 Box::new(ScriptedAdapter::new(
7520 0,
7521 AgentCapabilities::default(),
7522 event.clone(),
7523 )),
7524 None,
7525 );
7526 let second = AdapterHost::new(
7527 Box::new(ScriptedAdapter::new(
7528 1,
7529 AgentCapabilities::default(),
7530 [AgentEvent::TurnComplete { slot: 1 }],
7531 )),
7532 None,
7533 );
7534 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7535 relay.start().await.expect("start");
7536
7537 relay.pause();
7538 assert_eq!(
7539 relay.run_turn("paused", 0).await.expect("paused turn"),
7540 crate::relay::RelayDecision::Paused
7541 );
7542 assert!(relay.dispatches().is_empty());
7543
7544 relay.resume();
7545 relay.relay_mut().drop_agent(1).expect("drop reviewer");
7546 assert!(matches!(
7547 relay
7548 .run_turn("solo follow-up", 0)
7549 .await
7550 .expect("solo turn"),
7551 crate::relay::RelayDecision::Dispatch {
7552 slot: 0,
7553 can_stop: false,
7554 ..
7555 }
7556 ));
7557 assert_eq!(relay.dispatches().len(), 1);
7558 }
7559
7560 #[tokio::test]
7561 async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
7562 let first = AdapterHost::new(
7563 Box::new(ScriptedAdapter::new(
7564 0,
7565 AgentCapabilities::default(),
7566 [AgentEvent::TurnComplete { slot: 0 }],
7567 )),
7568 None,
7569 );
7570 let second = AdapterHost::new(
7571 Box::new(ScriptedAdapter::new(
7572 1,
7573 AgentCapabilities::default(),
7574 [AgentEvent::TurnComplete { slot: 1 }],
7575 )),
7576 None,
7577 );
7578 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7579 relay.set_roster_names(vec!["First".into(), "Second".into()]);
7580 relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
7581 relay.set_roster_launch_specs(vec![
7582 ("custom".into(), "owner".into()),
7583 ("custom".into(), "peer".into()),
7584 ]);
7585 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7586 let captured = std::sync::Arc::clone(&events);
7587 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7588 relay.start().await.expect("start");
7589 let slot = relay
7590 .add_agent(
7591 AdapterHost::new(
7592 Box::new(ScriptedAdapter::new(
7593 2,
7594 AgentCapabilities::default(),
7595 [AgentEvent::TurnComplete { slot: 2 }],
7596 )),
7597 None,
7598 ),
7599 "Reviewer",
7600 "reviewer.example",
7601 "reviewer --acp",
7602 )
7603 .await
7604 .expect("append agent");
7605 assert_eq!(slot, 2);
7606 assert_eq!(
7607 relay.relay().active_slots().collect::<Vec<_>>(),
7608 vec![0, 1, 2]
7609 );
7610 assert_eq!(
7611 relay
7612 .session_metadata()
7613 .get("agents")
7614 .and_then(|value| value.as_array())
7615 .map(Vec::len),
7616 Some(3)
7617 );
7618 relay.drop_agent(1).await.expect("drop middle peer");
7619 let metadata = relay.session_metadata();
7620 assert_eq!(
7621 metadata.get("agents"),
7622 Some(&serde_json::json!([
7623 {"slot": 0, "name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
7624 {"slot": 2, "name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
7625 ]))
7626 );
7627 assert!(
7628 events
7629 .lock()
7630 .expect("lock")
7631 .iter()
7632 .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
7633 );
7634 }
7635
7636 #[tokio::test]
7637 async fn relay_host_persists_coordinator_owned_runtime_metadata() {
7638 let path = unique_test_path("codeswarm-session-metadata", "json");
7639 let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
7640 let writer = metadata_store.buffered().expect("metadata writer");
7641 let first = AdapterHost::new(
7642 Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
7643 None,
7644 );
7645 let second = AdapterHost::new(
7646 Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
7647 None,
7648 );
7649 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7650 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7651 relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
7652 relay.set_roster_launch_specs(vec![
7653 ("custom".into(), "claude".into()),
7654 ("custom".into(), "codex".into()),
7655 ]);
7656 relay.set_session_metadata_writer(writer);
7657 relay.start().await.expect("start");
7658 relay.drop_agent(0).await.expect("drop first agent");
7659 relay.stop().await.expect("stop");
7660
7661 let loaded = metadata_store
7662 .read()
7663 .expect("read metadata")
7664 .expect("metadata snapshot");
7665 assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
7666 assert_eq!(
7667 loaded.get("agents"),
7668 Some(&serde_json::json!([{
7669 "slot": 1, "name": "Codex", "identity": "openai.com", "protocol": "custom",
7670 "command": "codex", "supports_load_session": false
7671 }]))
7672 );
7673 assert!(loaded.get("owner").is_none());
7674 let _ = std::fs::remove_file(path);
7675 }
7676
7677 #[tokio::test]
7678 async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
7679 let first = AdapterHost::new(
7680 Box::new(ScriptedAdapter::new(
7681 0,
7682 AgentCapabilities::default(),
7683 [
7684 AgentEvent::Text {
7685 slot: 0,
7686 text: "owner stream".into(),
7687 },
7688 AgentEvent::TurnComplete { slot: 0 },
7689 ],
7690 )),
7691 None,
7692 );
7693 let second = AdapterHost::new(
7694 Box::new(ScriptedAdapter::new(
7695 1,
7696 AgentCapabilities::default(),
7697 [
7698 AgentEvent::Text {
7699 slot: 1,
7700 text: "peer stream".into(),
7701 },
7702 AgentEvent::TurnComplete { slot: 1 },
7703 ],
7704 )),
7705 None,
7706 );
7707 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7708 relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
7709 relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
7710 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7711 let captured = std::sync::Arc::clone(&events);
7712 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7713 relay.start().await.expect("start");
7714
7715 relay.swap_agents(0, 1).expect("swap peers");
7716 assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
7717 assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
7718 relay.run_turn("task", 0).await.expect("swapped turn");
7719 let events = events.lock().expect("events");
7720 assert!(events.iter().any(|event| {
7721 matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
7722 }));
7723 assert!(relay.dispatches()[0].1.contains("You are Peer"));
7724 }
7725
7726 #[tokio::test]
7727 async fn relay_host_persists_all_active_agent_metadata_off_thread() {
7728 let path = unique_test_path("codeswarm-session-metadata", "json");
7729 let first = AdapterHost::new(
7730 Box::new(ScriptedAdapter::new(
7731 0,
7732 AgentCapabilities::default(),
7733 [AgentEvent::TurnComplete { slot: 0 }],
7734 )),
7735 None,
7736 );
7737 let second = AdapterHost::new(
7738 Box::new(ScriptedAdapter::new(
7739 1,
7740 AgentCapabilities::default(),
7741 [AgentEvent::TurnComplete { slot: 1 }],
7742 )),
7743 None,
7744 );
7745 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7746 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7747 relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
7748 relay.set_roster_launch_specs(vec![
7749 ("custom".into(), "claude".into()),
7750 ("custom".into(), "codex".into()),
7751 ]);
7752 let writer = SessionMetadataStore::open(&path)
7753 .buffered()
7754 .expect("metadata writer");
7755 relay.set_session_metadata_writer(writer);
7756 relay.start().await.expect("start");
7757 relay.stop().await.expect("stop");
7758 let loaded = SessionMetadataStore::open(&path)
7759 .read()
7760 .expect("read metadata")
7761 .expect("metadata snapshot");
7762 let agents = loaded
7763 .get("agents")
7764 .and_then(|value| value.as_array())
7765 .expect("agents");
7766 assert_eq!(agents.len(), 2);
7767 assert_eq!(agents[0]["identity"], "claude.com");
7768 assert_eq!(agents[1]["identity"], "openai.com");
7769 let _ = std::fs::remove_file(path);
7770 }
7771
7772 #[tokio::test]
7773 async fn relay_host_routes_unseen_public_context_to_next_agent() {
7774 let first = AdapterHost::new(
7775 Box::new(ScriptedAdapter::new(
7776 0,
7777 AgentCapabilities::default(),
7778 [
7779 AgentEvent::Text {
7780 slot: 0,
7781 text: "implemented the fix".into(),
7782 },
7783 AgentEvent::TurnComplete { slot: 0 },
7784 ],
7785 )),
7786 None,
7787 );
7788 let second = AdapterHost::new(
7789 Box::new(ScriptedAdapter::new(
7790 1,
7791 AgentCapabilities::default(),
7792 [AgentEvent::TurnComplete { slot: 1 }],
7793 )),
7794 None,
7795 );
7796 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7797 relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
7798 relay.start().await.expect("start");
7799 relay.run_turn("task", 0).await.expect("first turn");
7800 relay.run_turn("review this", 0).await.expect("review turn");
7801
7802 assert_eq!(relay.dispatches().len(), 2);
7803 assert_eq!(relay.dispatches()[0].0, 0);
7804 assert!(relay.dispatches()[0].1.contains("task"));
7805 assert!(relay.dispatches()[0].1.contains("You are Codex"));
7806 assert!(relay.dispatches()[0].1.contains("2. Qwen"));
7807 assert_eq!(relay.dispatches()[1].0, 1);
7808 assert!(relay.dispatches()[1].1.contains("review this"));
7809 let public = relay.dispatches()[1]
7810 .1
7811 .split_once("Public updates:\n")
7812 .map(|(_, updates)| updates)
7813 .expect("review receives public context");
7814 let header = public
7815 .lines()
7816 .find(|line| line.starts_with("Codex "))
7817 .expect("named previous agent");
7818 let timestamp = header
7819 .strip_prefix("Codex ")
7820 .and_then(|value| value.strip_suffix(':'))
7821 .expect("timestamped header");
7822 assert_eq!(timestamp.len(), 5);
7823 assert_eq!(timestamp.as_bytes()[2], b':');
7824 assert!(
7825 timestamp
7826 .bytes()
7827 .enumerate()
7828 .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
7829 );
7830 assert!(public.contains("implemented the fix"));
7831 assert!(!public.contains("Agent 0"));
7832 }
7833}