1use 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}