Skip to main content

reifydb_profiler/
layer.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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		// state_range shares the Flow category with flow::engine::apply, but it has no
415		// apply_time_us to override its duration with. If it took the apply branch the span
416		// would report 0us and every call site would look free; if it skipped interning the
417		// site the 11 call sites would collapse into one undifferentiated row, which is the
418		// whole reason the dimension exists. The site must land in dim 0 and the duration must
419		// come from the wall clock.
420		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		// A scan's cost tracks rows the storage fetched, not rows it returned: the timer probe
458		// asks for one row and sqlite hands back every tombstone in the prefix because LIMIT
459		// applies after the WHERE. Without these two counters the profile shows a slow probe and
460		// no way to tell an expensive scan from an expensive tombstone tax.
461		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}