1use 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 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 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(®istration);
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 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 if let Some(registration) = registration_op.borrow_mut().as_mut() {
143 if let Retry::Arm = registration.try_register(®istry, &worker_metrics) {
144 activator.activate_after(REGISTRATION_RETRY_INTERVAL);
145 }
146 }
147
148 let mut datum_vec = DatumVec::new();
151 move |frontiers| {
152 let mut frontier = Antichain::new();
155 for f in frontiers {
156 frontier.extend(f.frontier().iter().copied());
157 }
158 shared_frontier.borrow_mut().clone_from(&frontier);
165
166 {
169 let mut reg = registration_op.borrow_mut();
170 let Some(registration) = reg.as_mut() else {
171 ok_input.for_each(|_, _| {});
174 err_input.for_each(|_, _| {});
175 return;
176 };
177
178 if let Retry::Arm = registration.try_register(®istry, &worker_metrics) {
179 activator.activate_after(REGISTRATION_RETRY_INTERVAL);
181 }
182 }
183
184 let mut st = state.lock().expect("sink state mutex poisoned");
185
186 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 let collection = compute_state.expect_collection_mut(sink_id);
226 collection.sink_write_frontier = Some(sink_frontier);
227
228 let token: Rc<dyn Any> = registration;
236 Some(token)
237 }
238}
239
240const REGISTRATION_RETRY_INTERVAL: Duration = Duration::from_secs(1);
245
246const REGISTRATION_ESCALATE_AFTER: Duration = Duration::from_secs(60);
251
252#[derive(Debug, Clone, Copy, PartialEq, Eq)]
254enum Retry {
255 Arm,
257 Skip,
259}
260
261struct PendingRegistration {
266 label: String,
268 collector: SinkCollector,
269 handle: Option<Box<dyn Any + Send + Sync>>,
270 logged: bool,
271 terminated: bool,
273 first_collision: Option<Instant>,
275 escalated: bool,
277 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 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 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
359struct 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
394fn 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
436type RowKey = (
451 String,
452 Vec<(String, Option<String>)>,
453 Option<u64>,
454 Option<MetricKind>,
455 bool,
456 String,
457);
458
459type PublishedKey = (String, Vec<(String, String)>);
461type PublishedValue = (f64, MetricKind, String);
463
464#[derive(Default)]
475struct SinkState {
476 pending_ok: BTreeMap<Timestamp, BTreeMap<RowKey, i64>>,
478 pending_err: BTreeMap<Timestamp, i64>,
480 working: BTreeMap<RowKey, i64>,
482 published: BTreeMap<PublishedKey, PublishedValue>,
483 errors: i64,
486 frontier_ms: u64,
487 skipped: u64,
488 conflicts: u64,
489 collisions: u64,
490 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 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
521fn 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
536fn is_publishable_label_value(value: Option<&str>) -> bool {
540 matches!(value, Some(v) if !v.is_empty())
541}
542
543impl SinkState {
544 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 fn stage_err(&mut self, time: Timestamp, diff: i64) {
582 *self.pending_err.entry(time).or_default() += diff;
583 }
584
585 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 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
639fn 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
659fn rebuild_published(
675 working: &BTreeMap<RowKey, i64>,
676) -> (BTreeMap<PublishedKey, PublishedValue>, u64, u64) {
677 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 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
746fn 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
765fn 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#[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 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 const LABEL_A: &[(&str, Option<&str>)] = &[("a", Some("1"))];
934
935 fn key_m() -> PublishedKey {
936 ("m".into(), label_a())
937 }
938
939 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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); packer.push_dict_with(|_| {}); packer.push(Datum::Null); packer.push(Datum::String("")); packer.push(Datum::Null); packer.push(Datum::Null); }
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 #[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 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 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 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 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 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 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 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 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 assert_eq!(st.published[&key_m()].0, 5.0);
1317
1318 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 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 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 #[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 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(®istry, ComputeRuntimeRole::Solo).for_worker(0);
1375
1376 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(®istry, 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(®istry, &metrics), Retry::Arm);
1392 assert_eq!(counter_total(®istry, RETRIES), 1.0);
1393 assert_eq!(
1395 companion_gauge_count(®istry, FRONTIER, SINK_LABEL, LABEL),
1396 1
1397 );
1398
1399 assert_eq!(registration.try_register(®istry, &metrics), Retry::Skip);
1401 assert_eq!(counter_total(®istry, RETRIES), 1.0);
1402
1403 registration.next_attempt = None;
1405 assert_eq!(registration.try_register(®istry, &metrics), Retry::Arm);
1406 assert_eq!(counter_total(®istry, RETRIES), 2.0);
1407
1408 drop(incumbent);
1409 registration.next_attempt = None;
1410 assert_eq!(registration.try_register(®istry, &metrics), Retry::Skip);
1411 assert_eq!(counter_total(®istry, RETRIES), 2.0);
1412 assert_eq!(
1413 companion_gauge_count(®istry, FRONTIER, SINK_LABEL, LABEL),
1414 1
1415 );
1416
1417 assert_eq!(registration.try_register(®istry, &metrics), Retry::Skip);
1419 assert_eq!(counter_total(®istry, RETRIES), 2.0);
1420
1421 drop(registration);
1423 assert_eq!(
1424 companion_gauge_count(®istry, FRONTIER, SINK_LABEL, LABEL),
1425 0
1426 );
1427 }
1428
1429 #[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(®istry, 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 let registry_op = registry.clone();
1463 let metrics_op = metrics.clone();
1464 let registration = worker.dataflow::<Timestamp, _, _>(move |scope| {
1465 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(®istration);
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(®istry_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 assert_eq!(
1495 companion_gauge_count(®istry, FRONTIER, SINK_LABEL, LABEL),
1496 1
1497 );
1498
1499 for _ in 0..100 {
1501 if Rc::strong_count(®istration) == 1 {
1502 break;
1503 }
1504 worker.step();
1505 }
1506 assert_eq!(
1507 Rc::strong_count(®istration),
1508 1,
1509 "the operator closure should have dropped once its input closed"
1510 );
1511
1512 assert_eq!(
1514 companion_gauge_count(®istry, FRONTIER, SINK_LABEL, LABEL),
1515 1
1516 );
1517
1518 drop(registration);
1520 assert_eq!(
1521 companion_gauge_count(®istry, FRONTIER, SINK_LABEL, LABEL),
1522 0
1523 );
1524 }
1525
1526 #[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 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}