1use std::sync::Arc;
5
6use reifydb_runtime::context::clock::{Clock, Instant};
7use reifydb_value::reifydb_assertions;
8use tracing::{
9 Metadata, Subscriber,
10 span::{Attributes, Id, Record},
11 subscriber::Interest,
12};
13use tracing_subscriber::{Layer, layer::Context, registry::LookupSpan};
14
15use crate::{
16 callsite,
17 category::{CategorySet, ProfilerCategory},
18 intern::DimInterner,
19 record::{DIM_UNSET, DimIdx, MAX_DIMENSIONS, MinimalSpanRecord},
20 scope::{ProfilerScope, REGISTRY, ScopeState, active_scope},
21 sink::ProfilerSink,
22 spec::spec_for,
23 visit::SpecFields,
24};
25
26pub struct ProfilerLayer {
27 sink: Arc<dyn ProfilerSink>,
28 categories: CategorySet,
29 interner: Arc<DimInterner>,
30 ambient_scope: Arc<ScopeState>,
31 clock: Clock,
32}
33
34impl ProfilerLayer {
35 pub fn new(
36 sink: Arc<dyn ProfilerSink>,
37 categories: CategorySet,
38 interner: Arc<DimInterner>,
39 clock: Clock,
40 ) -> Self {
41 let ambient_scope = ProfilerScope::ambient("profiler.global", Arc::clone(&sink), &clock);
42 Self {
43 sink,
44 categories,
45 interner,
46 ambient_scope,
47 clock,
48 }
49 }
50}
51
52#[derive(Clone)]
53struct SpanExt {
54 category: ProfilerCategory,
55 scope: Arc<ScopeState>,
56 callsite_id: u64,
57 started_at: Instant,
58 child_us: u64,
59 fields: Option<SpecFields>,
60}
61
62fn metadata_callsite_id(metadata: &'static Metadata<'static>) -> u64 {
63 let ptr: *const Metadata<'static> = metadata;
64 ptr as usize as u64
65}
66
67fn discover_scope<S>(ctx: &Context<'_, S>, span_id: &Id) -> Option<Arc<ScopeState>>
68where
69 S: Subscriber + for<'a> LookupSpan<'a>,
70{
71 if let Some(span) = ctx.span(span_id) {
72 for ancestor in span.scope() {
73 let ext = ancestor.extensions();
74 if let Some(ent) = ext.get::<SpanExt>() {
75 return Some(Arc::clone(&ent.scope));
76 }
77 }
78 }
79 let id = active_scope()?;
80 REGISTRY.get(id)
81}
82
83impl<S> Layer<S> for ProfilerLayer
84where
85 S: Subscriber + for<'a> LookupSpan<'a>,
86{
87 fn register_callsite(&self, metadata: &'static Metadata<'static>) -> Interest {
88 if !metadata.is_span() {
89 return Interest::never();
90 }
91 let Some(c) = ProfilerCategory::from_span_name(metadata.name()) else {
92 return Interest::never();
93 };
94 let Some(max) = self.categories.level_for(c) else {
95 return Interest::never();
96 };
97 if max.admits(metadata.level()) {
98 Interest::always()
99 } else {
100 Interest::never()
101 }
102 }
103
104 fn enabled(&self, metadata: &Metadata<'_>, _ctx: Context<'_, S>) -> bool {
105 if !metadata.is_span() {
106 return false;
107 }
108 let Some(c) = ProfilerCategory::from_span_name(metadata.name()) else {
109 return false;
110 };
111 self.categories.level_for(c).map(|max| max.admits(metadata.level())).unwrap_or(false)
112 }
113
114 fn on_new_span(&self, attrs: &Attributes<'_>, id: &Id, ctx: Context<'_, S>) {
115 let metadata = attrs.metadata();
116 let Some(category) = self.resolve_admitted_category(metadata) else {
117 return;
118 };
119 let scope = self.discover_span_scope(&ctx, id);
120 let callsite_id = self.register_span_callsite(metadata);
121 let fields = self.extract_spec_fields(attrs, metadata);
122 self.insert_span_ext(&ctx, id, category, scope, callsite_id, fields);
123 }
124
125 fn on_record(&self, id: &Id, values: &Record<'_>, ctx: Context<'_, S>) {
126 let Some(span) = ctx.span(id) else {
127 return;
128 };
129 let mut ext = span.extensions_mut();
130 if let Some(entry) = ext.get_mut::<SpanExt>()
131 && let Some(fields) = entry.fields.as_mut()
132 {
133 values.record(fields);
134 }
135 }
136
137 fn on_close(&self, id: Id, ctx: Context<'_, S>) {
138 let Some(entry) = self.take_span_ext(&ctx, &id) else {
139 return;
140 };
141 let record = self.build_record(&entry);
142 self.charge_parent(&ctx, &id, record.duration_us as u64);
143 self.emit_record(&entry.scope, record);
144 }
145}
146
147impl ProfilerLayer {
148 #[inline]
149 fn resolve_admitted_category(&self, metadata: &'static Metadata<'static>) -> Option<ProfilerCategory> {
150 let category = ProfilerCategory::from_span_name(metadata.name())?;
151 let max = self.categories.level_for(category)?;
152 if max.admits(metadata.level()) {
153 Some(category)
154 } else {
155 None
156 }
157 }
158
159 #[inline]
160 fn register_span_callsite(&self, metadata: &'static Metadata<'static>) -> u64 {
161 let callsite_id = metadata_callsite_id(metadata);
162 callsite::register(callsite_id, metadata.name());
163 callsite_id
164 }
165
166 #[inline]
167 fn extract_spec_fields(
168 &self,
169 attrs: &Attributes<'_>,
170 metadata: &'static Metadata<'static>,
171 ) -> Option<SpecFields> {
172 let spec = spec_for(metadata.name())?;
173 let mut fields = SpecFields::new(spec);
174 attrs.record(&mut fields);
175 Some(fields)
176 }
177
178 #[inline]
179 fn build_record(&self, entry: &SpanExt) -> MinimalSpanRecord {
180 let mut record = MinimalSpanRecord::new(entry.category, entry.callsite_id, 0);
181 match entry.fields.as_ref().and_then(SpecFields::duration_override) {
182 Some(override_us) => {
183 record.duration_us = u32::try_from(override_us).unwrap_or(u32::MAX);
184 }
185 None => {
186 reifydb_assertions! {
187 let now = self.clock.instant();
188 assert!(
189 now >= entry.started_at,
190 "span close observed a clock instant before its start, so elapsed() would saturate \
191 to zero and silently report a missing duration for a wall-clock span; the \
192 profiler relies on a monotonic clock for span timing (now_us={}, started_us={})",
193 now.elapsed().as_micros(),
194 entry.started_at.elapsed().as_micros()
195 );
196 }
197 let elapsed = entry.started_at.elapsed().as_micros();
198 record.duration_us = u32::try_from(elapsed).unwrap_or(u32::MAX);
199 }
200 }
201 if let Some(fields) = &entry.fields {
202 let mut dims: [DimIdx; MAX_DIMENSIONS] = [DIM_UNSET; MAX_DIMENSIONS];
203 for (slot, label) in fields.dims().iter().enumerate() {
204 if !label.is_empty() {
205 dims[slot] = self.interner.intern(label);
206 }
207 }
208 record.dim_indices = dims;
209 record.extras = *fields.extras();
210 }
211 record.self_us =
212 u32::try_from((record.duration_us as u64).saturating_sub(entry.child_us)).unwrap_or(u32::MAX);
213 record
214 }
215
216 #[inline]
217 fn charge_parent<S>(&self, ctx: &Context<'_, S>, id: &Id, elapsed_us: u64)
218 where
219 S: Subscriber + for<'a> LookupSpan<'a>,
220 {
221 let Some(span) = ctx.span(id) else {
222 return;
223 };
224 let Some(parent) = span.parent() else {
225 return;
226 };
227 for ancestor in parent.scope() {
228 let mut ext = ancestor.extensions_mut();
229 if let Some(entry) = ext.get_mut::<SpanExt>() {
230 entry.child_us = entry.child_us.saturating_add(elapsed_us);
231 return;
232 }
233 }
234 }
235
236 #[inline]
237 fn emit_record(&self, scope: &ScopeState, record: MinimalSpanRecord) {
238 self.sink.on_span_record(&record);
239 scope.attach_interner(Arc::clone(&self.interner));
240 scope.push(record);
241 }
242}
243
244impl ProfilerLayer {
245 #[inline]
246 fn discover_span_scope<S>(&self, ctx: &Context<'_, S>, id: &Id) -> Arc<ScopeState>
247 where
248 S: Subscriber + for<'a> LookupSpan<'a>,
249 {
250 discover_scope(ctx, id).unwrap_or_else(|| Arc::clone(&self.ambient_scope))
251 }
252
253 #[inline]
254 fn insert_span_ext<S>(
255 &self,
256 ctx: &Context<'_, S>,
257 id: &Id,
258 category: ProfilerCategory,
259 scope: Arc<ScopeState>,
260 callsite_id: u64,
261 fields: Option<SpecFields>,
262 ) where
263 S: Subscriber + for<'a> LookupSpan<'a>,
264 {
265 let ext = SpanExt {
266 category,
267 scope,
268 callsite_id,
269 started_at: self.clock.instant(),
270 child_us: 0,
271 fields,
272 };
273 if let Some(span) = ctx.span(id) {
274 span.extensions_mut().insert(ext);
275 }
276 }
277
278 #[inline]
279 fn take_span_ext<S>(&self, ctx: &Context<'_, S>, id: &Id) -> Option<SpanExt>
280 where
281 S: Subscriber + for<'a> LookupSpan<'a>,
282 {
283 let span = ctx.span(id)?;
284 span.extensions_mut().remove::<SpanExt>()
285 }
286}
287
288#[cfg(test)]
289mod tests {
290 use std::sync::Arc;
291
292 use reifydb_runtime::{context::clock::Clock, sync::mutex::Mutex as StdMutex};
293 use tracing::{debug_span, subscriber::with_default, trace_span};
294 use tracing_subscriber::{Registry, layer::SubscriberExt};
295
296 use super::*;
297 use crate::{
298 category::ProfilerLevel, record::DIM_UNSET, scope::ProfilerScope, sink::ProfilerSink,
299 summary::ProfilerSummary,
300 };
301
302 #[derive(Default)]
303 struct RecordingSink {
304 records: StdMutex<Vec<MinimalSpanRecord>>,
305 summaries: StdMutex<Vec<ProfilerSummary>>,
306 }
307
308 impl ProfilerSink for RecordingSink {
309 fn on_span_record(&self, record: &MinimalSpanRecord) {
310 self.records.lock().push(*record);
311 }
312 fn on_scope_closed(&self, summary: &ProfilerSummary) {
313 self.summaries.lock().push(summary.clone());
314 }
315 fn on_scope_batch(&self, summary: &ProfilerSummary) {
316 self.summaries.lock().push(summary.clone());
317 }
318 }
319
320 fn build_layer(sink: Arc<dyn ProfilerSink>, categories: CategorySet) -> (ProfilerLayer, Arc<DimInterner>) {
321 let interner = Arc::new(DimInterner::new());
322 (ProfilerLayer::new(sink, categories, interner.clone(), Clock::Real), interner)
323 }
324
325 #[test]
326 fn out_of_scope_span_captured_via_ambient_scope() {
327 let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
328 let (layer, _interner) = build_layer(sink.clone(), CategorySet::all());
329 let subscriber = Registry::default().with(layer);
330 with_default(subscriber, || {
331 let span = debug_span!("flow::engine::apply", node_id = "n1", node_type = "map");
332 let _g = span.enter();
333 });
334 let recs = sink.records.lock();
335 assert_eq!(recs.len(), 1, "ambient scope must capture unscoped tracked spans for always-on profiling");
336 }
337
338 #[test]
339 fn admits_at_or_below_category_level() {
340 let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
341 let categories = CategorySet::empty().with_level(ProfilerCategory::Flow, ProfilerLevel::Debug);
342 let (layer, _interner) = build_layer(sink.clone(), categories);
343 let subscriber = Registry::default().with(layer);
344
345 with_default(subscriber, || {
346 let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
347 handle.run_sync(|| {
348 let trace_span =
349 trace_span!("flow::engine::apply", node_id = "n", node_type = "trace_op");
350 let _g1 = trace_span.enter();
351 let debug_span = debug_span!("flow::engine::process_batch", batch_size = 1u64);
352 let _g2 = debug_span.enter();
353 });
354 let _ = handle.finish();
355 });
356
357 let recs = sink.records.lock();
358 assert_eq!(
359 recs.len(),
360 1,
361 "DEBUG-level filter must admit process_batch (debug) but reject apply (trace), got: {:?}",
362 recs
363 );
364 }
365
366 #[test]
367 fn disabled_category_short_circuits() {
368 let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
369 let (layer, _interner) = build_layer(sink.clone(), CategorySet::empty().with(ProfilerCategory::Query));
370 let subscriber = Registry::default().with(layer);
371 with_default(subscriber, || {
372 let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
373 handle.run_sync(|| {
374 let _ = trace_span!("flow::engine::apply", node_id = "n", node_type = "m");
375 });
376 let _ = handle.finish();
377 });
378 assert!(sink.records.lock().is_empty());
379 }
380
381 #[test]
382 fn flow_apply_captures_fields_and_interns_dims() {
383 let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
384 let (layer, _interner) = build_layer(sink.clone(), CategorySet::all());
385 let subscriber = Registry::default().with(layer);
386 with_default(subscriber, || {
387 let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
388 handle.run_sync(|| {
389 let span = trace_span!(
390 "flow::engine::apply",
391 node_id = "n1",
392 node_type = "map",
393 input_rows = 10u64,
394 output_rows = 5u64,
395 apply_time_us = 200u64,
396 lock_wait_us = 3u64,
397 );
398 let _g = span.enter();
399 });
400 let _summary = handle.finish();
401 });
402 let recs = sink.records.lock();
403 assert_eq!(recs.len(), 1);
404 let rec = recs[0];
405 assert_eq!(rec.category(), ProfilerCategory::Flow);
406 assert_eq!(rec.duration_us, 200);
407 assert_eq!(rec.extras[0], 10);
408 assert_eq!(rec.extras[1], 5);
409 assert_eq!(rec.extras[2], 3);
410 }
411
412 #[test]
413 fn state_range_interns_its_site_without_borrowing_the_apply_duration_override() {
414 let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
421 let (layer, interner) = build_layer(sink.clone(), CategorySet::all());
422 let subscriber = Registry::default().with(layer);
423 with_default(subscriber, || {
424 let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
425 handle.run_sync(|| {
426 let span = trace_span!(
427 "flow::state::range_limited",
428 operator_id = 7u64,
429 site = "timer::hydrate_probe"
430 );
431 let _g = span.enter();
432 });
433 let _summary = handle.finish();
434 });
435 let recs = sink.records.lock();
436 assert_eq!(recs.len(), 1);
437 let rec = recs[0];
438 assert_eq!(rec.category(), ProfilerCategory::Flow);
439 assert_eq!(
440 interner.resolve(rec.dim_indices[0]).as_deref(),
441 Some("timer::hydrate_probe"),
442 "the call site must be interned into dim 0 so the profile table can split by it"
443 );
444 assert_ne!(
445 rec.dim_indices[0], DIM_UNSET,
446 "an unset dim would render as <no-dims> and merge every call site into one row"
447 );
448 assert_eq!(
449 interner.resolve(rec.dim_indices[1]).as_deref(),
450 Some("op7"),
451 "the operator must occupy dim 1 so one hot wheel cannot hide inside a site-wide average"
452 );
453 }
454
455 #[test]
456 fn state_range_carries_its_physical_row_counts_into_extras() {
457 let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
462 let (layer, _interner) = build_layer(sink.clone(), CategorySet::all());
463 let subscriber = Registry::default().with(layer);
464 with_default(subscriber, || {
465 let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
466 handle.run_sync(|| {
467 let span = trace_span!(
468 "flow::state::range_limited",
469 operator_id = 7u64,
470 site = "timer::hydrate_probe",
471 rows_fetched = 431u64,
472 rows_tombstoned = 430u64
473 );
474 let _g = span.enter();
475 });
476 let _summary = handle.finish();
477 });
478 let recs = sink.records.lock();
479 assert_eq!(recs.len(), 1);
480 let rec = recs[0];
481 assert_eq!(rec.extras[0], 431, "rows fetched from storage must reach extras[0]");
482 assert_eq!(rec.extras[1], 430, "of which tombstones must reach extras[1]");
483 }
484
485 #[test]
486 fn ancestor_walk_inherits_scope_across_spans() {
487 let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
488 let (layer, _interner) = build_layer(sink.clone(), CategorySet::all());
489 let subscriber = Registry::default().with(layer);
490 with_default(subscriber, || {
491 let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
492 handle.run_sync(|| {
493 let outer = debug_span!("flow::engine::process_batch", batch_size = 3u64);
494 let _g = outer.enter();
495 let inner = trace_span!("flow::engine::apply", node_id = "n", node_type = "filter");
496 let _g2 = inner.enter();
497 });
498 let _ = handle.finish();
499 });
500 let recs = sink.records.lock();
501 assert_eq!(recs.len(), 2);
502 }
503}