Skip to main content

risingwave_meta/barrier/
worker.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::mem::replace;
18use std::pin::pin;
19use std::sync::Arc;
20use std::time::Duration;
21
22use anyhow::anyhow;
23use arc_swap::ArcSwap;
24use futures::{TryFutureExt, pin_mut};
25use itertools::Itertools;
26use risingwave_common::catalog::DatabaseId;
27use risingwave_common::id::JobId;
28use risingwave_common::system_param::PAUSE_ON_NEXT_BOOTSTRAP_KEY;
29use risingwave_common::system_param::reader::SystemParamsRead;
30use risingwave_meta_model::WorkerId;
31use risingwave_pb::common::WorkerNode;
32use risingwave_pb::meta::Recovery;
33use risingwave_pb::meta::subscribe_response::{Info, Operation};
34use thiserror_ext::AsReport;
35use tokio::select;
36use tokio::sync::mpsc;
37use tokio::sync::oneshot::{Receiver, Sender};
38use tokio::task::JoinHandle;
39use tonic::Status;
40use tracing::{Instrument, debug, error, info, warn};
41
42use crate::barrier::checkpoint::{CheckpointControl, CheckpointControlEvent};
43use crate::barrier::complete_task::{BarrierCompleteOutput, CompletingTask};
44use crate::barrier::context::recovery::{RenderedDatabaseRuntimeInfo, render_runtime_info};
45use crate::barrier::context::{GlobalBarrierWorkerContext, GlobalBarrierWorkerContextImpl};
46use crate::barrier::info::InflightDatabaseInfo;
47use crate::barrier::rpc::{
48    DatabaseInitialBarrierCollector, database_partial_graphs, from_partial_graph_id,
49    merge_node_rpc_errors,
50};
51use crate::barrier::schedule::{MarkReadyOptions, PeriodicBarriers};
52use crate::barrier::{
53    BarrierManagerRequest, BarrierManagerStatus, BarrierWorkerRuntimeInfoSnapshot, Command,
54    CreateStreamingJobType, RecoveryReason, RescheduleContext, UpdateDatabaseBarrierRequest,
55    schedule,
56};
57use crate::controller::scale::{materialize_actor_assignments, preview_actor_assignments};
58use crate::error::MetaErrorInner;
59use crate::hummock::HummockManagerRef;
60use crate::manager::iceberg_compaction::IcebergCompactionManagerRef;
61use crate::manager::iceberg_pk_index_sink::IcebergPkIndexSinkManager;
62use crate::manager::sink_coordination::SinkCoordinatorManager;
63use crate::manager::{
64    ActiveStreamingWorkerChange, ActiveStreamingWorkerNodes, LocalNotification, MetaSrvEnv,
65    MetadataManager,
66};
67use crate::rpc::metrics::GLOBAL_META_METRICS;
68use crate::serving::ServingVnodeMappingRef;
69use crate::stream::{
70    GlobalRefreshManagerRef, ScaleControllerRef, SourceManagerRef, build_reschedule_commands,
71    rendered_layout_matches_current,
72};
73use crate::{MetaError, MetaResult};
74
75/// [`crate::barrier::worker::GlobalBarrierWorker`] sends barriers to all registered compute nodes and
76/// collect them, with monotonic increasing epoch numbers. On compute nodes, `LocalBarrierManager`
77/// in `risingwave_stream` crate will serve these requests and dispatch them to source actors.
78///
79/// Configuration change in our system is achieved by the mutation in the barrier. Thus,
80/// [`crate::barrier::worker::GlobalBarrierWorker`] provides a set of interfaces like a state machine,
81/// accepting [`crate::barrier::command::Command`] that carries info to build `Mutation`. To keep the consistency between
82/// barrier manager and meta store, some actions like "drop materialized view" or "create mv on mv"
83/// must be done in barrier manager transactional using [`crate::barrier::command::Command`].
84pub(super) struct GlobalBarrierWorker<C> {
85    /// Enable recovery or not when failover.
86    enable_recovery: bool,
87
88    /// The queue of scheduled barriers.
89    periodic_barriers: PeriodicBarriers,
90
91    /// Whether per database failure isolation is enabled in system parameters.
92    system_enable_per_database_isolation: bool,
93
94    pub(super) context: Arc<C>,
95
96    env: MetaSrvEnv,
97
98    checkpoint_control: CheckpointControl,
99
100    /// Command that has been collected but is still completing.
101    /// The join handle of the completing future is stored.
102    completing_task: CompletingTask,
103
104    request_rx: mpsc::UnboundedReceiver<BarrierManagerRequest>,
105
106    active_streaming_nodes: ActiveStreamingWorkerNodes,
107
108    partial_graph_manager: PartialGraphManager,
109}
110
111#[cfg(test)]
112mod tests {
113    use std::collections::HashMap;
114
115    use super::*;
116    use crate::barrier::RescheduleContext;
117    use crate::barrier::notifier::Notifier;
118
119    #[tokio::test]
120    async fn test_reschedule_intent_without_workers_notifies_start_failed() {
121        let env = MetaSrvEnv::for_test().await;
122        let database_id = DatabaseId::new(1);
123        let database_info =
124            InflightDatabaseInfo::empty(database_id, env.shared_actor_infos().clone());
125        let (notifier, started_rx) = Notifier::new();
126
127        let new_barrier = schedule::NewBarrier {
128            database_id,
129            command: Some((
130                Command::RescheduleIntent {
131                    context: RescheduleContext::empty(),
132                    reschedule_plan: None,
133                },
134                notifier,
135            )),
136            span: tracing::Span::none(),
137            checkpoint: false,
138        };
139
140        let result =
141            resolve_reschedule_intent(env, HashMap::new(), Some(&database_info), new_barrier);
142
143        assert!(matches!(result, Ok(None)));
144        let started = started_rx.await.expect("started notifier dropped");
145        assert!(started.is_err());
146    }
147}
148
149impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
150    pub(super) async fn new_inner(
151        env: MetaSrvEnv,
152        request_rx: mpsc::UnboundedReceiver<BarrierManagerRequest>,
153        context: Arc<C>,
154    ) -> Self {
155        let enable_recovery = env.opts.enable_recovery;
156
157        let active_streaming_nodes = ActiveStreamingWorkerNodes::uninitialized();
158
159        let partial_graph_manager = PartialGraphManager::uninitialized(env.clone());
160
161        let reader = env.system_params_reader().await;
162        let system_enable_per_database_isolation = reader.per_database_isolation();
163        // Load config will be performed in bootstrap phase.
164        let periodic_barriers = PeriodicBarriers::default();
165
166        let checkpoint_control = CheckpointControl::new(env.clone());
167        Self {
168            enable_recovery,
169            periodic_barriers,
170            system_enable_per_database_isolation,
171            context,
172            env,
173            checkpoint_control,
174            completing_task: CompletingTask::None,
175            request_rx,
176            active_streaming_nodes,
177            partial_graph_manager,
178        }
179    }
180}
181
182fn resolve_reschedule_intent(
183    env: MetaSrvEnv,
184    worker_nodes: HashMap<WorkerId, WorkerNode>,
185    database_info: Option<&InflightDatabaseInfo>,
186    mut new_barrier: schedule::NewBarrier,
187) -> MetaResult<Option<schedule::NewBarrier>> {
188    let Some((command, notifier)) = new_barrier.command.take() else {
189        return Ok(Some(new_barrier));
190    };
191
192    match command {
193        Command::RescheduleIntent {
194            context,
195            reschedule_plan,
196        } => {
197            if let Some(reschedule_plan) = reschedule_plan {
198                new_barrier.command = Some((
199                    Command::RescheduleIntent {
200                        context,
201                        reschedule_plan: Some(reschedule_plan),
202                    },
203                    notifier,
204                ));
205                return Ok(Some(new_barrier));
206            }
207            let span = tracing::info_span!(
208                "resolve_reschedule_intent",
209                database_id = %new_barrier.database_id
210            );
211            let reschedule_plan = {
212                let _guard = span.enter();
213                build_reschedule_from_context(
214                    &env,
215                    worker_nodes,
216                    new_barrier.database_id,
217                    context,
218                    database_info.ok_or_else(|| {
219                        anyhow!(
220                            "database {} not found when resolving reschedule intent",
221                            new_barrier.database_id
222                        )
223                    })?,
224                )
225            };
226            match reschedule_plan {
227                Ok(Some(reschedule_plan)) => {
228                    new_barrier.command = Some((
229                        Command::RescheduleIntent {
230                            context: RescheduleContext::empty(),
231                            reschedule_plan: Some(reschedule_plan),
232                        },
233                        notifier,
234                    ));
235                    Ok(Some(new_barrier))
236                }
237                Ok(None) => {
238                    // No-op intent: notify to unblock callers even though no barrier is injected.
239                    notifier.start().started();
240                    Ok(None)
241                }
242                Err(err) => {
243                    notifier.notify_start_failed(err);
244                    Ok(None)
245                }
246            }
247        }
248        _ => {
249            new_barrier.command = Some((command, notifier));
250            Ok(Some(new_barrier))
251        }
252    }
253}
254
255fn build_reschedule_from_context(
256    env: &MetaSrvEnv,
257    worker_nodes: HashMap<WorkerId, WorkerNode>,
258    database_id: DatabaseId,
259    context: RescheduleContext,
260    database_info: &InflightDatabaseInfo,
261) -> MetaResult<Option<crate::barrier::ReschedulePlan>> {
262    if worker_nodes.is_empty() {
263        return Err(anyhow!("no active streaming workers for reschedule").into());
264    }
265
266    if context.is_empty() {
267        return Ok(None);
268    }
269
270    // Barrier worker resolves this intent against a stable in-flight snapshot.
271    // Reuse the same fragment view for preview comparison and command building.
272    let all_prev_fragments = database_info
273        .fragment_infos()
274        .map(|fragment| (fragment.fragment_id, fragment))
275        .collect();
276
277    let previewed = preview_actor_assignments(&worker_nodes, &context.loaded)?;
278
279    if rendered_layout_matches_current(&previewed.fragments, &all_prev_fragments)? {
280        return Ok(None);
281    }
282
283    let actor_id_counter = env.actor_id_generator();
284    // Materialization only replaces preview actor ids with real ids. Worker
285    // placement, vnode ownership, and split assignment remain unchanged.
286    let rendered = materialize_actor_assignments(actor_id_counter, previewed);
287    let mut commands = build_reschedule_commands(rendered.fragments, context, all_prev_fragments)?;
288    Ok(commands.remove(&database_id))
289}
290
291impl GlobalBarrierWorker<GlobalBarrierWorkerContextImpl> {
292    /// Create a new [`crate::barrier::worker::GlobalBarrierWorker`].
293    #[expect(clippy::too_many_arguments)]
294    pub async fn new(
295        scheduled_barriers: schedule::ScheduledBarriers,
296        env: MetaSrvEnv,
297        metadata_manager: MetadataManager,
298        hummock_manager: HummockManagerRef,
299        serving_vnode_mapping: ServingVnodeMappingRef,
300        source_manager: SourceManagerRef,
301        sink_manager: SinkCoordinatorManager,
302        iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
303        iceberg_compaction_manager: IcebergCompactionManagerRef,
304        scale_controller: ScaleControllerRef,
305        request_rx: mpsc::UnboundedReceiver<BarrierManagerRequest>,
306        barrier_scheduler: schedule::BarrierScheduler,
307        refresh_manager: GlobalRefreshManagerRef,
308    ) -> Self {
309        let status = Arc::new(ArcSwap::new(Arc::new(BarrierManagerStatus::Starting)));
310
311        let context = Arc::new(GlobalBarrierWorkerContextImpl::new(
312            scheduled_barriers,
313            status,
314            metadata_manager,
315            hummock_manager,
316            serving_vnode_mapping,
317            source_manager,
318            scale_controller,
319            env.clone(),
320            barrier_scheduler,
321            refresh_manager,
322            sink_manager,
323            iceberg_pk_index_sink_manager,
324            iceberg_compaction_manager,
325        ));
326
327        Self::new_inner(env, request_rx, context).await
328    }
329
330    pub fn start(self) -> (JoinHandle<()>, Sender<()>) {
331        let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
332        let fut = (self.env.await_tree_reg())
333            .register_derived_root("Global Barrier Worker")
334            .instrument(self.run(shutdown_rx));
335        let join_handle = tokio::spawn(fut);
336
337        (join_handle, shutdown_tx)
338    }
339
340    /// Check whether we should pause on bootstrap from the system parameter and reset it.
341    async fn take_pause_on_bootstrap(&mut self) -> MetaResult<bool> {
342        let paused = self
343            .env
344            .system_params_reader()
345            .await
346            .pause_on_next_bootstrap()
347            || self.env.opts.pause_on_next_bootstrap_offline;
348
349        if paused {
350            warn!(
351                "The cluster will bootstrap with all data sources paused as specified by the system parameter `{}`. \
352                 It will now be reset to `false`. \
353                 To resume the data sources, either restart the cluster again or use `risectl meta resume`.",
354                PAUSE_ON_NEXT_BOOTSTRAP_KEY
355            );
356            self.env
357                .system_params_manager_impl_ref()
358                .set_param(PAUSE_ON_NEXT_BOOTSTRAP_KEY, Some("false".to_owned()))
359                .await?;
360        }
361        Ok(paused)
362    }
363
364    /// Start an infinite loop to take scheduled barriers and send them.
365    async fn run(mut self, shutdown_rx: Receiver<()>) {
366        tracing::info!(
367            "Starting barrier manager with: enable_recovery={}, in_flight_barrier_nums={}",
368            self.enable_recovery,
369            self.checkpoint_control.in_flight_barrier_nums,
370        );
371
372        if !self.enable_recovery {
373            let job_exist = self
374                .context
375                .metadata_manager
376                .catalog_controller
377                .has_any_streaming_jobs()
378                .await
379                .unwrap();
380            if job_exist {
381                panic!(
382                    "Some streaming jobs already exist in meta, please start with recovery enabled \
383                or clean up the metadata using `./risedev clean-data`"
384                );
385            }
386        }
387
388        {
389            // Bootstrap recovery. Here we simply trigger a recovery process to achieve the
390            // consistency.
391            // Even if there's no actor to recover, we still go through the recovery process to
392            // inject the first `Initial` barrier.
393            let span = tracing::info_span!("bootstrap_recovery");
394            crate::telemetry::report_event(
395                risingwave_pb::telemetry::TelemetryEventStage::Recovery,
396                "normal_recovery",
397                0,
398                None,
399                None,
400                None,
401            );
402
403            let paused = self.take_pause_on_bootstrap().await.unwrap_or(false);
404
405            // Keep the bootstrap recovery future boxed so the outer barrier worker future
406            // stays below clippy's large-futures threshold.
407            Box::pin(
408                self.recovery(paused, RecoveryReason::Bootstrap)
409                    .instrument(span),
410            )
411            .await;
412        }
413
414        Box::pin(self.run_inner(shutdown_rx)).await
415    }
416}
417
418impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
419    fn enable_per_database_isolation(&self) -> bool {
420        self.system_enable_per_database_isolation && {
421            if let Err(e) =
422                risingwave_common::license::Feature::DatabaseFailureIsolation.check_available()
423            {
424                warn!(error = %e.as_report(), "DatabaseFailureIsolation disabled by license");
425                false
426            } else {
427                true
428            }
429        }
430    }
431
432    async fn resolve_since_timestamp_snapshot_backfill(
433        &mut self,
434        new_barrier: &mut schedule::NewBarrier,
435    ) -> MetaResult<bool> {
436        let Some((
437            Command::CreateStreamingJob {
438                job_type:
439                    CreateStreamingJobType::SnapshotBackfill {
440                        snapshot_backfill_info,
441                        since_epoch: Some(since_epoch),
442                    },
443                ..
444            },
445            _,
446        )) = &mut new_barrier.command
447        else {
448            return Ok(true);
449        };
450        let since_timestamp_epoch = since_epoch.provided_since_epoch;
451
452        // Complete any inflight command first so the committed upstream epoch and
453        // table changelog view used for since_timestamp resolution cannot be stale.
454        match self.completing_task.wait_completing_task().await {
455            Ok(Some(output)) => self
456                .checkpoint_control
457                .ack_completed(&mut self.partial_graph_manager, output),
458            Ok(None) => {}
459            Err(err) => {
460                error!(
461                    err = %err.as_report(),
462                    "failed to wait completing task before resolving since_timestamp"
463                );
464                return Err(err);
465            }
466        }
467
468        match self
469            .context
470            .resolve_log_store_epoch(
471                snapshot_backfill_info
472                    .upstream_mv_table_id_to_backfill_epoch
473                    .keys()
474                    .copied(),
475                since_timestamp_epoch,
476            )
477            .await
478        {
479            Ok(upstream_log_epochs) => {
480                since_epoch.resolved = Some(upstream_log_epochs);
481                Ok(true)
482            }
483            Err(err) => {
484                error!(
485                    err = %err.as_report(),
486                    "failed to resolve log store epoch for since_timestamp"
487                );
488                let (_, notifier) = new_barrier
489                    .command
490                    .take()
491                    .expect("matched command should still exist");
492                notifier.notify_start_failed(err);
493                Ok(false)
494            }
495        }
496    }
497
498    pub(super) async fn run_inner(mut self, mut shutdown_rx: Receiver<()>) {
499        let (local_notification_tx, mut local_notification_rx) =
500            tokio::sync::mpsc::unbounded_channel();
501        self.env
502            .notification_manager()
503            .insert_local_sender(local_notification_tx);
504
505        // Start the event loop.
506        loop {
507            tokio::select! {
508                biased;
509
510                // Shutdown
511                _ = &mut shutdown_rx => {
512                    tracing::info!("Barrier manager is stopped");
513                    break;
514                }
515
516                request = self.request_rx.recv() => {
517                    if let Some(request) = request {
518                        match request {
519                            BarrierManagerRequest::GetBackfillProgress(result_tx) => {
520                                let progress = self.checkpoint_control.gen_backfill_progress();
521                                if result_tx.send(Ok(progress)).is_err() {
522                                    error!("failed to send get ddl progress");
523                                }
524                            }
525                            BarrierManagerRequest::GetFragmentBackfillProgress(result_tx) => {
526                                let progress =
527                                    self.checkpoint_control.gen_fragment_backfill_progress();
528                                if result_tx.send(Ok(progress)).is_err() {
529                                    error!("failed to send get fragment backfill progress");
530                                }
531                            }
532                            BarrierManagerRequest::GetCdcProgress(result_tx) => {
533                                let progress = self.checkpoint_control.gen_cdc_progress();
534                                if result_tx.send(Ok(progress)).is_err() {
535                                    error!("failed to send get ddl progress");
536                                }
537                            }
538                            // Handle adhoc recovery triggered by user.
539                            BarrierManagerRequest::AdhocRecovery(sender) => {
540                                self.adhoc_recovery().await;
541                                if sender.send(()).is_err() {
542                                    warn!("failed to notify finish of adhoc recovery");
543                                }
544                            }
545                            BarrierManagerRequest::UpdateDatabaseBarrier( UpdateDatabaseBarrierRequest {
546                                database_id,
547                                barrier_interval_ms,
548                                checkpoint_frequency,
549                                sender,
550                            }) => {
551                                self.periodic_barriers
552                                    .update_database_barrier(
553                                        database_id,
554                                        barrier_interval_ms,
555                                        checkpoint_frequency,
556                                    );
557                                if sender.send(()).is_err() {
558                                    warn!("failed to notify finish of update database barrier");
559                                }
560                            }
561                        }
562                    } else {
563                        tracing::info!("end of request stream. meta node may be shutting down. Stop global barrier manager");
564                        return;
565                    }
566                }
567
568                changed_worker = self.active_streaming_nodes.changed() => {
569                    #[cfg(debug_assertions)]
570                    {
571                        self.active_streaming_nodes.validate_change().await;
572                    }
573
574                    info!(?changed_worker, "worker changed");
575
576                    match changed_worker {
577                        ActiveStreamingWorkerChange::Add(node)
578                        | ActiveStreamingWorkerChange::Update(node) => {
579                            self.partial_graph_manager
580                                .add_worker(node, self.context.clone())
581                                .await;
582                        }
583                        ActiveStreamingWorkerChange::Remove(node) => {
584                            self.partial_graph_manager.remove_worker(node);
585                        }
586                    }
587                }
588
589                notification = local_notification_rx.recv() => {
590                    let notification = notification.unwrap();
591                    if let LocalNotification::SystemParamsChange(p) = notification {
592                        {
593                            self.periodic_barriers.set_sys_barrier_interval(Duration::from_millis(p.barrier_interval_ms() as u64));
594                            self.periodic_barriers
595                                .set_sys_checkpoint_frequency(p.checkpoint_frequency());
596                            self.system_enable_per_database_isolation = p.per_database_isolation();
597                        }
598                    }
599                }
600                complete_result = self
601                    .completing_task
602                    .next_completed_barrier(
603                        &mut self.periodic_barriers,
604                        &mut self.checkpoint_control,
605                        &mut self.partial_graph_manager,
606                        &self.context,
607                        &self.env,
608                ) => {
609                    match complete_result {
610                        Ok(output) => {
611                            self.checkpoint_control.ack_completed(&mut self.partial_graph_manager, output);
612                        }
613                        Err(e) => {
614                            self.failure_recovery(e).await;
615                        }
616                    }
617                },
618                event = self.checkpoint_control.next_event() => {
619                    let result: MetaResult<()> = try {
620                        match event {
621                            CheckpointControlEvent::EnteringInitializing(entering_initializing) => {
622                                let database_id = entering_initializing.database_id();
623                                let error = merge_node_rpc_errors(&format!("database {} reset", database_id), entering_initializing.action.0.iter().filter_map(|(worker_id, resp)| {
624                                    resp.root_err.as_ref().map(|root_err| {
625                                        (*worker_id, ScoredError {
626                                            error: Status::internal(&root_err.err_msg),
627                                            score: Score(root_err.score)
628                                        })
629                                    })
630                                }));
631                                Self::report_collect_failure(&self.env, &error);
632                                self.context.notify_creating_job_failed(Some(database_id), format!("{}", error.as_report())).await;
633                                let result: MetaResult<_> = try {
634                                    let runtime_info = self.context.reload_database_runtime_info(database_id).await.inspect_err(|err| {
635                                        warn!(%database_id, err = %err.as_report(), "reload runtime info failed");
636                                    })?;
637                                    let rendered_info = render_runtime_info(
638                                        self.env.actor_id_generator(),
639                                        &self.active_streaming_nodes,
640                                        &runtime_info.recovery_context,
641                                        database_id,
642                                    )
643                                    .inspect_err(|err: &MetaError| {
644                                        warn!(%database_id, err = %err.as_report(), "render runtime info failed");
645                                    })?;
646                                    if let Some(rendered_info) = rendered_info {
647                                        BarrierWorkerRuntimeInfoSnapshot::validate_database_info(
648                                            database_id,
649                                            &rendered_info.job_infos,
650                                            &self.active_streaming_nodes,
651                                            &rendered_info.stream_actors,
652                                            &runtime_info.state_table_committed_epochs,
653                                        )
654                                        .inspect_err(|err| {
655                                            warn!(%database_id, err = ?err.as_report(), "database runtime info failed validation");
656                                        })?;
657                                        Some((runtime_info, rendered_info))
658                                    } else {
659                                        None
660                                    }
661                                };
662                                match result {
663                                    Ok(Some((runtime_info, rendered_info))) => {
664                                        entering_initializing.enter(
665                                            runtime_info,
666                                            rendered_info,
667                                            &mut self.partial_graph_manager,
668                                        );
669                                    }
670                                    Ok(None) => {
671                                        info!(%database_id, "database removed after reloading empty runtime info");
672                                        entering_initializing.remove();
673                                        self.context
674                                            .refresh_table_refill_runtime_state_after_recovery()
675                                            .await?;
676                                        // Mark ready only after the refill runtime state is refreshed.
677                                        self.context.mark_ready(MarkReadyOptions::Database(database_id));
678                                    }
679                                    Err(e) => {
680                                        entering_initializing.fail_reload_runtime_info(e);
681                                    }
682                                }
683                            }
684                            CheckpointControlEvent::EnteringRunning(entering_running) => {
685                                let database_id = entering_running.database_id();
686                                entering_running.enter();
687                                self.context
688                                    .refresh_table_refill_runtime_state_after_recovery()
689                                    .await?;
690                                self.context
691                                    .mark_ready(MarkReadyOptions::Database(database_id));
692                            }
693                            CheckpointControlEvent::BatchRefreshTrigger { database_id, job_id } => {
694                                self.handle_batch_refresh_trigger(database_id, job_id).await?;
695                            }
696                        }
697                    };
698                    if let Err(e) = result {
699                        self.failure_recovery(e).await;
700                    }
701                }
702                event = self.partial_graph_manager.next_event(&self.context) => {
703                    let result: MetaResult<()> = try {
704                        match event {
705                            PartialGraphManagerEvent::Worker(_worker_id, WorkerEvent::WorkerConnected) => {
706                                // no handling on new worker connected event yet
707                            }
708                            PartialGraphManagerEvent::Worker(worker_id, WorkerEvent::WorkerError { err, affected_partial_graphs }) => {
709                                let failed_databases = self
710                                    .checkpoint_control
711                                    .databases_failed_at_worker_err(worker_id)
712                                    .chain(
713                                        affected_partial_graphs
714                                        .into_iter()
715                                        .map(|partial_graph_id| {
716                                            let (database_id, _) = from_partial_graph_id(partial_graph_id);
717                                            database_id
718                                        })
719                                    )
720                                    .collect::<HashSet<_>>();
721                                if !failed_databases.is_empty() {
722                                    if !self.enable_recovery {
723                                        panic!("control stream to worker {} failed but recovery not enabled: {}", worker_id, err.as_report());
724                                    }
725                                    if !self.enable_per_database_isolation() {
726                                        Err(err.clone())?;
727                                    }
728                                    Self::report_collect_failure(&self.env, &err);
729                                    for database_id in failed_databases {
730                                        if let Some(entering_recovery) = self.checkpoint_control.on_report_failure(database_id, &mut self.partial_graph_manager) {
731                                            warn!(%worker_id, %database_id, "database entering recovery on node failure");
732                                            self.context.abort_and_mark_blocked(Some(database_id), RecoveryReason::Failover(anyhow!("reset database: {}", database_id).into()));
733                                            self.context.notify_creating_job_failed(Some(database_id), format!("database {} reset due to node {} failure: {}", database_id, worker_id, err.as_report())).await;
734                                            // TODO: add log on blocking time
735                                            let output = self.completing_task.wait_completing_task().await?;
736                                            entering_recovery.enter(output, &mut self.partial_graph_manager);
737                                        }
738                                    }
739                                }  else {
740                                    warn!(%worker_id, "no barrier to collect from worker, ignore err");
741                                }
742                                continue;
743                            }
744                            PartialGraphManagerEvent::PartialGraph(partial_graph_id, event) => {
745                                let (database_id, _creating_job_id) = from_partial_graph_id(partial_graph_id);
746                                match event {
747                                    PartialGraphEvent::BarrierCollected(collected_barrier) => {
748                                        self.checkpoint_control.barrier_collected(partial_graph_id, collected_barrier, &mut self.periodic_barriers)?;
749                                    }
750                                    PartialGraphEvent::Error(worker_id) => {
751                                        if !self.enable_recovery {
752                                            panic!("database {database_id} failure reported from {worker_id} but recovery not enabled")
753                                        }
754                                        if !self.enable_per_database_isolation() {
755                                                Err(MetaError::from(anyhow!("database {database_id} report failure from {worker_id}")))?;
756                                            }
757                                        if let Some(entering_recovery) = self.checkpoint_control.on_report_failure(database_id, &mut self.partial_graph_manager) {
758                                            warn!(%database_id, "database entering recovery");
759                                            self.context.abort_and_mark_blocked(Some(database_id), RecoveryReason::Failover(anyhow!("reset database: {}", database_id).into()));
760                                            // TODO: add log on blocking time
761                                            let output = self.completing_task.wait_completing_task().await?;
762                                            entering_recovery.enter(output, &mut self.partial_graph_manager);
763                                        }
764                                    }
765                                    PartialGraphEvent::Reset(reset_resps) => {
766                                        self.checkpoint_control.on_partial_graph_reset(partial_graph_id, reset_resps);
767                                    }
768                                    PartialGraphEvent::Initialized => {
769                                        self.checkpoint_control.on_partial_graph_initialized(
770                                            partial_graph_id,
771                                            &mut self.partial_graph_manager,
772                                        )?;
773                                    }
774                                }
775                            }
776                        };
777                    };
778                    if let Err(e) = result {
779                        self.failure_recovery(e).await;
780                    }
781                }
782                new_barrier = self.periodic_barriers.next_barrier(&*self.context) => {
783                    let database_id = new_barrier.database_id;
784                    let mut new_barrier = if matches!(
785                        new_barrier.command,
786                        Some((Command::RescheduleIntent { .. }, _))
787                    ) {
788                        let env = self.env.clone();
789                        let worker_nodes = self
790                            .active_streaming_nodes
791                            .current()
792                            .iter()
793                            .map(|(worker_id, worker)| (*worker_id, worker.clone()))
794                            .collect();
795                        let database_info = self.checkpoint_control.database_info(database_id);
796                        match resolve_reschedule_intent(
797                            env,
798                            worker_nodes,
799                            database_info,
800                            new_barrier,
801                        ) {
802                            Ok(Some(new_barrier)) => new_barrier,
803                            Ok(None) => continue,
804                            Err(err) => {
805                                self.failure_recovery(err).await;
806                                continue;
807                            }
808                        }
809                    } else {
810                        new_barrier
811                    };
812                    match self
813                        .resolve_since_timestamp_snapshot_backfill(&mut new_barrier)
814                        .await
815                    {
816                        Ok(true) => {}
817                        Ok(false) => continue,
818                        Err(err) => {
819                            self.failure_recovery(err).await;
820                            continue;
821                        }
822                    }
823                    if let Err(e) = self.checkpoint_control.handle_new_barrier(
824                        new_barrier,
825                        &mut self.partial_graph_manager,
826                        self.active_streaming_nodes.current()
827                    ) {
828                        if !self.enable_recovery {
829                            panic!(
830                                "failed to inject barrier to some databases but recovery not enabled: {:?}", (
831                                    database_id,
832                                    e.as_report()
833                                )
834                            );
835                        }
836                        let result: MetaResult<_> = try {
837                            if !self.enable_per_database_isolation() {
838                                let err = anyhow!("failed to inject barrier to databases: {:?}", (database_id, e.as_report()));
839                                Err(MetaError::from(err))?;
840                            } else if let Some(entering_recovery) = self.checkpoint_control.on_report_failure(database_id, &mut self.partial_graph_manager) {
841                                warn!(%database_id, e = %e.as_report(),"database entering recovery on inject failure");
842                                self.context.abort_and_mark_blocked(Some(database_id), RecoveryReason::Failover(anyhow!(e).context("inject barrier failure").into()));
843                                // TODO: add log on blocking time
844                                let output = self.completing_task.wait_completing_task().await?;
845                                entering_recovery.enter(output, &mut self.partial_graph_manager);
846                            }
847                        };
848                        if let Err(e) = result {
849                            self.failure_recovery(e).await;
850                        }
851                    }
852                }
853            }
854        }
855    }
856}
857
858impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
859    /// We need to make sure there are no changes when doing recovery
860    pub async fn clear_on_err(&mut self, err: &MetaError) {
861        // join spawned completing command to finish no matter it succeeds or not.
862        match replace(&mut self.completing_task, CompletingTask::None) {
863            CompletingTask::None | CompletingTask::Err(_) => {}
864            CompletingTask::Completing {
865                epochs_to_ack,
866                join_handle,
867                ..
868            } => {
869                info!("waiting for completing command to finish in recovery");
870                match join_handle.await {
871                    Err(e) => {
872                        warn!(err = %e.as_report(), "failed to join completing task");
873                    }
874                    Ok(Err(e)) => {
875                        warn!(
876                            err = %e.as_report(),
877                            "failed to complete barrier during clear"
878                        );
879                    }
880                    Ok(Ok(hummock_version_stats)) => {
881                        self.checkpoint_control.ack_completed(
882                            &mut self.partial_graph_manager,
883                            BarrierCompleteOutput {
884                                epochs_to_ack,
885                                hummock_version_stats,
886                            },
887                        );
888                    }
889                }
890            }
891        };
892        self.partial_graph_manager.notify_all_err(err);
893    }
894}
895
896impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
897    /// Handle a batch refresh trigger: load metadata + log epochs, then start a logstore
898    /// consumption run for the given batch refresh job.
899    async fn handle_batch_refresh_trigger(
900        &mut self,
901        database_id: DatabaseId,
902        job_id: JobId,
903    ) -> MetaResult<()> {
904        // 1. Get the last committed epoch for this job (read-only).
905        let last_committed_epoch = self
906            .checkpoint_control
907            .get_batch_refresh_trigger_info(database_id, job_id);
908
909        // 2. Load context metadata + resolve log epochs asynchronously.
910        let context = self
911            .context
912            .load_batch_refresh_trigger_context(job_id, database_id, last_committed_epoch)
913            .await?;
914
915        // 3. Start the refresh run.
916        let started = self.checkpoint_control.start_batch_refresh_run(
917            database_id,
918            job_id,
919            &context,
920            self.active_streaming_nodes.current(),
921            self.env.actor_id_generator(),
922            &mut self.partial_graph_manager,
923        )?;
924
925        // 5. Update shared_actor_infos with the new fragment infos.
926        if started {
927            self.checkpoint_control
928                .apply_batch_refresh_fragment_infos(database_id, job_id);
929        }
930
931        Ok(())
932    }
933}
934
935impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
936    /// Set barrier manager status.
937    async fn failure_recovery(&mut self, err: MetaError) {
938        self.clear_on_err(&err).await;
939
940        if self.enable_recovery {
941            let span = tracing::info_span!(
942                "failure_recovery",
943                error = %err.as_report(),
944            );
945
946            crate::telemetry::report_event(
947                risingwave_pb::telemetry::TelemetryEventStage::Recovery,
948                "failure_recovery",
949                0,
950                None,
951                None,
952                None,
953            );
954
955            let reason = RecoveryReason::Failover(err);
956
957            // No need to clean dirty tables for barrier recovery,
958            // The foreground stream job should cleanup their own tables.
959            self.recovery(false, reason).instrument(span).await;
960        } else {
961            panic!(
962                "a streaming error occurred while recovery is disabled, aborting: {:?}",
963                err.as_report()
964            );
965        }
966    }
967
968    async fn adhoc_recovery(&mut self) {
969        let err = MetaErrorInner::AdhocRecovery.into();
970        self.clear_on_err(&err).await;
971
972        let span = tracing::info_span!(
973            "adhoc_recovery",
974            error = %err.as_report(),
975        );
976
977        crate::telemetry::report_event(
978            risingwave_pb::telemetry::TelemetryEventStage::Recovery,
979            "adhoc_recovery",
980            0,
981            None,
982            None,
983            None,
984        );
985
986        // No need to clean dirty tables for barrier recovery,
987        // The foreground stream job should cleanup their own tables.
988        self.recovery(false, RecoveryReason::Adhoc)
989            .instrument(span)
990            .await;
991    }
992}
993
994impl<C> GlobalBarrierWorker<C> {
995    /// Send barrier-complete-rpc and wait for responses from all CNs
996    pub(super) fn report_collect_failure(env: &MetaSrvEnv, error: &MetaError) {
997        // Record failure in event log.
998        use risingwave_pb::meta::event_log;
999        let event = event_log::EventCollectBarrierFail {
1000            error: error.to_report_string(),
1001        };
1002        env.event_log_manager_ref()
1003            .add_event_logs(vec![event_log::Event::CollectBarrierFail(event)]);
1004    }
1005}
1006
1007mod retry_strategy {
1008    use std::time::Duration;
1009
1010    use risingwave_common::util::retry::exponential_backoff;
1011    use tokio_retry::strategy::jitter;
1012
1013    // Retry base interval in milliseconds.
1014    const RECOVERY_RETRY_BASE_INTERVAL: u64 = 20;
1015    // Retry max interval.
1016    const RECOVERY_RETRY_MAX_INTERVAL: Duration = Duration::from_secs(5);
1017
1018    // MrCroxx: Use concrete type here to prevent unsolved compiler issue.
1019    // Feel free to replace the concrete type with TAIT after fixed.
1020
1021    // mod retry_backoff_future {
1022    //     use std::future::Future;
1023    //     use std::time::Duration;
1024    //
1025    //     use tokio::time::sleep;
1026    //
1027    //     pub(crate) type RetryBackoffFuture = impl Future<Output = ()> + Unpin + Send + 'static;
1028    //
1029    //     #[define_opaque(RetryBackoffFuture)]
1030    //     pub(super) fn get_retry_backoff_future(duration: Duration) -> RetryBackoffFuture {
1031    //         Box::pin(sleep(duration))
1032    //     }
1033    // }
1034    // pub(crate) use retry_backoff_future::*;
1035
1036    pub(crate) type RetryBackoffFuture = std::pin::Pin<Box<tokio::time::Sleep>>;
1037
1038    pub(crate) fn get_retry_backoff_future(duration: Duration) -> RetryBackoffFuture {
1039        Box::pin(tokio::time::sleep(duration))
1040    }
1041
1042    pub(crate) type RetryBackoffStrategy =
1043        impl Iterator<Item = RetryBackoffFuture> + Send + 'static;
1044
1045    /// Initialize a retry strategy for operation in recovery.
1046    #[inline(always)]
1047    pub(crate) fn get_retry_strategy() -> impl Iterator<Item = Duration> + Send + 'static {
1048        exponential_backoff(
1049            Duration::from_millis(RECOVERY_RETRY_BASE_INTERVAL),
1050            RECOVERY_RETRY_BASE_INTERVAL,
1051            RECOVERY_RETRY_MAX_INTERVAL,
1052        )
1053        .map(jitter)
1054    }
1055
1056    #[define_opaque(RetryBackoffStrategy)]
1057    pub(crate) fn get_retry_backoff_strategy() -> RetryBackoffStrategy {
1058        get_retry_strategy().map(get_retry_backoff_future)
1059    }
1060}
1061
1062pub(crate) use retry_strategy::*;
1063use risingwave_common::error::tonic::extra::{Score, ScoredError};
1064use risingwave_pb::meta::event_log::{Event, EventRecovery};
1065
1066use crate::barrier::partial_graph::{
1067    PartialGraphEvent, PartialGraphManager, PartialGraphManagerEvent, WorkerEvent,
1068};
1069
1070impl<C: GlobalBarrierWorkerContext> GlobalBarrierWorker<C> {
1071    /// Recovery the whole cluster from the latest epoch.
1072    ///
1073    /// If `paused_reason` is `Some`, all data sources (including connectors and DMLs) will be
1074    /// immediately paused after recovery, until the user manually resume them either by restarting
1075    /// the cluster or `risectl` command. Used for debugging purpose.
1076    ///
1077    /// Returns the new state of the barrier manager after recovery.
1078    pub async fn recovery(&mut self, is_paused: bool, recovery_reason: RecoveryReason) {
1079        // Clear all control streams to release resources (connections to compute nodes) first.
1080        self.partial_graph_manager.clear_worker();
1081
1082        let reason_str = match &recovery_reason {
1083            RecoveryReason::Bootstrap => "bootstrap".to_owned(),
1084            RecoveryReason::Failover(err) => {
1085                format!("failed over: {}", err.as_report())
1086            }
1087            RecoveryReason::Adhoc => "adhoc recovery".to_owned(),
1088        };
1089        self.context.abort_and_mark_blocked(None, recovery_reason);
1090
1091        self.recovery_inner(is_paused, reason_str).await;
1092        self.context.mark_ready(MarkReadyOptions::Global {
1093            blocked_databases: self.checkpoint_control.recovering_databases().collect(),
1094        });
1095    }
1096
1097    #[await_tree::instrument("recovery({recovery_reason})")]
1098    async fn recovery_inner(&mut self, is_paused: bool, recovery_reason: String) {
1099        let event_log_manager_ref = self.env.event_log_manager_ref();
1100
1101        tracing::info!("recovery start!");
1102        event_log_manager_ref.add_event_logs(vec![Event::Recovery(
1103            EventRecovery::global_recovery_start(recovery_reason.clone()),
1104        )]);
1105
1106        let retry_strategy = get_retry_strategy();
1107
1108        // We take retry into consideration because this is the latency user sees for a cluster to
1109        // get recovered.
1110        let recovery_timer = GLOBAL_META_METRICS
1111            .recovery_latency
1112            .with_label_values(&["global"])
1113            .start_timer();
1114
1115        let enable_per_database_isolation = self.enable_per_database_isolation();
1116
1117        let recovery_future = tokio_retry::Retry::spawn(retry_strategy, || async {
1118            self.env.stream_client_pool().invalidate_all();
1119            // Notify in every recovery retry so foreground creation handlers can either cancel an
1120            // early-failing job or re-register their finish notifier after recovery.
1121            self.context
1122                .notify_creating_job_failed(None, recovery_reason.clone())
1123                .await;
1124
1125            let runtime_info_snapshot = self
1126                .context
1127                .reload_runtime_info()
1128                .await?;
1129            let BarrierWorkerRuntimeInfoSnapshot {
1130                active_streaming_nodes,
1131                recovery_context,
1132                mut state_table_committed_epochs,
1133                mut state_table_log_epochs,
1134                mut mv_depended_subscriptions,
1135                mut creating_jobs,
1136                hummock_version_stats,
1137                database_infos,
1138                mut cdc_table_snapshot_splits,
1139            } = runtime_info_snapshot;
1140
1141            let mut partial_graph_manager = PartialGraphManager::recover(
1142                    self.env.clone(),
1143                    active_streaming_nodes.current(),
1144                    self.context.clone(),
1145                )
1146                .await;
1147            {
1148                let mut empty_databases = HashSet::new();
1149                let mut collected_databases = HashMap::new();
1150                let mut collecting_databases = HashMap::new();
1151                let mut failed_databases = HashMap::new();
1152                for &database_id in recovery_context.fragment_context.database_map.keys() {
1153                    let mut recoverer = partial_graph_manager.start_recover();
1154                    let result: MetaResult<_> = try {
1155                        let Some(rendered_info) = render_runtime_info(
1156                            self.env.actor_id_generator(),
1157                            &active_streaming_nodes,
1158                            &recovery_context,
1159                            database_id,
1160                        )
1161                            .inspect_err(|err: &MetaError| {
1162                                warn!(%database_id, err = %err.as_report(), "render runtime info failed");
1163                            })? else {
1164                            empty_databases.insert(database_id);
1165                            continue;
1166                        };
1167                        BarrierWorkerRuntimeInfoSnapshot::validate_database_info(
1168                            database_id,
1169                            &rendered_info.job_infos,
1170                            &active_streaming_nodes,
1171                            &rendered_info.stream_actors,
1172                            &state_table_committed_epochs,
1173                        )
1174                        .inspect_err(|err| {
1175                            warn!(%database_id, err = %err.as_report(), "rendered runtime info failed validation");
1176                        })?;
1177                        let RenderedDatabaseRuntimeInfo {
1178                            job_infos,
1179                            stream_actors,
1180                            mut source_splits,
1181                            batch_refresh,
1182                        } = rendered_info;
1183                        recoverer.inject_database_initial_barrier(
1184                            database_id,
1185                            job_infos,
1186                            &recovery_context.job_extra_info,
1187                            &mut state_table_committed_epochs,
1188                            &mut state_table_log_epochs,
1189                            &recovery_context.fragment_relations,
1190                            &stream_actors,
1191                            &mut source_splits,
1192                            &mut creating_jobs,
1193                            &mut mv_depended_subscriptions,
1194                            is_paused,
1195                            &hummock_version_stats,
1196                            &mut cdc_table_snapshot_splits,
1197                            batch_refresh,
1198                        )?
1199                    };
1200                    let collector = match result {
1201                        Ok(database) => {
1202                            DatabaseInitialBarrierCollector {
1203                                database_id,
1204                                initializing_partial_graphs: recoverer.all_initializing(),
1205                                database,
1206                            }
1207                        }
1208                        Err(e) => {
1209                            warn!(%database_id, e = %e.as_report(), "failed to inject database initial barrier");
1210                            assert!(failed_databases.insert(database_id, recoverer.failed()).is_none(), "non-duplicate");
1211                            continue;
1212                        }
1213                    };
1214                    if !collector.is_collected() {
1215                        assert!(collecting_databases.insert(database_id, collector).is_none());
1216                    } else {
1217                        warn!(%database_id, "database has no node to inject initial barrier");
1218                        assert!(collected_databases.insert(database_id, collector.finish()).is_none());
1219                    }
1220                }
1221                if !empty_databases.is_empty() {
1222                    info!(?empty_databases, "empty database in global recovery");
1223                }
1224                while !collecting_databases.is_empty() {
1225                    match partial_graph_manager.next_event(&self.context).await {
1226                        PartialGraphManagerEvent::Worker(_, WorkerEvent::WorkerConnected) => {
1227                            // not handle WorkerConnected yet
1228                        }
1229                        PartialGraphManagerEvent::Worker(worker_id, WorkerEvent::WorkerError { err, affected_partial_graphs }) => {
1230                            let affected_databases: HashSet<_> = affected_partial_graphs.into_iter().map(|partial_graph_id| {
1231                                let (database_id, _) = from_partial_graph_id(partial_graph_id);
1232                                database_id
1233                            }).collect();
1234                            warn!(%worker_id, err = %err.as_report(), "worker node failure during recovery");
1235                            for (failed_database_id, collector) in collecting_databases.extract_if(|database_id, collector| {
1236                                !collector.is_valid_after_worker_err(worker_id) || affected_databases.contains(database_id)
1237                            }) {
1238                                warn!(%failed_database_id, %worker_id, "database failed to recovery in global recovery due to worker node err");
1239                                let resetting_partial_graphs: HashSet<_> = collector.all_partial_graphs().collect();
1240                                partial_graph_manager.reset_partial_graphs(resetting_partial_graphs.iter().copied());
1241                                assert!(failed_databases.insert(failed_database_id, resetting_partial_graphs).is_none());
1242                            }
1243                        }
1244                        PartialGraphManagerEvent::PartialGraph(partial_graph_id, event) => {
1245                            match event {
1246                                PartialGraphEvent::BarrierCollected(_) => {
1247                                    unreachable!("no barrier collected event on initializing")
1248                                }
1249                                PartialGraphEvent::Reset(_) => {
1250                                    // Reset responses only carry diagnostic root errors. Recovery
1251                                    // itself only needs to track the partial graphs still resetting.
1252                                    let (database_id, _) =
1253                                        from_partial_graph_id(partial_graph_id);
1254                                    let resetting_partial_graphs = failed_databases
1255                                        .get_mut(&database_id)
1256                                        .expect("reset partial graph should belong to a failed database");
1257                                    assert!(
1258                                        resetting_partial_graphs.remove(&partial_graph_id),
1259                                        "partial graph {partial_graph_id} should be resetting"
1260                                    );
1261                                }
1262                                PartialGraphEvent::Error(worker_id) => {
1263                                    let (database_id, _) = from_partial_graph_id(partial_graph_id);
1264                                    if let Some(collector) = collecting_databases.remove(&database_id) {
1265                                        warn!(%database_id, %worker_id, "database reset during global recovery");
1266                                        let resetting_partial_graphs: HashSet<_> = collector.all_partial_graphs().collect();
1267                                        partial_graph_manager.reset_partial_graphs(resetting_partial_graphs.iter().copied());
1268                                        assert!(failed_databases.insert(database_id, resetting_partial_graphs).is_none());
1269                                    } else if let Some(database) = collected_databases.remove(&database_id) {
1270                                        warn!(%database_id, %worker_id, "database initialized but later reset during global recovery");
1271                                        let resetting_partial_graphs: HashSet<_> = database_partial_graphs(database_id, database.independent_checkpoint_job_controls.keys().copied()).collect();
1272                                        partial_graph_manager.reset_partial_graphs(resetting_partial_graphs.iter().copied());
1273                                        assert!(failed_databases.insert(database_id, resetting_partial_graphs).is_none());
1274                                    } else {
1275                                        assert!(failed_databases.contains_key(&database_id));
1276                                    }
1277                                }
1278                                PartialGraphEvent::Initialized => {
1279                                    let (database_id, _) = from_partial_graph_id(partial_graph_id);
1280                                    if failed_databases.contains_key(&database_id) {
1281                                        assert!(!collecting_databases.contains_key(&database_id));
1282                                        // ignore the lately initialized partial graph of failed database
1283                                        continue;
1284                                    }
1285                                    let Entry::Occupied(mut entry) = collecting_databases.entry(database_id) else {
1286                                        unreachable!("should exist")
1287                                    };
1288                                    let collector = entry.get_mut();
1289                                    collector.partial_graph_initialized(partial_graph_id);
1290                                    if collector.is_collected() {
1291                                        let collector = entry.remove();
1292                                        assert!(collected_databases.insert(database_id, collector.finish()).is_none());
1293                                    }
1294                                }
1295                            }
1296                        }
1297                    }
1298                }
1299                debug!("collected initial barrier");
1300                if !creating_jobs.is_empty() {
1301                    warn!(job_ids = ?creating_jobs.iter().collect_vec(), "unused recovered creating jobs in recovery");
1302                }
1303                if !mv_depended_subscriptions.is_empty() {
1304                    warn!(?mv_depended_subscriptions, "unused subscription infos in recovery");
1305                }
1306                if !state_table_committed_epochs.is_empty() {
1307                    warn!(?state_table_committed_epochs, "unused state table committed epoch in recovery");
1308                }
1309                if !enable_per_database_isolation && !failed_databases.is_empty() {
1310                    return Err(anyhow!(
1311                        "global recovery failed due to failure of databases {:?}",
1312                        failed_databases.keys().collect_vec()).into()
1313                    );
1314                }
1315                let checkpoint_control = CheckpointControl::recover(
1316                    collected_databases,
1317                    failed_databases,
1318                    hummock_version_stats,
1319                    self.env.clone(),
1320                );
1321
1322                let reader = self.env.system_params_reader().await;
1323                let checkpoint_frequency = reader.checkpoint_frequency();
1324                let barrier_interval = Duration::from_millis(reader.barrier_interval_ms() as u64);
1325                let periodic_barriers = PeriodicBarriers::new(
1326                    barrier_interval,
1327                    checkpoint_frequency,
1328                    database_infos,
1329                );
1330
1331                self.context
1332                    .refresh_table_refill_runtime_state_after_recovery()
1333                    .await?;
1334
1335                Ok((
1336                    active_streaming_nodes,
1337                    partial_graph_manager,
1338                    checkpoint_control,
1339                    periodic_barriers,
1340                ))
1341            }
1342        }.inspect_err(|err: &MetaError| {
1343            tracing::error!(error = %err.as_report(), "recovery failed");
1344            event_log_manager_ref.add_event_logs(vec![Event::Recovery(
1345                EventRecovery::global_recovery_failure(recovery_reason.clone(), err.to_report_string()),
1346            )]);
1347            GLOBAL_META_METRICS.recovery_failure_cnt.with_label_values(&["global"]).inc();
1348        }))
1349        .instrument(tracing::info_span!("recovery_attempt"));
1350
1351        let mut recover_txs = vec![];
1352        let mut update_barrier_requests = vec![];
1353        pin_mut!(recovery_future);
1354        let mut request_rx_closed = false;
1355        let new_state = loop {
1356            select! {
1357                biased;
1358                new_state = &mut recovery_future => {
1359                    break new_state.expect("Retry until recovery success.");
1360                }
1361                request = pin!(self.request_rx.recv()), if !request_rx_closed => {
1362                    let Some(request) = request else {
1363                        warn!("request rx channel closed during recovery");
1364                        request_rx_closed = true;
1365                        continue;
1366                    };
1367                    match request {
1368                        BarrierManagerRequest::GetBackfillProgress(tx) => {
1369                            let _ = tx.send(Err(anyhow!("cluster under recovery[{}]", recovery_reason).into()));
1370                        }
1371                        BarrierManagerRequest::GetFragmentBackfillProgress(tx) => {
1372                            let _ = tx.send(Err(anyhow!("cluster under recovery[{}]", recovery_reason).into()));
1373                        }
1374                        BarrierManagerRequest::GetCdcProgress(tx) => {
1375                            let _ = tx.send(Err(anyhow!("cluster under recovery[{}]", recovery_reason).into()));
1376                        }
1377                        BarrierManagerRequest::AdhocRecovery(tx) => {
1378                            recover_txs.push(tx);
1379                        }
1380                        BarrierManagerRequest::UpdateDatabaseBarrier(request) => {
1381                            update_barrier_requests.push(request);
1382                        }
1383                    }
1384                }
1385            }
1386        };
1387
1388        let duration = recovery_timer.stop_and_record();
1389
1390        (
1391            self.active_streaming_nodes,
1392            self.partial_graph_manager,
1393            self.checkpoint_control,
1394            self.periodic_barriers,
1395        ) = new_state;
1396
1397        tracing::info!("recovery success");
1398
1399        for UpdateDatabaseBarrierRequest {
1400            database_id,
1401            barrier_interval_ms,
1402            checkpoint_frequency,
1403            sender,
1404        } in update_barrier_requests
1405        {
1406            self.periodic_barriers.update_database_barrier(
1407                database_id,
1408                barrier_interval_ms,
1409                checkpoint_frequency,
1410            );
1411            let _ = sender.send(());
1412        }
1413
1414        for tx in recover_txs {
1415            let _ = tx.send(());
1416        }
1417
1418        let recovering_databases = self
1419            .checkpoint_control
1420            .recovering_databases()
1421            .map(|database| database.as_raw_id())
1422            .collect_vec();
1423        let running_databases = self
1424            .checkpoint_control
1425            .running_databases()
1426            .map(|database| database.as_raw_id())
1427            .collect_vec();
1428
1429        event_log_manager_ref.add_event_logs(vec![Event::Recovery(
1430            EventRecovery::global_recovery_success(
1431                recovery_reason.clone(),
1432                duration as f32,
1433                running_databases,
1434                recovering_databases,
1435            ),
1436        )]);
1437
1438        self.env
1439            .notification_manager()
1440            .notify_frontend_without_version(Operation::Update, Info::Recovery(Recovery {}));
1441        self.env
1442            .notification_manager()
1443            .notify_compute_without_version(Operation::Update, Info::Recovery(Recovery {}));
1444    }
1445}