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