Skip to main content

lattice_lsp/
diagnostics.rs

1//! Diagnostics routing: `textDocument/publishDiagnostics` → editor
2//! subscribers (Phase 4.1.d.i).
3//!
4//! ## Why broadcast
5//!
6//! A document may have multiple panes, multiple decoration
7//! consumers (gutter glyph provider + underline overlay + the
8//! `:diagnostics` buffer view), and -- once §5.10's event bus
9//! arrives -- arbitrary plugin subscribers. Each gets a clone
10//! of every event. `tokio::sync::broadcast` is the right shape:
11//!
12//! - Sender: the actor.
13//! - Receivers: each subscriber owns one; cheap to clone the
14//!   underlying queue.
15//! - Lagging consumers drop oldest first; for diagnostics this
16//!   is correct -- the LATEST publish supersedes any in-flight
17//!   older ones (per-version dropping is the editor's job).
18//!
19//! ## Why a typed event vs raw `PublishDiagnosticsParams`
20//!
21//! [`DiagnosticEvent`] flattens the LSP shape to what the editor
22//! actually consumes (uri / version / diagnostics) and adds the
23//! `server_id` so a multi-server setup (rust-analyzer + clippy
24//! linter bridge) can be disambiguated by the consumer without
25//! carrying server identity through every hop. Rest of the LSP
26//! payload is preserved verbatim via `Vec<Diagnostic>`.
27
28use std::sync::Arc;
29
30use lsp_types::{Diagnostic, PublishDiagnosticsParams, Uri};
31use tokio::sync::broadcast;
32
33/// Capacity of the diagnostics broadcast channel. 256 events
34/// per server fits a fast indexer's burst rate (rust-analyzer
35/// at startup may publish ~once per crate file). Lagged
36/// receivers drop oldest first; the editor reconciles by
37/// reading the URI's latest event and ignoring earlier ones
38/// for that URI.
39pub const DIAGNOSTICS_CHANNEL_CAPACITY: usize = 256;
40
41/// One diagnostics publish from the server.
42///
43/// `Arc<...>` on the heavier fields keeps the broadcast clone
44/// cheap: the channel internally buffers the last N events, so
45/// every subscriber that lags a beat ends up cloning the same
46/// `Arc`s rather than the whole `Vec<Diagnostic>`.
47#[derive(Debug, Clone)]
48pub struct DiagnosticEvent {
49    /// Server id (e.g. `"rust"`). Lets a multi-server setup
50    /// disambiguate without carrying the identity through every
51    /// hop.
52    pub server_id: Arc<str>,
53    /// URI the diagnostics apply to.
54    pub uri: Uri,
55    /// Doc version the server computed against. The editor
56    /// compares this with its own `DocSync::version(uri)` and
57    /// drops events older than the current version (avoids
58    /// stale diagnostics overwriting fresher state when an
59    /// edit raced with the publish).
60    pub version: Option<i32>,
61    /// Diagnostics list. Empty list means "the server cleared
62    /// this URI's diagnostics" -- a real, meaningful event,
63    /// don't filter it out.
64    pub diagnostics: Arc<[Diagnostic]>,
65}
66
67impl DiagnosticEvent {
68    /// Construct from the LSP `PublishDiagnosticsParams` shape.
69    pub fn from_lsp(server_id: Arc<str>, params: PublishDiagnosticsParams) -> Self {
70        Self {
71            server_id,
72            uri: params.uri,
73            version: params.version,
74            diagnostics: Arc::from(params.diagnostics.into_boxed_slice()),
75        }
76    }
77
78    /// True iff this event clears (rather than reports) the
79    /// URI's diagnostics. The editor uses this to drop the
80    /// decoration overlay rather than rendering an empty list.
81    pub fn is_clear(&self) -> bool {
82        self.diagnostics.is_empty()
83    }
84}
85
86/// The actor's send-side diagnostics bus. One per actor; shared
87/// across read-loop, supervisor, future telemetry hooks.
88#[derive(Clone)]
89pub struct DiagnosticsBus {
90    tx: broadcast::Sender<DiagnosticEvent>,
91}
92
93impl DiagnosticsBus {
94    pub fn new() -> Self {
95        let (tx, _rx) = broadcast::channel(DIAGNOSTICS_CHANNEL_CAPACITY);
96        Self { tx }
97    }
98
99    /// Subscribe to this bus. Returns a `Receiver` that yields
100    /// every event published after the call. Late subscribers
101    /// don't see older events.
102    pub fn subscribe(&self) -> broadcast::Receiver<DiagnosticEvent> {
103        self.tx.subscribe()
104    }
105
106    /// Publish an event to all subscribers. Drops silently if no
107    /// subscriber is listening (broadcast::send returns Err in
108    /// that case; we don't surface it -- "no listener" is a
109    /// supported state, not an error).
110    pub fn publish(&self, ev: DiagnosticEvent) {
111        let _ = self.tx.send(ev);
112    }
113
114    /// Approximate count of active subscribers. Useful for
115    /// telemetry / debug output; not load-bearing.
116    pub fn receiver_count(&self) -> usize {
117        self.tx.receiver_count()
118    }
119}
120
121impl Default for DiagnosticsBus {
122    fn default() -> Self {
123        Self::new()
124    }
125}
126
127impl std::fmt::Debug for DiagnosticsBus {
128    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
129        f.debug_struct("DiagnosticsBus")
130            .field("subscribers", &self.receiver_count())
131            .finish()
132    }
133}
134
135#[cfg(test)]
136mod tests {
137    use super::*;
138    use lsp_types::{Diagnostic, DiagnosticSeverity, Position as LspPosition, Range as LspRange};
139    use std::str::FromStr;
140
141    fn sample_diagnostic() -> Diagnostic {
142        Diagnostic {
143            range: LspRange {
144                start: LspPosition {
145                    line: 0,
146                    character: 0,
147                },
148                end: LspPosition {
149                    line: 0,
150                    character: 1,
151                },
152            },
153            severity: Some(DiagnosticSeverity::ERROR),
154            code: Some(lsp_types::NumberOrString::String("E0308".into())),
155            code_description: None,
156            source: Some("rustc".into()),
157            message: "type mismatch".into(),
158            related_information: None,
159            tags: None,
160            data: None,
161        }
162    }
163
164    #[test]
165    fn from_lsp_preserves_uri_version_and_diagnostics() {
166        let uri = Uri::from_str("file:///x.rs").unwrap();
167        let params = PublishDiagnosticsParams {
168            uri: uri.clone(),
169            diagnostics: vec![sample_diagnostic()],
170            version: Some(7),
171        };
172        let ev = DiagnosticEvent::from_lsp(Arc::from("rust"), params);
173        assert_eq!(ev.uri, uri);
174        assert_eq!(ev.version, Some(7));
175        assert_eq!(ev.diagnostics.len(), 1);
176        assert!(!ev.is_clear());
177    }
178
179    #[test]
180    fn empty_diagnostics_is_clear_event() {
181        let uri = Uri::from_str("file:///y.rs").unwrap();
182        let params = PublishDiagnosticsParams {
183            uri,
184            diagnostics: Vec::new(),
185            version: Some(2),
186        };
187        let ev = DiagnosticEvent::from_lsp(Arc::from("rust"), params);
188        assert!(ev.is_clear());
189    }
190
191    #[tokio::test]
192    async fn bus_fans_out_to_multiple_subscribers() {
193        let bus = DiagnosticsBus::new();
194        let mut a = bus.subscribe();
195        let mut b = bus.subscribe();
196        assert_eq!(bus.receiver_count(), 2);
197        let uri = Uri::from_str("file:///z.rs").unwrap();
198        bus.publish(DiagnosticEvent {
199            server_id: Arc::from("rust"),
200            uri: uri.clone(),
201            version: Some(1),
202            diagnostics: Arc::from(Vec::<Diagnostic>::new().into_boxed_slice()),
203        });
204        let got_a = a.recv().await.unwrap();
205        let got_b = b.recv().await.unwrap();
206        assert_eq!(got_a.uri, uri);
207        assert_eq!(got_b.uri, uri);
208    }
209
210    #[tokio::test]
211    async fn bus_publish_with_no_subscribers_is_silent() {
212        let bus = DiagnosticsBus::new();
213        let uri = Uri::from_str("file:///void.rs").unwrap();
214        bus.publish(DiagnosticEvent {
215            server_id: Arc::from("rust"),
216            uri,
217            version: None,
218            diagnostics: Arc::from(Vec::<Diagnostic>::new().into_boxed_slice()),
219        });
220        // No assertion; just must not panic.
221    }
222
223    #[tokio::test]
224    async fn late_subscriber_does_not_see_older_events() {
225        let bus = DiagnosticsBus::new();
226        let uri = Uri::from_str("file:///prior.rs").unwrap();
227        bus.publish(DiagnosticEvent {
228            server_id: Arc::from("rust"),
229            uri,
230            version: None,
231            diagnostics: Arc::from(Vec::<Diagnostic>::new().into_boxed_slice()),
232        });
233        // Subscribe AFTER the publish.
234        let mut rx = bus.subscribe();
235        // Recv with a tight timeout: should be empty.
236        let r = tokio::time::timeout(std::time::Duration::from_millis(50), rx.recv()).await;
237        assert!(r.is_err(), "subscriber should not see prior event");
238    }
239}