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