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::{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}