1mod bounded;
8mod event;
9pub mod export;
10pub mod replay;
11mod writer;
12
13pub use event::{
14 preview, ContentPolicy, EventKind, Limits, TraceOptions, MAX_CAPTURE_BYTES, MAX_CAPTURE_EVENTS,
15 MAX_RECORD_BYTES, SCHEMA_VERSION, TERMINAL_RESERVE,
16};
17
18use parking_lot::Mutex;
19use serde::Serialize;
20use std::fs::OpenOptions;
21use std::path::{Path, PathBuf};
22use std::sync::atomic::{AtomicBool, Ordering};
23use std::sync::mpsc::{sync_channel, SyncSender};
24use std::sync::{Arc, LazyLock};
25use std::thread::JoinHandle;
26use std::time::Instant;
27
28const QUEUE_CAPACITY: usize = 64;
31static ACTIVE: LazyLock<Mutex<Option<Arc<Recorder>>>> = LazyLock::new(|| Mutex::new(None));
32static ENABLED: AtomicBool = AtomicBool::new(false);
33static CONTENT: AtomicBool = AtomicBool::new(false);
34
35#[derive(Debug, thiserror::Error)]
36pub enum TraceError {
37 #[error("a trace session is already active")]
38 AlreadyActive,
39 #[error("cannot create trace {path}: {source}")]
40 Open {
41 path: PathBuf,
42 source: std::io::Error,
43 },
44 #[error("cannot start trace writer: {0}")]
45 Spawn(std::io::Error),
46 #[error("incomplete trace: {0}")]
47 Incomplete(String),
48}
49
50#[derive(Default)]
51struct Failure {
52 message: Mutex<Option<String>>,
53 reported: AtomicBool,
54}
55impl Failure {
56 fn set(&self, message: impl FnOnce() -> String) {
57 let mut failure = self.message.lock();
58 if failure.is_none() {
59 *failure = Some(message());
60 }
61 }
62}
63
64struct Record {
65 kind: EventKind,
66 elapsed_us: u128,
67 fields: Vec<u8>,
68}
69struct Recorder {
70 sender: Mutex<Option<SyncSender<Record>>>,
71 failure: Arc<Failure>,
72 started: Instant,
73 max_record: usize,
74}
75
76pub struct TraceSession {
78 recorder: Arc<Recorder>,
79 worker: Option<JoinHandle<()>>,
80}
81
82pub fn start(path: &Path, options: TraceOptions) -> Result<TraceSession, TraceError> {
83 if !options.limits.valid() {
84 return Err(TraceError::Incomplete("invalid capture limits".into()));
85 }
86 let mut active = ACTIVE.lock();
87 if active.is_some() {
88 return Err(TraceError::AlreadyActive);
89 }
90 let mut open = OpenOptions::new();
91 open.write(true).create_new(true);
92 #[cfg(unix)]
93 {
94 use std::os::unix::fs::OpenOptionsExt;
95 open.mode(0o600);
96 }
97 let file = open.open(path).map_err(|source| TraceError::Open {
98 path: path.to_path_buf(),
99 source,
100 })?;
101 let (sender, receiver) = sync_channel(QUEUE_CAPACITY);
102 let failure = Arc::new(Failure::default());
103 let writer_failure = Arc::clone(&failure);
104 let limits = options.limits;
105 let worker = std::thread::Builder::new()
106 .name("strop-trace".into())
107 .spawn(move || writer::run(file, receiver, writer_failure, limits))
108 .map_err(TraceError::Spawn)?;
109 let recorder = Arc::new(Recorder {
110 sender: Mutex::new(Some(sender)),
111 failure,
112 started: Instant::now(),
113 max_record: limits.record_bytes,
114 });
115 *active = Some(Arc::clone(&recorder));
116 CONTENT.store(options.content == ContentPolicy::Full, Ordering::Release);
117 ENABLED.store(true, Ordering::Release);
118 Ok(TraceSession {
119 recorder,
120 worker: Some(worker),
121 })
122}
123
124#[inline]
125pub fn enabled() -> bool {
126 ENABLED.load(Ordering::Relaxed)
127}
128#[inline]
129pub fn capture_content() -> bool {
130 enabled() && CONTENT.load(Ordering::Relaxed)
131}
132
133pub fn record_with<T: Serialize>(kind: EventKind, fields: impl FnOnce() -> T) {
135 if enabled() {
136 record(kind, &fields());
137 }
138}
139
140pub fn record<T: Serialize>(kind: EventKind, fields: &T) {
141 if !enabled() {
142 return;
143 }
144 let Some(recorder) = ACTIVE.lock().clone() else {
145 return;
146 };
147 let mut sender = recorder.sender.lock();
151 if sender.is_none() {
152 return;
153 }
154 if recorder.failure.message.lock().is_some() {
155 sender.take();
156 return;
157 }
158 let mut bytes = bounded::Bytes::new(recorder.max_record);
159 if serde_json::to_writer(&mut bytes, fields).is_err() {
160 recorder
161 .failure
162 .set(|| "record exceeds cap or cannot serialize".into());
163 sender.take();
164 return;
165 }
166 let record = Record {
167 kind,
168 elapsed_us: recorder.started.elapsed().as_micros(),
169 fields: bytes.into_vec(),
170 };
171 if sender
172 .as_ref()
173 .expect("checked sender")
174 .try_send(record)
175 .is_err()
176 {
177 recorder
178 .failure
179 .set(|| "capture queue full or writer unavailable".into());
180 sender.take();
181 }
182}
183
184pub fn mark_incomplete(message: &'static str) {
188 let Some(recorder) = ACTIVE.lock().clone() else {
189 return;
190 };
191 recorder.failure.set(|| message.into());
192 recorder.sender.lock().take();
193 CONTENT.store(false, Ordering::Release);
194}
195
196pub fn take_failure() -> Option<String> {
198 let recorder = ACTIVE.lock().clone()?;
199 let message = recorder.failure.message.lock().clone()?;
200 (!recorder.failure.reported.swap(true, Ordering::Relaxed)).then_some(message)
201}
202
203impl TraceSession {
204 pub fn finish(mut self) -> Result<(), TraceError> {
205 self.close()
206 }
207
208 fn close(&mut self) -> Result<(), TraceError> {
209 let Some(worker) = self.worker.take() else {
210 return Ok(());
211 };
212 {
213 let mut active = ACTIVE.lock();
214 if active
215 .as_ref()
216 .is_some_and(|value| Arc::ptr_eq(value, &self.recorder))
217 {
218 ENABLED.store(false, Ordering::Release);
219 CONTENT.store(false, Ordering::Release);
220 *active = None;
221 }
222 }
223 self.recorder.sender.lock().take();
224 if worker.join().is_err() {
225 self.recorder
226 .failure
227 .set(|| "writer thread panicked".into());
228 }
229 match self.recorder.failure.message.lock().clone() {
230 Some(error) => Err(TraceError::Incomplete(error)),
231 None => Ok(()),
232 }
233 }
234}
235impl Drop for TraceSession {
236 fn drop(&mut self) {
237 let _ = self.close();
240 }
241}
242
243#[cfg(test)]
244mod tests;