Skip to main content

reifydb_profiler/
layer.rs

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