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