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