Skip to main content

codeswarm_adapters/
adapters.rs

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