Skip to main content

nmbrs_runtime/
log_sink.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Asynchronous log file sink.
5//!
6//! See SRD-02 §"Display and Diagnostic Decoupling" for the
7//! design tenet. Producers (anyone calling `diag!()` /
8//! `observer::log()`) `try_send` a fully-formatted line into a
9//! bounded channel and return immediately. A dedicated
10//! `log-sink` OS thread is the only writer to the file —
11//! syscalls, locking, fsync stalls all happen there, never on
12//! a tokio worker.
13//!
14//! Overflow policy: bounded channel + `try_send`. If the sink
15//! falls behind (slow disk, full disk, NFS hang), producers
16//! drop the line and bump a `dropped_count` counter.
17//! Diagnostics never block the runtime; the dropped count is
18//! visible through the inspector endpoint and the post-run
19//! summary so the operator knows when log loss occurred.
20
21use std::fs::{File, OpenOptions};
22use std::io::{self, Write};
23use std::path::Path;
24use std::sync::OnceLock;
25use std::sync::atomic::{AtomicU64, Ordering};
26use std::sync::mpsc;
27use std::thread;
28
29/// Capacity of the bounded log channel. Sized for short bursts
30/// (every fiber emits one line at end-of-phase, plus periodic
31/// drain progress) without overusing memory: 4096 lines × ~256
32/// bytes ≈ 1 MB worst case in the queue.
33const LOG_CHANNEL_CAPACITY: usize = 4096;
34
35/// The single global log sink — set once via [`init`] when the
36/// session directory is known.
37static GLOBAL_LOG_SINK: OnceLock<LogSink> = OnceLock::new();
38
39/// Producer side of the async log sink.
40pub struct LogSink {
41    /// Bounded sender. Producers `try_send`; on overflow they
42    /// drop the line and bump [`Self::dropped_count`]. Never
43    /// blocks.
44    sender: mpsc::SyncSender<Vec<u8>>,
45    /// Lines dropped because the sink couldn't keep up. Visible
46    /// through the inspector and reported once at shutdown.
47    dropped_count: AtomicU64,
48}
49
50impl LogSink {
51    /// Try to enqueue a fully-formatted log line. Never blocks.
52    /// Returns `Ok(())` if accepted, `Err(())` if dropped.
53    // reason: the only failure is "dropped"; there is no error payload to
54    // carry, and callers branch on Ok/Err alone — a richer error type would
55    // be noise.
56    #[allow(clippy::result_unit_err)]
57    pub fn try_send(&self, line: Vec<u8>) -> Result<(), ()> {
58        match self.sender.try_send(line) {
59            Ok(()) => Ok(()),
60            Err(_) => {
61                self.dropped_count.fetch_add(1, Ordering::Relaxed);
62                Err(())
63            }
64        }
65    }
66
67    /// Count of lines dropped since startup. Useful as a health
68    /// signal in the inspector and at shutdown.
69    pub fn dropped_count(&self) -> u64 {
70        self.dropped_count.load(Ordering::Relaxed)
71    }
72}
73
74/// Initialize the global log sink with a target file. Called
75/// once by the runner after the session directory exists.
76/// Silently no-ops on a second call — the first session wins,
77/// matching the previous behavior of `set_log_file`.
78pub fn init(path: &Path) -> io::Result<()> {
79    let file = OpenOptions::new().create(true).append(true).open(path)?;
80    let (tx, rx) = mpsc::sync_channel::<Vec<u8>>(LOG_CHANNEL_CAPACITY);
81    spawn_writer_thread(file, rx);
82    let _ = GLOBAL_LOG_SINK.set(LogSink {
83        sender: tx,
84        dropped_count: AtomicU64::new(0),
85    });
86    Ok(())
87}
88
89/// Borrow the global log sink, if initialized. Used by the
90/// `log()` hot path and by the inspector to report the dropped
91/// count.
92pub fn global() -> Option<&'static LogSink> {
93    GLOBAL_LOG_SINK.get()
94}
95
96fn spawn_writer_thread(mut file: File, rx: mpsc::Receiver<Vec<u8>>) {
97    thread::Builder::new()
98        .name("log-sink".into())
99        .spawn(move || {
100            // recv() blocks the dedicated writer thread when
101            // the channel is empty — fine, this is not a tokio
102            // worker. SRD-02 §"No Blocking Primitives in Async
103            // Contexts" only forbids blocking *inside* tokio.
104            while let Ok(buf) = rx.recv() {
105                // Best-effort write. A failed write means the
106                // file system is unhealthy; we keep draining so
107                // upstream producers still see `try_send`
108                // success and the runtime keeps moving. Failed
109                // writes are silent today; an explicit
110                // last-write-error counter is a follow-up if
111                // needed.
112                let _ = file.write_all(&buf);
113            }
114            // Channel closed (all senders dropped). Flush
115            // pending writes to disk before the thread exits.
116            let _ = file.flush();
117            let _ = file.sync_all();
118        })
119        .expect("spawn log-sink thread");
120}