1mod 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
63pub 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
115pub 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 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
155pub 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 let _ = self.close();
199 }
200}
201
202#[cfg(test)]
203mod tests;