Skip to main content

mz_adapter/coord/
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//! Coordinator-installed metric sinks, the curated counterpart to `CREATE METRIC SINK`.
11//!
12//! A curated metric sink is a [`CURATED`] entry rendered on every replica, publishing its series
13//! into that replica's process-local Prometheus registry. Unlike a user's `CREATE METRIC SINK` it
14//! is not a catalog item: it gets a transient [`GlobalId`], targets one replica rather than a
15//! cluster, and is re-created from the static list on every boot. Modelling the curated set this
16//! way keeps it out of the catalog, so adding or removing a definition needs no builtin migration.
17//!
18//! Every replica means every replica of every cluster, user clusters included. Each definition is
19//! therefore a dataflow, with its arrangements, on customer compute, charged to that customer's
20//! cluster, and the cost scales with `CURATED`. `coord::introspection` already accepts this for its
21//! subscribes.
22//!
23//! `install_metric_sinks` installs every definition on a newly created replica
24//! (`bootstrap_metric_sinks` covers the replicas already present at startup), and
25//! `drop_metric_sinks` drops them before a replica is dropped. This mirrors
26//! [`crate::coord::introspection`], which installs introspection subscribes on the same triggers.
27//!
28//! The `disabled_metric_sinks` system var denies definitions by name, and
29//! `reconcile_metric_sinks` converges the installed set on it, tearing a denied sink down rather
30//! than only gating future installs.
31
32use std::collections::{BTreeMap, BTreeSet};
33
34use anyhow::bail;
35use mz_catalog::memory::objects::CatalogItem;
36use mz_cluster_client::ReplicaId;
37use mz_controller_types::ClusterId;
38use mz_ore::collections::CollectionExt;
39use mz_ore::{instrument, soft_panic_or_log};
40use mz_repr::optimize::OverrideFrom;
41use mz_repr::{CatalogItemId, GlobalId, RelationDesc};
42use mz_sql::catalog::SessionCatalog;
43use mz_sql::plan::{
44    HirRelationExpr, METRIC_SINK_CURATED_PREFIX_MARKER, Params, Plan, SubscribeFrom, SubscribePlan,
45    validate_metric_sink_desc, validate_metric_sink_prefix,
46};
47use mz_sql::session::user::{MZ_SYSTEM_ROLE_ID, RoleMetadata};
48use mz_sql::session::vars::{ENABLE_METRIC_SINK, SystemVars};
49use tracing::{Span, info, warn};
50
51use crate::catalog::Catalog;
52use crate::coord::{
53    Coordinator, Message, MetricSinkFinish, MetricSinkOptimize, MetricSinkStage, PlanValidity,
54    StageResult, Staged,
55};
56use crate::optimize::Optimize;
57use crate::optimize::dataflows::dataflow_import_id_bundle;
58use crate::{AdapterError, ExecuteResponse, optimize};
59
60/// A curated metric sink: SQL producing the canonical metric-sink columns, plus the name it is
61/// known by in logs.
62#[derive(Debug)]
63pub(super) struct CuratedMetricSink {
64    /// Stable identifier for the definition: used in logs, as the [`Coordinator::metric_sinks`] key,
65    /// and as the `sink` label on the health gauges (the `GlobalId` is transient, the name is not).
66    /// Must be unique within [`CURATED`].
67    name: &'static str,
68    /// A `SELECT` producing the canonical metric-sink columns (`metric_name`, `metric_type`,
69    /// `labels`, `value`, `help`), the contract `mz_sql::plan::validate_metric_sink_desc` checks.
70    ///
71    /// The query must read only introspection relations. A catalog-backed relation would put
72    /// envd's write frontier on the sink's emission path, which is exactly the coupling these
73    /// sinks exist to avoid: the sink would stall whenever envd did, taking the freshness signal
74    /// with it.
75    source_sql: &'static str,
76    /// Prepended to every row's `metric_name` to form the published name, exactly as a user's
77    /// `CREATE METRIC SINK ... WITH (PREFIX = ...)`. Every definition in [`CURATED`] uses
78    /// [`METRIC_SINK_CURATED_PREFIX_MARKER`], which user sinks are barred from, so nothing a user
79    /// publishes can collide with a curated family.
80    prefix: &'static str,
81}
82
83/// The curated metric sinks, installed on every replica.
84///
85/// Sources read the raw `..._raw` logging relations, not a derived view that re-aggregates them
86/// (`mz_dataflow_arrangement_sizes`) or a join view over the logs (`mz_dataflow_operator_dataflows`),
87/// for performance: those churn even on static data and their cost scales with the replica's
88/// dataflow-creation rate. A single-relation filter like `mz_compute_exports` carries no aggregation
89/// and is read freely.
90///
91/// Every family sums across workers, so a series carries no `worker_id`, and a multi-process replica
92/// reports one number per grouping key rather than one per process. The size families emit one series
93/// per dataflow, the errors family one per export.
94const CURATED: &[CuratedMetricSink] = &[
95    CuratedMetricSink {
96        name: "mz_metric_arrangement_sizes",
97        prefix: METRIC_SINK_CURATED_PREFIX_MARKER,
98        // Key on the export id from `mz_compute_exports`, not `mz_dataflow_global_ids`: a
99        // materialized view builds under a transient view id and appears only as an export, so a
100        // global-id label names a `t<N>` that maps to no catalog object and churns on every
101        // re-render.
102        //
103        // Logs are per operator, so map operator -> dataflow -> export id. Take the operator ->
104        // dataflow step off `mz_dataflow_addresses_per_worker` (`address[1]` is the dataflow id),
105        // which avoids the multi-way join behind `mz_dataflow_operator_dataflows`.
106        //
107        // `min(export_id)` collapses a multi-export dataflow to one series (lexicographic, so
108        // `min('u10', 'u2')` is `'u10'`: arbitrary but stable). The group-size hint stops that `min`
109        // from rendering the 8-level hierarchy, which would otherwise show up as tuning advice for
110        // the sink's own dataflow in `mz_expected_group_size_advice`.
111        //
112        // NOTE: counting an arrangement this sink reads is a feedback loop. Every change to it
113        // changes its logged size, the sink reads that on the next logging tick, and the dataflow
114        // re-runs every tick for the replica's lifetime (SQL-730). `NOT LIKE 't%'` keeps out the
115        // sink's own dataflow and the replica's introspection subscribes. `NOT LIKE 'si%'` keeps
116        // out the logging dataflow: `si<N>` names only the introspection source indexes it
117        // exports, and a plain system index prints `s<N>`. Only that filter excludes it under
118        // `INTROSPECTION DEBUGGING`, which registers the loggers before the logging dataflow is
119        // built, so its own operators get log rows. The cost is that transient dataflows'
120        // arrangements go unreported, so these families sum below what the replica holds.
121        //
122        // `f` joins each raw log separately to reuse its `(operator_id, worker_id)` index
123        // (`LogVariant::index_by`). One union would arrange all three logs' rows afresh.
124        source_sql: "
125WITH ex AS (
126    SELECT dataflow_id, min(export_id) AS export_id
127    FROM mz_introspection.mz_compute_exports
128    WHERE export_id NOT LIKE 't%' AND export_id NOT LIKE 'si%'
129    GROUP BY dataflow_id OPTIONS (AGGREGATE INPUT GROUP SIZE = 1)
130),
131oe AS (
132    SELECT a.id, a.worker_id, ex.export_id
133    FROM mz_introspection.mz_dataflow_addresses_per_worker a
134    JOIN ex ON ex.dataflow_id = a.address[1]
135),
136f AS (
137    SELECT 'arrangement_size_bytes'::text AS metric_name, oe.export_id
138    FROM mz_introspection.mz_arrangement_heap_size_raw r
139    JOIN oe ON r.operator_id = oe.id AND r.worker_id = oe.worker_id
140    UNION ALL
141    SELECT 'arrangement_records'::text, oe.export_id
142    FROM mz_introspection.mz_arrangement_records_raw r
143    JOIN oe ON r.operator_id = oe.id AND r.worker_id = oe.worker_id
144    UNION ALL
145    SELECT 'arrangement_batches'::text, oe.export_id
146    FROM mz_introspection.mz_arrangement_batches_raw r
147    JOIN oe ON r.operator_id = oe.id AND r.worker_id = oe.worker_id
148)
149SELECT metric_name, 'gauge'::text AS metric_type,
150       map_build(LIST[ROW('id', export_id)])::map[text=>text] AS labels,
151       count(*)::double precision AS value,
152       CASE metric_name
153           WHEN 'arrangement_size_bytes' THEN 'arrangement heap size in bytes'
154           WHEN 'arrangement_records' THEN 'number of records in arrangement heaps'
155           WHEN 'arrangement_batches' THEN 'number of batches in arrangements'
156       END AS help
157FROM f
158GROUP BY metric_name, export_id",
159    },
160    CuratedMetricSink {
161        name: "mz_metric_dataflow_errors",
162        prefix: METRIC_SINK_CURATED_PREFIX_MARKER,
163        // Raw log, not the `mz_compute_error_counts` view: the view joins the storage-managed
164        // `mz_internal.mz_compute_dependencies`, which `ensure_reads_only_logs` rejects. `count` is
165        // per-worker, so sum per export. `HAVING` drops the healthy ones.
166        //
167        // NOTE: direct errors only. The raw log attributes an error to the export that raised it. The
168        // view also forwards counts onto index-reuse exports, so a broken reuse-index reads 0 here and
169        // its errors show under the underlying export, undercounting against the view.
170        source_sql: "
171SELECT 'dataflow_error_count'::text AS metric_name, 'gauge'::text AS metric_type,
172       map_build(LIST[ROW('id', export_id::text)])::map[text=>text] AS labels,
173       sum(count)::double precision AS value, 'count of errors in the dataflow'::text AS help
174FROM mz_introspection.mz_compute_error_counts_raw
175GROUP BY export_id
176HAVING sum(count) > 0",
177    },
178];
179
180/// A [`CuratedMetricSink`] installed on one replica.
181#[derive(Debug)]
182pub(super) struct InstalledMetricSink {
183    /// The cluster the replica belongs to, needed to drop the sink's compute collection.
184    cluster_id: ClusterId,
185    /// The transient id of the sink's compute export.
186    sink_id: GlobalId,
187}
188
189/// A [`CuratedMetricSink`] planned once and shared across the replicas it installs on. See
190/// [`Coordinator::plan_metric_sink`].
191#[derive(Clone, Debug)]
192pub(super) struct PlannedMetricSink {
193    /// The shaped source query.
194    expr: HirRelationExpr,
195    /// The shape `expr` produces.
196    desc: RelationDesc,
197    /// The catalog items the source reads.
198    dependencies: BTreeSet<CatalogItemId>,
199}
200
201impl Coordinator {
202    /// Installs the curated metric sinks on all existing replicas.
203    pub(super) async fn bootstrap_metric_sinks(&mut self) {
204        for (cluster_id, replica_id) in self.all_cluster_replicas() {
205            self.install_metric_sinks(cluster_id, replica_id).await;
206        }
207    }
208
209    /// Installs the curated metric sinks on the given replica.
210    ///
211    /// Turning `enable_metric_sink` off stops installing on replicas created from then on. It does
212    /// not tear down what is already installed: those keep running until their replica is dropped
213    /// or envd restarts. A replica that merely reconnects re-renders them from the controller's
214    /// command history, so a replica restart does not clear them either.
215    pub(super) async fn install_metric_sinks(
216        &mut self,
217        cluster_id: ClusterId,
218        replica_id: ReplicaId,
219    ) {
220        if !ENABLE_METRIC_SINK.enabled(self.catalog().system_config()) {
221            return;
222        }
223
224        // TODO: Skip replicas created with introspection disabled. Their logging dataflows never
225        // run, so a `source_sql` reading introspection relations there never advances. That is not
226        // just wasted work: the sink publishes its input frontier as its write frontier, so a
227        // never-advancing input stalls the sink's frontier at its as-of and pins the read holds it
228        // takes on those collections for the replica's whole life (replica-local, released on
229        // drop). `coord::introspection` installs subscribes on the same triggers and has the same
230        // gap.
231        for definition in CURATED {
232            if metric_sink_denied(self.catalog().system_config(), definition.name) {
233                continue;
234            }
235            self.install_metric_sink(cluster_id, replica_id, definition)
236                .await;
237        }
238    }
239
240    /// Converges the installed curated sinks on `disabled_metric_sinks`.
241    ///
242    /// Reconciles the whole set rather than the delta.
243    pub(super) async fn reconcile_metric_sinks(&mut self) {
244        for entry in self.catalog().system_config().disabled_metric_sinks() {
245            if !CURATED.iter().any(|d| d.name == entry) {
246                warn!(
247                    name = %entry,
248                    "disabled_metric_sinks entry matches no curated definition"
249                );
250            }
251        }
252
253        let denied: Vec<_> = self
254            .metric_sinks
255            .keys()
256            .copied()
257            .filter(|(_, name)| metric_sink_denied(self.catalog().system_config(), name))
258            .collect();
259        for (replica_id, name) in denied {
260            self.drop_metric_sink(replica_id, name);
261        }
262
263        // Reinstall the full non-denied set on every replica. `install_metric_sink` is
264        // idempotent: it skips a definition already recorded in `metric_sinks`
265        // (`contains_key`, see `install_metric_sink`), so re-running the whole set only
266        // installs the ones a preceding `disabled_metric_sinks` edit un-denied.
267        for (cluster_id, replica_id) in self.all_cluster_replicas() {
268            self.install_metric_sinks(cluster_id, replica_id).await;
269        }
270    }
271
272    async fn install_metric_sink(
273        &mut self,
274        cluster_id: ClusterId,
275        replica_id: ReplicaId,
276        definition: &'static CuratedMetricSink,
277    ) {
278        // Cheap duplicate check before planning: if the definition is already installed on this
279        // replica, there is nothing to do. `metric_sink_finish` keeps a backstop for a double
280        // install still in flight (not yet recorded here).
281        if self
282            .metric_sinks
283            .contains_key(&(replica_id, definition.name))
284        {
285            return;
286        }
287
288        let Some(planned) = self.plan_metric_sink(definition) else {
289            return;
290        };
291
292        let (_, sink_id) = self.allocate_transient_id();
293        // Logged only once the definition is known good, so an abandoned install leaves no
294        // misleading "installing" line.
295        info!(%sink_id, %replica_id, name = definition.name, "installing metric sink");
296
297        let validity = PlanValidity::new(
298            &self.catalog,
299            planned.dependencies.clone(),
300            Some(cluster_id),
301            Some(replica_id),
302            RoleMetadata::new(MZ_SYSTEM_ROLE_ID),
303        );
304        let stage = MetricSinkStage::Optimize(MetricSinkOptimize {
305            validity,
306            definition,
307            sink_id,
308            expr: planned.expr.clone(),
309            desc: planned.desc.clone(),
310            cluster_id,
311            replica_id,
312        });
313        self.sequence_staged((), Span::current(), stage).await;
314    }
315
316    /// Plans a curated definition once, caching the result in [`Coordinator::metric_sink_plans`].
317    ///
318    /// The plan depends only on the catalog, never on the replica, so it is shared across every
319    /// replica the definition installs on rather than re-planned per replica. Curated sources read
320    /// only builtins (enforced by [`ensure_reads_only_logs`]), which do not change while envd runs,
321    /// so a cached plan stays valid for envd's lifetime. Returns `None` for an invalid definition,
322    /// having soft-panicked.
323    fn plan_metric_sink(
324        &mut self,
325        definition: &'static CuratedMetricSink,
326    ) -> Option<PlannedMetricSink> {
327        if let Some(planned) = self.metric_sink_plans.get(definition.name) {
328            return Some(planned.clone());
329        }
330
331        // A user sink's prefix is validated at plan time; a curated one has no such gate, so enforce
332        // the same contract here. A failure is a bug in our own definition, hence
333        // `soft_panic_or_log!`. User-vs-curated collisions need no check: the curated prefix is
334        // reserved against user sinks in `validate_user_metric_sink_prefix`.
335        //
336        // NOTE: curated definitions are not checked against each other; they stay disjoint by
337        // publishing distinct `metric_name`s under the shared curated prefix.
338        if let Err(err) = validate_metric_sink_prefix(definition.prefix) {
339            soft_panic_or_log!(
340                "invalid curated metric sink prefix (name={}): {err}",
341                definition.name
342            );
343            return None;
344        }
345
346        let catalog = self.catalog().for_system_session();
347        let (expr, desc, dependencies) = match definition.plan_source(&catalog) {
348            Ok(planned) => planned,
349            Err(err) => {
350                soft_panic_or_log!(
351                    "invalid curated metric sink (name={}): {err}",
352                    definition.name
353                );
354                return None;
355            }
356        };
357
358        // Enforce the introspection-only contract before any optimization work, against what the
359        // definition reads rather than how the optimizer imports it.
360        if let Err(err) = ensure_reads_only_logs(&self.catalog, &dependencies) {
361            soft_panic_or_log!(
362                "invalid curated metric sink (name={}): {err}",
363                definition.name
364            );
365            return None;
366        }
367
368        let planned = PlannedMetricSink {
369            expr,
370            desc,
371            dependencies,
372        };
373        self.metric_sink_plans
374            .insert(definition.name, planned.clone());
375        Some(planned)
376    }
377
378    #[instrument]
379    fn metric_sink_optimize(
380        &self,
381        stage: MetricSinkOptimize,
382    ) -> Result<StageResult<Box<MetricSinkStage>>, AdapterError> {
383        let MetricSinkOptimize {
384            mut validity,
385            definition,
386            sink_id,
387            expr,
388            desc,
389            cluster_id,
390            replica_id,
391        } = stage;
392
393        let compute_instance = self
394            .instance_snapshot(cluster_id)
395            .expect("compute instance exists");
396        // A transient id for the view the optimizer builds to shape the source rows, scoped to this
397        // dataflow. See `optimize::metric_sink::shape_metric_sink_source`.
398        let (_, view_id) = self.allocate_transient_id();
399
400        let optimizer_config = optimize::OptimizerConfig::from(self.catalog().system_config())
401            .override_from(&self.catalog.get_cluster(cluster_id).config.features())
402            .override_from(&self.cluster_scoped_optimizer_overrides(cluster_id));
403
404        let mut optimizer = optimize::metric_sink::Optimizer::new(
405            self.owned_catalog(),
406            compute_instance,
407            view_id,
408            sink_id,
409            optimizer_config,
410            self.optimizer_metrics(),
411        );
412        let catalog = self.owned_catalog();
413
414        let span = Span::current();
415        Ok(StageResult::Handle(mz_ore::task::spawn_blocking(
416            || "optimize metric sink",
417            move || {
418                span.in_scope(|| {
419                    let metric_sink = optimize::metric_sink::MetricSink::new(
420                        format!("metric-sink-{}-{replica_id}", definition.name),
421                        optimize::metric_sink::MetricSinkFrom::Query { expr, desc },
422                        definition.prefix.to_string(),
423                        Some(definition.name.to_string()),
424                    );
425
426                    // Both steps run inside one closure so either failure hits the same log.
427                    // `sequence_staged` has no session to report to for a coordinator-driven
428                    // install, so an error would otherwise vanish.
429                    let global_lir_plan = (|| {
430                        // MIR ⇒ MIR optimization (global)
431                        let global_mir_plan = optimizer.catch_unwind_optimize(metric_sink)?;
432                        // The optimizer imports indexes the SQL never named. Fold them into
433                        // validity so one dropped before the finish stage fails the recheck rather
434                        // than shipping a dataflow that imports a gone collection.
435                        let id_bundle =
436                            dataflow_import_id_bundle(global_mir_plan.df_desc(), cluster_id);
437                        let item_ids = id_bundle.iter().map(|id| catalog.resolve_item_id(&id));
438                        validity.extend_dependencies(&catalog, item_ids);
439                        // MIR ⇒ LIR lowering and LIR ⇒ LIR optimization (global)
440                        optimizer.catch_unwind_optimize(global_mir_plan)
441                    })()
442                    .inspect_err(|err| {
443                        soft_panic_or_log!(
444                            "curated metric sink failed to optimize (name={}): {err}",
445                            definition.name
446                        )
447                    })?;
448
449                    let stage = MetricSinkStage::Finish(MetricSinkFinish {
450                        validity,
451                        definition,
452                        sink_id,
453                        global_lir_plan,
454                        cluster_id,
455                        replica_id,
456                    });
457                    Ok(Box::new(stage))
458                })
459            },
460        )))
461    }
462
463    #[instrument]
464    async fn metric_sink_finish(
465        &mut self,
466        stage: MetricSinkFinish,
467    ) -> Result<StageResult<Box<MetricSinkStage>>, AdapterError> {
468        let MetricSinkFinish {
469            validity: _,
470            definition,
471            sink_id,
472            global_lir_plan,
473            cluster_id,
474            replica_id,
475        } = stage;
476
477        // `sequence_staged` rechecked validity before this stage ran, so the replica still exists.
478        // The coordinator handles one message at a time, so no replica drop runs between that check
479        // and the ship below.
480
481        // The metainfo is dropped rather than persisted: a curated sink is not a catalog item, so
482        // there is nothing for `mz_optimizer_notices` to hang its notices off.
483        let (mut df_desc, _df_meta) = global_lir_plan.unapply();
484
485        let id_bundle = dataflow_import_id_bundle(&df_desc, cluster_id);
486
487        // Backstop for the introspection-only contract; the real gate is `ensure_reads_only_logs`
488        // at install time. A log-only source imports only compute collections, so this should never
489        // fire, but a storage import would couple the sink's frontier to envd.
490        if !id_bundle.storage_ids.is_empty() {
491            soft_panic_or_log!(
492                "curated metric sink reads non-introspection relations (name={}): {:?}",
493                definition.name,
494                id_bundle.storage_ids
495            );
496            return Ok(StageResult::Response(ExecuteResponse::CreatedMetricSink));
497        }
498
499        // `reconcile_metric_sinks` only sees sinks already in `metric_sinks`, so a definition
500        // denied while its install was in flight would ship anyway without this recheck.
501        if metric_sink_denied(self.catalog().system_config(), definition.name) {
502            return Ok(StageResult::Response(ExecuteResponse::CreatedMetricSink));
503        }
504
505        // Hold a read on the imports across shipping, so their since cannot advance past the as-of
506        // just picked. Compute takes its own holds during `create_dataflow`.
507        let read_holds = self.acquire_read_holds(&id_bundle);
508        df_desc.set_as_of(read_holds.least_valid_read());
509
510        // Record the install just before shipping. A failed plan or optimize returns earlier, so it
511        // leaves no entry behind. `drop_metric_sinks` reads this entry to release the sink's
512        // instance-global collection state on replica drop. Recording before the ship is safe because
513        // the coordinator runs one message at a time with no await between the two, so no replica drop
514        // sees an entry whose dataflow has not shipped.
515        let install = InstalledMetricSink {
516            cluster_id,
517            sink_id,
518        };
519        if let Some(previous) = self
520            .metric_sinks
521            .insert((replica_id, definition.name), install)
522        {
523            // The key is already taken. `curated_names_are_unique` rules out a name collision,
524            // so this is the same definition installed twice: reconcile can start a second
525            // install while an earlier one is still in flight. Restore the first and abandon this
526            // one, else we leak the first's collection (unreachable to `drop_metric_sinks`) and
527            // register a second collector under the same `sink` label.
528            self.metric_sinks
529                .insert((replica_id, definition.name), previous);
530            info!(
531                %replica_id,
532                name = definition.name,
533                "abandoning metric sink install, already installed"
534            );
535            return Ok(StageResult::Response(ExecuteResponse::CreatedMetricSink));
536        }
537
538        self.ship_dataflow(df_desc, cluster_id, Some(replica_id))
539            .await;
540
541        drop(read_holds);
542        // Nobody is waiting on this: `StagedContext for ()` drops the result. Reuses the
543        // `CREATE METRIC SINK` response rather than adding a variant no client ever sees.
544        Ok(StageResult::Response(ExecuteResponse::CreatedMetricSink))
545    }
546
547    /// Drops the curated metric sinks installed on the given replica.
548    ///
549    /// Called before the replica itself is dropped. Dropping the replica would tear the sink
550    /// dataflows down anyway, but the controller's collection state for them is instance-global,
551    /// so it has to be released explicitly.
552    pub(super) fn drop_metric_sinks(&mut self, replica_id: ReplicaId) {
553        for name in metric_sinks_on_replica(&self.metric_sinks, replica_id) {
554            self.drop_metric_sink(replica_id, name);
555        }
556    }
557
558    /// Drops one curated metric sink, if it is installed.
559    fn drop_metric_sink(&mut self, replica_id: ReplicaId, name: &'static str) {
560        let Some(install) = self.metric_sinks.remove(&(replica_id, name)) else {
561            return;
562        };
563        let InstalledMetricSink {
564            cluster_id,
565            sink_id,
566        } = install;
567        info!(%sink_id, %replica_id, name, "dropping metric sink");
568
569        // The entry exists only for a shipped dataflow, so its collection is present and this
570        // drop succeeds. Result ignored: a failure during replica teardown is not worth a panic.
571        let _ = self
572            .controller
573            .compute
574            .drop_collections(cluster_id, vec![sink_id]);
575    }
576}
577
578/// Whether `disabled_metric_sinks` denies the curated definition called `name`.
579///
580/// An entry naming no definition is never asked about, so a stale or misspelled one is inert.
581fn metric_sink_denied(system_config: &SystemVars, name: &str) -> bool {
582    system_config
583        .disabled_metric_sinks()
584        .iter()
585        .any(|denied| denied == name)
586}
587
588/// The names of the definitions installed on `replica_id`, in key order.
589///
590/// The map is keyed replica-first, so a replica's installs are one contiguous range.
591fn metric_sinks_on_replica(
592    metric_sinks: &BTreeMap<(ReplicaId, &'static str), InstalledMetricSink>,
593    replica_id: ReplicaId,
594) -> Vec<&'static str> {
595    metric_sinks
596        .range((replica_id, "")..)
597        .take_while(|((id, _), _)| *id == replica_id)
598        .map(|((_, name), _)| *name)
599        .collect()
600}
601
602/// Enforces the introspection-only contract from [`CuratedMetricSink::source_sql`]: every relation
603/// the definition reads, walking views transitively, must be a log collection.
604///
605/// Checked here against what the definition reads rather than by import kind after optimization: the
606/// import split (storage vs index) depends on which indexes the target cluster happens to have, so
607/// it gives the same definition different verdicts on different clusters.
608fn ensure_reads_only_logs(
609    catalog: &Catalog,
610    dependencies: &BTreeSet<CatalogItemId>,
611) -> Result<(), anyhow::Error> {
612    let mut to_visit: Vec<_> = dependencies.iter().copied().collect();
613    let mut visited = BTreeSet::new();
614    while let Some(id) = to_visit.pop() {
615        if !visited.insert(id) {
616            continue;
617        }
618        let entry = catalog.get_entry(&id);
619        match entry.item() {
620            // The only data leaf allowed.
621            CatalogItem::Log(_) => {}
622            // Allowed only if everything it reads is, so walk its dependencies.
623            CatalogItem::View(_) => to_visit.extend(entry.uses()),
624            // No data dependency; a view over logs still references these.
625            CatalogItem::Type(_) | CatalogItem::Func(_) => {}
626            _ => bail!(
627                "curated metric sink reads {}, which is not an introspection log relation \
628                 (only logs and views over logs are allowed)",
629                catalog.resolve_full_name(entry.name(), None)
630            ),
631        }
632    }
633    Ok(())
634}
635
636impl CuratedMetricSink {
637    /// Plans `source_sql` against a session-less catalog, returning the query, its output shape,
638    /// and the catalog items it reads.
639    fn plan_source(
640        &self,
641        catalog: &dyn SessionCatalog,
642    ) -> Result<(HirRelationExpr, RelationDesc, BTreeSet<CatalogItemId>), anyhow::Error> {
643        // A definition is a single statement. Reject the count explicitly for a clear error.
644        let statements = mz_sql::parse::parse(self.source_sql)?;
645        if statements.len() != 1 {
646            bail!(
647                "source SQL must be exactly one statement, got {}",
648                statements.len()
649            );
650        }
651
652        // A metric sink's source is a continuously maintained dataflow, like a SUBSCRIBE, so plan it
653        // as one. A maintained lifetime folds any finishing into the expression (an ORDER BY over a
654        // maintained collection is dropped, a LIMIT becomes a TopK) rather than leaving it beside the
655        // query, so `MetricSinkFrom::Query` gets a self-contained expression whose arity matches its
656        // `desc`. This mirrors `coord::introspection`, which plans its specs as subscribes too.
657        let subscribe_sql = format!("SUBSCRIBE ({})", self.source_sql);
658        let parsed = mz_sql::parse::parse(&subscribe_sql)?.into_element();
659        let (stmt, resolved_ids) = mz_sql::names::resolve(catalog, parsed.ast)?;
660        let (plan, sql_impl_ids) =
661            mz_sql::plan::plan(None, catalog, stmt, &Params::empty(), &resolved_ids)?;
662        let Plan::Subscribe(SubscribePlan {
663            from: SubscribeFrom::Query { expr, desc },
664            ..
665        }) = plan
666        else {
667            bail!("source SQL must be a single SELECT");
668        };
669        validate_metric_sink_desc(&desc)?;
670
671        // Fold in ids from SQL-implemented function bodies. `plan` keeps them out of `resolved_ids`
672        // since a one-shot statement doesn't depend on a function's body, but a metric sink inlines
673        // that body into its dataflow, so the body's reads are real imports the gate must check.
674        let dependencies = resolved_ids
675            .items()
676            .chain(sql_impl_ids.items())
677            .copied()
678            .collect();
679        Ok((expr, desc, dependencies))
680    }
681}
682
683impl Staged for MetricSinkStage {
684    type Ctx = ();
685
686    fn validity(&mut self) -> &mut PlanValidity {
687        match self {
688            Self::Optimize(stage) => &mut stage.validity,
689            Self::Finish(stage) => &mut stage.validity,
690        }
691    }
692
693    async fn stage(
694        self,
695        coord: &mut Coordinator,
696        _ctx: &mut (),
697    ) -> Result<StageResult<Box<Self>>, AdapterError> {
698        match self {
699            Self::Optimize(stage) => coord.metric_sink_optimize(stage),
700            Self::Finish(stage) => coord.metric_sink_finish(stage).await,
701        }
702    }
703
704    fn message(self, _ctx: (), span: Span) -> Message {
705        Message::MetricSinkStageReady { span, stage: self }
706    }
707
708    fn cancel_enabled(&self) -> bool {
709        false
710    }
711}
712
713#[cfg(test)]
714mod tests {
715    use std::collections::{BTreeMap, BTreeSet};
716
717    use mz_catalog::memory::objects::CatalogItem;
718    use mz_cluster_client::ReplicaId;
719    use mz_controller_types::ClusterId;
720    use mz_repr::GlobalId;
721    use mz_sql::plan::{
722        METRIC_SINK_CURATED_PREFIX_MARKER, validate_metric_sink_prefix,
723        validate_user_metric_sink_prefix,
724    };
725    use mz_sql::session::vars::{DISABLED_METRIC_SINKS, SystemVars, Var, VarInput};
726
727    use crate::catalog::Catalog;
728    use crate::coord::metric_sink::{
729        CURATED, CuratedMetricSink, InstalledMetricSink, ensure_reads_only_logs,
730        metric_sink_denied, metric_sinks_on_replica,
731    };
732
733    /// `drop_metric_sinks` relies on this range scan returning exactly one replica's installs, with
734    /// no bleed into a neighbouring replica's contiguous range.
735    #[mz_ore::test]
736    fn metric_sinks_on_replica_scans_one_replica() {
737        let cluster = ClusterId::user(1).expect("valid cluster id");
738        let install = |sink_id| InstalledMetricSink {
739            cluster_id: cluster,
740            sink_id: GlobalId::Transient(sink_id),
741        };
742        let r = ReplicaId::User;
743
744        let mut sinks = BTreeMap::new();
745        sinks.insert((r(1), "a"), install(10));
746        sinks.insert((r(2), "a"), install(20));
747        sinks.insert((r(2), "b"), install(21));
748        sinks.insert((r(2), "c"), install(22));
749        sinks.insert((r(4), "a"), install(40));
750
751        // A replica with several installs: all of them, in key order, and nothing from r(1)/r(4).
752        assert_eq!(metric_sinks_on_replica(&sinks, r(2)), vec!["a", "b", "c"]);
753        // First and last replicas in the map: the scan stops at each boundary.
754        assert_eq!(metric_sinks_on_replica(&sinks, r(1)), vec!["a"]);
755        assert_eq!(metric_sinks_on_replica(&sinks, r(4)), vec!["a"]);
756        // A replica with no installs, whether ordered between present ones (the r(3) gap) or past
757        // the end, returns nothing rather than the next replica's range.
758        assert!(metric_sinks_on_replica(&sinks, r(3)).is_empty());
759        assert!(metric_sinks_on_replica(&sinks, r(5)).is_empty());
760    }
761
762    /// The var is a `Vec<Ident>`, so parsing follows the SQL identifier-list rules: surrounding
763    /// whitespace is tolerated, an unquoted name folds to lowercase, and a quoted name keeps its
764    /// case. Matching is otherwise exact (no prefix match).
765    #[mz_ore::test]
766    fn denylist_matches_names_leniently() {
767        let denied = |list: &str, name: &str| {
768            let mut vars = SystemVars::new();
769            vars.set(DISABLED_METRIC_SINKS.name(), VarInput::Flat(list))
770                .expect("valid denylist");
771            metric_sink_denied(&vars, name)
772        };
773
774        assert!(!denied("", "a"));
775        assert!(denied("a", "a"));
776        assert!(denied("a,b", "b"));
777        assert!(denied("  a , b  ", "a"));
778        // An unknown name denies nothing but is carried without error.
779        assert!(!denied("nope", "a"));
780        assert!(denied("nope,a", "a"));
781        // Exact match only: no prefix match.
782        assert!(!denied("a", "ab"));
783        assert!(!denied("ab", "a"));
784        // Unquoted names fold to lowercase; quoting pins the case.
785        assert!(denied("A", "a"));
786        assert!(!denied("\"A\"", "a"));
787
788        // An empty entry between commas is rejected by the identifier parser (unlike the old
789        // naive split, which silently dropped it).
790        let mut vars = SystemVars::new();
791        assert!(
792            vars.set(DISABLED_METRIC_SINKS.name(), VarInput::Flat("a,,b"))
793                .is_err()
794        );
795    }
796
797    #[mz_ore::test]
798    fn curated_prefixes_are_valid() {
799        for definition in CURATED {
800            validate_metric_sink_prefix(definition.prefix).unwrap_or_else(|err| {
801                panic!(
802                    "curated metric sink {:?} has an invalid prefix {:?}: {err}",
803                    definition.name, definition.prefix
804                )
805            });
806        }
807    }
808
809    #[mz_ore::test]
810    fn curated_prefixes_are_reserved_against_user_sinks() {
811        for definition in CURATED {
812            assert!(
813                definition
814                    .prefix
815                    .starts_with(METRIC_SINK_CURATED_PREFIX_MARKER),
816                "curated metric sink {:?} does not use the reserved curated prefix: {:?}",
817                definition.name,
818                definition.prefix
819            );
820            assert!(
821                validate_user_metric_sink_prefix(definition.prefix).is_err(),
822                "a user could claim curated metric sink {:?}'s prefix {:?}",
823                definition.name,
824                definition.prefix
825            );
826        }
827    }
828
829    /// The registry is keyed on the name, so a duplicate would make one definition's install
830    /// unreachable to teardown and both collectors collide on the `sink` label. Guarded at runtime
831    /// (`metric_sink_finish`) too, but caught here at build time before it can ship.
832    #[mz_ore::test]
833    fn curated_names_are_unique() {
834        let mut seen = BTreeSet::new();
835        for definition in CURATED {
836            assert!(
837                seen.insert(definition.name),
838                "duplicate curated metric sink name {:?}",
839                definition.name
840            );
841        }
842    }
843
844    /// Runs the two checks `plan_metric_sink` does at boot, where a failure soft-panics (a hard panic
845    /// under debug assertions, so a CI boot crash-loop). The gate is the load-bearing one: a source
846    /// can plan cleanly yet still read a storage-backed relation transitively (a builtin view joining
847    /// in an introspection source), which only `ensure_reads_only_logs` catches.
848    #[mz_ore::test(tokio::test)]
849    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `TLS_client_method`
850    async fn curated_definitions_plan() {
851        Catalog::with_debug(|catalog| async move {
852            let session_catalog = catalog.for_system_session();
853            for definition in CURATED {
854                let (_, _, dependencies) =
855                    definition
856                        .plan_source(&session_catalog)
857                        .unwrap_or_else(|err| {
858                            panic!(
859                                "curated metric sink {:?} does not plan: {err}",
860                                definition.name
861                            )
862                        });
863                ensure_reads_only_logs(&catalog, &dependencies).unwrap_or_else(|err| {
864                    panic!(
865                        "curated metric sink {:?} reads a non-introspection relation: {err}",
866                        definition.name
867                    )
868                });
869            }
870        })
871        .await
872    }
873
874    /// The five canonical columns, no finishing: the shape a definition must produce.
875    const VALID_SOURCE: &str = "SELECT 'n'::text AS metric_name, 'gauge'::text AS metric_type, \
876        NULL::map[text=>text] AS labels, NULL::double AS value, 'h'::text AS help";
877
878    /// `VALID_SOURCE` with an ORDER BY appended. Maintained-lifetime planning folds it away rather
879    /// than rejecting it, since ordering has no meaning for a continuously-consumed collection.
880    const ORDERED_SOURCE: &str = "SELECT 'n'::text AS metric_name, 'gauge'::text AS metric_type, \
881        NULL::map[text=>text] AS labels, NULL::double AS value, 'h'::text AS help ORDER BY 1";
882
883    /// `plan_source` accepts the canonical column contract (including a source with a finishing,
884    /// which maintained-lifetime planning folds in) and rejects a source missing the columns.
885    #[mz_ore::test(tokio::test)]
886    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `TLS_client_method`
887    async fn plan_source_enforces_the_metric_sink_contract() {
888        Catalog::with_debug(|catalog| async move {
889            let session_catalog = catalog.for_system_session();
890            let plan = |source_sql: &'static str| {
891                CuratedMetricSink {
892                    name: "test",
893                    source_sql,
894                    prefix: "mz_metric_sink_test_",
895                }
896                .plan_source(&session_catalog)
897            };
898
899            assert!(plan(VALID_SOURCE).is_ok());
900
901            // An ORDER BY is folded away by maintained-lifetime planning, not rejected.
902            assert!(plan(ORDERED_SOURCE).is_ok());
903
904            // Missing the canonical columns: rejected by `validate_metric_sink_desc`.
905            assert!(plan("SELECT 1 AS foo").is_err());
906
907            // Not exactly one statement: rejected by the explicit count guard.
908            assert!(plan("").is_err());
909            assert!(plan("SELECT 1; SELECT 2").is_err());
910        })
911        .await
912    }
913
914    /// A SQL-implemented builtin hides its reads: `pg_get_viewdef`'s body reads
915    /// `mz_catalog.mz_views`, which the dataflow imports but the statement's resolved ids omit.
916    /// `plan_source` must surface those reads so the gate rejects them.
917    #[mz_ore::test(tokio::test)]
918    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `TLS_client_method`
919    async fn ensure_reads_only_logs_sees_sql_impl_function_reads() {
920        Catalog::with_debug(|catalog| async move {
921            let session_catalog = catalog.for_system_session();
922            let (_, _, dependencies) = CuratedMetricSink {
923                name: "test",
924                source_sql: "SELECT pg_get_viewdef('x') AS metric_name, 'gauge'::text AS metric_type, \
925                    NULL::map[text=>text] AS labels, NULL::double AS value, 'h'::text AS help",
926                prefix: "mz_metric_sink_test_",
927            }
928            .plan_source(&session_catalog)
929            .expect("plans against the system catalog");
930            assert!(ensure_reads_only_logs(&catalog, &dependencies).is_err());
931        })
932        .await
933    }
934
935    /// The introspection-only contract: a log dependency is accepted, a storage-backed one is
936    /// rejected. Checked against what the definition reads, so the verdict does not depend on the
937    /// target cluster's index layout.
938    #[mz_ore::test(tokio::test)]
939    #[cfg_attr(miri, ignore)] // unsupported operation: can't call foreign function `TLS_client_method`
940    async fn ensure_reads_only_logs_accepts_logs_rejects_storage() {
941        Catalog::with_debug(|catalog| async move {
942            let log_id = catalog
943                .entries()
944                .find(|e| matches!(e.item(), CatalogItem::Log(_)))
945                .expect("debug catalog has a builtin log")
946                .id();
947            assert!(ensure_reads_only_logs(&catalog, &BTreeSet::from([log_id])).is_ok());
948
949            let storage_id = catalog
950                .entries()
951                .find(|e| matches!(e.item(), CatalogItem::Table(_) | CatalogItem::Source(_)))
952                .expect("debug catalog has a builtin table or source")
953                .id();
954            assert!(ensure_reads_only_logs(&catalog, &BTreeSet::from([storage_id])).is_err());
955        })
956        .await
957    }
958}