Skip to main content

mz_catalog/durable/
persist.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#[cfg(test)]
11mod tests;
12
13use std::cmp::max;
14use std::collections::{BTreeMap, VecDeque};
15use std::fmt::Debug;
16use std::str::FromStr;
17use std::sync::{Arc, LazyLock};
18use std::time::{Duration, Instant};
19
20use async_trait::async_trait;
21use differential_dataflow::lattice::Lattice;
22use futures::{FutureExt, StreamExt};
23use itertools::Itertools;
24use mz_audit_log::VersionedEvent;
25use mz_ore::cast::CastFrom;
26use mz_ore::metrics::MetricsFutureExt;
27use mz_ore::now::EpochMillis;
28use mz_ore::{
29    soft_assert_eq_no_log, soft_assert_eq_or_log, soft_assert_ne_or_log, soft_assert_no_log,
30    soft_assert_or_log, soft_panic_or_log,
31};
32use mz_persist_client::cfg::USE_CRITICAL_SINCE_CATALOG;
33use mz_persist_client::cli::admin::{CATALOG_FORCE_COMPACTION_FUEL, CATALOG_FORCE_COMPACTION_WAIT};
34use mz_persist_client::critical::{CriticalReaderId, Opaque, SinceHandle};
35use mz_persist_client::error::UpperMismatch;
36use mz_persist_client::read::{Listen, ListenEvent, ReadHandle};
37use mz_persist_client::write::WriteHandle;
38use mz_persist_client::{Diagnostics, PersistClient, ShardId};
39use mz_persist_types::codec_impls::UnitSchema;
40use mz_proto::{RustType, TryFromProtoError};
41use mz_repr::Diff;
42use mz_storage_client::controller::PersistEpoch;
43use mz_storage_types::StorageDiff;
44use mz_storage_types::sources::SourceData;
45use sha2::Digest;
46use timely::progress::{Antichain, Timestamp as TimelyTimestamp};
47use tracing::{debug, info, warn};
48use uuid::Uuid;
49
50use crate::durable::debug::{Collection, CollectionType, DebugCatalogState, Trace};
51use crate::durable::error::FenceError;
52use crate::durable::initialize::{
53    ENABLE_0DT_DEPLOYMENT_PANIC_AFTER_TIMEOUT, SYSTEM_CONFIG_SYNCED_KEY, USER_VERSION_KEY,
54    WITH_0DT_DEPLOYMENT_DDL_CHECK_INTERVAL, WITH_0DT_DEPLOYMENT_MAX_WAIT,
55};
56use crate::durable::metrics::Metrics;
57use crate::durable::objects::state_update::{
58    IntoStateUpdateKindJson, StateUpdate, StateUpdateKind, StateUpdateKindJson,
59    TryIntoStateUpdateKind,
60};
61use crate::durable::objects::{AuditLogKey, FenceToken, Snapshot};
62use crate::durable::transaction::TransactionBatch;
63use crate::durable::upgrade::upgrade;
64use crate::durable::{
65    BootstrapArgs, CATALOG_CONTENT_VERSION_KEY, CatalogError, DryRunTransaction,
66    DurableCatalogError, DurableCatalogState, Epoch, OpenableDurableCatalogState,
67    ReadOnlyDurableCatalogState, Transaction, initialize, persist_desc,
68};
69use crate::memory;
70
71/// New-type used to represent timestamps in persist.
72pub(crate) type Timestamp = mz_repr::Timestamp;
73
74/// The minimum value of an epoch.
75const MIN_EPOCH: Epoch = Epoch::new(1).expect("1 is non-zero");
76
77/// Human readable catalog shard name.
78const CATALOG_SHARD_NAME: &str = "catalog";
79
80/// [`CriticalReaderId`] for the catalog shard's own since hold, separate from
81/// [`PersistClient::CONTROLLER_CRITICAL_SINCE`].
82static CATALOG_CRITICAL_SINCE: LazyLock<CriticalReaderId> = LazyLock::new(|| {
83    "c55555555-6666-7777-8888-999999999999"
84        .parse()
85        .expect("valid CriticalReaderId")
86});
87
88/// Seed used to generate the persist shard ID for the catalog.
89const CATALOG_SEED: usize = 1;
90/// Legacy seed used to generate the persist shard ID for the upgrade shard. DO NOT REUSE.
91const _UPGRADE_SEED: usize = 2;
92/// Legacy seed used to generate the persist shard ID for builtin table migrations. DO NOT REUSE.
93pub const _BUILTIN_MIGRATION_SEED: usize = 3;
94/// Legacy seed used to generate the persist shard ID for the expression cache. DO NOT REUSE.
95pub const _EXPRESSION_CACHE_SEED: usize = 4;
96
97/// Durable catalog mode that dictates the effect of mutable operations.
98#[derive(Debug, Copy, Clone, Eq, PartialEq)]
99pub(crate) enum Mode {
100    /// Mutable operations are prohibited.
101    Readonly,
102    /// Mutable operations have an effect in-memory, but aren't persisted durably.
103    Savepoint,
104    /// Mutable operations have an effect in-memory and durably.
105    Writable,
106}
107
108/// Enum representing the fenced state of the catalog.
109#[derive(Debug)]
110pub(crate) enum FenceableToken {
111    /// The catalog is still initializing and learning about previously written fence tokens. This
112    /// state can be fenced if it encounters a larger deploy generation.
113    Initializing {
114        /// The largest fence token durably written to the catalog, if any.
115        durable_token: Option<FenceToken>,
116        /// This process's deploy generation.
117        current_deploy_generation: Option<u64>,
118    },
119    /// The current token has not been fenced.
120    Unfenced { current_token: FenceToken },
121    /// The current token has been fenced.
122    Fenced {
123        current_token: FenceToken,
124        fence_token: FenceToken,
125    },
126}
127
128impl FenceableToken {
129    /// Returns a new token.
130    fn new(current_deploy_generation: Option<u64>) -> Self {
131        Self::Initializing {
132            durable_token: None,
133            current_deploy_generation,
134        }
135    }
136
137    /// Returns the current token if it is not fenced, otherwise returns an error.
138    fn validate(&self) -> Result<Option<FenceToken>, FenceError> {
139        match self {
140            FenceableToken::Initializing { durable_token, .. } => Ok(durable_token.clone()),
141            FenceableToken::Unfenced { current_token, .. } => Ok(Some(current_token.clone())),
142            FenceableToken::Fenced {
143                current_token,
144                fence_token,
145            } => {
146                assert!(
147                    fence_token > current_token,
148                    "must be fenced by higher token; current={current_token:?}, fence={fence_token:?}"
149                );
150                if fence_token.deploy_generation > current_token.deploy_generation {
151                    Err(FenceError::DeployGeneration {
152                        current_generation: current_token.deploy_generation,
153                        fence_generation: fence_token.deploy_generation,
154                    })
155                } else {
156                    assert!(
157                        fence_token.epoch > current_token.epoch,
158                        "must be fenced by higher token; current={current_token:?}, fence={fence_token:?}"
159                    );
160                    Err(FenceError::Epoch {
161                        current_epoch: current_token.epoch,
162                        fence_epoch: fence_token.epoch,
163                    })
164                }
165            }
166        }
167    }
168
169    /// Returns the current token.
170    fn token(&self) -> Option<FenceToken> {
171        match self {
172            FenceableToken::Initializing { durable_token, .. } => durable_token.clone(),
173            FenceableToken::Unfenced { current_token, .. } => Some(current_token.clone()),
174            FenceableToken::Fenced { current_token, .. } => Some(current_token.clone()),
175        }
176    }
177
178    /// Returns `Err` if `token` fences out `self`, `Ok` otherwise.
179    fn maybe_fence(&mut self, token: FenceToken) -> Result<(), FenceError> {
180        match self {
181            FenceableToken::Initializing {
182                durable_token,
183                current_deploy_generation,
184                ..
185            } => {
186                match durable_token {
187                    Some(durable_token) => {
188                        *durable_token = max(durable_token.clone(), token.clone());
189                    }
190                    None => {
191                        *durable_token = Some(token.clone());
192                    }
193                }
194                if let Some(current_deploy_generation) = current_deploy_generation {
195                    if *current_deploy_generation < token.deploy_generation {
196                        *self = FenceableToken::Fenced {
197                            current_token: FenceToken {
198                                deploy_generation: *current_deploy_generation,
199                                epoch: token.epoch,
200                            },
201                            fence_token: token,
202                        };
203                        self.validate()?;
204                    }
205                }
206            }
207            FenceableToken::Unfenced { current_token } => {
208                if *current_token < token {
209                    *self = FenceableToken::Fenced {
210                        current_token: current_token.clone(),
211                        fence_token: token,
212                    };
213                    self.validate()?;
214                }
215            }
216            FenceableToken::Fenced { .. } => {
217                self.validate()?;
218            }
219        }
220
221        Ok(())
222    }
223
224    /// Returns a [`FenceableToken::Unfenced`] token and the updates to the catalog required to
225    /// transition to the `Unfenced` state if `self` is [`FenceableToken::Initializing`], otherwise
226    /// returns `None`.
227    fn generate_unfenced_token(
228        &self,
229        mode: Mode,
230    ) -> Result<Option<(Vec<(StateUpdateKind, Diff)>, FenceableToken)>, DurableCatalogError> {
231        let (durable_token, current_deploy_generation) = match self {
232            FenceableToken::Initializing {
233                durable_token,
234                current_deploy_generation,
235            } => (durable_token.clone(), current_deploy_generation.clone()),
236            FenceableToken::Unfenced { .. } | FenceableToken::Fenced { .. } => return Ok(None),
237        };
238
239        let mut fence_updates = Vec::with_capacity(2);
240
241        if let Some(durable_token) = &durable_token {
242            fence_updates.push((
243                StateUpdateKind::FenceToken(durable_token.clone()),
244                Diff::MINUS_ONE,
245            ));
246        }
247
248        let current_deploy_generation = current_deploy_generation
249            .or_else(|| durable_token.as_ref().map(|token| token.deploy_generation))
250            // We cannot initialize a catalog without a deploy generation.
251            .ok_or(DurableCatalogError::Uninitialized)?;
252        let mut current_epoch = durable_token
253            .map(|token| token.epoch)
254            .unwrap_or(MIN_EPOCH)
255            .get();
256        // Only writable catalogs attempt to increment the epoch.
257        if matches!(mode, Mode::Writable) {
258            current_epoch = current_epoch + 1;
259        }
260        let current_epoch = Epoch::new(current_epoch).expect("known to be non-zero");
261        let current_token = FenceToken {
262            deploy_generation: current_deploy_generation,
263            epoch: current_epoch,
264        };
265
266        fence_updates.push((
267            StateUpdateKind::FenceToken(current_token.clone()),
268            Diff::ONE,
269        ));
270
271        let current_fenceable_token = FenceableToken::Unfenced { current_token };
272
273        Ok(Some((fence_updates, current_fenceable_token)))
274    }
275}
276
277/// An error that can occur while executing [`PersistHandle::compare_and_append`].
278#[derive(Debug, thiserror::Error)]
279pub(crate) enum CompareAndAppendError {
280    #[error(transparent)]
281    Fence(#[from] FenceError),
282    /// Catalog encountered an upper mismatch when trying to write to the catalog: another
283    /// writer moved the upper between our snapshot of it and the write. Handled by the conflict
284    /// classification in the commit and advance paths (rebase over empty progress, surface
285    /// content conflicts as out-of-sync).
286    #[error(
287        "expected catalog upper {expected_upper:?} did not match actual catalog upper {actual_upper:?}"
288    )]
289    UpperMismatch {
290        expected_upper: Timestamp,
291        actual_upper: Timestamp,
292    },
293}
294
295impl CompareAndAppendError {
296    pub(crate) fn unwrap_fence_error(self) -> FenceError {
297        match self {
298            CompareAndAppendError::Fence(e) => e,
299            e @ CompareAndAppendError::UpperMismatch { .. } => {
300                panic!("unexpected upper mismatch: {e:?}")
301            }
302        }
303    }
304}
305
306impl From<UpperMismatch<Timestamp>> for CompareAndAppendError {
307    fn from(upper_mismatch: UpperMismatch<Timestamp>) -> Self {
308        Self::UpperMismatch {
309            expected_upper: antichain_to_timestamp(upper_mismatch.expected),
310            actual_upper: antichain_to_timestamp(upper_mismatch.current),
311        }
312    }
313}
314
315pub(crate) trait ApplyUpdate<T: IntoStateUpdateKindJson> {
316    /// Process and apply `update`.
317    ///
318    /// Returns `Some` if `update` should be cached in memory and `None` otherwise.
319    fn apply_update(
320        &mut self,
321        update: StateUpdate<T>,
322        current_fence_token: &mut FenceableToken,
323        metrics: &Arc<Metrics>,
324    ) -> Result<Option<StateUpdate<T>>, FenceError>;
325}
326
327/// A handle for interacting with the persist catalog shard.
328///
329/// The catalog shard is used in multiple different contexts, for example pre-open and post-open,
330/// but for all contexts the majority of the durable catalog's behavior is identical. This struct
331/// implements those behaviors that are identical while allowing the user to specify the different
332/// behaviors via generic parameters.
333///
334/// The behavior of the durable catalog can be different along one of two axes. The first is the
335/// format of each individual update, i.e. raw binary, the current protobuf version, previous
336/// protobuf versions, etc. The second axis is what to do with each individual update, for example
337/// before opening we cache all config updates but don't cache them after opening. These behaviors
338/// are customizable via the `T: TryIntoStateUpdateKind` and `U: ApplyUpdate<T>` generic parameters
339/// respectively.
340#[derive(Debug)]
341pub(crate) struct PersistHandle<T: TryIntoStateUpdateKind, U: ApplyUpdate<T>> {
342    /// The [`Mode`] that this catalog was opened in.
343    pub(crate) mode: Mode,
344    /// Since handle to control compaction.
345    since_handle: SinceHandle<SourceData, (), Timestamp, StorageDiff>,
346    /// Write handle to persist.
347    write_handle: WriteHandle<SourceData, (), Timestamp, StorageDiff>,
348    /// Listener to catalog changes.
349    listen: Listen<SourceData, (), Timestamp, StorageDiff>,
350    /// Handle for connecting to persist.
351    persist_client: PersistClient,
352    /// Catalog shard ID.
353    shard_id: ShardId,
354    /// Cache of the most recent catalog snapshot.
355    ///
356    /// We use a tuple instead of [`StateUpdate`] to make consolidation easier.
357    pub(crate) snapshot: Vec<(T, Timestamp, Diff)>,
358    /// Applies custom processing, filtering, and fencing for each individual update.
359    update_applier: U,
360    /// The current upper of the persist shard.
361    pub(crate) upper: Timestamp,
362    /// The fence token of the catalog, if one exists.
363    fenceable_token: FenceableToken,
364    /// The semantic version of the current binary.
365    catalog_content_version: semver::Version,
366    /// Flag to indicate if bootstrap is complete.
367    bootstrap_complete: bool,
368    /// Metrics for the persist catalog.
369    metrics: Arc<Metrics>,
370    /// Snapshot size at the last amortized consolidation, used by
371    /// [`Self::maybe_consolidate`] to decide when to consolidate. Initialized
372    /// lazily on the first call.
373    size_at_last_consolidation: Option<usize>,
374    /// Counts raw updates applied to this handle.
375    ///
376    /// This distinguishes empty upper progress from content. Memory updates are insufficient
377    /// because kinds such as ID allocators do not produce them.
378    updates_applied: u64,
379}
380
381impl<T: TryIntoStateUpdateKind, U: ApplyUpdate<T>> PersistHandle<T, U> {
382    /// Fetch the current upper of the catalog state.
383    #[mz_ore::instrument]
384    async fn current_upper(&mut self) -> Timestamp {
385        match self.mode {
386            Mode::Writable | Mode::Readonly => {
387                let upper = self.write_handle.fetch_recent_upper().await;
388                antichain_to_timestamp(upper.clone())
389            }
390            Mode::Savepoint => self.upper,
391        }
392    }
393
394    /// Appends `updates` iff the current global upper of the catalog is `self.upper`.
395    ///
396    /// Returns the next upper used to commit the transaction.
397    #[mz_ore::instrument]
398    pub(crate) async fn compare_and_append<S: IntoStateUpdateKindJson>(
399        &mut self,
400        updates: Vec<(S, Diff)>,
401        commit_ts: Timestamp,
402    ) -> Result<Timestamp, CompareAndAppendError> {
403        let updates = updates.into_iter().map(|(kind, diff)| {
404            let kind: StateUpdateKindJson = kind.into();
405            (
406                (Into::<SourceData>::into(kind), ()),
407                commit_ts,
408                diff.into_inner(),
409            )
410        });
411        let next_upper = commit_ts.step_forward();
412        // Upper mismatches are classified by the commit and advance callers.
413        self.compare_and_append_inner(updates, next_upper).await?;
414
415        self.sync(next_upper).await?;
416        Ok(next_upper)
417    }
418
419    /// Compare-and-append `updates` to the catalog shard, advancing the upper to `next_upper`.
420    ///
421    /// On success, updating `self.upper` is left to the caller. The caller can thus decide whether
422    /// or not it needs to sync the catalog.
423    ///
424    /// # Panics
425    ///
426    /// Panics if not in `Writable` mode.
427    /// Panics if `next_upper` is not greater than `self.upper`.
428    async fn compare_and_append_inner(
429        &mut self,
430        updates: impl IntoIterator<Item = ((SourceData, ()), Timestamp, StorageDiff)>,
431        next_upper: Timestamp,
432    ) -> Result<(), CompareAndAppendError> {
433        assert_eq!(self.mode, Mode::Writable);
434        assert!(
435            next_upper > self.upper,
436            "next_upper ({next_upper}) not greater than current upper ({})",
437            self.upper,
438        );
439
440        let res = self
441            .write_handle
442            .compare_and_append(
443                updates,
444                Antichain::from_elem(self.upper),
445                Antichain::from_elem(next_upper),
446            )
447            .await
448            .expect("invalid usage");
449
450        if let Err(e @ UpperMismatch { .. }) = res {
451            // Most likely we were fenced out.
452            // Sync to the current upper to detect that.
453            self.sync_to_current_upper().await?;
454            return Err(e.into());
455        }
456
457        // Lag the shard's upper by 1 to keep it readable.
458        let downgrade_to = Antichain::from_elem(next_upper.saturating_sub(1));
459
460        // The since handle gives us the ability to fence out other downgraders using an opaque token.
461        // (See the method documentation for details.)
462        // That's not needed here, so we use the since handle's opaque token to avoid any comparison
463        // failures.
464        let opaque = self.since_handle.opaque().clone();
465        let downgrade = self
466            .since_handle
467            .maybe_compare_and_downgrade_since(&opaque, (&opaque, &downgrade_to))
468            .await;
469        if let Some(Err(e)) = downgrade {
470            soft_panic_or_log!("found opaque value {e:?}, but expected {opaque:?}");
471        }
472
473        Ok(())
474    }
475
476    /// Accepts an upper mismatch caused only by empty progress.
477    ///
478    /// `updates_applied_before` must be captured before the compare-and-append, which synchronizes
479    /// this handle on mismatch.
480    fn classify_upper_mismatch(
481        &self,
482        updates_applied_before: u64,
483        actual_upper: Timestamp,
484    ) -> Result<(), DurableCatalogError> {
485        if self.updates_applied != updates_applied_before {
486            Err(DurableCatalogError::CatalogOutOfSync {
487                update_count: usize::cast_from(self.updates_applied - updates_applied_before),
488                upper: actual_upper,
489            })
490        } else {
491            Ok(())
492        }
493    }
494
495    /// Generates an iterator of [`StateUpdate`] that contain all unconsolidated updates to the
496    /// catalog state up to, and including, `as_of`.
497    #[mz_ore::instrument]
498    async fn snapshot_unconsolidated(&mut self) -> Vec<StateUpdate<StateUpdateKind>> {
499        let current_upper = self.current_upper().await;
500
501        let mut snapshot = Vec::new();
502        let mut read_handle = self.read_handle().await;
503        let as_of = as_of(&read_handle, current_upper);
504        let mut stream = Box::pin(
505            // We use `snapshot_and_stream` because it guarantees unconsolidated output.
506            read_handle
507                .snapshot_and_stream(Antichain::from_elem(as_of))
508                .await
509                .expect("we have advanced the restart_as_of by the since"),
510        );
511        while let Some(update) = stream.next().await {
512            snapshot.push(update)
513        }
514        read_handle.expire().await;
515        snapshot
516            .into_iter()
517            .map(Into::<StateUpdate<StateUpdateKindJson>>::into)
518            .map(|state_update| state_update.try_into().expect("kind decoding error"))
519            .collect()
520    }
521
522    /// Listen and apply all updates that are currently in persist.
523    ///
524    /// Returns an error if this instance has been fenced out.
525    #[mz_ore::instrument]
526    pub(crate) async fn sync_to_current_upper(&mut self) -> Result<(), FenceError> {
527        let upper = self.current_upper().await;
528        self.sync(upper).await
529    }
530
531    /// Listen and apply all updates up to `target_upper`.
532    ///
533    /// Returns an error if this instance has been fenced out.
534    #[mz_ore::instrument(level = "debug")]
535    pub(crate) async fn sync(&mut self, target_upper: Timestamp) -> Result<(), FenceError> {
536        self.metrics.syncs.inc();
537        let histogram = self.metrics.sync_latency_seconds.clone();
538        self.sync_inner(target_upper)
539            .wall_time()
540            .observe(histogram)
541            .await
542    }
543
544    #[mz_ore::instrument(level = "debug")]
545    async fn sync_inner(&mut self, target_upper: Timestamp) -> Result<(), FenceError> {
546        self.fenceable_token.validate()?;
547
548        // Savepoint catalogs do not yet know how to update themselves in response to concurrent
549        // writes from writer catalogs.
550        if self.mode == Mode::Savepoint {
551            self.upper = max(self.upper, target_upper);
552            return Ok(());
553        }
554
555        let mut updates: BTreeMap<_, Vec<_>> = BTreeMap::new();
556        let updates_applied_before = self.updates_applied;
557
558        // Reset the amortized consolidation tracker so it picks up the
559        // current snapshot size as its baseline.
560        self.size_at_last_consolidation = None;
561
562        while self.upper < target_upper {
563            let listen_events = self.listen.fetch_next().await;
564            for listen_event in listen_events {
565                match listen_event {
566                    ListenEvent::Progress(upper) => {
567                        debug!("synced up to {upper:?}");
568                        self.upper = antichain_to_timestamp(upper);
569                        // Attempt to apply updates in batches of a single timestamp. If another
570                        // catalog wrote a fence token at one timestamp and then updates in a new
571                        // format at a later timestamp, then we want to apply the fence token
572                        // before attempting to deserialize the new updates.
573                        while let Some((ts, updates)) = updates.pop_first() {
574                            assert!(ts < self.upper, "expected {} < {}", ts, self.upper);
575                            let updates = updates.into_iter().map(
576                                |update: StateUpdate<StateUpdateKindJson>| {
577                                    let kind =
578                                        T::try_from(update.kind).expect("kind decoding error");
579                                    StateUpdate {
580                                        kind,
581                                        ts: update.ts,
582                                        diff: update.diff,
583                                    }
584                                },
585                            );
586                            self.apply_updates(updates)?;
587                            self.maybe_consolidate();
588                        }
589                    }
590                    ListenEvent::Updates(batch_updates) => {
591                        for update in batch_updates {
592                            let update: StateUpdate<StateUpdateKindJson> = update.into();
593                            updates.entry(update.ts).or_default().push(update);
594                        }
595                    }
596                }
597            }
598        }
599        assert_eq!(updates, BTreeMap::new(), "all updates should be applied");
600        // Only consolidate when there are actual updates.
601        if self.updates_applied != updates_applied_before {
602            self.consolidate();
603        }
604        Ok(())
605    }
606
607    /// Apply a batch of updates and then consolidate the snapshot. This is the
608    /// typical entry point for callers that apply updates in a single batch.
609    ///
610    /// For hot loops that apply updates across many timestamps (e.g., `sync_inner`),
611    /// use `apply_updates` directly and call `consolidate()` periodically (e.g.,
612    /// on snapshot doubling) to bound memory while staying amortized O(N log N).
613    pub(crate) fn apply_updates_and_consolidate(
614        &mut self,
615        updates: impl IntoIterator<Item = StateUpdate<T>>,
616    ) -> Result<(), FenceError> {
617        self.apply_updates(updates)?;
618        self.consolidate();
619        Ok(())
620    }
621
622    /// Apply a batch of updates to the catalog state without consolidating.
623    ///
624    /// Does NOT consolidate the snapshot afterward. If you are calling this once,
625    /// prefer `apply_updates_and_consolidate`. This method exists for loops that
626    /// call it many times — consolidating per call would be O(K * N log N) instead
627    /// of O(N log N). Callers should consolidate periodically (e.g., on snapshot
628    /// doubling) to bound memory.
629    #[mz_ore::instrument(level = "debug")]
630    fn apply_updates(
631        &mut self,
632        updates: impl IntoIterator<Item = StateUpdate<T>>,
633    ) -> Result<(), FenceError> {
634        let mut updates: Vec<_> = updates
635            .into_iter()
636            .map(|StateUpdate { kind, ts, diff }| (kind, ts, diff))
637            .collect();
638
639        // This helps guarantee that for a single key, there is at most a single retraction and a
640        // single insertion per timestamp. Otherwise, we would need to match the retractions and
641        // insertions up by value and manually figure out what the end value should be.
642        differential_dataflow::consolidation::consolidate_updates(&mut updates);
643
644        // Updates must be applied in timestamp order. Within a timestamp retractions must be
645        // applied before insertions, or we might end up retracting the wrong value.
646        updates.sort_by(|(_, ts1, diff1), (_, ts2, diff2)| ts1.cmp(ts2).then(diff1.cmp(diff2)));
647
648        let mut errors = Vec::new();
649
650        for (kind, ts, diff) in updates {
651            if diff != Diff::ONE && diff != Diff::MINUS_ONE {
652                panic!("invalid update in consolidated trace: ({kind:?}, {ts:?}, {diff:?})");
653            }
654            self.updates_applied += 1;
655
656            match self.update_applier.apply_update(
657                StateUpdate { kind, ts, diff },
658                &mut self.fenceable_token,
659                &self.metrics,
660            ) {
661                Ok(Some(StateUpdate { kind, ts, diff })) => self.snapshot.push((kind, ts, diff)),
662                Ok(None) => {}
663                // Instead of returning immediately, we accumulate all the errors and return the one
664                // with the most information.
665                Err(err) => errors.push(err),
666            }
667        }
668
669        // Track the high-water mark of the unconsolidated snapshot size.
670        let len = i64::try_from(self.snapshot.len()).unwrap_or(i64::MAX);
671        if len > self.metrics.snapshot_max_entries.get() {
672            self.metrics.snapshot_max_entries.set(len);
673        }
674
675        errors.sort();
676        if let Some(err) = errors.into_iter().next() {
677            return Err(err);
678        }
679
680        Ok(())
681    }
682
683    /// Consolidate the snapshot if it has at least doubled in size since the
684    /// last consolidation. This amortizes the O(N log N) consolidation cost
685    /// over many small updates, keeping the total work O(N log N) rather than
686    /// O(K * N log N) for K timestamps.
687    fn maybe_consolidate(&mut self) {
688        let threshold = *self
689            .size_at_last_consolidation
690            // Use a minimum of 8 to avoid consolidating on every update when
691            // the snapshot is small or empty (since 0 * 2 = 0).
692            .get_or_insert_with(|| max(self.snapshot.len(), 8));
693        if self.snapshot.len() >= threshold * 2 {
694            self.consolidate();
695            self.size_at_last_consolidation = Some(self.snapshot.len());
696        }
697    }
698
699    #[mz_ore::instrument]
700    pub(crate) fn consolidate(&mut self) {
701        self.metrics.snapshot_consolidations.inc();
702        soft_assert_no_log!(
703            self.snapshot
704                .windows(2)
705                .all(|updates| updates[0].1 <= updates[1].1),
706            "snapshot should be sorted by timestamp, {:#?}",
707            self.snapshot
708        );
709
710        let new_ts = self
711            .snapshot
712            .last()
713            .map(|(_, ts, _)| *ts)
714            .unwrap_or_else(Timestamp::minimum);
715        for (_, ts, _) in &mut self.snapshot {
716            *ts = new_ts;
717        }
718        differential_dataflow::consolidation::consolidate_updates(&mut self.snapshot);
719    }
720
721    /// Execute and return the results of `f` on the current catalog trace.
722    ///
723    /// Will return an error if the catalog has been fenced out.
724    async fn with_trace<R>(
725        &mut self,
726        f: impl FnOnce(&Vec<(T, Timestamp, Diff)>) -> Result<R, CatalogError>,
727    ) -> Result<R, CatalogError> {
728        self.sync_to_current_upper().await?;
729        f(&self.snapshot)
730    }
731
732    /// Open a read handle to the catalog.
733    async fn read_handle(&self) -> ReadHandle<SourceData, (), Timestamp, StorageDiff> {
734        self.persist_client
735            .open_leased_reader(
736                self.shard_id,
737                Arc::new(persist_desc()),
738                Arc::new(UnitSchema::default()),
739                Diagnostics {
740                    shard_name: CATALOG_SHARD_NAME.to_string(),
741                    handle_purpose: "openable durable catalog state temporary reader".to_string(),
742                },
743                USE_CRITICAL_SINCE_CATALOG.get(self.persist_client.dyncfgs()),
744            )
745            .await
746            .expect("invalid usage")
747    }
748
749    /// Politely releases all external resources that can only be released in an async context.
750    async fn expire(self: Box<Self>) {
751        self.write_handle.expire().await;
752        self.listen.expire().await;
753    }
754}
755
756impl<U: ApplyUpdate<StateUpdateKind>> PersistHandle<StateUpdateKind, U> {
757    /// Execute and return the results of `f` on the current catalog snapshot.
758    ///
759    /// Will return an error if the catalog has been fenced out.
760    async fn with_snapshot<T>(
761        &mut self,
762        f: impl FnOnce(Snapshot) -> Result<T, CatalogError>,
763    ) -> Result<T, CatalogError> {
764        fn apply<K, V>(map: &mut BTreeMap<K, V>, key: &K, value: &V, diff: Diff)
765        where
766            K: Ord + Clone,
767            V: Ord + Clone + Debug,
768        {
769            let key = key.clone();
770            let value = value.clone();
771            if diff == Diff::ONE {
772                let prev = map.insert(key, value);
773                assert_eq!(
774                    prev, None,
775                    "values must be explicitly retracted before inserting a new value"
776                );
777            } else if diff == Diff::MINUS_ONE {
778                let prev = map.remove(&key);
779                assert_eq!(
780                    prev,
781                    Some(value),
782                    "retraction does not match existing value"
783                );
784            }
785        }
786
787        self.with_trace(|trace| {
788            let mut snapshot = Snapshot::empty();
789            for (kind, ts, diff) in trace {
790                let diff = *diff;
791                if diff != Diff::ONE && diff != Diff::MINUS_ONE {
792                    panic!("invalid update in consolidated trace: ({kind:?}, {ts:?}, {diff:?})");
793                }
794
795                match kind {
796                    StateUpdateKind::AuditLog(_key, ()) => {
797                        // Ignore for snapshots.
798                    }
799                    StateUpdateKind::Cluster(key, value) => {
800                        apply(&mut snapshot.clusters, key, value, diff);
801                    }
802                    StateUpdateKind::ClusterReplica(key, value) => {
803                        apply(&mut snapshot.cluster_replicas, key, value, diff);
804                    }
805                    StateUpdateKind::Comment(key, value) => {
806                        apply(&mut snapshot.comments, key, value, diff);
807                    }
808                    StateUpdateKind::Config(key, value) => {
809                        apply(&mut snapshot.configs, key, value, diff);
810                    }
811                    StateUpdateKind::Database(key, value) => {
812                        apply(&mut snapshot.databases, key, value, diff);
813                    }
814                    StateUpdateKind::DefaultPrivilege(key, value) => {
815                        apply(&mut snapshot.default_privileges, key, value, diff);
816                    }
817                    StateUpdateKind::FenceToken(_token) => {
818                        // Ignore for snapshots.
819                    }
820                    StateUpdateKind::IdAllocator(key, value) => {
821                        apply(&mut snapshot.id_allocator, key, value, diff);
822                    }
823                    StateUpdateKind::IntrospectionSourceIndex(key, value) => {
824                        apply(&mut snapshot.introspection_sources, key, value, diff);
825                    }
826                    StateUpdateKind::Item(key, value) => {
827                        apply(&mut snapshot.items, key, value, diff);
828                    }
829                    StateUpdateKind::NetworkPolicy(key, value) => {
830                        apply(&mut snapshot.network_policies, key, value, diff);
831                    }
832                    StateUpdateKind::Role(key, value) => {
833                        apply(&mut snapshot.roles, key, value, diff);
834                    }
835                    StateUpdateKind::Schema(key, value) => {
836                        apply(&mut snapshot.schemas, key, value, diff);
837                    }
838                    StateUpdateKind::Setting(key, value) => {
839                        apply(&mut snapshot.settings, key, value, diff);
840                    }
841                    StateUpdateKind::SourceReferences(key, value) => {
842                        apply(&mut snapshot.source_references, key, value, diff);
843                    }
844                    StateUpdateKind::SystemConfiguration(key, value) => {
845                        apply(&mut snapshot.system_configurations, key, value, diff);
846                    }
847                    StateUpdateKind::ClusterSystemConfiguration(key, value) => {
848                        apply(
849                            &mut snapshot.cluster_system_configurations,
850                            key,
851                            value,
852                            diff,
853                        );
854                    }
855                    StateUpdateKind::ReplicaSystemConfiguration(key, value) => {
856                        apply(
857                            &mut snapshot.replica_system_configurations,
858                            key,
859                            value,
860                            diff,
861                        );
862                    }
863                    StateUpdateKind::SystemObjectMapping(key, value) => {
864                        apply(&mut snapshot.system_object_mappings, key, value, diff);
865                    }
866                    StateUpdateKind::SystemPrivilege(key, value) => {
867                        apply(&mut snapshot.system_privileges, key, value, diff);
868                    }
869                    StateUpdateKind::StorageCollectionMetadata(key, value) => {
870                        apply(&mut snapshot.storage_collection_metadata, key, value, diff);
871                    }
872                    StateUpdateKind::UnfinalizedShard(key, ()) => {
873                        apply(&mut snapshot.unfinalized_shards, key, &(), diff);
874                    }
875                    StateUpdateKind::TxnWalShard((), value) => {
876                        apply(&mut snapshot.txn_wal_shard, &(), value, diff);
877                    }
878                    StateUpdateKind::RoleAuth(key, value) => {
879                        apply(&mut snapshot.role_auth, key, value, diff);
880                    }
881                }
882            }
883            f(snapshot)
884        })
885        .await
886    }
887
888    /// Generates an iterator of [`StateUpdate`] that contain all updates to the catalog
889    /// state.
890    ///
891    /// The output is fetched directly from persist instead of the in-memory cache.
892    ///
893    /// The output is consolidated and sorted by timestamp in ascending order.
894    #[mz_ore::instrument(level = "debug")]
895    async fn persist_snapshot(&self) -> impl Iterator<Item = StateUpdate> + DoubleEndedIterator {
896        let mut read_handle = self.read_handle().await;
897        let as_of = as_of(&read_handle, self.upper);
898        let snapshot = snapshot_binary(&mut read_handle, as_of, &self.metrics)
899            .await
900            .map(|update| update.try_into().expect("kind decoding error"));
901        read_handle.expire().await;
902        snapshot
903    }
904}
905
906/// Applies updates for an unopened catalog.
907#[derive(Debug)]
908pub(crate) struct UnopenedCatalogStateInner {
909    /// A cache of the config collection of the catalog.
910    configs: BTreeMap<String, u64>,
911    /// A cache of the settings collection of the catalog.
912    settings: BTreeMap<String, String>,
913}
914
915impl UnopenedCatalogStateInner {
916    fn new() -> UnopenedCatalogStateInner {
917        UnopenedCatalogStateInner {
918            configs: BTreeMap::new(),
919            settings: BTreeMap::new(),
920        }
921    }
922}
923
924impl ApplyUpdate<StateUpdateKindJson> for UnopenedCatalogStateInner {
925    fn apply_update(
926        &mut self,
927        update: StateUpdate<StateUpdateKindJson>,
928        current_fence_token: &mut FenceableToken,
929        _metrics: &Arc<Metrics>,
930    ) -> Result<Option<StateUpdate<StateUpdateKindJson>>, FenceError> {
931        if !update.kind.is_audit_log() && update.kind.is_always_deserializable() {
932            let kind = TryInto::try_into(&update.kind).expect("kind is known to be deserializable");
933            match (kind, update.diff) {
934                (StateUpdateKind::Config(key, value), Diff::ONE) => {
935                    let prev = self.configs.insert(key.key, value.value);
936                    assert_eq!(
937                        prev, None,
938                        "values must be explicitly retracted before inserting a new value"
939                    );
940                }
941                (StateUpdateKind::Config(key, value), Diff::MINUS_ONE) => {
942                    let prev = self.configs.remove(&key.key);
943                    assert_eq!(
944                        prev,
945                        Some(value.value),
946                        "retraction does not match existing value"
947                    );
948                }
949                (StateUpdateKind::Setting(key, value), Diff::ONE) => {
950                    let prev = self.settings.insert(key.name, value.value);
951                    assert_eq!(
952                        prev, None,
953                        "values must be explicitly retracted before inserting a new value"
954                    );
955                }
956                (StateUpdateKind::Setting(key, value), Diff::MINUS_ONE) => {
957                    let prev = self.settings.remove(&key.name);
958                    assert_eq!(
959                        prev,
960                        Some(value.value),
961                        "retraction does not match existing value"
962                    );
963                }
964                (StateUpdateKind::FenceToken(fence_token), Diff::ONE) => {
965                    current_fence_token.maybe_fence(fence_token)?;
966                }
967                _ => {}
968            }
969        }
970
971        Ok(Some(update))
972    }
973}
974
975/// A Handle to an unopened catalog stored in persist. The unopened catalog can serve `Config` data,
976/// `Setting` data, or the current epoch. All other catalog data may be un-migrated and should not
977/// be read until the catalog has been opened. The [`UnopenedPersistCatalogState`] is responsible
978/// for opening the catalog, see [`OpenableDurableCatalogState::open`] for more details.
979///
980/// Production users should call [`Self::expire`] before dropping an [`UnopenedPersistCatalogState`]
981/// so that it can expire its leases. If/when rust gets AsyncDrop, this will be done automatically.
982pub(crate) type UnopenedPersistCatalogState =
983    PersistHandle<StateUpdateKindJson, UnopenedCatalogStateInner>;
984
985impl UnopenedPersistCatalogState {
986    /// Create a new [`UnopenedPersistCatalogState`] to the catalog state associated with
987    /// `organization_id`.
988    ///
989    /// All usages of the persist catalog must go through this function. That includes the
990    /// catalog-debug tool, the adapter's catalog, etc.
991    #[mz_ore::instrument]
992    pub(crate) async fn new(
993        persist_client: PersistClient,
994        organization_id: Uuid,
995        version: semver::Version,
996        deploy_generation: Option<u64>,
997        metrics: Arc<Metrics>,
998    ) -> Result<UnopenedPersistCatalogState, DurableCatalogError> {
999        let catalog_shard_id = shard_id(organization_id, CATALOG_SEED);
1000        debug!(?catalog_shard_id, "new persist backed catalog state");
1001
1002        // Check the catalog shard version to ensure that we are compatible with the persist
1003        // data format. This lets us return an error gracefully, rather than panicking later in
1004        // persist.
1005        let version_in_catalog_shard =
1006            fetch_catalog_shard_version(&persist_client, catalog_shard_id).await;
1007        if let Some(version_in_catalog_shard) = version_in_catalog_shard {
1008            if !mz_persist_client::cfg::code_can_write_data(&version, &version_in_catalog_shard) {
1009                return Err(DurableCatalogError::IncompatiblePersistVersion {
1010                    found_version: version_in_catalog_shard,
1011                    catalog_version: version,
1012                });
1013            }
1014        }
1015
1016        let open_handles_start = Instant::now();
1017        info!("startup: envd serve: catalog init: open handles beginning");
1018        let since_handle = persist_client
1019            .open_critical_since(
1020                catalog_shard_id,
1021                CATALOG_CRITICAL_SINCE.clone(),
1022                Opaque::encode(&i64::MIN),
1023                Diagnostics {
1024                    shard_name: CATALOG_SHARD_NAME.to_string(),
1025                    handle_purpose: "durable catalog state critical since".to_string(),
1026                },
1027            )
1028            .await
1029            .expect("invalid usage");
1030
1031        let (mut write_handle, mut read_handle) = persist_client
1032            .open(
1033                catalog_shard_id,
1034                Arc::new(persist_desc()),
1035                Arc::new(UnitSchema::default()),
1036                Diagnostics {
1037                    shard_name: CATALOG_SHARD_NAME.to_string(),
1038                    handle_purpose: "durable catalog state handles".to_string(),
1039                },
1040                USE_CRITICAL_SINCE_CATALOG.get(persist_client.dyncfgs()),
1041            )
1042            .await
1043            .expect("invalid usage");
1044        info!(
1045            "startup: envd serve: catalog init: open handles complete in {:?}",
1046            open_handles_start.elapsed()
1047        );
1048
1049        // Commit an empty write at the minimum timestamp so the catalog is always readable.
1050        let upper = {
1051            const EMPTY_UPDATES: &[((SourceData, ()), Timestamp, StorageDiff)] = &[];
1052            let upper = Antichain::from_elem(Timestamp::minimum());
1053            let next_upper = Timestamp::minimum().step_forward();
1054            match write_handle
1055                .compare_and_append(EMPTY_UPDATES, upper, Antichain::from_elem(next_upper))
1056                .await
1057                .expect("invalid usage")
1058            {
1059                Ok(()) => next_upper,
1060                Err(mismatch) => antichain_to_timestamp(mismatch.current),
1061            }
1062        };
1063
1064        let snapshot_start = Instant::now();
1065        info!("startup: envd serve: catalog init: snapshot beginning");
1066        let as_of = as_of(&read_handle, upper);
1067        let snapshot: Vec<_> = snapshot_binary(&mut read_handle, as_of, &metrics)
1068            .await
1069            .map(|StateUpdate { kind, ts, diff }| (kind, ts, diff))
1070            .collect();
1071        let listen = read_handle
1072            .listen(Antichain::from_elem(as_of))
1073            .await
1074            .expect("invalid usage");
1075        info!(
1076            "startup: envd serve: catalog init: snapshot complete in {:?}",
1077            snapshot_start.elapsed()
1078        );
1079
1080        let mut handle = UnopenedPersistCatalogState {
1081            // Unopened catalogs are always writeable until they're opened in an explicit mode.
1082            mode: Mode::Writable,
1083            since_handle,
1084            write_handle,
1085            listen,
1086            persist_client,
1087            shard_id: catalog_shard_id,
1088            // Initialize empty in-memory state.
1089            snapshot: Vec::new(),
1090            update_applier: UnopenedCatalogStateInner::new(),
1091            upper,
1092            fenceable_token: FenceableToken::new(deploy_generation),
1093            catalog_content_version: version,
1094            bootstrap_complete: false,
1095            metrics,
1096            size_at_last_consolidation: None,
1097            updates_applied: 0,
1098        };
1099        // If the snapshot is not consolidated, and we see multiple epoch values while applying the
1100        // updates, then we might accidentally fence ourselves out.
1101        soft_assert_no_log!(
1102            snapshot.iter().all(|(_, _, diff)| *diff == Diff::ONE),
1103            "snapshot should be consolidated: {snapshot:#?}"
1104        );
1105
1106        let apply_start = Instant::now();
1107        info!("startup: envd serve: catalog init: apply updates beginning");
1108        let updates = snapshot
1109            .into_iter()
1110            .map(|(kind, ts, diff)| StateUpdate { kind, ts, diff });
1111        handle.apply_updates_and_consolidate(updates)?;
1112        info!(
1113            "startup: envd serve: catalog init: apply updates complete in {:?}",
1114            apply_start.elapsed()
1115        );
1116
1117        // Validate that the binary version of the current process is not less than any binary
1118        // version that has written to the catalog.
1119        // This condition is only checked once, right here. If a new process comes along with a
1120        // higher version, it must fence this process out with one of the existing fencing
1121        // mechanisms.
1122        if let Some(found_version) = handle.get_catalog_content_version().await? {
1123            // Use cmp_precedence() to ignore build metadata per SemVer 2.0.0 spec
1124            if handle
1125                .catalog_content_version
1126                .cmp_precedence(&found_version)
1127                == std::cmp::Ordering::Less
1128            {
1129                return Err(DurableCatalogError::IncompatiblePersistVersion {
1130                    found_version,
1131                    catalog_version: handle.catalog_content_version,
1132                });
1133            }
1134        }
1135
1136        Ok(handle)
1137    }
1138
1139    #[mz_ore::instrument]
1140    async fn open_inner(
1141        mut self,
1142        mode: Mode,
1143        initial_ts: Timestamp,
1144        bootstrap_args: &BootstrapArgs,
1145    ) -> Result<Box<dyn DurableCatalogState>, CatalogError> {
1146        // It would be nice to use `initial_ts` here, but it comes from the system clock, not the
1147        // timestamp oracle.
1148        let mut commit_ts = self.upper;
1149        self.mode = mode;
1150
1151        // Validate the current deploy generation.
1152        match (&self.mode, &self.fenceable_token) {
1153            (_, FenceableToken::Unfenced { .. } | FenceableToken::Fenced { .. }) => {
1154                return Err(DurableCatalogError::Internal(
1155                    "catalog should not have fenced before opening".to_string(),
1156                )
1157                .into());
1158            }
1159            (
1160                Mode::Writable | Mode::Savepoint,
1161                FenceableToken::Initializing {
1162                    current_deploy_generation: None,
1163                    ..
1164                },
1165            ) => {
1166                return Err(DurableCatalogError::Internal(format!(
1167                    "cannot open in mode '{:?}' without a deploy generation",
1168                    self.mode,
1169                ))
1170                .into());
1171            }
1172            _ => {}
1173        }
1174
1175        let read_only = matches!(self.mode, Mode::Readonly);
1176
1177        // Fence out previous catalogs.
1178        loop {
1179            self.sync_to_current_upper().await?;
1180            commit_ts = max(commit_ts, self.upper);
1181            let (fence_updates, current_fenceable_token) = self
1182                .fenceable_token
1183                .generate_unfenced_token(self.mode)?
1184                .ok_or_else(|| {
1185                    DurableCatalogError::Internal(
1186                        "catalog should not have fenced before opening".to_string(),
1187                    )
1188                })?;
1189            debug!(
1190                ?self.upper,
1191                ?self.fenceable_token,
1192                ?current_fenceable_token,
1193                "fencing previous catalogs"
1194            );
1195            if matches!(self.mode, Mode::Writable) {
1196                match self
1197                    .compare_and_append(fence_updates.clone(), commit_ts)
1198                    .await
1199                {
1200                    Ok(upper) => {
1201                        commit_ts = upper;
1202                    }
1203                    Err(CompareAndAppendError::Fence(e)) => return Err(e.into()),
1204                    Err(e @ CompareAndAppendError::UpperMismatch { .. }) => {
1205                        warn!("catalog write failed due to upper mismatch, retrying: {e:?}");
1206                        continue;
1207                    }
1208                }
1209            }
1210            self.fenceable_token = current_fenceable_token;
1211            break;
1212        }
1213
1214        if matches!(self.mode, Mode::Writable) {
1215            // One-time migration: The catalog previously used `CONTROLLER_CRITICAL_SINCE` for its
1216            // since handle. Now it uses its own `CATALOG_CRITICAL_SINCE`, to free
1217            // `CONTROLLER_CRITICAL_SINCE` up for the storage controller. The catalog and
1218            // controller handles differ in the `Opaque` codec, so we need a migration.
1219            //
1220            // TODO: Remove this once we don't support upgrading from v26 anymore.
1221            let mut controller_handle = self
1222                .persist_client
1223                .open_critical_since::<SourceData, (), Timestamp, StorageDiff>(
1224                    self.shard_id,
1225                    PersistClient::CONTROLLER_CRITICAL_SINCE,
1226                    Opaque::encode(&i64::MIN),
1227                    Diagnostics {
1228                        shard_name: CATALOG_SHARD_NAME.to_string(),
1229                        handle_purpose: "durable catalog state critical since (migration)"
1230                            .to_string(),
1231                    },
1232                )
1233                .await
1234                .expect("invalid usage");
1235
1236            let since = controller_handle.since().clone();
1237            let res = controller_handle
1238                .compare_and_downgrade_since(
1239                    &Opaque::encode(&i64::MIN),
1240                    (&Opaque::encode(&PersistEpoch::default()), &since),
1241                )
1242                .await;
1243            match res {
1244                Ok(_) => info!("migrated Opaque of catalog since handle"),
1245                Err(_) => { /* critical since was already migrated */ }
1246            }
1247        }
1248
1249        let is_initialized = self.is_initialized_inner();
1250        if !matches!(self.mode, Mode::Writable) && !is_initialized {
1251            return Err(CatalogError::Durable(DurableCatalogError::NotWritable(
1252                format!(
1253                    "catalog tables do not exist; will not create in {:?} mode",
1254                    self.mode
1255                ),
1256            )));
1257        }
1258        soft_assert_ne_or_log!(self.upper, Timestamp::minimum());
1259
1260        // Audit log entries are served from `mz_internal.mz_catalog_raw` via
1261        // the `mz_audit_events` materialized view, so they do not need to live
1262        // in the in-memory catalog snapshot. Drop them here and only keep the
1263        // count for metrics.
1264        let (audit_logs, snapshot): (Vec<_>, Vec<_>) = self
1265            .snapshot
1266            .into_iter()
1267            .partition(|(update, _, _)| update.is_audit_log());
1268        self.snapshot = snapshot;
1269        let audit_log_count = audit_logs.iter().map(|(_, _, diff)| diff).sum::<Diff>();
1270        drop(audit_logs);
1271
1272        // Perform data migrations.
1273        if is_initialized && !read_only {
1274            commit_ts = upgrade(&mut self, commit_ts).await?;
1275        }
1276
1277        debug!(
1278            ?is_initialized,
1279            ?self.upper,
1280            "initializing catalog state"
1281        );
1282        let mut catalog = PersistCatalogState {
1283            mode: self.mode,
1284            since_handle: self.since_handle,
1285            write_handle: self.write_handle,
1286            listen: self.listen,
1287            persist_client: self.persist_client,
1288            shard_id: self.shard_id,
1289            upper: self.upper,
1290            fenceable_token: self.fenceable_token,
1291            // Initialize empty in-memory state.
1292            snapshot: Vec::new(),
1293            update_applier: CatalogStateInner::new(),
1294            catalog_content_version: self.catalog_content_version,
1295            bootstrap_complete: false,
1296            metrics: self.metrics,
1297            size_at_last_consolidation: None,
1298            updates_applied: 0,
1299        };
1300        catalog.metrics.collection_entries.reset();
1301        // Normally, `collection_entries` is updated in `apply_updates`. The audit log updates skip
1302        // over that function so we manually update it here.
1303        catalog
1304            .metrics
1305            .collection_entries
1306            .with_label_values(&[&CollectionType::AuditLog.to_string()])
1307            .add(audit_log_count.into_inner());
1308        let updates = self.snapshot.into_iter().map(|(kind, ts, diff)| {
1309            let kind = TryIntoStateUpdateKind::try_into(kind).expect("kind decoding error");
1310            StateUpdate { kind, ts, diff }
1311        });
1312        catalog.apply_updates_and_consolidate(updates)?;
1313
1314        let catalog_content_version = catalog.catalog_content_version.to_string();
1315        let txn = if is_initialized {
1316            let mut txn = catalog.transaction_unchecked().await?;
1317
1318            // Ad-hoc migration: Initialize the `migration_version` expected by adapter to be
1319            // present in existing catalogs.
1320            //
1321            // Note: Need to exclude read-only catalog mode here, because in that mode all
1322            // transactions are expected to be no-ops.
1323            // TODO: remove this once we only support upgrades from version >= 0.164
1324            if txn.get_setting("migration_version".into()).is_none() && mode != Mode::Readonly {
1325                let old_version = txn.get_catalog_content_version();
1326                txn.set_setting("migration_version".into(), old_version.map(Into::into))?;
1327            }
1328
1329            // Opening the catalog with write intent fences out every previous
1330            // catalog owner, so all sessions served by previous owners are
1331            // dead. Reclaim the temporary items they owned here, before
1332            // anything else reads the catalog.
1333            //
1334            // NOTE: This only reclaims on the fence, which covers a
1335            // single-writer world where every crash is followed by some
1336            // process's writable open. Once several serving envds run
1337            // concurrently, a peer crash triggers no fence here, so
1338            // reclaiming its sessions' items needs the durable
1339            // envd-heartbeat mechanism described in the durable temporary
1340            // objects design doc.
1341            if mode != Mode::Readonly {
1342                txn.remove_ephemeral_items();
1343            }
1344
1345            txn.set_catalog_content_version(catalog_content_version)?;
1346            txn
1347        } else {
1348            soft_assert_eq_no_log!(
1349                catalog
1350                    .snapshot
1351                    .iter()
1352                    .filter(|(kind, _, _)| !matches!(kind, StateUpdateKind::FenceToken(_)))
1353                    .count(),
1354                0,
1355                "trace should not contain any updates for an uninitialized catalog: {:#?}",
1356                catalog.snapshot
1357            );
1358
1359            let mut txn = catalog.transaction_unchecked().await?;
1360            initialize::initialize(
1361                &mut txn,
1362                bootstrap_args,
1363                initial_ts.into(),
1364                catalog_content_version,
1365            )
1366            .await?;
1367            txn
1368        };
1369
1370        if read_only {
1371            let (txn_batch, _) = txn.into_parts()?;
1372            // The upper here doesn't matter because we are only applying the updates in memory.
1373            let updates = StateUpdate::from_txn_batch_ts(txn_batch, catalog.upper);
1374            catalog.apply_updates_and_consolidate(updates)?;
1375        } else {
1376            txn.commit_internal(commit_ts).await?;
1377        }
1378
1379        if matches!(catalog.mode, Mode::Writable) {
1380            let write_handle = catalog
1381                .persist_client
1382                .open_writer::<SourceData, (), Timestamp, i64>(
1383                    catalog.write_handle.shard_id(),
1384                    Arc::new(persist_desc()),
1385                    Arc::new(UnitSchema::default()),
1386                    Diagnostics {
1387                        shard_name: CATALOG_SHARD_NAME.to_string(),
1388                        handle_purpose: "compact catalog".to_string(),
1389                    },
1390                )
1391                .await
1392                .expect("invalid usage");
1393            let fuel = CATALOG_FORCE_COMPACTION_FUEL.handle(catalog.persist_client.dyncfgs());
1394            let wait = CATALOG_FORCE_COMPACTION_WAIT.handle(catalog.persist_client.dyncfgs());
1395            // We're going to gradually turn this on via dyncfgs. Run it in a task so that it
1396            // doesn't block startup.
1397            let _task = mz_ore::task::spawn(|| "catalog::force_shard_compaction", async move {
1398                let () =
1399                    mz_persist_client::cli::admin::dangerous_force_compaction_and_break_pushdown(
1400                        &write_handle,
1401                        || fuel.get(),
1402                        || wait.get(),
1403                    )
1404                    .await;
1405            });
1406        }
1407
1408        Ok(Box::new(catalog))
1409    }
1410
1411    /// Reports if the catalog state has been initialized.
1412    ///
1413    /// NOTE: This is the answer as of the last call to [`PersistHandle::sync`] or [`PersistHandle::sync_to_current_upper`],
1414    /// not necessarily what is currently in persist.
1415    #[mz_ore::instrument]
1416    fn is_initialized_inner(&self) -> bool {
1417        !self.update_applier.configs.is_empty()
1418    }
1419
1420    /// Get the current value of config `key`.
1421    ///
1422    /// Some configs need to be read before the catalog is opened for bootstrapping.
1423    #[mz_ore::instrument]
1424    async fn get_current_config(&mut self, key: &str) -> Result<Option<u64>, DurableCatalogError> {
1425        self.sync_to_current_upper().await?;
1426        Ok(self.update_applier.configs.get(key).cloned())
1427    }
1428
1429    /// Get the user version of this instance.
1430    ///
1431    /// The user version is used to determine if a migration is needed.
1432    #[mz_ore::instrument]
1433    pub(crate) async fn get_user_version(&mut self) -> Result<Option<u64>, DurableCatalogError> {
1434        self.get_current_config(USER_VERSION_KEY).await
1435    }
1436
1437    /// Get the current value of setting `name`.
1438    ///
1439    /// Some settings need to be read before the catalog is opened for bootstrapping.
1440    #[mz_ore::instrument]
1441    async fn get_current_setting(
1442        &mut self,
1443        name: &str,
1444    ) -> Result<Option<String>, DurableCatalogError> {
1445        self.sync_to_current_upper().await?;
1446        Ok(self.update_applier.settings.get(name).cloned())
1447    }
1448
1449    /// Get the catalog content version.
1450    ///
1451    /// The catalog content version is the semantic version of the most recent binary that wrote to
1452    /// the catalog.
1453    #[mz_ore::instrument]
1454    async fn get_catalog_content_version(
1455        &mut self,
1456    ) -> Result<Option<semver::Version>, DurableCatalogError> {
1457        let version = self
1458            .get_current_setting(CATALOG_CONTENT_VERSION_KEY)
1459            .await?;
1460        let version = version.map(|version| version.parse().expect("invalid version persisted"));
1461        Ok(version)
1462    }
1463}
1464
1465#[async_trait]
1466impl OpenableDurableCatalogState for UnopenedPersistCatalogState {
1467    #[mz_ore::instrument]
1468    async fn open_savepoint(
1469        mut self: Box<Self>,
1470        initial_ts: Timestamp,
1471        bootstrap_args: &BootstrapArgs,
1472    ) -> Result<Box<dyn DurableCatalogState>, CatalogError> {
1473        self.open_inner(Mode::Savepoint, initial_ts, bootstrap_args)
1474            .boxed()
1475            .await
1476    }
1477
1478    #[mz_ore::instrument]
1479    async fn open_read_only(
1480        mut self: Box<Self>,
1481        bootstrap_args: &BootstrapArgs,
1482    ) -> Result<Box<dyn DurableCatalogState>, CatalogError> {
1483        self.open_inner(Mode::Readonly, EpochMillis::MIN.into(), bootstrap_args)
1484            .boxed()
1485            .await
1486    }
1487
1488    #[mz_ore::instrument]
1489    async fn open(
1490        mut self: Box<Self>,
1491        initial_ts: Timestamp,
1492        bootstrap_args: &BootstrapArgs,
1493    ) -> Result<Box<dyn DurableCatalogState>, CatalogError> {
1494        self.open_inner(Mode::Writable, initial_ts, bootstrap_args)
1495            .boxed()
1496            .await
1497    }
1498
1499    #[mz_ore::instrument(level = "debug")]
1500    async fn open_debug(mut self: Box<Self>) -> Result<DebugCatalogState, CatalogError> {
1501        Ok(DebugCatalogState(*self))
1502    }
1503
1504    #[mz_ore::instrument]
1505    async fn is_initialized(&mut self) -> Result<bool, CatalogError> {
1506        self.sync_to_current_upper().await?;
1507        Ok(self.is_initialized_inner())
1508    }
1509
1510    #[mz_ore::instrument]
1511    async fn epoch(&mut self) -> Result<Epoch, CatalogError> {
1512        self.sync_to_current_upper().await?;
1513        self.fenceable_token
1514            .validate()?
1515            .map(|token| token.epoch)
1516            .ok_or(CatalogError::Durable(DurableCatalogError::Uninitialized))
1517    }
1518
1519    #[mz_ore::instrument]
1520    async fn get_deployment_generation(&mut self) -> Result<u64, CatalogError> {
1521        self.sync_to_current_upper().await?;
1522        self.fenceable_token
1523            .token()
1524            .map(|token| token.deploy_generation)
1525            .ok_or(CatalogError::Durable(DurableCatalogError::Uninitialized))
1526    }
1527
1528    #[mz_ore::instrument(level = "debug")]
1529    async fn get_0dt_deployment_max_wait(&mut self) -> Result<Option<Duration>, CatalogError> {
1530        let value = self
1531            .get_current_config(WITH_0DT_DEPLOYMENT_MAX_WAIT)
1532            .await?;
1533        match value {
1534            None => Ok(None),
1535            Some(millis) => Ok(Some(Duration::from_millis(millis))),
1536        }
1537    }
1538
1539    #[mz_ore::instrument(level = "debug")]
1540    async fn get_0dt_deployment_ddl_check_interval(
1541        &mut self,
1542    ) -> Result<Option<Duration>, CatalogError> {
1543        let value = self
1544            .get_current_config(WITH_0DT_DEPLOYMENT_DDL_CHECK_INTERVAL)
1545            .await?;
1546        match value {
1547            None => Ok(None),
1548            Some(millis) => Ok(Some(Duration::from_millis(millis))),
1549        }
1550    }
1551
1552    #[mz_ore::instrument(level = "debug")]
1553    async fn get_enable_0dt_deployment_panic_after_timeout(
1554        &mut self,
1555    ) -> Result<Option<bool>, CatalogError> {
1556        let value = self
1557            .get_current_config(ENABLE_0DT_DEPLOYMENT_PANIC_AFTER_TIMEOUT)
1558            .await?;
1559        match value {
1560            None => Ok(None),
1561            Some(0) => Ok(Some(false)),
1562            Some(1) => Ok(Some(true)),
1563            Some(v) => Err(
1564                DurableCatalogError::from(TryFromProtoError::UnknownEnumVariant(format!(
1565                    "{v} is not a valid boolean value"
1566                )))
1567                .into(),
1568            ),
1569        }
1570    }
1571
1572    #[mz_ore::instrument]
1573    async fn has_system_config_synced_once(&mut self) -> Result<bool, DurableCatalogError> {
1574        self.get_current_config(SYSTEM_CONFIG_SYNCED_KEY)
1575            .await
1576            .map(|value| value.map(|value| value > 0).unwrap_or(false))
1577    }
1578
1579    #[mz_ore::instrument]
1580    async fn trace_unconsolidated(&mut self) -> Result<Trace, CatalogError> {
1581        self.sync_to_current_upper().await?;
1582        if self.is_initialized_inner() {
1583            let snapshot = self.snapshot_unconsolidated().await;
1584            Ok(Trace::from_snapshot(snapshot))
1585        } else {
1586            Err(CatalogError::Durable(DurableCatalogError::Uninitialized))
1587        }
1588    }
1589
1590    #[mz_ore::instrument]
1591    async fn trace_consolidated(&mut self) -> Result<Trace, CatalogError> {
1592        self.sync_to_current_upper().await?;
1593        if self.is_initialized_inner() {
1594            let snapshot = self.current_snapshot().await?;
1595            Ok(Trace::from_snapshot(snapshot))
1596        } else {
1597            Err(CatalogError::Durable(DurableCatalogError::Uninitialized))
1598        }
1599    }
1600
1601    #[mz_ore::instrument(level = "debug")]
1602    async fn expire(self: Box<Self>) {
1603        self.expire().await
1604    }
1605}
1606
1607/// Applies updates for an opened catalog.
1608#[derive(Debug)]
1609struct CatalogStateInner {
1610    /// A trace of all catalog updates that can be consumed by some higher layer.
1611    updates: VecDeque<memory::objects::StateUpdate>,
1612}
1613
1614impl CatalogStateInner {
1615    fn new() -> CatalogStateInner {
1616        CatalogStateInner {
1617            updates: VecDeque::new(),
1618        }
1619    }
1620}
1621
1622impl ApplyUpdate<StateUpdateKind> for CatalogStateInner {
1623    fn apply_update(
1624        &mut self,
1625        update: StateUpdate<StateUpdateKind>,
1626        current_fence_token: &mut FenceableToken,
1627        metrics: &Arc<Metrics>,
1628    ) -> Result<Option<StateUpdate<StateUpdateKind>>, FenceError> {
1629        if let Some(collection_type) = update.kind.collection_type() {
1630            metrics
1631                .collection_entries
1632                .with_label_values(&[&collection_type.to_string()])
1633                .add(update.diff.into_inner());
1634        }
1635
1636        {
1637            let update: Option<memory::objects::StateUpdate> = (&update)
1638                .try_into()
1639                .expect("invalid persisted update: {update:#?}");
1640            if let Some(update) = update {
1641                self.updates.push_back(update);
1642            }
1643        }
1644
1645        match (update.kind, update.diff) {
1646            (StateUpdateKind::AuditLog(_, ()), _) => Ok(None),
1647            // Nothing to due for fence token retractions but wait for the next insertion.
1648            (StateUpdateKind::FenceToken(_), Diff::MINUS_ONE) => Ok(None),
1649            (StateUpdateKind::FenceToken(token), Diff::ONE) => {
1650                current_fence_token.maybe_fence(token)?;
1651                Ok(None)
1652            }
1653            (kind, diff) => Ok(Some(StateUpdate {
1654                kind,
1655                ts: update.ts,
1656                diff,
1657            })),
1658        }
1659    }
1660}
1661
1662/// A durable store of the catalog state using Persist as an implementation. The durable store can
1663/// serve any catalog data and transactionally modify catalog data.
1664///
1665/// Production users should call [`Self::expire`] before dropping a [`PersistCatalogState`]
1666/// so that it can expire its leases. If/when rust gets AsyncDrop, this will be done automatically.
1667type PersistCatalogState = PersistHandle<StateUpdateKind, CatalogStateInner>;
1668
1669impl PersistHandle<StateUpdateKind, CatalogStateInner> {
1670    /// Creates a transaction without validating pending catalog updates.
1671    async fn transaction_unchecked(&mut self) -> Result<Transaction<'_>, CatalogError> {
1672        self.metrics.transactions_started.inc();
1673        let snapshot = self.snapshot().await?;
1674        let commit_ts = self.upper;
1675        Transaction::new(self, snapshot, commit_ts)
1676    }
1677}
1678
1679#[async_trait]
1680impl ReadOnlyDurableCatalogState for PersistCatalogState {
1681    fn epoch(&self) -> Epoch {
1682        self.fenceable_token
1683            .token()
1684            .expect("opened catalog state must have an epoch")
1685            .epoch
1686    }
1687
1688    fn metrics(&self) -> &Metrics {
1689        &self.metrics
1690    }
1691
1692    #[mz_ore::instrument(level = "debug")]
1693    async fn expire(self: Box<Self>) {
1694        self.expire().await
1695    }
1696
1697    fn is_bootstrap_complete(&self) -> bool {
1698        self.bootstrap_complete
1699    }
1700
1701    async fn get_audit_logs(&mut self) -> Result<Vec<VersionedEvent>, CatalogError> {
1702        self.sync_to_current_upper().await?;
1703        let audit_logs: Vec<_> = self
1704            .persist_snapshot()
1705            .await
1706            .filter_map(
1707                |StateUpdate {
1708                     kind,
1709                     ts: _,
1710                     diff: _,
1711                 }| match kind {
1712                    StateUpdateKind::AuditLog(key, ()) => Some(key),
1713                    _ => None,
1714                },
1715            )
1716            .collect();
1717        let mut audit_logs: Vec<_> = audit_logs
1718            .into_iter()
1719            .map(RustType::from_proto)
1720            .map_ok(|key: AuditLogKey| key.event)
1721            .collect::<Result<_, _>>()?;
1722        audit_logs.sort_by(|a, b| a.sortable_id().cmp(&b.sortable_id()));
1723        Ok(audit_logs)
1724    }
1725
1726    #[mz_ore::instrument(level = "debug")]
1727    async fn get_next_id(&mut self, id_type: &str) -> Result<u64, CatalogError> {
1728        self.with_trace(|trace| {
1729            Ok(trace
1730                .into_iter()
1731                .rev()
1732                .filter_map(|(kind, _, _)| match kind {
1733                    StateUpdateKind::IdAllocator(key, value) if key.name == id_type => {
1734                        Some(value.next_id)
1735                    }
1736                    _ => None,
1737                })
1738                .next()
1739                .expect("must exist"))
1740        })
1741        .await
1742    }
1743
1744    #[mz_ore::instrument(level = "debug")]
1745    async fn get_deployment_generation(&mut self) -> Result<u64, CatalogError> {
1746        self.sync_to_current_upper().await?;
1747        Ok(self
1748            .fenceable_token
1749            .token()
1750            .expect("opened catalogs must have a token")
1751            .deploy_generation)
1752    }
1753
1754    #[mz_ore::instrument(level = "debug")]
1755    async fn snapshot(&mut self) -> Result<Snapshot, CatalogError> {
1756        self.with_snapshot(Ok).await
1757    }
1758
1759    #[mz_ore::instrument(level = "debug")]
1760    async fn sync_to_current_updates(
1761        &mut self,
1762    ) -> Result<Vec<memory::objects::StateUpdate>, CatalogError> {
1763        let upper = self.current_upper().await;
1764        self.sync_updates(upper).await
1765    }
1766
1767    #[mz_ore::instrument(level = "debug")]
1768    async fn sync_updates(
1769        &mut self,
1770        target_upper: mz_repr::Timestamp,
1771    ) -> Result<Vec<memory::objects::StateUpdate>, CatalogError> {
1772        self.sync(target_upper).await?;
1773        let mut updates = Vec::new();
1774        while let Some(update) = self.update_applier.updates.front() {
1775            if update.ts >= target_upper {
1776                break;
1777            }
1778
1779            let update = self
1780                .update_applier
1781                .updates
1782                .pop_front()
1783                .expect("peeked above");
1784            updates.push(update);
1785        }
1786        Ok(updates)
1787    }
1788
1789    #[mz_ore::instrument(level = "debug")]
1790    async fn ensure_not_out_of_sync(
1791        &mut self,
1792        target_upper: Timestamp,
1793    ) -> Result<(), CatalogError> {
1794        self.sync(target_upper).await?;
1795        let update_count = self
1796            .update_applier
1797            .updates
1798            .iter()
1799            .take_while(|update| update.ts < target_upper)
1800            .count();
1801        if update_count == 0 {
1802            Ok(())
1803        } else {
1804            Err(DurableCatalogError::CatalogOutOfSync {
1805                update_count,
1806                upper: target_upper,
1807            }
1808            .into())
1809        }
1810    }
1811
1812    async fn current_upper(&mut self) -> Timestamp {
1813        self.current_upper().await
1814    }
1815}
1816
1817#[async_trait]
1818#[allow(mismatched_lifetime_syntaxes)]
1819impl DurableCatalogState for PersistCatalogState {
1820    fn is_read_only(&self) -> bool {
1821        matches!(self.mode, Mode::Readonly)
1822    }
1823
1824    fn is_savepoint(&self) -> bool {
1825        matches!(self.mode, Mode::Savepoint)
1826    }
1827
1828    async fn mark_bootstrap_complete(&mut self) {
1829        self.bootstrap_complete = true;
1830        if matches!(self.mode, Mode::Writable) {
1831            self.since_handle
1832                .upgrade_version()
1833                .await
1834                .expect("invalid usage")
1835        }
1836    }
1837
1838    #[mz_ore::instrument(level = "debug")]
1839    async fn transaction(&mut self) -> Result<Transaction, CatalogError> {
1840        let mut txn = self.transaction_unchecked().await?;
1841        txn.ensure_not_out_of_sync().await?;
1842        Ok(txn)
1843    }
1844
1845    fn transaction_from_snapshot(
1846        &mut self,
1847        snapshot: Snapshot,
1848    ) -> Result<DryRunTransaction, CatalogError> {
1849        let commit_ts = self.upper;
1850        Transaction::new(self, snapshot, commit_ts).map(DryRunTransaction::new)
1851    }
1852
1853    #[mz_ore::instrument(level = "debug")]
1854    async fn allocate_id(
1855        &mut self,
1856        id_type: &str,
1857        amount: u64,
1858        commit_ts: Timestamp,
1859    ) -> Result<Vec<u64>, CatalogError> {
1860        let start = Instant::now();
1861        if amount == 0 {
1862            return Ok(Vec::new());
1863        }
1864        let mut txn = self.transaction_unchecked().await?;
1865        let ids = txn.get_and_increment_id_by(id_type.to_string(), amount)?;
1866        txn.commit_internal(commit_ts).await?;
1867        self.metrics
1868            .allocate_id_seconds
1869            .observe(start.elapsed().as_secs_f64());
1870        Ok(ids)
1871    }
1872
1873    #[mz_ore::instrument(level = "debug")]
1874    async fn commit_transaction(
1875        &mut self,
1876        txn_batch: TransactionBatch,
1877        commit_ts: Timestamp,
1878    ) -> Result<Timestamp, CatalogError> {
1879        async fn commit_transaction_inner(
1880            catalog: &mut PersistCatalogState,
1881            txn_batch: TransactionBatch,
1882            commit_ts: Timestamp,
1883        ) -> Result<Timestamp, CatalogError> {
1884            // If the transaction is empty then we don't error, even in read-only mode.
1885            // This is mostly for legacy reasons (i.e. with enough elbow grease this
1886            // behavior can be changed without breaking any fundamental assumptions).
1887            if catalog.mode == Mode::Readonly {
1888                let updates: Vec<_> = StateUpdate::from_txn_batch(txn_batch).collect();
1889                if !updates.is_empty() {
1890                    let collection_types: Vec<_> = updates
1891                        .iter()
1892                        .filter_map(|u| u.0.collection_type())
1893                        .collect();
1894                    return Err(DurableCatalogError::NotWritable(format!(
1895                        "cannot commit a transaction in a read-only catalog: \
1896                         {} updates across collections: {collection_types:?}",
1897                        updates.len(),
1898                    ))
1899                    .into());
1900                }
1901                return Ok(catalog.upper);
1902            }
1903
1904            // The handle is mutex-protected from transaction creation through commit, so only a
1905            // local programming error can change its cached upper here.
1906            assert_eq!(
1907                catalog.upper, txn_batch.upper,
1908                "the handle was mutated mid-transaction"
1909            );
1910
1911            let updates: Vec<_> = StateUpdate::from_txn_batch(txn_batch).collect();
1912            debug!("committing updates: {updates:?}");
1913
1914            // Empty upper progress does not invalidate the transaction.
1915            let mut commit_ts = max(commit_ts, catalog.upper);
1916
1917            let next_upper = match catalog.mode {
1918                Mode::Writable => loop {
1919                    let updates_applied_before = catalog.updates_applied;
1920                    match catalog.compare_and_append(updates.clone(), commit_ts).await {
1921                        Ok(next_upper) => break next_upper,
1922                        Err(CompareAndAppendError::Fence(e)) => return Err(e.into()),
1923                        Err(CompareAndAppendError::UpperMismatch { actual_upper, .. }) => {
1924                            // The mismatch synchronized the handle. Retry only if it applied no
1925                            // content.
1926                            catalog
1927                                .classify_upper_mismatch(updates_applied_before, actual_upper)?;
1928                            commit_ts = max(commit_ts, catalog.upper);
1929                        }
1930                    }
1931                },
1932                Mode::Savepoint => {
1933                    let updates = updates.into_iter().map(|(kind, diff)| StateUpdate {
1934                        kind,
1935                        ts: commit_ts,
1936                        diff,
1937                    });
1938                    catalog.apply_updates_and_consolidate(updates)?;
1939                    catalog.upper = commit_ts.step_forward();
1940                    catalog.upper
1941                }
1942                Mode::Readonly => unreachable!("handled above"),
1943            };
1944
1945            Ok(next_upper)
1946        }
1947        self.metrics.transaction_commits.inc();
1948        let histogram = self.metrics.transaction_commit_latency_seconds.clone();
1949        commit_transaction_inner(self, txn_batch, commit_ts)
1950            .wall_time()
1951            .observe(histogram)
1952            .await
1953    }
1954
1955    #[mz_ore::instrument(level = "debug")]
1956    async fn advance_upper(&mut self, new_upper: Timestamp) -> Result<(), CatalogError> {
1957        loop {
1958            if self.upper >= new_upper {
1959                // This does not consult Persist. It only revalidates a fence already cached by
1960                // this handle, since a sync can advance `upper` while recording the fence.
1961                self.fenceable_token.validate()?;
1962                return Ok(());
1963            }
1964
1965            match self.mode {
1966                Mode::Writable => {}
1967                Mode::Savepoint => {
1968                    self.upper = new_upper;
1969                    return Ok(());
1970                }
1971                Mode::Readonly => {
1972                    return Err(DurableCatalogError::NotWritable(
1973                        "cannot advance upper of a read-only catalog".into(),
1974                    )
1975                    .into());
1976                }
1977            }
1978
1979            let updates_applied_before = self.updates_applied;
1980            match self.compare_and_append_inner([], new_upper).await {
1981                Ok(()) => {
1982                    self.upper = new_upper;
1983                    // No sync needed since no data was written.
1984                    return Ok(());
1985                }
1986                Err(CompareAndAppendError::Fence(e)) => return Err(e.into()),
1987                Err(CompareAndAppendError::UpperMismatch { actual_upper, .. }) => {
1988                    // The mismatch synchronized the handle. Retry only if it applied no content.
1989                    self.classify_upper_mismatch(updates_applied_before, actual_upper)?;
1990                }
1991            }
1992        }
1993    }
1994
1995    fn shard_id(&self) -> ShardId {
1996        self.shard_id
1997    }
1998}
1999
2000/// Deterministically generate a shard ID for the given `organization_id` and `seed`.
2001pub fn shard_id(organization_id: Uuid, seed: usize) -> ShardId {
2002    let hash = sha2::Sha256::digest(format!("{organization_id}{seed}")).to_vec();
2003    soft_assert_eq_or_log!(hash.len(), 32, "SHA256 returns 32 bytes (256 bits)");
2004    let uuid = Uuid::from_slice(&hash[0..16]).expect("from_slice accepts exactly 16 bytes");
2005    ShardId::from_str(&format!("s{uuid}")).expect("known to be valid")
2006}
2007
2008/// Generates a timestamp for reading from `read_handle` that is as fresh as possible, given
2009/// `upper`.
2010fn as_of(
2011    read_handle: &ReadHandle<SourceData, (), Timestamp, StorageDiff>,
2012    upper: Timestamp,
2013) -> Timestamp {
2014    let since = read_handle.since().clone();
2015    let mut as_of = upper.checked_sub(1).unwrap_or_else(|| {
2016        panic!("catalog persist shard should be initialize, found upper: {upper:?}")
2017    });
2018    // We only downgrade the since after writing, and we always set the since to one less than the
2019    // upper.
2020    soft_assert_or_log!(
2021        since.less_equal(&as_of),
2022        "since={since:?}, as_of={as_of:?}; since must be less than or equal to as_of"
2023    );
2024    // This should be a no-op if the assert above passes, however if it doesn't then we'd like to
2025    // continue with a correct timestamp instead of entering a panic loop.
2026    as_of.advance_by(since.borrow());
2027    as_of
2028}
2029
2030/// Fetch the persist version of the catalog shard, if one exists. A version will not
2031/// exist if we are creating a brand-new environment.
2032async fn fetch_catalog_shard_version(
2033    persist_client: &PersistClient,
2034    catalog_shard_id: ShardId,
2035) -> Option<semver::Version> {
2036    let shard_state = persist_client
2037        .inspect_shard::<Timestamp>(&catalog_shard_id)
2038        .await
2039        .ok()?;
2040    let json_state = serde_json::to_value(shard_state).expect("state serialization error");
2041    let json_version = json_state
2042        .get("applier_version")
2043        .cloned()
2044        .expect("missing applier_version");
2045    let version = serde_json::from_value(json_version).expect("version deserialization error");
2046    Some(version)
2047}
2048
2049/// Generates an iterator of [`StateUpdate`] that contain all updates to the catalog
2050/// state up to, and including, `as_of`.
2051///
2052/// The output is consolidated and sorted by timestamp in ascending order.
2053#[mz_ore::instrument(level = "debug")]
2054async fn snapshot_binary(
2055    read_handle: &mut ReadHandle<SourceData, (), Timestamp, StorageDiff>,
2056    as_of: Timestamp,
2057    metrics: &Arc<Metrics>,
2058) -> impl Iterator<Item = StateUpdate<StateUpdateKindJson>> + DoubleEndedIterator + use<> {
2059    metrics.snapshots_taken.inc();
2060    let histogram = metrics.snapshot_latency_seconds.clone();
2061    snapshot_binary_inner(read_handle, as_of)
2062        .wall_time()
2063        .observe(histogram)
2064        .await
2065}
2066
2067/// Generates an iterator of [`StateUpdate`] that contain all updates to the catalog
2068/// state up to, and including, `as_of`.
2069///
2070/// The output is consolidated and sorted by timestamp in ascending order.
2071#[mz_ore::instrument(level = "debug")]
2072async fn snapshot_binary_inner(
2073    read_handle: &mut ReadHandle<SourceData, (), Timestamp, StorageDiff>,
2074    as_of: Timestamp,
2075) -> impl Iterator<Item = StateUpdate<StateUpdateKindJson>> + DoubleEndedIterator + use<> {
2076    let snapshot = read_handle
2077        .snapshot_and_fetch(Antichain::from_elem(as_of))
2078        .await
2079        .expect("we have advanced the restart_as_of by the since");
2080    soft_assert_no_log!(
2081        snapshot.iter().all(|(_, _, diff)| *diff == 1),
2082        "snapshot_and_fetch guarantees a consolidated result: {snapshot:#?}"
2083    );
2084    snapshot
2085        .into_iter()
2086        .map(Into::<StateUpdate<StateUpdateKindJson>>::into)
2087        .sorted_by(|a, b| Ord::cmp(&b.ts, &a.ts))
2088}
2089
2090/// Convert an [`Antichain<Timestamp>`] to a [`Timestamp`].
2091///
2092/// The correctness of this function relies on [`Timestamp`] being totally ordered and never
2093/// finalizing the catalog shard.
2094pub(crate) fn antichain_to_timestamp(antichain: Antichain<Timestamp>) -> Timestamp {
2095    antichain
2096        .into_option()
2097        .expect("we use a totally ordered time and never finalize the shard")
2098}
2099
2100// Debug methods used by the catalog-debug tool.
2101
2102impl Trace {
2103    /// Generates a [`Trace`] from snapshot.
2104    fn from_snapshot(snapshot: impl IntoIterator<Item = StateUpdate>) -> Trace {
2105        let mut trace = Trace::new();
2106        for StateUpdate { kind, ts, diff } in snapshot {
2107            match kind {
2108                StateUpdateKind::AuditLog(k, v) => trace.audit_log.values.push(((k, v), ts, diff)),
2109                StateUpdateKind::Cluster(k, v) => trace.clusters.values.push(((k, v), ts, diff)),
2110                StateUpdateKind::ClusterReplica(k, v) => {
2111                    trace.cluster_replicas.values.push(((k, v), ts, diff))
2112                }
2113                StateUpdateKind::Comment(k, v) => trace.comments.values.push(((k, v), ts, diff)),
2114                StateUpdateKind::Config(k, v) => trace.configs.values.push(((k, v), ts, diff)),
2115                StateUpdateKind::Database(k, v) => trace.databases.values.push(((k, v), ts, diff)),
2116                StateUpdateKind::DefaultPrivilege(k, v) => {
2117                    trace.default_privileges.values.push(((k, v), ts, diff))
2118                }
2119                StateUpdateKind::FenceToken(_) => {
2120                    // Fence token not included in trace.
2121                }
2122                StateUpdateKind::IdAllocator(k, v) => {
2123                    trace.id_allocator.values.push(((k, v), ts, diff))
2124                }
2125                StateUpdateKind::IntrospectionSourceIndex(k, v) => {
2126                    trace.introspection_sources.values.push(((k, v), ts, diff))
2127                }
2128                StateUpdateKind::Item(k, v) => trace.items.values.push(((k, v), ts, diff)),
2129                StateUpdateKind::NetworkPolicy(k, v) => {
2130                    trace.network_policies.values.push(((k, v), ts, diff))
2131                }
2132                StateUpdateKind::Role(k, v) => trace.roles.values.push(((k, v), ts, diff)),
2133                StateUpdateKind::Schema(k, v) => trace.schemas.values.push(((k, v), ts, diff)),
2134                StateUpdateKind::Setting(k, v) => trace.settings.values.push(((k, v), ts, diff)),
2135                StateUpdateKind::SourceReferences(k, v) => {
2136                    trace.source_references.values.push(((k, v), ts, diff))
2137                }
2138                StateUpdateKind::SystemConfiguration(k, v) => {
2139                    trace.system_configurations.values.push(((k, v), ts, diff))
2140                }
2141                StateUpdateKind::ClusterSystemConfiguration(k, v) => trace
2142                    .cluster_system_configurations
2143                    .values
2144                    .push(((k, v), ts, diff)),
2145                StateUpdateKind::ReplicaSystemConfiguration(k, v) => trace
2146                    .replica_system_configurations
2147                    .values
2148                    .push(((k, v), ts, diff)),
2149                StateUpdateKind::SystemObjectMapping(k, v) => {
2150                    trace.system_object_mappings.values.push(((k, v), ts, diff))
2151                }
2152                StateUpdateKind::SystemPrivilege(k, v) => {
2153                    trace.system_privileges.values.push(((k, v), ts, diff))
2154                }
2155                StateUpdateKind::StorageCollectionMetadata(k, v) => trace
2156                    .storage_collection_metadata
2157                    .values
2158                    .push(((k, v), ts, diff)),
2159                StateUpdateKind::UnfinalizedShard(k, ()) => {
2160                    trace.unfinalized_shards.values.push(((k, ()), ts, diff))
2161                }
2162                StateUpdateKind::TxnWalShard((), v) => {
2163                    trace.txn_wal_shard.values.push((((), v), ts, diff))
2164                }
2165                StateUpdateKind::RoleAuth(k, v) => trace.role_auth.values.push(((k, v), ts, diff)),
2166            }
2167        }
2168        trace
2169    }
2170}
2171
2172impl UnopenedPersistCatalogState {
2173    /// Manually update value of `key` in collection `T` to `value`.
2174    #[mz_ore::instrument]
2175    pub(crate) async fn debug_edit<T: Collection>(
2176        &mut self,
2177        key: T::Key,
2178        value: T::Value,
2179    ) -> Result<Option<T::Value>, CatalogError>
2180    where
2181        T::Key: PartialEq + Eq + Debug + Clone,
2182        T::Value: Debug + Clone,
2183    {
2184        let prev_value = loop {
2185            let key = key.clone();
2186            let value = value.clone();
2187            let snapshot = self.current_snapshot().await?;
2188            let trace = Trace::from_snapshot(snapshot);
2189            let collection_trace = T::collection_trace(trace);
2190            let prev_values: Vec<_> = collection_trace
2191                .values
2192                .into_iter()
2193                .filter(|((k, _), _, diff)| {
2194                    soft_assert_eq_or_log!(*diff, Diff::ONE, "trace is consolidated");
2195                    &key == k
2196                })
2197                .collect();
2198
2199            let prev_value = match &prev_values[..] {
2200                [] => None,
2201                [((_, v), _, _)] => Some(v.clone()),
2202                prev_values => panic!("multiple values found for key {key:?}: {prev_values:?}"),
2203            };
2204
2205            let mut updates: Vec<_> = prev_values
2206                .into_iter()
2207                .map(|((k, v), _, _)| (T::update(k, v), Diff::MINUS_ONE))
2208                .collect();
2209            updates.push((T::update(key, value), Diff::ONE));
2210            // We must fence out all other catalogs, if we haven't already, since we are writing.
2211            match self.fenceable_token.generate_unfenced_token(self.mode)? {
2212                Some((fence_updates, current_fenceable_token)) => {
2213                    updates.extend(fence_updates.clone());
2214                    match self.compare_and_append(updates, self.upper).await {
2215                        Ok(_) => {
2216                            self.fenceable_token = current_fenceable_token;
2217                            break prev_value;
2218                        }
2219                        Err(CompareAndAppendError::Fence(e)) => return Err(e.into()),
2220                        Err(e @ CompareAndAppendError::UpperMismatch { .. }) => {
2221                            warn!("catalog write failed due to upper mismatch, retrying: {e:?}");
2222                            continue;
2223                        }
2224                    }
2225                }
2226                None => {
2227                    self.compare_and_append(updates, self.upper)
2228                        .await
2229                        .map_err(|e| e.unwrap_fence_error())?;
2230                    break prev_value;
2231                }
2232            }
2233        };
2234        Ok(prev_value)
2235    }
2236
2237    /// Manually delete `key` from collection `T`.
2238    #[mz_ore::instrument]
2239    pub(crate) async fn debug_delete<T: Collection>(
2240        &mut self,
2241        key: T::Key,
2242    ) -> Result<(), CatalogError>
2243    where
2244        T::Key: PartialEq + Eq + Debug + Clone,
2245        T::Value: Debug,
2246    {
2247        loop {
2248            let key = key.clone();
2249            let snapshot = self.current_snapshot().await?;
2250            let trace = Trace::from_snapshot(snapshot);
2251            let collection_trace = T::collection_trace(trace);
2252            let mut retractions: Vec<_> = collection_trace
2253                .values
2254                .into_iter()
2255                .filter(|((k, _), _, diff)| {
2256                    soft_assert_eq_or_log!(*diff, Diff::ONE, "trace is consolidated");
2257                    &key == k
2258                })
2259                .map(|((k, v), _, _)| (T::update(k, v), Diff::MINUS_ONE))
2260                .collect();
2261
2262            // We must fence out all other catalogs, if we haven't already, since we are writing.
2263            match self.fenceable_token.generate_unfenced_token(self.mode)? {
2264                Some((fence_updates, current_fenceable_token)) => {
2265                    retractions.extend(fence_updates.clone());
2266                    match self.compare_and_append(retractions, self.upper).await {
2267                        Ok(_) => {
2268                            self.fenceable_token = current_fenceable_token;
2269                            break;
2270                        }
2271                        Err(CompareAndAppendError::Fence(e)) => return Err(e.into()),
2272                        Err(e @ CompareAndAppendError::UpperMismatch { .. }) => {
2273                            warn!("catalog write failed due to upper mismatch, retrying: {e:?}");
2274                            continue;
2275                        }
2276                    }
2277                }
2278                None => {
2279                    self.compare_and_append(retractions, self.upper)
2280                        .await
2281                        .map_err(|e| e.unwrap_fence_error())?;
2282                    break;
2283                }
2284            }
2285        }
2286        Ok(())
2287    }
2288
2289    /// Generates a [`Vec<StateUpdate>`] that contain all updates to the catalog
2290    /// state.
2291    ///
2292    /// The output is consolidated and sorted by timestamp in ascending order and the current upper.
2293    async fn current_snapshot(
2294        &mut self,
2295    ) -> Result<impl IntoIterator<Item = StateUpdate> + '_, CatalogError> {
2296        self.sync_to_current_upper().await?;
2297        self.consolidate();
2298        Ok(self.snapshot.iter().cloned().map(|(kind, ts, diff)| {
2299            let kind = TryIntoStateUpdateKind::try_into(kind).expect("kind decoding error");
2300            StateUpdate { kind, ts, diff }
2301        }))
2302    }
2303}