mz_catalog/durable/upgrade.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//! This module contains all the helpers and code paths for upgrading/migrating the `Catalog`.
11//!
12//! We facilitate migrations by keeping snapshots of the objects we previously stored, and relying
13//! entirely on these snapshots. These snapshots exist in the [`mz_catalog_protos`] crate in the
14//! form of `catalog-protos/protos/objects_vXX.proto`. By maintaining and relying on snapshots we
15//! don't have to worry about changes elsewhere in the codebase effecting our migrations because
16//! our application and serialization logic is decoupled, and the objects of the Catalog for a
17//! given version are "frozen in time".
18//!
19//! > **Note**: The protobuf snapshot files themselves live in a separate crate because it takes a
20//! relatively significant amount of time to codegen and build them. By placing them in
21//! a separate crate we don't have to pay this compile time cost when building the
22//! catalog, allowing for faster iteration.
23//!
24//! You cannot make any changes to the following message or anything that they depend on:
25//!
26//! - Config
27//! - Setting
28//! - FenceToken
29//! - AuditLog
30//!
31//! When you want to make a change to the `Catalog` you need to follow these steps:
32//!
33//! 1. Check the current [`CATALOG_VERSION`], make sure an `objects_v<CATALOG_VERSION>.proto` file
34//! exists. If one doesn't, copy and paste the current `objects.proto` file, renaming it to
35//! `objects_v<CATALOG_VERSION>.proto`.
36//! 2. Bump [`CATALOG_VERSION`] by one.
37//! 3. Make your changes to `objects.proto`.
38//! 4. Copy and paste `objects.proto`, naming the copy `objects_v<CATALOG_VERSION>.proto`. Update
39//! the package name of the `.proto` to `package objects_v<CATALOG_VERSION>;`
40//! 5. We should now have a copy of the protobuf objects as they currently exist, and a copy of
41//! how we want them to exist. For example, if the version of the Catalog before we made our
42//! changes was 15, we should now have `objects_v15.proto` and `objects_v16.proto`.
43//! 6. Rebuild Materialize which will error because the hashes stored in
44//! `src/catalog-protos/protos/hashes.json` have now changed. Update these to match the new
45//! hashes for objects.proto and `objects_v<CATALOG_VERSION>.proto`.
46//! 7. Add `v<CATALOG_VERSION>` to the call to the `objects!` macro in this file
47//! 8. Add a new file to `catalog/src/durable/upgrade` which is where we'll put the new migration
48//! path.
49//! 9. Write upgrade functions using the two versions of the protos we now have, e.g.
50//! `objects_v15.proto` and `objects_v16.proto`. In this migration code you __should not__
51//! import any defaults or constants from elsewhere in the codebase, because then a future
52//! change could then impact a previous migration.
53//! 10. Add an import for your new module to this file: mod v<CATALOG_VERSION-1>_to_v<CATALOG_VERSION>;
54//! 11. Call your upgrade function in [`run_upgrade()`].
55//! 12. Generate a test file for the new version:
56//! ```ignore
57//! cargo test --package mz-catalog --lib durable::upgrade::tests::generate_missing_encodings -- --ignored
58//! ```
59//!
60//! When in doubt, reach out to the Surfaces team, and we'll be more than happy to help :)
61
62pub mod json_compatible;
63#[cfg(test)]
64mod tests;
65
66use mz_ore::{soft_assert_eq_or_log, soft_assert_ne_or_log};
67use mz_repr::Diff;
68use paste::paste;
69#[cfg(test)]
70use proptest::prelude::*;
71#[cfg(test)]
72use proptest::strategy::ValueTree;
73#[cfg(test)]
74use proptest_derive::Arbitrary;
75use timely::progress::Timestamp as TimelyTimestamp;
76
77use crate::durable::initialize::USER_VERSION_KEY;
78use crate::durable::objects::serialization::proto;
79use crate::durable::objects::state_update::{
80 IntoStateUpdateKindJson, StateUpdate, StateUpdateKind, StateUpdateKindJson,
81};
82use crate::durable::persist::{Mode, Timestamp, UnopenedPersistCatalogState};
83use crate::durable::{CatalogError, DurableCatalogError};
84
85#[cfg(test)]
86const ENCODED_TEST_CASES: usize = 100;
87
88/// Generate per-version support code.
89///
90/// Here we have to deal with the fact that the pre-v79 objects had a protobuf-generated format,
91/// which gives them additional levels of indirection that the post-v79 objects don't have and thus
92/// requires slightly different code to be generated.
93macro_rules! objects {
94 ( [$( $x_old:ident ),*], [$( $x:ident ),*] ) => {
95 paste! {
96 $(
97 pub(crate) mod [<objects_ $x_old>] {
98 pub use mz_catalog_protos::[<objects_ $x_old>]::*;
99
100 use crate::durable::objects::state_update::StateUpdateKindJson;
101
102 impl From<StateUpdateKind> for StateUpdateKindJson {
103 fn from(value: StateUpdateKind) -> Self {
104 let kind = value.kind.expect("kind should be set");
105 // TODO: This requires that the json->proto->json roundtrips
106 // exactly, see database-issues#7179.
107 StateUpdateKindJson::from_serde(&kind)
108 }
109 }
110
111 impl From<StateUpdateKindJson> for StateUpdateKind {
112 fn from(value: StateUpdateKindJson) -> Self {
113 let kind: state_update_kind::Kind = value.to_serde();
114 StateUpdateKind { kind: Some(kind) }
115 }
116 }
117 }
118 )*
119
120 $(
121 pub(crate) mod [<objects_ $x>] {
122 pub use mz_catalog_protos::[<objects_ $x>]::*;
123
124 use crate::durable::objects::state_update::StateUpdateKindJson;
125
126 impl From<StateUpdateKind> for StateUpdateKindJson {
127 fn from(value: StateUpdateKind) -> Self {
128 // TODO: This requires that the json->proto->json roundtrips
129 // exactly, see database-issues#7179.
130 StateUpdateKindJson::from_serde(&value)
131 }
132 }
133
134 impl From<StateUpdateKindJson> for StateUpdateKind {
135 fn from(value: StateUpdateKindJson) -> Self {
136 value.to_serde()
137 }
138 }
139 }
140 )*
141
142 // Generate test helpers for each version.
143
144 #[cfg(test)]
145 #[derive(Debug, Arbitrary)]
146 enum AllVersionsStateUpdateKind {
147 $(
148 [<$x_old:upper>](crate::durable::upgrade::[<objects_ $x_old>]::StateUpdateKind),
149 )*
150 $(
151 [<$x:upper>](crate::durable::upgrade::[<objects_ $x>]::StateUpdateKind),
152 )*
153 }
154
155 #[cfg(test)]
156 impl AllVersionsStateUpdateKind {
157 #[cfg(test)]
158 fn arbitrary_vec(version: &str) -> Result<Vec<Self>, String> {
159 let mut runner = proptest::test_runner::TestRunner::deterministic();
160 std::iter::repeat(())
161 .filter_map(|_| {
162 AllVersionsStateUpdateKind::arbitrary(version, &mut runner)
163 .transpose()
164 })
165 .take(ENCODED_TEST_CASES)
166 .collect::<Result<_, _>>()
167 }
168
169 #[cfg(test)]
170 fn arbitrary(
171 version: &str,
172 runner: &mut proptest::test_runner::TestRunner,
173 ) -> Result<Option<Self>, String> {
174 match version {
175 $(
176 concat!("objects_", stringify!($x_old)) => {
177 let arbitrary_data =
178 crate::durable::upgrade
179 ::[<objects_ $x_old>]::StateUpdateKind::arbitrary()
180 .new_tree(runner)
181 .expect("unable to create arbitrary data")
182 .current();
183 // Skip over generated data where kind is None
184 // because they are not interesting or possible in
185 // production. Unfortunately any of the inner fields
186 // can still be None, which is also not possible in
187 // production.
188 // TODO(jkosh44) See if there's an arbitrary config
189 // that forces Some.
190 let arbitrary_data = if arbitrary_data.kind.is_some() {
191 Some(Self::[<$x_old:upper>](arbitrary_data))
192 } else {
193 None
194 };
195 Ok(arbitrary_data)
196 }
197 )*
198 $(
199 concat!("objects_", stringify!($x)) => {
200 let arbitrary_data =
201 crate::durable::upgrade
202 ::[<objects_ $x>]::StateUpdateKind::arbitrary()
203 .new_tree(runner)
204 .expect("unable to create arbitrary data")
205 .current();
206 Ok(Some(Self::[<$x:upper>](arbitrary_data)))
207 }
208 )*
209 _ => Err(format!("unrecognized version {version} add enum variant")),
210 }
211 }
212
213 #[cfg(test)]
214 fn try_from_raw(version: &str, raw: StateUpdateKindJson) -> Result<Self, String> {
215 match version {
216 $(
217 concat!("objects_", stringify!($x_old)) => Ok(Self::[<$x_old:upper>](raw.into())),
218 )*
219 $(
220 concat!("objects_", stringify!($x)) => Ok(Self::[<$x:upper>](raw.into())),
221
222 )*
223 _ => Err(format!("unrecognized version {version} add enum variant")),
224 }
225 }
226
227 #[cfg(test)]
228 fn raw(self) -> StateUpdateKindJson {
229 match self {
230 $(
231 Self::[<$x_old:upper>](kind) => kind.into(),
232 )*
233 $(
234 Self::[<$x:upper>](kind) => kind.into(),
235 )*
236 }
237 }
238 }
239 }
240 }
241}
242
243objects!(
244 [v74, v75, v76, v77, v78],
245 [v79, v80, v81, v82, v83, v84, v85, v86, v87, v88, v89, v90]
246);
247
248/// The current version of the `Catalog`.
249pub use mz_catalog_protos::CATALOG_VERSION;
250/// The minimum `Catalog` version number that we support migrating from.
251pub use mz_catalog_protos::MIN_CATALOG_VERSION;
252
253// Note(parkmycar): Ideally we wouldn't have to define these extra constants,
254// but const expressions aren't yet supported in match statements.
255const TOO_OLD_VERSION: u64 = MIN_CATALOG_VERSION - 1;
256const FUTURE_VERSION: u64 = CATALOG_VERSION + 1;
257
258mod v74_to_v75;
259mod v75_to_v76;
260mod v76_to_v77;
261mod v77_to_v78;
262mod v78_to_v79;
263mod v79_to_v80;
264mod v80_to_v81;
265mod v81_to_v82;
266mod v82_to_v83;
267mod v83_to_v84;
268mod v84_to_v85;
269mod v85_to_v86;
270mod v86_to_v87;
271mod v87_to_v88;
272mod v88_to_v89;
273mod v89_to_v90;
274
275/// Describes a single action to take during a migration from `V1` to `V2`.
276#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
277enum MigrationAction<V1: IntoStateUpdateKindJson, V2: IntoStateUpdateKindJson> {
278 /// Deletes the provided key.
279 #[allow(unused)]
280 Delete(V1),
281 /// Inserts the provided key-value pair. The key must not currently exist!
282 #[allow(unused)]
283 Insert(V2),
284 /// Update the key-value pair for the provided key.
285 #[allow(unused)]
286 Update(V1, V2),
287}
288
289impl<V1: IntoStateUpdateKindJson, V2: IntoStateUpdateKindJson> MigrationAction<V1, V2> {
290 /// Converts `self` into a `Vec<StateUpdate<StateUpdateKindBinary>>` that can be appended
291 /// to persist.
292 fn into_updates(self) -> Vec<(StateUpdateKindJson, Diff)> {
293 match self {
294 MigrationAction::Delete(kind) => {
295 vec![(kind.into(), Diff::MINUS_ONE)]
296 }
297 MigrationAction::Insert(kind) => {
298 vec![(kind.into(), Diff::ONE)]
299 }
300 MigrationAction::Update(old_kind, new_kind) => {
301 vec![
302 (old_kind.into(), Diff::MINUS_ONE),
303 (new_kind.into(), Diff::ONE),
304 ]
305 }
306 }
307 }
308}
309
310/// Upgrades the data in the catalog to version [`CATALOG_VERSION`].
311///
312/// Returns the current upper after all migrations have executed.
313#[mz_ore::instrument(name = "persist::upgrade", level = "debug")]
314pub(crate) async fn upgrade(
315 persist_handle: &mut UnopenedPersistCatalogState,
316 mut commit_ts: Timestamp,
317) -> Result<Timestamp, CatalogError> {
318 soft_assert_ne_or_log!(
319 persist_handle.upper,
320 Timestamp::minimum(),
321 "cannot upgrade uninitialized catalog"
322 );
323
324 // Consolidate to avoid migrating old state.
325 persist_handle.consolidate();
326 let mut version = persist_handle
327 .get_user_version()
328 .await?
329 .expect("initialized catalog must have a version");
330 // Run migrations until we're up-to-date.
331 while version < CATALOG_VERSION {
332 (version, commit_ts) = run_upgrade(persist_handle, version, commit_ts).await?;
333 }
334
335 Ok(commit_ts)
336}
337
338/// Determines which upgrade to run for the `version` and executes it.
339///
340/// Returns the new version and upper.
341async fn run_upgrade(
342 unopened_catalog_state: &mut UnopenedPersistCatalogState,
343 version: u64,
344 commit_ts: Timestamp,
345) -> Result<(u64, Timestamp), CatalogError> {
346 let incompatible = DurableCatalogError::IncompatibleDataVersion {
347 found_version: version,
348 min_catalog_version: MIN_CATALOG_VERSION,
349 catalog_version: CATALOG_VERSION,
350 }
351 .into();
352
353 match version {
354 ..=TOO_OLD_VERSION => Err(incompatible),
355
356 74 => {
357 run_versioned_upgrade(
358 unopened_catalog_state,
359 version,
360 commit_ts,
361 v74_to_v75::upgrade,
362 )
363 .await
364 }
365 75 => {
366 run_versioned_upgrade(
367 unopened_catalog_state,
368 version,
369 commit_ts,
370 v75_to_v76::upgrade,
371 )
372 .await
373 }
374 76 => {
375 run_versioned_upgrade(
376 unopened_catalog_state,
377 version,
378 commit_ts,
379 v76_to_v77::upgrade,
380 )
381 .await
382 }
383 77 => {
384 run_versioned_upgrade(
385 unopened_catalog_state,
386 version,
387 commit_ts,
388 v77_to_v78::upgrade,
389 )
390 .await
391 }
392 78 => {
393 run_versioned_upgrade(
394 unopened_catalog_state,
395 version,
396 commit_ts,
397 v78_to_v79::upgrade,
398 )
399 .await
400 }
401 79 => {
402 run_versioned_upgrade(
403 unopened_catalog_state,
404 version,
405 commit_ts,
406 v79_to_v80::upgrade,
407 )
408 .await
409 }
410 80 => {
411 run_versioned_upgrade(
412 unopened_catalog_state,
413 version,
414 commit_ts,
415 v80_to_v81::upgrade,
416 )
417 .await
418 }
419 81 => {
420 run_versioned_upgrade(
421 unopened_catalog_state,
422 version,
423 commit_ts,
424 v81_to_v82::upgrade,
425 )
426 .await
427 }
428 // v82→v83 is a one-shot byte-level repair, not a proto evolution. It
429 // needs raw access to the snapshot's diffs (which `run_versioned_upgrade`
430 // strips) so it plugs into `run_upgrade` directly.
431 82 => v82_to_v83::upgrade(unopened_catalog_state, commit_ts).await,
432 83 => v83_to_v84::upgrade(unopened_catalog_state, commit_ts).await,
433 84 => {
434 run_versioned_upgrade(
435 unopened_catalog_state,
436 version,
437 commit_ts,
438 v84_to_v85::upgrade,
439 )
440 .await
441 }
442 85 => {
443 run_versioned_upgrade(
444 unopened_catalog_state,
445 version,
446 commit_ts,
447 v85_to_v86::upgrade,
448 )
449 .await
450 }
451 86 => {
452 run_versioned_upgrade(
453 unopened_catalog_state,
454 version,
455 commit_ts,
456 v86_to_v87::upgrade,
457 )
458 .await
459 }
460 87 => {
461 run_versioned_upgrade(
462 unopened_catalog_state,
463 version,
464 commit_ts,
465 v87_to_v88::upgrade,
466 )
467 .await
468 }
469 88 => {
470 run_versioned_upgrade(
471 unopened_catalog_state,
472 version,
473 commit_ts,
474 v88_to_v89::upgrade,
475 )
476 .await
477 }
478 89 => {
479 run_versioned_upgrade(
480 unopened_catalog_state,
481 version,
482 commit_ts,
483 v89_to_v90::upgrade,
484 )
485 .await
486 }
487 // Up-to-date, no migration needed!
488 CATALOG_VERSION => Ok((CATALOG_VERSION, commit_ts)),
489 FUTURE_VERSION.. => Err(incompatible),
490 }
491}
492
493/// Runs `migration_logic` on the contents of the current catalog assuming a current version of
494/// `current_version`.
495///
496/// Returns the new version and upper.
497async fn run_versioned_upgrade<V1: IntoStateUpdateKindJson, V2: IntoStateUpdateKindJson>(
498 unopened_catalog_state: &mut UnopenedPersistCatalogState,
499 current_version: u64,
500 mut commit_ts: Timestamp,
501 migration_logic: impl FnOnce(Vec<V1>) -> Vec<MigrationAction<V1, V2>>,
502) -> Result<(u64, Timestamp), CatalogError> {
503 tracing::info!(current_version, "running versioned Catalog upgrade");
504
505 // 1. Use the V1 to deserialize the contents of the current snapshot.
506 let snapshot: Vec<_> = unopened_catalog_state
507 .snapshot
508 .iter()
509 .map(|(kind, ts, diff)| {
510 soft_assert_eq_or_log!(
511 *diff,
512 Diff::ONE,
513 "snapshot is consolidated, ({kind:?}, {ts:?}, {diff:?})"
514 );
515 V1::try_from(kind.clone()).expect("invalid catalog data persisted")
516 })
517 .collect();
518
519 // 2. Generate updates from version specific migration logic.
520 let migration_actions = migration_logic(snapshot);
521 let mut updates: Vec<_> = migration_actions
522 .into_iter()
523 .flat_map(|action| action.into_updates().into_iter())
524 .collect();
525 // Validate that we're not migrating an un-migratable collection.
526 for (update, _) in &updates {
527 if update.is_always_deserializable() {
528 panic!("migration to un-migratable collection: {update:?}\nall updates: {updates:?}");
529 }
530 }
531
532 // 3. Add a retraction for old version and insertion for new version into updates.
533 let next_version = current_version + 1;
534 let version_retraction = (version_update_kind(current_version), Diff::MINUS_ONE);
535 updates.push(version_retraction);
536 let version_insertion = (version_update_kind(next_version), Diff::ONE);
537 updates.push(version_insertion);
538
539 // 4. Apply migration to catalog.
540 if matches!(unopened_catalog_state.mode, Mode::Writable) {
541 commit_ts = unopened_catalog_state
542 .compare_and_append(updates, commit_ts)
543 .await
544 .map_err(|e| e.unwrap_fence_error())?;
545 } else {
546 let ts = commit_ts;
547 let updates = updates
548 .into_iter()
549 .map(|(kind, diff)| StateUpdate { kind, ts, diff });
550 commit_ts = commit_ts.step_forward();
551 unopened_catalog_state.apply_updates_and_consolidate(updates)?;
552 }
553
554 // 5. Consolidate snapshot to remove old versions.
555 unopened_catalog_state.consolidate();
556
557 Ok((next_version, commit_ts))
558}
559
560/// Generates a [`proto::StateUpdateKind`] to update the user version.
561fn version_update_kind(version: u64) -> StateUpdateKindJson {
562 // We can use the current version because Configs can never be migrated and are always wire
563 // compatible.
564 StateUpdateKind::Config(
565 proto::ConfigKey {
566 key: USER_VERSION_KEY.to_string(),
567 },
568 proto::ConfigValue { value: version },
569 )
570 .into()
571}