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}