Skip to main content

mz_adapter/coord/
hydration_history.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//! Durable history collection for completed object and replica hydration episodes.
11//!
12//! One sweep visits a single replica, installs a replica-targeted
13//! subscribe that diffs that replica's live hydration timestamps against the
14//! durable history tables, and appends what is missing through the timestamped
15//! OCC write path. Including each history table in its read expression is what
16//! makes the write idempotent across concurrent `environmentd` processes: two
17//! collectors that compute the same row race for one write timestamp, and the
18//! loser observes the winner's append through its own subscribe and finds
19//! nothing left to write.
20//!
21//! One replica is sampled per interval, so an environment with `N` eligible
22//! replicas revisits each one approximately every `N * interval`. Lowering the
23//! interval improves freshness at the cost of more replica dataflow installs.
24//!
25//! Collection is sampling, not an event log. Replica history records only the
26//! latest completed episode visible in a sweep. Intermediate episodes and
27//! intervals retracted before collection leave no evidence and are not recorded.
28//! See the design doc for the resulting semantics.
29
30use 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
62/// Longest a scheduler sleep may run before it rechecks the configuration.
63///
64/// Sleeping the whole interval would leave a dyncfg change ineffective until the
65/// old interval elapsed, so lowering the interval at runtime (which tests do)
66/// would not take effect for up to the previous interval.
67const SCHEDULE_RECHECK_CAP: Duration = Duration::from_secs(5);
68
69/// How often a disabled collector rechecks whether it was enabled.
70///
71/// A disabled collector has no pending collection work, so this can be much
72/// coarser than the enabled scheduler's recheck cadence.
73const DISABLED_RECHECK_INTERVAL: Duration = Duration::from_secs(60);
74
75/// Bound on one replica-targeted mutation.
76///
77/// Subscribe installation, replica-side progress, OCC conflict retries, and the
78/// external commit can all wait indefinitely. Exceeding the bound skips the
79/// step. The next sweep recomputes from current state.
80const MUTATION_TIMEOUT: Duration = Duration::from_secs(300);
81
82/// Rows retracted per retention step.
83///
84/// Retention has to be bounded: the OCC path refuses a selection larger than
85/// `max_result_size` before submitting any write, so an unbounded delete over a
86/// large backlog would fail identically on every sweep and never shrink the
87/// table. Retention repeats bounded batches across sweeps until it drains the
88/// fixed-cutoff backlog.
89const RETENTION_BATCH_SIZE: usize = 1000;
90
91/// Milliseconds until the next fire on this environment's own grid.
92///
93/// The grid has period `interval_ms` and is shifted by `offset`. When this
94/// environment's point in the current period has already passed, the next one is a
95/// full period later.
96fn 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
107/// Stable offset within `interval_ms` for one environment id.
108fn 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    /// Schedules the next hydration history sweep.
117    ///
118    /// Fires are aligned to interval boundaries so that they stay evenly spaced
119    /// across restarts, offset per environment so that a fleet-wide interval does
120    /// not make every environment sweep at the same instant, and each sleep is
121    /// capped so a configuration change is picked up promptly. Sweeps never
122    /// overlap: the next one is only scheduled once the previous one has finished
123    /// or failed.
124    ///
125    /// NOTE: Alignment reads the wall clock, so a test that freezes `NowFn` and
126    /// configures an interval longer than the recheck cap never reaches a
127    /// boundary and never fires.
128    pub(super) fn schedule_hydration_history_collection(&self) {
129        let interval =
130            HYDRATION_HISTORY_COLLECTION_INTERVAL.get(self.catalog().system_config().dyncfgs());
131
132        // A zero interval disables collection. Keep polling so that enabling it
133        // takes effect without an `environmentd` restart.
134        let (delay, fire) = if interval.is_zero() {
135            (DISABLED_RECHECK_INTERVAL, false)
136        } else {
137            // An absurd interval saturates rather than panicking. The setting is
138            // durable, so a panic here would recur on every restart.
139            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            // Best effort: the coordinator may be shutting down.
163            let _ = internal_cmd_tx.send(message);
164        });
165    }
166
167    /// A stable offset into the collection interval for this environment.
168    ///
169    /// Seeded from the full environment id, so it survives restarts but differs
170    /// between environments in the same organization. Without it every
171    /// environment would sweep on the same absolute grid, turning each boundary
172    /// into a fleet-wide burst of dataflow installs, oracle round trips and persist
173    /// writes.
174    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    /// Runs one sweep: collect from the next replica, then apply retention.
180    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        // Builtin tables are not writable in read-only mode, and a disabled
186        // collector must do no background work at all. Retention is part of the
187        // sweep, so disabling collection also suspends it once any in-flight
188        // sweep finishes. This makes zero a break-glass setting for the subsystem.
189        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                // Unmanaged replicas have no orchestrator status and are only
206                // used by tests. Their bounded mutation determines readiness.
207                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        // Retention runs on the catalog server, so that it keeps working when
219        // there are no user replicas to collect from at all. Without a replica
220        // there it is skipped, while collection still runs.
221        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            // Retention runs even when collection failed above. A replica that
241            // is crash-looping or slow must not be able to stop the table from
242            // shrinking back to its retention bound.
243            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        // Keep the sweep coordinator-owned so dropping the coordinator requests
255        // its abort. Fallible coordinator calls make concurrent shutdown safe.
256        self.hydration_history_sweep = Some(handle.abort_on_drop());
257    }
258
259    /// Assembles the sweep context, including the client it writes through.
260    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        // Background read-then-write always uses the frontend OCC path. This
271        // shared constructor field only controls session fallback, so the flag
272        // does not gate history collection.
273        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/// A replica eligible for one collection step.
303#[derive(Clone, Copy, Debug, Eq, PartialEq)]
304struct ReplicaTarget {
305    cluster_id: ClusterId,
306    replica_id: ReplicaId,
307    process_count: usize,
308}
309
310/// Picks the replica after `cursor`, wrapping around at the end.
311///
312/// `replicas` must be sorted ascending by replica id. Unsorted input still
313/// returns a replica but degenerates the rotation, revisiting some replicas and
314/// starving others.
315fn 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
323/// The rows this replica has completed that the history table is missing.
324///
325/// Aggregates every worker's row for an export, and records nothing until all of
326/// them have hydrated. One worker is not enough, because a materialized view's
327/// persist sink has a single active worker, `hash(sink_id) % workers`. Only that
328/// worker's reported output frontier is gated on the shard upper, so only its
329/// `hydrated_at` covers the initial snapshot write. Every other worker clears its
330/// sink write frontier and stamps at compute completion, which for a materialized
331/// view is before the data is durable. Taking `max` over a complete set of workers
332/// is therefore the only way to get a finish that means the same thing for every
333/// object, and it is the rule `mz_compute_hydration_times` already applies.
334///
335/// Completeness needs no configured worker count. The log carries a row per
336/// `(export_id, worker_id)` from installation with a null `hydrated_at`, so
337/// `count(*) = count(hydrated_at)` says every row visible at the OCC read
338/// timestamp has finished. Per-process logging clocks also determine Differential
339/// update timestamps, so a worker whose clock is ahead can be absent at that
340/// timestamp. A visible unfinished object is skipped and picked up by a later
341/// sweep.
342///
343/// A worker missing at the read timestamp cannot later change the episode key.
344/// Its logging clock stamps both the Differential update and `installed_at`, so
345/// late visibility means its installation stamp is later than the visible
346/// minimum. The anti-join therefore keeps matching the recorded row. A later
347/// `hydrated_at` can raise the aggregate's maximum, but history is not repaired
348/// after the episode key has been recorded.
349///
350/// The collector deliberately accepts this sampling race rather than depending on
351/// `ReplicaLocation::workers()`. A durable finish can therefore precede the latest
352/// worker's finish. A whole-replica restart resets the collection as a unit.
353///
354/// The interval spans workers, so it carries whatever skew there is between the
355/// process clocks that stamped its ends. Each process anchors its logging clock at
356/// its own `SystemTime`. That inflates a duration, and nothing here rejects a row
357/// for being inconsistent, which is deliberate: an ordering guard on cross-worker
358/// stamps rejects complete episodes permanently, since the log values never change.
359///
360/// Collection has no explicit batch bound. It returns at most one row per
361/// not-yet-recorded dataflow, and the OCC path rejects a result that exceeds
362/// `max_result_size` or `max_query_result_size`. At their 1 GiB defaults that
363/// ceiling only matters at millions of dataflows per replica.
364fn object_collection_sql(cluster_id: ClusterId, replica_id: ReplicaId, cutoff: &str) -> String {
365    // Interpolating into SQL is safe here: the ids are catalog-internal and the
366    // cutoff is an RFC 3339 timestamp we formatted ourselves. Nothing in this
367    // query comes from a user.
368    //
369    // NOTE: The cutoff and the anti-join sit outside the aggregate deliberately. As
370    // a `WHERE` clause either one drops not-yet-hydrated rows, which would make
371    // `count(*) = count(hydrated_at)` trivially true and hand back a compute-only
372    // finish for a materialized view whose active worker is still writing.
373    //
374    // NOTE: `hydrated_at` is the terminal stamp for a history episode. Nothing
375    // waits for the history row before proceeding. If the log gains a separate
376    // `written_at` stamp, only `hydrated_at` belongs in this completeness check.
377    // A materialized view being replaced can hydrate while it runs read-only, and
378    // may never write if the replacement is rolled back.
379    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
411/// Returns SQL for the latest completed compute hydration episode and its
412/// process resource peaks.
413///
414/// Episodes are connected components of export hydration intervals. An export
415/// that has not hydrated keeps its component open: whether it is slow or
416/// permanently stuck is unobservable, so an open component is simply an
417/// in-progress episode, recorded when (if) it completes. It blocks only its
418/// own component: an unhydrated export installed at or before a completed
419/// component's finish would extend that component, one installed later belongs
420/// to a later episode. The latest completed component disconnected from every
421/// open one is recorded. The monotonic history guard still admits an open
422/// episode once it completes, because its start lies after every recorded
423/// finish. Cross-process clock skew can break that ordering, in which case the
424/// guard suppresses the episode rather than misrecording it.
425///
426/// Collection also waits until every configured replica process has reported
427/// resource usage. The query itself narrates how each step works.
428fn replica_collection_sql(target: ReplicaTarget, cutoff: &str) -> String {
429    let ReplicaTarget {
430        cluster_id,
431        replica_id,
432        process_count,
433    } = target;
434    // Interpolating into SQL is safe here: the ids and process count are
435    // catalog-internal and the cutoff is an RFC 3339 timestamp we formatted.
436    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
569/// A bounded batch of history rows that have aged out.
570///
571/// Only rows with a `hydrated_at` age out. Every row written today has one, and
572/// a row without one would be immortal here, so an unfinished-episode
573/// representation needs a second age basis before it can be recorded.
574fn object_retention_sql(cutoff: &str) -> String {
575    // The LIMIT has to sit inside a subquery. A top-level LIMIT lands in the
576    // plan's `RowSetFinishing`, which this OCC stage cannot apply. Inside a
577    // derived table it lowers into the relation expression instead.
578    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
591/// A bounded batch of replica history rows that have aged out.
592fn 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
606/// What one sweep needs to run its mutations against the history table.
607struct 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    /// Rows finishing before their table's cutoff have aged out. Collection and
615    /// retention use the same cutoff for each table, so this sweep cannot
616    /// resurrect an episode its own retention step retracts.
617    /// Concurrent sweeps can have different cutoffs, making retention eventual.
618    object_cutoff: String,
619    replica_cutoff: String,
620}
621
622impl Sweep {
623    /// Appends completed object and replica episodes from one replica.
624    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    /// Retracts one bounded batch of aged-out rows.
656    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    /// Runs one mutation, leaving transient failures for the next sweep to retry.
691    ///
692    /// Replica loss, dependency replacement, and write races are logged rather
693    /// than propagated because the next sweep recomputes from current state.
694    ///
695    /// NOTE: A timed-out mutation can still commit. The write is submitted
696    /// before we wait for its answer, and a background write carries no
697    /// connection to cancel it with, so a `timeout` outcome says we stopped
698    /// waiting, not that nothing landed. Rows such a write commits afterwards
699    /// are never counted, which makes `rows_affected` a lower bound.
700    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            // A trailing replica can repeatedly certify a target only after the
763            // oracle has advanced past it. Each refused write raises the target,
764            // and the conflict loop can continue until this timeout fires.
765            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
790/// Plans `sql` as the read side of a mutation against `target_id`.
791///
792/// The statement is planned as a `SELECT` whose columns are already in the
793/// target table's order, so the mutation needs no assignments or projection.
794///
795/// The selection's column types are checked against the target table here. A
796/// user `INSERT ... SELECT` gets that from the planner, but a hand-built plan
797/// bypasses it, and a wrong type would be written into the shard verbatim and
798/// break every later read of a table that is deliberately never truncated.
799fn 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            // Nullability may be tighter than the column allows, only the
834            // scalar types have to agree.
835            .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    /// Every environment shares one interval, so the grid has to be shifted per
886    /// environment or the whole fleet sweeps at the same instant.
887    #[mz_ore::test]
888    fn fire_delay_is_offset_within_the_interval() {
889        let interval = 60_000;
890
891        // Before this environment's point in the period, we wait for it.
892        assert_eq!(
893            next_fire_delay(1_000, interval, 5_000),
894            Duration::from_millis(4_000)
895        );
896        // On it, we take the next period rather than firing twice.
897        assert_eq!(
898            next_fire_delay(5_000, interval, 5_000),
899            Duration::from_millis(interval)
900        );
901        // After it, the next period's point.
902        assert_eq!(
903            next_fire_delay(6_000, interval, 5_000),
904            Duration::from_millis(59_000)
905        );
906        // A zero offset is plain alignment, and never returns a zero delay.
907        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        // Region and ordinal are part of the seed, not just the organization.
917        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    /// A materialized view's finish is only durable on the sink's active worker, so
931    /// the query has to see every worker and take the latest stamp. Pinning a single
932    /// worker, or letting the cutoff or the anti-join filter rows before the
933    /// completeness check, silently reintroduces a finish that precedes the write.
934    #[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        // Both of these have to apply to the aggregate's output, not to the rows
950        // feeding it.
951        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    /// Replica episodes are connected components of object hydration intervals,
960    /// enumerated gaps-and-islands style. An episode still connected to an open
961    /// interval is skipped via a scalar comparison against the earliest open
962    /// install, and the latest remaining episode is recorded. The query must
963    /// also wait for every replica process before it snapshots process-local
964    /// high-water marks.
965    #[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        // The gaps-and-islands scaffolding: running coverage horizon, gap
978        // detection against the previous row's horizon, episode labels, and
979        // per-episode aggregation.
980        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        // The unfinished-export guard applies per episode. A replica-wide
988        // all-hydrated gate would lose a completed episode for good: once the
989        // in-progress one finishes, it is the latest and the earlier one is
990        // never recorded.
991        assert!(!sql.contains("bool_and(hydrated)"), "{sql}");
992        // The guard compares each episode against the earliest open install,
993        // one scalar row. A join against all open intervals is quadratic when
994        // many episodes coexist with many still-hydrating exports.
995        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}