Skip to content

Commit d2ca18c

Browse files
committed
adapter: address review comments on 0dt migrated-MV hydration
- Seed the caught-up gate's transitive-dependent walk from all of `new_builtin_collections`, not just the MVs, and require the write frontier to be within the allowed lag of `now` on the no-live-frontier path so a hydrated-then-stalled collection no longer passes the gate forever. - Document the `enable_0dt_hydrate_migrated_builtin_mvs` break-glass flag: it is read once at startup, so flipping it means changing the setting on the leader and restarting, and it is not an exact revert. - Correct the read-only write-enable comments to match the code: the `allow_writes_in_read_only` tripwire blocks promotion, exclusive shard ownership holds per (build version, deploy generation), a new builtin MV cannot be write-enabled, and cite PR #35402 for the v26.17 leader floor.
1 parent 1afc066 commit d2ca18c

7 files changed

Lines changed: 229 additions & 85 deletions

File tree

‎doc/developer/design/20251015_builtin_schema_migration.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,11 @@ Because this environment exclusively owns the replacement shard, the read-only p
8181
This lets a migrated builtin materialized view and its dependents hydrate before cut-over instead of all at once at cut-over.
8282
It is safe only for the self-owned replacement shard, never a shard the leader is still serving, which is why it applies to shard replacement and not schema evolution.
8383

84+
The force-write is conditional on two things.
85+
First, the leader must be at v26.17 or later, because every builtin materialized view reads the catalog shard and only leaders from that version on keep its frontier advancing with the current time; against an older leader the dataflow would sit at a stale frontier.
86+
Second, the `enable_0dt_hydrate_migrated_builtin_mvs` feature flag must be on; it exists as a break-glass revert.
87+
When either condition does not hold, the migrated materialized views and their dependents are instead excluded from the 0dt caught-up check, which is the older behaviour: promotion proceeds without them and they hydrate at cut-over.
88+
8489
A leader process performing shard replacement performs the same steps as in read-only mode.
8590
Additionally, it cleans up durable state written by earlier versions and/or deploy generations by:
8691
- arranging for the previous shards used by the migrated storage collections to be finalized

‎src/adapter-types/src/dyncfgs.rs‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -95,9 +95,11 @@ pub const ENABLE_0DT_HYDRATE_MIGRATED_BUILTIN_MVS: Config<bool> = Config::new(
9595
true,
9696
"Write-enable replacement-migrated builtin materialized views while read-only during a 0dt \
9797
deployment, so they hydrate before cut-over and keep gating promotion. Emergency break-glass \
98-
flag: disabling reverts to excluding migrated MVs (and their dependents) from the caught-up \
99-
check, which lets promotion proceed with them unhydrated and hydrate at cut-over instead. \
100-
Only takes effect when the leader is new enough for the write to make progress.",
98+
flag: disabling excludes migrated MVs (and their dependents) from the caught-up check again, \
99+
so promotion proceeds with them unhydrated. Not an exact revert: a collection with no live \
100+
leader frontier must be hydrated either way. Only takes effect when the leader is new enough \
101+
for the write to make progress, and is read once at startup, so changing it means setting it \
102+
on the leader and restarting the new deployment.",
101103
);
102104

103105
/// Enable logging of statement lifecycle events in mz_internal.mz_statement_lifecycle_history.

‎src/adapter/src/coord.rs‎

Lines changed: 30 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -250,10 +250,10 @@ mod validity;
250250
///
251251
/// Every builtin materialized view reads `mz_internal.mz_catalog_raw`, so its dataflow only makes
252252
/// progress up to the catalog shard's frontier. Holding that frontier at the current time is the
253-
/// leader's job, and leaders only started doing it in v26.17. Write-enable such an MV against an
254-
/// older leader and it sits at a stale frontier and never reports caught up, which blocks
255-
/// promotion outright instead of merely leaving the collection cold at cut-over. We still support
256-
/// upgrading from before v26.17, so that leader is a real case, not a hypothetical.
253+
/// leader's job, and leaders only started doing it in v26.17 (PR #35402). Write-enable such an MV
254+
/// against an older leader and it sits at a stale frontier and never reports caught up, which
255+
/// blocks promotion outright instead of merely leaving the collection cold at cut-over. We still
256+
/// support upgrading from before v26.17, so that leader is a real case, not a hypothetical.
257257
const MIN_LEADER_VERSION_FOR_MIGRATED_MV_WRITES: Version = Version::new(26, 17, 0);
258258

259259
/// A pool of pre-allocated user IDs to avoid per-DDL persist writes.
@@ -2799,8 +2799,9 @@ impl Coordinator {
27992799
self.ship_dataflow(df_desc, mview.cluster_id, mview.target_replica)
28002800
.await;
28012801

2802-
// If this is a replacement MV, it must remain read-only until the replacement
2803-
// gets applied.
2802+
// A pending `REPLACEMENT FOR` MV must stay read-only until
2803+
// `ALTER ... APPLY REPLACEMENT` swaps it in. Unrelated to the
2804+
// builtin-migration `Replacement` mechanism below.
28042805
if mview.replacement_target.is_none() {
28052806
let gid = mview.global_id_writes();
28062807
if hydrate_migrated_mvs
@@ -2811,6 +2812,10 @@ impl Coordinator {
28112812
// writing it while read-only hydrates the MV and its dependents before
28122813
// cut-over. An `Evolution`-migrated MV reuses the leader's live shard
28132814
// and must never reach here.
2815+
//
2816+
// A *new* builtin MV gets no such treatment: its shard allocation
2817+
// lives only in this read-only savepoint, so the promoted leader
2818+
// allocates a different shard and discards whatever we wrote.
28142819
self.controller
28152820
.compute
28162821
.allow_writes_in_read_only(mview.cluster_id, gid)
@@ -5187,32 +5192,31 @@ pub fn serve(
51875192

51885193
// A collection that can't advance its write frontier in read-only mode
51895194
// stalls its transitive dependents too, so exclude those from the caught-up
5190-
// check as well. That's new builtin MVs, whose fresh shard has no writer until
5191-
// this deployment promotes, plus migrated MVs whenever the leader is too old for
5192-
// them to write. An excluded dependent may still be hydrating right after
5193-
// promotion, a brief blip we accept because these MVs are small and get a writer
5194-
// at cut-over.
5195+
// check as well. That's every *new* builtin collection, whose fresh shard has no
5196+
// writer until this deployment promotes, plus migrated MVs whenever the leader is
5197+
// too old for them to write. An excluded dependent may still be hydrating right
5198+
// after promotion, a brief blip we accept because these collections are small and
5199+
// get a writer at cut-over.
51955200
//
5196-
// A migrated builtin *table* needs no such treatment even though a builtin MV can
5197-
// read one (`mz_clusters` joins `mz_cluster_replica_size_internal`):
5198-
// `read_only_mode_table_worker` keeps advancing the uppers of migrated tables, so
5199-
// an MV over one still catches up.
5200-
let new_builtin_mvs = new_builtin_collections
5201-
.iter()
5202-
.map(|global_id| {
5203-
catalog
5204-
.state()
5205-
.try_get_entry_by_global_id(global_id)
5206-
.expect("new builtin collections have catalog entries")
5207-
})
5208-
.filter(|entry| entry.is_materialized_view())
5209-
.map(|entry| entry.id());
5201+
// Seeded from all of `new_builtin_collections`, not just the MVs: a new builtin
5202+
// table or source has no read-only writer either (`register_table_collections`
5203+
// retains only *migrated* tables), so an MV reading one never advances past its
5204+
// empty frontier. A *migrated* table is the opposite case, even though a builtin
5205+
// MV can read one (`mz_clusters` joins `mz_cluster_replica_size_internal`):
5206+
// `read_only_mode_table_worker` keeps advancing migrated tables' uppers.
5207+
let new_builtin_items = new_builtin_collections.iter().map(|global_id| {
5208+
catalog
5209+
.state()
5210+
.try_get_entry_by_global_id(global_id)
5211+
.expect("new builtin collections have catalog entries")
5212+
.id()
5213+
});
52105214
let frozen_migrated_mvs = migrated_storage_collections_0dt
52115215
.iter()
52125216
.copied()
52135217
.filter(|_| !hydrate_migrated_mvs)
52145218
.filter(|id| catalog.state().get_entry(id).is_materialized_view());
5215-
let mut todo: Vec<_> = new_builtin_mvs.chain(frozen_migrated_mvs).collect();
5219+
let mut todo: Vec<_> = new_builtin_items.chain(frozen_migrated_mvs).collect();
52165220
while let Some(item_id) = todo.pop() {
52175221
let entry = catalog.state().get_entry(&item_id);
52185222
exclude_collections.extend(entry.global_ids());

‎src/adapter/src/coord/caught_up.rs‎

Lines changed: 31 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -271,6 +271,11 @@ impl Coordinator {
271271
// `snapshot_latest` requires that the collection consolidates to a
272272
// set. `mz_cluster_replica_frontiers` is a controller-managed builtin
273273
// written with ±1 diffs, so it satisfies that invariant.
274+
//
275+
// NOTE: these are the leader's frontiers only because we read the leader's shard. A
276+
// release that `Replacement`-migrates `mz_cluster_replica_frontiers` itself, or a test
277+
// forcing replacement across all builtins, hands us a shard we write ourselves, and the
278+
// lag check below then compares this deployment against itself.
274279
let live_frontiers = self
275280
.controller
276281
.storage_collections
@@ -634,12 +639,37 @@ impl Coordinator {
634639
self.controller.storage.collection_hydrated(id)?
635640
}
636641
};
642+
643+
// Also require the frontier to be within the allowed lag, the bound the
644+
// live-frontier path applies, with `now` standing in for the missing live
645+
// frontier. Hydration is one-shot: a collection that hydrated and then stalled
646+
// would otherwise satisfy this branch forever.
647+
//
648+
// NOTE: there is deliberately no `cutoff` escape hatch here. A frontier frozen
649+
// at the minimum is exactly what this gate must catch, so a collection stuck
650+
// here blocks promotion until `with_0dt_deployment_max_wait` elapses.
651+
let write_frontier_plus_allowed_lag = Antichain::from_iter(
652+
write_frontier
653+
.iter()
654+
.map(|t| t.step_forward_by(&allowed_lag)),
655+
);
656+
let within_lag = PartialOrder::less_equal(
657+
&Antichain::from_elem(now),
658+
&write_frontier_plus_allowed_lag,
659+
);
660+
637661
tracing::info!(
638662
?write_frontier,
639663
%collection_hydrated,
664+
%within_lag,
665+
?allowed_lag,
666+
?now,
640667
"collection {id} not in live frontiers"
641668
);
642-
if write_frontier.less_equal(&Timestamp::minimum()) || !collection_hydrated {
669+
if write_frontier.less_equal(&Timestamp::minimum())
670+
|| !collection_hydrated
671+
|| !within_lag
672+
{
643673
all_caught_up = false;
644674
}
645675
continue;

‎src/compute-client/src/controller.rs‎

Lines changed: 19 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1183,26 +1183,40 @@ impl ComputeController {
11831183

11841184
/// Like [`Self::allow_writes`], but takes effect even in read-only mode.
11851185
///
1186-
/// The caller must guarantee that no other environment writes the collection's output shard.
1186+
/// The caller must guarantee that no leader environment writes the collection's output shard.
11871187
/// In a 0dt deployment that means a shard this environment created for itself, the replacement
11881188
/// shard of a `Replacement`-migrated builtin collection, never one the leader is still serving
11891189
/// from. `Evolution` migrates in place and reuses the leader's shard, so it must not reach this
11901190
/// path. Violating the guarantee races two writers on one shard.
11911191
///
1192+
/// NOTE: ownership is exclusive per (build version, deploy generation), not per process: the
1193+
/// migration shard entry naming the shard is keyed by that pair and a read-only catalog open is
1194+
/// a savepoint, so two read-only processes of one generation both write it. Same shape as a
1195+
/// multi-replica materialized view, which the self-correcting persist sink tolerates (see the
1196+
/// `mz_compute::sink::materialized_view` module docs).
1197+
///
11921198
/// This is the compute-side counterpart to the storage controller's `force_writable` handling
11931199
/// of migrated storage collections. Migrated builtin tables are storage collections that
11941200
/// storage force-writes read-only; migrated builtin MVs are compute collections that only this
11951201
/// path can force-write. Both rest on the same guarantee (this environment exclusively owns the
11961202
/// replacement shard) but run on separate write paths, so each needs its own bypass.
1203+
///
1204+
/// NOTE: the replica-side handler enables persist compaction process-wide on the clusterd
1205+
/// (`ComputeState::handle_allow_writes`), which this path is the first to trigger inside a
1206+
/// read-only deployment.
11971207
pub fn allow_writes_in_read_only(
11981208
&mut self,
11991209
instance_id: ComputeInstanceId,
12001210
collection_id: GlobalId,
12011211
) -> Result<(), CollectionUpdateError> {
1202-
// Every builtin eligible for the read-only bypass has a system id. A non-system id means
1203-
// the caller's `Replacement`-only invariant broke, so degrade to the read-only no-op (cold
1204-
// collection at cut-over) rather than risk writing a shard the leader still serves. Mirrors
1205-
// the storage-side tripwire in `StorageController::register_introspection_collection`.
1212+
// Every builtin eligible for this bypass has a system id, so a non-system id means the
1213+
// caller's `Replacement`-only invariant broke. No-op rather than risk writing a shard the
1214+
// leader still serves. The collection then sits in the caught-up gate on an unwritten
1215+
// shard and blocks promotion, which is the loud, safe direction to fail.
1216+
//
1217+
// Storage asserts the same invariant hard, in
1218+
// `StorageController::register_introspection_collection`. The asymmetry is deliberate: a
1219+
// soft panic keeps a caller bug visible in CI and Sentry without downing production.
12061220
if self.read_only && !collection_id.is_system() {
12071221
soft_panic_or_log!(
12081222
"allow_writes_in_read_only called for non-system collection {collection_id}; \

0 commit comments

Comments
 (0)