Skip to main content

strop_trace/
lib.rs

1//! One opt-in diagnostic sink for all strop crates. Producers never perform file
2//! I/O or wait for the writer; an incomplete trace is always reported as such.
3mod event;
4mod writer;
5
6pub use event::{preview, ContentPolicy, EventKind, TraceOptions, SCHEMA_VERSION};
7
8use parking_lot::Mutex;
9use serde::Serialize;
10use std::fs::OpenOptions;
11use std::path::{Path, PathBuf};
12use std::sync::atomic::{AtomicBool, Ordering};
13use std::sync::mpsc::{sync_channel, SyncSender};
14use std::sync::{Arc, LazyLock};
15use std::thread::JoinHandle;
16use std::time::Instant;
17
18const QUEUE_CAPACITY: usize = 4096;
19static ACTIVE: LazyLock<Mutex<Option<Arc<Recorder>>>> = LazyLock::new(|| Mutex::new(None));
20static ENABLED: AtomicBool = AtomicBool::new(false);
21static CONTENT: AtomicBool = AtomicBool::new(false);
22
23#[derive(Debug, thiserror::Error)]
24pub enum TraceError {
25    #[error("a trace session is already active")]
26    AlreadyActive,
27    #[error("cannot create trace {path}: {source}")]
28    Open {
29        path: PathBuf,
30        source: std::io::Error,
31    },
32    #[error("cannot start trace writer: {0}")]
33    Spawn(std::io::Error),
34    #[error("incomplete trace: {0}")]
35    Incomplete(String),
36}
37
38#[derive(Default)]
39struct Failure {
40    message: Mutex<Option<String>>,
41    reported: AtomicBool,
42}
43impl Failure {
44    fn set(&self, message: impl FnOnce() -> String) {
45        let mut failure = self.message.lock();
46        if failure.is_none() {
47            *failure = Some(message());
48        }
49    }
50}
51
52struct Record {
53    kind: EventKind,
54    elapsed_us: u128,
55    fields: Vec<u8>,
56}
57struct Recorder {
58    sender: Mutex<Option<SyncSender<Record>>>,
59    failure: Arc<Failure>,
60    started: Instant,
61}
62
63/// Owning lifetime of a trace. Explicit finish reports errors; Drop still drains.
64pub struct TraceSession {
65    recorder: Arc<Recorder>,
66    worker: Option<JoinHandle<()>>,
67}
68
69pub fn start(path: &Path, options: TraceOptions) -> Result<TraceSession, TraceError> {
70    let mut active = ACTIVE.lock();
71    if active.is_some() {
72        return Err(TraceError::AlreadyActive);
73    }
74    let mut open = OpenOptions::new();
75    open.write(true).create_new(true);
76    #[cfg(unix)]
77    {
78        use std::os::unix::fs::OpenOptionsExt;
79        open.mode(0o600);
80    }
81    let file = open.open(path).map_err(|source| TraceError::Open {
82        path: path.to_path_buf(),
83        source,
84    })?;
85    let (sender, receiver) = sync_channel(QUEUE_CAPACITY);
86    let failure = Arc::new(Failure::default());
87    let writer_failure = Arc::clone(&failure);
88    let worker = std::thread::Builder::new()
89        .name("strop-trace".into())
90        .spawn(move || writer::run(file, receiver, writer_failure))
91        .map_err(TraceError::Spawn)?;
92    let recorder = Arc::new(Recorder {
93        sender: Mutex::new(Some(sender)),
94        failure,
95        started: Instant::now(),
96    });
97    *active = Some(Arc::clone(&recorder));
98    CONTENT.store(options.content == ContentPolicy::Full, Ordering::Release);
99    ENABLED.store(true, Ordering::Release);
100    Ok(TraceSession {
101        recorder,
102        worker: Some(worker),
103    })
104}
105
106#[inline]
107pub fn enabled() -> bool {
108    ENABLED.load(Ordering::Relaxed)
109}
110#[inline]
111pub fn capture_content() -> bool {
112    enabled() && CONTENT.load(Ordering::Relaxed)
113}
114
115/// Lazy producer: no payload construction or allocation when disabled.
116pub fn record_with<T: Serialize>(kind: EventKind, fields: impl FnOnce() -> T) {
117    if enabled() {
118        record(kind, &fields());
119    }
120}
121
122pub fn record<T: Serialize>(kind: EventKind, fields: &T) {
123    if !enabled() {
124        return;
125    }
126    let Some(recorder) = ACTIVE.lock().clone() else {
127        return;
128    };
129    let fields = match serde_json::to_vec(fields) {
130        Ok(fields) => fields,
131        Err(error) => {
132            recorder
133                .failure
134                .set(|| format!("event serialization failed: {error}"));
135            return;
136        }
137    };
138    // This lock protects queue admission only, never disk writes. Stamping under
139    // the same lock makes timestamps nondecreasing in the writer's receive order.
140    let sender = recorder.sender.lock();
141    if let Some(sender) = sender.as_ref() {
142        let record = Record {
143            kind,
144            elapsed_us: recorder.started.elapsed().as_micros(),
145            fields,
146        };
147        if let Err(error) = sender.try_send(record) {
148            recorder
149                .failure
150                .set(|| format!("event queue admission failed: {error}"));
151        }
152    }
153}
154
155/// Report once to the editor's status line; finish still returns the failure.
156pub fn take_failure() -> Option<String> {
157    let recorder = ACTIVE.lock().clone()?;
158    let message = recorder.failure.message.lock().clone()?;
159    (!recorder.failure.reported.swap(true, Ordering::Relaxed)).then_some(message)
160}
161
162impl TraceSession {
163    pub fn finish(mut self) -> Result<(), TraceError> {
164        self.close()
165    }
166
167    fn close(&mut self) -> Result<(), TraceError> {
168        let Some(worker) = self.worker.take() else {
169            return Ok(());
170        };
171        {
172            let mut active = ACTIVE.lock();
173            if active
174                .as_ref()
175                .is_some_and(|value| Arc::ptr_eq(value, &self.recorder))
176            {
177                ENABLED.store(false, Ordering::Release);
178                CONTENT.store(false, Ordering::Release);
179                *active = None;
180            }
181        }
182        self.recorder.sender.lock().take();
183        if worker.join().is_err() {
184            self.recorder
185                .failure
186                .set(|| "writer thread panicked".into());
187        }
188        match self.recorder.failure.message.lock().clone() {
189            Some(error) => Err(TraceError::Incomplete(error)),
190            None => Ok(()),
191        }
192    }
193}
194impl Drop for TraceSession {
195    fn drop(&mut self) {
196        // Explicit finish is the reporting boundary. Drop guarantees durability
197        // during unwinding without risking a second panic or corrupting the TUI.
198        let _ = self.close();
199    }
200}
201
202#[cfg(test)]
203mod tests;