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}