Skip to main content

mz_repr/explain/
tracing.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Tracing utilities for explainable plans.
11
12use std::fmt::{Debug, Display};
13use std::sync::Mutex;
14
15use tracing::{Level, span, subscriber};
16use tracing_core::{Interest, Metadata};
17use tracing_subscriber::{field, layer};
18
19use smallvec::SmallVec;
20
21/// A tracing layer used to accumulate a sequence of explainable plans.
22#[allow(missing_debug_implementations)]
23pub struct PlanTrace<T> {
24    /// A specific concrete path to find in this trace. If present,
25    /// [`PlanTrace::push`] will only collect traces if the current path is a
26    /// prefix of find.
27    filter: Option<SmallVec<[&'static str; 4]>>,
28    /// A path of segments identifying the spans in the current ancestor-or-self
29    /// chain. The current path is used when accumulating new `entries`.
30    path: Mutex<String>,
31    /// The first time when entering a span (None no span was entered yet).
32    start: Mutex<Option<std::time::Instant>>,
33    /// A path of times at which the spans in the current ancestor-or-self chain
34    /// were started. The duration since the last time is used when accumulating
35    /// new `entries`.
36    times: Mutex<Vec<std::time::Instant>>,
37    /// A sequence of entries associating for a specific plan type `T`.
38    entries: Mutex<Vec<TraceEntry<T>>>,
39}
40
41/// A struct created as a reflection of a [`trace_plan`] call.
42#[derive(Clone, Debug)]
43pub struct TraceEntry<T> {
44    /// The instant at which an entry was created.
45    ///
46    /// Used to impose global sorting when merging multiple `TraceEntry`
47    /// arrays in a single array.
48    pub instant: std::time::Instant,
49    /// The duration since the start of the enclosing span.
50    pub span_duration: std::time::Duration,
51    /// The duration since the start of the top-level span seen by the `PlanTrace`.
52    pub full_duration: std::time::Duration,
53    /// Ancestor chain of span names (root is first, parent is last).
54    pub path: String,
55    /// The plan produced this step.
56    pub plan: T,
57}
58
59/// Trace a fragment of type `T` to be emitted as part of an `EXPLAIN OPTIMIZER
60/// TRACE` output.
61///
62/// For best compatibility with the existing UI (which at the moment is the only
63/// sane way to look at such `EXPLAIN` traces), code instrumentation should
64/// adhere to the following constraints:
65///
66/// 1.  The plan type should be listed in the layers created in the
67///     `OptimizerTrace` constructor.
68/// 2.  Each `trace_plan` should be unique within it's enclosing span and should
69///     represent the result of the stage idenified by that span. In particular,
70///     this means that functions that call `trace_plan` more than once need to
71///     construct ad-hoc spans (see the iteration spans in the `Fixpoint`
72///     transform for example).
73///
74/// As a consequence of the second constraint, a sequence of paths such as
75/// ```text
76/// optimizer.foo.bar
77/// optimizer.foo.baz
78/// ```
79/// is not well-formed as it is missing the results of the prefix paths at the
80/// end:
81/// ```text
82/// optimizer.foo.bar
83/// optimizer.foo.baz
84/// optimizer.foo
85/// optimizer
86/// ```
87///
88/// Also, note that full paths can be repeated within a pipeline, but adjacent
89/// duplicates are interpreted as separete invocations. For example, the
90/// sub-sequence
91/// ```text
92/// ... // preceding stages
93/// optimizer.foo.bar // 1st call
94/// optimizer.foo.bar // 2nd call
95/// ... // following stages
96/// ```
97/// will be rendered by the UI as the following tree structure.
98/// ```text
99/// optimizer
100///   ... // following stages
101///   foo
102///     bar // 2nd call
103///     bar // 1st call
104///   ... // preceding stages
105/// ```
106pub fn trace_plan<T: Clone + 'static>(plan: &T) {
107    tracing::Span::current().with_subscriber(|(_id, subscriber)| {
108        if let Some(trace) = subscriber.downcast_ref::<PlanTrace<T>>() {
109            trace.push(plan)
110        }
111    });
112}
113
114/// Create a span identified by `segment` and trace `plan` in it.
115///
116/// This primitive is useful for instrumentic code, see this commit[^example]
117/// for an example.
118///
119/// [^example]: <https://github.com/MaterializeInc/materialize/commit/2ce93229>
120pub fn dbg_plan<S: Display, T: Clone + 'static>(segment: S, plan: &T) {
121    span!(target: "optimizer", Level::DEBUG, "segment", path.segment = %segment).in_scope(|| {
122        trace_plan(plan);
123    });
124}
125
126/// Create a span identified by `segment` and trace `misc` in it.
127///
128/// This primitive is useful for instrumentic code, see this commit[^example]
129/// for an example.
130///
131/// [^example]: <https://github.com/MaterializeInc/materialize/commit/2ce93229>
132pub fn dbg_misc<S: Display, T: Display>(segment: S, misc: T) {
133    span!(target: "optimizer", Level::DEBUG, "segment", path.segment = %segment).in_scope(|| {
134        trace_plan(&misc.to_string());
135    });
136}
137
138/// A helper struct for wrapping entries that represent the invocation context
139/// of a function or method call into an object that renders as their hash.
140///
141/// Useful when constructing path segments when instrumenting a function trace
142/// with additional debugging information.
143#[allow(missing_debug_implementations)]
144pub struct ContextHash(u64);
145
146impl ContextHash {
147    pub fn of<T: std::hash::Hash>(t: T) -> Self {
148        use std::collections::hash_map::DefaultHasher;
149        use std::hash::Hasher;
150
151        let mut h = DefaultHasher::new();
152        t.hash(&mut h);
153        ContextHash(h.finish())
154    }
155}
156
157impl Display for ContextHash {
158    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
159        write!(f, "{:x}", self.0 & 0xFFFFFFFu64) // show last 28 bits
160    }
161}
162
163/// A [`layer::Layer`] implementation for [`PlanTrace`].
164///
165/// Populates the `data` wrapped by the [`PlanTrace`] instance with
166/// [`TraceEntry`] values, one for each span with attached plan in its
167/// extensions map.
168impl<S, T> layer::Layer<S> for PlanTrace<T>
169where
170    S: subscriber::Subscriber,
171    T: 'static,
172{
173    fn on_new_span(
174        &self,
175        attrs: &span::Attributes<'_>,
176        _id: &span::Id,
177        _ctx: layer::Context<'_, S>,
178    ) {
179        // add segment to path
180        let mut path = self.path.lock().expect("path shouldn't be poisoned");
181        let segment = attrs.get_str("path.segment");
182        let segment = segment.unwrap_or_else(|| attrs.metadata().name().to_string());
183        if !path.is_empty() {
184            path.push('/');
185        }
186        path.push_str(segment.as_str());
187    }
188
189    fn on_enter(&self, _id: &span::Id, _ctx: layer::Context<'_, S>) {
190        let now = std::time::Instant::now();
191        // set start value on first ever on_enter
192        let mut start = self.start.lock().expect("start shouldn't be poisoned");
193        start.get_or_insert(now);
194        // push to time stack
195        let mut times = self.times.lock().expect("times shouldn't be poisoned");
196        times.push(now);
197    }
198
199    fn on_exit(&self, _id: &span::Id, _ctx: layer::Context<'_, S>) {
200        // truncate last segment from path
201        let mut path = self.path.lock().expect("path shouldn't be poisoned");
202        let new_len = path.rfind('/').unwrap_or(0);
203        path.truncate(new_len);
204        // pop from time stack
205        let mut times = self.times.lock().expect("times shouldn't be poisoned");
206        times.pop();
207    }
208}
209
210impl<S, T> layer::Filter<S> for PlanTrace<T>
211where
212    S: subscriber::Subscriber,
213    T: 'static + Clone,
214{
215    fn enabled(&self, meta: &Metadata<'_>, _cx: &layer::Context<'_, S>) -> bool {
216        self.is_enabled(meta)
217    }
218
219    fn callsite_enabled(&self, meta: &'static Metadata<'static>) -> Interest {
220        if self.is_enabled(meta) {
221            Interest::always()
222        } else {
223            Interest::never()
224        }
225    }
226}
227
228impl<T: 'static + Clone> PlanTrace<T> {
229    /// Create a new trace for plans of type `T` that will only accumulate
230    /// [`TraceEntry`] instances along the prefix of the given `path`.
231    pub fn new(filter: Option<SmallVec<[&'static str; 4]>>) -> Self {
232        Self {
233            filter,
234            path: Mutex::new(String::with_capacity(256)),
235            start: Mutex::new(None),
236            times: Mutex::new(Default::default()),
237            entries: Mutex::new(Default::default()),
238        }
239    }
240
241    /// Check if a subscriber layer of this kind will be interested in tracing
242    /// spans and events with the given metadata.
243    fn is_enabled(&self, meta: &Metadata<'_>) -> bool {
244        meta.is_span() && meta.target() == "optimizer"
245    }
246
247    /// Drain the trace data collected so far.
248    ///
249    /// Note that this method will mutate the internal state of the enclosing
250    /// [`PlanTrace`] even though its receiver is not `&mut self`. This quirk is
251    /// required because the tracing `Dispatch` does not have `downcast_mut` method.
252    pub fn drain_as_vec(&self) -> Vec<TraceEntry<T>> {
253        let mut entries = self.entries.lock().expect("entries shouldn't be poisoned");
254        entries.split_off(0)
255    }
256
257    /// Retrieve the trace data collected so far while leaving it in place.
258    pub fn collect_as_vec(&self) -> Vec<TraceEntry<T>> {
259        let entries = self.entries.lock().expect("entries shouldn't be poisoned");
260        (*entries).clone()
261    }
262
263    /// Find and return a clone of the [`TraceEntry`] for the given `path`.
264    pub fn find(&self, path: &str) -> Option<TraceEntry<T>>
265    where
266        T: Clone,
267    {
268        let entries = self.entries.lock().expect("entries shouldn't be poisoned");
269        entries.iter().find(|entry| entry.path == path).cloned()
270    }
271
272    /// Push a trace entry for the given `plan` to the current trace.
273    ///
274    /// This is a noop if
275    /// 1. the call is within a context without an enclosing span, or if
276    /// 2. [`PlanTrace::filter`] is set not equal to [`PlanTrace::current_path`].
277    fn push(&self, plan: &T)
278    where
279        T: Clone,
280    {
281        if let Some(current_path) = self.current_path() {
282            let times = self.times.lock().expect("times shouldn't be poisoned");
283            let start = self.start.lock().expect("start shouldn't is poisoned");
284            if let (Some(full_start), Some(span_start)) = (start.as_ref(), times.last()) {
285                let mut entries = self.entries.lock().expect("entries shouldn't be poisoned");
286                let time = std::time::Instant::now();
287                entries.push(TraceEntry {
288                    instant: time,
289                    span_duration: time.duration_since(*span_start),
290                    full_duration: time.duration_since(*full_start),
291                    path: current_path,
292                    plan: plan.clone(),
293                });
294            }
295        }
296    }
297
298    /// Helper method: get a copy of the current path.
299    ///
300    /// If [`PlanTrace::filter`] is set, this will also check the current path
301    /// against the `find` entry and return `None` if the two differ.
302    fn current_path(&self) -> Option<String> {
303        let path = self.path.lock().expect("path shouldn't be poisoned");
304        let path = path.as_str();
305        match self.filter.as_ref() {
306            Some(named_paths) => {
307                if named_paths.contains(&path) {
308                    Some(path.to_owned())
309                } else {
310                    None
311                }
312            }
313            None => Some(path.to_owned()),
314        }
315    }
316}
317
318/// Helper trait used to extract attributes of type `&'static str`.
319trait GetStr {
320    fn get_str(&self, key: &'static str) -> Option<String>;
321}
322
323impl<'a> GetStr for span::Attributes<'a> {
324    fn get_str(&self, key: &'static str) -> Option<String> {
325        let mut extract_str = ExtractStr::new(key);
326        self.record(&mut extract_str);
327        extract_str.val()
328    }
329}
330
331/// Helper struct that implements `field::Visit` and is used in the
332/// `GetStr::get_str` implementation for `span::Attributes`.
333struct ExtractStr {
334    key: &'static str,
335    val: Option<String>,
336}
337
338impl ExtractStr {
339    fn new(key: &'static str) -> Self {
340        Self { key, val: None }
341    }
342
343    fn val(self) -> Option<String> {
344        self.val
345    }
346}
347
348impl field::Visit for ExtractStr {
349    fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
350        if field.name() == self.key {
351            self.val = Some(value.to_string())
352        }
353    }
354
355    fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
356        if field.name() == self.key {
357            self.val = Some(format!("{value:?}"))
358        }
359    }
360}
361
362#[cfg(test)]
363mod test {
364    use mz_ore::instrument;
365    use tracing::dispatcher;
366    use tracing_subscriber::prelude::*;
367
368    use super::{PlanTrace, trace_plan};
369
370    #[mz_ore::test]
371    fn test_optimizer_trace() {
372        let subscriber = tracing_subscriber::registry().with(Some(PlanTrace::<String>::new(None)));
373        let dispatch = dispatcher::Dispatch::new(subscriber);
374
375        dispatcher::with_default(&dispatch, || {
376            optimize();
377        });
378
379        if let Some(trace) = dispatch.downcast_ref::<PlanTrace<String>>() {
380            let trace = trace.drain_as_vec();
381            assert_eq!(trace.len(), 5);
382            for (i, entry) in trace.into_iter().enumerate() {
383                let path = entry.path;
384                match i {
385                    0 => {
386                        assert_eq!(path, "optimize");
387                    }
388                    1 => {
389                        assert_eq!(path, "optimize/logical/my_optimization");
390                    }
391                    2 => {
392                        assert_eq!(path, "optimize/logical");
393                    }
394                    3 => {
395                        assert_eq!(path, "optimize/physical");
396                    }
397                    4 => {
398                        assert_eq!(path, "optimize");
399                    }
400                    _ => (),
401                }
402            }
403        }
404    }
405
406    #[instrument(level = "info")]
407    fn optimize() {
408        let mut plan = constant_plan(42);
409        trace_plan(&plan);
410        logical_optimizer(&mut plan);
411        physical_optimizer(&mut plan);
412        trace_plan(&plan);
413    }
414
415    #[instrument(level = "info", name = "logical")]
416    fn logical_optimizer(plan: &mut String) {
417        some_optimization(plan);
418        *plan = plan.replace("RawPlan", "LogicalPlan");
419        trace_plan(plan);
420    }
421
422    #[instrument(level = "info", name = "physical")]
423    fn physical_optimizer(plan: &mut String) {
424        *plan = plan.replace("LogicalPlan", "PhysicalPlan");
425        trace_plan(plan);
426    }
427
428    #[mz_ore::instrument(level = "debug", fields(path.segment ="my_optimization"))]
429    fn some_optimization(plan: &mut String) {
430        *plan = plan.replace("42", "47");
431        trace_plan(plan);
432    }
433
434    fn constant_plan(i: usize) -> String {
435        format!("RawPlan(#{})", i)
436    }
437}