1use std::pin::pin;
13use std::sync::Arc;
14use std::time::Duration;
15
16use mz_adapter::ResultExt;
17use mz_catalog::durable::{
18 BootstrapArgs, CatalogError, DurableCatalogState, Metrics, OpenableDurableCatalogState,
19};
20use mz_controller_types::ReplicaId;
21use mz_ore::channel::trigger;
22use mz_ore::exit;
23use mz_ore::halt;
24use mz_ore::str::separated;
25use mz_persist_client::PersistClient;
26use mz_repr::{CatalogItemId, Timestamp};
27use mz_sql::catalog::EnvironmentId;
28use tokio::sync::oneshot;
29use tracing::info;
30
31use crate::BUILD_INFO;
32use crate::deployment::state::DeploymentState;
33
34pub struct CatchupConfig {
36 pub boot_ts: Timestamp,
37 pub environment_id: EnvironmentId,
38 pub persist_client: PersistClient,
39 pub deploy_generation: u64,
40 pub deployment_state: DeploymentState,
41 pub catalog_metrics: Arc<Metrics>,
42 pub caught_up_max_wait: Duration,
43 pub ddl_check_interval: Duration,
44 pub panic_after_timeout: bool,
45 pub bootstrap_args: BootstrapArgs,
46}
47
48pub async fn preflight_0dt(
50 openable_adapter_storage: &mut dyn OpenableDurableCatalogState,
51 deploy_generation: u64,
52) -> Result<bool, CatalogError> {
53 if !openable_adapter_storage.is_initialized().await? {
54 info!("catalog not initialized; booting with writes allowed");
55 return Ok(false);
56 }
57
58 let catalog_generation = openable_adapter_storage.get_deployment_generation().await?;
59 info!(%catalog_generation, %deploy_generation, "catalog initialized");
60 if catalog_generation < deploy_generation {
61 info!("this deployment is a new generation; booting in read only mode");
62 Ok(true)
63 } else if catalog_generation == deploy_generation {
64 info!("this deployment is the current generation; booting with writes allowed");
65 Ok(false)
66 } else {
67 exit!(0, "this deployment has been fenced out");
68 }
69}
70
71pub fn spawn_catchup(
79 CatchupConfig {
80 boot_ts,
81 environment_id,
82 persist_client,
83 deploy_generation,
84 deployment_state,
85 catalog_metrics,
86 caught_up_max_wait,
87 ddl_check_interval,
88 panic_after_timeout,
89 bootstrap_args,
90 }: CatchupConfig,
91 mut caught_up_receiver: trigger::Receiver,
92 bootstrapped: oneshot::Receiver<(u64, u64)>,
93) {
94 mz_ore::task::spawn(|| "deployment_catchup", async move {
95 let mut skip_catchup = deployment_state.set_catching_up();
96
97 let initial_ids = tokio::select! {
98 biased;
99
100 () = &mut skip_catchup => None,
101 result = bootstrapped => match result {
102 Ok(ids) => Some(ids),
103 Err(_) => return,
104 },
105 };
106
107 if let Some((initial_next_user_item_id, initial_next_replica_id)) = initial_ids {
108 info!(
109 %initial_next_user_item_id,
110 %initial_next_replica_id,
111 ?caught_up_max_wait,
112 "waiting for deployment to be caught up"
113 );
114
115 let mut caught_up_max_wait_fut = pin!(tokio::time::sleep(caught_up_max_wait));
116
117 let mut check_ddl_changes_interval = tokio::time::interval(ddl_check_interval);
118 check_ddl_changes_interval
119 .set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
120
121 let mut should_skip_catchup = false;
122 loop {
123 tokio::select! {
124 biased;
125
126 () = &mut skip_catchup => {
127 info!("skipping waiting for deployment to catch up due to administrator request");
128 should_skip_catchup = true;
129 break;
130 }
131 () = &mut caught_up_receiver => {
132 info!("deployment caught up");
133 break;
134 }
135 () = &mut caught_up_max_wait_fut => {
136 if panic_after_timeout {
137 panic!("not caught up within {:?}", caught_up_max_wait);
138 }
139 info!("not caught up within {:?}, proceeding now", caught_up_max_wait);
140 break;
141 }
142 _ = check_ddl_changes_interval.tick() => {
143 check_ddl_changes(
144 boot_ts,
145 persist_client.clone(),
146 environment_id.clone(),
147 deploy_generation,
148 Arc::clone(&catalog_metrics),
149 bootstrap_args.clone(),
150 initial_next_user_item_id,
151 initial_next_replica_id,
152 )
153 .await;
154 }
155 }
156 }
157
158 if !should_skip_catchup {
161 check_ddl_changes(
162 boot_ts,
163 persist_client.clone(),
164 environment_id.clone(),
165 deploy_generation,
166 Arc::clone(&catalog_metrics),
167 bootstrap_args.clone(),
168 initial_next_user_item_id,
169 initial_next_replica_id,
170 )
171 .await;
172 }
173 } else {
174 info!("skipping bootstrap and catch-up due to administrator request");
175 }
176
177 let promoted = deployment_state.set_ready_to_promote();
179 info!("announced as ready to promote; waiting for promotion");
180 promoted.await;
181
182 info!("promoted; attempting takeover");
184
185 let openable_adapter_storage = mz_catalog::durable::persist_backed_catalog_state(
199 persist_client.clone(),
200 environment_id.organization_id(),
201 BUILD_INFO.semver_version(),
202 Some(deploy_generation),
203 Arc::clone(&catalog_metrics),
204 )
205 .await
206 .expect("incompatible catalog/persist version");
207
208 let _catalog = openable_adapter_storage
209 .open(boot_ts, &bootstrap_args)
210 .await
211 .unwrap_or_terminate("unexpected error while fencing out old deployment");
212
213 halt!("fenced out old deployment; rebooting as leader")
215 });
216}
217
218async fn check_ddl_changes(
226 boot_ts: Timestamp,
227 persist_client: PersistClient,
228 environment_id: EnvironmentId,
229 deploy_generation: u64,
230 catalog_metrics: Arc<Metrics>,
231 bootstrap_args: BootstrapArgs,
232 initial_next_user_item_id: u64,
233 initial_next_replica_id: u64,
234) {
235 let openable_adapter_storage = mz_catalog::durable::persist_backed_catalog_state(
236 persist_client,
237 environment_id.organization_id(),
238 BUILD_INFO.semver_version(),
239 Some(deploy_generation),
240 catalog_metrics,
241 )
242 .await
243 .expect("incompatible catalog/persist version");
244
245 let mut catalog = openable_adapter_storage
246 .open_savepoint(boot_ts, &bootstrap_args)
247 .await
248 .unwrap_or_terminate("can open in savepoint mode");
249
250 let _ = catalog
253 .sync_to_current_updates()
254 .await
255 .unwrap_or_terminate("unexpected error while draining initial catalog updates");
256 let tx = catalog
257 .transaction()
258 .await
259 .unwrap_or_terminate("unexpected error while getting transaction");
260
261 let new_replicas = tx
266 .get_cluster_replicas()
267 .filter_map(|replica| match replica.replica_id {
268 ReplicaId::User(id) if id >= initial_next_replica_id => Some(replica),
269 _ => None,
270 })
271 .collect::<Vec<_>>();
272
273 let new_objects = tx
274 .get_items()
275 .filter_map(|item| match item.id {
276 CatalogItemId::User(id) if id >= initial_next_user_item_id => Some(item),
277 _ => None,
278 })
279 .collect::<Vec<_>>();
280
281 if new_replicas.is_empty() && new_objects.is_empty() {
282 return;
283 }
284
285 let mut info_parts = Vec::new();
286
287 if !new_replicas.is_empty() {
288 let replicas = new_replicas.iter().map(|r| {
289 format!(
290 "{{replica_id: {}, replica_name: {}, cluster_id: {}}}",
291 r.replica_id, r.name, r.cluster_id
292 )
293 });
294 info_parts.push(format!("New replicas: [{}]", separated(", ", replicas)));
295 }
296
297 if !new_objects.is_empty() {
298 let objects = new_objects
299 .iter()
300 .map(|o| format!("{{object_id: {}, object_name: {}}}", o.id, o.name));
301 info_parts.push(format!("New objects: [{}]", separated(", ", objects)));
302 }
303
304 let extra_info = separated(". ", info_parts);
305
306 halt!(
307 "there have been DDL that we need to react to; rebooting in read-only mode. {}",
308 extra_info
309 )
310}
311
312pub async fn get_next_ids(
314 catalog: &mut dyn DurableCatalogState,
315) -> Result<(u64, u64), CatalogError> {
316 let snapshot = catalog.snapshot().await?;
318 let mut dry_run = catalog.transaction_from_snapshot(snapshot)?;
319 let tx = dry_run.transaction_mut();
320
321 fn next_user_id(iter: impl Iterator<Item = u64>) -> u64 {
323 iter.max().map(|id| id + 1).unwrap_or(0)
324 }
325
326 let next_user_item_id = next_user_id(tx.get_items().filter_map(|item| match item.id {
327 CatalogItemId::User(id) => Some(id),
328 _ => None,
329 }));
330
331 let next_replica_id = next_user_id(tx.get_cluster_replicas().filter_map(
332 |r| match r.replica_id {
333 ReplicaId::User(id) => Some(id),
334 ReplicaId::System(_) => None,
335 },
336 ));
337
338 Ok((next_user_item_id, next_replica_id))
339}
340
341#[cfg(test)]
342mod tests {
343 use super::*;
344 use mz_catalog::durable::{TestCatalogStateBuilder, test_bootstrap_args};
345 use mz_orchestratord::controller::materialize::generation::DeploymentStatus;
346 use mz_ore::metrics::MetricsRegistry;
347 use mz_ore::now::SYSTEM_TIME;
348 use mz_persist_client::PersistLocation;
349 use mz_persist_client::cache::PersistClientCache;
350 use mz_persist_client::cfg::PersistConfig;
351 use mz_persist_client::rpc::PubSubClientConnection;
352 use mz_repr::GlobalId;
353 use mz_sql::session::user::MZ_SYSTEM_ROLE_ID;
354
355 use crate::deployment::state::DeploymentStateHandle;
356
357 async fn setup() -> (
358 TestCatalogStateBuilder,
359 CatchupConfig,
360 DeploymentStateHandle,
361 ) {
362 let mut config = PersistConfig::new_for_tests();
363 config.build_version = BUILD_INFO.semver_version();
364 let cache = PersistClientCache::new(config, &MetricsRegistry::new(), |_, _| {
365 PubSubClientConnection::noop()
366 });
367 let persist_client = cache.open(PersistLocation::new_in_mem()).await.unwrap();
368 let environment_id = EnvironmentId::for_tests();
369 let metrics = Arc::new(Metrics::new(&MetricsRegistry::new()));
370 let builder = TestCatalogStateBuilder::new(persist_client.clone())
371 .with_organization_id(environment_id.organization_id())
372 .with_version(BUILD_INFO.semver_version())
373 .with_metrics(Arc::clone(&metrics))
374 .with_deploy_generation(0);
375 let boot_ts = SYSTEM_TIME().into();
376 let catalog = builder
377 .clone()
378 .unwrap_build()
379 .await
380 .open(boot_ts, &test_bootstrap_args())
381 .await
382 .unwrap();
383 catalog.expire().await;
384
385 let (deployment_state, handle) = DeploymentState::new();
386 let config = CatchupConfig {
387 boot_ts,
388 environment_id,
389 persist_client,
390 deploy_generation: 1,
391 deployment_state,
392 catalog_metrics: metrics,
393 caught_up_max_wait: Duration::from_secs(1),
394 ddl_check_interval: Duration::from_millis(10),
395 panic_after_timeout: false,
396 bootstrap_args: test_bootstrap_args(),
397 };
398 (builder, config, handle)
399 }
400
401 async fn wait_ready(handle: &DeploymentStateHandle) {
402 tokio::time::timeout(Duration::from_secs(10), async {
403 while handle.status() != DeploymentStatus::ReadyToPromote {
404 tokio::time::sleep(Duration::from_millis(10)).await;
405 }
406 })
407 .await
408 .unwrap();
409 }
410
411 #[mz_ore::test(tokio::test)]
412 async fn catchup_starts_after_bootstrap() {
413 let (builder, config, handle) = setup().await;
414 let mut openable = builder.with_deploy_generation(1).unwrap_build().await;
415 assert!(preflight_0dt(openable.as_mut(), 1).await.unwrap());
416 let (_trigger, receiver) = trigger::channel();
417 let (bootstrapped, bootstrapped_receiver) = oneshot::channel();
418 let metrics = Arc::clone(&config.catalog_metrics);
419 let before = metrics.transactions_started.get();
420 let caught_up_max_wait = config.caught_up_max_wait;
421 let (boot_ts, bootstrap_args) = (config.boot_ts, config.bootstrap_args.clone());
422
423 tokio::time::pause();
424 spawn_catchup(config, receiver, bootstrapped_receiver);
425 tokio::time::sleep(2 * caught_up_max_wait).await;
426 assert_eq!(metrics.transactions_started.get(), before);
427 assert_eq!(handle.status(), DeploymentStatus::Initializing);
428
429 let mut catalog = openable
430 .open_savepoint(boot_ts, &bootstrap_args)
431 .await
432 .unwrap();
433 bootstrapped
434 .send(get_next_ids(catalog.as_mut()).await.unwrap())
435 .unwrap();
436 wait_ready(&handle).await;
437 catalog.expire().await;
438 }
439
440 #[mz_ore::test(tokio::test)]
441 async fn skip_before_bootstrap() {
442 let (_, config, handle) = setup().await;
443 let (_trigger, receiver) = trigger::channel();
444 let (_bootstrapped, bootstrapped_receiver) = oneshot::channel();
445 spawn_catchup(config, receiver, bootstrapped_receiver);
446 tokio::time::timeout(Duration::from_secs(10), async {
447 while handle.try_skip_catchup().is_err() {
448 tokio::task::yield_now().await;
449 }
450 })
451 .await
452 .unwrap();
453 wait_ready(&handle).await;
454 }
455
456 #[mz_ore::test(tokio::test)]
457 async fn caught_up_before_task_starts() {
458 let (builder, mut config, handle) = setup().await;
459 let mut catalog = builder
460 .with_deploy_generation(1)
461 .unwrap_build()
462 .await
463 .open_savepoint(config.boot_ts, &config.bootstrap_args)
464 .await
465 .unwrap();
466 let (bootstrapped, bootstrapped_receiver) = oneshot::channel();
467 bootstrapped
468 .send(get_next_ids(catalog.as_mut()).await.unwrap())
469 .unwrap();
470 let (trigger, receiver) = trigger::channel();
471 drop(trigger);
472 config.caught_up_max_wait = Duration::ZERO;
473 config.panic_after_timeout = true;
474 spawn_catchup(config, receiver, bootstrapped_receiver);
475 wait_ready(&handle).await;
476 catalog.expire().await;
477 }
478
479 #[mz_ore::test(tokio::test)]
480 async fn baseline_uses_bootstrap_snapshot() {
481 let (builder, config, _) = setup().await;
482 let mut writer = builder
483 .clone()
484 .unwrap_build()
485 .await
486 .open(config.boot_ts, &config.bootstrap_args)
487 .await
488 .unwrap();
489 let mut catalog = builder
490 .with_deploy_generation(1)
491 .unwrap_build()
492 .await
493 .open_savepoint(config.boot_ts, &config.bootstrap_args)
494 .await
495 .unwrap();
496 let initial_ids = get_next_ids(catalog.as_mut()).await.unwrap();
497 assert!(!catalog.sync_to_current_updates().await.unwrap().is_empty());
498
499 writer.sync_to_current_updates().await.unwrap();
500 let mut tx = writer.transaction().await.unwrap();
501 let item_id = initial_ids.0 + 10;
502 let replica_id = initial_ids.1 + 10;
503 let schema_id = tx.get_schemas().find(|s| s.name == "public").unwrap().id;
504 tx.insert_item(
505 CatalogItemId::User(item_id),
506 20000,
507 GlobalId::User(item_id),
508 schema_id,
509 "new_item",
510 "CREATE TABLE new_item (a int)".into(),
511 MZ_SYSTEM_ROLE_ID,
512 Vec::new(),
513 Default::default(),
514 None,
515 )
516 .unwrap();
517 let replica = tx.get_cluster_replicas().next().unwrap();
518 tx.insert_cluster_replica_with_id(
519 replica.cluster_id,
520 ReplicaId::User(replica_id),
521 "new_replica",
522 replica.config,
523 replica.owner_id,
524 )
525 .unwrap();
526 let commit_ts = tx.upper();
527 let _ = tx.get_and_commit_op_updates();
528 tx.commit(commit_ts).await.unwrap();
529
530 assert_eq!(get_next_ids(catalog.as_mut()).await.unwrap(), initial_ids);
531 assert_eq!(
532 get_next_ids(writer.as_mut()).await.unwrap(),
533 (item_id + 1, replica_id + 1)
534 );
535 catalog.expire().await;
536 writer.expire().await;
537 }
538}