Skip to main content

lattice_runtime/
messages_subscriber.rs

1//! `MessagesLayer`: a `tracing::Layer` that fans every event
2//! into the App's `MessagesRing` + publishes a typed
3//! `MessagePushed` on the editor event bus.
4//!
5//! Per `docs/dev/architecture/design.md` §5.10.6, `*messages*` is
6//! the editor's audit log: every record `App::set_message`
7//! produces, *plus* every `tracing::*` event from the editor +
8//! plugins, flows through this single subscriber into one
9//! buffer. The subscriber:
10//!
11//! - **Captures every event**, irrespective of where it
12//!   originated (App code, `lattice-lsp`, future plugin host).
13//! - **Translates `tracing::Level` to
14//!   `lattice_grammar::EchoLevel`** for parity with the legacy
15//!   `set_message` records (same wire enum either way).
16//! - **Pushes to `MessagesRing`** for backlog seeding when the
17//!   user opens `*messages*` mid-session.
18//! - **Publishes `MessagePushed`** so per-tick drains in the
19//!   App (`drain_message_events`) can append to the buffer.
20//!
21//! ## Install once at App boot
22//!
23//! [`install_messages_subscriber`] calls
24//! `tracing_subscriber::registry().with(MessagesLayer { ...
25//! }).set_global_default()`. The global default can only be
26//! installed once per process — multi-App test setups must use
27//! the no-install path (a `MessagesLayer` can still be
28//! constructed and exercised directly for unit tests; only the
29//! global install is gated).
30//!
31//! ## Hot-path cost
32//!
33//! When no subscriber is installed, `tracing::info!` is a
34//! const-time atomic load → ~10ns per call (the tracing crate's
35//! commitment). When the layer is installed:
36//! `record_debug` visitor allocs ~one string + a `Mutex::lock`
37//! on `MessagesRing` + bus publish. Total ~hundreds of ns per
38//! event — acceptable because LSP / mode events fire at human
39//! cadence, not keystroke cadence (per §8.2 Background-class).
40
41use std::sync::{Arc, Mutex, OnceLock};
42
43use tracing::field::{Field, Visit};
44use tracing::{Event, Subscriber};
45use tracing_subscriber::Layer;
46use tracing_subscriber::filter::EnvFilter;
47use tracing_subscriber::layer::{Context, SubscriberExt};
48use tracing_subscriber::registry::LookupSpan;
49use tracing_subscriber::reload;
50
51use crate::events::EventBus;
52use crate::messages::{MessagePushed, MessageRecord, MessagesRing};
53
54/// `tracing::Layer` that bridges every event into the
55/// `MessagesRing` + bus. Cheap to construct; cloning the layer
56/// shares the underlying ring + bus handles via `Arc`.
57#[derive(Clone)]
58pub struct MessagesLayer {
59    ring: Arc<Mutex<MessagesRing>>,
60    bus: Arc<EventBus>,
61}
62
63impl MessagesLayer {
64    /// New layer bound to the given ring + bus handles. The
65    /// layer captures every event the installed subscriber
66    /// receives.
67    pub fn new(ring: Arc<Mutex<MessagesRing>>, bus: Arc<EventBus>) -> Self {
68        Self { ring, bus }
69    }
70
71    /// Push a record + publish the bus event. Public so unit
72    /// tests can drive the layer without standing up a global
73    /// subscriber.
74    pub fn emit(&self, record: MessageRecord) {
75        if let Ok(mut ring) = self.ring.lock() {
76            ring.push(record.clone());
77        }
78        self.bus.publish_typed(MessagePushed { record });
79    }
80}
81
82impl<S> Layer<S> for MessagesLayer
83where
84    S: Subscriber + for<'a> LookupSpan<'a>,
85{
86    fn on_event(&self, event: &Event<'_>, _ctx: Context<'_, S>) {
87        let meta = event.metadata();
88        let level = lattice_grammar::EchoLevel::from(*meta.level());
89
90        let mut visitor = MessageVisitor::default();
91        event.record(&mut visitor);
92        let text = visitor.finish();
93        // Skip events with NOTHING to show. Most `tracing` macros set the
94        // implicit `message` field, but a structured-only event would
95        // otherwise produce an empty entry.
96        if text.is_empty() {
97            return;
98        }
99
100        self.emit(MessageRecord {
101            timestamp: std::time::SystemTime::now(),
102            level,
103            text,
104        });
105    }
106}
107
108/// Visitor that renders an event's implicit `message` field **and its
109/// structured fields** into one line.
110///
111/// ## The fields were being thrown away
112///
113/// This used to keep `message` and discard every other field, so a
114/// diagnostic written as
115///
116/// ```ignore
117/// debug!(?axes, want_syntax, have_syntax, "cells_matrix_invalidated");
118/// ```
119///
120/// reached `*messages*` as the bare string `cells_matrix_invalidated`.
121/// The message name is the least informative part of such an event —
122/// the whole point of that line is *which axis differed and by how
123/// much* — so the instrumentation was answering a question nobody
124/// could read the answer to.
125///
126/// That is not hypothetical. The stale-highlight investigation added
127/// exactly that event to name the invalidation axis behind a wrong-colour
128/// render, then read a `*messages*` capture that could not show it, and
129/// stalled. `*messages*` is documented as "the canonical surface" for
130/// inspecting events without leaving the editor; it has to carry what the
131/// events actually say.
132///
133/// ## Cost
134///
135/// Formatting is proportional to field count and happens on the thread
136/// that logged. That is acceptable because the events carrying many
137/// fields are `debug!` / `trace!`, which the level filter drops before
138/// this layer sees them unless the user opted in (`--log-level debug`).
139/// At the default level the common case is an `info!` with no extra
140/// fields, which costs one `is_empty` check.
141#[derive(Default)]
142struct MessageVisitor {
143    message: String,
144    fields: String,
145}
146
147impl MessageVisitor {
148    /// `message field=value field2=value2` — `tracing`'s own fmt layout,
149    /// so a `*messages*` line and a stderr line read the same way.
150    fn finish(self) -> String {
151        if self.fields.is_empty() {
152            return self.message;
153        }
154        if self.message.is_empty() {
155            return self.fields;
156        }
157        format!("{} {}", self.message, self.fields)
158    }
159
160    fn push_field(&mut self, name: &str, value: &dyn std::fmt::Debug) {
161        use std::fmt::Write as _;
162        if !self.fields.is_empty() {
163            self.fields.push(' ');
164        }
165        // `write!` to a String is infallible; the result is ignored rather
166        // than unwrapped so a logging path can never panic.
167        let _ = write!(self.fields, "{name}={value:?}");
168    }
169}
170
171impl Visit for MessageVisitor {
172    fn record_str(&mut self, field: &Field, value: &str) {
173        if field.name() == "message" {
174            self.message = value.to_string();
175        } else {
176            // Strings render bare — `path=/tmp/x.org`, not `path="/tmp/x.org"`
177            // — because the quotes are noise in a transcript a human reads.
178            self.push_field(field.name(), &format_args!("{value}"));
179        }
180    }
181
182    fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
183        if field.name() == "message" {
184            if self.message.is_empty() {
185                // `format_args!` Debug-formats to the rendered
186                // output without quotes (its Debug impl is
187                // identical to its Display).
188                self.message = format!("{value:?}");
189            }
190        } else {
191            self.push_field(field.name(), value);
192        }
193    }
194}
195
196/// Tracks whether `install_messages_subscriber` has installed
197/// the global tracing subscriber for this process. Idempotent:
198/// repeated calls (e.g. multi-App tests sharing the process)
199/// no-op after the first install instead of panicking.
200static GLOBAL_INSTALLED: OnceLock<()> = OnceLock::new();
201
202/// Process-wide log-level override set by the CLI before
203/// `install_messages_subscriber` runs. `lattice-cli::main`
204/// computes the level from `-v`/`-q`/`--log-level` flags and
205/// calls [`set_boot_log_level`]; `editor_boot.rs` calls
206/// [`boot_log_level`] inside the install path.
207///
208/// `OnceLock<String>` (not an env var) because the
209/// `lattice-cli` crate denies `unsafe_code` and modern Rust
210/// flags `std::env::set_var` as unsafe (env mutations aren't
211/// thread-safe). A process-wide `OnceLock` is the
212/// equivalent-but-safe mechanism for set-once boot config.
213static BOOT_LOG_LEVEL: OnceLock<String> = OnceLock::new();
214
215/// CLI sets the boot-time log level before constructing the
216/// editor. Idempotent — second call no-ops. The first set
217/// wins (matching `tokio::main` -> `App::new` ordering).
218pub fn set_boot_log_level(level: impl Into<String>) {
219    let _ = BOOT_LOG_LEVEL.set(level.into());
220}
221
222/// `editor_boot` reads the boot-time log level when calling
223/// [`install_messages_subscriber`]. Returns `None` when the
224/// CLI didn't set one (typical for tests + library callers);
225/// boot falls back to `"info"`.
226pub fn boot_log_level() -> Option<String> {
227    BOOT_LOG_LEVEL.get().cloned()
228}
229
230/// Process-wide flag set by the CLI before App::new to tell
231/// the runtime whether to enable the fmt-to-stderr layer.
232/// `Some(true)` ⇒ enable; `Some(false)` ⇒ disable;
233/// `None` ⇒ runtime falls back to its default (currently
234/// `true` to preserve previous behaviour for library callers
235/// that don't set it).
236///
237/// Issue #36 (2026-05-22): TUI sets this to `false` because
238/// stderr IS the terminal it paints into. GPUI sets to
239/// `true` (stderr is a separate stream).
240static BOOT_STDERR_ENABLED: OnceLock<bool> = OnceLock::new();
241
242/// CLI sets the boot-time stderr-enabled flag before
243/// constructing the editor. Idempotent — second call no-ops.
244pub fn set_boot_stderr_enabled(enabled: bool) {
245    let _ = BOOT_STDERR_ENABLED.set(enabled);
246}
247
248/// `editor_boot` reads the boot-time stderr-enabled flag.
249/// Returns `None` when the CLI didn't set one (test paths
250/// and library callers); editor_boot falls back to `false`
251/// (safe — never accidentally corrupt a TUI screen).
252pub fn boot_stderr_enabled() -> Option<bool> {
253    BOOT_STDERR_ENABLED.get().copied()
254}
255
256/// Reload-handle for the `EnvFilter` that gates which events
257/// the `MessagesLayer` captures. Stored at install time so
258/// `:set messages.filter=...` can swap the filter live via
259/// [`reload_messages_filter`] without restarting the editor.
260/// `OnceLock<Option<...>>` so the "no filter wired" case (test
261/// paths that construct the layer directly) is observable.
262type FilterReloadHandle = reload::Handle<EnvFilter, tracing_subscriber::Registry>;
263static FILTER_HANDLE: OnceLock<FilterReloadHandle> = OnceLock::new();
264
265/// Install `MessagesLayer` as the global tracing subscriber,
266/// gated by an `EnvFilter` whose initial directive is
267/// `initial_filter`. The filter is reloadable -- live edits
268/// via [`reload_messages_filter`] swap the directive without
269/// re-installing the subscriber.
270///
271/// Idempotent: only the first call wins; subsequent calls
272/// return `false`. The first call's `ring` + `bus` are the
273/// ones every later event flows into; the first call's
274/// `initial_filter` seeds the filter handle.
275///
276/// `initial_filter` accepts the standard
277/// `tracing_subscriber::EnvFilter` directive syntax (`info`,
278/// `editor=info,lsp=debug`, ...). On parse failure the
279/// install falls back to `info`.
280///
281/// **Why a global subscriber:** `tracing` can only have one
282/// global default per process. Test isolation is handled by
283/// the layer's `ring`/`bus` Arcs — every test that wants its
284/// own messages stream constructs its own `MessagesLayer`
285/// (without installing globally) and exercises `on_event` /
286/// `emit` directly.
287pub fn install_messages_subscriber(
288    ring: Arc<Mutex<MessagesRing>>,
289    bus: Arc<EventBus>,
290    initial_filter: &str,
291    stderr_enabled: bool,
292) -> bool {
293    if GLOBAL_INSTALLED.get().is_some() {
294        return false;
295    }
296    let env_filter = EnvFilter::try_new(initial_filter)
297        .unwrap_or_else(|_| EnvFilter::try_new("info").expect("`info` is a valid EnvFilter spec"));
298    let (filter_layer, handle) = reload::Layer::new(env_filter);
299    let messages_layer = MessagesLayer::new(ring, bus);
300    // The `*messages*` buffer ALWAYS captures every event.
301    // The fmt layer (stderr writer) is OPTIONAL.
302    //
303    // Issue #36 (2026-05-22): TUI peers must NOT enable
304    // stderr — stderr IS the terminal ratatui paints into,
305    // so every `tracing::*` event blits a stray line over
306    // the screen until the next full redraw. The caller
307    // passes `stderr_enabled = false` for TUI; GPUI passes
308    // `true` (its stderr is a separate stream).
309    //
310    // `LATTICE_STDERR=1` (CLI flag) forces fmt ON regardless
311    // for users running TUI with `2>tracing.log` redirection.
312    //
313    // Either way, `*messages*` is the canonical surface —
314    // open with `:messages` to inspect events without leaving
315    // the editor.
316    let install_result = if stderr_enabled {
317        let fmt_layer = tracing_subscriber::fmt::layer().with_writer(std::io::stderr);
318        let subscriber = tracing_subscriber::registry()
319            .with(filter_layer)
320            .with(messages_layer)
321            .with(fmt_layer);
322        tracing::subscriber::set_global_default(subscriber)
323    } else {
324        let subscriber = tracing_subscriber::registry()
325            .with(filter_layer)
326            .with(messages_layer);
327        tracing::subscriber::set_global_default(subscriber)
328    };
329    match install_result {
330        Ok(()) => {
331            let _ = GLOBAL_INSTALLED.set(());
332            let _ = FILTER_HANDLE.set(handle);
333            true
334        }
335        Err(_) => {
336            // Another subscriber is already installed (perhaps
337            // by `RUST_LOG=...` env_logger setup in a downstream
338            // test or by a parallel App). Treat as success in
339            // the sense that we don't try again.
340            let _ = GLOBAL_INSTALLED.set(());
341            false
342        }
343    }
344}
345
346/// Live-swap the messages-layer filter directive. Returns
347/// `Err` when the directive fails to parse, when the global
348/// subscriber wasn't installed (test paths), or when the
349/// reload handle has been dropped. The App's option-change
350/// cascade calls this on `:set messages.filter=<spec>`.
351pub fn reload_messages_filter(spec: &str) -> Result<(), MessagesFilterReloadError> {
352    let new_filter = EnvFilter::try_new(spec).map_err(|e| MessagesFilterReloadError::Parse {
353        spec: spec.to_string(),
354        reason: e.to_string(),
355    })?;
356    let handle = FILTER_HANDLE
357        .get()
358        .ok_or(MessagesFilterReloadError::SubscriberNotInstalled)?;
359    handle
360        .modify(|f| *f = new_filter)
361        .map_err(|e| MessagesFilterReloadError::Reload(e.to_string()))
362}
363
364/// Why a [`reload_messages_filter`] call failed.
365#[derive(Debug, thiserror::Error)]
366pub enum MessagesFilterReloadError {
367    /// The directive didn't parse as `EnvFilter` syntax. The
368    /// typed-option validator already rejects bad strings at
369    /// `:set` time; reaching this variant means the validator
370    /// missed something.
371    #[error("messages.filter `{spec}` is not a valid filter directive: {reason}")]
372    Parse { spec: String, reason: String },
373    /// `install_messages_subscriber` was never called for this
374    /// process. Production boot always installs; test paths
375    /// that exercise the App without going through `App::new`
376    /// can observe this.
377    #[error("messages-mode tracing subscriber not installed; cannot reload filter")]
378    SubscriberNotInstalled,
379    /// `tracing_subscriber::reload::Handle::modify` returned
380    /// an error (typically because the underlying layer was
381    /// dropped). Should not happen in production but the
382    /// error type carries the message for diagnostics.
383    #[error("messages.filter reload failed: {0}")]
384    Reload(String),
385}
386
387#[cfg(test)]
388mod tests {
389    #![allow(clippy::unwrap_used)]
390    use super::*;
391    use tracing_subscriber::layer::SubscriberExt;
392
393    /// Layer captures `tracing::info!` text + level + emits a
394    /// `MessagePushed` event on the bus. Exercised against a
395    /// per-test subscriber via `with_default`, so this test
396    /// doesn't depend on `install_messages_subscriber` (which
397    /// can only run once per process).
398    #[test]
399    fn layer_captures_info_event_to_ring_and_bus() {
400        let ring = Arc::new(Mutex::new(MessagesRing::default()));
401        let bus = Arc::new(EventBus::new());
402        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<MessagePushed>();
403        bus.subscribe_typed(tx);
404        let layer = MessagesLayer::new(ring.clone(), bus);
405        let subscriber = tracing_subscriber::registry().with(layer);
406        tracing::subscriber::with_default(subscriber, || {
407            tracing::info!("hello world");
408        });
409        // Ring captured the record.
410        let ring_records = ring.lock().unwrap();
411        assert_eq!(ring_records.len(), 1);
412        let r = ring_records.records().front().unwrap();
413        assert_eq!(r.level, lattice_grammar::EchoLevel::Info);
414        assert_eq!(r.text, "hello world");
415        drop(ring_records);
416        // Bus delivered the event.
417        let evt = rx.try_recv().unwrap();
418        assert_eq!(evt.record.text, "hello world");
419        assert_eq!(evt.record.level, lattice_grammar::EchoLevel::Info);
420    }
421
422    /// A structured event's FIELDS reach `*messages*`, not just its name.
423    ///
424    /// They used to be discarded. A diagnostic written to name *which*
425    /// invalidation axis moved arrived as the bare string
426    /// `cells_matrix_invalidated`, which is the least informative part of
427    /// it — and `*messages*` is the surface the docs point at for reading
428    /// events without leaving the editor, so the answer existed and could
429    /// not be read. The stale-highlight investigation stalled on exactly
430    /// that.
431    #[test]
432    fn layer_renders_structured_fields_beside_the_message() {
433        let ring = Arc::new(Mutex::new(MessagesRing::default()));
434        let bus = Arc::new(EventBus::new());
435        let layer = MessagesLayer::new(ring.clone(), bus);
436        let subscriber = tracing_subscriber::registry().with(layer);
437        tracing::subscriber::with_default(subscriber, || {
438            tracing::info!(
439                want_syntax = 7u64,
440                have_syntax = 6u64,
441                "cells_matrix_invalidated"
442            );
443        });
444        let ring_records = ring.lock().unwrap();
445        let r = ring_records.records().front().unwrap();
446        assert!(
447            r.text.contains("want_syntax=7") && r.text.contains("have_syntax=6"),
448            "the fields are the whole diagnostic; got {:?}",
449            r.text
450        );
451        assert!(
452            r.text.starts_with("cells_matrix_invalidated"),
453            "the message still leads the line: {:?}",
454            r.text
455        );
456    }
457
458    /// A `&str` field renders bare. Quotes are noise in a transcript, and
459    /// paths are the common case (`path=/tmp/notes.org`).
460    #[test]
461    fn string_fields_render_without_quotes() {
462        let ring = Arc::new(Mutex::new(MessagesRing::default()));
463        let bus = Arc::new(EventBus::new());
464        let layer = MessagesLayer::new(ring.clone(), bus);
465        let subscriber = tracing_subscriber::registry().with(layer);
466        tracing::subscriber::with_default(subscriber, || {
467            tracing::info!(lang = "org", "attached");
468        });
469        let ring_records = ring.lock().unwrap();
470        let r = ring_records.records().front().unwrap();
471        assert_eq!(r.text, "attached lang=org");
472    }
473
474    /// An event with fields and NO message still lands, carrying its
475    /// fields.
476    ///
477    /// Supersedes `layer_drops_messageless_events`, which asserted the
478    /// opposite. Its stated reason was "don't pollute `*messages*` with
479    /// EMPTY rows" — a consequence of fields being discarded, not an
480    /// independent requirement. Now that fields render, the row is
481    /// `buffer_id=7` rather than blank, so the reason no longer holds and
482    /// dropping the event would be throwing away the only thing it said.
483    #[test]
484    fn a_field_only_event_is_not_dropped() {
485        let ring = Arc::new(Mutex::new(MessagesRing::default()));
486        let bus = Arc::new(EventBus::new());
487        let layer = MessagesLayer::new(ring.clone(), bus);
488        let subscriber = tracing_subscriber::registry().with(layer);
489        tracing::subscriber::with_default(subscriber, || {
490            tracing::info!(axis = "syntax");
491        });
492        let ring_records = ring.lock().unwrap();
493        assert_eq!(ring_records.len(), 1, "a field-only event must not vanish");
494        assert_eq!(ring_records.records().front().unwrap().text, "axis=syntax");
495    }
496
497    /// `*messages*` captures EVERY level the `EnvFilter` admits.
498    ///
499    /// This is a product requirement, not an implementation detail:
500    /// `--log-level debug` exists so debug output can be read **inside
501    /// the editor**, and `--stderr-logs` only ADDS the fmt layer — it
502    /// never diverts events away from here (`install_messages_subscriber`
503    /// installs this layer in both branches).
504    ///
505    /// 2026-08-16: I briefly filtered this layer to INFO+ to break a
506    /// feedback loop (a `debug!` per render became an edit to the rendered
507    /// `*messages*` buffer, which republished, which re-woke the workers).
508    /// That fixed the loop by deleting the feature. The loop is now fixed
509    /// where it belongs — the render cycle no longer logs per render, see
510    /// the worker `run` loops — and this test exists so the shortcut is
511    /// not taken again.
512    #[test]
513    fn layer_translates_every_level() {
514        let ring = Arc::new(Mutex::new(MessagesRing::default()));
515        let bus = Arc::new(EventBus::new());
516        let layer = MessagesLayer::new(ring.clone(), bus);
517        let subscriber = tracing_subscriber::registry().with(layer);
518        tracing::subscriber::with_default(subscriber, || {
519            tracing::trace!("t");
520            tracing::debug!("d");
521            tracing::info!("i");
522            tracing::warn!("w");
523            tracing::error!("e");
524        });
525        let ring = ring.lock().unwrap();
526        let levels: Vec<lattice_grammar::EchoLevel> =
527            ring.records().iter().map(|r| r.level).collect();
528        assert_eq!(
529            levels,
530            vec![
531                lattice_grammar::EchoLevel::Trace,
532                lattice_grammar::EchoLevel::Debug,
533                lattice_grammar::EchoLevel::Info,
534                lattice_grammar::EchoLevel::Warn,
535                lattice_grammar::EchoLevel::Error,
536            ]
537        );
538    }
539
540    /// Formatted-args messages render to their final string
541    /// (no quotes, no Debug noise).
542    #[test]
543    fn layer_renders_format_args_without_debug_quotes() {
544        let ring = Arc::new(Mutex::new(MessagesRing::default()));
545        let bus = Arc::new(EventBus::new());
546        let layer = MessagesLayer::new(ring.clone(), bus);
547        let subscriber = tracing_subscriber::registry().with(layer);
548        tracing::subscriber::with_default(subscriber, || {
549            tracing::info!("opened {} buffers", 42);
550        });
551        let ring = ring.lock().unwrap();
552        assert_eq!(ring.records().front().unwrap().text, "opened 42 buffers");
553    }
554}