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. If no meaningful correction is needed,\nend 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!(relay.dispatches()[0].1.contains("Do not use"));
6335 let lifecycle = events.lock().expect("events");
6336 let positions = lifecycle
6337 .iter()
6338 .filter_map(|event| match event {
6339 AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
6340 AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
6341 _ => None,
6342 })
6343 .collect::<Vec<_>>();
6344 assert_eq!(
6345 positions,
6346 [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
6347 );
6348 }
6349
6350 #[tokio::test]
6351 async fn failed_resume_does_not_stop_or_dispatch_to_healthy_peer() {
6352 let stops = Arc::new(AtomicUsize::new(0));
6353 let failed = FailingStartAdapter {
6354 slot: 0,
6355 stopped: stops.clone(),
6356 };
6357 let healthy = ScriptedAdapter::new(
6358 1,
6359 AgentCapabilities::default(),
6360 [
6361 AgentEvent::Text {
6362 slot: 1,
6363 text: "healthy response".into(),
6364 },
6365 AgentEvent::TurnComplete { slot: 1 },
6366 ],
6367 );
6368 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6369 let captured = events.clone();
6370 let mut relay = RelayHost::new(
6371 vec![
6372 AdapterHost::new(Box::new(failed), None),
6373 AdapterHost::new(Box::new(healthy), None),
6374 ],
6375 4,
6376 )
6377 .unwrap();
6378 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6379 relay.start_resuming().await.unwrap();
6380 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
6381 assert!(relay.dispatches().is_empty());
6382 assert_eq!(stops.load(Ordering::Relaxed), 1);
6383 assert!(
6384 events
6385 .lock()
6386 .unwrap()
6387 .iter()
6388 .any(|event| matches!(event, AgentEvent::Failed { slot: 0, .. }))
6389 );
6390 assert!(
6391 !events
6392 .lock()
6393 .unwrap()
6394 .iter()
6395 .any(|event| matches!(event, AgentEvent::Failed { slot: 1, .. }))
6396 );
6397 assert!(!relay.relay_mut().enqueue_human("do not retarget", Some(0)));
6398 assert!(
6399 relay
6400 .relay_mut()
6401 .enqueue_human("explicit healthy target", Some(1))
6402 );
6403 assert!(matches!(
6404 relay.run_turn("", 1).await.unwrap(),
6405 RelayDecision::Dispatch { slot: 1, .. }
6406 ));
6407 relay.stop().await.unwrap();
6408 }
6409
6410 #[tokio::test]
6411 async fn pair_strategy_wires_roles_into_non_direct_prompts() {
6412 let capabilities = AgentCapabilities::default();
6413 let first = ScriptedAdapter::new(
6414 0,
6415 capabilities.clone(),
6416 [
6417 AgentEvent::Text {
6418 slot: 0,
6419 text: "implemented".into(),
6420 },
6421 AgentEvent::TurnComplete { slot: 0 },
6422 AgentEvent::Text {
6423 slot: 0,
6424 text: format!("fixed review findings {STOP_TOKEN}"),
6425 },
6426 AgentEvent::TurnComplete { slot: 0 },
6427 ],
6428 );
6429 let second = ScriptedAdapter::new(
6430 1,
6431 capabilities,
6432 [
6433 AgentEvent::Text {
6434 slot: 1,
6435 text: "reviewed".into(),
6436 },
6437 AgentEvent::TurnComplete { slot: 1 },
6438 AgentEvent::Text {
6439 slot: 1,
6440 text: format!("approved {STOP_TOKEN}"),
6441 },
6442 AgentEvent::TurnComplete { slot: 1 },
6443 AgentEvent::Text {
6444 slot: 1,
6445 text: "new task".into(),
6446 },
6447 AgentEvent::TurnComplete { slot: 1 },
6448 ],
6449 );
6450 let hosts = vec![
6451 AdapterHost::new(Box::new(first), None),
6452 AdapterHost::new(Box::new(second), None),
6453 ];
6454 let mut relay = RelayHost::new(hosts, 4).expect("relay");
6455 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6456 relay.relay_mut().set_strategy(CollaborationStrategy::Pair);
6457 relay.start().await.expect("start");
6458 assert!(matches!(
6459 relay.run_turn("task", 0).await.expect("implementer turn"),
6460 RelayDecision::Dispatch {
6461 slot: 0,
6462 can_stop: false,
6463 ..
6464 }
6465 ));
6466 let implementer_prompt = &relay.dispatches()[0].1;
6467 assert!(implementer_prompt.contains("you are the implementer"));
6468 assert!(implementer_prompt.contains("pair reviewer will review the result next"));
6469 assert!(!implementer_prompt.contains("you are the reviewer"));
6470 assert!(implementer_prompt.contains("Do not use"));
6471 assert!(matches!(
6472 relay.run_turn("", 0).await.expect("reviewer turn"),
6473 RelayDecision::Dispatch {
6474 slot: 1,
6475 can_stop: true,
6476 ..
6477 }
6478 ));
6479 let reviewer_prompt = &relay.dispatches()[1].1;
6480 assert!(reviewer_prompt.contains("you are the reviewer"));
6481 assert!(reviewer_prompt.contains("Claude handed off"));
6482 assert!(reviewer_prompt.contains("concrete defects"));
6483 assert!(reviewer_prompt.contains("concise approval"));
6484 assert!(reviewer_prompt.contains(STOP_TOKEN));
6485 assert!(!reviewer_prompt.contains("you are the implementer"));
6486 relay.run_turn("", 0).await.unwrap();
6487 assert!(relay.dispatches()[2].1.contains("you are the implementer"));
6488 assert!(relay.dispatches()[2].1.contains("Do not use"));
6489 assert!(matches!(
6490 relay.run_turn("", 0).await.unwrap(),
6491 RelayDecision::Dispatch { slot: 1, .. }
6492 ));
6493 assert!(relay.dispatches()[3].1.contains("you are the reviewer"));
6494 assert!(relay.relay_mut().enqueue_human("new task", Some(1)));
6495 relay.run_turn("", 1).await.unwrap();
6496 assert!(relay.dispatches()[4].1.contains("you are the implementer"));
6497 }
6498
6499 #[tokio::test]
6500 async fn solo_roster_and_direct_prompts_omit_pair_roles() {
6501 let solo = ScriptedAdapter::new(
6502 0,
6503 AgentCapabilities::default(),
6504 [
6505 AgentEvent::Text {
6506 slot: 0,
6507 text: "solo".into(),
6508 },
6509 AgentEvent::TurnComplete { slot: 0 },
6510 ],
6511 );
6512 let mut solo_relay =
6513 RelayHost::new(vec![AdapterHost::new(Box::new(solo), None)], 4).expect("relay");
6514 solo_relay
6515 .relay_mut()
6516 .set_strategy(CollaborationStrategy::Pair);
6517 solo_relay.start().await.expect("start");
6518 solo_relay.run_turn("task", 0).await.expect("solo turn");
6519 assert!(!solo_relay.dispatches()[0].1.contains("Pair role"));
6520
6521 let roster_first = ScriptedAdapter::new(
6522 0,
6523 AgentCapabilities::default(),
6524 [AgentEvent::TurnComplete { slot: 0 }],
6525 );
6526 let roster_second = ScriptedAdapter::new(
6527 1,
6528 AgentCapabilities::default(),
6529 [AgentEvent::TurnComplete { slot: 1 }],
6530 );
6531 let mut roster = RelayHost::new(
6532 vec![
6533 AdapterHost::new(Box::new(roster_first), None),
6534 AdapterHost::new(Box::new(roster_second), None),
6535 ],
6536 4,
6537 )
6538 .expect("relay");
6539 roster.start().await.expect("start");
6540 roster.run_turn("task", 0).await.expect("first turn");
6541 roster.run_turn("", 0).await.expect("second turn");
6542 assert!(!roster.dispatches()[0].1.contains("Pair role"));
6543 assert!(!roster.dispatches()[1].1.contains("Pair role"));
6544
6545 let pair_first = ScriptedAdapter::new(
6546 0,
6547 AgentCapabilities::default(),
6548 [AgentEvent::TurnComplete { slot: 0 }],
6549 );
6550 let pair_second = ScriptedAdapter::new(
6551 1,
6552 AgentCapabilities::default(),
6553 [AgentEvent::TurnComplete { slot: 1 }],
6554 );
6555 let mut pair = RelayHost::new(
6556 vec![
6557 AdapterHost::new(Box::new(pair_first), None),
6558 AdapterHost::new(Box::new(pair_second), None),
6559 ],
6560 4,
6561 )
6562 .expect("relay");
6563 pair.relay_mut().set_strategy(CollaborationStrategy::Pair);
6564 assert_eq!(pair.relay_mut().enqueue_direct(1, "private"), Ok(true));
6565 pair.start().await.expect("start");
6566 assert!(matches!(
6567 pair.run_turn("ignored", 0).await.expect("direct turn"),
6568 RelayDecision::Dispatch {
6569 slot: 1,
6570 direct: true,
6571 ..
6572 }
6573 ));
6574 let direct_prompt = &pair.dispatches()[0].1;
6575 assert!(direct_prompt.contains("private"));
6576 assert!(!direct_prompt.contains("Pair role"));
6577 }
6578
6579 #[tokio::test]
6580 async fn relay_host_routes_around_a_usage_limited_agent() {
6581 let capabilities = AgentCapabilities::default();
6582 let first = ScriptedAdapter::new(
6583 0,
6584 capabilities.clone(),
6585 [
6586 AgentEvent::Text {
6587 slot: 0,
6588 text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
6589 credits or try again later."
6590 .into(),
6591 },
6592 AgentEvent::TurnComplete { slot: 0 },
6593 ],
6594 );
6595 let second = ScriptedAdapter::new(
6596 1,
6597 capabilities,
6598 [
6599 AgentEvent::Text {
6600 slot: 1,
6601 text: "review done".into(),
6602 },
6603 AgentEvent::TurnComplete { slot: 1 },
6604 ],
6605 );
6606 let hosts = vec![
6607 AdapterHost::new(Box::new(first), None),
6608 AdapterHost::new(Box::new(second), None),
6609 ];
6610 let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6611 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6612 let captured = std::sync::Arc::clone(&events);
6613 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6614 relay.start().await.expect("start");
6615 events.lock().expect("events").clear();
6616 assert!(matches!(
6617 relay.run_turn("task", 0).await.expect("limited turn"),
6618 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6619 ));
6620 assert!(
6621 events
6622 .lock()
6623 .expect("events")
6624 .iter()
6625 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6626 );
6627 assert!(matches!(
6629 relay.run_turn("", 0).await.expect("next turn"),
6630 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6631 ));
6632 assert!(relay.relay().is_limited(0));
6633 relay.reload(0).await.expect("reload");
6637 assert!(!relay.relay().is_limited(0));
6638 }
6639
6640 #[tokio::test]
6641 async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
6642 let limited = ScriptedAdapter::new(
6643 0,
6644 AgentCapabilities::default(),
6645 [AgentEvent::Failed {
6646 slot: 0,
6647 started: true,
6648 detail: "request failed: insufficient_quota".into(),
6649 }],
6650 );
6651 let healthy = ScriptedAdapter::new(
6652 1,
6653 AgentCapabilities::default(),
6654 [AgentEvent::TurnComplete { slot: 1 }],
6655 );
6656 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6657 let captured = std::sync::Arc::clone(&events);
6658 let mut relay = RelayHost::new(
6659 vec![
6660 AdapterHost::new(Box::new(limited), None),
6661 AdapterHost::new(Box::new(healthy), None),
6662 ],
6663 4,
6664 )
6665 .expect("relay");
6666 relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6667 relay.start().await.expect("start");
6668 events.lock().expect("events").clear();
6669
6670 assert!(matches!(
6671 relay.run_turn("task", 0).await.expect("limited failure"),
6672 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6673 ));
6674 assert!(relay.relay().is_limited(0));
6675 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
6676 {
6677 let events = events.lock().expect("events");
6678 assert!(
6679 events
6680 .iter()
6681 .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6682 );
6683 assert!(
6684 !events
6685 .iter()
6686 .any(|event| matches!(event, AgentEvent::Failed { .. }))
6687 );
6688 }
6689
6690 assert!(matches!(
6691 relay.run_turn("", 0).await.expect("healthy peer"),
6692 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6693 ));
6694 }
6695
6696 #[tokio::test]
6697 async fn relay_failure_is_skipped_for_one_batch_without_changing_the_roster() {
6698 let failed = ScriptedAdapter::new(
6699 0,
6700 AgentCapabilities::default(),
6701 [AgentEvent::Failed {
6702 slot: 0,
6703 started: true,
6704 detail: "connection lost".into(),
6705 }],
6706 );
6707 let healthy = ScriptedAdapter::new(
6708 1,
6709 AgentCapabilities::default(),
6710 [AgentEvent::TurnComplete { slot: 1 }],
6711 );
6712 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6713 let captured = std::sync::Arc::clone(&events);
6714 let mut relay = RelayHost::new(
6715 vec![
6716 AdapterHost::new(Box::new(failed), None),
6717 AdapterHost::new(Box::new(healthy), None),
6718 ],
6719 4,
6720 )
6721 .expect("relay");
6722 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6723 relay.start().await.expect("start");
6724
6725 assert!(matches!(
6726 relay.run_turn("task", 0).await.expect("handled failure"),
6727 crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6728 ));
6729 assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![0, 1]);
6730 assert!(relay.relay().is_limited(0));
6731 assert!(events.lock().expect("lock").iter().any(|event| {
6732 matches!(
6733 event,
6734 AgentEvent::Failed {
6735 slot: 0,
6736 started: true,
6737 ..
6738 }
6739 )
6740 }));
6741 assert!(matches!(
6742 relay.run_turn("", 0).await.expect("healthy peer"),
6743 crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6744 ));
6745 }
6746
6747 #[tokio::test]
6748 async fn codex_stop_does_not_skip_later_roster_reviewers() {
6749 let hosts = (0..3)
6750 .map(|slot| {
6751 AdapterHost::new(
6752 Box::new(ScriptedAdapter::new(
6753 slot,
6754 AgentCapabilities::default(),
6755 [
6756 AgentEvent::Text {
6757 slot,
6758 text: STOP_TOKEN.into(),
6759 },
6760 AgentEvent::TurnComplete { slot },
6761 ],
6762 )),
6763 None,
6764 )
6765 })
6766 .collect();
6767 let mut relay = RelayHost::new(hosts, 10).expect("relay");
6768 relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
6769 relay.start().await.expect("start");
6770 for expected in 0..3 {
6771 assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
6772 RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
6773 }
6774 assert_eq!(
6775 relay.run_turn("", 0).await.expect("complete"),
6776 RelayDecision::Complete
6777 );
6778 }
6779
6780 #[tokio::test]
6781 async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
6782 let first = ScriptedAdapter::new(
6783 0,
6784 AgentCapabilities::default(),
6785 [
6786 AgentEvent::Text {
6787 slot: 0,
6788 text: "done".into(),
6789 },
6790 AgentEvent::TurnComplete { slot: 0 },
6791 ],
6792 );
6793 let reviewer = ScriptedAdapter::new(
6794 1,
6795 AgentCapabilities::default(),
6796 [
6797 AgentEvent::Text {
6798 slot: 1,
6799 text: STOP_TOKEN.into(),
6800 },
6801 AgentEvent::TurnComplete { slot: 1 },
6802 ],
6803 );
6804 let mut relay = RelayHost::new(
6805 vec![
6806 AdapterHost::new(Box::new(first), None),
6807 AdapterHost::new(Box::new(reviewer), None),
6808 ],
6809 10,
6810 )
6811 .expect("relay");
6812 relay.start().await.expect("start");
6813 let first_decision = relay.run_turn("task", 0).await.expect("first");
6814 assert!(matches!(
6815 first_decision,
6816 RelayDecision::Dispatch { slot: 0, .. }
6817 ));
6818 let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
6819 assert!(matches!(
6820 reviewer_decision,
6821 RelayDecision::Dispatch {
6822 slot: 1,
6823 can_stop: true,
6824 ..
6825 }
6826 ));
6827 assert_eq!(
6828 relay.run_turn("", 0).await.expect("complete"),
6829 RelayDecision::Complete
6830 );
6831 }
6832
6833 #[tokio::test]
6834 async fn relay_stream_emits_text_and_thought_endings_before_tools() {
6835 let tool = AgentEvent::Tool {
6836 slot: 0,
6837 update: crate::ToolUpdate {
6838 id: "read".into(),
6839 title: "Read file".into(),
6840 status: ToolStatus::Running,
6841 detail: None,
6842 },
6843 };
6844 let updates = vec![
6845 AgentEvent::Thought {
6846 slot: 0,
6847 text: "Check the buffer. ✈".into(),
6848 },
6849 AgentEvent::Text {
6850 slot: 0,
6851 text: "Let me check.".into(),
6852 },
6853 tool.clone(),
6854 AgentEvent::Text {
6855 slot: 0,
6856 text: "[CODE".into(),
6857 },
6858 AgentEvent::Text {
6859 slot: 0,
6860 text: " is ordinary.".into(),
6861 },
6862 AgentEvent::Text {
6863 slot: 0,
6864 text: "[CODESWARM:".into(),
6865 },
6866 AgentEvent::Text {
6867 slot: 0,
6868 text: "STOP] Done.".into(),
6869 },
6870 tool,
6871 AgentEvent::TurnComplete { slot: 0 },
6872 ];
6873 let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
6874 let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
6875 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6876 let captured = std::sync::Arc::clone(&events);
6877 let mut relay = RelayHost::new(
6878 vec![
6879 AdapterHost::new(Box::new(first), None),
6880 AdapterHost::new(Box::new(reviewer), None),
6881 ],
6882 2,
6883 )
6884 .expect("relay");
6885 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6886 relay.start().await.unwrap();
6887 relay.run_turn("task", 0).await.unwrap();
6888 let captured = events.lock().unwrap();
6889 let visible: Vec<_> = captured
6890 .iter()
6891 .filter(|event| {
6892 matches!(
6893 event,
6894 AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
6895 )
6896 })
6897 .cloned()
6898 .collect();
6899 assert_eq!(
6900 visible,
6901 vec![
6902 updates[0].clone(),
6903 updates[1].clone(),
6904 updates[2].clone(),
6905 AgentEvent::Text {
6906 slot: 0,
6907 text: "[CODE is ordinary.".into()
6908 },
6909 AgentEvent::Text {
6910 slot: 0,
6911 text: " Done.".into()
6912 },
6913 updates[7].clone(),
6914 ]
6915 );
6916 }
6917
6918 #[tokio::test]
6919 async fn roster_handoff_routes_only_terminal_message_markers_and_refreshes_targets() {
6920 let text = |value: &str| AgentEvent::Text {
6921 slot: 0,
6922 text: value.into(),
6923 };
6924 let thought = || AgentEvent::Thought {
6925 slot: 0,
6926 text: "still checking".into(),
6927 };
6928 let tool = || AgentEvent::Tool {
6929 slot: 0,
6930 update: crate::ToolUpdate {
6931 id: "read".into(),
6932 title: "Read file".into(),
6933 status: ToolStatus::Running,
6934 detail: None,
6935 },
6936 };
6937 let cases = vec![
6938 (vec![text("result [CODESWARM:NEXT:3]\n ")], 2),
6939 (
6940 vec![text("result [CODESWARM:"), text("NEXT:"), text("3]")],
6941 2,
6942 ),
6943 (vec![text("result [CODESWARM:NEXT:3]"), text(" more")], 1),
6944 (vec![text("result [CODESWARM:NEXT:3]"), thought()], 1),
6945 (
6946 vec![text("result [CODESWARM:NEXT:3]"), tool(), text(" ")],
6947 1,
6948 ),
6949 (
6950 vec![text("result [CODESWARM:NEXT:"), thought(), text("3]")],
6951 1,
6952 ),
6953 (
6954 vec![text("result"), thought(), text("[CODESWARM:NEXT:3]")],
6955 2,
6956 ),
6957 (vec![text("result [CODESWARM:NEXT:1]")], 1),
6958 (vec![text("result [CODESWARM:NEXT:0]")], 1),
6959 (vec![text("result [CODESWARM:NEXT:99]")], 1),
6960 (
6961 vec![
6962 text("result [CODESWARM:NEXT:3]"),
6963 AgentEvent::UsageUpdated {
6964 slot: 0,
6965 usage: crate::UsageUpdate { used: 1, size: 100 },
6966 },
6967 ],
6968 2,
6969 ),
6970 ];
6971 for (mut updates, expected) in cases {
6972 updates.push(AgentEvent::TurnComplete { slot: 0 });
6973 let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates);
6974 let hosts = std::iter::once(AdapterHost::new(Box::new(first), None))
6975 .chain((1..3).map(|slot| {
6976 AdapterHost::new(
6977 Box::new(ScriptedAdapter::new(
6978 slot,
6979 AgentCapabilities::default(),
6980 [AgentEvent::TurnComplete { slot }],
6981 )),
6982 None,
6983 )
6984 }))
6985 .collect();
6986 let mut relay = RelayHost::new(hosts, 10).unwrap();
6987 relay.set_roster_names(vec!["Worker".into(), "Codex".into(), "Codex".into()]);
6988 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6989 let captured = events.clone();
6990 relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6991 relay.start().await.unwrap();
6992 relay.run_turn("task", 0).await.unwrap();
6993 let prompt = &relay.dispatches()[0].1;
6994 assert!(prompt.contains("[CODESWARM:NEXT:2] → Codex"));
6995 assert!(prompt.contains("[CODESWARM:NEXT:3] → Codex"));
6996 assert!(!prompt.contains("[CODESWARM:NEXT:1]"));
6997 relay.introduced.fill(true);
6999 relay.set_roster_names(vec!["Replacement".into(), "Codex".into(), "Codex".into()]);
7000 let next = relay.run_turn("", 0).await.unwrap();
7001 assert!(
7002 matches!(next, RelayDecision::Dispatch { slot, can_stop: false, .. } if slot == expected),
7003 "{next:?}"
7004 );
7005 let prompt = &relay.dispatches()[1].1;
7006 assert!(prompt.contains("[CODESWARM:NEXT:1] → Replacement"));
7007 assert!(prompt.contains("result"));
7008 let public = prompt
7009 .split("Public updates:\n")
7010 .nth(1)
7011 .unwrap()
7012 .split("\n\nDo not use")
7013 .next()
7014 .unwrap();
7015 assert!(!public.contains("[CODESWARM:NEXT:"));
7016 let visible = events
7017 .lock()
7018 .unwrap()
7019 .iter()
7020 .filter_map(|event| match event {
7021 AgentEvent::Text { text, .. } => Some(text.clone()),
7022 _ => None,
7023 })
7024 .collect::<String>();
7025 assert!(visible.contains("result"));
7026 assert!(!visible.contains("[CODESWARM:"), "{visible}");
7027 }
7028 }
7029
7030 #[tokio::test]
7031 async fn reviewer_stop_requires_a_terminal_marker_after_all_activity() {
7032 let text = |value: &str| AgentEvent::Text {
7033 slot: 1,
7034 text: value.into(),
7035 };
7036 let thought = || AgentEvent::Thought {
7037 slot: 1,
7038 text: "still checking".into(),
7039 };
7040 let tool = || AgentEvent::Tool {
7041 slot: 1,
7042 update: crate::ToolUpdate {
7043 id: "read".into(),
7044 title: "Read file".into(),
7045 status: ToolStatus::Running,
7046 detail: None,
7047 },
7048 };
7049 let cases = vec![
7050 (vec![text(&format!("done {STOP_TOKEN}"))], true),
7051 (vec![text(STOP_TOKEN), text("\n ")], true),
7052 (vec![text(STOP_TOKEN), text(" actually keep going")], false),
7053 (vec![text(STOP_TOKEN), thought()], false),
7054 (vec![text(STOP_TOKEN), tool()], false),
7055 (vec![text(STOP_TOKEN), tool(), text(" ")], false),
7056 (vec![text(STOP_TOKEN), tool(), text(STOP_TOKEN)], true),
7057 (vec![text("[CODESWARM:"), text("STOP]")], true),
7058 (vec![text("[CODESWARM:"), thought(), text("STOP]")], false),
7059 (
7060 vec![AgentEvent::Thought {
7061 slot: 1,
7062 text: STOP_TOKEN.into(),
7063 }],
7064 false,
7065 ),
7066 (
7067 vec![
7068 text(STOP_TOKEN),
7069 AgentEvent::UsageUpdated {
7070 slot: 1,
7071 usage: crate::UsageUpdate { used: 1, size: 100 },
7072 },
7073 ],
7074 true,
7075 ),
7076 ];
7077 for (mut events, stop) in cases {
7078 let first = ScriptedAdapter::new(
7079 0,
7080 AgentCapabilities::default(),
7081 [
7082 AgentEvent::Text {
7083 slot: 0,
7084 text: "initial response".into(),
7085 },
7086 AgentEvent::TurnComplete { slot: 0 },
7087 AgentEvent::TurnComplete { slot: 0 },
7088 ],
7089 );
7090 events.push(AgentEvent::TurnComplete { slot: 1 });
7091 let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), events.clone());
7092 let mut relay = RelayHost::new(
7093 vec![
7094 AdapterHost::new(Box::new(first), None),
7095 AdapterHost::new(Box::new(reviewer), None),
7096 ],
7097 4,
7098 )
7099 .unwrap();
7100 relay.start().await.unwrap();
7101 relay.run_turn("task", 0).await.unwrap();
7102 relay.run_turn("", 0).await.unwrap();
7103 let next = relay.run_turn("", 0).await.unwrap();
7104 assert_eq!(
7105 matches!(next, RelayDecision::Complete),
7106 stop,
7107 "events={events:?}"
7108 );
7109 relay.stop().await.unwrap();
7110 }
7111 }
7112
7113 #[tokio::test]
7114 async fn stop_token_is_filtered_from_streamed_ui_events() {
7115 let first = ScriptedAdapter::new(
7116 0,
7117 AgentCapabilities::default(),
7118 [
7119 AgentEvent::Text {
7120 slot: 0,
7121 text: format!("visible {STOP_TOKEN} trailing"),
7122 },
7123 AgentEvent::TurnComplete { slot: 0 },
7124 ],
7125 );
7126 let reviewer = ScriptedAdapter::new(
7127 1,
7128 AgentCapabilities::default(),
7129 [AgentEvent::TurnComplete { slot: 1 }],
7130 );
7131 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7132 let captured = std::sync::Arc::clone(&events);
7133 let mut relay = RelayHost::new(
7134 vec![
7135 AdapterHost::new(Box::new(first), None),
7136 AdapterHost::new(Box::new(reviewer), None),
7137 ],
7138 2,
7139 )
7140 .expect("relay");
7141 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7142 relay.start().await.expect("start");
7143 relay.run_turn("task", 0).await.expect("turn");
7144 let captured = events.lock().expect("lock");
7145 assert!(captured.iter().all(|event| match event {
7146 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7147 _ => true,
7148 }));
7149 let visible = captured
7150 .iter()
7151 .filter_map(|event| match event {
7152 AgentEvent::Text { text, .. } => Some(text.as_str()),
7153 _ => None,
7154 })
7155 .collect::<String>();
7156 assert_eq!(visible, "visible trailing");
7157 }
7158
7159 #[tokio::test]
7160 async fn token_only_reviewer_response_emits_visible_acknowledgment() {
7161 let first = ScriptedAdapter::new(
7162 0,
7163 AgentCapabilities::default(),
7164 [
7165 AgentEvent::Text {
7166 slot: 0,
7167 text: "done".into(),
7168 },
7169 AgentEvent::TurnComplete { slot: 0 },
7170 ],
7171 );
7172 let reviewer = ScriptedAdapter::new(
7173 1,
7174 AgentCapabilities::default(),
7175 [
7176 AgentEvent::Text {
7177 slot: 1,
7178 text: STOP_TOKEN.into(),
7179 },
7180 AgentEvent::TurnComplete { slot: 1 },
7181 ],
7182 );
7183 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7184 let captured = std::sync::Arc::clone(&events);
7185 let mut relay = RelayHost::new(
7186 vec![
7187 AdapterHost::new(Box::new(first), None),
7188 AdapterHost::new(Box::new(reviewer), None),
7189 ],
7190 4,
7191 )
7192 .expect("relay");
7193 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7194 relay.start().await.expect("start");
7195 relay.run_turn("task", 0).await.expect("first turn");
7196 relay.run_turn("", 0).await.expect("review turn");
7197 let captured = events.lock().expect("lock");
7198 assert!(captured.iter().any(|event| {
7199 matches!(
7200 event,
7201 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7202 )
7203 }));
7204 assert!(captured.iter().all(|event| match event {
7205 AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7206 _ => true,
7207 }));
7208 let acknowledgment = captured
7209 .iter()
7210 .position(|event| {
7211 matches!(
7212 event,
7213 AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7214 )
7215 })
7216 .expect("visible acknowledgment");
7217 let completion = captured
7218 .iter()
7219 .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
7220 .expect("reviewer completion");
7221 assert!(acknowledgment < completion);
7222 }
7223
7224 #[tokio::test]
7225 async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
7226 let first = ScriptedAdapter::new(
7227 0,
7228 AgentCapabilities::default(),
7229 [
7230 AgentEvent::Text {
7231 slot: 0,
7232 text: "done".into(),
7233 },
7234 AgentEvent::TurnComplete { slot: 0 },
7235 ],
7236 );
7237 let reviewer = ScriptedAdapter::new(
7238 1,
7239 AgentCapabilities::default(),
7240 [
7241 AgentEvent::Text {
7242 slot: 1,
7243 text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
7244 },
7245 AgentEvent::TurnComplete { slot: 1 },
7246 ],
7247 );
7248 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7249 let captured = std::sync::Arc::clone(&events);
7250 let mut relay = RelayHost::new(
7251 vec![
7252 AdapterHost::new(Box::new(first), None),
7253 AdapterHost::new(Box::new(reviewer), None),
7254 ],
7255 4,
7256 )
7257 .expect("relay");
7258 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7259 relay.start().await.expect("start");
7260 relay.run_turn("task", 0).await.expect("first turn");
7261 relay.run_turn("", 0).await.expect("review turn");
7262
7263 let visible = events
7264 .lock()
7265 .expect("lock")
7266 .iter()
7267 .filter_map(|event| match event {
7268 AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
7269 _ => None,
7270 })
7271 .collect::<String>();
7272 assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
7273 assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
7274 }
7275
7276 #[tokio::test]
7277 async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
7278 let first = AdapterHost::new(
7279 Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
7280 None,
7281 );
7282 let second = AdapterHost::new(
7283 Box::new(ScriptedAdapter::new(
7284 1,
7285 AgentCapabilities::default(),
7286 [AgentEvent::TurnComplete { slot: 1 }],
7287 )),
7288 None,
7289 );
7290 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7291 let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
7292 relay.set_event_sink(move |event| {
7293 if matches!(event, AgentEvent::Permission { .. }) {
7294 let _ = seen_sender.send(());
7295 }
7296 });
7297 relay.start().await.expect("start");
7298 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
7299 let answer = async move {
7300 seen_receiver.recv().await.expect("permission request");
7301 sender
7302 .send(super::RelayPermissionAnswer {
7303 slot: 0,
7304 request_id: "permission-1".into(),
7305 answer: PermissionAnswer::Selected {
7306 option_id: "allow".into(),
7307 },
7308 })
7309 .expect("queue permission answer");
7310 };
7311 tokio::time::timeout(std::time::Duration::from_millis(100), async {
7312 let ((), result) = tokio::join!(
7313 answer,
7314 relay.run_turn_with_permissions("task", 0, &mut receiver)
7315 );
7316 result
7317 })
7318 .await
7319 .expect("permission-gated turn should not deadlock")
7320 .expect("turn completes");
7321 }
7322
7323 #[tokio::test]
7324 async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
7325 let first = AdapterHost::new(
7326 Box::new(PendingAdapter {
7327 slot: 0,
7328 hang_on_cancel: false,
7329 }),
7330 None,
7331 );
7332 let second = AdapterHost::new(
7333 Box::new(ScriptedAdapter::new(
7334 1,
7335 AgentCapabilities::default(),
7336 [AgentEvent::TurnComplete { slot: 1 }],
7337 )),
7338 None,
7339 );
7340 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7341 relay.start().await.expect("start");
7342 let cancellation = relay.cancellation();
7343 let error = {
7344 let turn = relay.run_turn("task", 0);
7345 tokio::pin!(turn);
7346 cancellation.request();
7347 turn.await.expect_err("cancellation should stop turn")
7348 };
7349 assert!(error.to_string().contains("relay turn cancelled"));
7350
7351 assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
7352 relay
7353 .run_turn("", 1)
7354 .await
7355 .expect("replacement job reaches the selected peer");
7356 let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
7357 assert!(replacement.contains("replacement job"));
7358 assert!(replacement.contains("User "));
7359 assert!(replacement.contains(":\ntask"));
7360 let owner_updates = relay.relay_mut().unseen_context(0);
7361 assert!(owner_updates.contains("User "));
7362 assert!(owner_updates.contains(":\ntask"));
7363 assert!(owner_updates.contains(":\nreplacement job"));
7364 }
7365
7366 #[tokio::test]
7367 async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
7368 let first = AdapterHost::new(
7369 Box::new(PendingAdapter {
7370 slot: 0,
7371 hang_on_cancel: true,
7372 }),
7373 None,
7374 );
7375 let second = AdapterHost::new(
7376 Box::new(PendingAdapter {
7377 slot: 1,
7378 hang_on_cancel: false,
7379 }),
7380 None,
7381 );
7382 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7383 relay.start().await.expect("start");
7384 let cancellation = relay.cancellation();
7385 let turn = relay.run_turn("task", 0);
7386 tokio::pin!(turn);
7387 cancellation.request();
7388 let error = turn.await.expect_err("cancellation should stop turn");
7389 assert!(error.to_string().contains("timed out"));
7390 }
7391
7392 #[tokio::test]
7393 async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
7394 let event = [AgentEvent::TurnComplete { slot: 0 }];
7395 let first = AdapterHost::new(
7396 Box::new(ScriptedAdapter::new(
7397 0,
7398 AgentCapabilities::default(),
7399 event.clone(),
7400 )),
7401 None,
7402 );
7403 let second = AdapterHost::new(
7404 Box::new(ScriptedAdapter::new(
7405 1,
7406 AgentCapabilities::default(),
7407 [AgentEvent::TurnComplete { slot: 1 }],
7408 )),
7409 None,
7410 );
7411 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7412 relay.start().await.expect("start");
7413
7414 relay.pause();
7415 assert_eq!(
7416 relay.run_turn("paused", 0).await.expect("paused turn"),
7417 crate::relay::RelayDecision::Paused
7418 );
7419 assert!(relay.dispatches().is_empty());
7420
7421 relay.resume();
7422 relay.relay_mut().drop_agent(1).expect("drop reviewer");
7423 assert!(matches!(
7424 relay
7425 .run_turn("solo follow-up", 0)
7426 .await
7427 .expect("solo turn"),
7428 crate::relay::RelayDecision::Dispatch {
7429 slot: 0,
7430 can_stop: false,
7431 ..
7432 }
7433 ));
7434 assert_eq!(relay.dispatches().len(), 1);
7435 }
7436
7437 #[tokio::test]
7438 async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
7439 let first = AdapterHost::new(
7440 Box::new(ScriptedAdapter::new(
7441 0,
7442 AgentCapabilities::default(),
7443 [AgentEvent::TurnComplete { slot: 0 }],
7444 )),
7445 None,
7446 );
7447 let second = AdapterHost::new(
7448 Box::new(ScriptedAdapter::new(
7449 1,
7450 AgentCapabilities::default(),
7451 [AgentEvent::TurnComplete { slot: 1 }],
7452 )),
7453 None,
7454 );
7455 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7456 relay.set_roster_names(vec!["First".into(), "Second".into()]);
7457 relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
7458 relay.set_roster_launch_specs(vec![
7459 ("custom".into(), "owner".into()),
7460 ("custom".into(), "peer".into()),
7461 ]);
7462 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7463 let captured = std::sync::Arc::clone(&events);
7464 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7465 relay.start().await.expect("start");
7466 let slot = relay
7467 .add_agent(
7468 AdapterHost::new(
7469 Box::new(ScriptedAdapter::new(
7470 2,
7471 AgentCapabilities::default(),
7472 [AgentEvent::TurnComplete { slot: 2 }],
7473 )),
7474 None,
7475 ),
7476 "Reviewer",
7477 "reviewer.example",
7478 "reviewer --acp",
7479 )
7480 .await
7481 .expect("append agent");
7482 assert_eq!(slot, 2);
7483 assert_eq!(
7484 relay.relay().active_slots().collect::<Vec<_>>(),
7485 vec![0, 1, 2]
7486 );
7487 assert_eq!(
7488 relay
7489 .session_metadata()
7490 .get("agents")
7491 .and_then(|value| value.as_array())
7492 .map(Vec::len),
7493 Some(3)
7494 );
7495 relay.drop_agent(1).await.expect("drop middle peer");
7496 let metadata = relay.session_metadata();
7497 assert_eq!(
7498 metadata.get("agents"),
7499 Some(&serde_json::json!([
7500 {"slot": 0, "name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
7501 {"slot": 2, "name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
7502 ]))
7503 );
7504 assert!(
7505 events
7506 .lock()
7507 .expect("lock")
7508 .iter()
7509 .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
7510 );
7511 }
7512
7513 #[tokio::test]
7514 async fn relay_host_persists_coordinator_owned_runtime_metadata() {
7515 let path = unique_test_path("codeswarm-session-metadata", "json");
7516 let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
7517 let writer = metadata_store.buffered().expect("metadata writer");
7518 let first = AdapterHost::new(
7519 Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
7520 None,
7521 );
7522 let second = AdapterHost::new(
7523 Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
7524 None,
7525 );
7526 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7527 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7528 relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
7529 relay.set_roster_launch_specs(vec![
7530 ("custom".into(), "claude".into()),
7531 ("custom".into(), "codex".into()),
7532 ]);
7533 relay.set_session_metadata_writer(writer);
7534 relay.start().await.expect("start");
7535 relay.drop_agent(0).await.expect("drop first agent");
7536 relay.stop().await.expect("stop");
7537
7538 let loaded = metadata_store
7539 .read()
7540 .expect("read metadata")
7541 .expect("metadata snapshot");
7542 assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
7543 assert_eq!(
7544 loaded.get("agents"),
7545 Some(&serde_json::json!([{
7546 "slot": 1, "name": "Codex", "identity": "openai.com", "protocol": "custom",
7547 "command": "codex", "supports_load_session": false
7548 }]))
7549 );
7550 assert!(loaded.get("owner").is_none());
7551 let _ = std::fs::remove_file(path);
7552 }
7553
7554 #[tokio::test]
7555 async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
7556 let first = AdapterHost::new(
7557 Box::new(ScriptedAdapter::new(
7558 0,
7559 AgentCapabilities::default(),
7560 [
7561 AgentEvent::Text {
7562 slot: 0,
7563 text: "owner stream".into(),
7564 },
7565 AgentEvent::TurnComplete { slot: 0 },
7566 ],
7567 )),
7568 None,
7569 );
7570 let second = AdapterHost::new(
7571 Box::new(ScriptedAdapter::new(
7572 1,
7573 AgentCapabilities::default(),
7574 [
7575 AgentEvent::Text {
7576 slot: 1,
7577 text: "peer stream".into(),
7578 },
7579 AgentEvent::TurnComplete { slot: 1 },
7580 ],
7581 )),
7582 None,
7583 );
7584 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7585 relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
7586 relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
7587 let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7588 let captured = std::sync::Arc::clone(&events);
7589 relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7590 relay.start().await.expect("start");
7591
7592 relay.swap_agents(0, 1).expect("swap peers");
7593 assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
7594 assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
7595 relay.run_turn("task", 0).await.expect("swapped turn");
7596 let events = events.lock().expect("events");
7597 assert!(events.iter().any(|event| {
7598 matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
7599 }));
7600 assert!(relay.dispatches()[0].1.contains("You are Peer"));
7601 }
7602
7603 #[tokio::test]
7604 async fn relay_host_persists_all_active_agent_metadata_off_thread() {
7605 let path = unique_test_path("codeswarm-session-metadata", "json");
7606 let first = AdapterHost::new(
7607 Box::new(ScriptedAdapter::new(
7608 0,
7609 AgentCapabilities::default(),
7610 [AgentEvent::TurnComplete { slot: 0 }],
7611 )),
7612 None,
7613 );
7614 let second = AdapterHost::new(
7615 Box::new(ScriptedAdapter::new(
7616 1,
7617 AgentCapabilities::default(),
7618 [AgentEvent::TurnComplete { slot: 1 }],
7619 )),
7620 None,
7621 );
7622 let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7623 relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7624 relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
7625 relay.set_roster_launch_specs(vec![
7626 ("custom".into(), "claude".into()),
7627 ("custom".into(), "codex".into()),
7628 ]);
7629 let writer = SessionMetadataStore::open(&path)
7630 .buffered()
7631 .expect("metadata writer");
7632 relay.set_session_metadata_writer(writer);
7633 relay.start().await.expect("start");
7634 relay.stop().await.expect("stop");
7635 let loaded = SessionMetadataStore::open(&path)
7636 .read()
7637 .expect("read metadata")
7638 .expect("metadata snapshot");
7639 let agents = loaded
7640 .get("agents")
7641 .and_then(|value| value.as_array())
7642 .expect("agents");
7643 assert_eq!(agents.len(), 2);
7644 assert_eq!(agents[0]["identity"], "claude.com");
7645 assert_eq!(agents[1]["identity"], "openai.com");
7646 let _ = std::fs::remove_file(path);
7647 }
7648
7649 #[tokio::test]
7650 async fn relay_host_routes_unseen_public_context_to_next_agent() {
7651 let first = AdapterHost::new(
7652 Box::new(ScriptedAdapter::new(
7653 0,
7654 AgentCapabilities::default(),
7655 [
7656 AgentEvent::Text {
7657 slot: 0,
7658 text: "implemented the fix".into(),
7659 },
7660 AgentEvent::TurnComplete { slot: 0 },
7661 ],
7662 )),
7663 None,
7664 );
7665 let second = AdapterHost::new(
7666 Box::new(ScriptedAdapter::new(
7667 1,
7668 AgentCapabilities::default(),
7669 [AgentEvent::TurnComplete { slot: 1 }],
7670 )),
7671 None,
7672 );
7673 let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7674 relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
7675 relay.start().await.expect("start");
7676 relay.run_turn("task", 0).await.expect("first turn");
7677 relay.run_turn("review this", 0).await.expect("review turn");
7678
7679 assert_eq!(relay.dispatches().len(), 2);
7680 assert_eq!(relay.dispatches()[0].0, 0);
7681 assert!(relay.dispatches()[0].1.contains("task"));
7682 assert!(relay.dispatches()[0].1.contains("You are Codex"));
7683 assert!(relay.dispatches()[0].1.contains("2. Qwen"));
7684 assert_eq!(relay.dispatches()[1].0, 1);
7685 assert!(relay.dispatches()[1].1.contains("review this"));
7686 let public = relay.dispatches()[1]
7687 .1
7688 .split_once("Public updates:\n")
7689 .map(|(_, updates)| updates)
7690 .expect("review receives public context");
7691 let header = public
7692 .lines()
7693 .find(|line| line.starts_with("Codex "))
7694 .expect("named previous agent");
7695 let timestamp = header
7696 .strip_prefix("Codex ")
7697 .and_then(|value| value.strip_suffix(':'))
7698 .expect("timestamped header");
7699 assert_eq!(timestamp.len(), 5);
7700 assert_eq!(timestamp.as_bytes()[2], b':');
7701 assert!(
7702 timestamp
7703 .bytes()
7704 .enumerate()
7705 .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
7706 );
7707 assert!(public.contains("implemented the fix"));
7708 assert!(!public.contains("Agent 0"));
7709 }
7710}