1use std::collections::BTreeMap;
31use std::sync::Arc;
32use std::time::{Duration, Instant};
33
34use itertools::Itertools;
35use mz_adapter_types::dyncfgs::{
36 FRONTEND_READ_THEN_WRITE, HYDRATION_HISTORY_COLLECTION_INTERVAL,
37 HYDRATION_HISTORY_RETENTION_PERIOD, REPLICA_HYDRATION_HISTORY_RETENTION_PERIOD,
38};
39use mz_catalog::builtin::{
40 MZ_CATALOG_SERVER_CLUSTER, MZ_OBJECT_HYDRATION_HISTORY, MZ_REPLICA_HYDRATION_HISTORY,
41};
42use mz_cluster_client::ReplicaId;
43use mz_controller::clusters::{ClusterStatus, ReplicaLocation};
44use mz_controller_types::ClusterId;
45use mz_ore::cast::CastFrom;
46use mz_ore::collections::CollectionExt;
47use mz_ore::now::EpochMillis;
48use mz_ore::task;
49use mz_repr::CatalogItemId;
50use mz_sql::plan::{MutationKind, Params, Plan, ReadThenWritePlan};
51use sha2::{Digest, Sha256};
52use tracing::warn;
53
54use crate::catalog::Catalog;
55use crate::command::ExecuteResponse;
56use crate::coord::{Coordinator, Message};
57use crate::metrics::Metrics;
58use crate::peek_client::CoordinatorClient;
59use crate::session::Session;
60use crate::{AdapterError, PeekClient};
61
62const SCHEDULE_RECHECK_CAP: Duration = Duration::from_secs(5);
68
69const DISABLED_RECHECK_INTERVAL: Duration = Duration::from_secs(60);
74
75const MUTATION_TIMEOUT: Duration = Duration::from_secs(300);
81
82const RETENTION_BATCH_SIZE: usize = 1000;
90
91fn next_fire_delay(now: EpochMillis, interval_ms: EpochMillis, offset: EpochMillis) -> Duration {
97 debug_assert!(interval_ms > 0);
98 let this_period = (now - (now % interval_ms)).saturating_add(offset);
99 let next = if this_period > now {
100 this_period
101 } else {
102 this_period.saturating_add(interval_ms)
103 };
104 Duration::from_millis(next.saturating_sub(now))
105}
106
107fn environment_schedule_offset(environment_id: &str, interval_ms: EpochMillis) -> EpochMillis {
109 debug_assert!(interval_ms > 0);
110 let digest = Sha256::digest(environment_id);
111 let hash = u64::from_le_bytes(digest[..8].try_into().expect("SHA-256 digest has 32 bytes"));
112 hash % interval_ms
113}
114
115impl Coordinator {
116 pub(super) fn schedule_hydration_history_collection(&self) {
129 let interval =
130 HYDRATION_HISTORY_COLLECTION_INTERVAL.get(self.catalog().system_config().dyncfgs());
131
132 let (delay, fire) = if interval.is_zero() {
135 (DISABLED_RECHECK_INTERVAL, false)
136 } else {
137 let interval_ms = EpochMillis::try_from(interval.as_millis())
140 .unwrap_or(EpochMillis::MAX)
141 .max(1);
142 let remaining = next_fire_delay(
143 self.now(),
144 interval_ms,
145 self.hydration_history_schedule_offset(interval_ms),
146 );
147 if remaining <= SCHEDULE_RECHECK_CAP {
148 (remaining, true)
149 } else {
150 (SCHEDULE_RECHECK_CAP, false)
151 }
152 };
153
154 let internal_cmd_tx = self.internal_cmd_tx.clone();
155 task::spawn(|| "hydration_history_schedule", async move {
156 tokio::time::sleep(delay).await;
157 let message = if fire {
158 Message::HydrationHistoryRun
159 } else {
160 Message::HydrationHistorySchedule
161 };
162 let _ = internal_cmd_tx.send(message);
164 });
165 }
166
167 fn hydration_history_schedule_offset(&self, interval_ms: EpochMillis) -> EpochMillis {
175 let environment_id = self.catalog().state().config().environment_id.to_string();
176 environment_schedule_offset(&environment_id, interval_ms)
177 }
178
179 pub(super) fn run_hydration_history_collection(&mut self) {
181 let collection_interval = {
182 let dyncfgs = self.catalog().system_config().dyncfgs();
183 HYDRATION_HISTORY_COLLECTION_INTERVAL.get(dyncfgs)
184 };
185 if collection_interval.is_zero() || self.controller.read_only() {
190 self.schedule_hydration_history_collection();
191 return;
192 }
193
194 let replicas = self
195 .catalog()
196 .clusters()
197 .flat_map(|cluster| cluster.replicas())
198 .filter(|replica| replica.config.compute.logging.enabled())
199 .filter(|replica| match &replica.config.location {
200 ReplicaLocation::Managed(_) => {
201 self.cluster_replica_statuses
202 .get_cluster_replica_status(replica.cluster_id, replica.replica_id)
203 == ClusterStatus::Online
204 }
205 ReplicaLocation::Unmanaged(_) => true,
208 })
209 .map(|replica| ReplicaTarget {
210 cluster_id: replica.cluster_id,
211 replica_id: replica.replica_id,
212 process_count: replica.config.location.num_processes(),
213 })
214 .sorted_by_key(|replica| replica.replica_id)
215 .collect_vec();
216
217 let catalog = self.owned_catalog();
218 let catalog_server = catalog.resolve_builtin_cluster(&MZ_CATALOG_SERVER_CLUSTER);
222 let catalog_server_target = catalog_server
223 .replicas()
224 .next()
225 .map(|replica| (catalog_server.id, replica.replica_id));
226
227 let replica = next_replica(&replicas, self.hydration_history_replica_cursor);
228 if let Some(replica) = replica {
229 self.hydration_history_replica_cursor = Some(replica.replica_id);
230 }
231 let mut sweep = self.new_sweep(catalog);
232 let internal_cmd_tx = self.internal_cmd_tx.clone();
233
234 let handle = task::spawn(|| "hydration_history_sweep", async move {
235 let started = Instant::now();
236 if let Some(replica) = replica {
237 sweep.collect(replica).await;
238 }
239
240 if let Some((cluster_id, replica_id)) = catalog_server_target {
244 sweep.retain(cluster_id, replica_id).await;
245 }
246
247 sweep
248 .metrics
249 .hydration_history_sweep_duration_seconds
250 .observe(started.elapsed().as_secs_f64());
251 let _ = internal_cmd_tx.send(Message::HydrationHistorySchedule);
252 });
253
254 self.hydration_history_sweep = Some(handle.abort_on_drop());
257 }
258
259 fn new_sweep(&self, catalog: Arc<Catalog>) -> Sweep {
261 let now = self.now();
262 let cutoff = |retention: Duration| {
263 let retention_ms = u64::try_from(retention.as_millis()).unwrap_or(u64::MAX);
264 mz_ore::now::to_datetime(now.saturating_sub(retention_ms)).to_rfc3339()
265 };
266 let dyncfgs = catalog.system_config().dyncfgs();
267 let object_cutoff = cutoff(HYDRATION_HISTORY_RETENTION_PERIOD.get(dyncfgs));
268 let replica_cutoff = cutoff(REPLICA_HYDRATION_HISTORY_RETENTION_PERIOD.get(dyncfgs));
269 let build_version = catalog.state().config().build_info.human_version(None);
270 let client = PeekClient::new(
274 CoordinatorClient::Background {
275 tx: self.internal_cmd_tx.clone(),
276 metrics: self.metrics.clone(),
277 },
278 &catalog,
279 Arc::clone(&self.controller.storage_collections),
280 Arc::clone(&self.transient_id_gen),
281 self.optimizer_metrics.clone(),
282 self.persist_client.clone(),
283 self.statement_logging.create_frontend(build_version),
284 Arc::clone(&self.occ_write_semaphore),
285 FRONTEND_READ_THEN_WRITE.get(self.catalog().system_config().dyncfgs()),
286 self.group_commit_tx.clone(),
287 self.controller.read_only(),
288 );
289 Sweep {
290 client,
291 object_history_id: catalog.resolve_builtin_table(&MZ_OBJECT_HYDRATION_HISTORY),
292 replica_history_id: catalog.resolve_builtin_table(&MZ_REPLICA_HYDRATION_HISTORY),
293 catalog,
294 metrics: self.metrics.clone(),
295 wall_time: self.now_datetime(),
296 object_cutoff,
297 replica_cutoff,
298 }
299 }
300}
301
302#[derive(Clone, Copy, Debug, Eq, PartialEq)]
304struct ReplicaTarget {
305 cluster_id: ClusterId,
306 replica_id: ReplicaId,
307 process_count: usize,
308}
309
310fn next_replica(replicas: &[ReplicaTarget], cursor: Option<ReplicaId>) -> Option<ReplicaTarget> {
316 replicas
317 .iter()
318 .find(|replica| cursor.is_none_or(|cursor| replica.replica_id > cursor))
319 .or_else(|| replicas.first())
320 .copied()
321}
322
323fn object_collection_sql(cluster_id: ClusterId, replica_id: ReplicaId, cutoff: &str) -> String {
365 format!(
380 "SELECT
381 e.object_id,
382 '{cluster_id}'::text AS cluster_id,
383 '{replica_id}'::text AS replica_id,
384 e.installed_at,
385 e.started_at,
386 e.hydrated_at,
387 'hydrated'::text AS status
388 FROM (
389 SELECT
390 t.export_id AS object_id,
391 min(t.installed_at) AS installed_at,
392 min(t.started_at) AS started_at,
393 max(t.hydrated_at) AS hydrated_at
394 FROM mz_introspection.mz_compute_hydration_times_per_worker AS t
395 WHERE t.export_id NOT LIKE 'si%'
396 AND t.export_id NOT LIKE 't%'
397 GROUP BY t.export_id
398 HAVING count(*) = count(t.hydrated_at)
399 ) AS e
400 WHERE e.hydrated_at >= TIMESTAMPTZ '{cutoff}'
401 AND NOT EXISTS (
402 SELECT 1
403 FROM mz_internal.mz_object_hydration_history AS h
404 WHERE h.object_id = e.object_id
405 AND h.replica_id = '{replica_id}'::text
406 AND h.installed_at = e.installed_at
407 )"
408 )
409}
410
411fn replica_collection_sql(target: ReplicaTarget, cutoff: &str) -> String {
429 let ReplicaTarget {
430 cluster_id,
431 replica_id,
432 process_count,
433 } = target;
434 format!(
437 "WITH
438 -- One hydration interval per compute export: earliest install and
439 -- latest finish across its workers. Hydrated only once every worker
440 -- visible at this timestamp has finished.
441 objects AS (
442 SELECT
443 t.export_id AS object_id,
444 min(t.installed_at) AS installed_at,
445 max(t.hydrated_at) AS hydrated_at,
446 count(*) = count(t.hydrated_at) AS hydrated
447 FROM mz_introspection.mz_compute_hydration_times_per_worker AS t
448 WHERE t.export_id NOT LIKE 't%'
449 GROUP BY t.export_id
450 ),
451 -- Completed intervals in install order, each with the coverage
452 -- horizon: the latest finish among this and all earlier intervals.
453 covered AS (
454 SELECT
455 object_id,
456 installed_at,
457 hydrated_at,
458 max(hydrated_at) OVER (
459 ORDER BY installed_at, object_id
460 ROWS UNBOUNDED PRECEDING
461 ) AS covered_through
462 FROM objects
463 WHERE hydrated
464 ),
465 -- An interval starts a new episode when the horizon just before it
466 -- does not reach its install: for a moment, nothing was hydrating.
467 flagged AS (
468 SELECT
469 object_id,
470 installed_at,
471 hydrated_at,
472 lag(covered_through) OVER (
473 ORDER BY installed_at, object_id
474 ) IS NULL
475 OR lag(covered_through) OVER (
476 ORDER BY installed_at, object_id
477 ) < installed_at AS starts_episode
478 FROM covered
479 ),
480 -- Each interval belongs to the latest episode start at or before it.
481 labeled AS (
482 SELECT
483 object_id,
484 installed_at,
485 hydrated_at,
486 max(CASE WHEN starts_episode THEN installed_at END) OVER (
487 ORDER BY installed_at, object_id
488 ROWS UNBOUNDED PRECEDING
489 ) AS episode_started_at
490 FROM flagged
491 ),
492 -- One row per completed episode.
493 episodes AS (
494 SELECT
495 episode_started_at AS started_at,
496 max(hydrated_at) AS finished_at,
497 count(*) FILTER (WHERE object_id NOT LIKE 'si%')::uint8 AS object_count
498 FROM labeled
499 GROUP BY episode_started_at
500 ),
501 -- The earliest install of an export that has not hydrated yet.
502 open_min AS (
503 SELECT min(installed_at) AS v FROM objects WHERE NOT hydrated
504 ),
505 -- The episode to record: the latest one that finished before any
506 -- unhydrated export was installed. An episode finishing at or after
507 -- open_min contains that open interval and is still in progress.
508 -- Comparing against this one scalar, instead of joining episodes
509 -- with open intervals, avoids a cross product that is quadratic when
510 -- many episodes coexist with many still-hydrating exports.
511 episode AS (
512 SELECT e.started_at, e.finished_at, e.object_count
513 FROM episodes AS e, open_min AS o
514 WHERE o.v IS NULL OR e.finished_at < o.v
515 ORDER BY e.started_at DESC
516 LIMIT 1
517 ),
518 -- Process-lifetime resource high-water marks for each process.
519 resources AS (
520 SELECT
521 process_id,
522 max(value) FILTER (
523 WHERE source = 'cgroup' AND metric = 'memory_peak'
524 ) AS peak_memory_bytes,
525 coalesce(
526 max(value) FILTER (
527 WHERE source = 'statvfs' AND metric = 'fs_used_peak'
528 ),
529 max(value) FILTER (
530 WHERE source = 'cgroup' AND metric = 'swap_peak'
531 )
532 ) AS peak_disk_bytes
533 FROM mz_introspection.mz_cluster_replica_resource_usage
534 GROUP BY process_id
535 ),
536 -- The history rows to write, held back until every configured process
537 -- has reported resource usage and dropped once the episode has aged
538 -- past the retention cutoff.
539 candidate AS (
540 SELECT
541 '{replica_id}'::text AS replica_id,
542 '{cluster_id}'::text AS cluster_id,
543 e.started_at,
544 e.finished_at,
545 e.object_count,
546 r.peak_memory_bytes,
547 r.peak_disk_bytes,
548 'hydrated'::text AS status,
549 r.process_id
550 FROM episode AS e
551 CROSS JOIN resources AS r
552 WHERE (SELECT count(*) FROM resources) = {process_count}::uint8
553 AND e.finished_at >= TIMESTAMPTZ '{cutoff}'
554 )
555 -- Skip episodes the history already covers: a recorded row finishing
556 -- at or after this start is this episode, or overlaps it under
557 -- cross-process clock skew.
558 SELECT c.*
559 FROM candidate AS c
560 WHERE NOT EXISTS (
561 SELECT 1
562 FROM mz_internal.mz_replica_hydration_history AS h
563 WHERE h.replica_id = c.replica_id
564 AND h.finished_at >= c.started_at
565 )"
566 )
567}
568
569fn object_retention_sql(cutoff: &str) -> String {
575 format!(
579 "SELECT * FROM (
580 SELECT
581 object_id, cluster_id, replica_id, installed_at, started_at,
582 hydrated_at, status
583 FROM mz_internal.mz_object_hydration_history
584 WHERE hydrated_at < TIMESTAMPTZ '{cutoff}'
585 ORDER BY hydrated_at
586 LIMIT {RETENTION_BATCH_SIZE}
587 )"
588 )
589}
590
591fn replica_retention_sql(cutoff: &str) -> String {
593 format!(
594 "SELECT * FROM (
595 SELECT
596 replica_id, cluster_id, started_at, finished_at, object_count,
597 peak_memory_bytes, peak_disk_bytes, status, process_id
598 FROM mz_internal.mz_replica_hydration_history
599 WHERE finished_at < TIMESTAMPTZ '{cutoff}'
600 ORDER BY finished_at
601 LIMIT {RETENTION_BATCH_SIZE}
602 )"
603 )
604}
605
606struct Sweep {
608 client: PeekClient,
609 catalog: Arc<Catalog>,
610 object_history_id: CatalogItemId,
611 replica_history_id: CatalogItemId,
612 metrics: Metrics,
613 wall_time: chrono::DateTime<chrono::Utc>,
614 object_cutoff: String,
619 replica_cutoff: String,
620}
621
622impl Sweep {
623 async fn collect(&mut self, target: ReplicaTarget) {
625 let ReplicaTarget {
626 cluster_id,
627 replica_id,
628 ..
629 } = target;
630 let sql = object_collection_sql(cluster_id, replica_id, &self.object_cutoff);
631 let _ = self
632 .run(
633 "collection",
634 self.object_history_id,
635 cluster_id,
636 replica_id,
637 MutationKind::Insert,
638 &sql,
639 )
640 .await;
641
642 let sql = replica_collection_sql(target, &self.replica_cutoff);
643 let _ = self
644 .run(
645 "replica_collection",
646 self.replica_history_id,
647 cluster_id,
648 replica_id,
649 MutationKind::Insert,
650 &sql,
651 )
652 .await;
653 }
654
655 async fn retain(&mut self, cluster_id: ClusterId, replica_id: ReplicaId) {
657 let sql = object_retention_sql(&self.object_cutoff);
658 if let Some(deleted) = self
659 .run(
660 "retention",
661 self.object_history_id,
662 cluster_id,
663 replica_id,
664 MutationKind::Delete,
665 &sql,
666 )
667 .await
668 && deleted == RETENTION_BATCH_SIZE
669 {
670 self.metrics.hydration_history_retention_batch_full.inc();
671 }
672
673 let sql = replica_retention_sql(&self.replica_cutoff);
674 if let Some(deleted) = self
675 .run(
676 "replica_retention",
677 self.replica_history_id,
678 cluster_id,
679 replica_id,
680 MutationKind::Delete,
681 &sql,
682 )
683 .await
684 && deleted == RETENTION_BATCH_SIZE
685 {
686 self.metrics.hydration_history_retention_batch_full.inc();
687 }
688 }
689
690 async fn run(
701 &mut self,
702 step: &'static str,
703 history_id: CatalogItemId,
704 cluster_id: ClusterId,
705 replica_id: ReplicaId,
706 kind: MutationKind,
707 sql: &str,
708 ) -> Option<usize> {
709 let mutation = async {
710 let plan = plan_mutation(&self.catalog, history_id, kind, sql)?;
711 let mut session = Session::dummy();
712 session.start_transaction_single_stmt(self.wall_time);
713 let response = self
714 .client
715 .background_read_then_write(
716 &mut session,
717 plan,
718 cluster_id,
719 replica_id,
720 &self.catalog,
721 )
722 .await?;
723 Ok::<_, AdapterError>(response)
724 };
725 match tokio::time::timeout(MUTATION_TIMEOUT, mutation).await {
726 Ok(Ok(response)) => {
727 let (rows, action) = match (kind, response) {
728 (MutationKind::Insert, ExecuteResponse::Inserted(rows)) => (rows, "appended"),
729 (MutationKind::Update, ExecuteResponse::Updated(rows)) => (rows, "updated"),
730 (MutationKind::Delete, ExecuteResponse::Deleted(rows)) => (rows, "deleted"),
731 (_, response) => {
732 self.observe_mutation(step, "error");
733 mz_ore::soft_panic_or_log!(
734 "hydration history {step} returned an unexpected response: {response:?}"
735 );
736 return None;
737 }
738 };
739 self.metrics
740 .hydration_history_rows_affected
741 .with_label_values(&[action])
742 .inc_by(u64::cast_from(rows));
743 let outcome = if rows == 0 { "noop" } else { "success" };
744 self.observe_mutation(step, outcome);
745 Some(rows)
746 }
747 Ok(Err(error)) => {
748 self.observe_mutation(step, "error");
749 if step.ends_with("collection")
750 && matches!(&error, AdapterError::ReadThenWriteContention)
751 {
752 warn!(
753 %step, %cluster_id, %replica_id, %error,
754 "hydration history step failed, the replica's introspection frontier \
755 may be trailing the write frontier"
756 );
757 } else {
758 warn!(%step, %cluster_id, %replica_id, %error, "hydration history step failed");
759 }
760 None
761 }
762 Err(_) if step.ends_with("collection") => {
766 self.observe_mutation(step, "timeout");
767 warn!(
768 %step, %cluster_id, %replica_id,
769 "hydration history step timed out, \
770 the replica's introspection frontier may be trailing the write frontier"
771 );
772 None
773 }
774 Err(_) => {
775 self.observe_mutation(step, "timeout");
776 warn!(%step, %cluster_id, %replica_id, "hydration history step timed out");
777 None
778 }
779 }
780 }
781
782 fn observe_mutation(&self, operation: &str, outcome: &str) {
783 self.metrics
784 .hydration_history_mutations
785 .with_label_values(&[operation, outcome])
786 .inc();
787 }
788}
789
790fn plan_mutation(
800 catalog: &Arc<Catalog>,
801 target_id: CatalogItemId,
802 kind: MutationKind,
803 sql: &str,
804) -> Result<ReadThenWritePlan, AdapterError> {
805 let session_catalog = catalog.for_system_session();
806 let parsed = mz_sql::parse::parse(sql)
807 .map_err(AdapterError::from)?
808 .into_element();
809 let (stmt, resolved_ids) = mz_sql::names::resolve(&session_catalog, parsed.ast)?;
810 let (plan, _) = mz_sql::plan::plan(
811 None,
812 &session_catalog,
813 stmt,
814 &Params::empty(),
815 &resolved_ids,
816 )?;
817 let Plan::Select(select) = plan else {
818 return Err(AdapterError::Internal(
819 "hydration history query did not plan as SELECT".into(),
820 ));
821 };
822
823 let target_desc = catalog
824 .get_entry(&target_id)
825 .relation_desc_latest()
826 .expect("hydration history target is a table");
827 let selection_types = select.source.typ(&[], &BTreeMap::new()).column_types;
828 let target_types = &target_desc.typ().column_types;
829 let matches = selection_types.len() == target_types.len()
830 && selection_types
831 .iter()
832 .zip_eq(target_types)
833 .all(|(selected, target)| selected.scalar_type == target.scalar_type);
836 if !matches {
837 return Err(AdapterError::Internal(format!(
838 "hydration history query does not match the target table: \
839 selection {selection_types:?}, table {target_types:?}"
840 )));
841 }
842
843 Ok(ReadThenWritePlan {
844 id: target_id,
845 selection: select.source,
846 finishing: select.finishing,
847 assignments: BTreeMap::new(),
848 kind,
849 returning: Vec::new(),
850 })
851}
852
853#[cfg(test)]
854mod tests {
855 use super::*;
856
857 #[mz_ore::test]
858 fn replica_sweep_advances_and_wraps() {
859 let cluster = ClusterId::user(1).expect("valid cluster ID");
860 let replicas = [
861 ReplicaTarget {
862 cluster_id: cluster,
863 replica_id: ReplicaId::User(1),
864 process_count: 1,
865 },
866 ReplicaTarget {
867 cluster_id: cluster,
868 replica_id: ReplicaId::User(3),
869 process_count: 1,
870 },
871 ];
872
873 assert_eq!(next_replica(&replicas, None), Some(replicas[0]));
874 assert_eq!(
875 next_replica(&replicas, Some(ReplicaId::User(1))),
876 Some(replicas[1])
877 );
878 assert_eq!(
879 next_replica(&replicas, Some(ReplicaId::User(3))),
880 Some(replicas[0])
881 );
882 assert_eq!(next_replica(&[], None), None);
883 }
884
885 #[mz_ore::test]
888 fn fire_delay_is_offset_within_the_interval() {
889 let interval = 60_000;
890
891 assert_eq!(
893 next_fire_delay(1_000, interval, 5_000),
894 Duration::from_millis(4_000)
895 );
896 assert_eq!(
898 next_fire_delay(5_000, interval, 5_000),
899 Duration::from_millis(interval)
900 );
901 assert_eq!(
903 next_fire_delay(6_000, interval, 5_000),
904 Duration::from_millis(59_000)
905 );
906 assert_eq!(
908 next_fire_delay(59_999, interval, 0),
909 Duration::from_millis(1)
910 );
911 assert_eq!(
912 next_fire_delay(60_000, interval, 0),
913 Duration::from_millis(interval)
914 );
915
916 let one = environment_schedule_offset(
918 "aws-us-east-1-00000000-0000-0000-0000-000000000000-0",
919 interval,
920 );
921 let two = environment_schedule_offset(
922 "aws-us-west-1-00000000-0000-0000-0000-000000000000-1",
923 interval,
924 );
925 assert_eq!(one, 30_189);
926 assert_eq!(two, 38_252);
927 assert_ne!(one, two);
928 }
929
930 #[mz_ore::test]
935 fn collect_requires_every_worker() {
936 let cutoff = "1970-01-01T00:00:00+00:00";
937 let sql = object_collection_sql(
938 ClusterId::user(1).expect("valid cluster ID"),
939 ReplicaId::User(2),
940 cutoff,
941 );
942 assert!(
943 sql.contains("HAVING count(*) = count(t.hydrated_at)"),
944 "{sql}"
945 );
946 assert!(sql.contains("max(t.hydrated_at)"), "{sql}");
947 assert!(!sql.contains("worker_id"), "{sql}");
948
949 let aggregate_end = sql.find(") AS e").expect("aggregate subquery");
952 assert!(sql.find(cutoff).expect("cutoff") > aggregate_end, "{sql}");
953 assert!(
954 sql.find("NOT EXISTS").expect("anti-join") > aggregate_end,
955 "{sql}"
956 );
957 }
958
959 #[mz_ore::test]
966 fn replica_collection_uses_latest_completed_interval_island() {
967 let sql = replica_collection_sql(
968 ReplicaTarget {
969 cluster_id: ClusterId::user(1).expect("valid cluster ID"),
970 replica_id: ReplicaId::User(2),
971 process_count: 3,
972 },
973 "1970-01-01T00:00:00+00:00",
974 );
975 let normalized_sql = sql.split_whitespace().collect::<Vec<_>>().join(" ");
976
977 assert!(sql.contains("ROWS UNBOUNDED PRECEDING"), "{sql}");
981 assert!(sql.contains("lag(covered_through)"), "{sql}");
982 assert!(
983 sql.contains("CASE WHEN starts_episode THEN installed_at END"),
984 "{sql}"
985 );
986 assert!(sql.contains("GROUP BY episode_started_at"), "{sql}");
987 assert!(!sql.contains("bool_and(hydrated)"), "{sql}");
992 assert!(
996 normalized_sql.contains("WHERE o.v IS NULL OR e.finished_at < o.v"),
997 "{sql}"
998 );
999 assert!(!sql.contains("o.installed_at <= e.finished_at"), "{sql}");
1000 assert!(
1001 normalized_sql.contains("ORDER BY e.started_at DESC LIMIT 1"),
1002 "{sql}"
1003 );
1004 assert!(
1005 sql.contains("(SELECT count(*) FROM resources) = 3::uint8"),
1006 "{sql}"
1007 );
1008 assert!(sql.contains("WHERE t.export_id NOT LIKE 't%'"), "{sql}");
1009 assert!(!sql.contains("WHERE t.export_id LIKE 'u%'"), "{sql}");
1010 assert!(!sql.contains("mz_object_global_ids"), "{sql}");
1011 assert!(!sql.contains("mz_catalog.mz_objects"), "{sql}");
1012 assert!(
1013 normalized_sql
1014 .contains("max(value) FILTER ( WHERE source = 'cgroup' AND metric = 'memory_peak'"),
1015 "{sql}"
1016 );
1017 assert!(
1018 normalized_sql.contains(
1019 "max(value) FILTER ( WHERE source = 'statvfs' AND metric = 'fs_used_peak'"
1020 ),
1021 "{sql}"
1022 );
1023 assert!(
1024 normalized_sql
1025 .contains("max(value) FILTER ( WHERE source = 'cgroup' AND metric = 'swap_peak'"),
1026 "{sql}"
1027 );
1028 assert!(!normalized_sql.contains("sum(value)"), "{sql}");
1029 assert!(
1030 sql.contains("FROM mz_internal.mz_replica_hydration_history"),
1031 "{sql}"
1032 );
1033 assert!(sql.contains("h.finished_at >= c.started_at"), "{sql}");
1034 }
1035}