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