Skip to main content

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