Skip to main content

mz_environmentd/deployment/
preflight.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//! Preflight checks for deployments.
11
12use 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
34/// Configuration for catching up and promoting a read-only deployment.
35pub 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
48/// Returns whether this deployment should boot in read-only mode.
49pub 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
71/// Starts catching up and promoting a read-only deployment.
72///
73/// An administrative skip is accepted right away and promotes without waiting
74/// for bootstrap. Otherwise, catch-up checks and the catch-up timeout start
75/// once `bootstrapped` yields the ID baseline, which must come from the
76/// savepoint used to bootstrap the adapter. The task exits if `bootstrapped`
77/// is dropped.
78pub 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            // Check for DDL changes one last time before announcing as ready to
159            // promote.
160            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        // Announce that we're ready to promote.
178        let promoted = deployment_state.set_ready_to_promote();
179        info!("announced as ready to promote; waiting for promotion");
180        promoted.await;
181
182        // Take over the catalog.
183        info!("promoted; attempting takeover");
184
185        // NOTE: There _is_ a window where DDL can happen in the old
186        // environment, between checking above, us announcing as ready to
187        // promote, and cloud giving us the go-ahead signal. Its size
188        // depends on how quickly cloud will trigger promotion once we
189        // report as ready.
190        //
191        // We could add another check here, right before cutting over, but I
192        // think this requires changes in Cloud: with this additional check,
193        // it can now happen that cloud gives us the promote signal but we
194        // then notice there were changes and restart. Could would have to
195        // notice this and give us the promote signal again, once we're
196        // ready again.
197
198        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        // Reboot as the leader.
214        halt!("fenced out old deployment; rebooting as leader")
215    });
216}
217
218/// Check if there have been any DDL that create new collections or replicas,
219/// restart in read-only mode if so, in order to pick up those new items and
220/// start hydrating them before cutting over.
221///
222/// We do this by checking whether items or replicas with IDs above the highest
223/// ones in the catalog snapshot this deployment bootstrapped from were
224/// committed.
225async 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    // `transaction` rejects unapplied catalog content. This reader has no derived catalog to
251    // update, so discard the initial update stream before opening the transaction.
252    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    // We must explicitly check the catalog for these IDs since IDs can be
262    // allocated during sequencing/planning but not yet committed to the catalog.
263    // Furthermore, these IDs might never be committed to the catalog because
264    // their sequencing has been aborted.
265    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
312/// Returns the next user item and replica IDs based on existing catalog objects.
313pub async fn get_next_ids(
314    catalog: &mut dyn DurableCatalogState,
315) -> Result<(u64, u64), CatalogError> {
316    // Preserve the pending updates that adapter bootstrap must consume.
317    let snapshot = catalog.snapshot().await?;
318    let mut dry_run = catalog.transaction_from_snapshot(snapshot)?;
319    let tx = dry_run.transaction_mut();
320
321    // Allocator counters can be ahead of committed objects due to ID pooling.
322    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}