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}