cydonia 0.1.11

Desktop workspace for the ACP agents you run, keeping what they produce as files on your disk
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
//! One ACP session over a spawned agent subprocess.
//!
//! Two things here are load-bearing:
//! - All agent-side events flow through ONE channel. cacp's read loop awaits
//!   each notification before it reads the next frame, so a single channel
//!   preserves the exact wire order (a turn's final updates arrive before its
//!   result — separate channels lose that).
//! - The connection runs on its own tokio runtime. cacp spawns its read and
//!   write loops with `tokio::spawn`, and gpui's executor is smol's.

use crate::{
    agent::{mcp, serve},
    model::{
        media,
        session_preferences::{self, Choices},
        settings,
    },
};
use anyhow::{Result, anyhow};
use cacp::{
    AgentConn, Client, Direction, Error, Tap,
    client::{HistoryEntry, HistoryFork, HistoryRole},
    schema::{
        AuthenticateRequest, CancelNotification, ClientCapabilities, ContentBlock, EnvVariable,
        FileSystemCapabilities, HttpHeader, ImageContent, InitializeRequest, InitializeResponse,
        LoadSessionRequest, McpServer, McpServerHttp, McpServerStdio, NewSessionRequest,
        NewSessionResponse, PromptRequest, ReadTextFileRequest, ReadTextFileResponse,
        RequestPermissionRequest, RequestPermissionResponse, SessionConfigOptionValue, SessionId,
        SessionNotification, SessionUpdate, SetSessionConfigOptionRequest, SetSessionModeRequest,
        StopReason, WriteTextFileRequest, WriteTextFileResponse,
    },
};
use std::process::Stdio;
use std::{
    path::PathBuf,
    sync::{Arc, OnceLock},
    time::Duration,
};
use tokio::{
    io::{AsyncBufReadExt, BufReader},
    process::{Child, ChildStderr, Command},
    runtime::Runtime,
    sync::{Mutex, mpsc, oneshot},
};

/// Where the protocol tap writes, and the switch that echoes an agent's stderr
/// to ours on the way past — it reaches the transcript either way.
const DEBUG: &str = "CYDONIA_DEBUG";

/// What our own tool server is called to an agent. Also the namespace every
/// tool of ours comes back under, so the two are read from one place.
const SERVER: &str = "cydonia";

/// How long a `session/cancel` already on the wire is given to land before the
/// process is killed under it.
///
/// Short, because it is all this can buy. Closing the agent's stdin — the thing
/// that would actually let it wind itself up — is not reachable from here:
/// cacp's read loop holds a `Peer` of its own, so the write loop that owns
/// stdin outlives every handle this side can drop.
const SHUTDOWN_GRACE: Duration = Duration::from_millis(250);

/// The runtime every connection runs on, started on first use.
pub fn runtime() -> &'static Runtime {
    static RUNTIME: OnceLock<Runtime> = OnceLock::new();
    RUNTIME.get_or_init(|| Runtime::new().expect("failed to start the tokio runtime"))
}

/// Everything the agent side feeds into the frontend, in wire order.
pub enum Event {
    Update(SessionUpdate),
    /// The agent asks the user to authorize a tool call.
    Permission(RequestPermissionRequest, Reply<RequestPermissionResponse>),
    /// One line the agent wrote to its stderr. Not protocol — the runtime
    /// under it talking, or the shell that could not start it — so it arrives
    /// off its own task and only approximately in step with the rest.
    Stderr(String),
    /// The prompt turn settled: its stop reason, or the agent's error.
    TurnDone(Result<StopReason, Error>),
    /// The agent's read loop ended — the process died or closed its stdout.
    /// Last in wire order, so any final updates land before the frontend
    /// gives the session up.
    Closed,
}

/// The frontend's receiving end.
pub type Events = mpsc::UnboundedReceiver<Event>;

/// The agent side's sending end.
pub type Sender = mpsc::UnboundedSender<Event>;

/// Open the channel a session runs on.
///
/// The caller holds both halves and lends the sender to [`Session::spawn`],
/// rather than being given the pair back by it. A launch that fails is why:
/// the process can write to its stderr and die without ever answering
/// `initialize`, and a receiver created inside the launch would be dropped
/// with the error, taking the only account of what went wrong with it.
pub fn channel() -> (Sender, Events) {
    mpsc::unbounded_channel()
}

/// The answer half of a request the frontend has to make. Dropping it
/// declines the request rather than hanging the agent.
pub struct Reply<T>(oneshot::Sender<Result<T, Error>>);

impl<T> Reply<T> {
    pub fn send(self, value: T) {
        let _ = self.0.send(Ok(value));
    }
}

/// A live session: the connection, its identity, and the agent process.
pub struct Session {
    /// `Some` for the whole of a session's life — emptied only by
    /// [`Session::drop`], which has to take it to close it before the process.
    conn: Option<AgentConn>,
    tx: mpsc::UnboundedSender<Event>,
    /// The agent process. Held in an `Option` for the same reason as `conn`.
    child: Option<Child>,
    pub session_id: SessionId,
    pub init: InitializeResponse,
    pub response: NewSessionResponse,
    pub cwd: PathBuf,
    /// True when an existing session was loaded (history replayed as
    /// queued [`Event::Update`]s) instead of a fresh one created.
    pub loaded: bool,
    built_in_mcp: bool,
    history_fork: Option<Arc<Mutex<HistoryFork>>>,
}

/// How to open a session.
#[derive(Default)]
pub struct Launch {
    /// The session's working directory.
    pub cwd: PathBuf,
    /// A session to load instead of starting fresh.
    pub previous: Option<String>,
    pub history: Option<Vec<HistoryEntry>>,
    pub history_pending: bool,
    pub choices: Choices,
}

impl Launch {
    pub fn new(cwd: PathBuf) -> Self {
        Self {
            cwd,
            ..Default::default()
        }
    }
}

impl Session {
    /// Spawn `entry` over stdio, initialize, and open a session. The agent
    /// dies with the returned [`Session`].
    ///
    /// With [`Launch::previous`] set and the agent capable, `session/load`
    /// replays that session's history instead of starting fresh; a failed
    /// load (stale id, agent restart) falls back to a new session.
    pub async fn spawn(entry: &settings::Agent, launch: Launch, tx: Sender) -> Result<Self> {
        let mut command = Command::new(&entry.command);
        command.args(&entry.args).envs(&entry.env);
        // HTTP clients can inherit a system proxy that does not exempt IP
        // loopback addresses. Our MCP server must be reached directly. Merge
        // both spellings because agents differ in which one they honor.
        let bypass = loopback_bypass(["NO_PROXY", "no_proxy"].map(|key| {
            entry
                .env
                .get(key)
                .cloned()
                .or_else(|| std::env::var(key).ok())
        }));
        command.env("NO_PROXY", &bypass).env("no_proxy", &bypass);
        // An agent's diagnostics are not this app's to print. cacp leaves the
        // choice to the caller — "a TUI usually wants it captured and a CLI
        // usually does not" — and a desktop app that inherits them sprays a
        // node SDK's teardown chatter over whichever terminal happened to
        // launch it, about a shutdown the user asked for.
        //
        // Captured rather than discarded, though, because it is the only
        // account of a process that dies without ever speaking protocol: the
        // `#!/usr/bin/env node` shim that found no node, the package that
        // would not resolve. It reaches the transcript as the execution it is
        // — see [`Event::Stderr`]. `CYDONIA_DEBUG`, which already redirects
        // the protocol tap, still echoes the lines where a developer looks.
        command.stderr(Stdio::piped());

        let configured = mcp::servers();

        let (conn, mut child) =
            cacp::spawn(&mut command, Arc::new(Frontend(tx.clone())), debug_tap())
                .map_err(|e| anyhow!("failed to start {}: {}", entry.command, error_text(&e)))?;
        if let Some(stderr) = child.stderr.take() {
            runtime().spawn(drain(stderr, tx.clone()));
        }

        Self::open(conn, child, tx, launch, configured).await
    }

    async fn open(
        conn: AgentConn,
        child: Child,
        tx: mpsc::UnboundedSender<Event>,
        launch: Launch,
        configured: Vec<mcp::McpServer>,
    ) -> Result<Self> {
        let cwd = launch.cwd.clone();
        let init = conn
            .initialize(InitializeRequest::new(ClientCapabilities {
                fs: FileSystemCapabilities {
                    read_text_file: true,
                    write_text_file: true,
                    meta: None,
                },
                ..Default::default()
            }))
            .await
            .map_err(|e| anyhow!("initialize failed: {}", error_text(&e)))?;

        // Only now are the agent's MCP capabilities known, so remote
        // servers can be dropped for agents that can't reach them.
        let (mcp_servers, built_in_mcp) = acp_mcp_servers(&configured, &init, &cwd);

        let mut loaded = false;
        let mut response = None;
        if let Some(id) = launch
            .previous
            .clone()
            .filter(|_| init.agent_capabilities.load_session)
        {
            let load = || LoadSessionRequest {
                mcp_servers: mcp_servers.clone(),
                ..LoadSessionRequest::new(id.clone(), cwd.clone())
            };
            let result = match conn.load_session(load()).await {
                Err(e) if e.is_auth_required() => {
                    authenticate(&conn, &init).await?;
                    conn.load_session(load()).await
                }
                other => other,
            };
            // A failed load (stale id, agent state gone) falls through
            // to a fresh session rather than failing the launch.
            if let Ok(load_response) = result {
                response = Some(NewSessionResponse {
                    session_id: id.clone().into(),
                    modes: load_response.modes,
                    config_options: load_response.config_options,
                    meta: None,
                });
                loaded = true;
            }
        }

        let response = match response {
            Some(response) => response,
            None => {
                let new_session = || NewSessionRequest {
                    mcp_servers: mcp_servers.clone(),
                    ..NewSessionRequest::new(cwd.clone())
                };
                match conn.new_session(new_session()).await {
                    Ok(response) => response,
                    Err(e) if e.is_auth_required() => {
                        authenticate(&conn, &init).await?;
                        conn.new_session(new_session()).await.map_err(|e| {
                            anyhow!(
                                "session/new failed after authentication: {}",
                                error_text(&e)
                            )
                        })?
                    }
                    Err(e) => return Err(anyhow!("session/new failed: {}", error_text(&e))),
                }
            }
        };

        let history_fork = if !loaded || launch.history_pending {
            if let Some(mut history) = launch.history {
                if init.agent_capabilities.prompt_capabilities.image {
                    history = tokio::task::spawn_blocking(move || {
                        for entry in &mut history {
                            if entry.role == HistoryRole::User {
                                let paths: Vec<_> = entry
                                    .content
                                    .iter()
                                    .flat_map(|block| match block {
                                        ContentBlock::Text(text) => media::attached(&text.text),
                                        _ => Vec::new(),
                                    })
                                    .collect();
                                entry.content.extend(
                                    paths.iter().filter_map(|path| media::encode(path)).map(
                                        |(data, mime)| {
                                            ContentBlock::Image(ImageContent {
                                                data,
                                                mime_type: mime.to_owned(),
                                                uri: None,
                                                annotations: None,
                                                meta: None,
                                            })
                                        },
                                    ),
                                );
                            }
                        }
                        history
                    })
                    .await
                    .map_err(|e| anyhow!("fork history preparation failed: {e}"))?;
                }
                Some(Arc::new(Mutex::new(HistoryFork::restore(
                    conn.clone(),
                    response.clone(),
                    history,
                ))))
            } else {
                None
            }
        } else {
            None
        };

        let mut session = Self {
            conn: Some(conn),
            tx,
            child: Some(child),
            session_id: response.session_id.clone(),
            init,
            response,
            cwd,
            loaded,
            built_in_mcp,
            history_fork,
        };
        session.restore_choices(&launch.choices).await?;
        Ok(session)
    }

    /// Restore supported choices before the first queued prompt can run.
    async fn restore_choices(&mut self, choices: &Choices) -> Result<()> {
        if let Some(mode) = &choices.mode
            && let Some(modes) = &self.response.modes
            && modes
                .available_modes
                .iter()
                .any(|offered| offered.id.to_string() == *mode)
            && modes.current_mode_id.to_string() != *mode
        {
            self.conn()
                .set_session_mode(SetSessionModeRequest {
                    session_id: self.session_id.clone(),
                    mode_id: mode.clone().into(),
                    meta: None,
                })
                .await
                .map_err(|error| {
                    anyhow!("restoring session mode failed: {}", error_text(&error))
                })?;
            if let Some(modes) = &mut self.response.modes {
                modes.current_mode_id = mode.clone().into();
            }
        }
        let mut keys: Vec<_> = self
            .response
            .config_options
            .as_deref()
            .unwrap_or_default()
            .iter()
            .map(|option| {
                (
                    option.category != Some(cacp::schema::SessionConfigOptionCategory::Model),
                    option.id.to_string(),
                )
            })
            .collect();
        keys.sort();
        for (_, id) in keys {
            let Some(value) = choices.config.get(&id) else {
                continue;
            };
            let Some(option) = self
                .response
                .config_options
                .as_deref()
                .unwrap_or_default()
                .iter()
                .find(|option| option.id.to_string() == id)
            else {
                continue;
            };
            if !session_preferences::supports(option, value)
                || session_preferences::current(option) == *value
            {
                continue;
            }
            let response = self
                .conn()
                .set_session_config_option(SetSessionConfigOptionRequest {
                    session_id: self.session_id.clone(),
                    config_id: id.clone().into(),
                    value: value.clone(),
                    meta: None,
                })
                .await
                .map_err(|error| {
                    anyhow!(
                        "restoring session option {id} failed: {}",
                        error_text(&error)
                    )
                })?;
            self.response.config_options = Some(response.config_options);
        }
        Ok(())
    }

    /// A handle on the agent. Cheap to clone, and present for as long as
    /// anything can reach the session — the `Option` is [`Session::drop`]'s.
    fn conn(&self) -> AgentConn {
        self.conn.clone().expect("the session is being dropped")
    }

    /// Send a prompt turn. Its result arrives as [`Event::TurnDone`] —
    /// including a failure to send it at all.
    ///
    /// Pictures the message points at go along as image blocks when the agent
    /// takes them, read and encoded off the UI thread. An agent that does not
    /// still has their paths in the text.
    pub fn prompt(&self, content: &str) {
        let capabilities = &self.init.agent_capabilities.prompt_capabilities;
        let pictures = match capabilities.image {
            true => media::attached(content),
            false => Vec::new(),
        };
        let mut blocks = super::context::prompt(
            &self.cwd,
            self.built_in_mcp,
            capabilities.embedded_context,
            vec![content.to_owned().into()],
        );
        let session_id = self.session_id.clone();
        let conn = self.conn();
        let tx = self.tx.clone();
        let history_fork = self.history_fork.clone();
        runtime().spawn(async move {
            if !pictures.is_empty() {
                let images = tokio::task::spawn_blocking(move || {
                    pictures
                        .iter()
                        .filter_map(|path| media::encode(path))
                        .map(|(data, mime)| {
                            ContentBlock::Image(ImageContent {
                                data,
                                mime_type: mime.to_owned(),
                                uri: None,
                                annotations: None,
                                meta: None,
                            })
                        })
                        .collect::<Vec<_>>()
                })
                .await
                .unwrap_or_default();
                // After the message, ahead of the app's context at the end.
                blocks.splice(1..1, images);
            }
            let request = PromptRequest::new(session_id, blocks);
            let result = match history_fork {
                Some(fork) => fork.lock().await.prompt(request).await,
                None => conn.prompt(request).await,
            };
            let done = result.map(|response| response.stop_reason);
            let _ = tx.send(Event::TurnDone(done));
        });
    }

    /// Cancel the in-flight turn (`session/cancel`). The turn still ends
    /// with an [`Event::TurnDone`] carrying `StopReason::Cancelled`. Pending
    /// permission replies are the frontend's to answer `Cancelled`.
    pub fn cancel(&self) -> Result<(), Error> {
        self.conn().cancel(CancelNotification {
            session_id: self.session_id.clone(),
            meta: None,
        })
    }

    /// Switch the session mode (`session/set_mode`). Fire-and-forget:
    /// frontends validate the id against `response.modes` up front, and
    /// the agent's `CurrentModeUpdate` is the confirmation.
    pub fn set_mode(&self, mode_id: &str) {
        let request = SetSessionModeRequest {
            session_id: self.session_id.clone(),
            mode_id: mode_id.into(),
            meta: None,
        };
        let conn = self.conn();
        runtime().spawn(async move { conn.set_session_mode(request).await });
    }

    /// Set a session config option (`session/set_config_option`).
    /// Fire-and-forget like [`Self::set_mode`]: frontends validate
    /// against `response.config_options`, and the agent's
    /// `ConfigOptionUpdate` is the confirmation.
    pub fn set_config_option(&self, config_id: &str, value: SessionConfigOptionValue) {
        let request = SetSessionConfigOptionRequest {
            session_id: self.session_id.clone(),
            config_id: config_id.into(),
            value,
            meta: None,
        };
        let conn = self.conn();
        runtime().spawn(async move { conn.set_session_config_option(request).await });
    }
}

/// Convert the saved prefix to context without replaying tool executions.
pub fn history(items: &[artifact::session::chat::ChatItem]) -> Vec<HistoryEntry> {
    use artifact::session::chat::ChatItem;
    items
        .iter()
        .map(|item| {
            let (role, text) = match item {
                ChatItem::User(text) => (HistoryRole::User, text.clone()),
                ChatItem::Agent(text) => (HistoryRole::Agent, text.clone()),
                ChatItem::Thinking { text, .. } => {
                    (HistoryRole::Agent, format!("Prior reasoning: {text}"))
                }
                ChatItem::Tool { label, output, .. } => {
                    (HistoryRole::Tool, format!("{label}\n{output}"))
                }
                ChatItem::Process { command, output } => {
                    (HistoryRole::Tool, format!("{command}\n{output}"))
                }
                ChatItem::Notice { text, .. } => {
                    (HistoryRole::Tool, format!("Session notice: {text}"))
                }
            };
            HistoryEntry {
                role,
                content: vec![text.into()],
            }
        })
        .collect()
}

/// Serves what the agent asks of us: file access answered here, anything
/// the user has to see queued for the frontend.
struct Frontend(mpsc::UnboundedSender<Event>);

impl Client for Frontend {
    async fn session_update(&self, notification: SessionNotification) {
        let _ = self.0.send(Event::Update(notification.update));
    }

    async fn request_permission(
        &self,
        request: RequestPermissionRequest,
    ) -> Result<RequestPermissionResponse, Error> {
        let (tx, rx) = oneshot::channel();
        self.0
            .send(Event::Permission(request, Reply(tx)))
            .map_err(|_| Error::internal_error().data("the frontend is gone"))?;
        rx.await.unwrap_or_else(|_| Err(Error::method_not_found()))
    }

    async fn read_text_file(
        &self,
        request: ReadTextFileRequest,
    ) -> Result<ReadTextFileResponse, Error> {
        read_text_file(&request)
    }

    async fn write_text_file(
        &self,
        request: WriteTextFileRequest,
    ) -> Result<WriteTextFileResponse, Error> {
        std::fs::write(&request.path, &request.content)
            .map(|()| WriteTextFileResponse::default())
            .map_err(|e| io_error(&request.path, &e))
    }
}

/// Take the agent down: the connection first, then the process behind it.
///
/// `cacp::spawn` sets `kill_on_drop`, so letting the child field drop on its own
/// is an immediate SIGKILL — mid-request, if the agent was answering one.
/// [`crate::model::session::ChatSession::close`] sends `session/cancel` ahead
/// of this, and the pause here is what gives that notification time to be read.
///
/// It is not a clean shutdown, and cannot be until cacp can close an agent's
/// stdin: its read loop is handed a `Peer` by value, so the write loop holding
/// stdin lives as long as the agent does, whatever this side drops. Until then
/// the kill is the only exit and the agent's stderr is where the noise goes —
/// see [`DEBUG`].
impl Drop for Session {
    fn drop(&mut self) {
        let (Some(conn), Some(mut child)) = (self.conn.take(), self.child.take()) else {
            return;
        };
        drop(conn);
        runtime().spawn(async move {
            let _ = tokio::time::timeout(SHUTDOWN_GRACE, child.wait()).await;
        });
    }
}

/// cacp drops the client when its read loop ends, which is the only notice
/// the frontend gets that the agent is gone.
impl Drop for Frontend {
    fn drop(&mut self) {
        let _ = self.0.send(Event::Closed);
    }
}

// A GPUI (or any multi-threaded) frontend holds `Session` in its UI state
// and moves `Event` — reply included — across executor threads.
const _: () = {
    const fn assert_send<T: Send>() {}
    assert_send::<Session>();
    assert_send::<Event>();
};

fn loopback_bypass(existing: [Option<String>; 2]) -> String {
    let mut entries = Vec::new();
    for value in existing
        .iter()
        .flatten()
        .map(String::as_str)
        .chain(["localhost,127.0.0.1,::1"])
    {
        for entry in value
            .split(',')
            .map(str::trim)
            .filter(|entry| !entry.is_empty())
        {
            if !entries.contains(&entry) {
                entries.push(entry);
            }
        }
    }
    entries.join(",")
}

/// The enabled servers an agent can actually reach, in ACP's shape.
/// Remote servers are dropped for agents that don't advertise HTTP MCP
/// rather than being sent and failing.
///
/// Cydonia's own door goes first, when the agent can reach it. It is not in
/// `mcp.toml` and must not be — that file is the servers the user added, and
/// this one is not the user's to remove.
fn acp_mcp_servers(
    configured: &[mcp::McpServer],
    init: &InitializeResponse,
    cwd: &std::path::Path,
) -> (Vec<McpServer>, bool) {
    let http = init.agent_capabilities.mcp_capabilities.http;
    let ours = http.then(serve::url).flatten().map(|url| {
        McpServer::Http(McpServerHttp {
            name: SERVER.to_owned(),
            url,
            // Which project this session is. The tools then take no directory
            // at all — one a session's model had to supply is one it could
            // supply wrongly, about something already known here.
            headers: vec![{
                let (name, value) = serve::project(cwd);
                HttpHeader {
                    name: name.to_owned(),
                    value,
                    meta: None,
                }
            }],
            meta: None,
        })
    });
    let available = ours.is_some();
    let servers = ours
        .into_iter()
        .chain(
            configured
                .iter()
                .filter(|server| server.enabled)
                .filter_map(|server| match (&server.command, &server.url) {
                    (Some(command), _) => Some(McpServer::Stdio(McpServerStdio {
                        name: server.name.clone(),
                        command: command.into(),
                        args: server.args.clone(),
                        env: server
                            .env
                            .iter()
                            .map(|(name, value)| EnvVariable {
                                name: name.clone(),
                                value: value.clone(),
                                meta: None,
                            })
                            .collect(),
                        meta: None,
                    })),
                    (None, Some(url)) if http => Some(McpServer::Http(McpServerHttp {
                        name: server.name.clone(),
                        url: url.clone(),
                        headers: Vec::new(),
                        meta: None,
                    })),
                    _ => None,
                }),
        )
        .collect();
    (servers, available)
}

/// Try each advertised auth method in order. Non-interactive methods
/// (API keys read from the agent's env) fail fast when unset;
/// interactive ones (OAuth) block until the user completes the flow in
/// the browser the agent opens.
async fn authenticate(conn: &AgentConn, init: &InitializeResponse) -> Result<()> {
    if init.auth_methods.is_empty() {
        return Err(anyhow!(
            "authentication required, but the agent advertises no auth methods"
        ));
    }
    let mut failures = Vec::new();
    for method in &init.auth_methods {
        let request = AuthenticateRequest {
            method_id: method.id().clone(),
            meta: None,
        };
        match conn.authenticate(request).await {
            Ok(_) => return Ok(()),
            Err(e) => failures.push(format!("{}: {}", method.name(), error_text(&e))),
        }
    }
    Err(anyhow!("authentication failed — {}", failures.join("; ")))
}

/// Read the agent's stderr to its end, a line at a time, onto the channel the
/// rest of the session runs on.
///
/// Ends when the pipe closes, which a dead process is what does — so this
/// task is also what lets the channel close behind a launch that failed, and
/// the caller stop waiting on it.
async fn drain(stderr: ChildStderr, tx: Sender) {
    let echo = std::env::var_os(DEBUG).is_some();
    let mut lines = BufReader::new(stderr).lines();
    while let Ok(Some(line)) = lines.next_line().await {
        if echo {
            eprintln!("{line}");
        }
        if tx.send(Event::Stderr(line)).is_err() {
            return;
        }
    }
}

/// Spend a loaded session's replay, keeping only what is state rather than
/// transcript.
///
/// `session/load` replays the whole conversation before it answers, and the
/// client is holding that transcript already — the replay would arrive as a
/// second copy of what is on screen. Only updates can be queued at this point,
/// because nothing else is sent until we prompt.
///
/// What survives is the one thing in a replay that is not transcript: what the
/// conversation has already spent. Nothing in ACP asks for that — it arrives
/// as a notification or not at all — so dropping it with the rest is what
/// leaves a resumed session reading empty until its next turn. The last one
/// wins, and goes back on the channel the frontend is about to read.
///
/// Stderr is left where it is. It is not part of any replay: it is this
/// launch's own process talking, and it has as much right to the transcript
/// here as anywhere.
pub fn spend_replay(events: &mut Events, tx: &Sender) {
    let mut usage = None;
    let mut kept = Vec::new();
    while let Ok(event) = events.try_recv() {
        match event {
            Event::Update(SessionUpdate::UsageUpdate(update)) => usage = Some(update),
            Event::Update(_) => {}
            other => kept.push(other),
        }
    }
    if let Some(update) = usage {
        let _ = tx.send(Event::Update(SessionUpdate::UsageUpdate(update)));
    }
    for event in kept {
        let _ = tx.send(event);
    }
}

/// One-line rendering of a JSON-RPC error (`Display` dumps a JSON blob).
pub fn error_text(e: &Error) -> String {
    match &e.data {
        Some(data) => {
            let detail = data
                .as_str()
                .map(str::to_owned)
                .unwrap_or_else(|| data.to_string());
            format!("{} — {detail}", e.message)
        }
        None => e.message.clone(),
    }
}

/// Serve `fs/read_text_file`: whole file, or 1-based `line` + `limit` window.
fn read_text_file(request: &ReadTextFileRequest) -> Result<ReadTextFileResponse, Error> {
    let content =
        std::fs::read_to_string(&request.path).map_err(|e| io_error(&request.path, &e))?;
    let content = match (request.line, request.limit) {
        (None, None) => content,
        (line, limit) => {
            let skip = line.map(|l| l.saturating_sub(1) as usize).unwrap_or(0);
            let take = limit.map(|l| l as usize).unwrap_or(usize::MAX);
            content
                .lines()
                .skip(skip)
                .take(take)
                .collect::<Vec<_>>()
                .join("\n")
        }
    };
    Ok(ReadTextFileResponse {
        content,
        meta: None,
    })
}

fn io_error(path: &std::path::Path, e: &std::io::Error) -> Error {
    Error::internal_error().data(format!("{}: {e}", path.display()))
}

/// With `CYDONIA_DEBUG=<path>` set, append every JSON-RPC line to that file.
fn debug_tap() -> Option<Tap> {
    let path = std::env::var(DEBUG).ok()?;
    Some(Arc::new(move |direction: Direction, line: &str| {
        use std::io::Write;
        if let Ok(mut f) = std::fs::OpenOptions::new()
            .create(true)
            .append(true)
            .open(&path)
        {
            let _ = writeln!(f, "{direction:?}: {line}");
        }
    }))
}

#[cfg(test)]
#[path = "../../tests/unit/acp_proxy.rs"]
mod proxy_tests;