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 queued_events: VecDeque<AdapterResult<AgentEvent>>,
2508 tool_updates: BTreeMap<String, ToolUpdate>,
2509 stderr_task: Option<tokio::task::JoinHandle<String>>,
2510 terminals: BTreeMap<String, TerminalProcess>,
2511 next_terminal_id: u64,
2512}
2513
2514impl AcpAdapter {
2515 pub fn new(
2516 slot: RosterSlot,
2517 cwd: PathBuf,
2518 program: impl Into<String>,
2519 args: Vec<String>,
2520 ) -> Self {
2521 Self {
2522 slot,
2523 program: program.into(),
2524 args,
2525 cwd,
2526 child: None,
2527 reader: None,
2528 capabilities: AgentCapabilities::default(),
2529 modes: Vec::new(),
2530 models: Vec::new(),
2531 model_config_id: None,
2532 session_id: None,
2533 next_request_id: 1,
2534 prompt_request_id: None,
2535 queued_events: VecDeque::new(),
2536 tool_updates: BTreeMap::new(),
2537 stderr_task: None,
2538 terminals: BTreeMap::new(),
2539 next_terminal_id: 1,
2540 }
2541 }
2542
2543 pub fn with_session_id(
2544 slot: RosterSlot,
2545 cwd: PathBuf,
2546 program: impl Into<String>,
2547 args: Vec<String>,
2548 session_id: impl Into<String>,
2549 ) -> Self {
2550 let mut adapter = Self::new(slot, cwd, program, args);
2551 adapter.session_id = Some(session_id.into());
2552 adapter
2553 }
2554
2555 async fn request(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2556 self.request_with_timeout(method, params, std::time::Duration::from_secs(30))
2557 .await
2558 }
2559
2560 async fn request_with_timeout(
2561 &mut self,
2562 method: &str,
2563 params: Value,
2564 deadline: std::time::Duration,
2565 ) -> AdapterResult<Value> {
2566 tokio::time::timeout(deadline, self.request_inner(method, params))
2567 .await
2568 .map_err(|_| {
2569 AdapterError::Transport(format!(
2570 "ACP {method} timed out; reload the agent to retry"
2571 ))
2572 })?
2573 }
2574
2575 async fn request_inner(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2576 let request_id = self.next_request_id;
2577 self.next_request_id += 1;
2578 self.write_json(serde_json::json!({
2579 "jsonrpc": "2.0",
2580 "id": request_id,
2581 "method": method,
2582 "params": params,
2583 }))
2584 .await?;
2585 loop {
2586 let line = self.read_line().await?;
2587 let value: Value = match serde_json::from_str(&line) {
2588 Ok(value) => value,
2589 Err(_) => {
2590 continue;
2594 }
2595 };
2596 if self.reject_empty_permission_request(&value).await? {
2597 continue;
2598 }
2599 if self.handle_client_request(&value).await? {
2600 continue;
2601 }
2602 if value
2603 .get("id")
2604 .is_some_and(|id| rpc_id_to_string(id) == request_id.to_string())
2605 {
2606 if let Some(error) = value.get("error") {
2607 return Err(AdapterError::Protocol(error.to_string()));
2608 }
2609 return value
2610 .get("result")
2611 .cloned()
2612 .ok_or_else(|| AdapterError::Protocol("response has no result".into()));
2613 }
2614 if let Some(event) = parse_acp_value(self.slot, &value, &mut self.tool_updates)? {
2615 let event = if method == "session/load" {
2616 restored_history_event(event)
2617 } else {
2618 event
2619 };
2620 self.queued_events.push_back(Ok(event));
2621 }
2622 }
2623 }
2624
2625 async fn write_json(&mut self, value: Value) -> AdapterResult<()> {
2626 let child = self
2627 .child
2628 .as_mut()
2629 .ok_or_else(|| AdapterError::Transport("ACP agent is not running".into()))?;
2630 let stdin = child
2631 .stdin
2632 .as_mut()
2633 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdin".into()))?;
2634 stdin
2635 .write_all(value.to_string().as_bytes())
2636 .await
2637 .map_err(|error| AdapterError::Transport(error.to_string()))?;
2638 stdin
2639 .write_all(b"\n")
2640 .await
2641 .map_err(|error| AdapterError::Transport(error.to_string()))
2642 }
2643
2644 async fn reset_transport(&mut self) {
2647 let terminals = std::mem::take(&mut self.terminals);
2648 for terminal in terminals.values() {
2649 terminal.stop().await;
2650 }
2651 self.queued_events.clear();
2652 self.tool_updates.clear();
2653 if let Some(mut child) = self.child.take() {
2654 let _ = terminate_child(&mut child).await;
2655 }
2656 self.reader = None;
2657 if let Some(task) = self.stderr_task.take() {
2658 task.abort();
2659 }
2660 self.prompt_request_id = None;
2661 }
2662
2663 async fn reject_empty_permission_request(&mut self, value: &Value) -> AdapterResult<bool> {
2668 if value.get("method").and_then(Value::as_str) != Some("session/request_permission")
2669 || value.get("id").is_none()
2670 {
2671 return Ok(false);
2672 }
2673 let valid = value
2674 .get("params")
2675 .and_then(|params| params.get("options"))
2676 .and_then(Value::as_array)
2677 .is_some_and(|options| !options.is_empty());
2678 if valid {
2679 return Ok(false);
2680 }
2681 self.write_json(serde_json::json!({
2682 "jsonrpc": "2.0",
2683 "id": value.get("id").cloned().unwrap_or(Value::Null),
2684 "error": {
2685 "code": -32602,
2686 "message": "Permission request requires at least one option",
2687 },
2688 }))
2689 .await?;
2690 Ok(true)
2691 }
2692
2693 fn workspace_path(&self, path: &str) -> Result<PathBuf, String> {
2694 let root = self
2695 .cwd
2696 .canonicalize()
2697 .map_err(|error| format!("unable to resolve workspace: {error}"))?;
2698 let requested = Path::new(path);
2699 let candidate = if requested.is_absolute() {
2700 requested.to_path_buf()
2701 } else {
2702 root.join(requested)
2703 };
2704 let resolved = if !candidate.exists() {
2705 let parent = candidate
2706 .parent()
2707 .ok_or_else(|| "file path has no parent".to_owned())?
2708 .canonicalize()
2709 .map_err(|error| format!("unable to resolve parent directory: {error}"))?;
2710 parent.join(
2711 candidate
2712 .file_name()
2713 .ok_or_else(|| "file path has no filename".to_owned())?,
2714 )
2715 } else {
2716 candidate
2717 .canonicalize()
2718 .map_err(|error| format!("unable to resolve file path: {error}"))?
2719 };
2720 if !resolved.starts_with(&root) {
2721 return Err("file path is outside the project".into());
2722 }
2723 Ok(resolved)
2724 }
2725
2726 fn read_workspace_text(
2727 &self,
2728 path: &str,
2729 line: Option<i64>,
2730 limit: Option<i64>,
2731 ) -> Result<String, String> {
2732 if line.is_some_and(|line| line < 1) {
2733 return Err("line must be positive".into());
2734 }
2735 if limit.is_some_and(|limit| limit < 0) {
2736 return Err("limit must not be negative".into());
2737 }
2738 let path = self.workspace_path(path)?;
2739 let mut bytes = Vec::new();
2740 let mut source = match std::fs::File::open(path) {
2741 Ok(source) => source,
2742 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(String::new()),
2743 Err(error) => return Err(error.to_string()),
2744 };
2745 source
2746 .by_ref()
2747 .take((MAX_FILE_READ_BYTES as u64).saturating_add(1))
2748 .read_to_end(&mut bytes)
2749 .map_err(|error| error.to_string())?;
2750 bytes.truncate(MAX_FILE_READ_BYTES);
2751 let text = String::from_utf8_lossy(&bytes);
2752 if line.is_none() && limit.is_none() {
2753 return Ok(text.into_owned());
2754 }
2755 let start = line.map_or(0, |line| line as usize - 1);
2756 let limit = limit.unwrap_or(i64::MAX) as usize;
2757 let selected = text
2758 .split_inclusive('\n')
2759 .skip(start)
2760 .take(limit)
2761 .collect::<String>();
2762 if line.is_some() {
2763 Ok(selected.trim_end_matches('\n').to_owned())
2764 } else {
2765 Ok(selected)
2766 }
2767 }
2768
2769 fn write_workspace_text(&self, params: &Value) -> Result<(), String> {
2770 let path = params
2771 .get("path")
2772 .and_then(Value::as_str)
2773 .filter(|path| !path.is_empty())
2774 .ok_or("path must be a non-empty string")?;
2775 let content = params
2776 .get("content")
2777 .and_then(Value::as_str)
2778 .ok_or("content must be a string")?;
2779 let path = self.workspace_path(path)?;
2780 std::fs::write(path, content).map_err(|error| error.to_string())
2781 }
2782
2783 async fn terminal_create(&mut self, params: &Value) -> Result<Value, String> {
2784 let command = params
2785 .get("command")
2786 .and_then(Value::as_str)
2787 .filter(|command| !command.trim().is_empty())
2788 .ok_or_else(|| "terminal command is required".to_owned())?;
2789 let cwd = params.get("cwd").and_then(Value::as_str).unwrap_or(".");
2790 let cwd = self.workspace_path(cwd)?;
2791 if !cwd.is_dir() {
2792 return Err("terminal cwd is not a directory".into());
2793 }
2794 let mut process = Command::new(command);
2795 isolate_process_group(&mut process);
2796 if let Some(args) = params.get("args").and_then(Value::as_array) {
2797 process.args(args.iter().filter_map(Value::as_str));
2798 }
2799 process
2800 .current_dir(&cwd)
2801 .stdin(Stdio::null())
2802 .stdout(Stdio::piped())
2803 .stderr(Stdio::piped());
2804 if let Some(env) = params.get("env") {
2805 if let Some(entries) = env.as_array() {
2806 for entry in entries {
2807 if let (Some(name), Some(value)) = (
2808 entry.get("name").and_then(Value::as_str),
2809 entry.get("value").and_then(Value::as_str),
2810 ) {
2811 process.env(name, value);
2812 }
2813 }
2814 } else if let Some(entries) = env.as_object() {
2815 for (name, value) in entries {
2816 if let Some(value) = value.as_str() {
2817 process.env(name, value);
2818 }
2819 }
2820 }
2821 }
2822 let mut child = process.spawn().map_err(|error| error.to_string())?;
2823 let stdout = child.stdout.take();
2824 let stderr = child.stderr.take();
2825 let (Some(stdout), Some(stderr)) = (stdout, stderr) else {
2826 let _ = terminate_child(&mut child).await;
2827 return Err("terminal has no output pipes".into());
2828 };
2829 let output = Arc::new(Mutex::new(Vec::new()));
2830 let truncated = Arc::new(AtomicBool::new(false));
2831 let output_readers = Arc::new(AtomicUsize::new(2));
2832 let output_limit = params
2833 .get("outputByteLimit")
2834 .and_then(Value::as_u64)
2835 .map_or(MAX_TERMINAL_OUTPUT_BYTES, |limit| {
2836 usize::try_from(limit)
2837 .unwrap_or(MAX_TERMINAL_OUTPUT_BYTES)
2838 .min(MAX_TERMINAL_OUTPUT_BYTES)
2839 });
2840 tokio::spawn(drain_terminal_output(
2841 stdout,
2842 Arc::clone(&output),
2843 Arc::clone(&truncated),
2844 Arc::clone(&output_readers),
2845 output_limit,
2846 ));
2847 tokio::spawn(drain_terminal_output(
2848 stderr,
2849 Arc::clone(&output),
2850 Arc::clone(&truncated),
2851 Arc::clone(&output_readers),
2852 output_limit,
2853 ));
2854 let id = format!("terminal-{}", self.next_terminal_id);
2855 self.next_terminal_id = self.next_terminal_id.saturating_add(1);
2856 let state = TerminalProcess {
2857 child: Arc::new(AsyncMutex::new(Some(child))),
2858 output,
2859 truncated,
2860 output_readers,
2861 };
2862 self.terminals.insert(id.clone(), state);
2863 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2864 slot: self.slot,
2865 event: TerminalEvent::Created {
2866 id: id.clone(),
2867 command: std::iter::once(command)
2868 .chain(
2869 params
2870 .get("args")
2871 .and_then(Value::as_array)
2872 .into_iter()
2873 .flatten()
2874 .filter_map(Value::as_str),
2875 )
2876 .collect::<Vec<_>>()
2877 .join(" "),
2878 },
2879 }));
2880 Ok(serde_json::json!({"terminalId": id}))
2881 }
2882
2883 async fn terminal_output(&mut self, id: &str) -> Result<Value, String> {
2884 let terminal = self
2885 .terminals
2886 .get(id)
2887 .ok_or_else(|| "terminal not found".to_owned())?
2888 .clone();
2889 let output = terminal
2890 .output
2891 .lock()
2892 .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
2893 .unwrap_or_default();
2894 let exit_code = terminal.exit_code().await;
2895 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2896 slot: self.slot,
2897 event: TerminalEvent::Output {
2898 id: id.to_owned(),
2899 text: output.clone(),
2900 },
2901 }));
2902 let mut response = serde_json::json!({
2903 "output": output,
2904 "truncated": terminal.truncated.load(Ordering::Acquire),
2905 });
2906 if let Some(code) = exit_code {
2907 response["exitStatus"] = serde_json::json!({"exitCode": code});
2908 }
2909 Ok(response)
2910 }
2911
2912 async fn terminal_wait(&mut self, id: &str) -> Result<Value, String> {
2913 let terminal = self
2914 .terminals
2915 .get(id)
2916 .ok_or_else(|| "terminal not found".to_owned())?
2917 .clone();
2918 let exit_code = terminal.wait().await;
2919 self.queued_events.push_back(Ok(AgentEvent::Terminal {
2920 slot: self.slot,
2921 event: TerminalEvent::Exited {
2922 id: id.to_owned(),
2923 code: exit_code.unwrap_or(-1),
2924 },
2925 }));
2926 Ok(serde_json::json!({"exitCode": exit_code, "signal": Value::Null}))
2927 }
2928
2929 async fn handle_client_request(&mut self, value: &Value) -> AdapterResult<bool> {
2933 let Some(method) = value.get("method").and_then(Value::as_str) else {
2934 return Ok(false);
2935 };
2936 let Some(id) = value.get("id").cloned() else {
2937 return Ok(false);
2938 };
2939 if method == "session/request_permission" {
2942 return Ok(false);
2943 }
2944 let params = value.get("params").cloned().unwrap_or(Value::Null);
2945 let response = match method {
2946 "fs/read_text_file" => {
2947 let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2948 let line = params.get("line").and_then(Value::as_i64);
2949 let limit = params.get("limit").and_then(Value::as_i64);
2950 match self.read_workspace_text(path, line, limit) {
2951 Ok(content) => serde_json::json!({
2952 "jsonrpc": "2.0",
2953 "id": id,
2954 "result": {"content": content},
2955 }),
2956 Err(message) => serde_json::json!({
2957 "jsonrpc": "2.0",
2958 "id": id,
2959 "error": {"code": -32602, "message": message},
2960 }),
2961 }
2962 }
2963 "fs/write_text_file" => {
2964 let result = self.write_workspace_text(¶ms);
2965 match result {
2966 Ok(()) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}}),
2967 Err(message) => serde_json::json!({
2968 "jsonrpc": "2.0",
2969 "id": id,
2970 "error": {"code": -32602, "message": message},
2971 }),
2972 }
2973 }
2974 "terminal/create" => match self.terminal_create(¶ms).await {
2975 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2976 Err(message) => serde_json::json!({
2977 "jsonrpc": "2.0",
2978 "id": id,
2979 "error": {"code": -32602, "message": message},
2980 }),
2981 },
2982 "terminal/output" => {
2983 let terminal_id = params
2984 .get("terminalId")
2985 .and_then(Value::as_str)
2986 .unwrap_or("");
2987 match self.terminal_output(terminal_id).await {
2988 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2989 Err(message) => serde_json::json!({
2990 "jsonrpc": "2.0",
2991 "id": id,
2992 "error": {"code": -32602, "message": message},
2993 }),
2994 }
2995 }
2996 "terminal/wait_for_exit" => {
2997 let terminal_id = params
2998 .get("terminalId")
2999 .and_then(Value::as_str)
3000 .unwrap_or("");
3001 match self.terminal_wait(terminal_id).await {
3002 Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
3003 Err(message) => serde_json::json!({
3004 "jsonrpc": "2.0",
3005 "id": id,
3006 "error": {"code": -32602, "message": message},
3007 }),
3008 }
3009 }
3010 "terminal/kill" => {
3011 let terminal_id = params
3012 .get("terminalId")
3013 .and_then(Value::as_str)
3014 .unwrap_or("");
3015 if let Some(terminal) = self.terminals.get(terminal_id) {
3016 terminal.kill().await;
3017 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
3018 } else {
3019 serde_json::json!({
3020 "jsonrpc": "2.0",
3021 "id": id,
3022 "error": {"code": -32602, "message": "terminal not found"},
3023 })
3024 }
3025 }
3026 "terminal/release" => {
3027 let terminal_id = params
3028 .get("terminalId")
3029 .and_then(Value::as_str)
3030 .unwrap_or("");
3031 if let Some(terminal) = self.terminals.remove(terminal_id) {
3032 terminal.stop().await;
3033 self.queued_events.push_back(Ok(AgentEvent::Terminal {
3034 slot: self.slot,
3035 event: TerminalEvent::Released {
3036 id: terminal_id.to_owned(),
3037 },
3038 }));
3039 serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
3040 } else {
3041 serde_json::json!({
3042 "jsonrpc": "2.0",
3043 "id": id,
3044 "error": {"code": -32602, "message": "terminal not found"},
3045 })
3046 }
3047 }
3048 _ => serde_json::json!({
3049 "jsonrpc": "2.0",
3050 "id": id,
3051 "error": {"code": -32601, "message": format!("unsupported client method: {method}")},
3052 }),
3053 };
3054 self.write_json(response).await?;
3055 Ok(true)
3056 }
3057
3058 async fn read_line(&mut self) -> AdapterResult<String> {
3059 let reader = self
3060 .reader
3061 .as_mut()
3062 .ok_or_else(|| AdapterError::Transport("ACP agent has no stdout".into()))?;
3063 read_bounded_line(reader).await
3064 }
3065
3066 async fn start(&mut self) -> AdapterResult<()> {
3067 self.modes.clear();
3068 if self.child.is_some() {
3070 self.stop().await?;
3071 }
3072 let mut command = Command::new(&self.program);
3073 isolate_process_group(&mut command);
3074 command
3075 .args(&self.args)
3076 .current_dir(&self.cwd)
3077 .stdin(Stdio::piped())
3078 .stdout(Stdio::piped())
3079 .stderr(Stdio::piped())
3080 .env("CODESWARM_CWD", &self.cwd);
3081 if self.program.to_ascii_lowercase().contains("gemini")
3082 || self
3083 .args
3084 .iter()
3085 .any(|arg| arg.to_ascii_lowercase().contains("gemini"))
3086 {
3087 command.env("GEMINI_TELEMETRY_ENABLED", "false");
3088 }
3089 let mut child = command
3090 .spawn()
3091 .map_err(|error| AdapterError::Spawn(error.to_string()))?;
3092 let stdout = match child.stdout.take() {
3093 Some(stdout) => stdout,
3094 None => {
3095 let _ = terminate_child(&mut child).await;
3096 return Err(AdapterError::Transport("ACP agent has no stdout".into()));
3097 }
3098 };
3099 let stderr = match child.stderr.take() {
3100 Some(stderr) => stderr,
3101 None => {
3102 let _ = terminate_child(&mut child).await;
3103 return Err(AdapterError::Transport("ACP agent has no stderr".into()));
3104 }
3105 };
3106 self.child = Some(child);
3107 self.reader = Some(BufReader::new(stdout));
3108 self.stderr_task = Some(tokio::spawn(drain_bounded(stderr, 32 * 1024)));
3109
3110 let initialize = match self
3111 .request(
3112 "initialize",
3113 serde_json::json!({
3114 "protocolVersion": 1,
3115 "clientCapabilities": {
3116 "fs": {"readTextFile": true, "writeTextFile": true},
3117 "terminal": true,
3118 },
3119 "clientInfo": {
3120 "name": "CodeSwarm",
3121 "title": "CodeSwarm",
3122 "version": env!("CARGO_PKG_VERSION"),
3123 },
3124 }),
3125 )
3126 .await
3127 {
3128 Ok(value) => value,
3129 Err(error) => {
3130 let _ = self.stop().await;
3131 return Err(error);
3132 }
3133 };
3134 let agent_capabilities = initialize
3135 .get("agentCapabilities")
3136 .cloned()
3137 .unwrap_or(Value::Null);
3138 self.capabilities = AgentCapabilities {
3139 supports_cancel: true,
3140 supports_modes: true,
3141 supports_permissions: true,
3142 supports_terminals: true,
3143 supports_session_load: agent_capabilities
3144 .get("loadSession")
3145 .and_then(Value::as_bool)
3146 .unwrap_or(false),
3147 supports_models: false,
3148 };
3149 let session = if let Some(session_id) = self.session_id.clone() {
3150 if !self.capabilities.supports_session_load {
3151 let _ = self.stop().await;
3152 return Err(AdapterError::Unsupported("session/load"));
3153 }
3154 match self
3155 .request(
3156 "session/load",
3157 serde_json::json!({
3158 "cwd": self.cwd,
3159 "mcpServers": [],
3160 "sessionId": session_id,
3161 }),
3162 )
3163 .await
3164 {
3165 Ok(value) => value,
3166 Err(error) => {
3167 let _ = self.stop().await;
3168 return Err(error);
3169 }
3170 }
3171 } else {
3172 let session = match self
3173 .request(
3174 "session/new",
3175 serde_json::json!({"cwd": self.cwd, "mcpServers": []}),
3176 )
3177 .await
3178 {
3179 Ok(value) => value,
3180 Err(error) => {
3181 let _ = self.stop().await;
3182 return Err(error);
3183 }
3184 };
3185 self.session_id = session
3186 .get("sessionId")
3187 .and_then(Value::as_str)
3188 .map(str::to_owned);
3189 if self.session_id.is_none() {
3190 let _ = self.stop().await;
3191 return Err(AdapterError::Protocol(
3192 "session/new returned no sessionId".into(),
3193 ));
3194 }
3195 session
3196 };
3197 self.capabilities.supports_modes = false;
3198 if let Some(modes) = session.get("modes") {
3199 let available = modes
3200 .get("availableModes")
3201 .and_then(Value::as_array)
3202 .map(|modes| {
3203 modes
3204 .iter()
3205 .filter_map(|mode| {
3206 Some(Mode {
3207 id: mode.get("id")?.as_str()?.to_owned(),
3208 label: mode.get("name")?.as_str()?.to_owned(),
3209 })
3210 })
3211 .collect::<Vec<_>>()
3212 })
3213 .unwrap_or_default();
3214 self.modes = available.clone();
3215 self.capabilities.supports_modes = !available.is_empty();
3216 self.queued_events.push_back(Ok(AgentEvent::ModesReplaced {
3217 slot: self.slot,
3218 modes: available,
3219 current_mode: modes
3220 .get("currentModeId")
3221 .and_then(Value::as_str)
3222 .map(str::to_owned),
3223 }));
3224 }
3225 self.models.clear();
3226 self.model_config_id = None;
3227 let current_model =
3228 parse_model_config(&session).and_then(|(config_id, models, current)| {
3229 self.model_config_id = Some(config_id);
3230 self.models = models;
3231 current
3232 });
3233 self.capabilities.supports_models =
3234 self.model_config_id.is_some() && !self.models.is_empty();
3235 if let Some(config_id) = self.model_config_id.clone()
3236 && !self.models.is_empty()
3237 {
3238 self.queued_events.push_back(Ok(AgentEvent::ModelsReplaced {
3239 slot: self.slot,
3240 config_id,
3241 models: self.models.clone(),
3242 current_model,
3243 }));
3244 }
3245 self.queued_events.push_back(Ok(AgentEvent::Ready {
3246 slot: self.slot,
3247 capabilities: self.capabilities(),
3248 }));
3249 Ok(())
3250 }
3251}
3252
3253fn prompt_resource_paths(prompt: &str) -> Vec<String> {
3254 let characters = prompt.chars().collect::<Vec<_>>();
3255 let mut paths = Vec::new();
3256 let mut index = 0;
3257 while index < characters.len() {
3258 if characters[index] != '@' {
3259 index += 1;
3260 continue;
3261 }
3262 index += 1;
3263 let quoted = characters.get(index) == Some(&'"');
3264 if quoted {
3265 index += 1;
3266 }
3267 let start = index;
3268 while index < characters.len()
3269 && if quoted {
3270 characters[index] != '"'
3271 } else {
3272 !characters[index].is_whitespace()
3273 }
3274 {
3275 index += 1;
3276 }
3277 if index > start {
3278 paths.push(characters[start..index].iter().collect());
3279 }
3280 if quoted && index < characters.len() {
3281 index += 1;
3282 }
3283 }
3284 paths
3285}
3286
3287fn prompt_content_blocks(cwd: &Path, prompt: &str) -> Vec<Value> {
3288 let mut blocks = vec![serde_json::json!({"type": "text", "text": prompt})];
3289 for path in prompt_resource_paths(prompt) {
3290 if path.ends_with('/') {
3291 continue;
3292 }
3293 let Ok(resource) = resources::load(cwd, &path) else {
3294 continue;
3295 };
3296 let uri = format!("file://{}", resource.path.display());
3297 let resource_value = if let Some(text) = resource.text {
3298 serde_json::json!({
3299 "uri": uri,
3300 "text": text,
3301 "mimeType": resource.mime_type,
3302 })
3303 } else if let Some(data) = resource.data {
3304 serde_json::json!({
3305 "uri": uri,
3306 "blob": BASE64.encode(data),
3307 "mimeType": resource.mime_type,
3308 })
3309 } else {
3310 continue;
3311 };
3312 blocks.push(serde_json::json!({
3313 "type": "resource",
3314 "resource": resource_value,
3315 }));
3316 }
3317 blocks
3318}
3319
3320#[async_trait]
3321impl AgentAdapter for AcpAdapter {
3322 fn slot(&self) -> RosterSlot {
3323 self.slot
3324 }
3325
3326 fn session_id(&self) -> Option<String> {
3327 self.session_id.clone()
3328 }
3329
3330 fn protocol(&self) -> &'static str {
3331 "acp"
3332 }
3333
3334 fn needs_restart(&self) -> bool {
3335 self.reader.is_none() || self.child.is_none()
3336 }
3337
3338 fn capabilities(&self) -> AgentCapabilities {
3339 self.capabilities.clone()
3340 }
3341
3342 async fn start(&mut self) -> AdapterResult<()> {
3343 AcpAdapter::start(self).await
3346 }
3347
3348 async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
3349 if self.needs_restart() {
3350 return Err(AdapterError::Transport(
3351 "ACP agent transport is not running; reload the agent before retrying".into(),
3352 ));
3353 }
3354 let session_id = self
3355 .session_id
3356 .as_ref()
3357 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3358 self.tool_updates.clear();
3359 let request_id = self.next_request_id;
3360 self.next_request_id += 1;
3361 let prompt_blocks = prompt_content_blocks(&self.cwd, &prompt);
3362 let write_result = self
3363 .write_json(serde_json::json!({
3364 "jsonrpc": "2.0",
3365 "id": request_id,
3366 "method": "session/prompt",
3367 "params": {
3368 "sessionId": session_id,
3369 "prompt": prompt_blocks,
3370 },
3371 }))
3372 .await;
3373 if let Err(error) = write_result {
3374 self.reset_transport().await;
3375 return Err(error);
3376 }
3377 self.prompt_request_id = Some(request_id);
3378 Ok(())
3379 }
3380
3381 async fn cancel(&mut self) -> AdapterResult<bool> {
3382 let Some(session_id) = &self.session_id else {
3383 return Ok(false);
3384 };
3385 self.write_json(serde_json::json!({
3386 "jsonrpc": "2.0",
3387 "method": "session/cancel",
3388 "params": {"sessionId": session_id, "_meta": {}},
3389 }))
3390 .await?;
3391 let settled = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
3392 loop {
3393 match <Self as AgentAdapter>::next_event(self).await {
3394 Some(Ok(AgentEvent::TurnComplete { .. })) | None => break,
3395 Some(Ok(_)) => {}
3396 Some(Err(_)) => break,
3397 }
3398 }
3399 })
3400 .await
3401 .is_ok();
3402 if !settled {
3403 self.reload().await?;
3407 }
3408 Ok(true)
3409 }
3410
3411 async fn answer_permission(
3412 &mut self,
3413 request_id: String,
3414 answer: PermissionAnswer,
3415 ) -> AdapterResult<()> {
3416 let id = request_id
3417 .parse::<u64>()
3418 .map(Value::from)
3419 .unwrap_or_else(|_| Value::String(request_id));
3420 let outcome = match answer {
3421 PermissionAnswer::Selected { option_id } => {
3422 serde_json::json!({"outcome": "selected", "optionId": option_id})
3423 }
3424 PermissionAnswer::Cancelled => serde_json::json!({"outcome": "cancelled"}),
3425 };
3426 self.write_json(serde_json::json!({
3427 "jsonrpc": "2.0",
3428 "id": id,
3429 "result": {"outcome": outcome},
3433 }))
3434 .await
3435 }
3436
3437 async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
3438 let session_id = self
3439 .session_id
3440 .as_ref()
3441 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3442 let policy = match mode.as_str() {
3443 "plan" => "codeswarm:mode:plan",
3444 "default" | "manual" => "codeswarm:mode:manual",
3445 "accept-edits" => "codeswarm:mode:accept-edits",
3446 "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
3447 other => other,
3448 };
3449 let native_mode = crate::policy::resolve(policy, &self.modes)
3450 .map(|mode| mode.id)
3451 .unwrap_or(mode);
3452 let _ = self
3453 .request(
3454 "session/set_mode",
3455 serde_json::json!({"sessionId": session_id, "modeId": native_mode.clone()}),
3456 )
3457 .await?;
3458 self.queued_events.push_back(Ok(AgentEvent::ModeUpdated {
3459 slot: self.slot,
3460 current_mode: native_mode,
3461 }));
3462 Ok(())
3463 }
3464
3465 async fn set_model(&mut self, model: String) -> AdapterResult<()> {
3466 let session_id = self
3467 .session_id
3468 .clone()
3469 .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3470 let config_id = self
3471 .model_config_id
3472 .clone()
3473 .ok_or(AdapterError::Unsupported("set_model"))?;
3474 if !self.models.iter().any(|candidate| candidate.id == model) {
3475 return Err(AdapterError::Protocol(
3476 "model is not advertised by the agent".into(),
3477 ));
3478 }
3479 let _ = self
3480 .request(
3481 "session/set_config_option",
3482 serde_json::json!({
3483 "sessionId": session_id,
3484 "configId": config_id,
3485 "value": model,
3486 }),
3487 )
3488 .await?;
3489 Ok(())
3490 }
3491
3492 async fn reload(&mut self) -> AdapterResult<()> {
3493 let session_id = self
3498 .capabilities
3499 .supports_session_load
3500 .then(|| self.session_id.clone())
3501 .flatten();
3502 self.stop().await?;
3503 self.session_id = session_id.clone();
3504 let result = self.start().await;
3505 if result.is_err() {
3506 self.session_id = session_id;
3510 }
3511 result
3512 }
3513
3514 async fn stop(&mut self) -> AdapterResult<()> {
3515 let terminals = std::mem::take(&mut self.terminals);
3516 self.queued_events.clear();
3517 self.tool_updates.clear();
3518 for terminal in terminals.values() {
3519 terminal.stop().await;
3520 }
3521 if let Some(mut child) = self.child.take() {
3522 terminate_child(&mut child).await?;
3523 }
3524 self.reader = None;
3525 self.session_id = None;
3526 self.prompt_request_id = None;
3527 if let Some(task) = self.stderr_task.take() {
3528 task.abort();
3529 let _ = task.await;
3530 }
3531 Ok(())
3532 }
3533
3534 async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
3535 if let Some(event) = self.queued_events.pop_front() {
3536 return Some(event);
3537 }
3538 loop {
3539 let line = match self.read_line().await {
3540 Ok(line) => line,
3541 Err(error) => {
3542 self.reset_transport().await;
3546 return Some(Err(error));
3547 }
3548 };
3549 let value: Value = match serde_json::from_str(&line) {
3550 Ok(value) => value,
3551 Err(_) => continue,
3552 };
3553 match self.reject_empty_permission_request(&value).await {
3554 Ok(true) => continue,
3555 Ok(false) => {}
3556 Err(error) => return Some(Err(error)),
3557 }
3558 match self.handle_client_request(&value).await {
3559 Ok(true) => continue,
3560 Ok(false) => {}
3561 Err(error) => return Some(Err(error)),
3562 }
3563 match parse_acp_value(self.slot, &value, &mut self.tool_updates) {
3564 Ok(Some(event)) => {
3565 if let AgentEvent::ModelsReplaced {
3566 config_id, models, ..
3567 } = &event
3568 {
3569 self.model_config_id = Some(config_id.clone());
3570 self.models = models.clone();
3571 self.capabilities.supports_models = !models.is_empty();
3572 }
3573 return Some(Ok(event));
3574 }
3575 Ok(None) => {}
3576 Err(error) => return Some(Err(error)),
3577 }
3578 if value.get("id").is_some_and(|id| {
3579 self.prompt_request_id
3580 .is_some_and(|expected| rpc_id_to_string(id) == expected.to_string())
3581 }) {
3582 if let Some(error) = value.get("error") {
3583 self.prompt_request_id = None;
3584 return Some(Err(AdapterError::Protocol(error.to_string())));
3585 }
3586 self.prompt_request_id = None;
3587 if let Some(reason) = value
3588 .get("result")
3589 .and_then(|result| result.get("stopReason"))
3590 .and_then(Value::as_str)
3591 .filter(|reason| !matches!(*reason, "end_turn" | "cancelled"))
3592 {
3593 let detail = match reason {
3594 "max_tokens" => {
3595 "ACP turn stopped because the output token limit was reached".into()
3596 }
3597 "max_turn_requests" => {
3598 "ACP turn stopped because the tool-turn limit was reached".into()
3599 }
3600 other => format!("ACP turn stopped before responding: {other}"),
3601 };
3602 return Some(Ok(AgentEvent::Failed {
3603 slot: self.slot,
3604 started: true,
3605 detail,
3606 }));
3607 }
3608 return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3609 }
3610 }
3611 }
3612}
3613
3614#[cfg(test)]
3615fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3616 let value: Value =
3617 serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3618 parse_acp_value(slot, &value, &mut BTreeMap::new())
3619}
3620
3621fn parse_acp_value(
3622 slot: RosterSlot,
3623 value: &Value,
3624 tools: &mut BTreeMap<String, ToolUpdate>,
3625) -> AdapterResult<Option<AgentEvent>> {
3626 let method = value.get("method").and_then(Value::as_str);
3627 if method == Some("session/request_permission") {
3628 let params = value.get("params").cloned().unwrap_or(Value::Null);
3629 let request_id = value
3630 .get("id")
3631 .map(rpc_id_to_string)
3632 .unwrap_or_else(|| "permission".into());
3633 return Ok(parse_permission_event(
3634 slot,
3635 ¶ms,
3636 &request_id,
3637 params.get("options"),
3638 ));
3639 }
3640 if method != Some("session/update") {
3641 return Ok(None);
3642 }
3643 let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3644 return Ok(None);
3645 };
3646 let kind = update.get("sessionUpdate").and_then(Value::as_str);
3647 if kind == Some("config_option_update")
3648 && let Some((config_id, models, current_model)) = parse_model_config(update)
3649 {
3650 return Ok(Some(AgentEvent::ModelsReplaced {
3651 slot,
3652 config_id,
3653 models,
3654 current_model,
3655 }));
3656 }
3657 if kind == Some("request_permission") {
3658 let request_id = update
3659 .get("toolCall")
3660 .and_then(|tool| tool.get("toolCallId"))
3661 .and_then(Value::as_str)
3662 .unwrap_or("permission");
3663 return Ok(parse_permission_event(
3664 slot,
3665 update,
3666 request_id,
3667 update.get("options"),
3668 ));
3669 }
3670 if kind == Some("available_commands_update") {
3671 let commands = update
3672 .get("availableCommands")
3673 .and_then(Value::as_array)
3674 .map(|commands| {
3675 commands
3676 .iter()
3677 .filter_map(|command| {
3678 let name = command.get("name").and_then(Value::as_str)?.trim();
3679 (!name.is_empty()).then(|| AgentCommand {
3680 name: name.to_owned(),
3681 })
3682 })
3683 .collect::<Vec<_>>()
3684 })
3685 .unwrap_or_default();
3686 return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3687 }
3688 if kind == Some("current_mode_update") {
3689 if let Some(mode) = update
3690 .get("currentModeId")
3691 .and_then(Value::as_str)
3692 .filter(|mode| !mode.trim().is_empty())
3693 {
3694 return Ok(Some(AgentEvent::ModeUpdated {
3695 slot,
3696 current_mode: mode.to_owned(),
3697 }));
3698 }
3699 return Ok(None);
3700 }
3701 if kind == Some("usage_update") {
3702 let Some(used) = update.get("used").and_then(Value::as_u64) else {
3703 return Ok(None);
3704 };
3705 let Some(size) = update.get("size").and_then(Value::as_u64) else {
3706 return Ok(None);
3707 };
3708 return Ok(Some(AgentEvent::UsageUpdated {
3709 slot,
3710 usage: UsageUpdate { used, size },
3711 }));
3712 }
3713 if let Some(terminal) = parse_terminal_event(update, kind) {
3714 return Ok(Some(AgentEvent::Terminal {
3715 slot,
3716 event: terminal,
3717 }));
3718 }
3719 let text = update
3720 .get("content")
3721 .and_then(|content| content.get("text"))
3722 .and_then(Value::as_str)
3723 .map(str::to_owned);
3724 if kind == Some("user_message_chunk") {
3725 return Ok(text
3726 .filter(|text| !text.is_empty())
3727 .map(|text| AgentEvent::UserText { slot, text }));
3728 }
3729 if kind == Some("agent_message_chunk")
3730 && let Some(mode) = text
3731 .as_deref()
3732 .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3733 .map(str::trim)
3734 .filter(|mode| !mode.is_empty())
3735 {
3736 return Ok(Some(AgentEvent::ModesReplaced {
3741 slot,
3742 modes: vec![Mode {
3743 id: mode.to_owned(),
3744 label: mode.to_owned(),
3745 }],
3746 current_mode: Some(mode.to_owned()),
3747 }));
3748 }
3749 match (kind, text) {
3750 (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3751 Ok(Some(AgentEvent::Text { slot, text }))
3752 }
3753 (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3754 Ok(Some(AgentEvent::Thought { slot, text }))
3755 }
3756 (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3757 Ok(normalize_acp_tool(update, tools).map(|update| AgentEvent::Tool { slot, update }))
3758 }
3759 _ => Ok(None),
3760 }
3761}
3762
3763fn normalize_acp_tool(
3766 value: &Value,
3767 tools: &mut BTreeMap<String, ToolUpdate>,
3768) -> Option<ToolUpdate> {
3769 let id = value.get("toolCallId")?.as_str()?;
3770 if id.trim().is_empty() {
3771 return None;
3772 }
3773 if value.get("sessionUpdate").and_then(Value::as_str) == Some("tool_call") {
3774 tools.remove(id);
3775 }
3776 let tool = tools.entry(id.to_owned()).or_insert_with(|| ToolUpdate {
3777 id: id.to_owned(),
3778 title: "Tool call".into(),
3779 status: ToolStatus::Pending,
3780 detail: None,
3781 });
3782 if let Some(title) = value.get("title").and_then(Value::as_str) {
3783 tool.title = title.to_owned();
3784 }
3785 if let Some(status) =
3786 value
3787 .get("status")
3788 .and_then(Value::as_str)
3789 .and_then(|status| match status {
3790 "pending" => Some(ToolStatus::Pending),
3791 "in_progress" => Some(ToolStatus::Running),
3792 "completed" => Some(ToolStatus::Completed),
3793 "failed" => Some(ToolStatus::Failed),
3794 _ => None,
3795 })
3796 {
3797 tool.status = status;
3798 }
3799 if let Some(content) = value.get("content").and_then(Value::as_array) {
3800 let text = content
3801 .iter()
3802 .filter_map(|entry| match entry.get("type").and_then(Value::as_str) {
3803 Some("content") => entry.get("content")?.get("text")?.as_str(),
3804 Some("diff") => entry.get("newText")?.as_str(),
3805 _ => None,
3806 })
3807 .collect::<Vec<_>>()
3808 .join("\n");
3809 tool.detail = (!text.is_empty()).then_some(text);
3810 } else if let Some(output) = value.get("rawOutput").filter(|output| !output.is_null()) {
3811 tool.detail = Some(
3812 output
3813 .as_str()
3814 .map(str::to_owned)
3815 .unwrap_or_else(|| output.to_string()),
3816 );
3817 }
3818 Some(tool.clone())
3819}
3820
3821fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3822 let config = value
3823 .get("configOptions")?
3824 .as_array()?
3825 .iter()
3826 .find(|option| {
3827 option.get("category").and_then(Value::as_str) == Some("model")
3828 && matches!(
3829 option.get("type").and_then(Value::as_str),
3830 Some("select" | "enum")
3831 )
3832 })?;
3833 let config_id = config.get("id")?.as_str()?.to_owned();
3834 let models = config
3835 .get("options")?
3836 .as_array()?
3837 .iter()
3838 .filter_map(|option| {
3839 let id = option.get("value")?.as_str()?.to_owned();
3840 let label = option
3841 .get("name")
3842 .or_else(|| option.get("label"))
3843 .and_then(Value::as_str)
3844 .unwrap_or(&id)
3845 .to_owned();
3846 Some(Mode { id, label })
3847 })
3848 .collect::<Vec<_>>();
3849 (!models.is_empty()).then(|| {
3850 let current = config
3851 .get("currentValue")
3852 .and_then(Value::as_str)
3853 .map(str::to_owned);
3854 (config_id, models, current)
3855 })
3856}
3857
3858fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3863 let nested = value.get("terminal").unwrap_or(value);
3864 let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3865 let id = nested
3866 .get("terminalId")
3867 .or_else(|| nested.get("terminal_id"))
3868 .or_else(|| nested.get("id"))
3869 .and_then(Value::as_str)
3870 .unwrap_or("terminal")
3871 .to_owned();
3872 match kind {
3873 "terminal_created" | "terminal_create" | "terminal_started" => {
3874 let command = nested
3875 .get("command")
3876 .and_then(Value::as_str)
3877 .unwrap_or("")
3878 .to_owned();
3879 Some(TerminalEvent::Created { id, command })
3880 }
3881 "terminal_output" | "terminal_output_chunk" => {
3882 let text = nested
3883 .get("output")
3884 .or_else(|| nested.get("text"))
3885 .and_then(Value::as_str)
3886 .unwrap_or("")
3887 .to_owned();
3888 Some(TerminalEvent::Output { id, text })
3889 }
3890 "terminal_exited" | "terminal_exit" => {
3891 let code = nested
3892 .get("exitCode")
3893 .or_else(|| nested.get("exit_code"))
3894 .or_else(|| nested.get("code"))
3895 .and_then(Value::as_i64)
3896 .unwrap_or(0) as i32;
3897 Some(TerminalEvent::Exited { id, code })
3898 }
3899 "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3900 _ => None,
3901 }
3902}
3903
3904fn parse_permission_event(
3905 slot: RosterSlot,
3906 value: &Value,
3907 request_id: &str,
3908 options: Option<&Value>,
3909) -> Option<AgentEvent> {
3910 let tool = value.get("toolCall").unwrap_or(value);
3911 let title = tool
3912 .get("title")
3913 .and_then(Value::as_str)
3914 .unwrap_or("Agent requests permission")
3915 .to_owned();
3916 let (options, option_ids): (Vec<String>, Vec<String>) = options
3917 .and_then(Value::as_array)
3918 .map(|options| {
3919 options
3920 .iter()
3921 .filter_map(|option| {
3922 let label = option
3923 .get("name")
3924 .or_else(|| option.get("optionId"))
3925 .and_then(Value::as_str)?
3926 .to_owned();
3927 let option_id = option
3928 .get("optionId")
3929 .or_else(|| option.get("id"))
3930 .and_then(Value::as_str)
3931 .map(str::to_owned)
3932 .unwrap_or_else(|| label.clone());
3933 Some((label, option_id))
3934 })
3935 .unzip()
3936 })
3937 .unwrap_or_default();
3938 if options.is_empty() {
3939 return None;
3940 }
3941 Some(AgentEvent::Permission {
3942 slot,
3943 request: PermissionRequest {
3944 id: request_id.to_owned(),
3945 title,
3946 options,
3947 option_ids,
3948 },
3949 })
3950}
3951
3952fn rpc_id_to_string(value: &Value) -> String {
3953 value
3954 .as_str()
3955 .map(str::to_owned)
3956 .or_else(|| value.as_u64().map(|id| id.to_string()))
3957 .unwrap_or_else(|| value.to_string())
3958}
3959
3960#[cfg(test)]
3961mod tests {
3962 use super::{
3963 AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
3964 RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
3965 parse_model_config, prompt_content_blocks, read_bounded_line,
3966 };
3967 #[cfg(target_os = "linux")]
3968 use super::{isolate_process_group, terminate_child};
3969 use crate::TerminalEvent;
3970 use crate::{
3971 AdapterError, AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
3972 persistence::SessionMetadataStore,
3973 relay::{CollaborationStrategy, DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
3974 };
3975 use async_trait::async_trait;
3976 use serde_json::Value;
3977 use std::collections::VecDeque;
3978 use std::sync::{
3979 Arc, Mutex,
3980 atomic::{AtomicUsize, Ordering},
3981 };
3982
3983 fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
3984 let nonce = std::time::SystemTime::now()
3985 .duration_since(std::time::UNIX_EPOCH)
3986 .expect("clock")
3987 .as_nanos();
3988 std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
3989 }
3990
3991 #[test]
3992 fn malformed_file_writes_preserve_existing_content() {
3993 let root = unique_test_path("codeswarm-write-validation", "dir");
3994 std::fs::create_dir_all(&root).unwrap();
3995 let file = root.join("keep.txt");
3996 std::fs::write(&file, "valuable content").unwrap();
3997 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
3998 for content in [
3999 Value::Null,
4000 serde_json::json!(false),
4001 serde_json::json!(42),
4002 serde_json::json!([]),
4003 ] {
4004 assert!(
4005 adapter
4006 .write_workspace_text(
4007 &serde_json::json!({"path":"keep.txt", "content": content})
4008 )
4009 .is_err()
4010 );
4011 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4012 }
4013 assert!(
4014 adapter
4015 .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
4016 .is_err()
4017 );
4018 assert!(
4019 adapter
4020 .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
4021 .is_err()
4022 );
4023 assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4024 adapter
4025 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
4026 .unwrap();
4027 assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
4028 adapter
4029 .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
4030 .unwrap();
4031 assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
4032 std::fs::remove_dir_all(root).unwrap();
4033 }
4034
4035 #[tokio::test]
4036 async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
4037 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
4038 let mut adapter = AcpAdapter::new(
4039 0,
4040 std::env::current_dir().unwrap(),
4041 "sh",
4042 vec!["-c".into(), script.into()],
4043 );
4044 adapter.start().await.unwrap();
4045 let error = adapter
4046 .request_with_timeout(
4047 "session/set_mode",
4048 serde_json::json!({}),
4049 std::time::Duration::from_millis(10),
4050 )
4051 .await
4052 .unwrap_err();
4053 assert!(error.to_string().contains("session/set_mode timed out"));
4054 adapter.stop().await.unwrap();
4055 assert!(adapter.child.is_none());
4056 }
4057
4058 #[tokio::test]
4059 async fn goals_reach_every_roster_slot_without_native_goal_support() {
4060 use crate::goal::GoalCommand;
4061 let hosts = (0..3)
4062 .map(|slot| {
4063 AdapterHost::new(
4064 Box::new(ScriptedAdapter::new(
4065 slot,
4066 AgentCapabilities::default(),
4067 [
4068 AgentEvent::TurnComplete { slot },
4069 AgentEvent::TurnComplete { slot },
4070 ],
4071 )),
4072 None,
4073 )
4074 })
4075 .collect();
4076 let mut relay = RelayHost::new(hosts, 10).unwrap();
4077 relay.start().await.unwrap();
4078 let task = relay
4079 .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
4080 .unwrap()
4081 .unwrap();
4082 relay.relay_mut().enqueue_human(task, Some(0));
4083 for slot in 0..3 {
4084 relay.run_turn("", 0).await.unwrap();
4085 let (actual, prompt) = relay.dispatches().last().unwrap();
4086 assert_eq!(*actual, slot);
4087 assert!(prompt.contains("Active shared goal: Ship the settings screen"));
4088 }
4089 let snapshot = relay.session_metadata();
4090 let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
4091 assert!(restored.is_some());
4092 relay.restore_goal(restored);
4093 relay.reload(0).await.unwrap();
4094 relay.run_turn("", 0).await.unwrap();
4095 assert!(
4096 relay
4097 .dispatches()
4098 .last()
4099 .unwrap()
4100 .1
4101 .contains("Active shared goal: Ship the settings screen")
4102 );
4103 relay.apply_goal(GoalCommand::Done).unwrap();
4104 relay.run_turn("", 0).await.unwrap();
4105 assert!(
4106 relay
4107 .dispatches()
4108 .last()
4109 .unwrap()
4110 .1
4111 .contains("No active shared goal")
4112 );
4113 relay.apply_goal(GoalCommand::Clear).unwrap();
4114 assert!(relay.session_metadata().get("goal").unwrap().is_null());
4115 }
4116
4117 #[tokio::test]
4118 async fn replacement_agent_receives_task_after_public_journal_pruning() {
4119 let hosts = (0..2)
4120 .map(|slot| {
4121 AdapterHost::new(
4122 Box::new(ScriptedAdapter::new(
4123 slot,
4124 AgentCapabilities::default(),
4125 [
4126 AgentEvent::Text {
4127 slot,
4128 text: "progress".into(),
4129 },
4130 AgentEvent::TurnComplete { slot },
4131 AgentEvent::Text {
4132 slot,
4133 text: "more progress".into(),
4134 },
4135 AgentEvent::TurnComplete { slot },
4136 ],
4137 )),
4138 None,
4139 )
4140 })
4141 .collect();
4142 let mut relay = RelayHost::new(hosts, 10).unwrap();
4143 relay.start().await.unwrap();
4144 relay
4145 .relay_mut()
4146 .enqueue_human("Fix the login bug", Some(0));
4147 relay.run_turn("", 0).await.unwrap();
4148 relay.run_turn("", 0).await.unwrap();
4149 relay.run_turn("", 0).await.unwrap();
4150 relay.reload(1).await.unwrap();
4151 assert!(
4152 !relay
4153 .relay_mut()
4154 .unseen_context(1)
4155 .contains("Fix the login bug")
4156 );
4157 relay.run_turn("", 0).await.unwrap();
4158 assert!(
4159 relay
4160 .dispatches()
4161 .last()
4162 .unwrap()
4163 .1
4164 .contains("Shared task:\nFix the login bug")
4165 );
4166 }
4167
4168 #[cfg(target_os = "linux")]
4169 #[tokio::test]
4170 async fn termination_kills_only_the_verified_isolated_child_group() {
4171 use nix::unistd::{Pid, getpgid, getpgrp};
4172 use tokio::io::{AsyncBufReadExt, BufReader};
4173
4174 let own_group = getpgrp();
4175 let mut command = tokio::process::Command::new("sh");
4176 isolate_process_group(&mut command);
4177 command
4178 .arg("-c")
4179 .arg("sleep 60 & echo $!; wait")
4180 .stdout(std::process::Stdio::piped());
4181 let mut child = command.spawn().expect("spawn isolated shell");
4182 let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
4183 assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
4184 assert_ne!(leader, own_group);
4185
4186 let stdout = child.stdout.take().expect("child stdout");
4187 let mut lines = BufReader::new(stdout).lines();
4188 let descendant = lines
4189 .next_line()
4190 .await
4191 .expect("read descendant pid")
4192 .expect("descendant pid")
4193 .parse::<i32>()
4194 .expect("numeric descendant pid");
4195 let descendant = Pid::from_raw(descendant);
4196 assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
4197
4198 terminate_child(&mut child).await.expect("terminate group");
4199 for _ in 0..100 {
4200 if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
4201 return;
4202 }
4203 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
4204 }
4205 panic!("descendant {descendant} survived isolated group termination");
4206 }
4207
4208 #[test]
4209 fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
4210 assert_eq!(
4211 parse_command_line(r#"agent --name "local bridge" --flag 'two words'"#),
4212 Ok((
4213 "agent".into(),
4214 vec![
4215 "--name".into(),
4216 "local bridge".into(),
4217 "--flag".into(),
4218 "two words".into(),
4219 ]
4220 ),)
4221 );
4222 assert_eq!(
4223 parse_command_line(r#"agent "" escaped\ argument"#),
4224 Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
4225 );
4226 }
4227
4228 #[test]
4229 fn acp_prompt_expands_safe_at_path_resources() {
4230 let root = unique_test_path("codeswarm-prompt-resource", "dir");
4231 std::fs::create_dir_all(&root).expect("workspace");
4232 std::fs::write(root.join("note.md"), "resource text").expect("resource");
4233 let blocks = prompt_content_blocks(&root, "inspect @note.md");
4234 assert_eq!(blocks[0]["type"], "text");
4235 assert_eq!(blocks[0]["text"], "inspect @note.md");
4236 assert_eq!(blocks[1]["type"], "resource");
4237 assert_eq!(blocks[1]["resource"]["text"], "resource text");
4238 assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
4239 std::fs::remove_dir_all(root).expect("cleanup workspace");
4240 }
4241
4242 #[tokio::test]
4243 async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
4244 let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
4245 bytes.push(b'\n');
4246 let mut reader = tokio::io::BufReader::new(bytes.as_slice());
4247 assert!(matches!(
4248 read_bounded_line(&mut reader).await,
4249 Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
4250 ));
4251 }
4252
4253 #[test]
4254 fn rejects_malformed_configured_commands_before_spawn() {
4255 assert_eq!(
4256 parse_command_line("agent 'unfinished"),
4257 Err(super::CommandParseError::UnterminatedQuote)
4258 );
4259 assert_eq!(
4260 parse_command_line("agent\\"),
4261 Err(super::CommandParseError::TrailingEscape)
4262 );
4263 assert_eq!(
4264 parse_command_line(" \t"),
4265 Err(super::CommandParseError::Empty)
4266 );
4267 }
4268
4269 #[derive(Debug)]
4270 struct PendingAdapter {
4271 slot: usize,
4272 hang_on_cancel: bool,
4273 }
4274
4275 #[derive(Debug)]
4276 struct ConcurrentStartAdapter {
4277 slot: usize,
4278 barrier: Arc<tokio::sync::Barrier>,
4279 }
4280
4281 #[derive(Debug)]
4282 struct ReloadProbeAdapter {
4283 slot: usize,
4284 crashed: bool,
4285 reloaded: bool,
4286 events: VecDeque<AgentEvent>,
4287 prompts: Arc<Mutex<Vec<String>>>,
4288 }
4289
4290 #[async_trait]
4291 impl AgentAdapter for ReloadProbeAdapter {
4292 fn slot(&self) -> usize {
4293 self.slot
4294 }
4295
4296 fn display_name(&self) -> String {
4297 "Reload probe".into()
4298 }
4299
4300 fn protocol(&self) -> &'static str {
4301 "native"
4302 }
4303
4304 fn capabilities(&self) -> AgentCapabilities {
4305 AgentCapabilities {
4306 supports_modes: true,
4307 ..AgentCapabilities::default()
4308 }
4309 }
4310
4311 fn needs_restart(&self) -> bool {
4312 self.crashed && !self.reloaded
4313 }
4314
4315 async fn start(&mut self) -> super::AdapterResult<()> {
4316 self.events.push_back(AgentEvent::ModesReplaced {
4317 slot: self.slot,
4318 modes: vec![
4319 Mode {
4320 id: "codeswarm:mode:full-access".into(),
4321 label: "Auto pilot".into(),
4322 },
4323 Mode {
4324 id: "codeswarm:mode:plan".into(),
4325 label: "Plan".into(),
4326 },
4327 ],
4328 current_mode: Some("codeswarm:mode:full-access".into()),
4329 });
4330 self.events.push_back(AgentEvent::Ready {
4331 slot: self.slot,
4332 capabilities: self.capabilities(),
4333 });
4334 Ok(())
4335 }
4336
4337 async fn send_prompt(&mut self, prompt: String) -> super::AdapterResult<()> {
4338 self.prompts.lock().expect("prompts").push(prompt);
4339 if !self.crashed {
4340 self.crashed = true;
4341 self.events.push_back(AgentEvent::Failed {
4342 slot: self.slot,
4343 started: true,
4344 detail: "probe crashed".into(),
4345 });
4346 } else {
4347 self.events.push_back(AgentEvent::Text {
4348 slot: self.slot,
4349 text: "recovered".into(),
4350 });
4351 self.events
4352 .push_back(AgentEvent::TurnComplete { slot: self.slot });
4353 }
4354 Ok(())
4355 }
4356
4357 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4358 Ok(false)
4359 }
4360
4361 async fn answer_permission(
4362 &mut self,
4363 _request_id: String,
4364 _answer: PermissionAnswer,
4365 ) -> super::AdapterResult<()> {
4366 Err(super::AdapterError::Unsupported("permission answer"))
4367 }
4368
4369 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4370 Ok(())
4371 }
4372
4373 async fn reload(&mut self) -> super::AdapterResult<()> {
4374 self.reloaded = true;
4375 self.start().await
4376 }
4377
4378 async fn stop(&mut self) -> super::AdapterResult<()> {
4379 Ok(())
4380 }
4381
4382 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4383 self.events.pop_front().map(Ok)
4384 }
4385 }
4386
4387 #[async_trait]
4388 impl AgentAdapter for ConcurrentStartAdapter {
4389 fn slot(&self) -> usize {
4390 self.slot
4391 }
4392
4393 fn capabilities(&self) -> AgentCapabilities {
4394 AgentCapabilities::default()
4395 }
4396
4397 async fn start(&mut self) -> super::AdapterResult<()> {
4398 self.barrier.wait().await;
4399 Ok(())
4400 }
4401
4402 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4403 Ok(())
4404 }
4405
4406 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4407 Ok(true)
4408 }
4409
4410 async fn answer_permission(
4411 &mut self,
4412 _request_id: String,
4413 _answer: PermissionAnswer,
4414 ) -> super::AdapterResult<()> {
4415 Ok(())
4416 }
4417
4418 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4419 Ok(())
4420 }
4421
4422 async fn reload(&mut self) -> super::AdapterResult<()> {
4423 Ok(())
4424 }
4425
4426 async fn stop(&mut self) -> super::AdapterResult<()> {
4427 Ok(())
4428 }
4429
4430 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4431 std::future::pending().await
4432 }
4433 }
4434
4435 #[derive(Debug)]
4436 struct PermissionBlockingAdapter {
4437 slot: usize,
4438 phase: u8,
4439 }
4440
4441 #[async_trait]
4442 impl AgentAdapter for PermissionBlockingAdapter {
4443 fn slot(&self) -> usize {
4444 self.slot
4445 }
4446
4447 fn capabilities(&self) -> AgentCapabilities {
4448 AgentCapabilities {
4449 supports_permissions: true,
4450 ..AgentCapabilities::default()
4451 }
4452 }
4453
4454 async fn start(&mut self) -> super::AdapterResult<()> {
4455 Ok(())
4456 }
4457
4458 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4459 Ok(())
4460 }
4461
4462 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4463 Ok(true)
4464 }
4465
4466 async fn answer_permission(
4467 &mut self,
4468 request_id: String,
4469 answer: PermissionAnswer,
4470 ) -> super::AdapterResult<()> {
4471 if self.phase != 1 || request_id != "permission-1" {
4472 return Err(super::AdapterError::Protocol(
4473 "unexpected permission response".into(),
4474 ));
4475 }
4476 assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4477 self.phase = 2;
4478 Ok(())
4479 }
4480
4481 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4482 Ok(())
4483 }
4484
4485 async fn reload(&mut self) -> super::AdapterResult<()> {
4486 Ok(())
4487 }
4488
4489 async fn stop(&mut self) -> super::AdapterResult<()> {
4490 Ok(())
4491 }
4492
4493 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4494 match self.phase {
4495 0 => {
4496 self.phase = 1;
4497 Some(Ok(AgentEvent::Permission {
4498 slot: self.slot,
4499 request: crate::PermissionRequest {
4500 id: "permission-1".into(),
4501 title: "Allow?".into(),
4502 options: vec!["Allow".into()],
4503 option_ids: vec!["allow".into()],
4504 },
4505 }))
4506 }
4507 1 => std::future::pending().await,
4508 _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4509 }
4510 }
4511 }
4512
4513 #[async_trait]
4514 impl AgentAdapter for PendingAdapter {
4515 fn slot(&self) -> usize {
4516 self.slot
4517 }
4518
4519 fn capabilities(&self) -> AgentCapabilities {
4520 AgentCapabilities {
4521 supports_cancel: true,
4522 ..AgentCapabilities::default()
4523 }
4524 }
4525
4526 async fn start(&mut self) -> super::AdapterResult<()> {
4527 Ok(())
4528 }
4529
4530 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4531 Ok(())
4532 }
4533
4534 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4535 if self.hang_on_cancel {
4536 return std::future::pending().await;
4537 }
4538 Ok(true)
4539 }
4540
4541 async fn answer_permission(
4542 &mut self,
4543 _request_id: String,
4544 _answer: PermissionAnswer,
4545 ) -> super::AdapterResult<()> {
4546 Ok(())
4547 }
4548
4549 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4550 Ok(())
4551 }
4552
4553 async fn reload(&mut self) -> super::AdapterResult<()> {
4554 Ok(())
4555 }
4556
4557 async fn stop(&mut self) -> super::AdapterResult<()> {
4558 Ok(())
4559 }
4560
4561 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4562 std::future::pending().await
4563 }
4564 }
4565
4566 #[derive(Debug)]
4567 struct StopTrackingAdapter {
4568 slot: usize,
4569 stopped: Arc<AtomicUsize>,
4570 fail_stop: bool,
4571 }
4572
4573 #[derive(Debug)]
4574 struct ModeOrderAdapter {
4575 slot: usize,
4576 log: Arc<Mutex<Vec<String>>>,
4577 phase: u8,
4578 }
4579
4580 #[derive(Debug)]
4581 struct StartupAcpAdapter {
4582 slot: usize,
4583 events: std::collections::VecDeque<AgentEvent>,
4584 }
4585
4586 impl StartupAcpAdapter {
4587 fn new(slot: usize) -> Self {
4588 Self {
4589 slot,
4590 events: [
4591 AgentEvent::ModesReplaced {
4592 slot,
4593 modes: vec![Mode {
4594 id: "full-access".into(),
4595 label: "Auto pilot".into(),
4596 }],
4597 current_mode: Some("full-access".into()),
4598 },
4599 AgentEvent::Ready {
4600 slot,
4601 capabilities: AgentCapabilities {
4602 supports_modes: true,
4603 ..AgentCapabilities::default()
4604 },
4605 },
4606 ]
4607 .into(),
4608 }
4609 }
4610 }
4611
4612 #[async_trait]
4613 impl AgentAdapter for StartupAcpAdapter {
4614 fn slot(&self) -> usize {
4615 self.slot
4616 }
4617
4618 fn protocol(&self) -> &'static str {
4619 "acp"
4620 }
4621
4622 fn capabilities(&self) -> AgentCapabilities {
4623 AgentCapabilities {
4624 supports_modes: true,
4625 ..AgentCapabilities::default()
4626 }
4627 }
4628
4629 async fn start(&mut self) -> super::AdapterResult<()> {
4630 Ok(())
4631 }
4632
4633 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4634 Ok(())
4635 }
4636
4637 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4638 Ok(true)
4639 }
4640
4641 async fn answer_permission(
4642 &mut self,
4643 _request_id: String,
4644 _answer: PermissionAnswer,
4645 ) -> super::AdapterResult<()> {
4646 Ok(())
4647 }
4648
4649 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4650 Ok(())
4651 }
4652
4653 async fn reload(&mut self) -> super::AdapterResult<()> {
4654 self.events = Self::new(self.slot).events;
4655 Ok(())
4656 }
4657
4658 async fn stop(&mut self) -> super::AdapterResult<()> {
4659 Ok(())
4660 }
4661
4662 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4663 self.events.pop_front().map(Ok)
4664 }
4665 }
4666
4667 #[async_trait]
4668 impl AgentAdapter for ModeOrderAdapter {
4669 fn slot(&self) -> usize {
4670 self.slot
4671 }
4672
4673 fn capabilities(&self) -> AgentCapabilities {
4674 AgentCapabilities {
4675 supports_modes: true,
4676 ..AgentCapabilities::default()
4677 }
4678 }
4679
4680 async fn start(&mut self) -> super::AdapterResult<()> {
4681 self.log.lock().expect("log").push("start".into());
4682 Ok(())
4683 }
4684
4685 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4686 self.log.lock().expect("log").push("prompt".into());
4687 Ok(())
4688 }
4689
4690 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4691 Ok(true)
4692 }
4693
4694 async fn answer_permission(
4695 &mut self,
4696 _request_id: String,
4697 _answer: PermissionAnswer,
4698 ) -> super::AdapterResult<()> {
4699 Ok(())
4700 }
4701
4702 async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4703 self.log.lock().expect("log").push(format!("mode:{mode}"));
4704 Ok(())
4705 }
4706
4707 async fn reload(&mut self) -> super::AdapterResult<()> {
4708 self.log.lock().expect("log").push("reload".into());
4709 self.phase = 0;
4710 Ok(())
4711 }
4712
4713 async fn stop(&mut self) -> super::AdapterResult<()> {
4714 Ok(())
4715 }
4716
4717 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4718 match self.phase {
4719 0 => {
4720 self.phase = 1;
4721 Some(Ok(AgentEvent::ModesReplaced {
4722 slot: self.slot,
4723 modes: vec![Mode {
4724 id: "yolo".into(),
4725 label: "YOLO".into(),
4726 }],
4727 current_mode: None,
4728 }))
4729 }
4730 1 => {
4731 self.phase = 2;
4732 Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4733 }
4734 _ => std::future::pending().await,
4735 }
4736 }
4737 }
4738
4739 #[async_trait]
4740 impl AgentAdapter for StopTrackingAdapter {
4741 fn slot(&self) -> usize {
4742 self.slot
4743 }
4744
4745 fn capabilities(&self) -> AgentCapabilities {
4746 AgentCapabilities::default()
4747 }
4748
4749 async fn start(&mut self) -> super::AdapterResult<()> {
4750 Ok(())
4751 }
4752
4753 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4754 Ok(())
4755 }
4756
4757 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4758 Ok(false)
4759 }
4760
4761 async fn answer_permission(
4762 &mut self,
4763 _request_id: String,
4764 _answer: PermissionAnswer,
4765 ) -> super::AdapterResult<()> {
4766 Ok(())
4767 }
4768
4769 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4770 Ok(())
4771 }
4772
4773 async fn reload(&mut self) -> super::AdapterResult<()> {
4774 Ok(())
4775 }
4776
4777 async fn stop(&mut self) -> super::AdapterResult<()> {
4778 self.stopped.fetch_add(1, Ordering::Relaxed);
4779 if self.fail_stop {
4780 Err(super::AdapterError::Transport("stop failed".into()))
4781 } else {
4782 Ok(())
4783 }
4784 }
4785
4786 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4787 None
4788 }
4789 }
4790
4791 #[derive(Debug)]
4795 struct FailingStartAdapter {
4796 slot: usize,
4797 stopped: Arc<AtomicUsize>,
4798 }
4799
4800 #[async_trait]
4801 impl AgentAdapter for FailingStartAdapter {
4802 fn slot(&self) -> usize {
4803 self.slot
4804 }
4805
4806 fn capabilities(&self) -> AgentCapabilities {
4807 AgentCapabilities::default()
4808 }
4809
4810 async fn start(&mut self) -> super::AdapterResult<()> {
4811 Err(super::AdapterError::Spawn("startup failed".into()))
4812 }
4813
4814 async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4815 Ok(())
4816 }
4817
4818 async fn cancel(&mut self) -> super::AdapterResult<bool> {
4819 Ok(false)
4820 }
4821
4822 async fn answer_permission(
4823 &mut self,
4824 _request_id: String,
4825 _answer: PermissionAnswer,
4826 ) -> super::AdapterResult<()> {
4827 Ok(())
4828 }
4829
4830 async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4831 Ok(())
4832 }
4833
4834 async fn reload(&mut self) -> super::AdapterResult<()> {
4835 Ok(())
4836 }
4837
4838 async fn stop(&mut self) -> super::AdapterResult<()> {
4839 self.stopped.fetch_add(1, Ordering::Relaxed);
4840 Ok(())
4841 }
4842
4843 async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4844 None
4845 }
4846 }
4847
4848 #[tokio::test]
4849 async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4850 let stopped = Arc::new(AtomicUsize::new(0));
4851 let relay = RelayHost::new(
4852 vec![
4853 AdapterHost::new(
4854 Box::new(StopTrackingAdapter {
4855 slot: 0,
4856 stopped: Arc::clone(&stopped),
4857 fail_stop: true,
4858 }),
4859 None,
4860 ),
4861 AdapterHost::new(
4862 Box::new(StopTrackingAdapter {
4863 slot: 1,
4864 stopped: Arc::clone(&stopped),
4865 fail_stop: false,
4866 }),
4867 None,
4868 ),
4869 ],
4870 4,
4871 )
4872 .expect("relay");
4873 let mut relay = relay;
4874
4875 let error = relay.stop().await.expect_err("first stop failure");
4876 assert!(error.to_string().contains("stop failed"));
4877 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4878 }
4879
4880 #[tokio::test]
4881 async fn relay_start_isolates_a_failed_adapter_and_keeps_healthy_peers() {
4882 let stopped = Arc::new(AtomicUsize::new(0));
4883 let mut relay = RelayHost::new(
4884 vec![
4885 AdapterHost::new(
4886 Box::new(StopTrackingAdapter {
4887 slot: 0,
4888 stopped: Arc::clone(&stopped),
4889 fail_stop: false,
4890 }),
4891 None,
4892 ),
4893 AdapterHost::new(
4894 Box::new(FailingStartAdapter {
4895 slot: 1,
4896 stopped: Arc::clone(&stopped),
4897 }),
4898 None,
4899 ),
4900 ],
4901 4,
4902 )
4903 .expect("relay");
4904
4905 relay.start().await.expect("healthy peer remains available");
4906 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0]);
4907 assert_eq!(stopped.load(Ordering::Relaxed), 1);
4908 relay.stop().await.unwrap();
4909 assert_eq!(stopped.load(Ordering::Relaxed), 2);
4910 }
4911
4912 #[test]
4913 fn parses_acp_text_without_ui_dependency() {
4914 let event = parse_acp_notification(
4915 2,
4916 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4917 )
4918 .expect("valid ACP")
4919 .expect("text event");
4920 assert_eq!(
4921 event,
4922 AgentEvent::Text {
4923 slot: 2,
4924 text: "hello".into(),
4925 }
4926 );
4927 }
4928
4929 #[test]
4930 fn parses_acp_state_notifications_at_the_adapter_boundary() {
4931 let commands = parse_acp_notification(
4932 3,
4933 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
4934 )
4935 .expect("valid ACP")
4936 .expect("commands event");
4937 assert_eq!(
4938 commands,
4939 AgentEvent::CommandsReplaced {
4940 slot: 3,
4941 commands: vec![crate::AgentCommand {
4942 name: "review".into()
4943 }]
4944 }
4945 );
4946
4947 let mode = parse_acp_notification(
4948 3,
4949 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
4950 )
4951 .expect("valid ACP")
4952 .expect("mode event");
4953 assert_eq!(
4954 mode,
4955 AgentEvent::ModeUpdated {
4956 slot: 3,
4957 current_mode: "review".into()
4958 }
4959 );
4960
4961 let usage = parse_acp_notification(
4962 3,
4963 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
4964 )
4965 .expect("valid ACP")
4966 .expect("usage event");
4967 assert_eq!(
4968 usage,
4969 AgentEvent::UsageUpdated {
4970 slot: 3,
4971 usage: crate::UsageUpdate {
4972 used: 4200,
4973 size: 128000
4974 }
4975 }
4976 );
4977
4978 let models = parse_acp_notification(
4979 3,
4980 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"}]}]}}}"#,
4981 )
4982 .expect("valid ACP")
4983 .expect("models event");
4984 assert!(matches!(
4985 models,
4986 AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
4987 if models.len() == 2 && current_model.as_deref() == Some("smart")
4988 ));
4989 assert_eq!(
4990 parse_model_config(&serde_json::json!({
4991 "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
4992 })),
4993 None
4994 );
4995
4996 let user = parse_acp_notification(
4997 3,
4998 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
4999 )
5000 .expect("valid ACP")
5001 .expect("user event");
5002 assert_eq!(
5003 user,
5004 AgentEvent::UserText {
5005 slot: 3,
5006 text: "context".into()
5007 }
5008 );
5009 }
5010
5011 #[test]
5012 fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
5013 let event = parse_acp_notification(
5014 0,
5015 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
5016 )
5017 .expect("valid ACP")
5018 .expect("mode event");
5019 assert!(matches!(
5020 event,
5021 AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
5022 if mode == "yolo" && modes[0].id == "yolo"
5023 ));
5024 }
5025
5026 #[test]
5027 fn parses_native_agy_text_without_acp_bridge() {
5028 let event = parse_agy_line(
5029 1,
5030 r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
5031 )
5032 .expect("valid stream-json")
5033 .expect("text event");
5034 assert_eq!(
5035 event,
5036 AgentEvent::Text {
5037 slot: 1,
5038 text: "hello".into(),
5039 }
5040 );
5041 }
5042
5043 #[test]
5044 fn parses_tool_lifecycle_from_each_protocol() {
5045 let agy = parse_agy_line(
5046 1,
5047 r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
5048 )
5049 .expect("valid native tool")
5050 .expect("tool event");
5051 assert!(matches!(
5052 agy,
5053 AgentEvent::Tool {
5054 update: crate::ToolUpdate {
5055 status: ToolStatus::Completed,
5056 ..
5057 },
5058 ..
5059 }
5060 ));
5061
5062 let acp = parse_acp_notification(
5063 1,
5064 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
5065 )
5066 .expect("valid ACP tool")
5067 .expect("tool event");
5068 assert!(matches!(
5069 acp,
5070 AgentEvent::Tool {
5071 update: crate::ToolUpdate {
5072 status: ToolStatus::Failed,
5073 ..
5074 },
5075 ..
5076 }
5077 ));
5078 }
5079
5080 #[test]
5081 fn parses_terminal_lifecycle_from_acp_and_native_events() {
5082 let created = parse_acp_notification(
5083 0,
5084 r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
5085 )
5086 .expect("valid ACP terminal")
5087 .expect("terminal event");
5088 assert_eq!(
5089 created,
5090 AgentEvent::Terminal {
5091 slot: 0,
5092 event: TerminalEvent::Created {
5093 id: "term-1".into(),
5094 command: "cargo test".into(),
5095 },
5096 }
5097 );
5098 let output = parse_agy_line(
5099 1,
5100 r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
5101 )
5102 .expect("valid native terminal")
5103 .expect("terminal event");
5104 assert_eq!(
5105 output,
5106 AgentEvent::Terminal {
5107 slot: 1,
5108 event: TerminalEvent::Output {
5109 id: "term-1".into(),
5110 text: "ok\n".into(),
5111 },
5112 }
5113 );
5114 let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
5115 .expect("valid native release")
5116 .expect("terminal event");
5117 assert!(matches!(
5118 released,
5119 AgentEvent::Terminal {
5120 event: TerminalEvent::Released { id },
5121 ..
5122 } if id == "term-1"
5123 ));
5124 }
5125
5126 #[test]
5127 fn parses_acp_permission_requests() {
5128 let event = parse_acp_notification(
5129 0,
5130 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"}]}}}"#,
5131 )
5132 .expect("valid permission")
5133 .expect("permission event");
5134 assert!(matches!(
5135 event,
5136 AgentEvent::Permission { request, .. }
5137 if request.id == "t1"
5138 && request.title == "Write file"
5139 && request.options == ["Allow once", "Reject"]
5140 && request.option_ids == ["allow-once", "reject"]
5141 ));
5142 }
5143
5144 #[test]
5145 fn parses_acp_permission_request_as_json_rpc_request() {
5146 let event = parse_acp_notification(
5147 2,
5148 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
5149 )
5150 .expect("valid permission request")
5151 .expect("permission event");
5152 assert!(matches!(
5153 event,
5154 AgentEvent::Permission { request, .. }
5155 if request.id == "17"
5156 && request.title == "Write file"
5157 && request.options == ["allow-once", "reject"]
5158 && request.option_ids == ["allow-once", "reject"]
5159 ));
5160 }
5161
5162 #[tokio::test]
5163 async fn native_adapter_explicitly_rejects_permission_answers() {
5164 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5165 assert_eq!(
5166 adapter
5167 .answer_permission(
5168 "request".into(),
5169 PermissionAnswer::Selected {
5170 option_id: "allow".into()
5171 },
5172 )
5173 .await,
5174 Err(super::AdapterError::Unsupported("permission answer"))
5175 );
5176 }
5177
5178 #[tokio::test]
5179 async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
5180 let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5181 adapter
5182 .set_mode("full-access".into())
5183 .await
5184 .expect("auto-pilot alias");
5185 assert!(matches!(
5186 adapter.next_event().await,
5187 Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
5188 ));
5189 }
5190
5191 #[tokio::test]
5192 async fn native_turns_receive_a_twenty_four_hour_timeout() {
5193 let script_path = unique_test_path("codeswarm-native-timeout", "sh");
5194 std::fs::write(
5195 &script_path,
5196 r#"#!/bin/sh
5197seen=0
5198while [ "$#" -gt 0 ]; do
5199 case "$1" in
5200 --print-timeout)
5201 shift
5202 [ "$1" = "1440m" ] || exit 2
5203 seen=$((seen + 1))
5204 ;;
5205 esac
5206 shift
5207done
5208[ "$seen" = 1 ] || exit 3
5209printf '%s\n' '{"event":"result","result":{"status":"SUCCESS","response":"timeout accepted"}}'
5210"#,
5211 )
5212 .unwrap();
5213 let mut adapter = AgyAdapter::with_session_id(
5214 0,
5215 std::env::current_dir().unwrap(),
5216 format!("sh {}", script_path.display()),
5217 "saved-session",
5218 );
5219 adapter.start().await.unwrap();
5220 adapter.next_event().await.unwrap().unwrap();
5221 adapter.next_event().await.unwrap().unwrap();
5222 for prompt in ["first task", "follow-up task"] {
5223 adapter.send_prompt(prompt.into()).await.unwrap();
5224 assert!(
5225 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "timeout accepted")
5226 );
5227 assert!(matches!(
5228 adapter.next_event().await,
5229 Some(Ok(AgentEvent::TurnComplete { .. }))
5230 ));
5231 }
5232 adapter.stop().await.unwrap();
5233 std::fs::remove_file(script_path).unwrap();
5234 }
5235
5236 #[tokio::test]
5237 async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
5238 let script_path = unique_test_path("codeswarm-agy-session", "sh");
5239 std::fs::write(
5240 &script_path,
5241 "#!/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",
5242 )
5243 .expect("write native test script");
5244 let mut adapter = AgyAdapter::new(
5245 0,
5246 std::env::current_dir().expect("cwd"),
5247 format!("sh {}", script_path.display()),
5248 );
5249 adapter.start().await.expect("start native adapter");
5250 assert!(adapter.next_event().await.is_some());
5252 assert!(adapter.next_event().await.is_some());
5253 adapter
5254 .send_prompt("first".into())
5255 .await
5256 .expect("first prompt");
5257 while !matches!(
5258 adapter.next_event().await,
5259 Some(Ok(AgentEvent::TurnComplete { .. }))
5260 ) {}
5261 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5262 adapter
5263 .send_prompt("follow up".into())
5264 .await
5265 .expect("follow-up prompt");
5266 while !matches!(
5267 adapter.next_event().await,
5268 Some(Ok(AgentEvent::TurnComplete { .. }))
5269 ) {}
5270 assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5271 adapter.stop().await.expect("stop native adapter");
5272 std::fs::remove_file(script_path).expect("cleanup native script");
5273 }
5274
5275 #[tokio::test]
5276 async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
5277 let script_path = unique_test_path("codeswarm-agy-failure", "sh");
5278 std::fs::write(
5279 &script_path,
5280 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
5281 )
5282 .expect("write native test script");
5283 let mut adapter = AgyAdapter::new(
5284 0,
5285 std::env::current_dir().expect("cwd"),
5286 format!("sh {}", script_path.display()),
5287 );
5288 adapter.start().await.expect("start native adapter");
5289 assert!(adapter.next_event().await.is_some());
5290 assert!(adapter.next_event().await.is_some());
5291 adapter.send_prompt("fail".into()).await.expect("prompt");
5292 assert!(matches!(
5293 adapter.next_event().await,
5294 Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
5295 if detail == "agent failed"
5296 ));
5297 adapter.stop().await.expect("stop native adapter");
5298 std::fs::remove_file(script_path).expect("cleanup native script");
5299 }
5300
5301 #[tokio::test]
5302 async fn native_crash_reaps_process_and_retries_on_next_prompt() {
5303 let script_path = unique_test_path("codeswarm-agy-retry", "sh");
5304 let marker_path = unique_test_path("codeswarm-agy-retry-marker", "txt");
5305 std::fs::write(
5306 &script_path,
5307 format!(
5308 "#!/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",
5309 marker_path.display(),
5310 marker_path.display(),
5311 marker_path.display(),
5312 ),
5313 )
5314 .expect("write retry script");
5315 let mut adapter = AgyAdapter::new(
5316 0,
5317 std::env::current_dir().expect("cwd"),
5318 format!("sh {}", script_path.display()),
5319 );
5320 adapter.start().await.expect("start native adapter");
5321 assert!(adapter.next_event().await.is_some());
5322 assert!(adapter.next_event().await.is_some());
5323 adapter
5324 .send_prompt("first".into())
5325 .await
5326 .expect("first prompt");
5327 assert!(matches!(
5328 adapter.next_event().await,
5329 Some(Ok(AgentEvent::Failed { detail, .. })) if detail == "first crash"
5330 ));
5331 adapter
5332 .send_prompt("retry".into())
5333 .await
5334 .expect("retry prompt starts a fresh process");
5335 assert!(matches!(
5336 adapter.next_event().await,
5337 Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5338 ));
5339 assert!(matches!(
5340 adapter.next_event().await,
5341 Some(Ok(AgentEvent::TurnComplete { .. }))
5342 ));
5343 adapter.stop().await.expect("stop native adapter");
5344 std::fs::remove_file(script_path).expect("cleanup retry script");
5345 std::fs::remove_file(marker_path).expect("cleanup retry marker");
5346 }
5347
5348 #[tokio::test]
5349 async fn acp_adapter_initializes_session_and_completes_a_prompt() {
5350 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"}}'"#;
5351 let cwd = std::env::current_dir().expect("cwd");
5352 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5353 adapter.start().await.expect("initialize");
5354 assert!(matches!(
5355 adapter.next_event().await,
5356 Some(Ok(AgentEvent::ModesReplaced { .. }))
5357 ));
5358 assert!(matches!(
5359 adapter.next_event().await,
5360 Some(Ok(AgentEvent::Ready { .. }))
5361 ));
5362 adapter.send_prompt("hello".into()).await.expect("prompt");
5363 assert!(matches!(
5364 adapter.next_event().await,
5365 Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
5366 ));
5367 assert!(matches!(
5368 adapter.next_event().await,
5369 Some(Ok(AgentEvent::TurnComplete { .. }))
5370 ));
5371 }
5372
5373 #[tokio::test]
5374 async fn acp_output_token_limit_is_reported_as_a_failed_turn() {
5375 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"}}'"#;
5376 let cwd = std::env::current_dir().expect("cwd");
5377 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5378 adapter.start().await.expect("initialize");
5379 assert!(matches!(
5380 adapter.next_event().await,
5381 Some(Ok(AgentEvent::Ready { .. }))
5382 ));
5383 adapter.send_prompt("hello".into()).await.expect("prompt");
5384 assert!(matches!(
5385 adapter.next_event().await,
5386 Some(Ok(AgentEvent::Failed {
5387 slot: 1,
5388 started: true,
5389 detail,
5390 })) if detail.contains("output token limit")
5391 ));
5392 }
5393
5394 #[tokio::test]
5395 async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
5396 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"}}'"#;
5397 let cwd = std::env::current_dir().expect("cwd");
5398 let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5399 adapter.start().await.expect("initialize");
5400 assert!(adapter.next_event().await.is_some());
5401 assert!(adapter.next_event().await.is_some());
5402
5403 for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
5404 adapter.send_prompt(prompt.into()).await.expect("prompt");
5405 assert!(matches!(
5406 adapter.next_event().await,
5407 Some(Ok(AgentEvent::Text { text, .. })) if text == expected
5408 ));
5409 assert!(matches!(
5410 adapter.next_event().await,
5411 Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
5412 ));
5413 }
5414 adapter.stop().await.expect("stop");
5415 }
5416
5417 #[tokio::test]
5418 async fn empty_acp_mode_catalog_disables_mode_control() {
5419 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":[]}}}'"#;
5420 let cwd = std::env::current_dir().expect("cwd");
5421 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5422 adapter.start().await.expect("initialize");
5423 assert!(!adapter.capabilities().supports_modes);
5424 assert!(matches!(
5425 adapter.next_event().await,
5426 Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
5427 ));
5428 assert!(matches!(
5429 adapter.next_event().await,
5430 Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
5431 ));
5432 adapter.stop().await.expect("stop");
5433 }
5434
5435 #[tokio::test]
5436 async fn acp_models_are_discovered_live_and_changed_through_session_config() {
5437 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"#;
5438 let cwd = std::env::current_dir().expect("cwd");
5439 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5440 adapter.start().await.expect("initialize");
5441 assert!(adapter.capabilities().supports_models);
5442 assert!(matches!(
5443 adapter.next_event().await,
5444 Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
5445 if config_id == "model"
5446 && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
5447 && current_model.as_deref() == Some("fast")
5448 ));
5449 assert!(matches!(
5450 adapter.next_event().await,
5451 Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
5452 ));
5453 adapter.set_model("smart".into()).await.expect("set model");
5454 assert!(adapter.set_model("invented".into()).await.is_err());
5455 adapter.stop().await.expect("stop");
5456 }
5457
5458 #[tokio::test]
5459 async fn acp_mode_change_is_acknowledged_without_provider_notification() {
5460 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":{}}'"#;
5461 let cwd = std::env::current_dir().expect("cwd");
5462 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5463 adapter.start().await.expect("initialize");
5464 adapter
5465 .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
5466 .await
5467 .expect("set mode");
5468 assert!(matches!(
5469 adapter.next_event().await,
5470 Some(Ok(AgentEvent::ModesReplaced { .. }))
5471 ));
5472 assert!(matches!(
5473 adapter.next_event().await,
5474 Some(Ok(AgentEvent::Ready { .. }))
5475 ));
5476 assert!(matches!(
5477 adapter.next_event().await,
5478 Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
5479 ));
5480 adapter.stop().await.expect("stop");
5481 }
5482
5483 #[tokio::test]
5484 async fn acp_reload_preserves_a_loadable_session_id() {
5485 let cwd = std::env::current_dir().expect("cwd");
5486 let mut adapter = AcpAdapter::with_session_id(
5487 0,
5488 cwd,
5489 "__codeswarm_missing_acp_for_reload_test__",
5490 Vec::new(),
5491 "saved-session",
5492 );
5493 adapter.capabilities.supports_session_load = true;
5494 assert!(adapter.reload().await.is_err());
5498 assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
5499 }
5500
5501 #[tokio::test]
5502 async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
5503 let cwd = std::env::current_dir().expect("cwd");
5504 let mut adapter = AcpAdapter::with_session_id(
5505 0,
5506 cwd,
5507 "__codeswarm_missing_nonloadable_acp__",
5508 Vec::new(),
5509 "stale-session",
5510 );
5511 adapter.capabilities.supports_session_load = false;
5512 assert!(adapter.reload().await.is_err());
5513 assert_eq!(adapter.session_id, None);
5514 }
5515
5516 #[tokio::test]
5517 async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
5518 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"}}'"#;
5519 let cwd = std::env::current_dir().expect("cwd");
5520 let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5521 adapter.start().await.expect("initialize");
5522 assert!(matches!(
5523 adapter.next_event().await,
5524 Some(Ok(AgentEvent::Ready { .. }))
5525 ));
5526 adapter.send_prompt("hello".into()).await.expect("prompt");
5527 assert!(matches!(
5528 adapter.next_event().await,
5529 Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
5530 ));
5531 assert!(matches!(
5532 adapter.next_event().await,
5533 Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
5534 ));
5535 }
5536
5537 #[test]
5538 fn acp_tool_patches_preserve_fields_and_honor_explicit_replacements() {
5539 let mut tools = std::collections::BTreeMap::new();
5540 let first = serde_json::json!({"sessionUpdate":"tool_call", "toolCallId":"read", "title":"Read config", "status":"in_progress",
5541 "content":[{"type":"content", "content":{"type":"text", "text":"old output"}}]});
5542 let initial = super::normalize_acp_tool(&first, &mut tools).unwrap();
5543 assert_eq!(initial.detail.as_deref(), Some("old output"));
5544 let completed = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call_update","toolCallId":"read","status":"completed"}), &mut tools).unwrap();
5545 assert_eq!(completed.title, "Read config");
5546 assert_eq!(completed.detail.as_deref(), Some("old output"));
5547 assert_eq!(completed.status, ToolStatus::Completed);
5548 let malformed = super::normalize_acp_tool(
5549 &serde_json::json!({"toolCallId":"read","title":3,"status":"unknown","content":null}),
5550 &mut tools,
5551 )
5552 .unwrap();
5553 assert_eq!(malformed, completed);
5554 let replaced = super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read","content":[false,{"type":"content","content":{"type":"text","text":"new output"}}]}), &mut tools).unwrap();
5555 assert_eq!(replaced.detail.as_deref(), Some("new output"));
5556 let cleared = super::normalize_acp_tool(
5557 &serde_json::json!({"toolCallId":"read","content":[]}),
5558 &mut tools,
5559 )
5560 .unwrap();
5561 assert_eq!(cleared.detail, None);
5562 let raw = super::normalize_acp_tool(
5563 &serde_json::json!({"toolCallId":"read","rawOutput":{"ok":true}}),
5564 &mut tools,
5565 )
5566 .unwrap();
5567 assert_eq!(raw.detail.as_deref(), Some("{\"ok\":true}"));
5568 let fresh = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call","toolCallId":"read","title":"New call"}), &mut tools).unwrap();
5569 assert_eq!(fresh.status, ToolStatus::Pending);
5570 assert_eq!(fresh.detail, None);
5571 for invalid in [
5572 serde_json::json!({}),
5573 serde_json::json!({"toolCallId":7}),
5574 serde_json::json!({"toolCallId":" "}),
5575 ] {
5576 assert!(super::normalize_acp_tool(&invalid, &mut tools).is_none());
5577 }
5578 assert_eq!(tools.len(), 1);
5579 super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read "}), &mut tools).unwrap();
5581 assert_eq!(tools.len(), 2);
5582 }
5583
5584 #[tokio::test]
5585 async fn acp_tool_status_only_notifications_retain_name_and_output() {
5586 let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'
5587read _; echo '{"id":2,"result":{"sessionId":"s"}}'
5588read _
5589echo '{"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"}}]}}}'
5590echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"r","status":"completed"}}}'
5591echo '{"id":3,"result":{"stopReason":"end_turn"}}'"#;
5592 let mut adapter = AcpAdapter::new(
5593 0,
5594 std::env::current_dir().unwrap(),
5595 "sh",
5596 vec!["-c".into(), script.into()],
5597 );
5598 adapter.start().await.unwrap();
5599 adapter.next_event().await.unwrap().unwrap();
5600 adapter.send_prompt("read".into()).await.unwrap();
5601 for status in [ToolStatus::Running, ToolStatus::Completed] {
5602 let Some(Ok(AgentEvent::Tool { update, .. })) = adapter.next_event().await else {
5603 panic!("tool event");
5604 };
5605 assert_eq!(update.status, status);
5606 assert_eq!(update.title, "Read config");
5607 assert_eq!(update.detail.as_deref(), Some("file content"));
5608 }
5609 assert!(matches!(
5610 adapter.next_event().await,
5611 Some(Ok(AgentEvent::TurnComplete { .. }))
5612 ));
5613 adapter.stop().await.unwrap();
5614 }
5615
5616 #[tokio::test]
5617 async fn acp_reload_discards_old_queued_events_and_catalogs() {
5618 let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"id":2,"result":{"sessionId":"new"}}'"#;
5619 let mut adapter = AcpAdapter::new(
5620 0,
5621 std::env::current_dir().unwrap(),
5622 "sh",
5623 vec!["-c".into(), script.into()],
5624 );
5625 adapter.start().await.unwrap();
5626 adapter.queued_events.push_back(Ok(AgentEvent::Text {
5627 slot: 0,
5628 text: "stale".into(),
5629 }));
5630 adapter.modes = vec![Mode {
5631 id: "stale".into(),
5632 label: "Stale".into(),
5633 }];
5634 adapter.next_request_id = 1;
5636 adapter.reload().await.unwrap();
5637 assert!(adapter.modes.is_empty());
5638 assert_eq!(adapter.queued_events.len(), 1);
5639 assert!(matches!(
5640 adapter.next_event().await,
5641 Some(Ok(AgentEvent::Ready { .. }))
5642 ));
5643 adapter.stop().await.unwrap();
5644 assert!(adapter.queued_events.is_empty());
5645 }
5646
5647 #[tokio::test]
5648 async fn acp_load_replays_history_without_starting_a_turn() {
5649 let script = r#"
5650read _
5651echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'
5652read request
5653case "$request" in *session/load*) ;; *) exit 2;; esac
5654echo '{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"text":"old question"}}}}'
5655echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"old answer"}}}}'
5656echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_thought_chunk","content":{"text":"old reasoning"}}}}'
5657echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call","toolCallId":"old-tool","title":"Read","status":"in_progress"}}}'
5658echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"old-tool","title":"Read","status":"completed"}}}'
5659echo '{"jsonrpc":"2.0","id":2,"result":{}}'
5660read request
5661case "$request" in *session/prompt*) ;; *) exit 3;; esac
5662echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"new answer"}}}}'
5663echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
5664"#;
5665 let mut adapter = AcpAdapter::with_session_id(
5666 2,
5667 std::env::current_dir().unwrap(),
5668 "sh",
5669 vec!["-c".into(), script.into()],
5670 "saved",
5671 );
5672 adapter.start().await.unwrap();
5673 let mut state = crate::SessionState::new(3);
5674 for _ in 0..5 {
5675 let event = adapter.next_event().await.unwrap().unwrap();
5676 assert!(matches!(&event, AgentEvent::History { slot: 2, .. }));
5677 crate::reduce(&mut state, event);
5678 assert_eq!(state.active_slot, None);
5679 }
5680 assert!(matches!(
5681 adapter.next_event().await,
5682 Some(Ok(AgentEvent::Ready { slot: 2, .. }))
5683 ));
5684 adapter.send_prompt("new question".into()).await.unwrap();
5685 assert!(
5686 matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "new answer")
5687 );
5688 assert!(matches!(
5689 adapter.next_event().await,
5690 Some(Ok(AgentEvent::TurnComplete { slot: 2 }))
5691 ));
5692 adapter.stop().await.unwrap();
5693 }
5694
5695 #[tokio::test]
5696 async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
5697 let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
5698 let cwd = std::env::current_dir().expect("cwd");
5699 let mut adapter = AcpAdapter::with_session_id(
5700 0,
5701 cwd,
5702 "sh",
5703 vec!["-c".into(), script.into()],
5704 "existing-session",
5705 );
5706 adapter.start().await.expect("load existing session");
5707 assert!(matches!(
5708 adapter.next_event().await,
5709 Some(Ok(AgentEvent::Ready { .. }))
5710 ));
5711 }
5712
5713 #[tokio::test]
5714 async fn acp_start_failure_reaps_transport_process() {
5715 let mut adapter = AcpAdapter::new(
5720 0,
5721 std::env::current_dir().expect("cwd"),
5722 "sh",
5723 vec!["-c".into(), "printf 'not-json\\n'".into()],
5724 );
5725 assert!(adapter.start().await.is_err());
5726 assert!(adapter.child.is_none());
5727 assert!(adapter.reader.is_none());
5728 }
5729
5730 #[tokio::test]
5731 async fn acp_transport_crash_is_reloaded_before_the_next_prompt() {
5732 let marker = unique_test_path("codeswarm-acp-retry", "count");
5733 let script = format!(
5734 r#"count=0
5735if [ -f '{0}' ]; then count=$(cat '{0}'); fi
5736count=$((count + 1))
5737printf '%s' "$count" > '{0}'
5738while IFS= read -r request; do
5739 id=$(printf '%s' "$request" | sed -n 's/.*"id":\([0-9][0-9]*\).*/\1/p')
5740 case "$request" in
5741 *initialize*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"agentCapabilities":{{"loadSession":true}}}}}}' ;;
5742 *session/new*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"sessionId":"saved-session"}}}}' ;;
5743 *session/load*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{}}}}' ;;
5744 *session/prompt*)
5745 if [ "$count" = 1 ]; then exit 0; fi
5746 printf '%s\n' '{{"jsonrpc":"2.0","method":"session/update","params":{{"update":{{"sessionUpdate":"agent_message_chunk","content":{{"text":"recovered"}}}}}}}}'
5747 printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"stopReason":"end_turn"}}}}'
5748 ;;
5749 esac
5750done
5751"#,
5752 marker.display()
5753 );
5754 let mut adapter = AcpAdapter::new(
5755 0,
5756 std::env::current_dir().expect("cwd"),
5757 "sh",
5758 vec!["-c".into(), script],
5759 );
5760 adapter.start().await.expect("initial ACP startup");
5761 assert!(matches!(
5762 adapter.next_event().await,
5763 Some(Ok(AgentEvent::Ready { .. }))
5764 ));
5765 adapter
5766 .send_prompt("first".into())
5767 .await
5768 .expect("first prompt");
5769 assert!(matches!(
5770 adapter.next_event().await,
5771 Some(Err(AdapterError::Transport(_)))
5772 ));
5773 assert!(adapter.child.is_none());
5774 assert!(adapter.reader.is_none());
5775 assert_eq!(adapter.session_id(), Some("saved-session".into()));
5776
5777 adapter.reload().await.expect("reload ACP transport");
5778 assert!(matches!(
5779 adapter.next_event().await,
5780 Some(Ok(AgentEvent::Ready { .. }))
5781 ));
5782 adapter
5783 .send_prompt("retry".into())
5784 .await
5785 .expect("retry prompt");
5786 assert!(matches!(
5787 adapter.next_event().await,
5788 Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5789 ));
5790 assert!(matches!(
5791 adapter.next_event().await,
5792 Some(Ok(AgentEvent::TurnComplete { .. }))
5793 ));
5794 adapter.stop().await.expect("stop ACP");
5795 std::fs::remove_file(marker).expect("cleanup marker");
5796 }
5797
5798 #[tokio::test]
5799 async fn coordinator_reload_replays_context_and_reintroduces_a_crashed_slot() {
5800 let prompts = Arc::new(Mutex::new(Vec::new()));
5801 let healthy = ScriptedAdapter::new(
5802 0,
5803 AgentCapabilities::default(),
5804 [
5805 AgentEvent::Text {
5806 slot: 0,
5807 text: "peer context".into(),
5808 },
5809 AgentEvent::TurnComplete { slot: 0 },
5810 ],
5811 );
5812 let probe = ReloadProbeAdapter {
5813 slot: 1,
5814 crashed: false,
5815 reloaded: false,
5816 events: VecDeque::new(),
5817 prompts: Arc::clone(&prompts),
5818 };
5819 let mut relay = RelayHost::new(
5820 vec![
5821 AdapterHost::new(Box::new(healthy), None),
5822 AdapterHost::new(Box::new(probe), None),
5823 ],
5824 8,
5825 )
5826 .expect("relay");
5827 relay.start().await.expect("start");
5828 relay
5829 .run_turn("original task", 0)
5830 .await
5831 .expect("first turn");
5832 relay.run_turn("", 0).await.expect("crashed turn");
5833 relay.relay_mut().enqueue_human("retry", Some(1));
5834 relay.run_turn("", 0).await.expect("reloaded turn");
5835
5836 let prompts = prompts.lock().expect("prompts");
5837 let retry = prompts.last().expect("retry prompt");
5838 assert!(retry.contains("You are Reload probe"), "{retry}");
5839 assert!(retry.contains("original task"), "{retry}");
5840 assert!(retry.contains("peer context"), "{retry}");
5841 assert!(retry.contains("retry"), "{retry}");
5842 }
5843
5844 #[tokio::test]
5845 async fn acp_adapter_answers_permission_json_rpc_requests() {
5846 let path = std::env::temp_dir().join(format!(
5847 "codeswarm-permission-answer-{}",
5848 std::process::id()
5849 ));
5850 let script = format!(
5851 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"}}}}'"#,
5852 path.display()
5853 );
5854 let mut adapter = AcpAdapter::new(
5855 0,
5856 std::env::current_dir().expect("cwd"),
5857 "sh",
5858 vec!["-c".into(), script],
5859 );
5860 adapter.start().await.expect("start ACP");
5861 assert!(matches!(
5862 adapter.next_event().await,
5863 Some(Ok(AgentEvent::Ready { .. }))
5864 ));
5865 adapter.send_prompt("do it".into()).await.expect("prompt");
5866 assert!(matches!(
5867 adapter.next_event().await,
5868 Some(Ok(AgentEvent::Permission { request, .. }))
5869 if request.id == "9"
5870 && request.options == ["Allow once"]
5871 && request.option_ids == ["allow-once"]
5872 ));
5873 adapter
5874 .answer_permission(
5875 "9".into(),
5876 PermissionAnswer::Selected {
5877 option_id: "allow-once".into(),
5878 },
5879 )
5880 .await
5881 .expect("permission answer");
5882 assert!(matches!(
5883 adapter.next_event().await,
5884 Some(Ok(AgentEvent::TurnComplete { .. }))
5885 ));
5886 let answer: Value = serde_json::from_str(
5887 &std::fs::read_to_string(&path).expect("captured permission answer"),
5888 )
5889 .expect("valid JSON-RPC answer");
5890 assert_eq!(answer["id"], 9);
5891 assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
5892 assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
5893 std::fs::remove_file(path).expect("cleanup");
5894 }
5895
5896 #[test]
5897 fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
5898 let event = parse_acp_notification(
5899 0,
5900 r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
5901 )
5902 .expect("valid JSON-RPC request");
5903 assert!(event.is_none());
5904 }
5905
5906 #[tokio::test]
5907 async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
5908 let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
5909 std::fs::write(
5910 &script_path,
5911 "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
5912 )
5913 .expect("write native test script");
5914 let mut adapter = AgyAdapter::new(
5915 0,
5916 std::env::current_dir().expect("cwd"),
5917 format!("sh {}", script_path.display()),
5918 );
5919 adapter.start().await.expect("start native adapter");
5920 assert!(adapter.next_event().await.is_some());
5921 assert!(adapter.next_event().await.is_some());
5922 adapter
5923 .send_prompt("continue".into())
5924 .await
5925 .expect("prompt");
5926 assert!(matches!(
5927 adapter.next_event().await,
5928 Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
5929 ));
5930 assert!(matches!(
5931 adapter.next_event().await,
5932 Some(Ok(AgentEvent::TurnComplete { .. }))
5933 ));
5934 adapter.stop().await.expect("stop native adapter");
5935 std::fs::remove_file(script_path).expect("cleanup native script");
5936 }
5937
5938 #[test]
5939 fn acp_workspace_file_access_is_root_bound_and_size_limited() {
5940 let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
5941 let _ = std::fs::remove_dir_all(&root);
5942 std::fs::create_dir_all(&root).expect("workspace");
5943 std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
5944 let outside =
5945 std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
5946 std::fs::write(&outside, "secret").expect("outside file");
5947 let link = root.join("outside-link");
5948 #[cfg(unix)]
5949 std::os::unix::fs::symlink(&outside, &link).expect("symlink");
5950 let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5951
5952 assert_eq!(
5953 adapter
5954 .read_workspace_text("inside.txt", Some(2), Some(1))
5955 .expect("read inside"),
5956 "two"
5957 );
5958 std::fs::write(
5959 root.join("large.txt"),
5960 vec![b'x'; MAX_FILE_READ_BYTES + 1024],
5961 )
5962 .expect("large file");
5963 let bounded = adapter
5964 .read_workspace_text("large.txt", None, None)
5965 .expect("bounded read");
5966 assert!(bounded.len() <= MAX_FILE_READ_BYTES);
5967 #[cfg(unix)]
5968 {
5969 std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
5970 .expect("internal symlink");
5971 assert_eq!(
5972 adapter
5973 .read_workspace_text("inside-link", None, None)
5974 .expect("read internal symlink"),
5975 "one\ntwo\nthree\n"
5976 );
5977 }
5978 assert!(adapter.workspace_path("../codeswarm-outside").is_err());
5979 assert!(
5980 adapter
5981 .workspace_path(&outside.display().to_string())
5982 .is_err()
5983 );
5984 #[cfg(unix)]
5985 assert!(adapter.workspace_path("outside-link").is_err());
5986 #[cfg(unix)]
5987 std::fs::remove_file(link).expect("cleanup symlink");
5988 #[cfg(unix)]
5989 std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
5990 std::fs::remove_file(outside).expect("cleanup outside");
5991 std::fs::remove_dir_all(root).expect("cleanup workspace");
5992 }
5993
5994 #[tokio::test]
5995 async fn running_terminal_output_omits_exit_status_until_completion() {
5996 let root = unique_test_path("codeswarm-terminal-output", "dir");
5997 std::fs::create_dir_all(&root).expect("workspace");
5998 let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5999 let result = adapter
6000 .terminal_create(&serde_json::json!({
6001 "command": "sh",
6002 "args": ["-c", "sleep 0.2; printf done"],
6003 "cwd": ".",
6004 }))
6005 .await
6006 .expect("terminal create");
6007 let id = result["terminalId"].as_str().expect("terminal id");
6008 let output = adapter.terminal_output(id).await.expect("terminal output");
6009 assert!(output.get("exitStatus").is_none());
6010 if let Some(terminal) = adapter.terminals.remove(id) {
6011 terminal.stop().await;
6012 }
6013 std::fs::remove_dir_all(root).expect("cleanup workspace");
6014 }
6015
6016 #[tokio::test]
6017 async fn acp_adapter_answers_workspace_read_requests() {
6018 let root =
6019 std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
6020 let _ = std::fs::remove_dir_all(&root);
6021 std::fs::create_dir_all(&root).expect("workspace");
6022 let source = root.join("inside.txt");
6023 let answer = root.join("answer.json");
6024 std::fs::write(&source, "workspace content").expect("source");
6025 let script = format!(
6026 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"}}}}'"#,
6027 source.display(),
6028 answer.display(),
6029 );
6030 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6031 adapter.start().await.expect("start ACP");
6032 assert!(matches!(
6033 adapter.next_event().await,
6034 Some(Ok(AgentEvent::Ready { .. }))
6035 ));
6036 adapter.send_prompt("read it".into()).await.expect("prompt");
6037 assert!(matches!(
6038 adapter.next_event().await,
6039 Some(Ok(AgentEvent::TurnComplete { .. }))
6040 ));
6041 let response: Value =
6042 serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
6043 .expect("response JSON");
6044 assert_eq!(response["id"], 9);
6045 assert_eq!(response["result"]["content"], "workspace content");
6046 adapter.stop().await.expect("stop ACP");
6047 std::fs::remove_dir_all(root).expect("cleanup workspace");
6048 }
6049
6050 #[tokio::test]
6051 async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
6052 let root =
6053 std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
6054 let _ = std::fs::remove_dir_all(&root);
6055 std::fs::create_dir_all(&root).expect("workspace");
6056 let create_request = serde_json::json!({
6057 "jsonrpc": "2.0",
6058 "id": 9,
6059 "method": "terminal/create",
6060 "params": {
6061 "sessionId": "s1",
6062 "command": "sh",
6063 "args": ["-c", "sleep 0.1; printf terminal-ok"],
6064 "cwd": ".",
6065 },
6066 });
6067 let wait_request = serde_json::json!({
6068 "jsonrpc": "2.0",
6069 "id": 10,
6070 "method": "terminal/wait_for_exit",
6071 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6072 });
6073 let output_request = serde_json::json!({
6074 "jsonrpc": "2.0",
6075 "id": 11,
6076 "method": "terminal/output",
6077 "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6078 });
6079 let create_answer = root.join("create-answer.json");
6080 let wait_answer = root.join("wait-answer.json");
6081 let output_answer = root.join("output-answer.json");
6082 let script = format!(
6083 "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\"}}}}'",
6084 create_request,
6085 create_answer.display(),
6086 wait_request,
6087 wait_answer.display(),
6088 output_request,
6089 output_answer.display(),
6090 );
6091 let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6092 adapter.start().await.expect("start ACP");
6093 assert!(matches!(
6094 adapter.next_event().await,
6095 Some(Ok(AgentEvent::Ready { .. }))
6096 ));
6097 adapter
6098 .send_prompt("run terminal".into())
6099 .await
6100 .expect("prompt");
6101 let mut saw_complete = false;
6102 for _ in 0..6 {
6103 match adapter.next_event().await {
6104 Some(Ok(AgentEvent::TurnComplete { .. })) => {
6105 saw_complete = true;
6106 break;
6107 }
6108 Some(_) => {}
6109 None => break,
6110 }
6111 }
6112 assert!(saw_complete, "terminal requests should not stall ACP");
6113 let create: Value = serde_json::from_str(
6114 &std::fs::read_to_string(&create_answer).expect("captured create response"),
6115 )
6116 .expect("create JSON");
6117 assert_eq!(create["result"]["terminalId"], "terminal-1");
6118 let output: Value = serde_json::from_str(
6119 &std::fs::read_to_string(&output_answer).expect("captured output response"),
6120 )
6121 .expect("output JSON");
6122 assert!(
6123 output["result"]["output"]
6124 .as_str()
6125 .unwrap_or_default()
6126 .contains("terminal-ok"),
6127 "output response: {output}"
6128 );
6129 adapter.stop().await.expect("stop ACP");
6130 std::fs::remove_dir_all(root).expect("cleanup workspace");
6131 }
6132
6133 #[tokio::test]
6134 async fn host_reduces_and_persists_adapter_events() {
6135 let path =
6136 std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
6137 let adapter = ScriptedAdapter::new(
6138 0,
6139 AgentCapabilities::default(),
6140 [AgentEvent::Text {
6141 slot: 0,
6142 text: "hello".into(),
6143 }],
6144 );
6145 let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
6146 host.start().await.expect("start");
6147 host.next_effects()
6148 .await
6149 .expect("event")
6150 .expect("valid event");
6151 assert_eq!(host.state.public_text[0].1, "hello");
6152 assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
6153 std::fs::remove_file(path).expect("cleanup");
6154 }
6155
6156 #[tokio::test]
6157 async fn relay_applies_default_policy_before_the_first_prompt() {
6158 let first_log = Arc::new(Mutex::new(Vec::new()));
6159 let second_log = Arc::new(Mutex::new(Vec::new()));
6160 let first = AdapterHost::new(
6161 Box::new(ModeOrderAdapter {
6162 slot: 0,
6163 log: Arc::clone(&first_log),
6164 phase: 0,
6165 }),
6166 None,
6167 );
6168 let second = AdapterHost::new(
6169 Box::new(ModeOrderAdapter {
6170 slot: 1,
6171 log: Arc::clone(&second_log),
6172 phase: 0,
6173 }),
6174 None,
6175 );
6176 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6177 relay.start().await.expect("start and synchronize policy");
6178 relay.run_turn("task", 0).await.expect("first turn");
6179 {
6180 let log = first_log.lock().expect("log");
6181 assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
6182 }
6183 assert_eq!(
6184 second_log.lock().expect("log").as_slice(),
6185 ["start", "mode:yolo"]
6186 );
6187 let added_log = Arc::new(Mutex::new(Vec::new()));
6188 relay
6189 .add_agent(
6190 AdapterHost::new(
6191 Box::new(ModeOrderAdapter {
6192 slot: 2,
6193 log: Arc::clone(&added_log),
6194 phase: 0,
6195 }),
6196 None,
6197 ),
6198 "Added",
6199 "added.example",
6200 "added-agent",
6201 )
6202 .await
6203 .expect("add with synchronized policy");
6204 assert_eq!(
6205 added_log.lock().expect("log").as_slice(),
6206 ["start", "mode:yolo"]
6207 );
6208 relay.drop_agent(2).await.expect("drop added agent");
6209 added_log.lock().expect("log").clear();
6210 relay
6211 .reload(2)
6212 .await
6213 .expect("reload with synchronized policy");
6214 assert_eq!(
6215 added_log.lock().expect("log").as_slice(),
6216 ["reload", "mode:yolo"]
6217 );
6218 }
6219
6220 #[tokio::test]
6221 async fn acp_roster_is_ready_before_any_prompt_is_sent() {
6222 let hosts = (0..2)
6223 .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
6224 .collect::<Vec<_>>();
6225 let startup_events = Arc::new(Mutex::new(Vec::new()));
6226 let captured = Arc::clone(&startup_events);
6227 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6228 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6229
6230 relay.start().await.expect("complete startup handshake");
6231
6232 assert!(relay.dispatches().is_empty());
6233 let ready_slots = startup_events
6234 .lock()
6235 .expect("events")
6236 .iter()
6237 .filter_map(|event| match event {
6238 AgentEvent::Ready { slot, .. } => Some(*slot),
6239 _ => None,
6240 })
6241 .collect::<Vec<_>>();
6242 assert_eq!(ready_slots, vec![0, 1]);
6243 }
6244
6245 #[tokio::test]
6246 async fn independent_roster_adapters_start_concurrently() {
6247 let barrier = Arc::new(tokio::sync::Barrier::new(2));
6248 let hosts = (0..2)
6249 .map(|slot| {
6250 AdapterHost::new(
6251 Box::new(ConcurrentStartAdapter {
6252 slot,
6253 barrier: Arc::clone(&barrier),
6254 }),
6255 None,
6256 )
6257 })
6258 .collect::<Vec<_>>();
6259 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6260 tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
6261 .await
6262 .expect("startup should not serialize barrier participants")
6263 .expect("startup succeeds");
6264 }
6265
6266 #[tokio::test]
6267 async fn relay_host_dispatches_turns_sequentially() {
6268 let capabilities = AgentCapabilities {
6269 supports_cancel: true,
6270 ..AgentCapabilities::default()
6271 };
6272 let first = ScriptedAdapter::new(
6273 0,
6274 capabilities.clone(),
6275 [
6276 AgentEvent::Text {
6277 slot: 0,
6278 text: "first".into(),
6279 },
6280 AgentEvent::TurnComplete { slot: 0 },
6281 ],
6282 );
6283 let second = ScriptedAdapter::new(
6284 1,
6285 capabilities,
6286 [
6287 AgentEvent::Text {
6288 slot: 1,
6289 text: "review".into(),
6290 },
6291 AgentEvent::TurnComplete { slot: 1 },
6292 ],
6293 );
6294 let hosts = vec![
6295 AdapterHost::new(Box::new(first), None),
6296 AdapterHost::new(Box::new(second), None),
6297 ];
6298 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6299 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6300 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6301 let captured = std::sync::Arc::clone(&events);
6302 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6303 relay.start().await.expect("start");
6304 events.lock().expect("events").clear();
6305 assert!(matches!(
6306 relay.run_turn("task", 0).await.expect("first turn"),
6307 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6308 ));
6309 assert!(matches!(
6310 relay.run_turn("first", 0).await.expect("second turn"),
6311 crate::relay::RelayDecision::Dispatch {
6312 slot: 1,
6313 can_stop: true,
6314 ..
6315 }
6316 ));
6317 assert_eq!(
6318 relay
6319 .dispatches()
6320 .iter()
6321 .map(|(slot, _)| *slot)
6322 .collect::<Vec<_>>(),
6323 [0, 1]
6324 );
6325 assert!(relay.dispatches()[0].1.contains("You are Claude"));
6326 assert!(
6327 relay.dispatches()[0]
6328 .1
6329 .contains("CodeSwarm roster (ordered)")
6330 );
6331 assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
6332 assert!(relay.dispatches()[0].1.contains("2. Codex"));
6333 assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
6334 assert!(
6335 relay.dispatches()[1]
6336 .1
6337 .contains("stops all other agents and ends the entire automated relay")
6338 );
6339 assert!(relay.dispatches()[1].1.contains("Use it with extreme care"));
6340 assert!(
6341 relay.dispatches()[1]
6342 .1
6343 .contains("If there is any uncertainty, do not use it")
6344 );
6345 assert!(relay.dispatches()[0].1.contains("Do not use"));
6346 let lifecycle = events.lock().expect("events");
6347 let positions = lifecycle
6348 .iter()
6349 .filter_map(|event| match event {
6350 AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
6351 AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
6352 _ => None,
6353 })
6354 .collect::<Vec<_>>();
6355 assert_eq!(
6356 positions,
6357 [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
6358 );
6359 }
6360
6361 #[tokio::test]
6362 async fn failed_resume_does_not_stop_or_dispatch_to_healthy_peer() {
6363 let stops = Arc::new(AtomicUsize::new(0));
6364 let failed = FailingStartAdapter {
6365 slot: 0,
6366 stopped: stops.clone(),
6367 };
6368 let healthy = ScriptedAdapter::new(
6369 1,
6370 AgentCapabilities::default(),
6371 [
6372 AgentEvent::Text {
6373 slot: 1,
6374 text: "healthy response".into(),
6375 },
6376 AgentEvent::TurnComplete { slot: 1 },
6377 ],
6378 );
6379 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6380 let captured = events.clone();
6381 let mut relay = RelayHost::new(
6382 vec![
6383 AdapterHost::new(Box::new(failed), None),
6384 AdapterHost::new(Box::new(healthy), None),
6385 ],
6386 4,
6387 )
6388 .unwrap();
6389 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6390 relay.start_resuming().await.unwrap();
6391 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
6392 assert!(relay.dispatches().is_empty());
6393 assert_eq!(stops.load(Ordering::Relaxed), 1);
6394 assert!(
6395 events
6396 .lock()
6397 .unwrap()
6398 .iter()
6399 .any(|event| matches!(event, AgentEvent::Failed { slot: 0, .. }))
6400 );
6401 assert!(
6402 !events
6403 .lock()
6404 .unwrap()
6405 .iter()
6406 .any(|event| matches!(event, AgentEvent::Failed { slot: 1, .. }))
6407 );
6408 assert!(!relay.relay_mut().enqueue_human("do not retarget", Some(0)));
6409 assert!(
6410 relay
6411 .relay_mut()
6412 .enqueue_human("explicit healthy target", Some(1))
6413 );
6414 assert!(matches!(
6415 relay.run_turn("", 1).await.unwrap(),
6416 RelayDecision::Dispatch { slot: 1, .. }
6417 ));
6418 relay.stop().await.unwrap();
6419 }
6420
6421 #[tokio::test]
6422 async fn pair_strategy_wires_roles_into_non_direct_prompts() {
6423 let capabilities = AgentCapabilities::default();
6424 let first = ScriptedAdapter::new(
6425 0,
6426 capabilities.clone(),
6427 [
6428 AgentEvent::Text {
6429 slot: 0,
6430 text: "implemented".into(),
6431 },
6432 AgentEvent::TurnComplete { slot: 0 },
6433 AgentEvent::Text {
6434 slot: 0,
6435 text: format!("fixed review findings {STOP_TOKEN}"),
6436 },
6437 AgentEvent::TurnComplete { slot: 0 },
6438 ],
6439 );
6440 let second = ScriptedAdapter::new(
6441 1,
6442 capabilities,
6443 [
6444 AgentEvent::Text {
6445 slot: 1,
6446 text: "reviewed".into(),
6447 },
6448 AgentEvent::TurnComplete { slot: 1 },
6449 AgentEvent::Text {
6450 slot: 1,
6451 text: format!("approved {STOP_TOKEN}"),
6452 },
6453 AgentEvent::TurnComplete { slot: 1 },
6454 AgentEvent::Text {
6455 slot: 1,
6456 text: "new task".into(),
6457 },
6458 AgentEvent::TurnComplete { slot: 1 },
6459 ],
6460 );
6461 let hosts = vec![
6462 AdapterHost::new(Box::new(first), None),
6463 AdapterHost::new(Box::new(second), None),
6464 ];
6465 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6466 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6467 relay.relay_mut().set_strategy(CollaborationStrategy::Pair);
6468 relay.start().await.expect("start");
6469 assert!(matches!(
6470 relay.run_turn("task", 0).await.expect("implementer turn"),
6471 RelayDecision::Dispatch {
6472 slot: 0,
6473 can_stop: false,
6474 ..
6475 }
6476 ));
6477 let implementer_prompt = &relay.dispatches()[0].1;
6478 assert!(implementer_prompt.contains("you are the implementer"));
6479 assert!(implementer_prompt.contains("pair reviewer will review the result next"));
6480 assert!(!implementer_prompt.contains("you are the reviewer"));
6481 assert!(implementer_prompt.contains("Do not use"));
6482 assert!(matches!(
6483 relay.run_turn("", 0).await.expect("reviewer turn"),
6484 RelayDecision::Dispatch {
6485 slot: 1,
6486 can_stop: true,
6487 ..
6488 }
6489 ));
6490 let reviewer_prompt = &relay.dispatches()[1].1;
6491 assert!(reviewer_prompt.contains("you are the reviewer"));
6492 assert!(reviewer_prompt.contains("Claude handed off"));
6493 assert!(reviewer_prompt.contains("concrete defects"));
6494 assert!(reviewer_prompt.contains("concise approval"));
6495 assert!(reviewer_prompt.contains(STOP_TOKEN));
6496 assert!(!reviewer_prompt.contains("you are the implementer"));
6497 relay.run_turn("", 0).await.unwrap();
6498 assert!(relay.dispatches()[2].1.contains("you are the implementer"));
6499 assert!(relay.dispatches()[2].1.contains("Do not use"));
6500 assert!(matches!(
6501 relay.run_turn("", 0).await.unwrap(),
6502 RelayDecision::Dispatch { slot: 1, .. }
6503 ));
6504 assert!(relay.dispatches()[3].1.contains("you are the reviewer"));
6505 assert!(relay.relay_mut().enqueue_human("new task", Some(1)));
6506 relay.run_turn("", 1).await.unwrap();
6507 assert!(relay.dispatches()[4].1.contains("you are the implementer"));
6508 }
6509
6510 #[tokio::test]
6511 async fn solo_roster_and_direct_prompts_omit_pair_roles() {
6512 let solo = ScriptedAdapter::new(
6513 0,
6514 AgentCapabilities::default(),
6515 [
6516 AgentEvent::Text {
6517 slot: 0,
6518 text: "solo".into(),
6519 },
6520 AgentEvent::TurnComplete { slot: 0 },
6521 ],
6522 );
6523 let mut solo_relay =
6524 RelayHost::new(vec![AdapterHost::new(Box::new(solo), None)], 4).expect("relay");
6525 solo_relay
6526 .relay_mut()
6527 .set_strategy(CollaborationStrategy::Pair);
6528 solo_relay.start().await.expect("start");
6529 solo_relay.run_turn("task", 0).await.expect("solo turn");
6530 assert!(!solo_relay.dispatches()[0].1.contains("Pair role"));
6531
6532 let roster_first = ScriptedAdapter::new(
6533 0,
6534 AgentCapabilities::default(),
6535 [AgentEvent::TurnComplete { slot: 0 }],
6536 );
6537 let roster_second = ScriptedAdapter::new(
6538 1,
6539 AgentCapabilities::default(),
6540 [AgentEvent::TurnComplete { slot: 1 }],
6541 );
6542 let mut roster = RelayHost::new(
6543 vec![
6544 AdapterHost::new(Box::new(roster_first), None),
6545 AdapterHost::new(Box::new(roster_second), None),
6546 ],
6547 4,
6548 )
6549 .expect("relay");
6550 roster.start().await.expect("start");
6551 roster.run_turn("task", 0).await.expect("first turn");
6552 roster.run_turn("", 0).await.expect("second turn");
6553 assert!(!roster.dispatches()[0].1.contains("Pair role"));
6554 assert!(!roster.dispatches()[1].1.contains("Pair role"));
6555
6556 let pair_first = ScriptedAdapter::new(
6557 0,
6558 AgentCapabilities::default(),
6559 [AgentEvent::TurnComplete { slot: 0 }],
6560 );
6561 let pair_second = ScriptedAdapter::new(
6562 1,
6563 AgentCapabilities::default(),
6564 [AgentEvent::TurnComplete { slot: 1 }],
6565 );
6566 let mut pair = RelayHost::new(
6567 vec![
6568 AdapterHost::new(Box::new(pair_first), None),
6569 AdapterHost::new(Box::new(pair_second), None),
6570 ],
6571 4,
6572 )
6573 .expect("relay");
6574 pair.relay_mut().set_strategy(CollaborationStrategy::Pair);
6575 assert_eq!(pair.relay_mut().enqueue_direct(1, "private"), Ok(true));
6576 pair.start().await.expect("start");
6577 assert!(matches!(
6578 pair.run_turn("ignored", 0).await.expect("direct turn"),
6579 RelayDecision::Dispatch {
6580 slot: 1,
6581 direct: true,
6582 ..
6583 }
6584 ));
6585 let direct_prompt = &pair.dispatches()[0].1;
6586 assert!(direct_prompt.contains("private"));
6587 assert!(!direct_prompt.contains("Pair role"));
6588 }
6589
6590 #[tokio::test]
6591 async fn relay_host_routes_around_a_usage_limited_agent() {
6592 let capabilities = AgentCapabilities::default();
6593 let first = ScriptedAdapter::new(
6594 0,
6595 capabilities.clone(),
6596 [
6597 AgentEvent::Text {
6598 slot: 0,
6599 text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
6600 credits or try again later."
6601 .into(),
6602 },
6603 AgentEvent::TurnComplete { slot: 0 },
6604 ],
6605 );
6606 let second = ScriptedAdapter::new(
6607 1,
6608 capabilities,
6609 [
6610 AgentEvent::Text {
6611 slot: 1,
6612 text: "review done".into(),
6613 },
6614 AgentEvent::TurnComplete { slot: 1 },
6615 ],
6616 );
6617 let hosts = vec![
6618 AdapterHost::new(Box::new(first), None),
6619 AdapterHost::new(Box::new(second), None),
6620 ];
6621 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6622 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6623 let captured = std::sync::Arc::clone(&events);
6624 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6625 relay.start().await.expect("start");
6626 events.lock().expect("events").clear();
6627 assert!(matches!(
6628 relay.run_turn("task", 0).await.expect("limited turn"),
6629 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6630 ));
6631 assert!(
6632 events
6633 .lock()
6634 .expect("events")
6635 .iter()
6636 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6637 );
6638 assert!(matches!(
6640 relay.run_turn("", 0).await.expect("next turn"),
6641 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6642 ));
6643 assert!(relay.relay().is_limited(0));
6644 relay.reload(0).await.expect("reload");
6648 assert!(!relay.relay().is_limited(0));
6649 }
6650
6651 #[tokio::test]
6652 async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
6653 let limited = ScriptedAdapter::new(
6654 0,
6655 AgentCapabilities::default(),
6656 [AgentEvent::Failed {
6657 slot: 0,
6658 started: true,
6659 detail: "request failed: insufficient_quota".into(),
6660 }],
6661 );
6662 let healthy = ScriptedAdapter::new(
6663 1,
6664 AgentCapabilities::default(),
6665 [AgentEvent::TurnComplete { slot: 1 }],
6666 );
6667 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6668 let captured = std::sync::Arc::clone(&events);
6669 let mut relay = RelayHost::new(
6670 vec![
6671 AdapterHost::new(Box::new(limited), None),
6672 AdapterHost::new(Box::new(healthy), None),
6673 ],
6674 4,
6675 )
6676 .expect("relay");
6677 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6678 relay.start().await.expect("start");
6679 events.lock().expect("events").clear();
6680
6681 assert!(matches!(
6682 relay.run_turn("task", 0).await.expect("limited failure"),
6683 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6684 ));
6685 assert!(relay.relay().is_limited(0));
6686 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
6687 {
6688 let events = events.lock().expect("events");
6689 assert!(
6690 events
6691 .iter()
6692 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6693 );
6694 assert!(
6695 !events
6696 .iter()
6697 .any(|event| matches!(event, AgentEvent::Failed { .. }))
6698 );
6699 }
6700
6701 assert!(matches!(
6702 relay.run_turn("", 0).await.expect("healthy peer"),
6703 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6704 ));
6705 }
6706
6707 #[tokio::test]
6708 async fn relay_failure_is_skipped_for_one_batch_without_changing_the_roster() {
6709 let failed = ScriptedAdapter::new(
6710 0,
6711 AgentCapabilities::default(),
6712 [AgentEvent::Failed {
6713 slot: 0,
6714 started: true,
6715 detail: "connection lost".into(),
6716 }],
6717 );
6718 let healthy = ScriptedAdapter::new(
6719 1,
6720 AgentCapabilities::default(),
6721 [AgentEvent::TurnComplete { slot: 1 }],
6722 );
6723 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6724 let captured = std::sync::Arc::clone(&events);
6725 let mut relay = RelayHost::new(
6726 vec![
6727 AdapterHost::new(Box::new(failed), None),
6728 AdapterHost::new(Box::new(healthy), None),
6729 ],
6730 4,
6731 )
6732 .expect("relay");
6733 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6734 relay.start().await.expect("start");
6735
6736 assert!(matches!(
6737 relay.run_turn("task", 0).await.expect("handled failure"),
6738 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6739 ));
6740 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![0, 1]);
6741 assert!(relay.relay().is_limited(0));
6742 assert!(events.lock().expect("lock").iter().any(|event| {
6743 matches!(
6744 event,
6745 AgentEvent::Failed {
6746 slot: 0,
6747 started: true,
6748 ..
6749 }
6750 )
6751 }));
6752 assert!(matches!(
6753 relay.run_turn("", 0).await.expect("healthy peer"),
6754 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6755 ));
6756 }
6757
6758 #[tokio::test]
6759 async fn codex_stop_does_not_skip_later_roster_reviewers() {
6760 let hosts = (0..3)
6761 .map(|slot| {
6762 AdapterHost::new(
6763 Box::new(ScriptedAdapter::new(
6764 slot,
6765 AgentCapabilities::default(),
6766 [
6767 AgentEvent::Text {
6768 slot,
6769 text: STOP_TOKEN.into(),
6770 },
6771 AgentEvent::TurnComplete { slot },
6772 ],
6773 )),
6774 None,
6775 )
6776 })
6777 .collect();
6778 let mut relay = RelayHost::new(hosts, 10).expect("relay");
6779 relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
6780 relay.start().await.expect("start");
6781 for expected in 0..3 {
6782 assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
6783 RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
6784 }
6785 assert_eq!(
6786 relay.run_turn("", 0).await.expect("complete"),
6787 RelayDecision::Complete
6788 );
6789 }
6790
6791 #[tokio::test]
6792 async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
6793 let first = ScriptedAdapter::new(
6794 0,
6795 AgentCapabilities::default(),
6796 [
6797 AgentEvent::Text {
6798 slot: 0,
6799 text: "done".into(),
6800 },
6801 AgentEvent::TurnComplete { slot: 0 },
6802 ],
6803 );
6804 let reviewer = ScriptedAdapter::new(
6805 1,
6806 AgentCapabilities::default(),
6807 [
6808 AgentEvent::Text {
6809 slot: 1,
6810 text: STOP_TOKEN.into(),
6811 },
6812 AgentEvent::TurnComplete { slot: 1 },
6813 ],
6814 );
6815 let mut relay = RelayHost::new(
6816 vec![
6817 AdapterHost::new(Box::new(first), None),
6818 AdapterHost::new(Box::new(reviewer), None),
6819 ],
6820 10,
6821 )
6822 .expect("relay");
6823 relay.start().await.expect("start");
6824 let first_decision = relay.run_turn("task", 0).await.expect("first");
6825 assert!(matches!(
6826 first_decision,
6827 RelayDecision::Dispatch { slot: 0, .. }
6828 ));
6829 let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
6830 assert!(matches!(
6831 reviewer_decision,
6832 RelayDecision::Dispatch {
6833 slot: 1,
6834 can_stop: true,
6835 ..
6836 }
6837 ));
6838 assert_eq!(
6839 relay.run_turn("", 0).await.expect("complete"),
6840 RelayDecision::Complete
6841 );
6842 }
6843
6844 #[tokio::test]
6845 async fn relay_stream_emits_text_and_thought_endings_before_tools() {
6846 let tool = AgentEvent::Tool {
6847 slot: 0,
6848 update: crate::ToolUpdate {
6849 id: "read".into(),
6850 title: "Read file".into(),
6851 status: ToolStatus::Running,
6852 detail: None,
6853 },
6854 };
6855 let updates = vec![
6856 AgentEvent::Thought {
6857 slot: 0,
6858 text: "Check the buffer. ✈".into(),
6859 },
6860 AgentEvent::Text {
6861 slot: 0,
6862 text: "Let me check.".into(),
6863 },
6864 tool.clone(),
6865 AgentEvent::Text {
6866 slot: 0,
6867 text: "[CODE".into(),
6868 },
6869 AgentEvent::Text {
6870 slot: 0,
6871 text: " is ordinary.".into(),
6872 },
6873 AgentEvent::Text {
6874 slot: 0,
6875 text: "[CODESWARM:".into(),
6876 },
6877 AgentEvent::Text {
6878 slot: 0,
6879 text: "STOP] Done.".into(),
6880 },
6881 tool,
6882 AgentEvent::TurnComplete { slot: 0 },
6883 ];
6884 let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
6885 let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
6886 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6887 let captured = std::sync::Arc::clone(&events);
6888 let mut relay = RelayHost::new(
6889 vec![
6890 AdapterHost::new(Box::new(first), None),
6891 AdapterHost::new(Box::new(reviewer), None),
6892 ],
6893 2,
6894 )
6895 .expect("relay");
6896 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6897 relay.start().await.unwrap();
6898 relay.run_turn("task", 0).await.unwrap();
6899 let captured = events.lock().unwrap();
6900 let visible: Vec<_> = captured
6901 .iter()
6902 .filter(|event| {
6903 matches!(
6904 event,
6905 AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
6906 )
6907 })
6908 .cloned()
6909 .collect();
6910 assert_eq!(
6911 visible,
6912 vec![
6913 updates[0].clone(),
6914 updates[1].clone(),
6915 updates[2].clone(),
6916 AgentEvent::Text {
6917 slot: 0,
6918 text: "[CODE is ordinary.".into()
6919 },
6920 AgentEvent::Text {
6921 slot: 0,
6922 text: " Done.".into()
6923 },
6924 updates[7].clone(),
6925 ]
6926 );
6927 }
6928
6929 #[tokio::test]
6930 async fn roster_handoff_routes_only_terminal_message_markers_and_refreshes_targets() {
6931 let text = |value: &str| AgentEvent::Text {
6932 slot: 0,
6933 text: value.into(),
6934 };
6935 let thought = || AgentEvent::Thought {
6936 slot: 0,
6937 text: "still checking".into(),
6938 };
6939 let tool = || AgentEvent::Tool {
6940 slot: 0,
6941 update: crate::ToolUpdate {
6942 id: "read".into(),
6943 title: "Read file".into(),
6944 status: ToolStatus::Running,
6945 detail: None,
6946 },
6947 };
6948 let cases = vec![
6949 (vec![text("result [CODESWARM:NEXT:3]\n ")], 2),
6950 (
6951 vec![text("result [CODESWARM:"), text("NEXT:"), text("3]")],
6952 2,
6953 ),
6954 (vec![text("result [CODESWARM:NEXT:3]"), text(" more")], 1),
6955 (vec![text("result [CODESWARM:NEXT:3]"), thought()], 1),
6956 (
6957 vec![text("result [CODESWARM:NEXT:3]"), tool(), text(" ")],
6958 1,
6959 ),
6960 (
6961 vec![text("result [CODESWARM:NEXT:"), thought(), text("3]")],
6962 1,
6963 ),
6964 (
6965 vec![text("result"), thought(), text("[CODESWARM:NEXT:3]")],
6966 2,
6967 ),
6968 (vec![text("result [CODESWARM:NEXT:1]")], 1),
6969 (vec![text("result [CODESWARM:NEXT:0]")], 1),
6970 (vec![text("result [CODESWARM:NEXT:99]")], 1),
6971 (
6972 vec![
6973 text("result [CODESWARM:NEXT:3]"),
6974 AgentEvent::UsageUpdated {
6975 slot: 0,
6976 usage: crate::UsageUpdate { used: 1, size: 100 },
6977 },
6978 ],
6979 2,
6980 ),
6981 ];
6982 for (mut updates, expected) in cases {
6983 updates.push(AgentEvent::TurnComplete { slot: 0 });
6984 let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates);
6985 let hosts = std::iter::once(AdapterHost::new(Box::new(first), None))
6986 .chain((1..3).map(|slot| {
6987 AdapterHost::new(
6988 Box::new(ScriptedAdapter::new(
6989 slot,
6990 AgentCapabilities::default(),
6991 [AgentEvent::TurnComplete { slot }],
6992 )),
6993 None,
6994 )
6995 }))
6996 .collect();
6997 let mut relay = RelayHost::new(hosts, 10).unwrap();
6998 relay.set_roster_names(vec!["Worker".into(), "Codex".into(), "Codex".into()]);
6999 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
7000 let captured = events.clone();
7001 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
7002 relay.start().await.unwrap();
7003 relay.run_turn("task", 0).await.unwrap();
7004 let prompt = &relay.dispatches()[0].1;
7005 assert!(prompt.contains("[CODESWARM:NEXT:2] → Codex"));
7006 assert!(prompt.contains("[CODESWARM:NEXT:3] → Codex"));
7007 assert!(!prompt.contains("[CODESWARM:NEXT:1]"));
7008 relay.introduced.fill(true);
7010 relay.set_roster_names(vec!["Replacement".into(), "Codex".into(), "Codex".into()]);
7011 let next = relay.run_turn("", 0).await.unwrap();
7012 assert!(
7013 matches!(next, RelayDecision::Dispatch { slot, can_stop: false, .. } if slot == expected),
7014 "{next:?}"
7015 );
7016 let prompt = &relay.dispatches()[1].1;
7017 assert!(prompt.contains("[CODESWARM:NEXT:1] → Replacement"));
7018 assert!(prompt.contains("result"));
7019 let public = prompt
7020 .split("Public updates:\n")
7021 .nth(1)
7022 .unwrap()
7023 .split("\n\nDo not use")
7024 .next()
7025 .unwrap();
7026 assert!(!public.contains("[CODESWARM:NEXT:"));
7027 let visible = events
7028 .lock()
7029 .unwrap()
7030 .iter()
7031 .filter_map(|event| match event {
7032 AgentEvent::Text { text, .. } => Some(text.clone()),
7033 _ => None,
7034 })
7035 .collect::<String>();
7036 assert!(visible.contains("result"));
7037 assert!(!visible.contains("[CODESWARM:"), "{visible}");
7038 }
7039 }
7040
7041 #[tokio::test]
7042 async fn reviewer_stop_requires_a_terminal_marker_after_all_activity() {
7043 let text = |value: &str| AgentEvent::Text {
7044 slot: 1,
7045 text: value.into(),
7046 };
7047 let thought = || AgentEvent::Thought {
7048 slot: 1,
7049 text: "still checking".into(),
7050 };
7051 let tool = || AgentEvent::Tool {
7052 slot: 1,
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(&format!("done {STOP_TOKEN}"))], true),
7062 (vec![text(STOP_TOKEN), text("\n ")], true),
7063 (vec![text(STOP_TOKEN), text(" actually keep going")], false),
7064 (vec![text(STOP_TOKEN), thought()], false),
7065 (vec![text(STOP_TOKEN), tool()], false),
7066 (vec![text(STOP_TOKEN), tool(), text(" ")], false),
7067 (vec![text(STOP_TOKEN), tool(), text(STOP_TOKEN)], true),
7068 (vec![text("[CODESWARM:"), text("STOP]")], true),
7069 (vec![text("[CODESWARM:"), thought(), text("STOP]")], false),
7070 (
7071 vec![AgentEvent::Thought {
7072 slot: 1,
7073 text: STOP_TOKEN.into(),
7074 }],
7075 false,
7076 ),
7077 (
7078 vec![
7079 text(STOP_TOKEN),
7080 AgentEvent::UsageUpdated {
7081 slot: 1,
7082 usage: crate::UsageUpdate { used: 1, size: 100 },
7083 },
7084 ],
7085 true,
7086 ),
7087 ];
7088 for (mut events, stop) in cases {
7089 let first = ScriptedAdapter::new(
7090 0,
7091 AgentCapabilities::default(),
7092 [
7093 AgentEvent::Text {
7094 slot: 0,
7095 text: "initial response".into(),
7096 },
7097 AgentEvent::TurnComplete { slot: 0 },
7098 AgentEvent::TurnComplete { slot: 0 },
7099 ],
7100 );
7101 events.push(AgentEvent::TurnComplete { slot: 1 });
7102 let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), events.clone());
7103 let mut relay = RelayHost::new(
7104 vec![
7105 AdapterHost::new(Box::new(first), None),
7106 AdapterHost::new(Box::new(reviewer), None),
7107 ],
7108 4,
7109 )
7110 .unwrap();
7111 relay.start().await.unwrap();
7112 relay.run_turn("task", 0).await.unwrap();
7113 relay.run_turn("", 0).await.unwrap();
7114 let next = relay.run_turn("", 0).await.unwrap();
7115 assert_eq!(
7116 matches!(next, RelayDecision::Complete),
7117 stop,
7118 "events={events:?}"
7119 );
7120 relay.stop().await.unwrap();
7121 }
7122 }
7123
7124 #[tokio::test]
7125 async fn stop_token_is_filtered_from_streamed_ui_events() {
7126 let first = ScriptedAdapter::new(
7127 0,
7128 AgentCapabilities::default(),
7129 [
7130 AgentEvent::Text {
7131 slot: 0,
7132 text: format!("visible {STOP_TOKEN} trailing"),
7133 },
7134 AgentEvent::TurnComplete { slot: 0 },
7135 ],
7136 );
7137 let reviewer = ScriptedAdapter::new(
7138 1,
7139 AgentCapabilities::default(),
7140 [AgentEvent::TurnComplete { slot: 1 }],
7141 );
7142 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7143 let captured = std::sync::Arc::clone(&events);
7144 let mut relay = RelayHost::new(
7145 vec![
7146 AdapterHost::new(Box::new(first), None),
7147 AdapterHost::new(Box::new(reviewer), None),
7148 ],
7149 2,
7150 )
7151 .expect("relay");
7152 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7153 relay.start().await.expect("start");
7154 relay.run_turn("task", 0).await.expect("turn");
7155 let captured = events.lock().expect("lock");
7156 assert!(captured.iter().all(|event| match event {
7157 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7158 _ => true,
7159 }));
7160 let visible = captured
7161 .iter()
7162 .filter_map(|event| match event {
7163 AgentEvent::Text { text, .. } => Some(text.as_str()),
7164 _ => None,
7165 })
7166 .collect::<String>();
7167 assert_eq!(visible, "visible trailing");
7168 }
7169
7170 #[tokio::test]
7171 async fn token_only_reviewer_response_emits_visible_acknowledgment() {
7172 let first = ScriptedAdapter::new(
7173 0,
7174 AgentCapabilities::default(),
7175 [
7176 AgentEvent::Text {
7177 slot: 0,
7178 text: "done".into(),
7179 },
7180 AgentEvent::TurnComplete { slot: 0 },
7181 ],
7182 );
7183 let reviewer = ScriptedAdapter::new(
7184 1,
7185 AgentCapabilities::default(),
7186 [
7187 AgentEvent::Text {
7188 slot: 1,
7189 text: STOP_TOKEN.into(),
7190 },
7191 AgentEvent::TurnComplete { slot: 1 },
7192 ],
7193 );
7194 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7195 let captured = std::sync::Arc::clone(&events);
7196 let mut relay = RelayHost::new(
7197 vec![
7198 AdapterHost::new(Box::new(first), None),
7199 AdapterHost::new(Box::new(reviewer), None),
7200 ],
7201 4,
7202 )
7203 .expect("relay");
7204 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7205 relay.start().await.expect("start");
7206 relay.run_turn("task", 0).await.expect("first turn");
7207 relay.run_turn("", 0).await.expect("review turn");
7208 let captured = events.lock().expect("lock");
7209 assert!(captured.iter().any(|event| {
7210 matches!(
7211 event,
7212 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7213 )
7214 }));
7215 assert!(captured.iter().all(|event| match event {
7216 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7217 _ => true,
7218 }));
7219 let acknowledgment = captured
7220 .iter()
7221 .position(|event| {
7222 matches!(
7223 event,
7224 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7225 )
7226 })
7227 .expect("visible acknowledgment");
7228 let completion = captured
7229 .iter()
7230 .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
7231 .expect("reviewer completion");
7232 assert!(acknowledgment < completion);
7233 }
7234
7235 #[tokio::test]
7236 async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
7237 let first = ScriptedAdapter::new(
7238 0,
7239 AgentCapabilities::default(),
7240 [
7241 AgentEvent::Text {
7242 slot: 0,
7243 text: "done".into(),
7244 },
7245 AgentEvent::TurnComplete { slot: 0 },
7246 ],
7247 );
7248 let reviewer = ScriptedAdapter::new(
7249 1,
7250 AgentCapabilities::default(),
7251 [
7252 AgentEvent::Text {
7253 slot: 1,
7254 text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
7255 },
7256 AgentEvent::TurnComplete { slot: 1 },
7257 ],
7258 );
7259 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7260 let captured = std::sync::Arc::clone(&events);
7261 let mut relay = RelayHost::new(
7262 vec![
7263 AdapterHost::new(Box::new(first), None),
7264 AdapterHost::new(Box::new(reviewer), None),
7265 ],
7266 4,
7267 )
7268 .expect("relay");
7269 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7270 relay.start().await.expect("start");
7271 relay.run_turn("task", 0).await.expect("first turn");
7272 relay.run_turn("", 0).await.expect("review turn");
7273
7274 let visible = events
7275 .lock()
7276 .expect("lock")
7277 .iter()
7278 .filter_map(|event| match event {
7279 AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
7280 _ => None,
7281 })
7282 .collect::<String>();
7283 assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
7284 assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
7285 }
7286
7287 #[tokio::test]
7288 async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
7289 let first = AdapterHost::new(
7290 Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
7291 None,
7292 );
7293 let second = AdapterHost::new(
7294 Box::new(ScriptedAdapter::new(
7295 1,
7296 AgentCapabilities::default(),
7297 [AgentEvent::TurnComplete { slot: 1 }],
7298 )),
7299 None,
7300 );
7301 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7302 let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
7303 relay.set_event_sink(move |event| {
7304 if matches!(event, AgentEvent::Permission { .. }) {
7305 let _ = seen_sender.send(());
7306 }
7307 });
7308 relay.start().await.expect("start");
7309 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
7310 let answer = async move {
7311 seen_receiver.recv().await.expect("permission request");
7312 sender
7313 .send(super::RelayPermissionAnswer {
7314 slot: 0,
7315 request_id: "permission-1".into(),
7316 answer: PermissionAnswer::Selected {
7317 option_id: "allow".into(),
7318 },
7319 })
7320 .expect("queue permission answer");
7321 };
7322 tokio::time::timeout(std::time::Duration::from_millis(100), async {
7323 let ((), result) = tokio::join!(
7324 answer,
7325 relay.run_turn_with_permissions("task", 0, &mut receiver)
7326 );
7327 result
7328 })
7329 .await
7330 .expect("permission-gated turn should not deadlock")
7331 .expect("turn completes");
7332 }
7333
7334 #[tokio::test]
7335 async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
7336 let first = AdapterHost::new(
7337 Box::new(PendingAdapter {
7338 slot: 0,
7339 hang_on_cancel: false,
7340 }),
7341 None,
7342 );
7343 let second = AdapterHost::new(
7344 Box::new(ScriptedAdapter::new(
7345 1,
7346 AgentCapabilities::default(),
7347 [AgentEvent::TurnComplete { slot: 1 }],
7348 )),
7349 None,
7350 );
7351 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7352 relay.start().await.expect("start");
7353 let cancellation = relay.cancellation();
7354 let error = {
7355 let turn = relay.run_turn("task", 0);
7356 tokio::pin!(turn);
7357 cancellation.request();
7358 turn.await.expect_err("cancellation should stop turn")
7359 };
7360 assert!(error.to_string().contains("relay turn cancelled"));
7361
7362 assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
7363 relay
7364 .run_turn("", 1)
7365 .await
7366 .expect("replacement job reaches the selected peer");
7367 let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
7368 assert!(replacement.contains("replacement job"));
7369 assert!(replacement.contains("User "));
7370 assert!(replacement.contains(":\ntask"));
7371 let owner_updates = relay.relay_mut().unseen_context(0);
7372 assert!(owner_updates.contains("User "));
7373 assert!(owner_updates.contains(":\ntask"));
7374 assert!(owner_updates.contains(":\nreplacement job"));
7375 }
7376
7377 #[tokio::test]
7378 async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
7379 let first = AdapterHost::new(
7380 Box::new(PendingAdapter {
7381 slot: 0,
7382 hang_on_cancel: true,
7383 }),
7384 None,
7385 );
7386 let second = AdapterHost::new(
7387 Box::new(PendingAdapter {
7388 slot: 1,
7389 hang_on_cancel: false,
7390 }),
7391 None,
7392 );
7393 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7394 relay.start().await.expect("start");
7395 let cancellation = relay.cancellation();
7396 let turn = relay.run_turn("task", 0);
7397 tokio::pin!(turn);
7398 cancellation.request();
7399 let error = turn.await.expect_err("cancellation should stop turn");
7400 assert!(error.to_string().contains("timed out"));
7401 }
7402
7403 #[tokio::test]
7404 async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
7405 let event = [AgentEvent::TurnComplete { slot: 0 }];
7406 let first = AdapterHost::new(
7407 Box::new(ScriptedAdapter::new(
7408 0,
7409 AgentCapabilities::default(),
7410 event.clone(),
7411 )),
7412 None,
7413 );
7414 let second = AdapterHost::new(
7415 Box::new(ScriptedAdapter::new(
7416 1,
7417 AgentCapabilities::default(),
7418 [AgentEvent::TurnComplete { slot: 1 }],
7419 )),
7420 None,
7421 );
7422 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7423 relay.start().await.expect("start");
7424
7425 relay.pause();
7426 assert_eq!(
7427 relay.run_turn("paused", 0).await.expect("paused turn"),
7428 crate::relay::RelayDecision::Paused
7429 );
7430 assert!(relay.dispatches().is_empty());
7431
7432 relay.resume();
7433 relay.relay_mut().drop_agent(1).expect("drop reviewer");
7434 assert!(matches!(
7435 relay
7436 .run_turn("solo follow-up", 0)
7437 .await
7438 .expect("solo turn"),
7439 crate::relay::RelayDecision::Dispatch {
7440 slot: 0,
7441 can_stop: false,
7442 ..
7443 }
7444 ));
7445 assert_eq!(relay.dispatches().len(), 1);
7446 }
7447
7448 #[tokio::test]
7449 async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
7450 let first = AdapterHost::new(
7451 Box::new(ScriptedAdapter::new(
7452 0,
7453 AgentCapabilities::default(),
7454 [AgentEvent::TurnComplete { slot: 0 }],
7455 )),
7456 None,
7457 );
7458 let second = AdapterHost::new(
7459 Box::new(ScriptedAdapter::new(
7460 1,
7461 AgentCapabilities::default(),
7462 [AgentEvent::TurnComplete { slot: 1 }],
7463 )),
7464 None,
7465 );
7466 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7467 relay.set_roster_names(vec!["First".into(), "Second".into()]);
7468 relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
7469 relay.set_roster_launch_specs(vec![
7470 ("custom".into(), "owner".into()),
7471 ("custom".into(), "peer".into()),
7472 ]);
7473 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7474 let captured = std::sync::Arc::clone(&events);
7475 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7476 relay.start().await.expect("start");
7477 let slot = relay
7478 .add_agent(
7479 AdapterHost::new(
7480 Box::new(ScriptedAdapter::new(
7481 2,
7482 AgentCapabilities::default(),
7483 [AgentEvent::TurnComplete { slot: 2 }],
7484 )),
7485 None,
7486 ),
7487 "Reviewer",
7488 "reviewer.example",
7489 "reviewer --acp",
7490 )
7491 .await
7492 .expect("append agent");
7493 assert_eq!(slot, 2);
7494 assert_eq!(
7495 relay.relay().active_slots().collect::<Vec<_>>(),
7496 vec![0, 1, 2]
7497 );
7498 assert_eq!(
7499 relay
7500 .session_metadata()
7501 .get("agents")
7502 .and_then(|value| value.as_array())
7503 .map(Vec::len),
7504 Some(3)
7505 );
7506 relay.drop_agent(1).await.expect("drop middle peer");
7507 let metadata = relay.session_metadata();
7508 assert_eq!(
7509 metadata.get("agents"),
7510 Some(&serde_json::json!([
7511 {"slot": 0, "name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
7512 {"slot": 2, "name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
7513 ]))
7514 );
7515 assert!(
7516 events
7517 .lock()
7518 .expect("lock")
7519 .iter()
7520 .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
7521 );
7522 }
7523
7524 #[tokio::test]
7525 async fn relay_host_persists_coordinator_owned_runtime_metadata() {
7526 let path = unique_test_path("codeswarm-session-metadata", "json");
7527 let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
7528 let writer = metadata_store.buffered().expect("metadata writer");
7529 let first = AdapterHost::new(
7530 Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
7531 None,
7532 );
7533 let second = AdapterHost::new(
7534 Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
7535 None,
7536 );
7537 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7538 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7539 relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
7540 relay.set_roster_launch_specs(vec![
7541 ("custom".into(), "claude".into()),
7542 ("custom".into(), "codex".into()),
7543 ]);
7544 relay.set_session_metadata_writer(writer);
7545 relay.start().await.expect("start");
7546 relay.drop_agent(0).await.expect("drop first agent");
7547 relay.stop().await.expect("stop");
7548
7549 let loaded = metadata_store
7550 .read()
7551 .expect("read metadata")
7552 .expect("metadata snapshot");
7553 assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
7554 assert_eq!(
7555 loaded.get("agents"),
7556 Some(&serde_json::json!([{
7557 "slot": 1, "name": "Codex", "identity": "openai.com", "protocol": "custom",
7558 "command": "codex", "supports_load_session": false
7559 }]))
7560 );
7561 assert!(loaded.get("owner").is_none());
7562 let _ = std::fs::remove_file(path);
7563 }
7564
7565 #[tokio::test]
7566 async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
7567 let first = AdapterHost::new(
7568 Box::new(ScriptedAdapter::new(
7569 0,
7570 AgentCapabilities::default(),
7571 [
7572 AgentEvent::Text {
7573 slot: 0,
7574 text: "owner stream".into(),
7575 },
7576 AgentEvent::TurnComplete { slot: 0 },
7577 ],
7578 )),
7579 None,
7580 );
7581 let second = AdapterHost::new(
7582 Box::new(ScriptedAdapter::new(
7583 1,
7584 AgentCapabilities::default(),
7585 [
7586 AgentEvent::Text {
7587 slot: 1,
7588 text: "peer stream".into(),
7589 },
7590 AgentEvent::TurnComplete { slot: 1 },
7591 ],
7592 )),
7593 None,
7594 );
7595 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7596 relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
7597 relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
7598 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7599 let captured = std::sync::Arc::clone(&events);
7600 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7601 relay.start().await.expect("start");
7602
7603 relay.swap_agents(0, 1).expect("swap peers");
7604 assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
7605 assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
7606 relay.run_turn("task", 0).await.expect("swapped turn");
7607 let events = events.lock().expect("events");
7608 assert!(events.iter().any(|event| {
7609 matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
7610 }));
7611 assert!(relay.dispatches()[0].1.contains("You are Peer"));
7612 }
7613
7614 #[tokio::test]
7615 async fn relay_host_persists_all_active_agent_metadata_off_thread() {
7616 let path = unique_test_path("codeswarm-session-metadata", "json");
7617 let first = AdapterHost::new(
7618 Box::new(ScriptedAdapter::new(
7619 0,
7620 AgentCapabilities::default(),
7621 [AgentEvent::TurnComplete { slot: 0 }],
7622 )),
7623 None,
7624 );
7625 let second = AdapterHost::new(
7626 Box::new(ScriptedAdapter::new(
7627 1,
7628 AgentCapabilities::default(),
7629 [AgentEvent::TurnComplete { slot: 1 }],
7630 )),
7631 None,
7632 );
7633 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7634 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7635 relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
7636 relay.set_roster_launch_specs(vec![
7637 ("custom".into(), "claude".into()),
7638 ("custom".into(), "codex".into()),
7639 ]);
7640 let writer = SessionMetadataStore::open(&path)
7641 .buffered()
7642 .expect("metadata writer");
7643 relay.set_session_metadata_writer(writer);
7644 relay.start().await.expect("start");
7645 relay.stop().await.expect("stop");
7646 let loaded = SessionMetadataStore::open(&path)
7647 .read()
7648 .expect("read metadata")
7649 .expect("metadata snapshot");
7650 let agents = loaded
7651 .get("agents")
7652 .and_then(|value| value.as_array())
7653 .expect("agents");
7654 assert_eq!(agents.len(), 2);
7655 assert_eq!(agents[0]["identity"], "claude.com");
7656 assert_eq!(agents[1]["identity"], "openai.com");
7657 let _ = std::fs::remove_file(path);
7658 }
7659
7660 #[tokio::test]
7661 async fn relay_host_routes_unseen_public_context_to_next_agent() {
7662 let first = AdapterHost::new(
7663 Box::new(ScriptedAdapter::new(
7664 0,
7665 AgentCapabilities::default(),
7666 [
7667 AgentEvent::Text {
7668 slot: 0,
7669 text: "implemented the fix".into(),
7670 },
7671 AgentEvent::TurnComplete { slot: 0 },
7672 ],
7673 )),
7674 None,
7675 );
7676 let second = AdapterHost::new(
7677 Box::new(ScriptedAdapter::new(
7678 1,
7679 AgentCapabilities::default(),
7680 [AgentEvent::TurnComplete { slot: 1 }],
7681 )),
7682 None,
7683 );
7684 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7685 relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
7686 relay.start().await.expect("start");
7687 relay.run_turn("task", 0).await.expect("first turn");
7688 relay.run_turn("review this", 0).await.expect("review turn");
7689
7690 assert_eq!(relay.dispatches().len(), 2);
7691 assert_eq!(relay.dispatches()[0].0, 0);
7692 assert!(relay.dispatches()[0].1.contains("task"));
7693 assert!(relay.dispatches()[0].1.contains("You are Codex"));
7694 assert!(relay.dispatches()[0].1.contains("2. Qwen"));
7695 assert_eq!(relay.dispatches()[1].0, 1);
7696 assert!(relay.dispatches()[1].1.contains("review this"));
7697 let public = relay.dispatches()[1]
7698 .1
7699 .split_once("Public updates:\n")
7700 .map(|(_, updates)| updates)
7701 .expect("review receives public context");
7702 let header = public
7703 .lines()
7704 .find(|line| line.starts_with("Codex "))
7705 .expect("named previous agent");
7706 let timestamp = header
7707 .strip_prefix("Codex ")
7708 .and_then(|value| value.strip_suffix(':'))
7709 .expect("timestamped header");
7710 assert_eq!(timestamp.len(), 5);
7711 assert_eq!(timestamp.as_bytes()[2], b':');
7712 assert!(
7713 timestamp
7714 .bytes()
7715 .enumerate()
7716 .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
7717 );
7718 assert!(public.contains("implemented the fix"));
7719 assert!(!public.contains("Agent 0"));
7720 }
7721}