Skip to main content

lattice_ai/mcp/
notifications.rs

1//! IDE-protocol I6: server-initiated notifications (server → agent).
2//!
3//! A task subscribes to the generic event bus (`SelectionsChanged`), coalesces
4//! bursts (cursor moves fire at ~30 Hz under a held key), and broadcasts
5//! `selection_changed` — plus a `didChangeActiveEditor` when the active buffer
6//! changes — to every connected agent through the server's broadcast sender
7//! ([`crate::mcp::server::ClaudeCodeServerHandle::notify_sender`]). Each connection's
8//! forwarder relays the frame to its WS writer; a lagged connection skips
9//! dropped frames (coalescing — latest wins).
10//!
11//! Crate-owned per `feedback_mode_owns_its_surface`: the host only publishes
12//! the generic events. The editor thread pays nothing new — it already
13//! `publish`es `SelectionsChanged`; the coalescing + framing run off-thread.
14//!
15//! **The frame SHAPE is PROVISIONAL** until validated against a live `claude`
16//! CLI — server→agent notification formats are less documented than the tool
17//! replies. In particular `selection.start/end.character` carries the editor's
18//! **byte** offset within the line (lattice's `Position.byte`); the VS Code
19//! contract is a UTF-16 character offset. Carrying it verbatim mirrors the I3
20//! selection-encoding caveat; the conversion lands once a live CLI pins it.
21
22use std::path::Path;
23use std::sync::Arc;
24
25use lattice_agent::EditorStateHandle;
26use lattice_protocol::ids::DocumentId;
27use lattice_protocol::position::Position;
28use lattice_protocol::{Event, EventKind, SelectionSet};
29use lattice_runtime::{EventBus, EventFilter, SubscriptionTarget};
30use serde_json::json;
31use tokio::sync::broadcast;
32
33/// Subscribe to `SelectionsChanged` and spawn the notification task. The task
34/// coalesces bursts and broadcasts notification frames through `notify_tx`.
35pub fn spawn_notifier(
36    bus: &Arc<EventBus>,
37    notify_tx: broadcast::Sender<String>,
38    cache: EditorStateHandle,
39    rt: &tokio::runtime::Handle,
40) {
41    let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Event>();
42    bus.subscribe(
43        EventFilter::kinds(vec![EventKind::SelectionsChanged]),
44        SubscriptionTarget::Channel(tx),
45    );
46    rt.spawn(async move {
47        let mut last_active: Option<DocumentId> = None;
48        while let Some(first) = rx.recv().await {
49            // Coalesce a burst — keep only the latest selection (cursor moves
50            // fire per keystroke; the agent only needs where the cursor is now).
51            let mut latest = first;
52            while let Ok(next) = rx.try_recv() {
53                latest = next;
54            }
55            let Event::SelectionsChanged { id, selections, .. } = &latest else {
56                continue; // the filter guarantees this, but stay total
57            };
58
59            // Resolve the buffer's path from the read cache (the same cache the
60            // read tools answer from).
61            let path = {
62                let guard = cache.lock().unwrap_or_else(|e| e.into_inner());
63                guard.open_buffers.get(id).and_then(|b| b.path.clone())
64            };
65
66            // On an active-buffer change, announce it first.
67            if last_active != Some(*id) {
68                last_active = Some(*id);
69                let _ = notify_tx.send(did_change_active_editor_frame(path.as_deref()));
70            }
71            let _ = notify_tx.send(selection_changed_frame(selections, path.as_deref()));
72        }
73    });
74}
75
76/// Build the `selection_changed` notification frame (PROVISIONAL shape).
77pub fn selection_changed_frame(selections: &SelectionSet, path: Option<&Path>) -> String {
78    let sel = selections.primary();
79    let (start, end) = ordered(sel.anchor, sel.head);
80    json!({
81        "jsonrpc": "2.0",
82        "method": "selection_changed",
83        "params": {
84            "filePath": path.map(display),
85            "selection": {
86                "start": { "line": start.line, "character": start.byte },
87                "end": { "line": end.line, "character": end.byte },
88                "isEmpty": start == end,
89            }
90        }
91    })
92    .to_string()
93}
94
95/// Build the `didChangeActiveEditor` notification frame (PROVISIONAL shape).
96pub fn did_change_active_editor_frame(path: Option<&Path>) -> String {
97    json!({
98        "jsonrpc": "2.0",
99        "method": "didChangeActiveEditor",
100        "params": { "filePath": path.map(display) }
101    })
102    .to_string()
103}
104
105/// Build the `at_mentioned` notification frame (PROVISIONAL shape) — the user
106/// pushing the current file + selected line range into the agent's context via
107/// `:claude-send` / `@`.
108pub fn at_mentioned_frame(selections: &SelectionSet, path: Option<&Path>) -> String {
109    let sel = selections.primary();
110    let (start, end) = ordered(sel.anchor, sel.head);
111    json!({
112        "jsonrpc": "2.0",
113        "method": "at_mentioned",
114        "params": {
115            "filePath": path.map(display),
116            "lineStart": start.line,
117            "lineEnd": end.line,
118        }
119    })
120    .to_string()
121}
122
123fn display(p: &Path) -> String {
124    p.display().to_string()
125}
126
127/// Order two positions so `start <= end` (a selection may be anchored either
128/// way — the head can sit before the anchor).
129fn ordered(a: Position, b: Position) -> (Position, Position) {
130    if (a.line, a.byte) <= (b.line, b.byte) {
131        (a, b)
132    } else {
133        (b, a)
134    }
135}
136
137#[cfg(test)]
138mod tests {
139    #![allow(clippy::unwrap_used)]
140    use super::*;
141
142    fn sel(anchor: (u32, u32), head: (u32, u32)) -> SelectionSet {
143        SelectionSet::single(lattice_protocol::selection::Selection {
144            anchor: Position {
145                line: anchor.0,
146                byte: anchor.1,
147            },
148            head: Position {
149                line: head.0,
150                byte: head.1,
151            },
152            visual: None,
153        })
154    }
155
156    #[test]
157    fn selection_changed_frame_carries_method_path_and_ordered_range() {
158        let frame = selection_changed_frame(&sel((1, 2), (3, 4)), Some(Path::new("/work/a.rs")));
159        let v: serde_json::Value = serde_json::from_str(&frame).unwrap();
160        assert_eq!(v["method"], "selection_changed");
161        assert_eq!(v["params"]["filePath"], "/work/a.rs");
162        assert_eq!(v["params"]["selection"]["start"]["line"], 1);
163        assert_eq!(v["params"]["selection"]["end"]["line"], 3);
164        assert_eq!(v["params"]["selection"]["isEmpty"], false);
165    }
166
167    #[test]
168    fn selection_range_is_ordered_even_when_anchored_backwards() {
169        // head before anchor → start must still be the earlier position.
170        let frame = selection_changed_frame(&sel((5, 0), (2, 0)), None);
171        let v: serde_json::Value = serde_json::from_str(&frame).unwrap();
172        assert_eq!(v["params"]["selection"]["start"]["line"], 2);
173        assert_eq!(v["params"]["selection"]["end"]["line"], 5);
174        assert_eq!(v["params"]["filePath"], serde_json::Value::Null);
175    }
176
177    #[test]
178    fn empty_selection_is_flagged() {
179        let frame = selection_changed_frame(&sel((1, 1), (1, 1)), None);
180        let v: serde_json::Value = serde_json::from_str(&frame).unwrap();
181        assert_eq!(v["params"]["selection"]["isEmpty"], true);
182    }
183
184    #[test]
185    fn did_change_active_editor_frame_carries_method_and_path() {
186        let frame = did_change_active_editor_frame(Some(Path::new("/work/b.rs")));
187        let v: serde_json::Value = serde_json::from_str(&frame).unwrap();
188        assert_eq!(v["method"], "didChangeActiveEditor");
189        assert_eq!(v["params"]["filePath"], "/work/b.rs");
190    }
191
192    #[test]
193    fn at_mentioned_frame_carries_method_path_and_line_range() {
194        let frame = at_mentioned_frame(&sel((2, 0), (6, 0)), Some(Path::new("/work/c.rs")));
195        let v: serde_json::Value = serde_json::from_str(&frame).unwrap();
196        assert_eq!(v["method"], "at_mentioned");
197        assert_eq!(v["params"]["filePath"], "/work/c.rs");
198        assert_eq!(v["params"]["lineStart"], 2);
199        assert_eq!(v["params"]["lineEnd"], 6);
200    }
201}