Skip to main content

mz_compute/sink/
metric_sink.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//! Render arm for `MetricSinkConnection`.
11//!
12//! A metric sink funnels every row of its source collection to one worker, folds
13//! it into a [`SinkState`], and exposes that state to the process's Prometheus registry through
14//! a [`SinkCollector`]. `SinkState` is shared between the timely operator (the sole writer) and
15//! `SinkCollector::collect` (the reader, invoked from whatever thread scrapes the registry) via
16//! `Arc<Mutex<_>>`. Both sides only ever hold the lock across a short, synchronous section: the
17//! operator has no `await` points (it is a synchronous `builder_rc` operator), and the collector
18//! only clones out the data it needs to build `MetricFamily` protos before releasing the lock.
19//!
20//! The planner (`optimize::metric_sink::shape_metric_sink_source`) does the row-wise shaping:
21//! it coalesces `labels`/`help` to their identity element and computes the `metric_kind`
22//! and `name_valid` columns `extract_row` reads below, so this module no longer parses
23//! `metric_type` strings or validates `metric_name` itself. Dedup, collision detection, and
24//! family-conflict counting stay here because they need the cross-row state of the fold.
25
26use std::any::Any;
27use std::cell::RefCell;
28use std::collections::BTreeMap;
29use std::rc::Rc;
30use std::sync::{Arc, Mutex};
31use std::time::{Duration, Instant};
32
33use differential_dataflow::{Hashable, VecCollection};
34use mz_compute_types::sinks::{ComputeSinkDesc, MetricSinkConnection};
35use mz_ore::cast::{CastFrom, CastLossy};
36use mz_ore::metrics::MetricsRegistry;
37use mz_repr::{ColumnName, Datum, DatumVec, Diff, GlobalId, RelationDesc, Row, Timestamp};
38use mz_storage_types::controller::CollectionMetadata;
39use mz_timely_util::probe::{Handle, ProbeNotify};
40use prometheus::core::{Collector, Desc};
41use prometheus::proto::{
42    Counter as ProtoCounter, Gauge as ProtoGauge, LabelPair, Metric as ProtoMetric, MetricFamily,
43    MetricType,
44};
45use prometheus::{Gauge, Opts};
46use timely::dataflow::channels::pact::Exchange;
47use timely::dataflow::operators::generic::builder_rc::OperatorBuilder;
48use timely::progress::Antichain;
49
50use crate::metrics::WorkerMetrics;
51use crate::render::StartSignal;
52use crate::render::errors::DataflowErrorSer;
53use crate::render::sinks::SinkRender;
54
55impl<'scope> SinkRender<'scope> for MetricSinkConnection {
56    fn render_sink(
57        &self,
58        compute_state: &mut crate::compute_state::ComputeState,
59        sink: &ComputeSinkDesc<CollectionMetadata>,
60        sink_id: GlobalId,
61        _as_of: Antichain<Timestamp>,
62        _start_signal: StartSignal,
63        sinked_collection: VecCollection<'scope, Timestamp, Row, Diff>,
64        err_collection: VecCollection<'scope, Timestamp, DataflowErrorSer, Diff>,
65        output_probe: &Handle<Timestamp>,
66    ) -> Option<Rc<dyn Any>> {
67        let cols = ColumnIndices::resolve(&sink.from_desc);
68
69        let scope = sinked_collection.scope();
70        let worker_id = scope.index();
71        // The registry is process-local, so every row must land on the same worker or the
72        // series would be split across processes. Which worker is chosen doesn't matter, only
73        // that all workers agree, so hash the sink's own id.
74        //
75        // Routing by metric key instead would spread the fold across workers, but each process
76        // has its own registry, so one metric family could then be split across processes' scrape
77        // outputs. Partition-by-key is a possible future refinement.
78        let active_worker_id = usize::cast_from(sink_id.hashed()) % scope.peers();
79
80        let ok_stream = sinked_collection
81            .inner
82            .probe_notify_with(vec![output_probe.clone()]);
83        let err_stream = err_collection.inner;
84
85        let state = Arc::new(Mutex::new(SinkState::default()));
86
87        // Only the active worker registers a collector; the `MetricsRegistry` is process-wide.
88        //
89        // NOTE: a `Desc` id collision here is expected, not a logic error, and only for curated
90        // sinks. A curated sink keys its gauges on its stable name, so every incarnation shares one
91        // `Desc` id (each still gets its own transient, boot-reused `GlobalId`). The registering
92        // worker is `sink_id.hashed() % peers`, so two incarnations usually land on different
93        // workers, and `Worker::reconcile` has no cross-worker barrier: on a restart the new one can
94        // register before the other worker drops the old. User sinks keep a durable `GlobalId`, so
95        // both incarnations hash to the same worker, which unregisters the old collector before
96        // registering the new: `Worker::reconcile` clears a non-retained dataflow's sink token
97        // (dropping this guard) before applying the replacement `CreateDataflow`, so they never
98        // collide. Registration is fallible, retried until the predecessor drops.
99        let registration = Rc::new(RefCell::new((worker_id == active_worker_id).then(|| {
100            PendingRegistration::new(
101                self.label.clone(),
102                SinkCollector::new(&self.label, Arc::clone(&state)),
103            )
104        })));
105        let registration_op = Rc::clone(&registration);
106        let registry = compute_state.metrics_registry.clone();
107        let worker_metrics = compute_state.metrics.clone();
108
109        let mut op = OperatorBuilder::new(format!("MetricSink({sink_id})"), scope.clone());
110        let mut ok_input = op.new_input(
111            ok_stream,
112            Exchange::new(move |_: &(Row, Timestamp, Diff)| u64::cast_from(active_worker_id)),
113        );
114        let mut err_input = op.new_input(
115            err_stream,
116            Exchange::new(move |_: &(DataflowErrorSer, Timestamp, Diff)| {
117                u64::cast_from(active_worker_id)
118            }),
119        );
120
121        // The frontier this worker reports to the controller. Every worker reports, not just the
122        // active one: the controller meets the per-worker frontiers, so a worker that never
123        // advanced would pin the input's since forever.
124        let sink_frontier = Rc::new(RefCell::new(Antichain::from_elem(Timestamp::MIN)));
125        let shared_frontier = Rc::clone(&sink_frontier);
126
127        let operator_info = op.operator_info();
128        op.build(move |_capabilities| {
129            let activator = scope.activator_for(operator_info.address);
130
131            // Register now, at build time, not on first activation: activation waits on the input
132            // frontier, so a sink with a slow input would publish no series until it happened to be
133            // scheduled, and a scrape in that gap sees fewer sinks than exist. A collision retries
134            // through the activation path below.
135            //
136            // NOTE: these retries ride on operator activations, so they stop if both inputs close
137            // before a collision clears (e.g. a constant-folded source): timely drops the operator
138            // closure, leaving the token holding a handle-less `PendingRegistration`, so the sink
139            // never publishes. A collision outlasting a finite source isn't expected for these
140            // curated sinks, so we accept the gap rather than pin an activation open until the
141            // handle is set.
142            if let Some(registration) = registration_op.borrow_mut().as_mut() {
143                if let Retry::Arm = registration.try_register(&registry, &worker_metrics) {
144                    activator.activate_after(REGISTRATION_RETRY_INTERVAL);
145                }
146            }
147
148            // Recycled across activations: unpacking a row into `Datum`s otherwise allocates a
149            // fresh `Vec` per row on this hot path.
150            let mut datum_vec = DatumVec::new();
151            move |frontiers| {
152                // Combined ok+err input frontier. A timestamp is closed once neither input can
153                // still produce data at it, so folding a closed time observes all of its diffs.
154                let mut frontier = Antichain::new();
155                for f in frontiers {
156                    frontier.extend(f.frontier().iter().copied());
157                }
158                // NOTE: the reported frontier is published here, before the active worker folds
159                // through it below (`st.integrate` at the end of this closure). This is only safe
160                // because the active worker must complete that fold in the same activation that
161                // observed the frontier. An early return added on the active-worker path between
162                // here and `st.integrate` would advertise progress the sink has not folded,
163                // downgrading the input's since ahead of the data actually consumed.
164                shared_frontier.borrow_mut().clone_from(&frontier);
165
166                // `registration` is `Some` only on the active worker. Scoped so the borrow ends
167                // before the fold.
168                {
169                    let mut reg = registration_op.borrow_mut();
170                    let Some(registration) = reg.as_mut() else {
171                        // Non-active worker: all input is exchanged to the active worker, so nothing
172                        // arrives here. Drain to avoid being rescheduled forever.
173                        ok_input.for_each(|_, _| {});
174                        err_input.for_each(|_, _| {});
175                        return;
176                    };
177
178                    if let Retry::Arm = registration.try_register(&registry, &worker_metrics) {
179                        // The input may be quiescent, so nothing else would reschedule us.
180                        activator.activate_after(REGISTRATION_RETRY_INTERVAL);
181                    }
182                }
183
184                let mut st = state.lock().expect("sink state mutex poisoned");
185
186                // Buffer every incoming update under its timestamp. The input is not
187                // consolidated and timely does not guarantee that all diffs at a timestamp arrive
188                // in one activation, so nothing is folded into `working` until the timestamp is
189                // closed.
190                ok_input.for_each(|_, data| {
191                    for (row, time, diff) in data.drain(..) {
192                        let datums = datum_vec.borrow_with(&row);
193                        let (name, metric_kind, name_valid, labels, value, help) =
194                            extract_row(&cols, &datums);
195                        st.stage_ok(
196                            name,
197                            metric_kind,
198                            name_valid,
199                            &labels,
200                            value,
201                            help,
202                            time,
203                            diff.into_inner(),
204                        );
205                    }
206                });
207                err_input.for_each(|_, data| {
208                    for (_err, time, diff) in data.drain(..) {
209                        st.stage_err(time, diff.into_inner());
210                    }
211                });
212
213                st.integrate(&frontier);
214                st.frontier_ms = frontier
215                    .as_option()
216                    .map(|t| u64::from(*t))
217                    .unwrap_or(u64::MAX);
218                st.publish_if_healthy();
219            }
220        });
221
222        // Report frontier updates to the `ComputeState`. A metric sink writes to the metrics
223        // registry rather than to a collection, so its "write" frontier is the input frontier it
224        // has folded through.
225        let collection = compute_state.expect_collection_mut(sink_id);
226        collection.sink_write_frontier = Some(sink_frontier);
227
228        // The guard is returned as the sink token, not kept in the operator closure. This operator
229        // has no outputs, so timely drops its closure the moment both input frontiers close (e.g. a
230        // constant-folded source), and a guard held there would unregister a still-installed sink.
231        // On the token it drops with the dataflow instead, freeing the `Desc` id for a retrying
232        // successor and letting a torn-down sink keep serving its now-frozen series until the drain.
233        //
234        // A non-active worker's guard is empty and harmless.
235        let token: Rc<dyn Any> = registration;
236        Some(token)
237    }
238}
239
240/// How long to wait before retrying a collector registration that collided.
241///
242/// Fixed, no backoff: a cross-worker `AlreadyReg` collision clears once the predecessor's dataflow
243/// is torn down, so it converges without a growing wait.
244const REGISTRATION_RETRY_INTERVAL: Duration = Duration::from_secs(1);
245
246/// How long a collision may persist before it soft-panics.
247///
248/// The retry expects a predecessor to drop within a few seconds. A collision outlasting this one
249/// isn't dropping, so soft-panic once to surface it, then keep retrying in case it clears.
250const REGISTRATION_ESCALATE_AFTER: Duration = Duration::from_secs(60);
251
252/// Whether a registration attempt left a collision the caller must arm a retry for.
253#[derive(Debug, Clone, Copy, PartialEq, Eq)]
254enum Retry {
255    /// Arm a retry activation [`REGISTRATION_RETRY_INTERVAL`] from now.
256    Arm,
257    /// Nothing to arm: registered, permanently abandoned, or a retry is already armed.
258    Skip,
259}
260
261/// Active-worker collector registration that tolerates a transient descriptor-id collision with an
262/// incarnation of this sink that has not been torn down yet.
263///
264/// Holds the guard once registered, so dropping this unregisters the collector.
265struct PendingRegistration {
266    /// The sink's label, for the log line on first failure.
267    label: String,
268    collector: SinkCollector,
269    handle: Option<Box<dyn Any + Send + Sync>>,
270    logged: bool,
271    /// Set once a non-collision error soft-panicked, so we stop retrying.
272    terminated: bool,
273    /// When the collision streak began, for the escalation bound.
274    first_collision: Option<Instant>,
275    /// Set once the escalation soft-panic has fired, so it fires only once.
276    escalated: bool,
277    /// Earliest instant a new attempt may run. `Some` means a retry is already armed, so
278    /// activations before it are no-ops. This paces attempts to one per
279    /// [`REGISTRATION_RETRY_INTERVAL`] regardless of how often the input reschedules the operator.
280    next_attempt: Option<Instant>,
281}
282
283impl PendingRegistration {
284    fn new(label: String, collector: SinkCollector) -> Self {
285        PendingRegistration {
286            label,
287            collector,
288            handle: None,
289            logged: false,
290            terminated: false,
291            first_collision: None,
292            escalated: false,
293            next_attempt: None,
294        }
295    }
296
297    /// Attempts to register the collector, at most once per [`REGISTRATION_RETRY_INTERVAL`].
298    ///
299    /// Returns [`Retry::Arm`] only when a fresh collision was just seen and the caller must arm a
300    /// retry activation; [`Retry::Skip`] once registered, permanently abandoned, or while a retry
301    /// armed by an earlier collision is still pending. The gate matters because a busy input
302    /// reschedules the operator far more often than the retry interval, and each bare attempt
303    /// clones the collector and takes the registry write lock.
304    ///
305    /// An `AlreadyReg` collision retries until the predecessor drops the `Desc` id; after
306    /// [`REGISTRATION_ESCALATE_AFTER`] it soft-panics once to surface a predecessor that never
307    /// drops, then keeps retrying. Any other error is a logic error, so it soft-panics and stops.
308    fn try_register(&mut self, registry: &MetricsRegistry, metrics: &WorkerMetrics) -> Retry {
309        if self.handle.is_some() || self.terminated {
310            return Retry::Skip;
311        }
312        let now = Instant::now();
313        if let Some(next) = self.next_attempt {
314            if now < next {
315                return Retry::Skip;
316            }
317        }
318        match registry.try_register_collector_with_dropper(self.collector.clone()) {
319            Ok(handle) => {
320                self.handle = Some(handle);
321                self.next_attempt = None;
322                Retry::Skip
323            }
324            Err(prometheus::Error::AlreadyReg) => {
325                metrics.inc_metric_sink_registration_retries();
326                let first = *self.first_collision.get_or_insert(now);
327                // Only the first failure logs; the counter carries the ongoing signal.
328                if !self.logged {
329                    self.logged = true;
330                    tracing::info!(
331                        sink = %self.label,
332                        "metric sink collector registration collided, retrying"
333                    );
334                }
335                if !self.escalated && now.duration_since(first) >= REGISTRATION_ESCALATE_AFTER {
336                    self.escalated = true;
337                    mz_ore::soft_panic_or_log!(
338                        "metric sink {} collector registration still colliding after {:?}; a \
339                         predecessor incarnation has not dropped its registration",
340                        self.label,
341                        REGISTRATION_ESCALATE_AFTER
342                    );
343                }
344                self.next_attempt = Some(now + REGISTRATION_RETRY_INTERVAL);
345                Retry::Arm
346            }
347            Err(err) => {
348                self.terminated = true;
349                mz_ore::soft_panic_or_log!(
350                    "metric sink {} collector registration failed: {err}",
351                    self.label
352                );
353                Retry::Skip
354            }
355        }
356    }
357}
358
359/// Column indices resolved once from the sink's source relation.
360///
361/// The source relation exposes `metric_name`, `labels`, `value`, and `help` of the required types,
362/// and `shape_metric_sink_source` adds the `metric_kind` and `name_valid` columns this reads.
363/// `resolve` panics if a column is missing. No tree caller enforces this column contract yet; the
364/// SQL planner will. `metric_name`, `value`, and the values of the `labels` map may still be
365/// `Datum::Null` (a `map[text=>text]` has no per-value nullability); the columns themselves are
366/// non-null by construction. Column position within the row is unconstrained.
367struct ColumnIndices {
368    metric_name: usize,
369    labels: usize,
370    value: usize,
371    help: usize,
372    metric_kind: usize,
373    name_valid: usize,
374}
375
376impl ColumnIndices {
377    fn resolve(desc: &RelationDesc) -> Self {
378        let idx = |name: &str| {
379            desc.get_by_name(&ColumnName::from(name))
380                .expect("column existence validated by the SQL planner")
381                .0
382        };
383        ColumnIndices {
384            metric_name: idx("metric_name"),
385            labels: idx("labels"),
386            value: idx("value"),
387            help: idx("help"),
388            metric_kind: idx("metric_kind"),
389            name_valid: idx("name_valid"),
390        }
391    }
392}
393
394/// Extracts `(metric_name, metric_kind, name_valid, sorted labels, value, help)` from one shaped
395/// source row.
396///
397/// The planner already did the row-wise shaping.
398///
399/// A null label value stays `None` rather than being coerced to a string. Neither a null nor a
400/// genuine empty-string value is a representable Prometheus label (an empty value reads as absent),
401/// so a row carrying either is skipped downstream instead of published.
402///
403/// Strings borrow from `datums`, so the caller must own what it needs
404/// (see `SinkState::stage_ok`) before the row backing `datums` is dropped.
405fn extract_row<'a>(
406    cols: &ColumnIndices,
407    datums: &[Datum<'a>],
408) -> (
409    &'a str,
410    Option<MetricKind>,
411    bool,
412    Vec<(&'a str, Option<&'a str>)>,
413    Option<f64>,
414    &'a str,
415) {
416    let metric_name = match datums[cols.metric_name] {
417        Datum::Null => "",
418        d => d.unwrap_str(),
419    };
420    let metric_kind = MetricKind::from_datum(datums[cols.metric_kind]);
421    let name_valid = matches!(datums[cols.name_valid], Datum::True);
422    let mut labels: Vec<(&str, Option<&str>)> = datums[cols.labels]
423        .unwrap_map()
424        .iter()
425        .map(|(k, v)| (k, (!v.is_null()).then(|| v.unwrap_str())))
426        .collect();
427    labels.sort();
428    let value = match datums[cols.value] {
429        Datum::Null => None,
430        d => Some(d.unwrap_float64()),
431    };
432    let help = datums[cols.help].unwrap_str();
433    (metric_name, metric_kind, name_valid, labels, value, help)
434}
435
436/// Full identity of one source row: metric name, sorted labels, value, metric kind, name
437/// validity, and help.
438///
439/// A null value is its own distinct row identity, not a stand-in for any particular number,
440/// so it is kept apart from every `Some(_)` identity rather than coerced to
441/// one. A null label value is likewise distinct from a `''` value, so `{a => NULL}` and
442/// `{a => ''}` retract only against their own inserts even though both are unpublishable.
443/// The name and labels lead the tuple so that a `BTreeMap<RowKey, _>` keeps all rows of one
444/// `(metric_name, labels)` series adjacent.
445///
446/// `metric_kind` is the planner's classification (`None` for any unsupported `metric_type`), not
447/// the raw string. Two source rows that differ only in which unsupported type they carry (e.g.
448/// `"histogram"` vs. `"summary"`) now share one identity instead of two. That only affects the
449/// granularity of the `skipped` count for rows that are never published either way.
450type RowKey = (
451    String,
452    Vec<(String, Option<String>)>,
453    Option<u64>,
454    Option<MetricKind>,
455    bool,
456    String,
457);
458
459/// Key into [`SinkState::published`]: a metric name paired with its sorted label vector.
460type PublishedKey = (String, Vec<(String, String)>);
461/// Value in [`SinkState::published`]: the series' value, kind, and help string.
462type PublishedValue = (f64, MetricKind, String);
463
464/// Working and published metric state for one metric sink.
465///
466/// Because the input is not consolidated and a timestamp's diffs may span several operator
467/// activations, incoming updates are buffered by timestamp in `pending_ok`/`pending_err` and only
468/// folded into `working` once the input frontier has closed that timestamp. `working` accumulates
469/// a signed multiplicity per full row identity, and a row is live iff its accumulated diff is
470/// positive.
471///
472/// `published` is what the collector exposes and is rebuilt from the live set of
473/// `working`.
474#[derive(Default)]
475struct SinkState {
476    /// Ok-collection updates awaiting their timestamp closing, accumulated per identity.
477    pending_ok: BTreeMap<Timestamp, BTreeMap<RowKey, i64>>,
478    /// Err-collection diffs awaiting their timestamp closing.
479    pending_err: BTreeMap<Timestamp, i64>,
480    /// Accumulated multiplicity per row identity over all closed timestamps.
481    working: BTreeMap<RowKey, i64>,
482    published: BTreeMap<PublishedKey, PublishedValue>,
483    /// Net count of live errors on the sink's input. Can rise and fall as errors are
484    /// retracted. Publication is frozen while this is nonzero.
485    errors: i64,
486    frontier_ms: u64,
487    skipped: u64,
488    conflicts: u64,
489    collisions: u64,
490    /// Count of live `(metric_name, labels)` groups whose only live rows carry a null `value`,
491    /// so the series is currently absent from `published` (a gap) rather than published as some
492    /// number.
493    null_values: u64,
494}
495
496#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
497enum MetricKind {
498    Gauge,
499    Counter,
500}
501
502impl MetricKind {
503    /// Recovers the classification the planner's `metric_kind` column already computed (`0` =
504    /// gauge, `1` = counter).
505    fn from_datum(d: Datum) -> Option<Self> {
506        match d {
507            Datum::Int32(0) => Some(MetricKind::Gauge),
508            Datum::Int32(1) => Some(MetricKind::Counter),
509            _ => None,
510        }
511    }
512
513    fn proto_type(self) -> MetricType {
514        match self {
515            MetricKind::Gauge => MetricType::GAUGE,
516            MetricKind::Counter => MetricType::COUNTER,
517        }
518    }
519}
520
521/// Matches Prometheus's label name grammar: `[a-zA-Z_][a-zA-Z0-9_]*`.
522///
523/// Unlike the metric-name grammar (`name_valid`, computed in MIR), this stays in Rust: it
524/// applies to every key of the `labels` map, an unbounded per-row collection that doesn't fit a
525/// scalar `Map` expression.
526fn is_valid_label_name(name: &str) -> bool {
527    let mut chars = name.chars();
528    match chars.next() {
529        Some(c) if c.is_ascii_alphabetic() || c == '_' => {
530            chars.all(|c| c.is_ascii_alphanumeric() || c == '_')
531        }
532        _ => false,
533    }
534}
535
536/// Whether a label value can be published as-is. A null (`None`) has no value to encode, and an
537/// empty string reads as absent to Prometheus, so `{a => ''}` would fold into `{}`. Both make the
538/// row that carries them unpublishable.
539fn is_publishable_label_value(value: Option<&str>) -> bool {
540    matches!(value, Some(v) if !v.is_empty())
541}
542
543impl SinkState {
544    /// Buffers one ok-collection update under its timestamp.
545    ///
546    /// `diff` follows differential dataflow sign conventions and is accumulated per full row
547    /// identity, so an un-consolidated input carrying the same identity twice sums rather than
548    /// being mistaken for a second live row. Takes borrowed strings (see `extract_row`) and owns
549    /// them only here, once, at the point the identity is committed to `pending_ok`.
550    fn stage_ok(
551        &mut self,
552        metric_name: &str,
553        metric_kind: Option<MetricKind>,
554        name_valid: bool,
555        labels: &[(&str, Option<&str>)],
556        value: Option<f64>,
557        help: &str,
558        time: Timestamp,
559        diff: i64,
560    ) {
561        let key = (
562            metric_name.to_string(),
563            labels
564                .iter()
565                .map(|&(k, v)| (k.to_string(), v.map(str::to_string)))
566                .collect(),
567            value.map(f64::to_bits),
568            metric_kind,
569            name_valid,
570            help.to_string(),
571        );
572        *self
573            .pending_ok
574            .entry(time)
575            .or_default()
576            .entry(key)
577            .or_default() += diff;
578    }
579
580    /// Buffers one err-collection diff under its timestamp.
581    fn stage_err(&mut self, time: Timestamp, diff: i64) {
582        *self.pending_err.entry(time).or_default() += diff;
583    }
584
585    /// Folds every buffered timestamp the `frontier` has closed into `working` and `errors`.
586    ///
587    /// A timestamp is closed once the combined ok+err frontier can no longer produce data at it,
588    /// which guarantees all of its diffs are already buffered. Accumulated entries that reach a
589    /// multiplicity of zero are dropped. `skipped` is recomputed here from the live set, so it
590    /// counts the input rows currently dropped for an unsupported type or invalid name.
591    fn integrate(&mut self, frontier: &Antichain<Timestamp>) {
592        let closed_ok: Vec<Timestamp> = self
593            .pending_ok
594            .keys()
595            .filter(|t| !frontier.less_equal(t))
596            .copied()
597            .collect();
598        for time in closed_ok {
599            let rows = self.pending_ok.remove(&time).expect("key from keys()");
600            for (key, diff) in rows {
601                *self.working.entry(key).or_default() += diff;
602            }
603        }
604
605        let closed_err: Vec<Timestamp> = self
606            .pending_err
607            .keys()
608            .filter(|t| !frontier.less_equal(t))
609            .copied()
610            .collect();
611        for time in closed_err {
612            self.errors += self.pending_err.remove(&time).expect("key from keys()");
613        }
614
615        self.working.retain(|_, acc| *acc != 0);
616        self.skipped = count_skipped(&self.working);
617    }
618
619    /// Rebuilds `published` and the `collisions`/`conflicts` counts from the live set of
620    /// `working`, but only while the dataflow is free of live errors. While `errors > 0`,
621    /// publication stays frozen at the last healthy snapshot. `working` keeps integrating closed
622    /// timestamps in the meantime, so the next healthy publish reflects everything that happened
623    /// during the freeze.
624    ///
625    // NOTE: this is a full O(n) rebuild over the entire live set on every healthy activation, not an
626    // incremental update. Consider revisiting with incremental maintenance if a sink's series
627    // count grows large enough for the per-activation scan to matter.
628    fn publish_if_healthy(&mut self) {
629        if self.errors == 0 {
630            let (published, collisions, null_values) = rebuild_published(&self.working);
631            self.published = published;
632            self.collisions = collisions;
633            self.null_values = null_values;
634            self.conflicts = count_conflicts(&self.published);
635        }
636    }
637}
638
639/// Counts live working rows dropped for an unsupported `metric_type`, an invalid Prometheus
640/// metric or label name, or a null or empty label value.
641fn count_skipped(working: &BTreeMap<RowKey, i64>) -> u64 {
642    let mut skipped = 0u64;
643    for ((_name, labels, _bits, metric_kind, name_valid, _help), acc) in working {
644        if *acc <= 0 {
645            continue;
646        }
647        let unsupported = metric_kind.is_none();
648        let invalid = !name_valid
649            || !labels
650                .iter()
651                .all(|(k, v)| is_publishable_label_value(v.as_deref()) && is_valid_label_name(k));
652        if unsupported || invalid {
653            skipped += 1;
654        }
655    }
656    skipped
657}
658
659/// Collapses the live, representable rows of `working` into one published entry per
660/// `(metric_name, labels)` series and counts colliding and null-suppressed series.
661///
662/// A row whose label set is not representable (an invalid label name, or a null or empty label
663/// value) is dropped here and counted by [`count_skipped`] instead.
664///
665/// A series collides when more than one distinct live non-null value exists for its
666/// `(metric_name, labels)`: a genuine conflict of two source rows, unlike an ordinary
667/// value update whose old row is retracted and new row inserted within the same closed
668/// timestamp, which leaves a single live value. When a series collides, the winner is chosen
669/// deterministically as the row with the numerically smallest value, breaking ties by metric type
670/// then help. A null `value` carries no number to compare or publish: a series whose live rows
671/// are all null-valued is absent from `published` (a gap) and counted in the returned
672/// `null_values` instead of `collisions`. A series with at least one live non-null value publishes
673/// normally and is not counted in `null_values`, even if null-valued rows are also live for it.
674fn rebuild_published(
675    working: &BTreeMap<RowKey, i64>,
676) -> (BTreeMap<PublishedKey, PublishedValue>, u64, u64) {
677    // `working` orders rows by `(name, labels, ...)`, so all rows of one series are adjacent.
678    let mut grouped: BTreeMap<PublishedKey, Vec<(Option<f64>, MetricKind, String)>> =
679        BTreeMap::new();
680    for ((name, labels, bits, metric_kind, name_valid, help), acc) in working {
681        if *acc <= 0 {
682            continue;
683        }
684        let Some(kind) = metric_kind else {
685            continue;
686        };
687        // A null or empty label value has no representable Prometheus label (an empty value reads
688        // as absent, folding `{a => ''}` into `{}`), so the whole row is unpublishable. Drop it
689        // rather than emit a series the source never had.
690        let Some(labels) = labels
691            .iter()
692            .map(|(k, v)| {
693                is_publishable_label_value(v.as_deref()).then(|| {
694                    (
695                        k.clone(),
696                        v.as_ref().expect("publishable is non-null").clone(),
697                    )
698                })
699            })
700            .collect::<Option<Vec<_>>>()
701        else {
702            continue;
703        };
704        if !name_valid || !labels.iter().all(|(k, _)| is_valid_label_name(k)) {
705            continue;
706        }
707        grouped.entry((name.clone(), labels)).or_default().push((
708            bits.map(f64::from_bits),
709            *kind,
710            help.clone(),
711        ));
712    }
713
714    let mut published = BTreeMap::new();
715    let mut collisions = 0u64;
716    let mut null_values = 0u64;
717    for (key, candidates) in grouped {
718        let mut non_null: Vec<PublishedValue> = candidates
719            .into_iter()
720            .filter_map(|(value, kind, help)| value.map(|v| (v, kind, help)))
721            .collect();
722        if non_null.is_empty() {
723            null_values += 1;
724            continue;
725        }
726        let mut distinct: Vec<u64> = non_null.iter().map(|(v, _, _)| v.to_bits()).collect();
727        distinct.sort_unstable();
728        distinct.dedup();
729        if distinct.len() > 1 {
730            collisions += 1;
731        }
732        non_null.sort_by(|a, b| {
733            a.0.total_cmp(&b.0)
734                .then(a.1.cmp(&b.1))
735                .then_with(|| a.2.cmp(&b.2))
736        });
737        let winner = non_null
738            .into_iter()
739            .next()
740            .expect("checked non-empty above");
741        published.insert(key, winner);
742    }
743    (published, collisions, null_values)
744}
745
746/// Counts published series whose own `metric_type`/`help` disagree with their family's winning
747/// type/help. See [`build_families`] for how the winner is chosen.
748fn count_conflicts(published: &BTreeMap<PublishedKey, PublishedValue>) -> u64 {
749    let mut conflicts = 0u64;
750    let mut winner: Option<(&str, MetricKind, &str)> = None;
751    for ((name, _labels), (_value, kind, help)) in published {
752        winner = match winner {
753            Some((n, k, h)) if n == name.as_str() => {
754                if k != *kind || h != help.as_str() {
755                    conflicts += 1;
756                }
757                Some((n, k, h))
758            }
759            _ => Some((name.as_str(), *kind, help.as_str())),
760        };
761    }
762    conflicts
763}
764
765/// Groups `published` by metric name into one `MetricFamily` per name, since Prometheus requires
766/// a single type and help string per family. Within a group, the entry with the
767/// lexicographically smallest label vector wins the family's type and help string;
768/// `BTreeMap`'s `(name, labels)` key ordering already sorts each group that way, so the
769/// first entry seen for a given name is that winner.
770///
771/// `MetricFamily.type_` and `Metric.{gauge,counter}` are protobuf wrapper types
772/// (`EnumOrUnknown`/`MessageField`).
773fn build_families(published: &BTreeMap<PublishedKey, PublishedValue>) -> Vec<MetricFamily> {
774    let mut families = Vec::new();
775    let mut group_name: Option<&str> = None;
776    let mut family: Option<MetricFamily> = None;
777    let mut family_kind = MetricKind::Gauge;
778
779    for ((name, labels), (value, kind, help)) in published {
780        if group_name != Some(name.as_str()) {
781            if let Some(f) = family.take() {
782                families.push(f);
783            }
784            let mut mf = MetricFamily::new();
785            mf.name = Some(name.clone());
786            mf.help = Some(help.clone());
787            mf.type_ = Some(kind.proto_type().into());
788            family = Some(mf);
789            family_kind = *kind;
790            group_name = Some(name.as_str());
791        }
792
793        let mut metric = ProtoMetric::new();
794        metric.label = labels
795            .iter()
796            .map(|(k, v)| {
797                let mut lp = LabelPair::new();
798                lp.name = Some(k.clone());
799                lp.value = Some(v.clone());
800                lp
801            })
802            .collect();
803        match family_kind {
804            MetricKind::Gauge => {
805                let mut g = ProtoGauge::new();
806                g.value = Some(*value);
807                metric.gauge = Some(g).into();
808            }
809            MetricKind::Counter => {
810                let mut c = ProtoCounter::new();
811                c.value = Some(*value);
812                metric.counter = Some(c).into();
813            }
814        }
815        family
816            .as_mut()
817            .expect("initialized above for the first entry of every group")
818            .metric
819            .push(metric);
820    }
821    if let Some(f) = family.take() {
822        families.push(f);
823    }
824    families
825}
826
827/// A `prometheus::core::Collector` that exposes a metric sink's [`SinkState`].
828///
829/// The companion gauges (`mz_compute_metric_sink_*`) are declared statically, each carrying a `sink`
830/// const label so that per-sink series get distinct `Desc` ids on registration. The user-defined
831/// series are entirely dynamic: their names come from the sink's source query, so they are built
832/// directly as [`MetricFamily`] protos in `collect` and are not declared via `desc`. Prometheus's
833/// registry only uses `desc` for registration-time collision detection, not to validate the
834/// output of `collect`, so this is safe.
835#[derive(Clone)]
836struct SinkCollector {
837    state: Arc<Mutex<SinkState>>,
838    frontier_gauge: Gauge,
839    errors_gauge: Gauge,
840    skipped_gauge: Gauge,
841    conflicts_gauge: Gauge,
842    collisions_gauge: Gauge,
843    null_values_gauge: Gauge,
844}
845
846impl SinkCollector {
847    fn new(label: &str, state: Arc<Mutex<SinkState>>) -> Self {
848        let gauge = |name: &str, help: &str| {
849            Gauge::with_opts(Opts::new(name, help).const_label("sink", label))
850                .expect("static metric sink companion gauge options are valid")
851        };
852        SinkCollector {
853            state,
854            frontier_gauge: gauge(
855                "mz_compute_metric_sink_frontier_ms",
856                "The metric sink's input frontier, in milliseconds since the epoch.",
857            ),
858            errors_gauge: gauge(
859                "mz_compute_metric_sink_errors",
860                "The number of live errors on the metric sink's input.",
861            ),
862            skipped_gauge: gauge(
863                "mz_compute_metric_sink_skipped",
864                "The number of input rows skipped for an unsupported metric type, an invalid name, or a null or empty label value.",
865            ),
866            conflicts_gauge: gauge(
867                "mz_compute_metric_sink_conflicts",
868                "The number of published series whose type or help disagree with their family's chosen type or help.",
869            ),
870            collisions_gauge: gauge(
871                "mz_compute_metric_sink_collisions",
872                "The number of series with more than one distinct live value for the same metric name and labels.",
873            ),
874            null_values_gauge: gauge(
875                "mz_compute_metric_sink_null_values",
876                "The number of series currently suppressed because their value is null.",
877            ),
878        }
879    }
880}
881
882impl Collector for SinkCollector {
883    fn desc(&self) -> Vec<&Desc> {
884        let mut descs = Vec::with_capacity(6);
885        descs.extend(self.frontier_gauge.desc());
886        descs.extend(self.errors_gauge.desc());
887        descs.extend(self.skipped_gauge.desc());
888        descs.extend(self.conflicts_gauge.desc());
889        descs.extend(self.collisions_gauge.desc());
890        descs.extend(self.null_values_gauge.desc());
891        descs
892    }
893
894    fn collect(&self) -> Vec<MetricFamily> {
895        let mut families = {
896            let state = self.state.lock().expect("sink state mutex poisoned");
897            self.frontier_gauge.set(f64::cast_lossy(state.frontier_ms));
898            self.errors_gauge.set(f64::cast_lossy(state.errors));
899            self.skipped_gauge.set(f64::cast_lossy(state.skipped));
900            self.conflicts_gauge.set(f64::cast_lossy(state.conflicts));
901            self.collisions_gauge.set(f64::cast_lossy(state.collisions));
902            self.null_values_gauge
903                .set(f64::cast_lossy(state.null_values));
904            build_families(&state.published)
905        };
906
907        families.extend(self.frontier_gauge.collect());
908        families.extend(self.errors_gauge.collect());
909        families.extend(self.skipped_gauge.collect());
910        families.extend(self.conflicts_gauge.collect());
911        families.extend(self.collisions_gauge.collect());
912        families.extend(self.null_values_gauge.collect());
913        families
914    }
915}
916
917#[cfg(test)]
918mod tests {
919    use super::*;
920    use crate::metrics::ComputeMetrics;
921    use crate::server::ComputeRuntimeRole;
922
923    /// A frontier that has closed every timestamp strictly below `bound`.
924    fn frontier(bound: u64) -> Antichain<Timestamp> {
925        Antichain::from_elem(Timestamp::from(bound))
926    }
927
928    fn label_a() -> Vec<(String, String)> {
929        vec![("a".into(), "1".into())]
930    }
931
932    /// `label_a()`, borrowed: what `stage_ok` now takes (see `extract_row`).
933    const LABEL_A: &[(&str, Option<&str>)] = &[("a", Some("1"))];
934
935    fn key_m() -> PublishedKey {
936        ("m".into(), label_a())
937    }
938
939    /// Mirrors the shaped relation `shape_metric_sink_source` builds: `labels`/`help` are
940    /// non-null by construction, `metric_name`/`value` stay nullable, and `metric_kind`/
941    /// `name_valid` are the planner's computed classification columns.
942    fn shaped_desc() -> RelationDesc {
943        use mz_repr::SqlScalarType;
944
945        RelationDesc::builder()
946            .with_column("metric_name", SqlScalarType::String.nullable(true))
947            .with_column(
948                "labels",
949                SqlScalarType::Map {
950                    value_type: Box::new(SqlScalarType::String),
951                    custom_id: None,
952                }
953                .nullable(false),
954            )
955            .with_column("value", SqlScalarType::Float64.nullable(true))
956            .with_column("help", SqlScalarType::String.nullable(false))
957            .with_column("metric_kind", SqlScalarType::Int32.nullable(true))
958            .with_column("name_valid", SqlScalarType::Bool.nullable(true))
959            .finish()
960    }
961
962    /// Stages one gauge update for the `(m, {a:1})` series.
963    fn stage_m(st: &mut SinkState, value: f64, time: u64, diff: i64) {
964        st.stage_ok(
965            "m",
966            Some(MetricKind::Gauge),
967            true,
968            LABEL_A,
969            Some(value),
970            "h",
971            Timestamp::from(time),
972            diff,
973        );
974    }
975
976    /// Stages one gauge update with a null `value` for the `(m, {a:1})` series.
977    fn stage_m_null(st: &mut SinkState, time: u64, diff: i64) {
978        st.stage_ok(
979            "m",
980            Some(MetricKind::Gauge),
981            true,
982            LABEL_A,
983            None,
984            "h",
985            Timestamp::from(time),
986            diff,
987        );
988    }
989
990    #[mz_ore::test]
991    fn fold_and_publish() {
992        let mut st = SinkState::default();
993        stage_m(&mut st, 2.0, 0, 1);
994        st.integrate(&frontier(1));
995        st.publish_if_healthy();
996        assert_eq!(st.published.len(), 1);
997        assert_eq!(st.published[&key_m()].0, 2.0);
998
999        // Unsupported type (metric_kind = None) is skipped and counted.
1000        st.stage_ok("h1", None, true, &[], Some(1.0), "h", Timestamp::from(1), 1);
1001        st.integrate(&frontier(2));
1002        assert_eq!(st.skipped, 1);
1003
1004        // Error freezes publication. The update at time 2 retracts the old value and inserts the
1005        // new one, both fold into `working` while frozen.
1006        st.errors = 1;
1007        stage_m(&mut st, 2.0, 2, -1);
1008        stage_m(&mut st, 9.0, 2, 1);
1009        st.integrate(&frontier(3));
1010        st.publish_if_healthy();
1011        assert_eq!(st.published[&key_m()].0, 2.0);
1012
1013        // Recovery republishes the integrated value.
1014        st.errors = 0;
1015        st.publish_if_healthy();
1016        assert_eq!(st.published[&key_m()].0, 9.0);
1017        assert_eq!(st.collisions, 0);
1018    }
1019
1020    #[mz_ore::test]
1021    fn value_update_split_across_activations_no_collision() {
1022        let mut st = SinkState::default();
1023        // Establish the series at value 5.
1024        stage_m(&mut st, 5.0, 0, 1);
1025        st.integrate(&frontier(1));
1026        st.publish_if_healthy();
1027        assert_eq!(st.published[&key_m()].0, 5.0);
1028
1029        // A value update 5 -> 9 at time 1 arrives insert-first, split across two activations. The
1030        // timestamp stays open until both diffs are buffered.
1031        stage_m(&mut st, 9.0, 1, 1);
1032        st.integrate(&frontier(1));
1033        st.publish_if_healthy();
1034        assert_eq!(st.collisions, 0);
1035        stage_m(&mut st, 5.0, 1, -1);
1036
1037        // Close the timestamp: the series is present at value 9 with no collision.
1038        st.integrate(&frontier(2));
1039        st.publish_if_healthy();
1040        assert_eq!(st.published[&key_m()].0, 9.0);
1041        assert_eq!(st.collisions, 0);
1042    }
1043
1044    #[mz_ore::test]
1045    fn duplicate_multiplicity_consolidates() {
1046        let mut st = SinkState::default();
1047        // The same identity at multiplicity 2 consolidates to a single live row.
1048        stage_m(&mut st, 5.0, 0, 1);
1049        stage_m(&mut st, 5.0, 0, 1);
1050        st.integrate(&frontier(1));
1051        st.publish_if_healthy();
1052        assert_eq!(st.published[&key_m()].0, 5.0);
1053        assert_eq!(st.collisions, 0);
1054
1055        // A second, distinct live value for the same series is a genuine collision.
1056        stage_m(&mut st, 7.0, 1, 1);
1057        st.integrate(&frontier(2));
1058        st.publish_if_healthy();
1059        assert_eq!(st.collisions, 1);
1060        // The smallest value wins deterministically.
1061        assert_eq!(st.published[&key_m()].0, 5.0);
1062    }
1063
1064    #[mz_ore::test]
1065    fn no_publish_before_time_closed() {
1066        let mut st = SinkState::default();
1067        // An update at time 5 must not appear while the frontier still allows data at time 5.
1068        stage_m(&mut st, 2.0, 5, 1);
1069        st.integrate(&frontier(5));
1070        st.publish_if_healthy();
1071        assert!(st.published.is_empty());
1072
1073        // Once the frontier advances past time 5, the update publishes.
1074        st.integrate(&frontier(6));
1075        st.publish_if_healthy();
1076        assert_eq!(st.published[&key_m()].0, 2.0);
1077    }
1078
1079    #[mz_ore::test]
1080    fn null_value_gaps_series() {
1081        let mut st = SinkState::default();
1082        // A null-valued row for (m,{a}) at a closed time: no series, counted in null_values.
1083        stage_m_null(&mut st, 1, 1);
1084        st.integrate(&frontier(2));
1085        st.publish_if_healthy();
1086        assert!(!st.published.contains_key(&key_m()));
1087        assert_eq!(st.null_values, 1);
1088
1089        // A later non-null value republishes the series (gap closes) and clears the count, even
1090        // though the null-valued row is still live alongside it.
1091        stage_m(&mut st, 5.0, 3, 1);
1092        st.integrate(&frontier(4));
1093        st.publish_if_healthy();
1094        assert_eq!(st.published[&key_m()].0, 5.0);
1095        assert_eq!(st.null_values, 0);
1096    }
1097
1098    #[mz_ore::test]
1099    fn null_labels_become_empty() {
1100        let mut st = SinkState::default();
1101        // An empty label vector (the shaped relation's `{}` for a source row with no labels)
1102        // keys and publishes correctly.
1103        st.stage_ok(
1104            "m",
1105            Some(MetricKind::Gauge),
1106            true,
1107            &[],
1108            Some(1.0),
1109            "h",
1110            Timestamp::from(1),
1111            1,
1112        );
1113        st.integrate(&frontier(2));
1114        st.publish_if_healthy();
1115        assert_eq!(st.published[&("m".into(), vec![])].0, 1.0);
1116    }
1117
1118    #[mz_ore::test]
1119    fn extract_row_normalizes_null_datums() {
1120        let desc = shaped_desc();
1121        let cols = ColumnIndices::resolve(&desc);
1122
1123        let mut row = Row::default();
1124        {
1125            let mut packer = row.packer();
1126            packer.push(Datum::Null); // metric_name
1127            packer.push_dict_with(|_| {}); // labels: always non-null by construction
1128            packer.push(Datum::Null); // value
1129            packer.push(Datum::String("")); // help: always non-null by construction
1130            packer.push(Datum::Null); // metric_kind: defensively treated as unsupported
1131            packer.push(Datum::Null); // name_valid: defensively treated as invalid
1132        }
1133
1134        let datums: Vec<Datum> = row.iter().collect();
1135        let (name, metric_kind, name_valid, labels, value, help) = extract_row(&cols, &datums);
1136        assert_eq!(name, "");
1137        assert_eq!(metric_kind, None);
1138        assert!(!name_valid);
1139        assert_eq!(labels, Vec::new());
1140        assert_eq!(value, None);
1141        assert_eq!(help, "");
1142    }
1143
1144    /// A null map value must survive extraction as `None`; unwrapping it as a string panicked the
1145    /// worker and took down the replica.
1146    #[mz_ore::test]
1147    fn extract_row_keeps_null_label_values() {
1148        let desc = shaped_desc();
1149        let cols = ColumnIndices::resolve(&desc);
1150
1151        let mut row = Row::default();
1152        {
1153            let mut packer = row.packer();
1154            packer.push(Datum::String("m"));
1155            packer.push_dict_with(|row| {
1156                row.push(Datum::String("bad"));
1157                row.push(Datum::Null);
1158                row.push(Datum::String("good"));
1159                row.push(Datum::String("1"));
1160            });
1161            packer.push(Datum::Float64(1.0.into()));
1162            packer.push(Datum::String("h"));
1163            packer.push(Datum::Int32(0));
1164            packer.push(Datum::True);
1165        }
1166
1167        let datums: Vec<Datum> = row.iter().collect();
1168        let (_name, _metric_kind, _name_valid, labels, _value, _help) = extract_row(&cols, &datums);
1169        assert_eq!(labels, vec![("bad", None), ("good", Some("1"))]);
1170    }
1171
1172    #[mz_ore::test]
1173    fn null_or_empty_label_value_skips_row() {
1174        let mut st = SinkState::default();
1175        // A null label value has no representable label set: publishes nothing, counts as skipped.
1176        st.stage_ok(
1177            "m",
1178            Some(MetricKind::Gauge),
1179            true,
1180            &[("a", None)],
1181            Some(1.0),
1182            "h",
1183            Timestamp::from(1),
1184            1,
1185        );
1186        st.integrate(&frontier(2));
1187        st.publish_if_healthy();
1188        assert!(st.published.is_empty());
1189        assert_eq!(st.skipped, 1);
1190
1191        // An empty label value folds to `{}` in Prometheus, so it is skipped too rather than
1192        // published as a bare `{}` series. It is a distinct identity from the null row, so both
1193        // stay live and count: `skipped` reaches 2, not 1.
1194        st.stage_ok(
1195            "m",
1196            Some(MetricKind::Gauge),
1197            true,
1198            &[("a", Some(""))],
1199            Some(1.0),
1200            "h",
1201            Timestamp::from(3),
1202            1,
1203        );
1204        st.integrate(&frontier(4));
1205        st.publish_if_healthy();
1206        assert!(st.published.is_empty());
1207        assert_eq!(st.skipped, 2);
1208    }
1209
1210    fn pkey(name: &str, labels: &[(&str, &str)]) -> PublishedKey {
1211        (
1212            name.to_string(),
1213            labels
1214                .iter()
1215                .map(|&(k, v)| (k.to_string(), v.to_string()))
1216                .collect(),
1217        )
1218    }
1219
1220    #[mz_ore::test]
1221    fn build_families_groups_by_name_and_kind() {
1222        let published: BTreeMap<PublishedKey, PublishedValue> = BTreeMap::from([
1223            (
1224                pkey("http_requests", &[("code", "200")]),
1225                (5.0, MetricKind::Counter, "requests".to_string()),
1226            ),
1227            (
1228                pkey("http_requests", &[("code", "500")]),
1229                (2.0, MetricKind::Counter, "requests".to_string()),
1230            ),
1231            (
1232                pkey("temp_celsius", &[]),
1233                (21.5, MetricKind::Gauge, "temperature".to_string()),
1234            ),
1235        ]);
1236
1237        let families = build_families(&published);
1238
1239        // One family per metric name, in `BTreeMap` (name) order.
1240        assert_eq!(families.len(), 2);
1241
1242        let requests = &families[0];
1243        assert_eq!(requests.name(), "http_requests");
1244        assert_eq!(requests.help(), "requests");
1245        let metrics = requests.get_metric();
1246        assert_eq!(metrics.len(), 2);
1247        // Metrics keep the `BTreeMap` label order, and land in the counter oneof.
1248        assert_eq!(metrics[0].get_label()[0].value(), "200");
1249        assert_eq!(metrics[0].get_counter().value(), 5.0);
1250        assert_eq!(metrics[1].get_label()[0].value(), "500");
1251        assert_eq!(metrics[1].get_counter().value(), 2.0);
1252
1253        let temp = &families[1];
1254        assert_eq!(temp.name(), "temp_celsius");
1255        let temp_metrics = temp.get_metric();
1256        assert_eq!(temp_metrics.len(), 1);
1257        assert_eq!(temp_metrics[0].get_gauge().value(), 21.5);
1258    }
1259
1260    #[mz_ore::test]
1261    fn count_conflicts_flags_type_and_help_disagreement() {
1262        // The family winner is the smallest-label entry. Here `m`'s winner is the `[a=1]` gauge with
1263        // help `h1`; the other two disagree on kind, then help.
1264        let published: BTreeMap<PublishedKey, PublishedValue> = BTreeMap::from([
1265            (
1266                pkey("m", &[("a", "1")]),
1267                (1.0, MetricKind::Gauge, "h1".to_string()),
1268            ),
1269            (
1270                pkey("m", &[("b", "2")]),
1271                (2.0, MetricKind::Counter, "h1".to_string()),
1272            ),
1273            (
1274                pkey("m", &[("c", "3")]),
1275                (3.0, MetricKind::Gauge, "h2".to_string()),
1276            ),
1277            (
1278                pkey("other", &[]),
1279                (1.0, MetricKind::Gauge, "h".to_string()),
1280            ),
1281        ]);
1282        assert_eq!(count_conflicts(&published), 2);
1283
1284        // A family whose entries all agree has no conflicts.
1285        let consistent: BTreeMap<PublishedKey, PublishedValue> = BTreeMap::from([
1286            (
1287                pkey("m", &[("a", "1")]),
1288                (1.0, MetricKind::Gauge, "h".to_string()),
1289            ),
1290            (
1291                pkey("m", &[("b", "2")]),
1292                (2.0, MetricKind::Gauge, "h".to_string()),
1293            ),
1294        ]);
1295        assert_eq!(count_conflicts(&consistent), 0);
1296    }
1297
1298    #[mz_ore::test]
1299    fn err_stream_freezes_and_recovers() {
1300        let mut st = SinkState::default();
1301        // Establish a healthy value.
1302        stage_m(&mut st, 5.0, 0, 1);
1303        st.integrate(&frontier(1));
1304        st.publish_if_healthy();
1305        assert_eq!(st.published[&key_m()].0, 5.0);
1306
1307        // An error appears, buffered through `stage_err` (not by setting `errors` directly). A value
1308        // update lands in the same window.
1309        st.stage_err(Timestamp::from(1), 1);
1310        stage_m(&mut st, 5.0, 1, -1);
1311        stage_m(&mut st, 9.0, 1, 1);
1312        st.integrate(&frontier(2));
1313        assert_eq!(st.errors, 1);
1314        st.publish_if_healthy();
1315        // Publication frozen at the last healthy value while erroring.
1316        assert_eq!(st.published[&key_m()].0, 5.0);
1317
1318        // The error is retracted; net errors returns to 0 and publication recovers to the value
1319        // integrated during the freeze.
1320        st.stage_err(Timestamp::from(2), -1);
1321        st.integrate(&frontier(3));
1322        assert_eq!(st.errors, 0);
1323        st.publish_if_healthy();
1324        assert_eq!(st.published[&key_m()].0, 9.0);
1325    }
1326
1327    /// Sums a counter family's samples across the registry's scrape output.
1328    fn counter_total(registry: &MetricsRegistry, name: &str) -> f64 {
1329        registry
1330            .gather()
1331            .iter()
1332            .filter(|family| family.name() == name)
1333            .flat_map(|family| family.get_metric())
1334            .map(|metric| metric.get_counter().value())
1335            .sum()
1336    }
1337
1338    /// Counts the samples of family `metric` whose `label_key` label equals `label`.
1339    fn companion_gauge_count(
1340        registry: &MetricsRegistry,
1341        metric: &str,
1342        label_key: &str,
1343        label: &str,
1344    ) -> usize {
1345        registry
1346            .gather()
1347            .iter()
1348            .filter(|family| family.name() == metric)
1349            .flat_map(|family| family.get_metric())
1350            .filter(|metric| {
1351                metric
1352                    .get_label()
1353                    .iter()
1354                    .any(|l| l.name() == label_key && l.value() == label)
1355            })
1356            .count()
1357    }
1358
1359    /// The register/retry/unregister ordering the operator runs on every activation, with a stand-in
1360    /// for the previous incarnation of the same curated sink. Two incarnations share a `Desc` id
1361    /// because a curated sink's companion gauges are keyed on its stable name, so the new one can
1362    /// only register once the old one's handle drops.
1363    #[mz_ore::test]
1364    fn pending_registration_retries_until_predecessor_drops() {
1365        const LABEL: &str = "mz_curated_example";
1366        const RETRIES: &str = "mz_compute_metric_sink_registration_retries_total";
1367        // A live collector emits exactly one frontier gauge, so its presence stands in for
1368        // "collector registered" across the register/drop lifecycle below.
1369        const FRONTIER: &str = "mz_compute_metric_sink_frontier_ms";
1370        const SINK_LABEL: &str = "sink";
1371
1372        let registry = MetricsRegistry::new();
1373        let metrics =
1374            ComputeMetrics::register_with(&registry, ComputeRuntimeRole::Solo).for_worker(0);
1375
1376        // The predecessor: an incarnation of the same sink still registered.
1377        let incumbent = registry.register_collector_with_dropper(SinkCollector::new(
1378            LABEL,
1379            Arc::new(Mutex::new(SinkState::default())),
1380        ));
1381        assert_eq!(
1382            companion_gauge_count(&registry, FRONTIER, SINK_LABEL, LABEL),
1383            1
1384        );
1385
1386        let mut registration = PendingRegistration::new(
1387            LABEL.to_string(),
1388            SinkCollector::new(LABEL, Arc::new(Mutex::new(SinkState::default()))),
1389        );
1390
1391        assert_eq!(registration.try_register(&registry, &metrics), Retry::Arm);
1392        assert_eq!(counter_total(&registry, RETRIES), 1.0);
1393        // The failed attempt registered nothing, so the incumbent is still the only series.
1394        assert_eq!(
1395            companion_gauge_count(&registry, FRONTIER, SINK_LABEL, LABEL),
1396            1
1397        );
1398
1399        // A retry is armed, so an activation before the interval elapses attempts nothing.
1400        assert_eq!(registration.try_register(&registry, &metrics), Retry::Skip);
1401        assert_eq!(counter_total(&registry, RETRIES), 1.0);
1402
1403        // Clearing `next_attempt` stands in for the retry interval elapsing.
1404        registration.next_attempt = None;
1405        assert_eq!(registration.try_register(&registry, &metrics), Retry::Arm);
1406        assert_eq!(counter_total(&registry, RETRIES), 2.0);
1407
1408        drop(incumbent);
1409        registration.next_attempt = None;
1410        assert_eq!(registration.try_register(&registry, &metrics), Retry::Skip);
1411        assert_eq!(counter_total(&registry, RETRIES), 2.0);
1412        assert_eq!(
1413            companion_gauge_count(&registry, FRONTIER, SINK_LABEL, LABEL),
1414            1
1415        );
1416
1417        // Already registered: a further activation is a no-op, not a second registration.
1418        assert_eq!(registration.try_register(&registry, &metrics), Retry::Skip);
1419        assert_eq!(counter_total(&registry, RETRIES), 2.0);
1420
1421        // Dropping the registration is what frees the `Desc` id for the next incarnation.
1422        drop(registration);
1423        assert_eq!(
1424            companion_gauge_count(&registry, FRONTIER, SINK_LABEL, LABEL),
1425            0
1426        );
1427    }
1428
1429    /// The registration follows the sink token, not the operator closure. Builds a real timely
1430    /// operator with `render_sink`'s wiring (register at build time, keep only a clone of the
1431    /// registration in the operator closure, return the original as the sink token) around a
1432    /// constant input that closes at once. When timely retires the operator and drops its closure
1433    /// the collector's series must survive, because the token still holds the registration; only
1434    /// dropping the token unregisters. Building a real operator, rather than dropping an `Rc` by
1435    /// hand, is what proves the closure does not carry the registering handle.
1436    ///
1437    /// It stops short of driving `render_sink` itself, which would need a full `ComputeState`.
1438    #[mz_ore::test]
1439    fn registration_survives_input_close_in_dataflow() {
1440        use differential_dataflow::input::Input;
1441        use timely::WorkerConfig;
1442        use timely::communication::Allocator;
1443        use timely::dataflow::channels::pact::Pipeline;
1444        use timely::worker::Worker as TimelyWorker;
1445
1446        const LABEL: &str = "mz_curated_example";
1447        const FRONTIER: &str = "mz_compute_metric_sink_frontier_ms";
1448        const SINK_LABEL: &str = "sink";
1449
1450        let registry = MetricsRegistry::new();
1451        let metrics =
1452            ComputeMetrics::register_with(&registry, ComputeRuntimeRole::Solo).for_worker(0);
1453
1454        let mut worker = TimelyWorker::new(
1455            WorkerConfig::default(),
1456            Allocator::Thread(Default::default()),
1457            Some(Instant::now()),
1458        );
1459
1460        // `MetricsRegistry` and `WorkerMetrics` are cheap shared handles, so the clones the operator
1461        // registers through and the originals the assertions scrape are the same registry.
1462        let registry_op = registry.clone();
1463        let metrics_op = metrics.clone();
1464        let registration = worker.dataflow::<Timestamp, _, _>(move |scope| {
1465            // Dropping the input handle immediately closes the input, standing in for a
1466            // constant-folded source whose frontier closes as soon as the dataflow starts.
1467            let (_input, collection) = scope.new_collection::<Row, Diff>();
1468            let registration = Rc::new(RefCell::new(Some(PendingRegistration::new(
1469                LABEL.to_string(),
1470                SinkCollector::new(LABEL, Arc::new(Mutex::new(SinkState::default()))),
1471            ))));
1472            let registration_op = Rc::clone(&registration);
1473
1474            let scope_for_activator = scope.clone();
1475            let mut op = OperatorBuilder::new("MetricSinkTest".to_string(), scope.clone());
1476            let mut input = op.new_input(collection.inner, Pipeline);
1477            let operator_info = op.operator_info();
1478            op.build(move |_caps| {
1479                let activator = scope_for_activator.activator_for(operator_info.address);
1480                if let Some(reg) = registration_op.borrow_mut().as_mut() {
1481                    if let Retry::Arm = reg.try_register(&registry_op, &metrics_op) {
1482                        activator.activate_after(REGISTRATION_RETRY_INTERVAL);
1483                    }
1484                }
1485                move |_frontiers| {
1486                    input.for_each(|_, _| {});
1487                }
1488            });
1489
1490            registration
1491        });
1492
1493        // The build-time attempt registered the collector.
1494        assert_eq!(
1495            companion_gauge_count(&registry, FRONTIER, SINK_LABEL, LABEL),
1496            1
1497        );
1498
1499        // Step until timely retires the operator and drops its closure, releasing its `Rc` clone.
1500        for _ in 0..100 {
1501            if Rc::strong_count(&registration) == 1 {
1502                break;
1503            }
1504            worker.step();
1505        }
1506        assert_eq!(
1507            Rc::strong_count(&registration),
1508            1,
1509            "the operator closure should have dropped once its input closed"
1510        );
1511
1512        // The registration rode on the token, not the closure, so the series is still live.
1513        assert_eq!(
1514            companion_gauge_count(&registry, FRONTIER, SINK_LABEL, LABEL),
1515            1
1516        );
1517
1518        // Dropping the token is what unregisters.
1519        drop(registration);
1520        assert_eq!(
1521            companion_gauge_count(&registry, FRONTIER, SINK_LABEL, LABEL),
1522            0
1523        );
1524    }
1525
1526    /// Every companion gauge carries the label the collector was built with. A curated sink passes
1527    /// its stable name here, so its health series stay identifiable across boots even though the
1528    /// sink's `GlobalId` is transient.
1529    #[mz_ore::test]
1530    fn collector_labels_gauges_with_sink_label() {
1531        let state = Arc::new(Mutex::new(SinkState::default()));
1532        let collector = SinkCollector::new("mz_curated_example", state);
1533
1534        let families = collector.collect();
1535        // The point is that every emitted family carries the label, so guard only that there is
1536        // something to check. Pinning the exact count churns on every added gauge without telling
1537        // the next reader whether they broke labelling or just added a family.
1538        assert!(!families.is_empty());
1539        for family in &families {
1540            for metric in family.get_metric() {
1541                let sink = metric
1542                    .get_label()
1543                    .iter()
1544                    .find(|l| l.name() == "sink")
1545                    .expect("sink label present");
1546                assert_eq!(sink.value(), "mz_curated_example");
1547            }
1548        }
1549    }
1550}