1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
use crate::buffer::{BufferMode, SegmentWriter};
use crate::flush_loop::run_flush_loop;
use crate::handle::{ControlCommand, Dial9Handle, InstallGlobalHandleError};
use crate::primitives::sync::{Arc, Mutex};
use crate::primitives::{sync::mpsc, thread::JoinHandle};
use crate::recorder::SoleRecorderGuard;
use crate::shared_state::SharedState;
use std::time::Duration;
/// The background worker thread and its stop signal.
///
/// Present only when a segment-processing pipeline is configured.
#[cfg(feature = "pipeline")]
pub(crate) struct WorkerHandle {
shutdown: Option<tokio::sync::oneshot::Sender<Duration>>,
thread: Option<JoinHandle<()>>,
}
#[cfg(feature = "pipeline")]
impl WorkerHandle {
/// Wrap the worker's shutdown sender and join handle.
pub(crate) fn new(
shutdown: tokio::sync::oneshot::Sender<Duration>,
thread: JoinHandle<()>,
) -> Self {
Self {
shutdown: Some(shutdown),
thread: Some(thread),
}
}
}
/// Owns the recording state: the [`Dial9Handle`], the flush thread, and (with
/// the `pipeline` feature) the background worker.
///
/// This is an RAII guard: dropping it flushes remaining events, seals the final
/// segment, and stops the worker. For a bounded drain of the background worker
/// (symbolize, compress, upload) call [`graceful_shutdown`](Self::graceful_shutdown)
/// instead.
pub struct Recorder {
handle: Dial9Handle,
flush_thread: Option<JoinHandle<()>>,
/// Hooks run once, with the handle, on the first `enable()`.
recording_start_hooks: Mutex<Vec<RecordingStartHook>>,
/// Held while this is the process's recorder. Dropping it frees the slot.
sole_recorder: Option<SoleRecorderGuard>,
#[cfg(feature = "pipeline")]
worker: Option<WorkerHandle>,
}
impl std::fmt::Debug for Recorder {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Recorder")
.field("enabled", &self.handle.is_enabled())
.field("recording", &self.flush_thread.is_some())
.finish_non_exhaustive()
}
}
/// A hook run once, with the live [`Dial9Handle`], when the recorder first
/// enables recording.
pub type RecordingStartHook = Box<dyn FnOnce(&Dial9Handle) + Send>;
impl Recorder {
/// Create a recorder from an existing handle and flush thread.
pub(crate) fn new(handle: Dial9Handle, flush_thread: Option<JoinHandle<()>>) -> Self {
Self {
handle,
flush_thread,
recording_start_hooks: Mutex::new(Vec::new()),
sole_recorder: None,
#[cfg(feature = "pipeline")]
worker: None,
}
}
/// Hold the process's recorder slot for this recorder's lifetime.
pub(crate) fn hold_process(&mut self, guard: crate::recorder::SoleRecorderGuard) {
self.sole_recorder = Some(guard);
}
/// Install the one-shot hooks to run on the first `enable()`.
pub(crate) fn set_recording_start_hooks(&self, hooks: Vec<RecordingStartHook>) {
*self.recording_start_hooks.lock().unwrap() = hooks;
}
/// Start recording over `shared`: build the recording [`Dial9Handle`], spawn
/// the flush thread that drains the bus into `writer`, and own its lifecycle.
///
/// The flush-thread control channel is created and owned internally; reach
/// the handle via [`handle`](Self::handle).
///
/// `thread_init` runs once on the flush thread before the loop and returns
/// a teardown closure run after it — use it to register/unregister the
/// thread with a runtime's profiler.
///
/// Pass `None` for `flush_metrics_sink` to discard flush metrics.
pub(crate) fn start<M, Init, Teardown>(
shared: Arc<SharedState>,
writer: SegmentWriter<M>,
flush_metrics_sink: Option<metrique::writer::BoxEntrySink>,
thread_init: Init,
) -> Self
where
M: BufferMode + Send + 'static,
Init: FnOnce() -> Teardown + Send + 'static,
Teardown: FnOnce(),
{
let (control_tx, control_rx) = mpsc::sync_channel(1);
let handle = Dial9Handle::enabled(shared.clone(), control_tx);
let flush_metrics_sink =
flush_metrics_sink.unwrap_or_else(metrique::writer::sink::DevNullSink::boxed);
let flush_thread = crate::primitives::thread::spawn_named("dial9-flush", move || {
// The flush thread is latency-tolerant; lower its priority.
#[cfg(target_os = "linux")]
// SAFETY: nice() is a simple syscall with no memory-safety
// implications; lowering priority is always permitted unprivileged.
unsafe {
let _ = libc::nice(10);
}
let teardown = thread_init();
run_flush_loop(control_rx, &shared, &flush_metrics_sink, writer);
teardown();
});
Self::new(handle, Some(flush_thread))
}
/// The recording handle for this recorder.
pub fn handle(&self) -> &Dial9Handle {
&self.handle
}
/// Publish this recorder's handle as the process-global one.
///
/// When set, [`Dial9Handle::current`] resolves on every thread in the
/// process. When not set, it resolves only on threads a runtime integration
/// has installed a handle on.
///
/// Returns [`InstallGlobalHandleError`] and changes nothing if another handle
/// is already installed: two live globals would split one process's events
/// across two traces. A recorder clears its own when it stops, so a later
/// install succeeds.
///
/// ```no_run
/// use dial9_core::buffer::MemoryBuffer;
/// use dial9_core::handle::Dial9Handle;
/// use dial9_core::recorder::recorder;
/// use dial9_trace_format::TraceEvent;
///
/// #[derive(TraceEvent)]
/// struct Tick {
/// #[traceevent(timestamp)]
/// timestamp_ns: u64,
/// }
///
/// let rec = recorder(MemoryBuffer::new(1 << 20)?).build();
/// rec.install_global_handle()?;
///
/// std::thread::spawn(|| {
/// // reachable here, with no handle plumbed in
/// Dial9Handle::current().record_event(Tick { timestamp_ns: 0 });
/// });
/// # Ok::<_, Box<dyn std::error::Error>>(())
/// ```
pub fn install_global_handle(&self) -> Result<(), InstallGlobalHandleError> {
crate::handle::set_global_handle(self.handle.clone())
}
/// Attach the background worker to this recorder, so its lifecycle is tied
/// to the recorder's (drained on `graceful_shutdown`, stopped on drop).
#[cfg(feature = "pipeline")]
pub(crate) fn attach_worker(&mut self, worker: WorkerHandle) {
self.worker = Some(worker);
}
crate::test_util_pub! {
/// The shared recording state.
fn shared(&self) -> Option<&Arc<SharedState>> {
self.handle.shared()
}
}
/// Monotonic start time of the recorder in nanoseconds.
pub fn start_time(&self) -> Option<u64> {
self.shared().map(|s| s.start_time_ns())
}
/// Enable recording.
pub fn enable(&self) {
self.handle.enable();
// Run the one-shot start hooks now that the handle is live and
// recording. Draining leaves them run-once across repeated enables.
let hooks = std::mem::take(&mut *self.recording_start_hooks.lock().unwrap());
for hook in hooks {
hook(&self.handle);
}
}
/// Disable recording.
pub fn disable(&self) {
self.handle.disable();
}
/// Flush remaining events, seal the final segment, and join the flush thread.
///
/// Call this before dropping any runtime state that owns worker threads, so
/// that their thread-local buffers have already been flushed to the central
/// collector.
pub(crate) fn stop_flush_thread(&mut self) {
// Clear the global before the blocking flush below, otherwise other threads
// keep resolving it and recording into buffers that nothing will drain.
if let Some(shared) = self.handle.shared() {
crate::handle::clear_global_handle_for(shared);
}
// Drain the calling thread's local buffer — it won't get a thread-stop
// hook, so any unflushed events would be lost otherwise.
if let Some(shared) = self.handle.shared() {
crate::encoder::drain_to_collector(&shared.collector);
}
// Tell the flush thread to do a final flush + finalize, then exit.
let (ack_tx, ack_rx) = mpsc::sync_channel(0);
if let Some(tx) = self.handle.control_tx()
&& tx.send(ControlCommand::FinalizeAndStop(ack_tx)).is_ok()
{
let _ = ack_rx.recv();
}
if let Some(t) = self.flush_thread.take() {
let _ = t.join();
}
// Stop is permanent from here: recording off, enable() and new
// attaches refused. Nothing drains sources once the flush thread is
// gone, so release them and whatever they own.
if let Some(shared) = self.handle.shared() {
shared.mark_stopped();
shared.clear_sources();
}
// Runtime threads drop their handle in a thread-stop hook, but the
// thread that attached the runtime gets no such hook and would hold a
// handle to a stopped recorder for the rest of its life.
crate::handle::clear_tl_handle();
}
/// Flush remaining events, seal the final segment, and (with `pipeline`)
/// wait for the background worker to drain within `timeout`.
///
/// Call this after any runtime that owns worker threads has been dropped, so
/// their thread-local buffers have already been flushed. Consumes the
/// recorder so `Drop` becomes a no-op.
///
/// Failures during draining are logged.
pub fn graceful_shutdown(mut self, timeout: Duration) {
// `timeout` only bounds the worker drain, which exists under `pipeline`.
#[cfg(not(feature = "pipeline"))]
let _ = timeout;
// 1. Flush + finalize the last segment.
self.stop_flush_thread();
// 2. Signal the worker to drain, then join it.
#[cfg(feature = "pipeline")]
if let Some(w) = &mut self.worker {
if let Some(tx) = w.shutdown.take() {
let _ = tx.send(timeout);
}
if let Some(t) = w.thread.take()
&& let Err(e) = t.join()
{
tracing::error!(target: "dial9", panic = ?e, "worker thread panicked during shutdown");
}
}
}
}
impl Drop for Recorder {
fn drop(&mut self) {
// 1. Flush + finalize. Idempotent, so a prior graceful_shutdown/stop is fine.
self.stop_flush_thread();
// 2. Hard shutdown: drop the sender without sending — the worker sees a
// closed channel and exits without draining. For a graceful drain, call
// graceful_shutdown() instead.
#[cfg(feature = "pipeline")]
if let Some(w) = &mut self.worker {
w.shutdown.take();
}
}
}