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::{DimIdx, MAX_EXTRAS, MinimalSpanRecord},
20	scope::{ProfilerScope, REGISTRY, ScopeState, active_scope},
21	sink::ProfilerSink,
22	visit::FlowApplyFields,
23};
24
25pub struct ProfilerLayer {
26	sink: Arc<dyn ProfilerSink>,
27	categories: CategorySet,
28	interner: Arc<DimInterner>,
29	ambient_scope: Arc<ScopeState>,
30	clock: Clock,
31}
32
33impl ProfilerLayer {
34	pub fn new(
35		sink: Arc<dyn ProfilerSink>,
36		categories: CategorySet,
37		interner: Arc<DimInterner>,
38		clock: Clock,
39	) -> Self {
40		let ambient_scope = ProfilerScope::ambient("profiler.global", Arc::clone(&sink), &clock);
41		Self {
42			sink,
43			categories,
44			interner,
45			ambient_scope,
46			clock,
47		}
48	}
49}
50
51#[derive(Clone)]
52struct SpanExt {
53	category: ProfilerCategory,
54	scope: Arc<ScopeState>,
55	callsite_id: u64,
56	started_at: Instant,
57	flow_fields: Option<FlowApplyFields>,
58}
59
60fn metadata_callsite_id(metadata: &'static Metadata<'static>) -> u64 {
61	let ptr: *const Metadata<'static> = metadata;
62	ptr as usize as u64
63}
64
65fn discover_scope<S>(ctx: &Context<'_, S>, span_id: &Id) -> Option<Arc<ScopeState>>
66where
67	S: Subscriber + for<'a> LookupSpan<'a>,
68{
69	if let Some(span) = ctx.span(span_id) {
70		for ancestor in span.scope() {
71			let ext = ancestor.extensions();
72			if let Some(ent) = ext.get::<SpanExt>() {
73				return Some(Arc::clone(&ent.scope));
74			}
75		}
76	}
77	let id = active_scope()?;
78	REGISTRY.get(id)
79}
80
81impl<S> Layer<S> for ProfilerLayer
82where
83	S: Subscriber + for<'a> LookupSpan<'a>,
84{
85	fn register_callsite(&self, metadata: &'static Metadata<'static>) -> Interest {
86		if !metadata.is_span() {
87			return Interest::never();
88		}
89		let Some(c) = ProfilerCategory::from_span_name(metadata.name()) else {
90			return Interest::never();
91		};
92		let Some(max) = self.categories.level_for(c) else {
93			return Interest::never();
94		};
95		if max.admits(metadata.level()) {
96			Interest::always()
97		} else {
98			Interest::never()
99		}
100	}
101
102	fn enabled(&self, metadata: &Metadata<'_>, _ctx: Context<'_, S>) -> bool {
103		if !metadata.is_span() {
104			return false;
105		}
106		let Some(c) = ProfilerCategory::from_span_name(metadata.name()) else {
107			return false;
108		};
109		self.categories.level_for(c).map(|max| max.admits(metadata.level())).unwrap_or(false)
110	}
111
112	fn on_new_span(&self, attrs: &Attributes<'_>, id: &Id, ctx: Context<'_, S>) {
113		let metadata = attrs.metadata();
114		let Some(category) = self.resolve_admitted_category(metadata) else {
115			return;
116		};
117		let scope = self.discover_span_scope(&ctx, id);
118		let callsite_id = self.register_span_callsite(metadata);
119		let flow_fields = self.extract_flow_fields(category, attrs, metadata);
120		self.insert_span_ext(&ctx, id, category, scope, callsite_id, flow_fields);
121	}
122
123	fn on_record(&self, id: &Id, values: &Record<'_>, ctx: Context<'_, S>) {
124		let Some(span) = ctx.span(id) else {
125			return;
126		};
127		let mut ext = span.extensions_mut();
128		if let Some(entry) = ext.get_mut::<SpanExt>()
129			&& let Some(fields) = entry.flow_fields.as_mut()
130		{
131			values.record(fields);
132		}
133	}
134
135	fn on_close(&self, id: Id, ctx: Context<'_, S>) {
136		let Some(entry) = self.take_span_ext(&ctx, &id) else {
137			return;
138		};
139		let record = self.build_record(&entry);
140		self.emit_record(&entry.scope, record);
141	}
142}
143
144impl ProfilerLayer {
145	#[inline]
146	fn resolve_admitted_category(&self, metadata: &'static Metadata<'static>) -> Option<ProfilerCategory> {
147		let category = ProfilerCategory::from_span_name(metadata.name())?;
148		let max = self.categories.level_for(category)?;
149		if max.admits(metadata.level()) {
150			Some(category)
151		} else {
152			None
153		}
154	}
155
156	#[inline]
157	fn register_span_callsite(&self, metadata: &'static Metadata<'static>) -> u64 {
158		let callsite_id = metadata_callsite_id(metadata);
159		callsite::register(callsite_id, metadata.name());
160		callsite_id
161	}
162
163	#[inline]
164	fn extract_flow_fields(
165		&self,
166		category: ProfilerCategory,
167		attrs: &Attributes<'_>,
168		metadata: &'static Metadata<'static>,
169	) -> Option<FlowApplyFields> {
170		if category == ProfilerCategory::Flow && metadata.name() == "flow::engine::apply" {
171			let mut fields = FlowApplyFields::default();
172			attrs.record(&mut fields);
173			Some(fields)
174		} else {
175			None
176		}
177	}
178
179	#[inline]
180	fn build_record(&self, entry: &SpanExt) -> MinimalSpanRecord {
181		let mut record = MinimalSpanRecord::new(entry.category, entry.callsite_id, 0);
182		match (entry.category, &entry.flow_fields) {
183			(ProfilerCategory::Flow, Some(f)) => {
184				record.duration_us = u32::try_from(f.apply_time_us).unwrap_or(u32::MAX);
185				let mut dims: [DimIdx; 2] = [0, 0];
186				if !f.node_type.is_empty() {
187					dims[0] = self.interner.intern(&f.node_type);
188				}
189				if !f.node_id.is_empty() {
190					dims[1] = self.interner.intern(&f.node_id);
191				}
192				record.dim_indices = dims;
193				let mut extras = [0u64; MAX_EXTRAS];
194				extras[0] = f.input_rows;
195				extras[1] = f.output_rows;
196				extras[2] = f.lock_wait_us;
197				record.extras = extras;
198			}
199			_ => {
200				reifydb_assertions! {
201					let now = self.clock.instant();
202					assert!(
203						now >= entry.started_at,
204						"span close observed a clock instant before its start, so elapsed() would saturate \
205						 to zero and silently report a missing duration for an untracked-flow span; the \
206						 profiler relies on a monotonic clock for span timing (now_us={}, started_us={})",
207						now.elapsed().as_micros(),
208						entry.started_at.elapsed().as_micros()
209					);
210				}
211				let elapsed = entry.started_at.elapsed().as_micros();
212				record.duration_us = u32::try_from(elapsed).unwrap_or(u32::MAX);
213			}
214		}
215		record
216	}
217
218	#[inline]
219	fn emit_record(&self, scope: &ScopeState, record: MinimalSpanRecord) {
220		self.sink.on_span_record(&record);
221		scope.attach_interner(Arc::clone(&self.interner));
222		scope.push(record);
223	}
224}
225
226impl ProfilerLayer {
227	#[inline]
228	fn discover_span_scope<S>(&self, ctx: &Context<'_, S>, id: &Id) -> Arc<ScopeState>
229	where
230		S: Subscriber + for<'a> LookupSpan<'a>,
231	{
232		discover_scope(ctx, id).unwrap_or_else(|| Arc::clone(&self.ambient_scope))
233	}
234
235	#[inline]
236	fn insert_span_ext<S>(
237		&self,
238		ctx: &Context<'_, S>,
239		id: &Id,
240		category: ProfilerCategory,
241		scope: Arc<ScopeState>,
242		callsite_id: u64,
243		flow_fields: Option<FlowApplyFields>,
244	) where
245		S: Subscriber + for<'a> LookupSpan<'a>,
246	{
247		let ext = SpanExt {
248			category,
249			scope,
250			callsite_id,
251			started_at: self.clock.instant(),
252			flow_fields,
253		};
254		if let Some(span) = ctx.span(id) {
255			span.extensions_mut().insert(ext);
256		}
257	}
258
259	#[inline]
260	fn take_span_ext<S>(&self, ctx: &Context<'_, S>, id: &Id) -> Option<SpanExt>
261	where
262		S: Subscriber + for<'a> LookupSpan<'a>,
263	{
264		let span = ctx.span(id)?;
265		span.extensions_mut().remove::<SpanExt>()
266	}
267}
268
269#[cfg(test)]
270mod tests {
271	use std::sync::Arc;
272
273	use reifydb_runtime::{context::clock::Clock, sync::mutex::Mutex as StdMutex};
274	use tracing::{debug_span, subscriber::with_default, trace_span};
275	use tracing_subscriber::{Registry, layer::SubscriberExt};
276
277	use super::*;
278	use crate::{category::ProfilerLevel, scope::ProfilerScope, sink::ProfilerSink, summary::ProfilerSummary};
279
280	#[derive(Default)]
281	struct RecordingSink {
282		records: StdMutex<Vec<MinimalSpanRecord>>,
283		summaries: StdMutex<Vec<ProfilerSummary>>,
284	}
285
286	impl ProfilerSink for RecordingSink {
287		fn on_span_record(&self, record: &MinimalSpanRecord) {
288			self.records.lock().push(*record);
289		}
290		fn on_scope_closed(&self, summary: &ProfilerSummary) {
291			self.summaries.lock().push(summary.clone());
292		}
293		fn on_scope_batch(&self, summary: &ProfilerSummary) {
294			self.summaries.lock().push(summary.clone());
295		}
296	}
297
298	fn build_layer(sink: Arc<dyn ProfilerSink>, categories: CategorySet) -> (ProfilerLayer, Arc<DimInterner>) {
299		let interner = Arc::new(DimInterner::new());
300		(ProfilerLayer::new(sink, categories, interner.clone(), Clock::Real), interner)
301	}
302
303	#[test]
304	fn out_of_scope_span_captured_via_ambient_scope() {
305		let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
306		let (layer, _interner) = build_layer(sink.clone(), CategorySet::all());
307		let subscriber = Registry::default().with(layer);
308		with_default(subscriber, || {
309			let span = debug_span!("flow::engine::apply", node_id = "n1", node_type = "map");
310			let _g = span.enter();
311		});
312		let recs = sink.records.lock();
313		assert_eq!(recs.len(), 1, "ambient scope must capture unscoped tracked spans for always-on profiling");
314	}
315
316	#[test]
317	fn admits_at_or_below_category_level() {
318		let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
319		let categories = CategorySet::empty().with_level(ProfilerCategory::Flow, ProfilerLevel::Debug);
320		let (layer, _interner) = build_layer(sink.clone(), categories);
321		let subscriber = Registry::default().with(layer);
322
323		with_default(subscriber, || {
324			let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
325			handle.run_sync(|| {
326				let trace_span =
327					trace_span!("flow::engine::apply", node_id = "n", node_type = "trace_op");
328				let _g1 = trace_span.enter();
329				let debug_span = debug_span!("flow::engine::process_batch", batch_size = 1u64);
330				let _g2 = debug_span.enter();
331			});
332			let _ = handle.finish();
333		});
334
335		let recs = sink.records.lock();
336		assert_eq!(
337			recs.len(),
338			1,
339			"DEBUG-level filter must admit process_batch (debug) but reject apply (trace), got: {:?}",
340			recs
341		);
342	}
343
344	#[test]
345	fn disabled_category_short_circuits() {
346		let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
347		let (layer, _interner) = build_layer(sink.clone(), CategorySet::empty().with(ProfilerCategory::Query));
348		let subscriber = Registry::default().with(layer);
349		with_default(subscriber, || {
350			let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
351			handle.run_sync(|| {
352				let _ = trace_span!("flow::engine::apply", node_id = "n", node_type = "m");
353			});
354			let _ = handle.finish();
355		});
356		assert!(sink.records.lock().is_empty());
357	}
358
359	#[test]
360	fn flow_apply_captures_fields_and_interns_dims() {
361		let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
362		let (layer, _interner) = build_layer(sink.clone(), CategorySet::all());
363		let subscriber = Registry::default().with(layer);
364		with_default(subscriber, || {
365			let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
366			handle.run_sync(|| {
367				let span = trace_span!(
368					"flow::engine::apply",
369					node_id = "n1",
370					node_type = "map",
371					input_rows = 10u64,
372					output_rows = 5u64,
373					apply_time_us = 200u64,
374					lock_wait_us = 3u64,
375				);
376				let _g = span.enter();
377			});
378			let _summary = handle.finish();
379		});
380		let recs = sink.records.lock();
381		assert_eq!(recs.len(), 1);
382		let rec = recs[0];
383		assert_eq!(rec.category(), ProfilerCategory::Flow);
384		assert_eq!(rec.duration_us, 200);
385		assert_eq!(rec.extras[0], 10);
386		assert_eq!(rec.extras[1], 5);
387		assert_eq!(rec.extras[2], 3);
388	}
389
390	#[test]
391	fn ancestor_walk_inherits_scope_across_spans() {
392		let sink: Arc<RecordingSink> = Arc::new(RecordingSink::default());
393		let (layer, _interner) = build_layer(sink.clone(), CategorySet::all());
394		let subscriber = Registry::default().with(layer);
395		with_default(subscriber, || {
396			let handle = ProfilerScope::start_with_sink("scope", sink.clone(), Clock::Real);
397			handle.run_sync(|| {
398				let outer = debug_span!("flow::engine::process_batch", batch_size = 3u64);
399				let _g = outer.enter();
400				let inner = trace_span!("flow::engine::apply", node_id = "n", node_type = "filter");
401				let _g2 = inner.enter();
402			});
403			let _ = handle.finish();
404		});
405		let recs = sink.records.lock();
406		assert_eq!(recs.len(), 2);
407	}
408}