Skip to main content

trusty_common/
log_buffer.rs

1//! Bounded in-memory ring buffer of recent tracing log lines.
2//!
3//! Why: Operators debugging a running daemon want the last N log lines
4//!      without SSHing to the box, tailing a file, or restarting with a
5//!      different `RUST_LOG`. A small in-process ring buffer lets the daemon
6//!      serve recent logs over HTTP (`GET /logs/tail`) at near-zero cost and
7//!      with no file I/O. The cap keeps memory bounded on a long-running
8//!      process.
9//! What: [`LogBuffer`] is a thread-safe `VecDeque<String>` capped at a fixed
10//!      capacity; the oldest line is evicted on overflow. [`LogBufferLayer`]
11//!      is a `tracing_subscriber::Layer` that formats every event into one
12//!      line and pushes it onto the buffer. The HTTP handler drains the tail.
13//! Test: see the `tests` module — capacity eviction, tail semantics, and a
14//!      layer-integration test that emits events through a real subscriber.
15//!
16//! [`LogBuffer`]: crate::log_buffer::LogBuffer
17//! [`LogBufferLayer`]: crate::log_buffer::LogBufferLayer
18
19use std::collections::VecDeque;
20use std::fmt::Write as _;
21use std::sync::{Arc, Mutex};
22
23use tracing::field::{Field, Visit};
24use tracing_subscriber::Layer;
25use tracing_subscriber::layer::Context;
26
27/// Default ring-buffer capacity (lines). Sized so a daemon retains a few
28/// minutes of INFO-level chatter while costing well under 1 MB of RAM.
29pub const DEFAULT_LOG_CAPACITY: usize = 1000;
30
31/// Thread-safe, bounded ring buffer of formatted log lines.
32///
33/// Why: shared between the tracing `Layer` (writer) and the HTTP handler
34///      (reader); both hold cheap `Arc` clones of the same underlying deque.
35/// What: wraps `Arc<Mutex<VecDeque<String>>>`. `push` appends and evicts the
36///      oldest line past capacity; `tail` snapshots the most recent N lines.
37/// Test: `capacity_evicts_oldest`, `tail_returns_last_n`.
38#[derive(Clone, Debug)]
39pub struct LogBuffer {
40    inner: Arc<Mutex<VecDeque<String>>>,
41    capacity: usize,
42}
43
44impl LogBuffer {
45    /// Create an empty buffer with the given line capacity.
46    ///
47    /// Why: callers (daemon startup) choose the cap; tests use a tiny one.
48    /// What: allocates a `VecDeque` with `capacity.max(1)` reserved slots so a
49    ///      zero capacity is treated as one (a zero-cap ring is useless and
50    ///      would panic on the eviction arithmetic).
51    /// Test: `capacity_evicts_oldest`.
52    #[must_use]
53    pub fn new(capacity: usize) -> Self {
54        let capacity = capacity.max(1);
55        Self {
56            inner: Arc::new(Mutex::new(VecDeque::with_capacity(capacity))),
57            capacity,
58        }
59    }
60
61    /// Append a line, evicting the oldest entry when at capacity.
62    ///
63    /// Why: a tracing `Layer` calls this on every event; it must never panic
64    ///      or block long. A poisoned mutex (a prior panic while logging) is
65    ///      recovered via `into_inner` so logging itself never cascades a
66    ///      panic into the daemon.
67    /// What: pushes `line` to the back; if length now exceeds `capacity`,
68    ///      pops the front.
69    /// Test: `capacity_evicts_oldest`.
70    pub fn push(&self, line: String) {
71        let mut guard = match self.inner.lock() {
72            Ok(g) => g,
73            Err(poisoned) => poisoned.into_inner(),
74        };
75        guard.push_back(line);
76        while guard.len() > self.capacity {
77            guard.pop_front();
78        }
79    }
80
81    /// Snapshot the most recent `n` lines (or all, when `n` exceeds the
82    /// current length).
83    ///
84    /// Why: the `/logs/tail` handler returns these as a JSON array. Cloning
85    ///      under the lock keeps the critical section short and lets the
86    ///      caller serialise without holding the mutex.
87    /// What: returns a `Vec<String>` of at most `n` lines, oldest-first.
88    /// Test: `tail_returns_last_n`, `tail_all_when_n_exceeds_len`.
89    #[must_use]
90    pub fn tail(&self, n: usize) -> Vec<String> {
91        let guard = match self.inner.lock() {
92            Ok(g) => g,
93            Err(poisoned) => poisoned.into_inner(),
94        };
95        let skip = guard.len().saturating_sub(n);
96        guard.iter().skip(skip).cloned().collect()
97    }
98
99    /// Total number of lines currently buffered.
100    ///
101    /// Why: the `/logs/tail` response reports `total` so callers can tell
102    ///      whether the buffer has wrapped.
103    /// What: returns the deque length.
104    /// Test: `tail_returns_last_n` asserts `len` after pushes.
105    #[must_use]
106    pub fn len(&self) -> usize {
107        match self.inner.lock() {
108            Ok(g) => g.len(),
109            Err(poisoned) => poisoned.into_inner().len(),
110        }
111    }
112
113    /// Whether the buffer holds no lines.
114    ///
115    /// Why: clippy requires `is_empty` alongside `len`; also a convenient
116    ///      readiness check in tests.
117    /// What: returns `len() == 0`.
118    /// Test: covered by `capacity_evicts_oldest`.
119    #[must_use]
120    pub fn is_empty(&self) -> bool {
121        self.len() == 0
122    }
123
124    /// Remove and return every buffered line, oldest first (#7390).
125    ///
126    /// Why: [`capture_logs`] reuses ONE buffer for the whole test binary, so it
127    /// needs to empty it around each capture rather than read a tail whose age
128    /// it cannot tell. Test-only: nothing in the daemon path empties the ring.
129    /// What: drains the deque under the same poison-tolerant lock as
130    /// [`LogBuffer::push`].
131    /// Test: `capture_logs_sees_events_from_the_body_only`.
132    #[cfg(test)]
133    fn take(&self) -> Vec<String> {
134        let mut guard = match self.inner.lock() {
135            Ok(g) => g,
136            Err(poisoned) => poisoned.into_inner(),
137        };
138        guard.drain(..).collect()
139    }
140}
141
142/// `tracing_subscriber::Layer` that mirrors every event into a [`LogBuffer`].
143///
144/// Why: wiring this layer into the subscriber means the daemon's normal
145///      `tracing::info!` / `warn!` calls are captured for `/logs/tail` with
146///      no extra call sites — the buffer stays in lock-step with stderr.
147/// What: on each event, formats
148///      `<YYYY-MM-DD HH:MM:SS> [<level> <target>] <message> k=v …` into a
149///      single line and pushes it. The leading local-time timestamp (issue
150///      #846) lets the dashboard log view show when each line was emitted.
151///      Level/target/fields are collected via a lightweight `Visit`
152///      implementation.
153/// Test: `layer_captures_events` installs the layer on a real subscriber and
154///      asserts an emitted event lands in the buffer.
155pub struct LogBufferLayer {
156    buffer: LogBuffer,
157}
158
159impl LogBufferLayer {
160    /// Wrap a [`LogBuffer`] as a tracing layer.
161    ///
162    /// Why: the daemon constructs the buffer first (so it can also hand a
163    ///      clone to its HTTP state) and then builds the layer around it.
164    /// What: stores a clone of the buffer handle.
165    /// Test: `layer_captures_events`.
166    #[must_use]
167    pub fn new(buffer: LogBuffer) -> Self {
168        Self { buffer }
169    }
170}
171
172/// Field visitor that accumulates an event's message and key/value fields
173/// into a single human-readable string.
174///
175/// Why: tracing events expose their data only through the `Visit` callback;
176///      we render it to text once so the buffer stores plain `String`s.
177/// What: the canonical `message` field becomes the line body; every other
178///      field is appended as ` key=value`.
179/// Test: exercised indirectly by `layer_captures_events`.
180struct LineVisitor {
181    message: String,
182    fields: String,
183}
184
185impl LineVisitor {
186    fn new() -> Self {
187        Self {
188            message: String::new(),
189            fields: String::new(),
190        }
191    }
192}
193
194impl Visit for LineVisitor {
195    fn record_debug(&mut self, field: &Field, value: &dyn std::fmt::Debug) {
196        if field.name() == "message" {
197            // `{:?}` on the message preserves it without surrounding quotes
198            // for string payloads in practice; use Display-ish formatting.
199            let _ = write!(self.message, "{value:?}");
200        } else {
201            let _ = write!(self.fields, " {}={value:?}", field.name());
202        }
203    }
204
205    fn record_str(&mut self, field: &Field, value: &str) {
206        if field.name() == "message" {
207            self.message.push_str(value);
208        } else {
209            let _ = write!(self.fields, " {}={value}", field.name());
210        }
211    }
212}
213
214impl<S: tracing::Subscriber> Layer<S> for LogBufferLayer {
215    fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
216        let meta = event.metadata();
217        let mut visitor = LineVisitor::new();
218        event.record(&mut visitor);
219        // Trim the leading `"` artefact that `{:?}` adds for the message when
220        // the payload was a quoted string literal — keep lines readable.
221        let message = visitor.message.trim_matches('"');
222        // Prepend a local-time timestamp (issue #846) so the dashboard log
223        // view shows per-line timing. We use `chrono::Local` directly rather
224        // than `tracing_subscriber::fmt::time::LocalTime`, which is unsound in
225        // multithreaded programs.
226        let ts = chrono::Local::now().format("%Y-%m-%d %H:%M:%S");
227        let line = format!(
228            "{} [{} {}] {}{}",
229            ts,
230            meta.level(),
231            meta.target(),
232            message,
233            visitor.fields
234        );
235        self.buffer.push(line);
236    }
237}
238
239/// The one capture subscriber the crate's tests install, the buffer it writes
240/// to, and the lock that gives one capturing test at a time exclusive use of it.
241///
242/// Why it is a static rather than a subscriber built per call: `tracing`'s
243/// per-callsite interest and its global maximum level are recomputed, process
244/// wide, every time a `Dispatch` is created — over the dispatchers alive at that
245/// instant. A capture subscriber that lives only for one test is therefore racing
246/// every other test that builds one, and the loser's callsites stay cached as
247/// "no subscriber wants this" while its own subscriber is installed. A build of
248/// this crate's suite failed once in six runs that way, with a `warn!` that ran
249/// producing zero captured lines (#7390). One `Dispatch` that is created once and
250/// never dropped is always in that recomputation, so the cached answer can no
251/// longer go stale. This is the tracing-subscriber exception to the crate's
252/// no-global-state rule, and it is compiled only into the test build.
253#[cfg(test)]
254static CAPTURE: std::sync::LazyLock<(tracing::Dispatch, LogBuffer, Mutex<()>)> =
255    std::sync::LazyLock::new(|| {
256        use tracing_subscriber::layer::SubscriberExt as _;
257
258        let buffer = LogBuffer::new(256);
259        let dispatch = tracing::Dispatch::new(
260            tracing_subscriber::registry().with(LogBufferLayer::new(buffer.clone())),
261        );
262        (dispatch, buffer, Mutex::new(()))
263    });
264
265/// Run `body` under a subscriber that captures every event, and return the
266/// captured lines alongside the body's value (#7365, #7390).
267///
268/// Why: the LEVEL a line carries is a behavioural contract in this crate — a
269/// self-healing condition reported at warn reads to an operator as a fault
270/// needing repair, which is the defect #7365 and #7390 each fixed. [`LogBuffer`]
271/// already renders the level into the line, so a level regression becomes an
272/// ordinary assertion instead of something only a human reading logs would
273/// notice. It lives here, beside the buffer it drives, so the modules under test
274/// share one capture rather than each growing a copy.
275/// What: takes [`CAPTURE`]'s lock, empties its buffer, runs `body` with
276/// [`CAPTURE`]'s dispatcher as this thread's default, and returns `body`'s value
277/// with everything the buffer collected, oldest first. Only the locked thread
278/// has that dispatcher as its default, so nothing another test logs concurrently
279/// reaches these lines. A panicking `body` poisons the lock and the next caller
280/// recovers it.
281/// Test: `capture_logs_sees_events_from_the_body_only` below, plus
282/// `search_index_reconcile.rs::{a_resolvable_conflict_is_logged_at_info_without_operator_advice,
283/// an_unrecoverable_refusal_still_warns,
284/// an_unanswered_create_the_registry_confirms_emits_no_warning}` and
285/// `search_index_confirm.rs::an_exhausted_confirm_deadline_is_still_a_warning`.
286#[cfg(test)]
287pub(crate) fn capture_logs<T>(body: impl FnOnce() -> T) -> (T, Vec<String>) {
288    let (dispatch, buffer, lock) = &*CAPTURE;
289    let _held = lock.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
290    let _ = buffer.take();
291    let value = tracing::dispatcher::with_default(dispatch, body);
292    (value, buffer.take())
293}
294
295#[cfg(test)]
296mod tests {
297    use super::*;
298
299    /// [`capture_logs`] returns what the body logged, and only that (#7390).
300    ///
301    /// Why: every level assertion in this crate reads its evidence from here, so
302    /// a capture that silently returned nothing would turn a "no warnings"
303    /// assertion into one that passes for the wrong reason — which is exactly
304    /// how the shared buffer's predecessor failed one run in six. Running it
305    /// twice also pins the emptying: the second capture must not see the first
306    /// one's line.
307    /// What: captures one `warn!` and one `info!`, then captures again from a
308    /// body that logs nothing.
309    /// Test: itself.
310    #[test]
311    fn capture_logs_sees_events_from_the_body_only() {
312        let ((), lines) = capture_logs(|| {
313            tracing::warn!("first line");
314            tracing::info!("second line");
315        });
316
317        assert_eq!(lines.len(), 2, "both events must be captured: {lines:?}");
318        assert!(lines[0].contains("WARN") && lines[0].contains("first line"));
319        assert!(lines[1].contains("INFO") && lines[1].contains("second line"));
320
321        let ((), after) = capture_logs(|| {});
322        assert!(
323            after.is_empty(),
324            "the buffer is emptied per capture, saw {after:?}"
325        );
326    }
327
328    #[test]
329    fn capacity_evicts_oldest() {
330        let buf = LogBuffer::new(3);
331        assert!(buf.is_empty());
332        for i in 0..5 {
333            buf.push(format!("line {i}"));
334        }
335        // Capacity 3 → only the last three survive.
336        assert_eq!(buf.len(), 3);
337        assert_eq!(buf.tail(10), vec!["line 2", "line 3", "line 4"]);
338    }
339
340    #[test]
341    fn tail_returns_last_n() {
342        let buf = LogBuffer::new(100);
343        for i in 0..10 {
344            buf.push(format!("l{i}"));
345        }
346        assert_eq!(buf.len(), 10);
347        assert_eq!(buf.tail(3), vec!["l7", "l8", "l9"]);
348    }
349
350    #[test]
351    fn tail_all_when_n_exceeds_len() {
352        let buf = LogBuffer::new(100);
353        buf.push("only".to_string());
354        assert_eq!(buf.tail(50), vec!["only"]);
355        assert_eq!(buf.tail(0), Vec::<String>::new());
356    }
357
358    #[test]
359    fn zero_capacity_treated_as_one() {
360        let buf = LogBuffer::new(0);
361        buf.push("a".to_string());
362        buf.push("b".to_string());
363        assert_eq!(buf.tail(10), vec!["b"]);
364    }
365
366    #[test]
367    fn layer_captures_events() {
368        use tracing_subscriber::layer::SubscriberExt;
369
370        let buffer = LogBuffer::new(10);
371        let subscriber = tracing_subscriber::registry().with(LogBufferLayer::new(buffer.clone()));
372        tracing::subscriber::with_default(subscriber, || {
373            tracing::info!(answer = 42, "hello from test");
374        });
375        let lines = buffer.tail(10);
376        assert_eq!(lines.len(), 1, "expected one captured line, got {lines:?}");
377        let line = &lines[0];
378        assert!(line.contains("hello from test"), "line was: {line}");
379        assert!(line.contains("answer=42"), "line was: {line}");
380        assert!(line.contains("INFO"), "line was: {line}");
381
382        // Issue #846: every line is prefixed with a `YYYY-MM-DD HH:MM:SS `
383        // local-time timestamp. Lock in the shape without over-fitting to a
384        // specific clock value: 4-digit year, then '-', then a space-delimited
385        // time component, with the level appearing after the timestamp.
386        let bytes = line.as_bytes();
387        assert!(
388            bytes.len() >= 19,
389            "line too short to hold a timestamp: {line}"
390        );
391        assert!(
392            bytes[0..4].iter().all(u8::is_ascii_digit),
393            "expected a 4-digit year prefix, line was: {line}"
394        );
395        assert_eq!(
396            bytes[4], b'-',
397            "expected '-' after the year, line was: {line}"
398        );
399        // The timestamp must come before the level/target bracket.
400        let ts_end = line
401            .find(" [")
402            .expect("expected a ' [' after the timestamp");
403        let ts = &line[..ts_end];
404        assert_eq!(
405            ts.len(),
406            19,
407            "timestamp should be exactly 'YYYY-MM-DD HH:MM:SS', got: {ts:?}"
408        );
409    }
410}