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}