1#[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
71pub(crate) type Timestamp = mz_repr::Timestamp;
73
74const MIN_EPOCH: Epoch = Epoch::new(1).expect("1 is non-zero");
76
77const CATALOG_SHARD_NAME: &str = "catalog";
79
80static CATALOG_CRITICAL_SINCE: LazyLock<CriticalReaderId> = LazyLock::new(|| {
83 "c55555555-6666-7777-8888-999999999999"
84 .parse()
85 .expect("valid CriticalReaderId")
86});
87
88const CATALOG_SEED: usize = 1;
90const _UPGRADE_SEED: usize = 2;
92pub const _BUILTIN_MIGRATION_SEED: usize = 3;
94pub const _EXPRESSION_CACHE_SEED: usize = 4;
96
97#[derive(Debug, Copy, Clone, Eq, PartialEq)]
99pub(crate) enum Mode {
100 Readonly,
102 Savepoint,
104 Writable,
106}
107
108#[derive(Debug)]
110pub(crate) enum FenceableToken {
111 Initializing {
114 durable_token: Option<FenceToken>,
116 current_deploy_generation: Option<u64>,
118 },
119 Unfenced { current_token: FenceToken },
121 Fenced {
123 current_token: FenceToken,
124 fence_token: FenceToken,
125 },
126}
127
128impl FenceableToken {
129 fn new(current_deploy_generation: Option<u64>) -> Self {
131 Self::Initializing {
132 durable_token: None,
133 current_deploy_generation,
134 }
135 }
136
137 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 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 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 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 .ok_or(DurableCatalogError::Uninitialized)?;
252 let mut current_epoch = durable_token
253 .map(|token| token.epoch)
254 .unwrap_or(MIN_EPOCH)
255 .get();
256 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#[derive(Debug, thiserror::Error)]
279pub(crate) enum CompareAndAppendError {
280 #[error(transparent)]
281 Fence(#[from] FenceError),
282 #[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 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#[derive(Debug)]
341pub(crate) struct PersistHandle<T: TryIntoStateUpdateKind, U: ApplyUpdate<T>> {
342 pub(crate) mode: Mode,
344 since_handle: SinceHandle<SourceData, (), Timestamp, StorageDiff>,
346 write_handle: WriteHandle<SourceData, (), Timestamp, StorageDiff>,
348 listen: Listen<SourceData, (), Timestamp, StorageDiff>,
350 persist_client: PersistClient,
352 shard_id: ShardId,
354 pub(crate) snapshot: Vec<(T, Timestamp, Diff)>,
358 update_applier: U,
360 pub(crate) upper: Timestamp,
362 fenceable_token: FenceableToken,
364 catalog_content_version: semver::Version,
366 bootstrap_complete: bool,
368 metrics: Arc<Metrics>,
370 size_at_last_consolidation: Option<usize>,
374 updates_applied: u64,
379}
380
381impl<T: TryIntoStateUpdateKind, U: ApplyUpdate<T>> PersistHandle<T, U> {
382 #[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 #[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 self.compare_and_append_inner(updates, next_upper).await?;
414
415 self.sync(next_upper).await?;
416 Ok(next_upper)
417 }
418
419 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 self.sync_to_current_upper().await?;
454 return Err(e.into());
455 }
456
457 let downgrade_to = Antichain::from_elem(next_upper.saturating_sub(1));
459
460 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 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 #[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 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 #[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 #[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 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 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 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 if self.updates_applied != updates_applied_before {
602 self.consolidate();
603 }
604 Ok(())
605 }
606
607 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 #[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 differential_dataflow::consolidation::consolidate_updates(&mut updates);
643
644 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 Err(err) => errors.push(err),
666 }
667 }
668
669 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 fn maybe_consolidate(&mut self) {
688 let threshold = *self
689 .size_at_last_consolidation
690 .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 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 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 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 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 }
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 }
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 #[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#[derive(Debug)]
908pub(crate) struct UnopenedCatalogStateInner {
909 configs: BTreeMap<String, u64>,
911 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
975pub(crate) type UnopenedPersistCatalogState =
983 PersistHandle<StateUpdateKindJson, UnopenedCatalogStateInner>;
984
985impl UnopenedPersistCatalogState {
986 #[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 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 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 mode: Mode::Writable,
1083 since_handle,
1084 write_handle,
1085 listen,
1086 persist_client,
1087 shard_id: catalog_shard_id,
1088 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 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 if let Some(found_version) = handle.get_catalog_content_version().await? {
1123 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 let mut commit_ts = self.upper;
1149 self.mode = mode;
1150
1151 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 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 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(_) => { }
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 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 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 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 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 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 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 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 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 #[mz_ore::instrument]
1416 fn is_initialized_inner(&self) -> bool {
1417 !self.update_applier.configs.is_empty()
1418 }
1419
1420 #[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 #[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 #[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 #[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#[derive(Debug)]
1609struct CatalogStateInner {
1610 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 (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
1662type PersistCatalogState = PersistHandle<StateUpdateKind, CatalogStateInner>;
1668
1669impl PersistHandle<StateUpdateKind, CatalogStateInner> {
1670 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 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 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 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 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 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 return Ok(());
1985 }
1986 Err(CompareAndAppendError::Fence(e)) => return Err(e.into()),
1987 Err(CompareAndAppendError::UpperMismatch { actual_upper, .. }) => {
1988 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
2000pub 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
2008fn 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 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 as_of.advance_by(since.borrow());
2027 as_of
2028}
2029
2030async 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#[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#[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
2090pub(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
2100impl Trace {
2103 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 }
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 #[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 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 #[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 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 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}