Skip to main content

mz_adapter/catalog/open/
builtin_schema_migration.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//! Support for migrating the schemas of builtin storage collections.
11//!
12//! If a version upgrade changes the schema of a builtin collection that's made durable in persist,
13//! that persist shard's schema must be migrated accordingly. The migration must happen in a way
14//! that's compatible with 0dt upgrades: Read-only environments need to be able to read the
15//! collections with the new schema, without interfering with the leader environment's continued
16//! use of the old schema.
17//!
18//! Two migration mechanisms are provided:
19//!
20//!  * [`Mechanism::Evolution`] uses persist's schema evolution support to evolve the persist
21//!    shard's schema in-place. Only works for backward-compatible changes.
22//!  * [`Mechanism::Replacement`] creates a new shard to serve the builtin collection in the new
23//!    version. Works for all schema changes but discards existing data.
24//!
25//! Which mechanism to use is selected through entries in the `MIGRATIONS` list. In general, the
26//! `Evolution` mechanism should be used when possible, as it avoids data loss.
27//!
28//! For more context and details on the implementation, see
29//! `doc/developer/design/20251015_builtin_schema_migration.md`.
30
31use std::collections::{BTreeMap, BTreeSet};
32use std::sync::{Arc, LazyLock};
33
34use anyhow::bail;
35use futures::FutureExt;
36use futures::future::BoxFuture;
37use mz_build_info::{BuildInfo, DUMMY_BUILD_INFO};
38use mz_catalog::builtin::{
39    BUILTIN_LOOKUP, Builtin, Fingerprint, MZ_CATALOG_RAW, MZ_CATALOG_RAW_DESCRIPTION,
40    MZ_CLUSTER_REPLICA_FRONTIERS_DESCRIPTION, MZ_OBJECT_ARRANGEMENT_SIZE_HISTORY_DESCRIPTION,
41    MZ_OBJECT_HYDRATION_HISTORY, MZ_OBJECT_HYDRATION_HISTORY_DESCRIPTION,
42    MZ_REPLICA_HYDRATION_HISTORY, MZ_REPLICA_HYDRATION_HISTORY_DESCRIPTION,
43    MZ_STORAGE_USAGE_BY_SHARD, MZ_STORAGE_USAGE_BY_SHARD_DESCRIPTION,
44    RUNTIME_ALTERABLE_FINGERPRINT_SENTINEL,
45};
46use mz_catalog::config::BuiltinItemMigrationConfig;
47use mz_catalog::durable::objects::SystemObjectUniqueIdentifier;
48use mz_catalog::durable::{SystemObjectDescription, SystemObjectMapping, Transaction};
49use mz_catalog::memory::error::{Error, ErrorKind};
50use mz_ore::soft_assert_or_log;
51use mz_persist_client::cfg::USE_CRITICAL_SINCE_CATALOG;
52use mz_persist_client::critical::{Opaque, SinceHandle};
53use mz_persist_client::read::ReadHandle;
54use mz_persist_client::schema::CaESchema;
55use mz_persist_client::write::WriteHandle;
56use mz_persist_client::{Diagnostics, PersistClient};
57use mz_persist_types::ShardId;
58use mz_persist_types::codec_impls::{ShardIdSchema, UnitSchema};
59use mz_persist_types::schema::backward_compatible;
60use mz_repr::namespaces::{MZ_CATALOG_SCHEMA, MZ_INTERNAL_SCHEMA};
61use mz_repr::{CatalogItemId, GlobalId, Timestamp};
62use mz_sql::catalog::{CatalogItemType, NameReference};
63use mz_storage_client::controller::StorageTxn;
64use mz_storage_types::StorageDiff;
65use mz_storage_types::sources::SourceData;
66use semver::Version;
67use timely::progress::Antichain;
68use tracing::{debug, info};
69
70use crate::catalog::migrate::get_migration_version;
71
72/// Builtin schema migrations required to upgrade to the current build version.
73///
74/// Migration steps for old versions must be retained around according to the upgrade policy. For
75/// example, if we support upgrading one major version at a time, the release of version `N.0.0`
76/// can delete all migration steps with versions before `(N-1).0.0`.
77///
78/// Exception: when a builtin's `SystemObjectDescription` changes — e.g. a builtin table is
79/// converted to a materialized view (see `migrate_builtin_tables_to_mvs`), or a builtin is
80/// renamed or removed — existing steps naming the old description must be removed regardless
81/// of version, because `validate_migration_steps` panics on steps that don't resolve to a
82/// current builtin. This is safe only if a `Replacement` step for the new description is added
83/// at the conversion version: every environment that needed the removed steps upgrades from an
84/// even older version, so the new replacement subsumes them.
85///
86/// Smallest supported version: 0.147.0
87static MIGRATIONS: LazyLock<Vec<MigrationStep>> = LazyLock::new(|| {
88    vec![
89        MigrationStep::replacement(
90            "0.149.0",
91            CatalogItemType::Source,
92            MZ_INTERNAL_SCHEMA,
93            "mz_sink_statistics_raw",
94        ),
95        MigrationStep::replacement(
96            "0.149.0",
97            CatalogItemType::Source,
98            MZ_INTERNAL_SCHEMA,
99            "mz_source_statistics_raw",
100        ),
101        MigrationStep::evolution(
102            "0.159.0",
103            CatalogItemType::Source,
104            MZ_INTERNAL_SCHEMA,
105            "mz_cluster_replica_metrics_history",
106        ),
107        MigrationStep::replacement(
108            "26.18.0-dev.0",
109            CatalogItemType::MaterializedView,
110            MZ_CATALOG_SCHEMA,
111            "mz_databases",
112        ),
113        MigrationStep::replacement(
114            "26.19.0-dev.0",
115            CatalogItemType::MaterializedView,
116            MZ_CATALOG_SCHEMA,
117            "mz_schemas",
118        ),
119        MigrationStep::replacement(
120            "26.19.0-dev.0",
121            CatalogItemType::MaterializedView,
122            MZ_CATALOG_SCHEMA,
123            "mz_role_members",
124        ),
125        MigrationStep::replacement(
126            "26.19.0-dev.0",
127            CatalogItemType::MaterializedView,
128            MZ_INTERNAL_SCHEMA,
129            "mz_network_policies",
130        ),
131        MigrationStep::replacement(
132            "26.19.0-dev.0",
133            CatalogItemType::MaterializedView,
134            MZ_INTERNAL_SCHEMA,
135            "mz_network_policy_rules",
136        ),
137        MigrationStep::replacement(
138            "26.19.0-dev.0",
139            CatalogItemType::MaterializedView,
140            MZ_INTERNAL_SCHEMA,
141            "mz_cluster_workload_classes",
142        ),
143        MigrationStep::replacement(
144            "26.19.0-dev.0",
145            CatalogItemType::MaterializedView,
146            MZ_INTERNAL_SCHEMA,
147            "mz_internal_cluster_replicas",
148        ),
149        MigrationStep::replacement(
150            "26.19.0-dev.0",
151            CatalogItemType::MaterializedView,
152            MZ_INTERNAL_SCHEMA,
153            "mz_pending_cluster_replicas",
154        ),
155        MigrationStep::replacement(
156            "26.20.0-dev.0",
157            CatalogItemType::MaterializedView,
158            MZ_CATALOG_SCHEMA,
159            "mz_materialized_views",
160        ),
161        MigrationStep::replacement(
162            "26.22.0-dev.0",
163            CatalogItemType::MaterializedView,
164            MZ_CATALOG_SCHEMA,
165            "mz_connections",
166        ),
167        MigrationStep::replacement(
168            "26.22.0-dev.0",
169            CatalogItemType::MaterializedView,
170            MZ_CATALOG_SCHEMA,
171            "mz_secrets",
172        ),
173        MigrationStep::replacement(
174            "26.27.0-dev.0",
175            CatalogItemType::MaterializedView,
176            MZ_CATALOG_SCHEMA,
177            "mz_sources",
178        ),
179        MigrationStep::replacement(
180            "26.29.0-dev.0",
181            CatalogItemType::MaterializedView,
182            MZ_CATALOG_SCHEMA,
183            "mz_indexes",
184        ),
185        MigrationStep::replacement(
186            "26.29.0-dev.0",
187            CatalogItemType::MaterializedView,
188            MZ_CATALOG_SCHEMA,
189            "mz_roles",
190        ),
191        MigrationStep::replacement(
192            "26.29.0-dev.0",
193            CatalogItemType::MaterializedView,
194            MZ_CATALOG_SCHEMA,
195            "mz_role_parameters",
196        ),
197        // Required because we added `mz_cluster_replica_size_internal_ind` builtin
198        // index without bumping mz_indexes. make_mz_indexes inlines the builtin-index
199        // set as VALUES, so any add/remove changes its SQL fingerprint and requires
200        // an explicit replacement step.
201        MigrationStep::replacement(
202            "26.30.0-dev.0",
203            CatalogItemType::MaterializedView,
204            MZ_CATALOG_SCHEMA,
205            "mz_indexes",
206        ),
207        MigrationStep::replacement(
208            "26.30.0-dev.0",
209            CatalogItemType::MaterializedView,
210            MZ_CATALOG_SCHEMA,
211            "mz_clusters",
212        ),
213        MigrationStep::replacement(
214            "26.30.0-dev.0",
215            CatalogItemType::MaterializedView,
216            MZ_CATALOG_SCHEMA,
217            "mz_cluster_replicas",
218        ),
219        MigrationStep::replacement(
220            "26.30.0-dev.0",
221            CatalogItemType::MaterializedView,
222            MZ_INTERNAL_SCHEMA,
223            "mz_cluster_schedules",
224        ),
225        MigrationStep::replacement(
226            "26.30.0-dev.0",
227            CatalogItemType::MaterializedView,
228            MZ_CATALOG_SCHEMA,
229            "mz_default_privileges",
230        ),
231        MigrationStep::replacement(
232            "26.30.0-dev.0",
233            CatalogItemType::MaterializedView,
234            MZ_CATALOG_SCHEMA,
235            "mz_system_privileges",
236        ),
237        // The mz_cluster_replicas MV definition changed in 26.31.0-dev (the
238        // `availability_zone` column now aggregates the durable
239        // `availability_zones` list).
240        MigrationStep::replacement(
241            "26.31.0-dev.0",
242            CatalogItemType::MaterializedView,
243            MZ_CATALOG_SCHEMA,
244            "mz_cluster_replicas",
245        ),
246        MigrationStep::replacement(
247            "26.32.0-dev.0",
248            CatalogItemType::MaterializedView,
249            MZ_INTERNAL_SCHEMA,
250            "mz_comments",
251        ),
252        // Required because we added the console cluster-utilization overview builtin
253        // indexes (overview/_3h/_24h). make_mz_indexes inlines the builtin-index set
254        // as VALUES, so any add/remove changes its SQL fingerprint and requires an
255        // explicit replacement step.
256        MigrationStep::replacement(
257            "26.32.0-dev.0",
258            CatalogItemType::MaterializedView,
259            MZ_CATALOG_SCHEMA,
260            "mz_indexes",
261        ),
262        MigrationStep::replacement(
263            "26.33.0-dev.0",
264            CatalogItemType::MaterializedView,
265            MZ_CATALOG_SCHEMA,
266            "mz_audit_events",
267        ),
268        // Required because we added the `mz_cluster_reconfigurations_ind` and
269        // `mz_cluster_auto_scaling_strategies_ind` builtin indexes without
270        // bumping mz_indexes. make_mz_indexes inlines the builtin-index set as
271        // VALUES, so any add/remove changes its SQL fingerprint and requires an
272        // explicit replacement step.
273        //
274        // NOTE: this version must stay at the workspace's current dev version
275        // until this change ships in a release. A dev version orders below its
276        // release, so a step pinned to an older dev version is skipped when
277        // upgrading from that release onward, and the fingerprint check then
278        // panics at catalog open.
279        MigrationStep::replacement(
280            "26.34.0-dev.0",
281            CatalogItemType::MaterializedView,
282            MZ_CATALOG_SCHEMA,
283            "mz_indexes",
284        ),
285        // Converting mz_postgres_sources / mz_kafka_sources from builtin tables
286        // to materialized views changes their catalog fingerprint, so both need
287        // an explicit replacement step. See the NOTE above: this version must
288        // stay at the workspace's current dev version until the change ships.
289        MigrationStep::replacement(
290            "26.34.0-dev.0",
291            CatalogItemType::MaterializedView,
292            MZ_INTERNAL_SCHEMA,
293            "mz_postgres_sources",
294        ),
295        MigrationStep::replacement(
296            "26.34.0-dev.0",
297            CatalogItemType::MaterializedView,
298            MZ_CATALOG_SCHEMA,
299            "mz_kafka_sources",
300        ),
301        // Converting the four mz_*_source_tables from builtin tables to
302        // materialized views changes their catalog fingerprint, so each needs an
303        // explicit replacement step.
304        MigrationStep::replacement(
305            "26.37.0-dev.0",
306            CatalogItemType::MaterializedView,
307            MZ_INTERNAL_SCHEMA,
308            "mz_postgres_source_tables",
309        ),
310        MigrationStep::replacement(
311            "26.37.0-dev.0",
312            CatalogItemType::MaterializedView,
313            MZ_INTERNAL_SCHEMA,
314            "mz_mysql_source_tables",
315        ),
316        MigrationStep::replacement(
317            "26.37.0-dev.0",
318            CatalogItemType::MaterializedView,
319            MZ_INTERNAL_SCHEMA,
320            "mz_sql_server_source_tables",
321        ),
322        MigrationStep::replacement(
323            "26.37.0-dev.0",
324            CatalogItemType::MaterializedView,
325            MZ_INTERNAL_SCHEMA,
326            "mz_kafka_source_tables",
327        ),
328        // `mz_type_pg_metadata` gained a trailing `typsend` column. See the NOTE
329        // above: this version must stay at the workspace's current dev version
330        // until the change ships.
331        MigrationStep::replacement(
332            "26.40.0-dev.0",
333            CatalogItemType::Table,
334            MZ_INTERNAL_SCHEMA,
335            "mz_type_pg_metadata",
336        ),
337        // Converting the connection-detail builtin tables to materialized views
338        // changes their catalog fingerprint, so each needs an explicit
339        // replacement step.
340        MigrationStep::replacement(
341            "26.37.0-dev.0",
342            CatalogItemType::MaterializedView,
343            MZ_CATALOG_SCHEMA,
344            "mz_kafka_connections",
345        ),
346        MigrationStep::replacement(
347            "26.37.0-dev.0",
348            CatalogItemType::MaterializedView,
349            MZ_CATALOG_SCHEMA,
350            "mz_ssh_tunnel_connections",
351        ),
352        MigrationStep::replacement(
353            "26.37.0-dev.0",
354            CatalogItemType::MaterializedView,
355            MZ_INTERNAL_SCHEMA,
356            "mz_aws_connections",
357        ),
358        MigrationStep::replacement(
359            "26.37.0-dev.0",
360            CatalogItemType::MaterializedView,
361            MZ_CATALOG_SCHEMA,
362            "mz_aws_privatelink_connections",
363        ),
364        // The three sink tables became materialized views, which moves their
365        // fingerprints. The old `0.160.0` `mz_sinks` step had to go at the same
366        // time, because it names the `Table` description and that no longer
367        // resolves. Nothing is lost: anything that needed the old step upgrades
368        // from further back than this one, so this one covers it too.
369        MigrationStep::replacement(
370            "26.38.0-dev.0",
371            CatalogItemType::MaterializedView,
372            MZ_CATALOG_SCHEMA,
373            "mz_sinks",
374        ),
375        MigrationStep::replacement(
376            "26.38.0-dev.0",
377            CatalogItemType::MaterializedView,
378            MZ_CATALOG_SCHEMA,
379            "mz_kafka_sinks",
380        ),
381        MigrationStep::replacement(
382            "26.38.0-dev.0",
383            CatalogItemType::MaterializedView,
384            MZ_CATALOG_SCHEMA,
385            "mz_iceberg_sinks",
386        ),
387        // The mz_cluster_reconfigurations MV definition changed (the `changes`
388        // diff now includes the `arrangement_compression` dimension).
389        MigrationStep::replacement(
390            "26.38.0-rc.2",
391            CatalogItemType::MaterializedView,
392            MZ_INTERNAL_SCHEMA,
393            "mz_cluster_reconfigurations",
394        ),
395        // The mz_audit_events MV gained a `metric-sink` arm in its object_type
396        // CASE, changing its SQL fingerprint, so it needs an explicit
397        // replacement step. See the NOTE above: this version must stay at the
398        // workspace's current dev version until the change ships.
399        MigrationStep::replacement(
400            "26.39.0-dev.0",
401            CatalogItemType::MaterializedView,
402            MZ_CATALOG_SCHEMA,
403            "mz_audit_events",
404        ),
405        // Required because we added the `mz_metric_sinks_ind` builtin index.
406        // make_mz_indexes inlines the builtin-index set as VALUES, so any add or
407        // remove changes its SQL fingerprint and requires an explicit
408        // replacement. See the NOTE above: this version must stay at the
409        // workspace's current dev version until the change ships. The addition
410        // is unshipped, so this step supersedes the earlier `26.39.0-dev.0`
411        // `mz_indexes` step (which covered `mz_object_graph_edges_ind`): a
412        // replacement recreates `mz_indexes` from its current definition, so a
413        // single step at the current dev version covers every earlier change too.
414        MigrationStep::replacement(
415            "26.40.0-dev.0",
416            CatalogItemType::MaterializedView,
417            MZ_CATALOG_SCHEMA,
418            "mz_indexes",
419        ),
420        // Converting mz_tables and mz_views from builtin tables to
421        // materialized views over mz_catalog_raw changes their catalog
422        // fingerprints, so each needs an explicit replacement step. See the
423        // NOTE above: this version must stay at the workspace's current dev
424        // version until the change ships.
425        MigrationStep::replacement(
426            "26.39.0-dev.0",
427            CatalogItemType::MaterializedView,
428            MZ_CATALOG_SCHEMA,
429            "mz_tables",
430        ),
431        MigrationStep::replacement(
432            "26.39.0-dev.0",
433            CatalogItemType::MaterializedView,
434            MZ_CATALOG_SCHEMA,
435            "mz_views",
436        ),
437        // Required because we added the `mz_cluster_replica_resource_usage` builtin log.
438        // make_mz_indexes and make_mz_sources inline the builtin-log set as
439        // VALUES, so adding one changes both MVs' SQL fingerprints. See the NOTE
440        // above: this version must stay at the workspace's current dev version
441        // until the change ships.
442        MigrationStep::replacement(
443            "26.40.0-dev.0",
444            CatalogItemType::MaterializedView,
445            MZ_CATALOG_SCHEMA,
446            "mz_indexes",
447        ),
448        MigrationStep::replacement(
449            "26.40.0-dev.0",
450            CatalogItemType::MaterializedView,
451            MZ_CATALOG_SCHEMA,
452            "mz_sources",
453        ),
454        // Converting mz_object_dependencies from a builtin table to a
455        // materialized view over mz_catalog_raw changes its catalog
456        // fingerprint, so it needs an explicit replacement step. See the
457        // NOTE above: this version must stay at the workspace's current dev
458        // version until the change ships.
459        MigrationStep::replacement(
460            "26.41.0-dev.0",
461            CatalogItemType::MaterializedView,
462            MZ_INTERNAL_SCHEMA,
463            "mz_object_dependencies",
464        ),
465        // Adding the `swap_bytes` column to `mz_cluster_replica_metrics_history`
466        // is a backward-compatible append, so the shard schema can evolve in
467        // place. See the NOTE above: this version must stay at the workspace's
468        // current dev version until the change ships.
469        MigrationStep::evolution(
470            "26.43.0-dev.0",
471            CatalogItemType::Source,
472            MZ_INTERNAL_SCHEMA,
473            "mz_cluster_replica_metrics_history",
474        ),
475        MigrationStep::evolution(
476            "26.43.0-dev.0",
477            CatalogItemType::Table,
478            MZ_INTERNAL_SCHEMA,
479            "mz_replica_hydration_history",
480        ),
481        // Converting mz_object_global_ids from a builtin table to a
482        // materialized view over mz_catalog_raw changes its catalog
483        // fingerprint.
484        MigrationStep::replacement(
485            "26.44.0-dev.0",
486            CatalogItemType::MaterializedView,
487            MZ_INTERNAL_SCHEMA,
488            "mz_object_global_ids",
489        ),
490        // Converting mz_source_references from a builtin table to a
491        // materialized view over mz_catalog_raw changes its catalog
492        // fingerprint.
493        MigrationStep::replacement(
494            "26.45.0-dev.0",
495            CatalogItemType::MaterializedView,
496            MZ_INTERNAL_SCHEMA,
497            "mz_source_references",
498        ),
499    ]
500});
501
502/// A migration required to upgrade past a specific version.
503#[derive(Clone, Debug)]
504struct MigrationStep {
505    /// The build version that requires this migration.
506    version: Version,
507    /// The object that requires migration.
508    object: SystemObjectDescription,
509    /// The migration mechanism to be used.
510    mechanism: Mechanism,
511}
512
513impl MigrationStep {
514    /// Helper to construct an `Evolution` migration step.
515    fn evolution(version: &str, type_: CatalogItemType, schema: &str, name: &str) -> Self {
516        Self {
517            version: Version::parse(version).expect("valid"),
518            object: SystemObjectDescription {
519                schema_name: schema.into(),
520                object_type: type_,
521                object_name: name.into(),
522            },
523            mechanism: Mechanism::Evolution,
524        }
525    }
526
527    /// Helper to construct a `Replacement` migration step.
528    fn replacement(version: &str, type_: CatalogItemType, schema: &str, name: &str) -> Self {
529        Self {
530            version: Version::parse(version).expect("valid"),
531            object: SystemObjectDescription {
532                schema_name: schema.into(),
533                object_type: type_,
534                object_name: name.into(),
535            },
536            mechanism: Mechanism::Replacement,
537        }
538    }
539}
540
541/// The mechanism to use to migrate the schema of a builtin collection.
542#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
543#[allow(dead_code)]
544enum Mechanism {
545    /// Persist schema evolution.
546    ///
547    /// Keeps existing contents but only works for schema changes that are backward compatible
548    /// according to [`backward_compatible`].
549    Evolution,
550    /// Shard replacement.
551    ///
552    /// Works for arbitrary schema changes but loses existing contents.
553    Replacement,
554}
555
556/// The result of a builtin schema migration.
557pub(super) struct MigrationResult {
558    /// IDs of items whose shards have been replaced using the `Replacement` mechanism.
559    pub replaced_items: BTreeSet<CatalogItemId>,
560    /// A cleanup action to take once the migration has been made durable.
561    pub cleanup_action: BoxFuture<'static, ()>,
562}
563
564impl Default for MigrationResult {
565    fn default() -> Self {
566        Self {
567            replaced_items: Default::default(),
568            cleanup_action: async {}.boxed(),
569        }
570    }
571}
572
573/// Run builtin schema migrations.
574///
575/// This is the entry point used by adapter when opening the catalog. It uses the hardcoded
576/// `BUILTINS` and `MIGRATIONS` lists to initialize the lists of available builtins and required
577/// migrations, respectively.
578pub(super) async fn run(
579    build_info: &BuildInfo,
580    deploy_generation: u64,
581    txn: &mut Transaction<'_>,
582    config: BuiltinItemMigrationConfig,
583) -> Result<MigrationResult, Error> {
584    // Sanity check to ensure we're not touching durable state in read-only mode.
585    assert_eq!(config.read_only, txn.is_savepoint());
586
587    // Tests may provide a dummy build info that confuses the migration step selection logic. Skip
588    // migrations if we observe this build info.
589    if *build_info == DUMMY_BUILD_INFO {
590        return Ok(MigrationResult::default());
591    }
592
593    let Some(durable_version) = get_migration_version(txn) else {
594        // New catalog; nothing to do.
595        return Ok(MigrationResult::default());
596    };
597    let build_version = build_info.semver_version();
598
599    let collection_metadata = txn.get_collection_metadata();
600    let system_objects = txn
601        .get_system_object_mappings()
602        .map(|m| {
603            let object = m.description;
604            let global_id = m.unique_identifier.global_id;
605            let shard_id = collection_metadata.get(&global_id).copied();
606            let Some((_, builtin)) = BUILTIN_LOOKUP.get(&object) else {
607                panic!("missing builtin {object:?}");
608            };
609            let info = ObjectInfo {
610                global_id,
611                shard_id,
612                builtin,
613                fingerprint: m.unique_identifier.fingerprint,
614            };
615            (object, info)
616        })
617        .collect();
618
619    let migration_shard = txn.get_builtin_migration_shard().expect("must exist");
620
621    let migration = Migration {
622        source_version: durable_version.clone(),
623        target_version: build_version.clone(),
624        deploy_generation,
625        system_objects,
626        migration_shard,
627        config,
628    };
629
630    let result = migration.run(&MIGRATIONS).await.map_err(|e| {
631        Error::new(ErrorKind::FailedBuiltinSchemaMigration {
632            last_seen_version: durable_version.to_string(),
633            this_version: build_version.to_string(),
634            cause: e.to_string(),
635        })
636    })?;
637
638    result.apply(txn);
639
640    let replaced_items = txn
641        .get_system_object_mappings()
642        .map(|m| m.unique_identifier)
643        .filter(|ids| result.new_shards.contains_key(&ids.global_id))
644        .map(|ids| ids.catalog_id)
645        .collect();
646
647    Ok(MigrationResult {
648        replaced_items,
649        cleanup_action: result.cleanup_action,
650    })
651}
652
653/// Result produced by `Migration::run`.
654struct MigrationRunResult {
655    new_shards: BTreeMap<GlobalId, ShardId>,
656    new_fingerprints: BTreeMap<SystemObjectDescription, String>,
657    shards_to_finalize: BTreeSet<ShardId>,
658    cleanup_action: BoxFuture<'static, ()>,
659}
660
661impl Default for MigrationRunResult {
662    fn default() -> Self {
663        Self {
664            new_shards: BTreeMap::new(),
665            new_fingerprints: BTreeMap::new(),
666            shards_to_finalize: BTreeSet::new(),
667            cleanup_action: async {}.boxed(),
668        }
669    }
670}
671
672impl MigrationRunResult {
673    /// Apply this migration result to the given transaction.
674    fn apply(&self, txn: &mut Transaction<'_>) {
675        // Update collection metadata.
676        let replaced_ids = self.new_shards.keys().copied().collect();
677        let old_metadata = txn.delete_collection_metadata(replaced_ids);
678        txn.insert_collection_metadata(self.new_shards.clone())
679            .expect("inserting unique shards IDs after deleting existing entries");
680
681        // Register shards for finalization.
682        let mut unfinalized_shards: BTreeSet<_> =
683            old_metadata.into_iter().map(|(_, sid)| sid).collect();
684        unfinalized_shards.extend(self.shards_to_finalize.iter().copied());
685        txn.insert_unfinalized_shards(unfinalized_shards)
686            .expect("cannot fail");
687
688        // Update fingerprints.
689        let mappings = txn
690            .get_system_object_mappings()
691            .filter_map(|m| {
692                let fingerprint = self.new_fingerprints.get(&m.description)?;
693                Some(SystemObjectMapping {
694                    description: m.description,
695                    unique_identifier: SystemObjectUniqueIdentifier {
696                        catalog_id: m.unique_identifier.catalog_id,
697                        global_id: m.unique_identifier.global_id,
698                        fingerprint: fingerprint.clone(),
699                    },
700                })
701            })
702            .collect();
703        txn.set_system_object_mappings(mappings)
704            .expect("filtered existing mappings remain unique");
705    }
706}
707
708/// Information about a system object required to run a `Migration`.
709#[derive(Clone, Debug)]
710struct ObjectInfo {
711    global_id: GlobalId,
712    shard_id: Option<ShardId>,
713    builtin: &'static Builtin<NameReference>,
714    fingerprint: String,
715}
716
717/// Context of a builtin schema migration.
718struct Migration {
719    /// The version we are migrating from.
720    ///
721    /// Same as the build version of the most recent leader process that successfully performed
722    /// migrations.
723    source_version: Version,
724    /// The version we are migration to.
725    ///
726    /// Same as the build version of this process.
727    target_version: Version,
728    /// The deploy generation of this process.
729    deploy_generation: u64,
730    /// Information about all objects in the system.
731    system_objects: BTreeMap<SystemObjectDescription, ObjectInfo>,
732    /// The ID of the migration shard.
733    migration_shard: ShardId,
734    /// Additional configuration.
735    config: BuiltinItemMigrationConfig,
736}
737
738/// Whether `builtin` participates in a forced migration using `mechanism`.
739fn participates_in_forced_migration(
740    builtin: &Builtin<NameReference>,
741    mechanism: Mechanism,
742) -> bool {
743    use Builtin::*;
744    match builtin {
745        // A forced replacement allocates a fresh shard, which discards the
746        // table's contents. Exclude the tables whose contents are the point:
747        // storage usage is retained for billing, and hydration histories cannot
748        // be rebuilt from any other source.
749        //
750        // Hydration history tables take part in a forced `Evolution`, which
751        // keeps the rows. They have to: dev upgrades force one for every object,
752        // and a table left out of the plan never gets its new schema registered.
753        // `update_fingerprints` then panics at open as soon as the desc changes.
754        // See the tripwire in `validate_migration_steps` for how to give up the
755        // replacement exemption deliberately.
756        Table(table) => {
757            **table != *MZ_STORAGE_USAGE_BY_SHARD
758                && (mechanism != Mechanism::Replacement
759                    || (**table != *MZ_OBJECT_HYDRATION_HISTORY
760                        && **table != *MZ_REPLICA_HYDRATION_HISTORY))
761        }
762        MaterializedView(..) => true,
763        Source(source) => **source != *MZ_CATALOG_RAW,
764        Log(..) | View(..) | Type(..) | Func(..) | Index(..) | Connection(..) => false,
765    }
766}
767
768impl Migration {
769    async fn run(self, steps: &[MigrationStep]) -> anyhow::Result<MigrationRunResult> {
770        info!(
771            deploy_generation = %self.deploy_generation,
772            "running builtin schema migration: {} -> {}",
773            self.source_version, self.target_version
774        );
775
776        self.validate_migration_steps(steps);
777
778        // Version-based migration filter fails for dev versions, see for example
779        // https://github.com/MaterializeInc/database-issues/issues/11335
780        let force_migration = if self.source_version != self.target_version
781            && self.source_version.pre.as_str().starts_with("dev")
782            && self.config.force_migration.is_none()
783        {
784            Some("evolution".to_string())
785        } else {
786            self.config.force_migration.clone()
787        };
788
789        let (force, plan) = match force_migration.as_deref() {
790            None => (false, self.plan_migration(steps)),
791            Some("evolution") => (true, self.plan_forced_migration(Mechanism::Evolution)),
792            Some("replacement") => (true, self.plan_forced_migration(Mechanism::Replacement)),
793            Some(other) => panic!("unknown force migration mechanism: {other}"),
794        };
795        let plan = self.drop_shardless(plan);
796
797        if self.source_version == self.target_version && !force {
798            info!("skipping migration: already at target version");
799            return Ok(MigrationRunResult::default());
800        } else if self.source_version > self.target_version {
801            bail!("downgrade not supported");
802        }
803
804        // In leader mode, upgrade the version of the migration shard to the target version.
805        // This fences out any readers at lower versions.
806        if !self.config.read_only {
807            self.upgrade_migration_shard_version().await;
808        }
809
810        info!("executing migration plan: {plan:?}");
811
812        self.migrate_evolve(&plan.evolve).await?;
813        let new_shards = self.migrate_replace(&plan.replace).await?;
814
815        let mut migrated_objects = BTreeSet::new();
816        migrated_objects.extend(plan.evolve);
817        migrated_objects.extend(plan.replace);
818
819        let new_fingerprints = self.update_fingerprints(&migrated_objects)?;
820
821        let (shards_to_finalize, cleanup_action) = self.cleanup().await?;
822
823        Ok(MigrationRunResult {
824            new_shards,
825            new_fingerprints,
826            shards_to_finalize,
827            cleanup_action,
828        })
829    }
830
831    /// Sanity check the given migration steps.
832    ///
833    /// If any of these checks fail, that's a bug in Materialize, and we panic immediately.
834    fn validate_migration_steps(&self, steps: &[MigrationStep]) {
835        for step in steps {
836            assert!(
837                step.version <= self.target_version,
838                "migration step version greater than target version: {} > {}",
839                step.version,
840                self.target_version,
841            );
842
843            let object = &step.object;
844
845            // `mz_storage_usage_by_shard` cannot be migrated for multiple reasons. Firstly, it would
846            // cause the table to be truncated because the contents are not also stored in the durable
847            // catalog. Secondly, we prune `mz_storage_usage_by_shard` of old events in the background
848            // on startup. The correctness of that pruning relies on there being no other retractions
849            // to `mz_storage_usage_by_shard`.
850            //
851            // TODO: Confirm the above reasoning, it might be outdated?
852            assert_ne!(
853                &*MZ_STORAGE_USAGE_BY_SHARD_DESCRIPTION, object,
854                "mz_storage_usage_by_shard cannot be migrated or else the table will be truncated"
855            );
856
857            // Same hazard as `mz_storage_usage_by_shard`: the startup pruner
858            // (`Coordinator::prune_arrangement_sizes_history_on_startup`) assumes it is
859            // the only source of retractions, so migration-driven truncation would break it.
860            assert_ne!(
861                &*MZ_OBJECT_ARRANGEMENT_SIZE_HISTORY_DESCRIPTION, object,
862                "mz_object_arrangement_size_history cannot be migrated or else the table will be truncated"
863            );
864
865            // Unlike the two tables above, there is no correctness hazard here, a
866            // truncation would only lose history. This is a tripwire so that the
867            // loss is chosen rather than stumbled into: if a schema change to
868            // this table is worth clearing it for, remove this assert along with
869            // the exemption in `plan_forced_migration`, and say in the release
870            // notes that the history restarts.
871            //
872            // Only `Replacement` loses the rows. An `Evolution` keeps the shard
873            // and its contents, so it is the path a schema change to this table
874            // should take and we let it through.
875            if step.mechanism == Mechanism::Replacement {
876                assert_ne!(
877                    &*MZ_OBJECT_HYDRATION_HISTORY_DESCRIPTION, object,
878                    "replacing mz_object_hydration_history clears it, see the comment above"
879                );
880                assert_ne!(
881                    &*MZ_REPLICA_HYDRATION_HISTORY_DESCRIPTION, object,
882                    "replacing mz_replica_hydration_history clears it, see the comment above"
883                );
884            }
885
886            // `mz_catalog_raw` cannot be migrated because it contains the durable catalog and it
887            // wouldn't be very durable if we allowed it to be truncated.
888            assert_ne!(
889                &*MZ_CATALOG_RAW_DESCRIPTION, object,
890                "mz_catalog_raw cannot be migrated"
891            );
892
893            // The 0dt caught-up gate reads the leader's `mz_cluster_replica_frontiers` shard for
894            // the live frontiers it checks every collection against. Migrating it via `Replacement`
895            // hands us a fresh shard we write ourselves, so the gate would compare us against
896            // ourselves instead of against the leader.
897            assert_ne!(
898                &*MZ_CLUSTER_REPLICA_FRONTIERS_DESCRIPTION, object,
899                "mz_cluster_replica_frontiers cannot be migrated or else the 0dt caught-up gate loses its live-frontier reference"
900            );
901
902            let Some(object_info) = self.system_objects.get(object) else {
903                panic!("migration step for non-existent builtin: {object:?}");
904            };
905
906            let builtin = object_info.builtin;
907            use Builtin::*;
908            assert!(
909                matches!(builtin, Table(..) | Source(..) | MaterializedView(..)),
910                "schema migration not supported for builtin: {builtin:?}",
911            );
912        }
913    }
914
915    /// Select for each object to migrate the appropriate migration mechanism.
916    fn plan_migration(&self, steps: &[MigrationStep]) -> Plan {
917        // Ignore any steps at versions before `source_version`.
918        let steps = steps.iter().filter(|s| s.version > self.source_version);
919
920        // Select a mechanism for each object, according to the requested migrations:
921        //  * If any `Replacement` was requested, use `Replacement`.
922        //  * Otherwise, (i.e. only `Evolution` was requested), use `Evolution`.
923        let mut by_object = BTreeMap::new();
924        for step in steps {
925            if let Some(entry) = by_object.get_mut(&step.object) {
926                *entry = match (step.mechanism, *entry) {
927                    (Mechanism::Evolution, Mechanism::Evolution) => Mechanism::Evolution,
928                    (Mechanism::Replacement, _) | (_, Mechanism::Replacement) => {
929                        Mechanism::Replacement
930                    }
931                };
932            } else {
933                by_object.insert(step.object.clone(), step.mechanism);
934            }
935        }
936
937        let mut plan = Plan::default();
938        for (object, mechanism) in by_object {
939            match mechanism {
940                Mechanism::Evolution => plan.evolve.push(object),
941                Mechanism::Replacement => plan.replace.push(object),
942            }
943        }
944
945        plan
946    }
947
948    /// Plan a forced migration of all objects using the given mechanism.
949    fn plan_forced_migration(&self, mechanism: Mechanism) -> Plan {
950        let objects = self
951            .system_objects
952            .iter()
953            .filter(|(_, info)| participates_in_forced_migration(info.builtin, mechanism))
954            .map(|(object, _)| object.clone())
955            .collect();
956
957        let mut plan = Plan::default();
958        match mechanism {
959            Mechanism::Evolution => plan.evolve = objects,
960            Mechanism::Replacement => plan.replace = objects,
961        }
962
963        plan
964    }
965
966    /// Drop objects that have no shard registered from the plan.
967    ///
968    /// A builtin without a shard does not exist in persist yet: either it was added in the target
969    /// version, or the source version predates it entirely (an upgrade from before the builtin
970    /// was introduced can reach an `Evolution` step registered for a later version). In both
971    /// cases the leader allocates the shard during bootstrap, at the target schema, and there is
972    /// nothing to evolve or replace. Applied to every plan, forced or versioned, so that both
973    /// paths agree; a read-only process that reached `migrate_evolve_one` with such an object
974    /// would otherwise bail out and crash-loop.
975    fn drop_shardless(&self, mut plan: Plan) -> Plan {
976        let has_shard = |object: &SystemObjectDescription| {
977            self.system_objects
978                .get(object)
979                .is_none_or(|info| info.shard_id.is_some())
980        };
981        plan.evolve.retain(has_shard);
982        plan.replace.retain(has_shard);
983        plan
984    }
985
986    /// Upgrade the migration shard to the target version.
987    async fn upgrade_migration_shard_version(&self) {
988        let persist = &self.config.persist_client;
989        let diagnostics = Diagnostics {
990            shard_name: "builtin_migration".to_string(),
991            handle_purpose: format!("migration shard upgrade @ {}", self.target_version),
992        };
993
994        persist
995            .upgrade_version::<migration_shard::Key, ShardId, Timestamp, StorageDiff>(
996                self.migration_shard,
997                diagnostics,
998            )
999            .await
1000            .expect("valid usage");
1001    }
1002
1003    /// Migrate the given objects using the `Evolution` mechanism.
1004    async fn migrate_evolve(&self, objects: &[SystemObjectDescription]) -> anyhow::Result<()> {
1005        for object in objects {
1006            self.migrate_evolve_one(object).await?;
1007        }
1008        Ok(())
1009    }
1010
1011    async fn migrate_evolve_one(&self, object: &SystemObjectDescription) -> anyhow::Result<()> {
1012        let persist = &self.config.persist_client;
1013
1014        let Some(object_info) = self.system_objects.get(object) else {
1015            bail!("missing builtin {object:?}");
1016        };
1017        let id = object_info.global_id;
1018
1019        let Some(shard_id) = object_info.shard_id else {
1020            // No shard is registered for this builtin. In leader mode, this is fine, we'll
1021            // register the shard during bootstrap. In read-only mode, we might be racing with the
1022            // leader to register the shard and it's unclear what sort of confusion can arise from
1023            // that -- better to bail out in this case.
1024            if self.config.read_only {
1025                bail!("missing shard ID for builtin {object:?} ({id})");
1026            } else {
1027                return Ok(());
1028            }
1029        };
1030
1031        let target_desc = match object_info.builtin {
1032            Builtin::Table(table) => &table.desc,
1033            Builtin::Source(source) => &source.desc,
1034            Builtin::MaterializedView(mv) => &mv.desc,
1035            _ => bail!("not a storage collection: {object:?}"),
1036        };
1037
1038        let diagnostics = Diagnostics {
1039            shard_name: id.to_string(),
1040            handle_purpose: format!("builtin schema migration @ {}", self.target_version),
1041        };
1042        let source_schema = persist
1043            .latest_schema::<SourceData, (), Timestamp, StorageDiff>(shard_id, diagnostics.clone())
1044            .await
1045            .expect("valid usage");
1046
1047        info!(?object, %id, %shard_id, ?source_schema, ?target_desc, "migrating by evolution");
1048
1049        if self.config.read_only {
1050            // In read-only mode, only check that the new schema is backward compatible.
1051            // We'll register it when/if we restart in leader mode.
1052            if let Some((_, source_desc, _)) = &source_schema {
1053                let old = mz_persist_types::columnar::data_type::<SourceData>(source_desc)?;
1054                let new = mz_persist_types::columnar::data_type::<SourceData>(target_desc)?;
1055                if backward_compatible(&old, &new).is_none() {
1056                    bail!(
1057                        "incompatible schema evolution for {object:?}: \
1058                         {source_desc:?} -> {target_desc:?}"
1059                    );
1060                }
1061            }
1062
1063            return Ok(());
1064        }
1065
1066        let (mut schema_id, mut source_desc) = match source_schema {
1067            Some((schema_id, source_desc, _)) => (schema_id, source_desc),
1068            None => {
1069                // If no schema was previously registered, simply try to register the new one. This
1070                // might fail due to a concurrent registration, in which case we'll fall back to
1071                // `compare_and_evolve_schema`.
1072
1073                debug!(%id, %shard_id, "no previous schema found; registering initial one");
1074                let schema_id = persist
1075                    .register_schema::<SourceData, (), Timestamp, StorageDiff>(
1076                        shard_id,
1077                        target_desc,
1078                        &UnitSchema,
1079                        diagnostics.clone(),
1080                    )
1081                    .await
1082                    .expect("valid usage");
1083                if schema_id.is_some() {
1084                    return Ok(());
1085                }
1086
1087                debug!(%id, %shard_id, "schema registration failed; falling back to CaES");
1088                let (schema_id, source_desc, _) = persist
1089                    .latest_schema::<SourceData, (), Timestamp, StorageDiff>(
1090                        shard_id,
1091                        diagnostics.clone(),
1092                    )
1093                    .await
1094                    .expect("valid usage")
1095                    .expect("known to exist");
1096
1097                (schema_id, source_desc)
1098            }
1099        };
1100
1101        loop {
1102            // Evolving the schema might fail if another process evolved the schema concurrently,
1103            // in which case we need to retry. Most likely the other process evolved the schema to
1104            // our own target schema and the second try is a no-op.
1105
1106            debug!(%id, %shard_id, %schema_id, ?source_desc, ?target_desc, "attempting CaES");
1107            let result = persist
1108                .compare_and_evolve_schema::<SourceData, (), Timestamp, StorageDiff>(
1109                    shard_id,
1110                    schema_id,
1111                    target_desc,
1112                    &UnitSchema,
1113                    diagnostics.clone(),
1114                )
1115                .await
1116                .expect("valid usage");
1117
1118            match result {
1119                CaESchema::Ok(schema_id) => {
1120                    debug!(%id, %shard_id, %schema_id, "schema evolved successfully");
1121                    break;
1122                }
1123                CaESchema::Incompatible => bail!(
1124                    "incompatible schema evolution for {object:?}: \
1125                     {source_desc:?} -> {target_desc:?}"
1126                ),
1127                CaESchema::ExpectedMismatch {
1128                    schema_id: new_id,
1129                    key,
1130                    val: UnitSchema,
1131                } => {
1132                    schema_id = new_id;
1133                    source_desc = key;
1134                }
1135            }
1136        }
1137
1138        Ok(())
1139    }
1140
1141    /// Migrate the given objects using the `Replacement` mechanism.
1142    async fn migrate_replace(
1143        &self,
1144        objects: &[SystemObjectDescription],
1145    ) -> anyhow::Result<BTreeMap<GlobalId, ShardId>> {
1146        if objects.is_empty() {
1147            return Ok(Default::default());
1148        }
1149
1150        let diagnostics = Diagnostics {
1151            shard_name: "builtin_migration".to_string(),
1152            handle_purpose: format!("builtin schema migration @ {}", self.target_version),
1153        };
1154        let (mut persist_write, mut persist_read) =
1155            self.open_migration_shard(diagnostics.clone()).await;
1156
1157        let mut ids_to_replace = BTreeSet::new();
1158        for object in objects {
1159            if let Some(info) = self.system_objects.get(object) {
1160                ids_to_replace.insert(info.global_id);
1161            } else {
1162                bail!("missing id for builtin {object:?}");
1163            }
1164        }
1165
1166        info!(?objects, ?ids_to_replace, "migrating by replacement");
1167
1168        // Fetch replacement shard IDs from the migration shard, or insert new ones if none exist.
1169        // This can fail due to writes by concurrent processes, so we need to retry.
1170        let replaced_shards = loop {
1171            if let Some(shards) = self
1172                .try_get_or_insert_replacement_shards(
1173                    &ids_to_replace,
1174                    &mut persist_write,
1175                    &mut persist_read,
1176                )
1177                .await?
1178            {
1179                break shards;
1180            }
1181        };
1182
1183        Ok(replaced_shards)
1184    }
1185
1186    /// Try to get or insert replacement shards for the given IDs into the migration shard, at
1187    /// `target_version` and `deploy_generation`.
1188    ///
1189    /// This method looks for existing entries in the migration shards and returns those if they
1190    /// are present. Otherwise it generates new shard IDs and tries to insert them.
1191    ///
1192    /// The result of this call is `None` if no existing entries were found and inserting new ones
1193    /// failed because of a concurrent write to the migration shard. In this case, the caller is
1194    /// expected to retry.
1195    async fn try_get_or_insert_replacement_shards(
1196        &self,
1197        ids_to_replace: &BTreeSet<GlobalId>,
1198        persist_write: &mut WriteHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1199        persist_read: &mut ReadHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1200    ) -> anyhow::Result<Option<BTreeMap<GlobalId, ShardId>>> {
1201        let upper = persist_write.fetch_recent_upper().await;
1202        let write_ts = *upper.as_option().expect("migration shard not sealed");
1203
1204        let mut ids_to_replace = ids_to_replace.clone();
1205        let mut replaced_shards = BTreeMap::new();
1206
1207        // Another process might already have done a shard replacement at our version and
1208        // generation, in which case we can directly reuse the replacement shards.
1209        //
1210        // Note that we can't assume that the previous process had the same `ids_to_replace` as we
1211        // do. The set of migrations to run depends on both the source and the target version, and
1212        // the migration shard is not keyed by source version. The previous writer might have seen
1213        // a different source version, if there was a concurrent migration by a leader process.
1214        if let Some(read_ts) = write_ts.step_back() {
1215            let pred = |key: &migration_shard::Key| {
1216                key.build_version == self.target_version
1217                    && key.deploy_generation == Some(self.deploy_generation)
1218            };
1219            if let Some(entries) = read_migration_shard(persist_read, read_ts, pred).await {
1220                for (key, shard_id) in entries {
1221                    let id = GlobalId::System(key.global_id);
1222                    if ids_to_replace.remove(&id) {
1223                        replaced_shards.insert(id, shard_id);
1224                    }
1225                }
1226
1227                debug!(
1228                    %read_ts, ?replaced_shards, ?ids_to_replace,
1229                    "found existing entries in migration shard",
1230                );
1231            }
1232
1233            if ids_to_replace.is_empty() {
1234                return Ok(Some(replaced_shards));
1235            }
1236        }
1237
1238        // Generate new shard IDs and attempt to insert them into the migration shard. If we get a
1239        // CaA failure at `write_ts` that means a concurrent process has inserted in the meantime
1240        // and we need to re-check the migration shard contents.
1241        let mut updates = Vec::new();
1242        for id in ids_to_replace {
1243            let shard_id = ShardId::new();
1244            replaced_shards.insert(id, shard_id);
1245
1246            let GlobalId::System(global_id) = id else {
1247                bail!("attempt to migrate a non-system collection: {id}");
1248            };
1249            let key = migration_shard::Key {
1250                global_id,
1251                build_version: self.target_version.clone(),
1252                deploy_generation: Some(self.deploy_generation),
1253            };
1254            updates.push(((key, shard_id), write_ts, 1));
1255        }
1256
1257        let upper = Antichain::from_elem(write_ts);
1258        let new_upper = Antichain::from_elem(write_ts.step_forward());
1259        debug!(%write_ts, "attempting insert into migration shard");
1260        let result = persist_write
1261            .compare_and_append(updates, upper, new_upper)
1262            .await
1263            .expect("valid usage");
1264
1265        match result {
1266            Ok(()) => {
1267                debug!(
1268                    %write_ts, ?replaced_shards,
1269                    "successfully inserted into migration shard"
1270                );
1271                Ok(Some(replaced_shards))
1272            }
1273            Err(_mismatch) => Ok(None),
1274        }
1275    }
1276
1277    /// Open writer and reader for the migration shard.
1278    async fn open_migration_shard(
1279        &self,
1280        diagnostics: Diagnostics,
1281    ) -> (
1282        WriteHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1283        ReadHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1284    ) {
1285        let persist = &self.config.persist_client;
1286
1287        persist
1288            .open(
1289                self.migration_shard,
1290                Arc::new(migration_shard::KeySchema),
1291                Arc::new(ShardIdSchema),
1292                diagnostics,
1293                USE_CRITICAL_SINCE_CATALOG.get(persist.dyncfgs()),
1294            )
1295            .await
1296            .expect("valid usage")
1297    }
1298
1299    /// Open a [`SinceHandle`] for the migration shard.
1300    async fn open_migration_shard_since(
1301        &self,
1302        diagnostics: Diagnostics,
1303    ) -> SinceHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff> {
1304        self.config
1305            .persist_client
1306            .open_critical_since(
1307                self.migration_shard,
1308                // TODO: We may need to use a different critical reader
1309                // id for this if we want to be able to introspect it via SQL.
1310                PersistClient::CONTROLLER_CRITICAL_SINCE,
1311                Opaque::encode(&i64::MIN),
1312                diagnostics.clone(),
1313            )
1314            .await
1315            .expect("valid usage")
1316    }
1317
1318    /// Update the fingerprints for `migrated_items`.
1319    ///
1320    /// Returns the new fingerprints. Also asserts that the current fingerprints of all other
1321    /// system items match their builtin definitions.
1322    fn update_fingerprints(
1323        &self,
1324        migrated_items: &BTreeSet<SystemObjectDescription>,
1325    ) -> anyhow::Result<BTreeMap<SystemObjectDescription, String>> {
1326        let mut new_fingerprints = BTreeMap::new();
1327        for (object, object_info) in &self.system_objects {
1328            let id = object_info.global_id;
1329            let builtin = object_info.builtin;
1330
1331            let fingerprint = builtin.fingerprint();
1332            if fingerprint == object_info.fingerprint {
1333                continue; // fingerprint unchanged, nothing to do
1334            }
1335
1336            // Fingerprint mismatch is expected for a migrated item.
1337            let migrated = migrated_items.contains(object);
1338            // Some builtin types have schemas but no durable state. No migration needed for those.
1339            let ephemeral = matches!(
1340                builtin,
1341                Builtin::Log(_) | Builtin::View(_) | Builtin::Index(_),
1342            );
1343
1344            if migrated || ephemeral {
1345                new_fingerprints.insert(object.clone(), fingerprint);
1346            } else if builtin.runtime_alterable() {
1347                // Runtime alterable builtins have no meaningful builtin fingerprint, and a
1348                // sentinel value stored in the catalog.
1349                assert_eq!(
1350                    object_info.fingerprint, RUNTIME_ALTERABLE_FINGERPRINT_SENTINEL,
1351                    "fingerprint mismatch for runtime-alterable builtin {object:?} ({id})",
1352                );
1353            } else {
1354                panic!(
1355                    "fingerprint mismatch for builtin {builtin:?} ({id}): {} != {}",
1356                    fingerprint, object_info.fingerprint,
1357                );
1358            }
1359        }
1360
1361        Ok(new_fingerprints)
1362    }
1363
1364    /// Perform cleanup of migration state, i.e. the migration shard.
1365    ///
1366    /// Returns a list of shards to finalize, and a `Future` that must be run after the shard
1367    /// finalization has been durably enqueued. The `Future` is used to remove entries from the
1368    /// migration shard only after we know the respective shards will be finalized. Removing
1369    /// entries immediately would risk leaking the shards.
1370    ///
1371    /// We only perform cleanup in leader mode, to keep the durable state changes made by read-only
1372    /// processes a minimal as possible. Given that Materialize doesn't support version downgrades,
1373    /// it is safe to assume that any state for versions below the `target_version` is not needed
1374    /// anymore and can be cleaned up.
1375    ///
1376    /// Note that it is fine for cleanup to sometimes fail or be skipped. The size of the migration
1377    /// shard should always be pretty small, so keeping migration state around for longer isn't a
1378    /// concern. As a result, we can keep the logic simple here and skip doing cleanup in response
1379    /// to transient failures, instead of retrying.
1380    async fn cleanup(&self) -> anyhow::Result<(BTreeSet<ShardId>, BoxFuture<'static, ()>)> {
1381        let noop_action = async {}.boxed();
1382        let noop_result = (BTreeSet::new(), noop_action);
1383
1384        if self.config.read_only {
1385            return Ok(noop_result);
1386        }
1387
1388        let diagnostics = Diagnostics {
1389            shard_name: "builtin_migration".to_string(),
1390            handle_purpose: "builtin schema migration cleanup".into(),
1391        };
1392        let (mut persist_write, mut persist_read) =
1393            self.open_migration_shard(diagnostics.clone()).await;
1394        let mut persist_since = self.open_migration_shard_since(diagnostics.clone()).await;
1395
1396        let upper = persist_write.fetch_recent_upper().await.clone();
1397        let write_ts = *upper.as_option().expect("migration shard not sealed");
1398        let Some(read_ts) = write_ts.step_back() else {
1399            return Ok(noop_result);
1400        };
1401
1402        // Collect old entries to remove.
1403        let pred = |key: &migration_shard::Key| key.build_version < self.target_version;
1404        let Some(stale_entries) = read_migration_shard(&mut persist_read, read_ts, pred).await
1405        else {
1406            return Ok(noop_result);
1407        };
1408
1409        debug!(
1410            ?stale_entries,
1411            "cleaning migration shard up to version {}", self.target_version,
1412        );
1413
1414        let current_shards: BTreeMap<_, _> = self
1415            .system_objects
1416            .values()
1417            .filter_map(|o| o.shard_id.map(|shard_id| (o.global_id, shard_id)))
1418            .collect();
1419
1420        let mut shards_to_finalize = BTreeSet::new();
1421        let mut retractions = Vec::new();
1422        for (key, shard_id) in stale_entries {
1423            // The migration shard contains both shards created during aborted upgrades and shards
1424            // created during successful upgrades. The latter may still be in use, so we have to
1425            // check and only finalize those that aren't anymore.
1426            let gid = GlobalId::System(key.global_id);
1427            if current_shards.get(&gid) != Some(&shard_id) {
1428                shards_to_finalize.insert(shard_id);
1429            }
1430
1431            retractions.push(((key, shard_id), write_ts, -1));
1432        }
1433
1434        let cleanup_action = async move {
1435            if !retractions.is_empty() {
1436                let new_upper = Antichain::from_elem(write_ts.step_forward());
1437                let result = persist_write
1438                    .compare_and_append(retractions, upper, new_upper)
1439                    .await
1440                    .expect("valid usage");
1441                match result {
1442                    Ok(()) => debug!("cleaned up migration shard"),
1443                    Err(mismatch) => debug!(?mismatch, "migration shard cleanup failed"),
1444                }
1445            }
1446        }
1447        .boxed();
1448
1449        // Downgrade the since, to enable some compaction.
1450        let o = persist_since.opaque().clone();
1451        let new_since = Antichain::from_elem(read_ts);
1452        let result = persist_since
1453            .maybe_compare_and_downgrade_since(&o, (&o, &new_since))
1454            .await;
1455        soft_assert_or_log!(result.is_none_or(|r| r.is_ok()), "opaque mismatch");
1456
1457        Ok((shards_to_finalize, cleanup_action))
1458    }
1459}
1460
1461/// Read the migration shard at the given timestamp, returning all entries that match the given
1462/// predicate.
1463///
1464/// Returns `None` if the migration shard contains no matching entries, or if it isn't readable at
1465/// `read_ts`.
1466async fn read_migration_shard<P>(
1467    persist_read: &mut ReadHandle<migration_shard::Key, ShardId, Timestamp, StorageDiff>,
1468    read_ts: Timestamp,
1469    predicate: P,
1470) -> Option<Vec<(migration_shard::Key, ShardId)>>
1471where
1472    P: Fn(&migration_shard::Key) -> bool,
1473{
1474    let as_of = Antichain::from_elem(read_ts);
1475    let updates = persist_read.snapshot_and_fetch(as_of).await.ok()?;
1476
1477    assert!(
1478        updates.iter().all(|(_, _, diff)| *diff == 1),
1479        "migration shard contains invalid diffs: {updates:?}",
1480    );
1481
1482    let entries: Vec<_> = updates
1483        .into_iter()
1484        .map(|(data, _, _)| data)
1485        .filter(move |(key, _)| predicate(key))
1486        .collect();
1487
1488    (!entries.is_empty()).then_some(entries)
1489}
1490
1491/// A plan to migrate between two versions.
1492#[derive(Debug, Default)]
1493struct Plan {
1494    /// Objects to migrate using the `Evolution` mechanism.
1495    evolve: Vec<SystemObjectDescription>,
1496    /// Objects to migrate using the `Replacement` mechanism.
1497    replace: Vec<SystemObjectDescription>,
1498}
1499
1500/// Types and persist codec impls for the migration shard used by the `Replacement` mechanism.
1501mod migration_shard {
1502    use std::fmt;
1503    use std::str::FromStr;
1504
1505    use arrow::array::{StringArray, StringBuilder};
1506    use bytes::{BufMut, Bytes};
1507    use mz_persist_types::Codec;
1508    use mz_persist_types::codec_impls::{
1509        SimpleColumnarData, SimpleColumnarDecoder, SimpleColumnarEncoder,
1510    };
1511    use mz_persist_types::columnar::Schema;
1512    use mz_persist_types::stats::NoneStats;
1513    use semver::Version;
1514    use serde::{Deserialize, Serialize};
1515
1516    #[derive(Debug, Clone, Eq, Ord, PartialEq, PartialOrd, Serialize, Deserialize)]
1517    pub(super) struct Key {
1518        pub(super) global_id: u64,
1519        pub(super) build_version: Version,
1520        // Versions < 26.0 didn't include the deploy generation. As long as we still might
1521        // encounter migration shard entries that don't have it, we need to keep this an `Option`
1522        // and keep supporting both key formats.
1523        pub(super) deploy_generation: Option<u64>,
1524    }
1525
1526    impl fmt::Display for Key {
1527        fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1528            if self.deploy_generation.is_some() {
1529                // current format
1530                let s = serde_json::to_string(self).expect("JSON serializable");
1531                f.write_str(&s)
1532            } else {
1533                // pre-26.0 format
1534                write!(f, "{}-{}", self.global_id, self.build_version)
1535            }
1536        }
1537    }
1538
1539    impl FromStr for Key {
1540        type Err = String;
1541
1542        fn from_str(s: &str) -> Result<Self, String> {
1543            // current format
1544            if let Ok(key) = serde_json::from_str(s) {
1545                return Ok(key);
1546            };
1547
1548            // pre-26.0 format
1549            let parts: Vec<_> = s.splitn(2, '-').collect();
1550            let &[global_id, build_version] = parts.as_slice() else {
1551                return Err(format!("invalid Key '{s}'"));
1552            };
1553            let global_id = global_id.parse::<u64>().map_err(|e| e.to_string())?;
1554            let build_version = build_version
1555                .parse::<Version>()
1556                .map_err(|e| e.to_string())?;
1557            Ok(Key {
1558                global_id,
1559                build_version,
1560                deploy_generation: None,
1561            })
1562        }
1563    }
1564
1565    impl Default for Key {
1566        fn default() -> Self {
1567            Self {
1568                global_id: Default::default(),
1569                build_version: Version::new(0, 0, 0),
1570                deploy_generation: Some(0),
1571            }
1572        }
1573    }
1574
1575    impl Codec for Key {
1576        type Schema = KeySchema;
1577        type Storage = ();
1578
1579        fn codec_name() -> String {
1580            "TableKey".into()
1581        }
1582
1583        fn encode<B: BufMut>(&self, buf: &mut B) {
1584            buf.put(self.to_string().as_bytes())
1585        }
1586
1587        fn decode<'a>(buf: &'a [u8], _schema: &KeySchema) -> Result<Self, String> {
1588            let s = str::from_utf8(buf).map_err(|e| e.to_string())?;
1589            s.parse()
1590        }
1591
1592        fn encode_schema(_schema: &KeySchema) -> Bytes {
1593            Bytes::new()
1594        }
1595
1596        fn decode_schema(buf: &Bytes) -> Self::Schema {
1597            assert_eq!(*buf, Bytes::new());
1598            KeySchema
1599        }
1600    }
1601
1602    impl SimpleColumnarData for Key {
1603        type ArrowBuilder = StringBuilder;
1604        type ArrowColumn = StringArray;
1605
1606        fn goodbytes(builder: &Self::ArrowBuilder) -> usize {
1607            builder.values_slice().len()
1608        }
1609
1610        fn push(&self, builder: &mut Self::ArrowBuilder) {
1611            builder.append_value(&self.to_string());
1612        }
1613
1614        fn push_null(builder: &mut Self::ArrowBuilder) {
1615            builder.append_null();
1616        }
1617
1618        fn read(&mut self, idx: usize, column: &Self::ArrowColumn) {
1619            *self = column.value(idx).parse().expect("valid Key");
1620        }
1621    }
1622
1623    #[derive(Debug, PartialEq)]
1624    pub(super) struct KeySchema;
1625
1626    impl Schema<Key> for KeySchema {
1627        type ArrowColumn = StringArray;
1628        type Statistics = NoneStats;
1629        type Decoder = SimpleColumnarDecoder<Key>;
1630        type Encoder = SimpleColumnarEncoder<Key>;
1631
1632        fn encoder(&self) -> anyhow::Result<SimpleColumnarEncoder<Key>> {
1633            Ok(SimpleColumnarEncoder::default())
1634        }
1635
1636        fn decoder(&self, col: StringArray) -> anyhow::Result<SimpleColumnarDecoder<Key>> {
1637            Ok(SimpleColumnarDecoder::new(col))
1638        }
1639    }
1640}
1641
1642#[cfg(test)]
1643#[path = "builtin_schema_migration_tests.rs"]
1644mod tests;