Skip to main content

lattice_ai/mcp/
server.rs

1//! The IDE server supervisor: an idle tokio task that starts/stops the
2//! loopback WebSocket listener on command.
3//!
4//! Mirrors `lattice_lsp`'s supervisor shape: a single task owns the
5//! listener lifecycle; a clone-able [`ClaudeCodeServerHandle`] drives it
6//! via a non-blocking `cmd_tx` send. The ex-command `apply` closures
7//! (`:claude-code-start` / `:claude-code-stop`) and `claude-code-mode`'s
8//! `on_activate` (I5) hold the handle and call `start` / `stop`.
9//!
10//! All protocol I/O happens off the editor thread: the supervisor task
11//! and one task per connection run on the IDE runtime. The editor thread
12//! only ever calls `start` / `stop` (a channel send) and reads the
13//! wait-free [`ServerState`] snapshot.
14
15use std::path::PathBuf;
16use std::sync::Arc;
17use std::sync::atomic::{AtomicUsize, Ordering};
18
19use tokio::sync::Notify;
20
21use arc_swap::ArcSwap;
22use futures::{SinkExt, StreamExt};
23use tokio::net::{TcpListener, TcpStream};
24use tokio::sync::{broadcast, mpsc};
25use tokio::task::JoinHandle;
26use tokio_tungstenite::tungstenite::Message as WsMessage;
27
28use lattice_diff::ProgrammaticDiffBus;
29use lattice_mode::inbound::InboundBus;
30use lattice_runtime::EventBus;
31
32use lattice_agent::{EditorAccess, EditorStateHandle, EditorWriteRequest};
33
34use crate::mcp::dispatch::{self, DispatchContext, Outgoing};
35use crate::mcp::error::Result;
36use crate::mcp::lockfile::{Lockfile, LockfileContents};
37use crate::mcp::reads::ReadContext;
38use crate::mcp::{auth, transport};
39
40/// IDE name reported in the discovery lockfile.
41const IDE_NAME: &str = "Lattice";
42
43/// I6: capacity of the server-initiated notification broadcast. A connection
44/// that falls this far behind a burst skips the dropped frames (acceptable —
45/// selection notifications coalesce to "latest wins").
46const NOTIFY_CAPACITY: usize = 64;
47
48/// Static config the server binds with.
49#[derive(Debug, Clone)]
50pub struct ServerConfig {
51    /// Absolute paths advertised as workspace folders.
52    pub workspace_folders: Vec<String>,
53    /// Directory the discovery lockfile is written into (`~/.claude/ide`
54    /// in production; a temp dir in tests).
55    pub lock_dir: PathBuf,
56}
57
58/// Wait-free snapshot of server state for status reads (headerline, I7).
59#[derive(Debug, Clone, Default)]
60pub struct ServerState {
61    /// Whether the listener is currently bound + accepting.
62    pub running: bool,
63    /// The bound loopback port, when running.
64    pub port: Option<u16>,
65}
66
67/// Command to the supervisor task.
68enum ServerCmd {
69    /// I5.1: the pre-bound listener (bound synchronously in [`start`] so the
70    /// caller learns the port immediately) + the auth token. The supervisor
71    /// writes the discovery lockfile, wraps the listener for tokio, and runs
72    /// the accept loop.
73    ///
74    /// [`start`]: ClaudeCodeServerHandle::start
75    Start {
76        listener: std::net::TcpListener,
77        token: String,
78    },
79    Stop,
80}
81
82/// Handle to the supervisor. Clone-able + `Send + Sync`: the ex-command
83/// `apply` closures and `claude-code-mode`'s `on_activate` hold one and
84/// drive `start` / `stop`. Mirrors `lattice_lsp::LspSupervisorHandle`.
85#[derive(Clone)]
86pub struct ClaudeCodeServerHandle {
87    cmd_tx: mpsc::UnboundedSender<ServerCmd>,
88    state: Arc<ArcSwap<ServerState>>,
89    /// The crate-owned read cache (subscription set up at spawn). Shared
90    /// into the [`DispatchContext`] built by `install_read_services`.
91    cache: EditorStateHandle,
92    /// Workspace folders from the config (for `getWorkspaceFolders`).
93    workspace_folders: Vec<String>,
94    /// The live dispatch context connections read. Starts with no generic
95    /// services (cache + config only); `install_services` upgrades it once
96    /// boot has wired the buffer-store / diagnostics handles + the write bus.
97    dispatch_ctx: Arc<ArcSwap<DispatchContext>>,
98    /// I6: broadcast channel for server-initiated notification frames. The
99    /// notification task (`notifications.rs`) + `:claude-send` publish frames
100    /// here; each connection subscribes a receiver and forwards frames to its
101    /// WS writer. A dropped/lagged receiver is pruned by the channel itself —
102    /// no manual connection registry.
103    notify_tx: broadcast::Sender<String>,
104    /// I7: the connection counter + status wake (drives the modeline segment).
105    signals: StatusSignals,
106    /// I7: buffers showing the `claude-code` status segment (the agent
107    /// terminals). `claude-code-mode`'s `on_activate` registers its buffer; the
108    /// Guard unregisters on deactivate. The status publisher reads this set.
109    ide_buffers: crate::mcp::status::IdeBuffers,
110    /// D-fix.6 follow-up: the shared pending-review tracker. Held so
111    /// `install_services` can re-seat it on the rebuilt dispatch context (it
112    /// must be the SAME instance the publisher reads + `openDiff` increments).
113    review: crate::mcp::status::ReviewHandle,
114    /// D-fix.6 follow-up: the transient `@sent` echo tracker; pinged by
115    /// `:claude-send` via [`Self::ping_mention`].
116    mention: crate::mcp::status::MentionHandle,
117}
118
119impl ClaudeCodeServerHandle {
120    /// Start the server and return the bound loopback **port**, or `None` on
121    /// failure. Idempotent: a second call while already running returns the
122    /// existing port without re-binding.
123    ///
124    /// I5.1: the listener is pre-bound *synchronously here* (not async on the
125    /// supervisor) so `:claude` learns the port immediately and can inject
126    /// `CLAUDE_CODE_SSE_PORT` into the agent's environment before spawning it.
127    /// The bind + the supervisor's subsequent lockfile write are one-shot (a
128    /// user command, never the render loop), so the brief sync work is fine.
129    pub fn start(&self) -> Option<u16> {
130        let current = self.state.load();
131        if current.running {
132            return current.port; // idempotent — already bound
133        }
134        let listener = match std::net::TcpListener::bind(("127.0.0.1", 0)) {
135            Ok(l) => l,
136            Err(e) => {
137                tracing::debug!(error = %e, "claude-code: pre-bind failed");
138                return None;
139            }
140        };
141        let port = listener.local_addr().ok()?.port();
142        if let Err(e) = listener.set_nonblocking(true) {
143            tracing::debug!(error = %e, "claude-code: set_nonblocking failed");
144            return None;
145        }
146        let token = match auth::generate_token() {
147            Ok(t) => t,
148            Err(e) => {
149                tracing::debug!(error = %e, "claude-code: token generation failed");
150                return None;
151            }
152        };
153        // Optimistic running state so a re-entrant `start()` doesn't double-bind;
154        // the supervisor rolls it back if it can't take the listener over.
155        self.state.store(Arc::new(ServerState {
156            running: true,
157            port: Some(port),
158        }));
159        self.signals.fire(); // repaint the status segment (running + port)
160        if self
161            .cmd_tx
162            .send(ServerCmd::Start { listener, token })
163            .is_err()
164        {
165            // Supervisor gone — roll the optimistic state back.
166            self.state.store(Arc::new(ServerState::default()));
167            return None;
168        }
169        Some(port)
170    }
171
172    /// Request the server stop (unbind + unlink lockfile + drop conns).
173    /// Idempotent, non-blocking.
174    pub fn stop(&self) {
175        let _ = self.cmd_tx.send(ServerCmd::Stop);
176    }
177
178    /// Current state snapshot (wait-free `Arc` load).
179    pub fn snapshot(&self) -> Arc<ServerState> {
180        self.state.load_full()
181    }
182
183    /// I6: broadcast a server-initiated notification frame to every connected
184    /// agent. A no-op (the frame is dropped) when no connections are
185    /// subscribed. Used by `:claude-send`; the notification task uses its own
186    /// [`Self::notify_sender`] clone.
187    pub fn notify(&self, frame: String) {
188        let _ = self.notify_tx.send(frame);
189    }
190
191    /// D-fix.6 follow-up: flash the transient `@sent` echo on the modeline —
192    /// called by `:claude-send` after it broadcasts an at-mention. The echo
193    /// clears itself after a few seconds (the status publisher's timed wake).
194    pub fn ping_mention(&self) {
195        self.mention.ping();
196    }
197
198    /// I6.1: a clone of the broadcast sender for the notification task to
199    /// publish `selection_changed` / `didChangeActiveEditor` frames through.
200    pub fn notify_sender(&self) -> broadcast::Sender<String> {
201        self.notify_tx.clone()
202    }
203
204    /// I7: number of currently-connected agents. Surfaced in the
205    /// `claude-code-mode` status. Counted explicitly (a `ConnGuard` per
206    /// connection) so it is exact the instant a connection ends, rather than
207    /// lagging the broadcast receiver drop.
208    pub fn connection_count(&self) -> usize {
209        self.signals.conn_count.load(Ordering::Relaxed)
210    }
211
212    /// I7: register `buf` (an agent terminal) to show the `claude-code` status
213    /// segment, and wake the publisher to paint it immediately. Called from
214    /// `claude-code-mode`'s `on_activate`.
215    pub fn register_status_buffer(&self, buf: lattice_core::BufferId) {
216        self.ide_buffers
217            .lock()
218            .unwrap_or_else(|e| e.into_inner())
219            .insert(buf);
220        self.signals.fire();
221    }
222
223    /// I7: stop showing the status on `buf` (the mode's Guard `Drop` on
224    /// deactivate). The publisher clears the element on its next wake.
225    pub fn unregister_status_buffer(&self, buf: lattice_core::BufferId) {
226        self.ide_buffers
227            .lock()
228            .unwrap_or_else(|e| e.into_inner())
229            .remove(&buf);
230        self.signals.fire();
231    }
232
233    /// BC.3b: the crate-owned read cache. `install()` clones it to build the
234    /// inbound handler ([`lattice_agent::make_handler`]) — the per-item closure
235    /// the generic `boot.inbound` bus drains, which maps write requests against
236    /// the same cache the read tools snapshot.
237    pub fn read_cache(&self) -> EditorStateHandle {
238        self.cache.clone()
239    }
240
241    /// I2.2 + I3.2 / BC.3b: seat the generic read handles + the write bus into
242    /// the live dispatch context (connections read it wait-free). `writes` is
243    /// the generic [`InboundBus`] built by `install()` via `boot.inbound` (which
244    /// owns the channel, the per-tick drain, and the off-keystroke wake — the
245    /// drain's registration token rides `boot.into_registrations()` into the
246    /// Editor). AG-2b: the bus now rides inside the `EditorAccess` the write
247    /// tools call through, not a separate `DispatchContext` field. The server
248    /// handle is spawned before boot wires these handles, so this upgrade runs
249    /// once, from the subsystem's `install(boot)`.
250    pub fn install_services(
251        &self,
252        buffer_store: Option<lattice_mode::BufferStoreHandle>,
253        diagnostics: Option<lattice_lsp::modes::DiagnosticsQueryHandle>,
254        writes: InboundBus<EditorWriteRequest>,
255        diff: Option<ProgrammaticDiffBus>,
256    ) {
257        self.dispatch_ctx.store(Arc::new(DispatchContext {
258            // Shared template; `serve_connection` clones this and stamps each
259            // connection's own `conn_id` (D-fix.6).
260            conn_id: 0,
261            reads: ReadContext {
262                editor: EditorAccess::new(
263                    self.cache.clone(),
264                    buffer_store,
265                    self.workspace_folders.clone(),
266                    Some(writes),
267                ),
268                diagnostics,
269            },
270            diff,
271            // D-fix.6 follow-up: re-seat the SAME review tracker so the badge
272            // count is shared across the boot rebuild.
273            review: self.review.clone(),
274        }));
275    }
276}
277
278/// Spawn the supervisor task on `rt`. Returns immediately; the task stays
279/// idle until `start()`.
280pub fn spawn(
281    config: ServerConfig,
282    event_bus: Arc<EventBus>,
283    rt: &tokio::runtime::Handle,
284) -> ClaudeCodeServerHandle {
285    let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
286    let state = Arc::new(ArcSwap::from_pointee(ServerState::default()));
287    // I2.1: the read cache subscribes to the generic event bus here.
288    let cache = lattice_agent::state_cache::spawn_read_cache(&event_bus, rt);
289    // I7 / D-fix.6 follow-up: the status wake + the shared pending-review
290    // tracker. Created up front so the dispatch context, the publisher, and
291    // `install_services`'s rebuilt context all share ONE tracker (the count is
292    // global across connections).
293    let signals = StatusSignals::new();
294    let review = crate::mcp::status::ReviewState::new(signals.changed.clone());
295    // D-fix.6 follow-up: the transient `@sent` echo tracker (its `ping` wakes
296    // the same status `changed`).
297    let mention = crate::mcp::status::MentionState::new(signals.changed.clone());
298    // I2.2: start with a deps-less dispatch context (cache + config only);
299    // `install_read_services` upgrades it once boot wires the generic
300    // buffer-store / diagnostics handles.
301    let dispatch_ctx = Arc::new(ArcSwap::from_pointee(DispatchContext {
302        // Shared template; per-connection `conn_id` stamped in
303        // `serve_connection` (D-fix.6).
304        conn_id: 0,
305        reads: ReadContext {
306            // I3.2 wires the real inbound bus (via `EditorAccess`'s `writes`
307            // field) here; until then write tools report a graceful "not
308            // initialized".
309            editor: EditorAccess::new(cache.clone(), None, config.workspace_folders.clone(), None),
310            diagnostics: None,
311        },
312        // I4: the openDiff bus is wired by `install_services`; until then
313        // `openDiff` reports a graceful "not initialized".
314        diff: None,
315        review: review.clone(),
316    }));
317    // I6: the server-initiated notification broadcast. Bounded — a lagged
318    // connection skips dropped frames (coalescing is fine for selection
319    // notifications, where only the latest matters).
320    let (notify_tx, _) = broadcast::channel::<String>(NOTIFY_CAPACITY);
321    // I6.1: the notification task — coalesces SelectionsChanged + broadcasts
322    // selection_changed / didChangeActiveEditor frames. Crate-owned (reads the
323    // same generic event bus + read cache the read tools use).
324    crate::mcp::notifications::spawn_notifier(&event_bus, notify_tx.clone(), cache.clone(), rt);
325    // I7: the modeline status segment. The publisher republishes running/port +
326    // conn-count + project + pending-review badge to each registered IDE buffer
327    // when the wake fires (start/stop, a connection open/close, a review
328    // begin/end, a buffer register/unregister).
329    let ide_buffers: crate::mcp::status::IdeBuffers =
330        Arc::new(std::sync::Mutex::new(std::collections::HashSet::new()));
331    crate::mcp::status::spawn_status_publisher(
332        event_bus.clone(),
333        state.clone(),
334        signals.conn_count.clone(),
335        // The project the agent runs for — the workspace basename, static
336        // for the server's lifetime.
337        crate::mcp::status::project_name(&config.workspace_folders),
338        review.count_handle(),
339        mention.clone(),
340        ide_buffers.clone(),
341        signals.changed.clone(),
342        rt,
343    );
344    rt.spawn(supervisor_main(
345        config.clone(),
346        cmd_rx,
347        state.clone(),
348        dispatch_ctx.clone(),
349        notify_tx.clone(),
350        signals.clone(),
351    ));
352    ClaudeCodeServerHandle {
353        cmd_tx,
354        state,
355        cache,
356        workspace_folders: config.workspace_folders,
357        dispatch_ctx,
358        notify_tx,
359        signals,
360        ide_buffers,
361        review,
362        mention,
363    }
364}
365
366/// I7: the live connection counter + the "status changed" wake, shared between
367/// the accept path (which bumps the count), the handle (start/stop + buffer
368/// (un)register fire the wake), and the status publisher (reads both).
369#[derive(Clone)]
370struct StatusSignals {
371    conn_count: Arc<AtomicUsize>,
372    changed: Arc<Notify>,
373}
374
375impl StatusSignals {
376    fn new() -> Self {
377        Self {
378            conn_count: Arc::new(AtomicUsize::new(0)),
379            changed: Arc::new(Notify::new()),
380        }
381    }
382
383    /// Wake the status publisher. `notify_one` (not `notify_waiters`) so a fire
384    /// that lands *before* the single publisher task parks on `notified()`
385    /// stores a permit and isn't lost — the publisher then wakes immediately on
386    /// its next await. There is exactly one publisher task, so one permit is
387    /// enough; bursts coalesce (each wake re-reads the live state).
388    fn fire(&self) {
389        self.changed.notify_one();
390    }
391}
392
393/// Decrements the live connection count + fires the status wake when a
394/// connection task ends. `Drop` runs on a normal end, an error, or a panic, so
395/// the count can never leak high. Held by `serve_connection`.
396struct ConnGuard(StatusSignals);
397
398impl Drop for ConnGuard {
399    fn drop(&mut self) {
400        self.0.conn_count.fetch_sub(1, Ordering::SeqCst);
401        self.0.fire();
402    }
403}
404
405/// A bound listener + its lockfile. Dropping aborts the accept loop and
406/// (via the lockfile's `Drop`) unlinks the discovery file.
407struct RunningServer {
408    _lockfile: Lockfile,
409    accept_task: JoinHandle<()>,
410    /// I7: clean teardown on stop/quit. Live connections hold a receiver
411    /// subscribed off this sender; dropping it (when the server stops) makes
412    /// their `recv()` return `Closed`, which the connection's read loop selects
413    /// on to close the socket — so `:claude-code-stop` actually disconnects the
414    /// agent instead of leaving it functional against a stopped server.
415    _shutdown_tx: broadcast::Sender<()>,
416}
417
418impl Drop for RunningServer {
419    fn drop(&mut self) {
420        self.accept_task.abort();
421    }
422}
423
424async fn supervisor_main(
425    config: ServerConfig,
426    mut cmd_rx: mpsc::UnboundedReceiver<ServerCmd>,
427    state: Arc<ArcSwap<ServerState>>,
428    dispatch_ctx: Arc<ArcSwap<DispatchContext>>,
429    notify_tx: broadcast::Sender<String>,
430    signals: StatusSignals,
431) {
432    let mut running: Option<RunningServer> = None;
433    while let Some(cmd) = cmd_rx.recv().await {
434        match cmd {
435            ServerCmd::Start { listener, token } => {
436                if running.is_some() {
437                    continue; // idempotent (start() also guards via state)
438                }
439                match start_accepting(
440                    &config,
441                    listener,
442                    token,
443                    dispatch_ctx.clone(),
444                    notify_tx.clone(),
445                    signals.clone(),
446                ) {
447                    Ok(server) => {
448                        // `start()` already published running+port optimistically.
449                        running = Some(server);
450                        tracing::info!("claude-code IDE server accepting");
451                    }
452                    Err(e) => {
453                        // Roll back the optimistic running state set by `start()`.
454                        state.store(Arc::new(ServerState::default()));
455                        tracing::debug!(error = %e, "claude-code IDE server failed to start");
456                    }
457                }
458            }
459            ServerCmd::Stop => {
460                if running.take().is_some() {
461                    // RunningServer::drop aborts accept + unlinks lockfile.
462                    state.store(Arc::new(ServerState::default()));
463                    signals.fire(); // hide the status segment (server stopped)
464                    tracing::info!("claude-code IDE server stopped");
465                }
466            }
467        }
468    }
469}
470
471/// I5.1: take over the pre-bound listener — write the discovery lockfile, wrap
472/// the listener for tokio, and spawn the accept loop. Runs on the supervisor
473/// task (inside the IDE runtime), so `from_std` + `tokio::spawn` have a runtime
474/// context. The std listener was already set non-blocking in [`start`].
475///
476/// [`start`]: ClaudeCodeServerHandle::start
477fn start_accepting(
478    config: &ServerConfig,
479    std_listener: std::net::TcpListener,
480    token: String,
481    dispatch_ctx: Arc<ArcSwap<DispatchContext>>,
482    notify_tx: broadcast::Sender<String>,
483    signals: StatusSignals,
484) -> Result<RunningServer> {
485    let port = std_listener.local_addr()?.port();
486    let lockfile = Lockfile::write(
487        &config.lock_dir,
488        port,
489        &LockfileContents {
490            pid: std::process::id(),
491            workspace_folders: config.workspace_folders.clone(),
492            ide_name: IDE_NAME.to_string(),
493            transport: "ws".to_string(),
494            auth_token: token.clone(),
495            running_in_windows: false,
496        },
497    )?;
498    let listener = TcpListener::from_std(std_listener)?;
499    // I7: the shutdown signal — held here (in `RunningServer`) so it lives as
500    // long as the server runs; the accept loop subscribes a receiver per
501    // connection. Dropping `RunningServer` on stop drops this sender → live
502    // connections see `Closed` and close.
503    let (shutdown_tx, _) = broadcast::channel::<()>(1);
504    let accept_task = tokio::spawn(accept_loop(
505        listener,
506        token,
507        dispatch_ctx,
508        notify_tx,
509        shutdown_tx.clone(),
510        signals,
511    ));
512    Ok(RunningServer {
513        _lockfile: lockfile,
514        accept_task,
515        _shutdown_tx: shutdown_tx,
516    })
517}
518
519async fn accept_loop(
520    listener: TcpListener,
521    token: String,
522    dispatch_ctx: Arc<ArcSwap<DispatchContext>>,
523    notify_tx: broadcast::Sender<String>,
524    shutdown_tx: broadcast::Sender<()>,
525    signals: StatusSignals,
526) {
527    // D-fix.6: monotonic per-connection id. The accept loop is a single
528    // sequential task, so a plain counter (no atomic) is race-free; `0` is
529    // reserved for the shared/boot context + non-IDE diff producers, so start
530    // at 1. Wraps after u64::MAX connections (never reached in practice).
531    let mut next_conn_id: u64 = 1;
532    loop {
533        match listener.accept().await {
534            Ok((stream, _addr)) => {
535                let conn_id = next_conn_id;
536                next_conn_id = next_conn_id.wrapping_add(1).max(1);
537                let token = token.clone();
538                let ctx = dispatch_ctx.clone();
539                // I6: each connection subscribes its own broadcast receiver.
540                let notify_rx = notify_tx.subscribe();
541                // I7: and a shutdown receiver — closed when the server stops.
542                let shutdown_rx = shutdown_tx.subscribe();
543                // I7: bump the live connection count + wake the status segment;
544                // the `ConnGuard` decrements + wakes again when this connection
545                // ends (drop runs on a normal end, error, or panic).
546                signals.conn_count.fetch_add(1, Ordering::SeqCst);
547                signals.fire();
548                let conn_guard = ConnGuard(signals.clone());
549                tokio::spawn(async move {
550                    if let Err(e) = serve_connection(
551                        stream,
552                        token,
553                        ctx,
554                        conn_id,
555                        notify_rx,
556                        shutdown_rx,
557                        conn_guard,
558                    )
559                    .await
560                    {
561                        tracing::debug!(error = %e, "claude-code connection ended with error");
562                    }
563                });
564            }
565            Err(e) => {
566                tracing::debug!(error = %e, "claude-code accept error");
567                // Avoid a busy-spin if accept keeps failing.
568                tokio::task::yield_now().await;
569            }
570        }
571    }
572}
573
574async fn serve_connection(
575    stream: TcpStream,
576    token: String,
577    dispatch_ctx: Arc<ArcSwap<DispatchContext>>,
578    // D-fix.6: this connection's unique id, stamped into the per-connection
579    // dispatch context so `openDiff` tags its diff and the close tools scope
580    // teardown to THIS session.
581    conn_id: u64,
582    mut notify_rx: broadcast::Receiver<String>,
583    mut shutdown_rx: broadcast::Receiver<()>,
584    // I7: held for the connection's lifetime; its `Drop` decrements the live
585    // connection count + wakes the status segment when this connection ends.
586    _conn_guard: ConnGuard,
587) -> Result<()> {
588    let ws = transport::accept(stream, &token).await?;
589    let (mut write, mut read) = ws.split();
590    // Load the current dispatch context once per connection (installed at
591    // boot before any start, so it carries the generic read services), and
592    // stamp THIS connection's id onto it (D-fix.6). Cheap clone — all fields
593    // are Arc-based handles or small values.
594    let mut ctx = (*dispatch_ctx.load_full()).clone();
595    ctx.conn_id = conn_id;
596
597    // I6: one WS sink can't be written from two tasks, so a single outbound
598    // channel funnels BOTH request responses (from the read loop) and pushed
599    // server-initiated notifications (from the broadcast) to one writer task.
600    let (out_tx, mut out_rx) = mpsc::unbounded_channel::<String>();
601
602    // Writer: drain the outbound channel → WS. Ends when every `out_tx` clone
603    // drops (read loop done + forwarder gone) or the peer write fails.
604    let writer = tokio::spawn(async move {
605        while let Some(payload) = out_rx.recv().await {
606            if write.send(WsMessage::Text(payload)).await.is_err() {
607                break; // peer gone
608            }
609        }
610    });
611
612    // Forwarder: broadcast → outbound. A `Lagged` receiver (fell behind a
613    // burst) skips the dropped frames — coalescing is the intended behaviour
614    // for selection notifications (latest wins).
615    let notif_out = out_tx.clone();
616    let forwarder = tokio::spawn(async move {
617        loop {
618            match notify_rx.recv().await {
619                Ok(frame) => {
620                    if notif_out.send(frame).is_err() {
621                        break; // writer gone
622                    }
623                }
624                Err(broadcast::error::RecvError::Lagged(_)) => continue,
625                Err(broadcast::error::RecvError::Closed) => break,
626            }
627        }
628    });
629
630    // Read loop: handle incoming MCP frames; responses ride the same writer.
631    // The blocking `openDiff` tool is dispatched on its OWN task (tracked in
632    // `blocking`) so its unbounded, no-timeout verdict-wait never blocks the
633    // loop — the loop keeps polling `read.next()` + `shutdown_rx`, so a closed
634    // socket or `:claude-code-stop` is observed promptly even while a review is
635    // pending, and the agent can't be stranded forever. Every non-blocking tool
636    // stays inline + ordered. Capture the result so teardown runs on every path.
637    let mut blocking: tokio::task::JoinSet<()> = tokio::task::JoinSet::new();
638    let result: Result<()> = 'read: loop {
639        let frame = tokio::select! {
640            maybe = read.next() => match maybe {
641                Some(Ok(f)) => f,
642                Some(Err(e)) => break 'read Err(e.into()),
643                None => break 'read Ok(()), // peer closed the socket
644            },
645            // I7: the server stopped — the shutdown sender dropped, so
646            // `recv()` returns `Closed`. Either arm closes this connection.
647            _ = shutdown_rx.recv() => break 'read Ok(()),
648        };
649        if frame.is_close() {
650            break 'read Ok(());
651        }
652        // MCP frames are JSON text; ignore binary / ping / pong (pings are
653        // auto-ponged by the stream's read machinery).
654        let Ok(text) = frame.to_text() else {
655            continue;
656        };
657        if dispatch::is_blocking_tool_call(text.as_bytes()) {
658            // The verdict-wait posts its response (correlated by request id)
659            // when the user resolves the diff — or when teardown below rejects
660            // it. Meanwhile the loop moves straight on to the next frame.
661            let ctx = ctx.clone();
662            let out = out_tx.clone();
663            let text = text.to_owned();
664            blocking.spawn(async move {
665                for outgoing in dispatch::dispatch_frame(text.as_bytes(), &ctx).await {
666                    if let Some(payload) = serialize_outgoing(&outgoing) {
667                        let _ = out.send(payload);
668                    }
669                }
670            });
671        } else {
672            for outgoing in dispatch::dispatch_frame(text.as_bytes(), &ctx).await {
673                let Some(payload) = serialize_outgoing(&outgoing) else {
674                    continue;
675                };
676                if out_tx.send(payload).is_err() {
677                    break 'read Ok(()); // writer gone — connection is finished
678                }
679            }
680        }
681    };
682
683    // Teardown (runs on socket close, error, or stop). A connection that drops
684    // mid-review must not strand the agent or leak diff panes: reject + close
685    // every diff THIS connection opened (idempotent — a no-op when it opened
686    // none), so the host fires `DIFF_REJECTED` + tears the transient panes down.
687    // Then abort any in-flight `openDiff` handlers (their reply can't reach a
688    // dead socket anyway) so they release their `out_tx` clones, and drain the
689    // writer.
690    {
691        // AG-2b: the write bus rides inside `reads.editor`; a `None` bus
692        // (test harness, or a server never fully installed) degrades to a
693        // graceful no-op inside `close_all_diff_tabs`, so this is safe to
694        // always spawn.
695        let editor = ctx.reads.editor.clone();
696        tokio::spawn(async move {
697            let _ = crate::mcp::writes::close_all_diff_tabs(&editor, conn_id).await;
698        });
699    }
700    blocking.abort_all();
701    drop(out_tx);
702    forwarder.abort();
703    let _ = writer.await;
704    result
705}
706
707/// Serialize one outbound frame for the WS writer. A serialization failure
708/// (never expected for a well-formed `Response`/`Notification`) drops just
709/// that frame rather than tearing the connection down.
710fn serialize_outgoing(outgoing: &Outgoing) -> Option<String> {
711    match outgoing {
712        Outgoing::Response(r) => serde_json::to_string(r).ok(),
713        Outgoing::Notification(n) => serde_json::to_string(n).ok(),
714    }
715}
716
717#[cfg(test)]
718mod tests {
719    #![allow(clippy::unwrap_used)]
720    use super::*;
721
722    #[tokio::test]
723    async fn install_services_seats_read_handles_and_bus() {
724        use lattice_agent::make_handler;
725        use lattice_mode::inbound::make_inbound;
726        use tokio::sync::Notify;
727
728        let config = ServerConfig {
729            workspace_folders: vec![],
730            lock_dir: std::env::temp_dir(),
731        };
732        let handle = spawn(
733            config,
734            Arc::new(EventBus::new()),
735            &tokio::runtime::Handle::current(),
736        );
737        // The generic inbound bus, as `install()` builds it via `boot.inbound`.
738        let (bus, mut drain) = make_inbound::<EditorWriteRequest, _>(
739            Arc::new(Notify::new()),
740            make_handler(handle.read_cache()),
741        );
742        handle.install_services(None, None, bus, None);
743
744        // AG-2b: the write bus now rides inside the seated `EditorAccess`
745        // (no separate `DispatchContext::writes` field). Prove it's wired by
746        // sending a write through the live context and observing it reach
747        // the drain, rather than the graceful "not initialized" error a
748        // `None` bus would produce.
749        let editor = handle.dispatch_ctx.load().reads.editor.clone();
750        let call = tokio::spawn(async move {
751            editor
752                .open_file(std::path::PathBuf::from("/a.rs"), None)
753                .await
754        });
755        let mut effects = Vec::new();
756        for _ in 0..50 {
757            effects = drain();
758            if !effects.is_empty() {
759                break;
760            }
761            tokio::time::sleep(std::time::Duration::from_millis(5)).await;
762        }
763        assert_eq!(
764            effects.len(),
765            1,
766            "the write bus is installed and reaches the drain"
767        );
768        assert!(
769            call.await.unwrap().is_ok(),
770            "the write completes successfully"
771        );
772    }
773
774    /// I5.1a: `start()` pre-binds synchronously and returns the bound port; a
775    /// second call while running is idempotent (same port, no re-bind).
776    #[tokio::test]
777    async fn start_returns_port_and_is_idempotent() {
778        let handle = spawn(
779            ServerConfig {
780                workspace_folders: vec![],
781                lock_dir: std::env::temp_dir(),
782            },
783            Arc::new(EventBus::new()),
784            &tokio::runtime::Handle::current(),
785        );
786        let p1 = handle.start().expect("start binds + returns a port");
787        let p2 = handle.start().expect("idempotent re-start returns a port");
788        assert_eq!(p1, p2, "second start returns the same port (no re-bind)");
789        let snap = handle.snapshot();
790        assert_eq!(snap.port, Some(p1), "snapshot reflects the bound port");
791        assert!(snap.running, "server is running after start");
792        handle.stop();
793    }
794}