Skip to main content

mz_compute/
metrics.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
10use std::sync::{Arc, Mutex};
11
12use mz_compute_client::metrics::{CommandMetrics, HistoryMetrics};
13use mz_ore::cast::CastFrom;
14use mz_ore::metric;
15use mz_ore::metrics::{
16    MakeCollectorOpts, MetricTag, MetricVisibility, MetricsRegistry, UIntGauge, raw,
17};
18use mz_repr::{GlobalId, SharedRow};
19use prometheus::core::{AtomicF64, GenericCounter};
20use prometheus::proto::LabelPair;
21use prometheus::{Histogram, HistogramVec, IntCounter};
22
23/// Metrics exposed by compute replicas.
24//
25// Most of the metrics here use the `raw` implementations, rather than the `DeleteOnDrop` wrappers
26// because their labels are fixed throughout the lifetime of the replica process. For example, any
27// metric labeled only by `worker_id` can be `raw` since the number of workers cannot change.
28//
29// Metrics that are labelled by a dimension that can change throughout the lifetime of the process
30// (such as `collection_id`) MUST NOT use the `raw` metric types and must use the `DeleteOnDrop`
31// types instead, to avoid memory leaks.
32#[derive(Clone, Debug)]
33pub struct ComputeMetrics {
34    // Optional workload class label to apply to all metrics in registry.
35    workload_class: Arc<Mutex<Option<String>>>,
36
37    // command history
38    history_command_count: raw::UIntGaugeVec,
39    history_dataflow_count: raw::UIntGaugeVec,
40
41    // reconciliation
42    reconciliation_reused_dataflows_count_total: raw::IntCounterVec,
43    reconciliation_replaced_dataflows_count_total: raw::IntCounterVec,
44
45    // arrangements
46    arrangement_maintenance_seconds_total: raw::CounterVec,
47    arrangement_maintenance_active_info: raw::UIntGaugeVec,
48
49    // logging
50    logging_records_total: raw::IntCounterVec,
51
52    // timings
53    //
54    // Note that this particular metric unfortunately takes some care to
55    // interpret. It measures the duration of step_or_park calls, which
56    // undesirably includes the parking. This is probably fine because we
57    // regularly send progress information through persist sources, which likely
58    // means the parking is capped at a second or two in practice. It also
59    // doesn't do anything to let you pinpoint _which_ operator or worker isn't
60    // yielding, but it should hopefully alert us when there is something to
61    // look at.
62    timely_step_duration_seconds: HistogramVec,
63    logging_step_duration_seconds: HistogramVec,
64    persist_peek_seconds: HistogramVec,
65    handle_command_duration_seconds: HistogramVec,
66
67    // Index peek timing phases (per-cluster, no worker label)
68    index_peek_total_seconds: Histogram,
69    index_peek_seek_fulfillment_seconds: Histogram,
70    index_peek_error_scan_seconds: Histogram,
71    index_peek_cursor_setup_seconds: Histogram,
72    index_peek_row_iteration_seconds: Histogram,
73    index_peek_row_iteration_rows: Histogram,
74    index_peek_result_sort_seconds: Histogram,
75    index_peek_result_sort_rows: Histogram,
76    index_peek_frontier_check_seconds: Histogram,
77    index_peek_row_collection_seconds: Histogram,
78    index_peek_walks_total: raw::IntCounterVec,
79    index_peek_stashed_total: IntCounter,
80    index_peek_permit_queue_depth: UIntGauge,
81    index_peek_permit_wait_seconds: Histogram,
82    index_peek_offload_seconds: Histogram,
83
84    // memory usage
85    shared_row_heap_capacity_bytes: raw::UIntGaugeVec,
86
87    // replica expiration
88    replica_expiration_timestamp_seconds: raw::UIntGaugeVec,
89    replica_expiration_remaining_seconds: raw::GaugeVec,
90
91    // collections
92    collection_count: raw::UIntGaugeVec,
93
94    // subscribes
95    subscribe_snapshots_skipped_total: IntCounter,
96
97    // metric sinks
98    metric_sink_registration_retries_total: IntCounter,
99}
100
101/// Applies the per-role const label to `opts`, unless `role` is `Solo`.
102///
103/// The two named roles (maintenance, interactive) each get a distinct `role` label so a second
104/// compute runtime in the same process registers a distinct series rather than colliding with the
105/// first. `Solo` omits the label so a single-runtime deployment registers exactly as it did before
106/// a second runtime existed.
107fn with_role(
108    mut opts: MakeCollectorOpts,
109    role: crate::server::ComputeRuntimeRole,
110) -> MakeCollectorOpts {
111    if let Some(label) = role.label() {
112        opts.opts = opts.opts.const_label("role", label);
113    }
114    opts
115}
116
117impl ComputeMetrics {
118    /// Registers the compute metrics for `role` into `registry`.
119    ///
120    /// The two named roles carry a `role` const label so that a second compute runtime in the same
121    /// process registers a distinct series rather than colliding with the first. `Solo` carries no
122    /// such label.
123    pub fn register_with(
124        registry: &MetricsRegistry,
125        role: crate::server::ComputeRuntimeRole,
126    ) -> Self {
127        let workload_class = Arc::new(Mutex::new(None));
128        let mut index_peek_row_buckets =
129            prometheus::exponential_buckets(1.0, 2.0, 25).expect("valid parameters");
130        index_peek_row_buckets.insert(0, 0.0);
131
132        // Apply a `workload_class` label to all metrics in the registry when we
133        // have a known workload class.
134        //
135        // The postprocessor rewrites every metric in the whole registry, so only the maintenance
136        // runtime registers it. A second registration from the interactive runtime would push the
137        // label twice onto each metric and produce a duplicate-label scrape error.
138        if role.owns_process_globals() {
139            registry.register_postprocessor({
140                let workload_class = Arc::clone(&workload_class);
141                move |metrics| {
142                    let workload_class: Option<String> =
143                        workload_class.lock().expect("lock poisoned").clone();
144                    let Some(workload_class) = workload_class else {
145                        return;
146                    };
147                    for metric in metrics {
148                        for metric in metric.mut_metric() {
149                            let mut label = LabelPair::default();
150                            label.set_name("workload_class".into());
151                            label.set_value(workload_class.clone());
152
153                            let mut labels = metric.take_label();
154                            labels.push(label);
155                            metric.set_label(labels);
156                        }
157                    }
158                }
159            });
160        }
161
162        Self {
163            workload_class,
164            history_command_count: registry.register(with_role(metric!(
165                name: "mz_compute_replica_history_command_count",
166                help: "The number of commands in the replica's command history.",
167                var_labels: ["worker_id", "command_type"],
168            ), role)),
169            history_dataflow_count: registry.register(with_role(metric!(
170                name: "mz_compute_replica_history_dataflow_count",
171                help: "The number of dataflows in the replica's command history.",
172                var_labels: ["worker_id"],
173                visibility: MetricVisibility::Public,
174                tags: [MetricTag::Compute],
175            ), role)),
176            reconciliation_reused_dataflows_count_total: registry.register(with_role(metric!(
177                name: "mz_compute_reconciliation_reused_dataflows_count_total",
178                help: "The total number of dataflows that were reused during compute reconciliation.",
179                var_labels: ["worker_id"],
180            ), role)),
181            reconciliation_replaced_dataflows_count_total: registry.register(with_role(metric!(
182                name: "mz_compute_reconciliation_replaced_dataflows_count_total",
183                help: "The total number of dataflows that were replaced during compute reconciliation.",
184                var_labels: ["worker_id", "reason"],
185            ), role)),
186            arrangement_maintenance_seconds_total: registry.register(with_role(metric!(
187                name: "mz_arrangement_maintenance_seconds_total",
188                help: "The total time spent maintaining arrangements.",
189                var_labels: ["worker_id"],
190                visibility: MetricVisibility::Public,
191                tags: [MetricTag::Compute],
192            ), role)),
193            arrangement_maintenance_active_info: registry.register(with_role(metric!(
194                name: "mz_arrangement_maintenance_active_info",
195                help: "Whether maintenance is currently occuring.",
196                var_labels: ["worker_id"],
197            ), role)),
198            timely_step_duration_seconds: registry.register(with_role(metric!(
199                name: "mz_timely_step_duration_seconds",
200                help: "The time spent in each compute step_or_park call",
201                const_labels: {"cluster" => "compute"},
202                var_labels: ["worker_id"],
203                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 32.0),
204            ), role)),
205            logging_step_duration_seconds: registry.register(with_role(metric!(
206                name: "mz_compute_logging_step_duration_seconds",
207                help: "The time spent in each scheduling of the logging dataflow.",
208                var_labels: ["worker_id"],
209                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_016, 8.0),
210            ), role)),
211            logging_records_total: registry.register(with_role(metric!(
212                name: "mz_compute_logging_records_total",
213                help: "The number of log records handed to the logging dataflow, by log.",
214                var_labels: ["worker_id", "log"],
215            ), role)),
216            shared_row_heap_capacity_bytes: registry.register(with_role(metric!(
217                name: "mz_dataflow_shared_row_heap_capacity_bytes",
218                help: "The heap capacity of the shared row.",
219                var_labels: ["worker_id"],
220            ), role)),
221            persist_peek_seconds: registry.register(with_role(metric!(
222                name: "mz_persist_peek_seconds",
223                help: "Time spent in (experimental) Persist fast-path peeks.",
224                var_labels: ["worker_id"],
225                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
226            ), role)),
227            handle_command_duration_seconds: registry.register(with_role(metric!(
228                name: "mz_cluster_handle_command_duration_seconds",
229                help: "Time spent in handling commands.",
230                const_labels: {"cluster" => "compute"},
231                var_labels: ["worker_id", "command_type"],
232                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
233            ), role)),
234            index_peek_total_seconds: registry.register(with_role(metric!(
235                name: "mz_index_peek_total_seconds",
236                help: "Time one visit to an index peek spent on the timely worker. A peek whose walk was offloaded contributes only the inline slice that offloaded it, and its time away from the worker is `mz_index_peek_offload_seconds`.",
237                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
238            ), role)),
239            index_peek_seek_fulfillment_seconds: registry.register(with_role(metric!(
240                name: "mz_index_peek_seek_fulfillment_seconds",
241                help: "Time in seek_fulfillment method including frontier checks and data collection.",
242                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
243            ), role)),
244            index_peek_error_scan_seconds: registry.register(with_role(metric!(
245                name: "mz_index_peek_error_scan_seconds",
246                help: "Time scanning the error trace for errors, summed over the slices the scan was cut into and observed only for scans that find no error.",
247                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
248            ), role)),
249            index_peek_cursor_setup_seconds: registry.register(with_role(metric!(
250                name: "mz_index_peek_cursor_setup_seconds",
251                help: "Time opening the trace cursor and sorting the literal constraints, excluding the seek to those literals.",
252                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
253            ), role)),
254            index_peek_row_iteration_seconds: registry.register(with_role(metric!(
255                name: "mz_index_peek_row_iteration_seconds",
256                help: "Time iterating rows, seeking the cursor to the literal constraints, and evaluating MFP, summed over the slices the walk was cut into.",
257                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
258            ), role)),
259            index_peek_row_iteration_rows: registry.register(with_role(metric!(
260                name: "mz_index_peek_row_iteration_rows",
261                help: "Number of arrangement rows evaluated by the index peek result iterator.",
262                buckets: index_peek_row_buckets.clone(),
263            ), role)),
264            index_peek_result_sort_seconds: registry.register(with_role(metric!(
265                name: "mz_index_peek_result_sort_seconds",
266                help: "Time thinning intermediate results down to the rows a peek's finishing needs.",
267                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
268            ), role)),
269            index_peek_result_sort_rows: registry.register(with_role(metric!(
270                name: "mz_index_peek_result_sort_rows",
271                help: "Number of intermediate result rows handed to thinning during peek collection, summed across the times it ran.",
272                buckets: index_peek_row_buckets,
273            ), role)),
274            index_peek_frontier_check_seconds: registry.register(with_role(metric!(
275                name: "mz_index_peek_frontier_check_seconds",
276                help: "Time checking trace frontiers.",
277                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
278            ), role)),
279            index_peek_row_collection_seconds: registry.register(with_role(metric!(
280                name: "mz_index_peek_row_collection_seconds",
281                help: "Time constructing RowCollection from peek results, including converting the row counts the scan produced.",
282                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
283            ), role)),
284            index_peek_walks_total: registry.register(with_role(metric!(
285                name: "mz_index_peek_walks_total",
286                help: "The number of index peek walks that reached an outcome, by the substrate they ended on: `inline` on the timely worker, `offloaded` away from it.",
287                var_labels: ["substrate"],
288            ), role)),
289            index_peek_stashed_total: registry.register(with_role(metric!(
290                name: "mz_index_peek_stashed_total",
291                help: "The number of index peek walks that answered with a handle to the peek response stash, always a subset of the `offloaded` substrate of `mz_index_peek_walks_total`.",
292            ), role)),
293            index_peek_permit_queue_depth: registry.register(with_role(metric!(
294                name: "mz_index_peek_permit_queue_depth",
295                help: "The number of offloaded index peek walks waiting for a permit to run.",
296            ), role)),
297            index_peek_permit_wait_seconds: registry.register(with_role(metric!(
298                name: "mz_index_peek_permit_wait_seconds",
299                help: "Time an offloaded index peek walk waited for a permit, observed only for walks that were admitted.",
300                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
301            ), role)),
302            index_peek_offload_seconds: registry.register(with_role(metric!(
303                name: "mz_index_peek_offload_seconds",
304                help: "Wall-clock time an offloaded index peek walk spent away from the timely worker, including the wait for a permit.",
305                buckets: mz_ore::stats::histogram_seconds_buckets(0.000_128, 8.0),
306            ), role)),
307            replica_expiration_timestamp_seconds: registry.register(with_role(metric!(
308                name: "mz_dataflow_replica_expiration_timestamp_seconds",
309                help: "The replica expiration timestamp in seconds since epoch.",
310                var_labels: ["worker_id"],
311            ), role)),
312            replica_expiration_remaining_seconds: registry.register(with_role(metric!(
313                name: "mz_dataflow_replica_expiration_remaining_seconds",
314                help: "The remaining seconds until replica expiration. Can go negative, can lag behind.",
315                var_labels: ["worker_id"],
316            ), role)),
317            collection_count: registry.register(with_role(metric!(
318                name: "mz_compute_collection_count",
319                help: "The number and hydration status of maintained compute collections.",
320                var_labels: ["worker_id", "type", "hydrated"],
321            ), role)),
322            subscribe_snapshots_skipped_total: registry.register(with_role(metric!(
323                name: "mz_subscribe_snapshots_skipped_total",
324                help: "The number of collection snapshots that were skipped by the subscribe snapshot optimization.",
325            ), role)),
326            metric_sink_registration_retries_total: registry.register(with_role(metric!(
327                name: "mz_compute_metric_sink_registration_retries_total",
328                help: "The number of times a metric sink failed to register its collector and scheduled a retry.",
329            ), role)),
330        }
331    }
332
333    /// Sets the workload class for the compute metrics.
334    pub fn set_workload_class(&self, workload_class: Option<String>) {
335        let mut guard = self.workload_class.lock().expect("lock poisoned");
336        *guard = workload_class
337    }
338
339    pub fn for_worker(&self, worker_id: usize) -> WorkerMetrics {
340        let worker = worker_id.to_string();
341        let arrangement_maintenance_seconds_total = self
342            .arrangement_maintenance_seconds_total
343            .with_label_values(&[&worker]);
344        let arrangement_maintenance_active_info = self
345            .arrangement_maintenance_active_info
346            .with_label_values(&[&worker]);
347        let timely_step_duration_seconds = self
348            .timely_step_duration_seconds
349            .with_label_values(&[&worker]);
350        let persist_peek_seconds = self.persist_peek_seconds.with_label_values(&[&worker]);
351        let handle_command_duration_seconds = CommandMetrics::build(|typ| {
352            self.handle_command_duration_seconds
353                .with_label_values(&[worker.as_ref(), typ])
354        });
355        let index_peek_total_seconds = self.index_peek_total_seconds.clone();
356        let index_peek_seek_fulfillment_seconds = self.index_peek_seek_fulfillment_seconds.clone();
357        let index_peek_error_scan_seconds = self.index_peek_error_scan_seconds.clone();
358        let index_peek_cursor_setup_seconds = self.index_peek_cursor_setup_seconds.clone();
359        let index_peek_row_iteration_seconds = self.index_peek_row_iteration_seconds.clone();
360        let index_peek_row_iteration_rows = self.index_peek_row_iteration_rows.clone();
361        let index_peek_result_sort_seconds = self.index_peek_result_sort_seconds.clone();
362        let index_peek_result_sort_rows = self.index_peek_result_sort_rows.clone();
363        let index_peek_frontier_check_seconds = self.index_peek_frontier_check_seconds.clone();
364        let index_peek_row_collection_seconds = self.index_peek_row_collection_seconds.clone();
365        let index_peek_walks_inline = self.index_peek_walks_total.with_label_values(&["inline"]);
366        let index_peek_walks_offloaded = self
367            .index_peek_walks_total
368            .with_label_values(&["offloaded"]);
369        let index_peek_stashed_total = self.index_peek_stashed_total.clone();
370        let index_peek_permit_queue_depth = self.index_peek_permit_queue_depth.clone();
371        let index_peek_permit_wait_seconds = self.index_peek_permit_wait_seconds.clone();
372        let index_peek_offload_seconds = self.index_peek_offload_seconds.clone();
373        let replica_expiration_timestamp_seconds = self
374            .replica_expiration_timestamp_seconds
375            .with_label_values(&[&worker]);
376        let replica_expiration_remaining_seconds = self
377            .replica_expiration_remaining_seconds
378            .with_label_values(&[&worker]);
379        let shared_row_heap_capacity_bytes = self
380            .shared_row_heap_capacity_bytes
381            .with_label_values(&[&worker]);
382
383        WorkerMetrics {
384            worker_label: worker,
385            metrics: self.clone(),
386            arrangement_maintenance_seconds_total,
387            arrangement_maintenance_active_info,
388            timely_step_duration_seconds,
389            persist_peek_seconds,
390            handle_command_duration_seconds,
391            index_peek_total_seconds,
392            index_peek_seek_fulfillment_seconds,
393            index_peek_error_scan_seconds,
394            index_peek_cursor_setup_seconds,
395            index_peek_row_iteration_seconds,
396            index_peek_row_iteration_rows,
397            index_peek_result_sort_seconds,
398            index_peek_result_sort_rows,
399            index_peek_frontier_check_seconds,
400            index_peek_row_collection_seconds,
401            index_peek_walks_inline,
402            index_peek_walks_offloaded,
403            index_peek_stashed_total,
404            index_peek_permit_queue_depth,
405            index_peek_permit_wait_seconds,
406            index_peek_offload_seconds,
407            replica_expiration_timestamp_seconds,
408            replica_expiration_remaining_seconds,
409            shared_row_heap_capacity_bytes,
410        }
411    }
412}
413
414/// Per-worker metrics of the logging dataflow.
415#[derive(Clone, Debug)]
416pub(crate) struct LoggingMetrics {
417    pub(crate) timely_records_total: IntCounter,
418    pub(crate) reachability_records_total: IntCounter,
419    pub(crate) differential_records_total: IntCounter,
420    pub(crate) compute_records_total: IntCounter,
421    /// Duration of each scheduling of the logging dataflow as a whole.
422    pub(crate) step_duration_seconds: Histogram,
423}
424
425/// Per-worker metrics.
426#[derive(Clone, Debug)]
427pub struct WorkerMetrics {
428    worker_label: String,
429    metrics: ComputeMetrics,
430
431    /// The amount of time spent in arrangement maintenance.
432    pub(crate) arrangement_maintenance_seconds_total: GenericCounter<AtomicF64>,
433    /// 1 if this worker is currently doing maintenance.
434    ///
435    /// If maintenance turns out to take a very long time, this will allow us
436    /// to gain a sense that Materialize is stuck on maintenance before the
437    /// maintenance completes
438    pub(crate) arrangement_maintenance_active_info: UIntGauge,
439    /// Histogram of Timely step timings.
440    pub(crate) timely_step_duration_seconds: Histogram,
441    /// Histogram of persist peek durations.
442    pub(crate) persist_peek_seconds: Histogram,
443    /// Histogram of command handling durations.
444    pub(crate) handle_command_duration_seconds: CommandMetrics<Histogram>,
445    /// Histogram of total index peek durations.
446    pub(crate) index_peek_total_seconds: Histogram,
447    /// Histogram of index peek seek_fulfillment durations.
448    pub(crate) index_peek_seek_fulfillment_seconds: Histogram,
449    /// Histogram of index peek error scan durations.
450    pub(crate) index_peek_error_scan_seconds: Histogram,
451    /// Histogram of index peek cursor setup durations.
452    pub(crate) index_peek_cursor_setup_seconds: Histogram,
453    /// Histogram of index peek row iteration durations.
454    pub(crate) index_peek_row_iteration_seconds: Histogram,
455    /// Histogram of index peek rows processed by the result iterator.
456    pub(crate) index_peek_row_iteration_rows: Histogram,
457    /// Histogram of index peek result sort durations.
458    pub(crate) index_peek_result_sort_seconds: Histogram,
459    /// Histogram of index peek rows sorted across all result sort operations.
460    pub(crate) index_peek_result_sort_rows: Histogram,
461    /// Histogram of index peek frontier check durations.
462    pub(crate) index_peek_frontier_check_seconds: Histogram,
463    /// Histogram of index peek row collection construction durations.
464    pub(crate) index_peek_row_collection_seconds: Histogram,
465    /// Counts index peek walks that ran on the timely worker.
466    ///
467    /// Both substrate series are resolved when the worker's metrics are built, so each exists at
468    /// zero before its first walk. A series at zero says the offload never engaged, where an absent
469    /// series says nothing.
470    pub(crate) index_peek_walks_inline: IntCounter,
471    /// Counts index peek walks that ran away from the timely worker.
472    pub(crate) index_peek_walks_offloaded: IntCounter,
473    /// Counts index peek walks that answered from the peek response stash.
474    ///
475    /// Resolved when the worker's metrics are built, so it reports zero before the first stashed
476    /// answer rather than being absent. Whether a peek reached the stash has no other signal.
477    pub(crate) index_peek_stashed_total: IntCounter,
478    /// How many offloaded index peek walks are waiting for a permit.
479    pub(crate) index_peek_permit_queue_depth: UIntGauge,
480    /// Histogram of how long an offloaded index peek walk waited for its permit.
481    pub(crate) index_peek_permit_wait_seconds: Histogram,
482    /// Histogram of how long an offloaded index peek walk was away from the worker.
483    pub(crate) index_peek_offload_seconds: Histogram,
484    /// The timestamp of replica expiration.
485    pub(crate) replica_expiration_timestamp_seconds: UIntGauge,
486    /// Remaining seconds until replica expiration.
487    pub(crate) replica_expiration_remaining_seconds: raw::Gauge,
488    /// Heap capacity of the shared row.
489    shared_row_heap_capacity_bytes: UIntGauge,
490}
491
492impl WorkerMetrics {
493    pub(crate) fn for_logging(&self) -> LoggingMetrics {
494        let records_total = |log| {
495            self.metrics
496                .logging_records_total
497                .with_label_values(&[self.worker_label.as_ref(), log])
498        };
499        LoggingMetrics {
500            timely_records_total: records_total("timely"),
501            reachability_records_total: records_total("reachability"),
502            differential_records_total: records_total("differential"),
503            compute_records_total: records_total("compute"),
504            step_duration_seconds: self
505                .metrics
506                .logging_step_duration_seconds
507                .with_label_values(&[&self.worker_label]),
508        }
509    }
510
511    pub fn for_history(&self) -> HistoryMetrics<UIntGauge> {
512        let command_counts = CommandMetrics::build(|typ| {
513            self.metrics
514                .history_command_count
515                .with_label_values(&[self.worker_label.as_ref(), typ])
516        });
517        let dataflow_count = self
518            .metrics
519            .history_dataflow_count
520            .with_label_values(&[&self.worker_label]);
521
522        HistoryMetrics {
523            command_counts,
524            dataflow_count,
525        }
526    }
527
528    /// Record the reconciliation result for a single dataflow.
529    ///
530    /// Reconciliation is recorded as successful if the given properties all hold. Otherwise it is
531    /// recorded as unsuccessful, with a reason based on the first property that does not hold.
532    ///
533    /// The properties are:
534    ///  * compatible: The old and new dataflow descriptions are compatible.
535    ///  * uncompacted: Collections currently installed for the dataflow exports have not been
536    ///                 allowed to compact beyond that new dataflow as-of.
537    ///  * subscribe_free: The dataflow does not export a subscribe sink.
538    ///  * copy_to_free: The dataflow does not export a copy-to sink.
539    ///  * dependencies_retained: All local inputs to the dataflow were retained by compute
540    ///                           reconciliation.
541    pub fn record_dataflow_reconciliation(
542        &self,
543        compatible: bool,
544        uncompacted: bool,
545        subscribe_free: bool,
546        copy_to_free: bool,
547        dependencies_retained: bool,
548    ) {
549        if !compatible {
550            self.metrics
551                .reconciliation_replaced_dataflows_count_total
552                .with_label_values(&[self.worker_label.as_ref(), "incompatible"])
553                .inc();
554        } else if !uncompacted {
555            self.metrics
556                .reconciliation_replaced_dataflows_count_total
557                .with_label_values(&[self.worker_label.as_ref(), "compacted"])
558                .inc();
559        } else if !subscribe_free {
560            self.metrics
561                .reconciliation_replaced_dataflows_count_total
562                .with_label_values(&[self.worker_label.as_ref(), "subscribe"])
563                .inc();
564        } else if !copy_to_free {
565            self.metrics
566                .reconciliation_replaced_dataflows_count_total
567                .with_label_values(&[self.worker_label.as_ref(), "copy-to"])
568                .inc();
569        } else if !dependencies_retained {
570            self.metrics
571                .reconciliation_replaced_dataflows_count_total
572                .with_label_values(&[self.worker_label.as_ref(), "dependency"])
573                .inc();
574        } else {
575            self.metrics
576                .reconciliation_reused_dataflows_count_total
577                .with_label_values(&[&self.worker_label])
578                .inc();
579        }
580    }
581
582    /// Record the heap capacity of the shared row.
583    pub fn record_shared_row_metrics(&self) {
584        let binding = SharedRow::get();
585        self.shared_row_heap_capacity_bytes
586            .set(u64::cast_from(binding.byte_capacity()));
587    }
588
589    /// Increase the count of maintained collections.
590    fn inc_collection_count(&self, collection_type: &str, hydrated: bool) {
591        let hydrated = if hydrated { "1" } else { "0" };
592        self.metrics
593            .collection_count
594            .with_label_values(&[self.worker_label.as_ref(), collection_type, hydrated])
595            .inc();
596    }
597
598    /// Decrease the count of maintained collections.
599    fn dec_collection_count(&self, collection_type: &str, hydrated: bool) {
600        let hydrated = if hydrated { "1" } else { "0" };
601        self.metrics
602            .collection_count
603            .with_label_values(&[self.worker_label.as_ref(), collection_type, hydrated])
604            .dec();
605    }
606
607    pub fn inc_subscribe_snapshot_optimization(&self) {
608        self.metrics.subscribe_snapshots_skipped_total.inc()
609    }
610
611    /// Increment the count of metric sink collector registrations that collided and will be
612    /// retried.
613    ///
614    /// Unlabeled on purpose: this is an "is registration contending" signal, and the sink's
615    /// identity comes from the log line the sink emits on its first failure.
616    pub fn inc_metric_sink_registration_retries(&self) {
617        self.metrics.metric_sink_registration_retries_total.inc()
618    }
619
620    /// Sets the workload class for the compute metrics.
621    pub fn set_workload_class(&self, workload_class: Option<String>) {
622        self.metrics.set_workload_class(workload_class);
623    }
624
625    pub fn for_collection(&self, id: GlobalId) -> CollectionMetrics {
626        CollectionMetrics::new(id, self.clone())
627    }
628}
629
630/// Collection metrics.
631///
632/// Note that these metrics do _not_ have a `collection_id` label. We avoid introducing
633/// per-collection, per-worker metrics because the number of resulting time series would
634/// potentially be huge. Instead we count classes of collections, such as hydrated collections.
635#[derive(Clone, Debug)]
636pub struct CollectionMetrics {
637    metrics: WorkerMetrics,
638    collection_type: &'static str,
639    collection_hydrated: bool,
640}
641
642impl CollectionMetrics {
643    pub fn new(collection_id: GlobalId, metrics: WorkerMetrics) -> Self {
644        let collection_type = match collection_id {
645            GlobalId::System(_) => "system",
646            GlobalId::IntrospectionSourceIndex(_) => "log",
647            GlobalId::User(_) => "user",
648            GlobalId::Transient(_) => "transient",
649            GlobalId::Explain => "explain",
650        };
651        let collection_hydrated = false;
652
653        metrics.inc_collection_count(collection_type, collection_hydrated);
654
655        Self {
656            metrics,
657            collection_type,
658            collection_hydrated,
659        }
660    }
661
662    /// Record this collection as hydration.
663    pub fn record_collection_hydrated(&mut self) {
664        if self.collection_hydrated {
665            return;
666        }
667
668        self.metrics
669            .dec_collection_count(self.collection_type, false);
670        self.metrics
671            .inc_collection_count(self.collection_type, true);
672        self.collection_hydrated = true;
673    }
674}
675
676impl Drop for CollectionMetrics {
677    fn drop(&mut self) {
678        self.metrics
679            .dec_collection_count(self.collection_type, self.collection_hydrated);
680    }
681}
682
683#[cfg(test)]
684mod tests {
685    use std::collections::BTreeSet;
686
687    use mz_ore::metrics::MetricsRegistry;
688
689    use super::ComputeMetrics;
690    use crate::server::ComputeRuntimeRole;
691
692    /// The `Solo` (single-runtime) role registers exactly as compute did before a second runtime
693    /// existed: no metric carries a `role` label, so single-runtime dashboards and exact-match
694    /// alerts are byte-unchanged.
695    #[mz_ore::test]
696    fn solo_runtime_omits_role_label() {
697        let registry = MetricsRegistry::new();
698        let metrics = ComputeMetrics::register_with(&registry, ComputeRuntimeRole::Solo);
699        // Instantiate the per-worker children so the `*Vec` families emit rows to inspect.
700        let _worker = metrics.for_worker(0);
701
702        for family in registry.gather() {
703            for metric in family.get_metric() {
704                for label in metric.get_label() {
705                    assert_ne!(
706                        label.name(),
707                        "role",
708                        "solo metric {} unexpectedly carries a role label",
709                        family.name(),
710                    );
711                }
712            }
713        }
714    }
715
716    /// The two named roles each carry their own `role` label, so two runtimes in one process
717    /// register distinct series rather than colliding. Registering both on one registry also
718    /// exercises the non-collision that lets them coexist.
719    #[mz_ore::test]
720    fn named_roles_carry_distinct_role_label() {
721        let registry = MetricsRegistry::new();
722        let maintenance = ComputeMetrics::register_with(&registry, ComputeRuntimeRole::Maintenance);
723        let interactive = ComputeMetrics::register_with(&registry, ComputeRuntimeRole::Interactive);
724        let _maintenance_worker = maintenance.for_worker(0);
725        let _interactive_worker = interactive.for_worker(0);
726
727        let mut roles = BTreeSet::new();
728        for family in registry.gather() {
729            for metric in family.get_metric() {
730                let role = metric
731                    .get_label()
732                    .iter()
733                    .find(|label| label.name() == "role")
734                    .unwrap_or_else(|| panic!("metric {} missing a role label", family.name()));
735                roles.insert(role.value().to_string());
736            }
737        }
738
739        assert!(roles.contains("maintenance"), "roles seen: {roles:?}");
740        assert!(roles.contains("interactive"), "roles seen: {roles:?}");
741    }
742}