Skip to main content

risingwave_meta/barrier/checkpoint/
control.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::future::{Future, poll_fn};
18use std::ops::Bound::{Excluded, Unbounded};
19use std::sync::atomic::AtomicU32;
20use std::task::Poll;
21
22use anyhow::anyhow;
23use fail::fail_point;
24use itertools::Itertools;
25use risingwave_common::catalog::{DatabaseId, TableId};
26use risingwave_common::id::JobId;
27use risingwave_common::metrics::{LabelGuardedHistogram, LabelGuardedIntGauge};
28use risingwave_common::util::epoch::{Epoch, EpochPair};
29use risingwave_common::util::stream_graph_visitor::visit_stream_node_cont;
30use risingwave_meta_model::WorkerId;
31use risingwave_pb::common::WorkerNode;
32use risingwave_pb::hummock::HummockVersionStats;
33use risingwave_pb::id::{FragmentId, PartialGraphId};
34use risingwave_pb::stream_plan::stream_node::NodeBody;
35use risingwave_pb::stream_plan::{DispatcherType as PbDispatcherType, PbSubscriptionUpstreamInfo};
36use risingwave_pb::stream_service::BarrierCompleteResponse;
37use risingwave_pb::stream_service::streaming_control_stream_response::ResetPartialGraphResponse;
38use tracing::{debug, warn};
39use uuid::Uuid;
40
41use crate::barrier::cdc_progress::CdcProgress;
42use crate::barrier::checkpoint::independent_job::{
43    BatchRefreshJobTriggerContext, IndependentCheckpointJob, IndependentCheckpointJobControl,
44};
45use crate::barrier::checkpoint::recovery::{
46    DatabaseRecoveringState, DatabaseStatusAction, EnterInitializing, EnterRunning,
47    RecoveringStateAction,
48};
49use crate::barrier::checkpoint::state::{ApplyCommandInfo, BarrierWorkerState};
50use crate::barrier::complete_task::{BarrierCompleteOutput, CompleteBarrierTask};
51use crate::barrier::info::{InflightDatabaseInfo, SharedActorInfos};
52use crate::barrier::partial_graph::{CollectedBarrier, PartialGraphManager, PartialGraphStat};
53use crate::barrier::progress::TrackingJob;
54use crate::barrier::rpc::{from_partial_graph_id, to_partial_graph_id};
55use crate::barrier::schedule::{NewBarrier, PeriodicBarriers};
56use crate::barrier::utils::{BarrierItemCollector, collect_independent_job_commit_epoch_info};
57use crate::barrier::{
58    BackfillProgress, Command, CreateStreamingJobType, FragmentBackfillProgress, Reschedule,
59};
60use crate::controller::fragment::InflightFragmentInfo;
61use crate::controller::scale::{build_no_shuffle_fragment_graph_edges, find_no_shuffle_graphs};
62use crate::manager::MetaSrvEnv;
63use crate::notification::Notifier;
64
65fn fragment_has_online_unreschedulable_scan(fragment: &InflightFragmentInfo) -> bool {
66    let mut has_unreschedulable_scan = false;
67    visit_stream_node_cont(&fragment.nodes, |node| {
68        if let Some(NodeBody::StreamScan(stream_scan)) = node.node_body.as_ref() {
69            let scan_type = stream_scan.stream_scan_type();
70            if !scan_type.is_reschedulable(true) {
71                has_unreschedulable_scan = true;
72                return false;
73            }
74        }
75        true
76    });
77    has_unreschedulable_scan
78}
79
80fn collect_fragment_upstream_fragment_ids(
81    fragment: &InflightFragmentInfo,
82    upstream_fragment_ids: &mut HashSet<FragmentId>,
83) {
84    visit_stream_node_cont(&fragment.nodes, |node| {
85        if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref() {
86            upstream_fragment_ids.insert(merge.upstream_fragment_id);
87        }
88        true
89    });
90}
91
92use crate::model::ActorId;
93use crate::rpc::metrics::GLOBAL_META_METRICS;
94use crate::{MetaError, MetaResult};
95
96pub(crate) struct CheckpointControl {
97    pub(crate) env: MetaSrvEnv,
98    pub(super) databases: HashMap<DatabaseId, DatabaseCheckpointControlStatus>,
99    pub(super) hummock_version_stats: HummockVersionStats,
100    /// The maximum number of pending barriers in each partial graph.
101    pub(crate) in_flight_barrier_nums: usize,
102}
103
104impl CheckpointControl {
105    pub fn new(env: MetaSrvEnv) -> Self {
106        Self {
107            in_flight_barrier_nums: env.opts.in_flight_barrier_nums,
108            env,
109            databases: Default::default(),
110            hummock_version_stats: Default::default(),
111        }
112    }
113
114    pub(crate) fn recover(
115        databases: HashMap<DatabaseId, DatabaseCheckpointControl>,
116        failed_databases: HashMap<DatabaseId, HashSet<PartialGraphId>>, /* `database_id` -> set of resetting partial graph ids */
117        hummock_version_stats: HummockVersionStats,
118        env: MetaSrvEnv,
119    ) -> Self {
120        env.shared_actor_infos()
121            .retain_databases(databases.keys().chain(failed_databases.keys()).cloned());
122        Self {
123            in_flight_barrier_nums: env.opts.in_flight_barrier_nums,
124            env,
125            databases: databases
126                .into_iter()
127                .map(|(database_id, control)| {
128                    (
129                        database_id,
130                        DatabaseCheckpointControlStatus::Running(control),
131                    )
132                })
133                .chain(failed_databases.into_iter().map(
134                    |(database_id, resetting_partial_graphs)| {
135                        (
136                            database_id,
137                            DatabaseCheckpointControlStatus::Recovering(
138                                DatabaseRecoveringState::new_resetting(
139                                    database_id,
140                                    resetting_partial_graphs,
141                                ),
142                            ),
143                        )
144                    },
145                ))
146                .collect(),
147            hummock_version_stats,
148        }
149    }
150
151    pub(crate) fn ack_completed(
152        &mut self,
153        partial_graph_manager: &mut PartialGraphManager,
154        output: BarrierCompleteOutput,
155    ) {
156        self.hummock_version_stats = output.hummock_version_stats;
157        for (database_id, (command_prev_epoch, independent_job_epochs)) in output.epochs_to_ack {
158            self.databases
159                .get_mut(&database_id)
160                .expect("should exist")
161                .expect_running("should have wait for completing command before enter recovery")
162                .ack_completed(
163                    partial_graph_manager,
164                    command_prev_epoch,
165                    independent_job_epochs,
166                );
167        }
168    }
169
170    pub(crate) fn next_complete_barrier_task(
171        &mut self,
172        periodic_barriers: &mut PeriodicBarriers,
173        partial_graph_manager: &mut PartialGraphManager,
174    ) -> Option<CompleteBarrierTask> {
175        let mut task = None;
176        for database in self.databases.values_mut() {
177            let Some(database) = database.running_state_mut() else {
178                continue;
179            };
180            database.next_complete_barrier_task(
181                periodic_barriers,
182                partial_graph_manager,
183                &mut task,
184                &self.hummock_version_stats,
185            );
186        }
187        task
188    }
189
190    pub(crate) fn barrier_collected(
191        &mut self,
192        partial_graph_id: PartialGraphId,
193        collected_barrier: CollectedBarrier<'_>,
194        periodic_barriers: &mut PeriodicBarriers,
195    ) -> MetaResult<()> {
196        let (database_id, _) = from_partial_graph_id(partial_graph_id);
197        let database_status = self.databases.get_mut(&database_id).expect("should exist");
198        match database_status {
199            DatabaseCheckpointControlStatus::Running(database) => {
200                database.barrier_collected(partial_graph_id, collected_barrier, periodic_barriers)
201            }
202            DatabaseCheckpointControlStatus::Recovering(_) => {
203                if cfg!(debug_assertions) {
204                    panic!(
205                        "receive collected barrier {:?} on recovering database {} from partial graph {}",
206                        collected_barrier, database_id, partial_graph_id
207                    );
208                } else {
209                    warn!(?collected_barrier, %partial_graph_id, "ignore collected barrier on recovering database");
210                }
211                Ok(())
212            }
213        }
214    }
215
216    pub(crate) fn recovering_databases(&self) -> impl Iterator<Item = DatabaseId> + '_ {
217        self.databases.iter().filter_map(|(database_id, database)| {
218            database.running_state().is_none().then_some(*database_id)
219        })
220    }
221
222    pub(crate) fn running_databases(&self) -> impl Iterator<Item = DatabaseId> + '_ {
223        self.databases.iter().filter_map(|(database_id, database)| {
224            database.running_state().is_some().then_some(*database_id)
225        })
226    }
227
228    pub(crate) fn database_info(&self, database_id: DatabaseId) -> Option<&InflightDatabaseInfo> {
229        self.databases
230            .get(&database_id)
231            .and_then(|database| database.running_state())
232            .map(|database| &database.database_info)
233    }
234
235    /// return Some(failed `database_id` -> `err`)
236    pub(crate) fn handle_new_barrier(
237        &mut self,
238        new_barrier: NewBarrier,
239        partial_graph_manager: &mut PartialGraphManager,
240        worker_nodes: &HashMap<WorkerId, WorkerNode>,
241    ) -> MetaResult<()> {
242        let NewBarrier {
243            database_id,
244            command,
245            span,
246            checkpoint,
247        } = new_barrier;
248
249        if let Some((mut command, notifier)) = command {
250            if let &mut Command::CreateStreamingJob {
251                ref mut cross_db_snapshot_backfill_info,
252                ref info,
253                ..
254            } = &mut command
255            {
256                for (table_id, snapshot_epoch) in
257                    &mut cross_db_snapshot_backfill_info.upstream_mv_table_id_to_backfill_epoch
258                {
259                    for database in self.databases.values() {
260                        if let Some(database) = database.running_state()
261                            && database.database_info.contains_job(table_id.as_job_id())
262                        {
263                            if let Some(committed_epoch) = database.committed_epoch {
264                                *snapshot_epoch = Some(committed_epoch);
265                            }
266                            break;
267                        }
268                    }
269                    if snapshot_epoch.is_none() {
270                        let table_id = *table_id;
271                        warn!(
272                            ?cross_db_snapshot_backfill_info,
273                            ?table_id,
274                            ?info,
275                            "database of cross db upstream table not found"
276                        );
277                        let err: MetaError =
278                            anyhow!("database of cross db upstream table {} not found", table_id)
279                                .into();
280                        notifier.notify_start_failed(err);
281
282                        return Ok(());
283                    }
284                }
285            }
286
287            let database = match self.databases.entry(database_id) {
288                Entry::Occupied(entry) => entry
289                    .into_mut()
290                    .expect_running("should not have command when not running"),
291                Entry::Vacant(entry) => match &command {
292                    Command::CreateStreamingJob { info, job_type, .. } => {
293                        let CreateStreamingJobType::Normal = job_type else {
294                            if cfg!(debug_assertions) {
295                                panic!(
296                                    "unexpected first job of type {job_type:?} with info {info:?}"
297                                );
298                            } else {
299                                notifier.notify_start_failed(anyhow!("unexpected job_type {job_type:?} for first job {} in database {database_id}", info.streaming_job.id()).into());
300                                return Ok(());
301                            }
302                        };
303                        let new_database = DatabaseCheckpointControl::new(
304                            database_id,
305                            self.env.shared_actor_infos().clone(),
306                        );
307                        let adder = partial_graph_manager.add_partial_graph(
308                            to_partial_graph_id(database_id, None),
309                            new_database.term_id(),
310                            DatabaseCheckpointControlMetrics::new(database_id),
311                        );
312                        adder.added();
313                        entry
314                            .insert(DatabaseCheckpointControlStatus::Running(new_database))
315                            .expect_running("just initialized as running")
316                    }
317                    Command::Flush
318                    | Command::Pause
319                    | Command::Resume
320                    | Command::DropStreamingJobs { .. }
321                    | Command::DropSubscription { .. } => {
322                        notifier.start().started();
323                        warn!(?command, "skip command for empty database");
324                        return Ok(());
325                    }
326                    Command::RescheduleIntent { .. }
327                    | Command::ReplaceStreamJob(_)
328                    | Command::SourceChangeSplit(_)
329                    | Command::Throttle { .. }
330                    | Command::CreateSubscription { .. }
331                    | Command::AlterSubscriptionRetention { .. }
332                    | Command::ConnectorPropsChange(_)
333                    | Command::Refresh { .. }
334                    | Command::ListFinish { .. }
335                    | Command::LoadFinish { .. }
336                    | Command::FinishRefresh { .. }
337                    | Command::ResetSource { .. }
338                    | Command::ResumeBackfill { .. }
339                    | Command::InjectSourceOffsets { .. } => {
340                        if cfg!(debug_assertions) {
341                            panic!(
342                                "new database graph info can only be created for normal creating streaming job, but get command: {} {:?}",
343                                database_id, command
344                            )
345                        } else {
346                            warn!(%database_id, ?command, "database does not exist while handling the command");
347                            notifier.notify_start_failed(anyhow!("database {database_id} does not exist while handling command {command:?}").into());
348                            return Ok(());
349                        }
350                    }
351                },
352            };
353
354            database.handle_new_barrier(
355                Some((command, notifier)),
356                checkpoint,
357                span,
358                partial_graph_manager,
359                &self.hummock_version_stats,
360                worker_nodes,
361            )
362        } else {
363            let database = match self.databases.entry(database_id) {
364                Entry::Occupied(entry) => entry.into_mut(),
365                Entry::Vacant(_) => {
366                    // If it does not exist in the HashMap yet, it means that the first streaming
367                    // job has not been created, and we do not need to send a barrier.
368                    return Ok(());
369                }
370            };
371            let Some(database) = database.running_state_mut() else {
372                // Skip new barrier for database which is not running.
373                return Ok(());
374            };
375            if partial_graph_manager.pending_barrier_num(database.partial_graph_id)
376                >= self.in_flight_barrier_nums
377            {
378                // Skip new barrier with no explicit command when the database should pause inject additional barrier
379                return Ok(());
380            }
381            database.handle_new_barrier(
382                None,
383                checkpoint,
384                span,
385                partial_graph_manager,
386                &self.hummock_version_stats,
387                worker_nodes,
388            )
389        }
390    }
391
392    pub(crate) fn gen_backfill_progress(&self) -> HashMap<JobId, BackfillProgress> {
393        let mut progress = HashMap::new();
394        for status in self.databases.values() {
395            let Some(database_checkpoint_control) = status.running_state() else {
396                continue;
397            };
398            // Progress of normal backfill
399            progress.extend(
400                database_checkpoint_control
401                    .database_info
402                    .gen_backfill_progress(),
403            );
404            // Progress of independent checkpoint jobs
405            for (job_id, job) in &database_checkpoint_control.independent_checkpoint_job_controls {
406                if let Some(p) = job.gen_backfill_progress() {
407                    progress.insert(*job_id, p);
408                }
409            }
410        }
411        progress
412    }
413
414    pub(crate) fn gen_fragment_backfill_progress(&self) -> Vec<FragmentBackfillProgress> {
415        let mut progress = Vec::new();
416        for status in self.databases.values() {
417            let Some(database_checkpoint_control) = status.running_state() else {
418                continue;
419            };
420            progress.extend(
421                database_checkpoint_control
422                    .database_info
423                    .gen_fragment_backfill_progress(),
424            );
425            for job in database_checkpoint_control
426                .independent_checkpoint_job_controls
427                .values()
428            {
429                progress.extend(job.gen_fragment_backfill_progress());
430            }
431        }
432        progress
433    }
434
435    pub(crate) fn gen_cdc_progress(&self) -> HashMap<JobId, CdcProgress> {
436        let mut progress = HashMap::new();
437        for status in self.databases.values() {
438            let Some(database_checkpoint_control) = status.running_state() else {
439                continue;
440            };
441            // Progress of normal backfill
442            progress.extend(database_checkpoint_control.database_info.gen_cdc_progress());
443        }
444        progress
445    }
446
447    pub(crate) fn databases_failed_at_worker_err(
448        &mut self,
449        worker_id: WorkerId,
450    ) -> impl Iterator<Item = DatabaseId> + '_ {
451        self.databases
452            .iter_mut()
453            .filter_map(
454                move |(database_id, database_status)| match database_status {
455                    DatabaseCheckpointControlStatus::Running(control) => {
456                        if !control.is_valid_after_worker_err(worker_id) {
457                            Some(*database_id)
458                        } else {
459                            None
460                        }
461                    }
462                    DatabaseCheckpointControlStatus::Recovering(state) => {
463                        if !state.is_valid_after_worker_err(worker_id) {
464                            Some(*database_id)
465                        } else {
466                            None
467                        }
468                    }
469                },
470            )
471    }
472
473    // ── Batch refresh trigger helpers (delegating to DatabaseCheckpointControl) ──
474
475    pub(crate) fn get_batch_refresh_trigger_info(
476        &self,
477        database_id: DatabaseId,
478        job_id: JobId,
479    ) -> u64 {
480        let database = self
481            .databases
482            .get(&database_id)
483            .and_then(|s| s.running_state())
484            .expect("database should be running for batch refresh trigger");
485        database.get_batch_refresh_trigger_info(job_id)
486    }
487
488    pub(crate) fn start_batch_refresh_run(
489        &mut self,
490        database_id: DatabaseId,
491        job_id: JobId,
492        context: &BatchRefreshJobTriggerContext,
493        worker_nodes: &HashMap<WorkerId, WorkerNode>,
494        actor_id_counter: &AtomicU32,
495        partial_graph_manager: &mut PartialGraphManager,
496    ) -> MetaResult<bool> {
497        let database = self
498            .databases
499            .get_mut(&database_id)
500            .and_then(|s| s.running_state_mut())
501            .expect("database should be running");
502        database.start_batch_refresh_run(
503            job_id,
504            context,
505            worker_nodes,
506            actor_id_counter,
507            partial_graph_manager,
508        )
509    }
510
511    pub(crate) fn apply_batch_refresh_fragment_infos(
512        &mut self,
513        database_id: DatabaseId,
514        job_id: JobId,
515    ) {
516        // This is called synchronously after `start_batch_refresh_run` returns `true`, with no
517        // barrier-worker event in between that could transition the job to Resetting.
518        let database = self
519            .databases
520            .get_mut(&database_id)
521            .and_then(|s| s.running_state_mut())
522            .expect("database should be running");
523        let br_job = match database
524            .independent_checkpoint_job_controls
525            .get(&job_id)
526            .expect("job should exist")
527            .running()
528            .expect("job should be running")
529        {
530            IndependentCheckpointJob::BatchRefresh(job) => job,
531            _ => panic!("expected batch refresh job"),
532        };
533        if let Some(fragment_infos) = br_job.fragment_infos() {
534            database
535                .database_info
536                .shared_actor_infos
537                .upsert(database_id, fragment_infos.values().map(|f| (f, job_id)));
538        }
539    }
540}
541
542pub(crate) enum CheckpointControlEvent<'a> {
543    EnteringInitializing(DatabaseStatusAction<'a, EnterInitializing>),
544    EnteringRunning(DatabaseStatusAction<'a, EnterRunning>),
545    /// A batch refresh job is idle and its upstream has advanced past the refresh interval.
546    /// Carries owned values so the async handler can call into context without borrowing self.
547    BatchRefreshTrigger {
548        database_id: DatabaseId,
549        job_id: JobId,
550    },
551}
552
553impl CheckpointControl {
554    pub(crate) fn on_partial_graph_reset(
555        &mut self,
556        partial_graph_id: PartialGraphId,
557        reset_resps: HashMap<WorkerId, ResetPartialGraphResponse>,
558    ) {
559        let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
560        match self.databases.get_mut(&database_id).expect("should exist") {
561            DatabaseCheckpointControlStatus::Running(database) => {
562                if let Some(independent_job_id) = independent_job_id {
563                    match database
564                        .independent_checkpoint_job_controls
565                        .remove(&independent_job_id)
566                    {
567                        Some(independent_job) => {
568                            database
569                                .pending_independent_job_subscriptions_to_drop
570                                .extend(independent_job.on_partial_graph_reset());
571                        }
572                        None => {
573                            if cfg!(debug_assertions) {
574                                panic!(
575                                    "receive reset partial graph resp on non-existing independent job {independent_job_id} in database {database_id}"
576                                )
577                            }
578                            warn!(
579                                %database_id,
580                                %independent_job_id,
581                                "ignore reset partial graph resp on non-existing independent job on running database"
582                            );
583                        }
584                    }
585                } else {
586                    unreachable!("should not receive reset database resp when database running")
587                }
588            }
589            DatabaseCheckpointControlStatus::Recovering(state) => {
590                state.on_partial_graph_reset(partial_graph_id, reset_resps);
591            }
592        }
593    }
594
595    pub(crate) fn on_partial_graph_initialized(
596        &mut self,
597        partial_graph_id: PartialGraphId,
598        partial_graph_manager: &mut PartialGraphManager,
599    ) -> MetaResult<()> {
600        let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
601        match self.databases.get_mut(&database_id).expect("should exist") {
602            DatabaseCheckpointControlStatus::Running(database) => {
603                let Some(independent_job_id) = independent_job_id else {
604                    unreachable!("database partial graph should not initialize when running")
605                };
606                let job = database
607                    .independent_checkpoint_job_controls
608                    .get_mut(&independent_job_id)
609                    .expect("independent job should exist");
610                match job.running_mut().expect("job should be running") {
611                    IndependentCheckpointJob::BatchRefresh(job) => {
612                        job.on_log_store_initialized(partial_graph_manager)
613                    }
614                    IndependentCheckpointJob::CreatingStreamingJob(_) => {
615                        unreachable!("creating streaming job should not initialize when running")
616                    }
617                    IndependentCheckpointJob::IcebergV3(_) => {
618                        unreachable!("Iceberg V3 jobs do not wait for graph initialization")
619                    }
620                }
621            }
622            DatabaseCheckpointControlStatus::Recovering(state) => {
623                state.partial_graph_initialized(partial_graph_id);
624                Ok(())
625            }
626        }
627    }
628
629    pub(crate) fn next_event(
630        &mut self,
631    ) -> impl Future<Output = CheckpointControlEvent<'_>> + Send + '_ {
632        let mut this = Some(self);
633        poll_fn(move |cx| {
634            let Some(this_mut) = this.as_mut() else {
635                unreachable!("should not be polled after poll ready")
636            };
637            for (&database_id, database_status) in &mut this_mut.databases {
638                match database_status {
639                    DatabaseCheckpointControlStatus::Running(database) => {
640                        // Check if any idle batch refresh job should start a refresh run.
641                        if let Some(committed_epoch) = database.committed_epoch {
642                            for (job_id, job) in &database.independent_checkpoint_job_controls {
643                                if let Some(IndependentCheckpointJob::BatchRefresh(br_job)) =
644                                    job.running()
645                                    && br_job.should_start_refresh(committed_epoch)
646                                {
647                                    let job_id = *job_id;
648                                    let _ = this.take().expect("checked Some");
649                                    return Poll::Ready(
650                                        CheckpointControlEvent::BatchRefreshTrigger {
651                                            database_id,
652                                            job_id,
653                                        },
654                                    );
655                                }
656                            }
657                        }
658                    }
659                    DatabaseCheckpointControlStatus::Recovering(state) => {
660                        let poll_result = state.poll_next_event(cx);
661                        if let Poll::Ready(action) = poll_result {
662                            let this = this.take().expect("checked Some");
663                            return Poll::Ready(match action {
664                                RecoveringStateAction::EnterInitializing(reset_workers) => {
665                                    CheckpointControlEvent::EnteringInitializing(
666                                        this.new_database_status_action(
667                                            database_id,
668                                            EnterInitializing(reset_workers),
669                                        ),
670                                    )
671                                }
672                                RecoveringStateAction::EnterRunning => {
673                                    CheckpointControlEvent::EnteringRunning(
674                                        this.new_database_status_action(database_id, EnterRunning),
675                                    )
676                                }
677                            });
678                        }
679                    }
680                }
681            }
682            Poll::Pending
683        })
684    }
685}
686
687pub(crate) enum DatabaseCheckpointControlStatus {
688    Running(DatabaseCheckpointControl),
689    Recovering(DatabaseRecoveringState),
690}
691
692impl DatabaseCheckpointControlStatus {
693    fn running_state(&self) -> Option<&DatabaseCheckpointControl> {
694        match self {
695            DatabaseCheckpointControlStatus::Running(state) => Some(state),
696            DatabaseCheckpointControlStatus::Recovering(_) => None,
697        }
698    }
699
700    fn running_state_mut(&mut self) -> Option<&mut DatabaseCheckpointControl> {
701        match self {
702            DatabaseCheckpointControlStatus::Running(state) => Some(state),
703            DatabaseCheckpointControlStatus::Recovering(_) => None,
704        }
705    }
706
707    fn expect_running(&mut self, reason: &'static str) -> &mut DatabaseCheckpointControl {
708        match self {
709            DatabaseCheckpointControlStatus::Running(state) => state,
710            DatabaseCheckpointControlStatus::Recovering(_) => {
711                panic!("should be at running: {}", reason)
712            }
713        }
714    }
715}
716
717pub(in crate::barrier) struct DatabaseCheckpointControlMetrics {
718    barrier_latency: LabelGuardedHistogram,
719    in_flight_barrier_nums: LabelGuardedIntGauge,
720    all_barrier_nums: LabelGuardedIntGauge,
721}
722
723impl DatabaseCheckpointControlMetrics {
724    pub(in crate::barrier) fn new(database_id: DatabaseId) -> Self {
725        let database_id_str = database_id.to_string();
726        let barrier_latency = GLOBAL_META_METRICS
727            .barrier_latency
728            .with_guarded_label_values(&[&database_id_str]);
729        let in_flight_barrier_nums = GLOBAL_META_METRICS
730            .in_flight_barrier_nums
731            .with_guarded_label_values(&[&database_id_str]);
732        let all_barrier_nums = GLOBAL_META_METRICS
733            .all_barrier_nums
734            .with_guarded_label_values(&[&database_id_str]);
735        Self {
736            barrier_latency,
737            in_flight_barrier_nums,
738            all_barrier_nums,
739        }
740    }
741}
742
743impl PartialGraphStat for DatabaseCheckpointControlMetrics {
744    fn observe_barrier_latency(&self, _epoch: EpochPair, barrier_latency_secs: f64) {
745        self.barrier_latency.observe(barrier_latency_secs);
746    }
747
748    fn observe_barrier_num(&self, inflight_barrier_num: usize, collected_barrier_num: usize) {
749        self.in_flight_barrier_nums.set(inflight_barrier_num as _);
750        self.all_barrier_nums
751            .set((inflight_barrier_num + collected_barrier_num) as _);
752    }
753}
754
755/// Controls the concurrent execution of commands.
756pub(in crate::barrier) struct DatabaseCheckpointControl {
757    pub(super) database_id: DatabaseId,
758    /// Identifies the current recovery incarnation of this database.
759    pub(super) term_id: String,
760    partial_graph_id: PartialGraphId,
761    pub(super) state: BarrierWorkerState,
762
763    finishing_jobs_collector:
764        BarrierItemCollector<JobId, (Vec<BarrierCompleteResponse>, TrackingJob), ()>,
765    /// The barrier that are completing.
766    completing_barrier: Option<EpochPair>,
767
768    committed_epoch: Option<u64>,
769
770    /// `None` while the database has no streaming job, so that a frozen timestamp does not
771    /// render as an ever-growing barrier pending time.
772    last_committed_barrier_time: Option<LabelGuardedIntGauge>,
773
774    pub(super) database_info: InflightDatabaseInfo,
775    pub independent_checkpoint_job_controls: HashMap<JobId, IndependentCheckpointJobControl>,
776    pub(super) pending_independent_job_subscriptions_to_drop: Vec<PbSubscriptionUpstreamInfo>,
777}
778
779impl DatabaseCheckpointControl {
780    fn new(database_id: DatabaseId, shared_actor_infos: SharedActorInfos) -> Self {
781        Self {
782            database_id,
783            term_id: Uuid::new_v4().to_string(),
784            partial_graph_id: to_partial_graph_id(database_id, None),
785            state: BarrierWorkerState::new(),
786            finishing_jobs_collector: BarrierItemCollector::new(false),
787            completing_barrier: None,
788            committed_epoch: None,
789            last_committed_barrier_time: None,
790            database_info: InflightDatabaseInfo::empty(database_id, shared_actor_infos),
791            independent_checkpoint_job_controls: Default::default(),
792            pending_independent_job_subscriptions_to_drop: Default::default(),
793        }
794    }
795
796    pub(crate) fn recovery(
797        database_id: DatabaseId,
798        term_id: String,
799        state: BarrierWorkerState,
800        committed_epoch: u64,
801        database_info: InflightDatabaseInfo,
802        independent_checkpoint_job_controls: HashMap<JobId, IndependentCheckpointJobControl>,
803    ) -> Self {
804        Self {
805            database_id,
806            term_id,
807            partial_graph_id: to_partial_graph_id(database_id, None),
808            state,
809            finishing_jobs_collector: BarrierItemCollector::new(false),
810            completing_barrier: None,
811            committed_epoch: Some(committed_epoch),
812            last_committed_barrier_time: None,
813            database_info,
814            independent_checkpoint_job_controls,
815            pending_independent_job_subscriptions_to_drop: Default::default(),
816        }
817    }
818
819    pub(super) fn term_id(&self) -> &str {
820        &self.term_id
821    }
822
823    pub(crate) fn is_valid_after_worker_err(&self, worker_id: WorkerId) -> bool {
824        !self.database_info.contains_worker(worker_id as _)
825            && self
826                .independent_checkpoint_job_controls
827                .values()
828                .all(|job| {
829                    job.fragment_infos()
830                        .map(|fragment_infos| {
831                            !InflightFragmentInfo::contains_worker(
832                                fragment_infos.values(),
833                                worker_id,
834                            )
835                        })
836                        .unwrap_or(true)
837                })
838    }
839
840    /// Enqueue a barrier command
841    fn enqueue_command(&mut self, epoch: EpochPair, independent_jobs_to_wait: HashSet<JobId>) {
842        let prev_epoch = epoch.prev;
843        tracing::trace!(prev_epoch, ?independent_jobs_to_wait, "enqueue command");
844        if !independent_jobs_to_wait.is_empty() {
845            self.finishing_jobs_collector
846                .enqueue(epoch, independent_jobs_to_wait, ());
847        }
848    }
849
850    /// Change the state of this `prev_epoch` to `Completed`. Return continuous nodes
851    /// with `Completed` starting from first node [`Completed`..`InFlight`) and remove them.
852    fn barrier_collected(
853        &mut self,
854        partial_graph_id: PartialGraphId,
855        collected_barrier: CollectedBarrier<'_>,
856        periodic_barriers: &mut PeriodicBarriers,
857    ) -> MetaResult<()> {
858        let prev_epoch = collected_barrier.epoch.prev;
859        tracing::trace!(
860            prev_epoch,
861            partial_graph_id = %partial_graph_id,
862            "barrier collected"
863        );
864        let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
865        assert_eq!(self.database_id, database_id);
866        if let Some(independent_job_id) = independent_job_id {
867            let job = self
868                .independent_checkpoint_job_controls
869                .get_mut(&independent_job_id)
870                .expect("should exist");
871            let should_force_checkpoint = job.collect(collected_barrier);
872            if should_force_checkpoint {
873                periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
874            }
875        }
876        Ok(())
877    }
878}
879
880impl DatabaseCheckpointControl {
881    fn collect_backfill_pinned_upstream_tables(&self) -> HashSet<TableId> {
882        self.independent_checkpoint_job_controls
883            .values()
884            .flat_map(|job| job.pinned_upstream_tables().iter().copied())
885            .collect()
886    }
887
888    fn collect_no_shuffle_fragment_relations_for_reschedule_check(
889        &self,
890    ) -> Vec<(FragmentId, FragmentId)> {
891        let mut no_shuffle_relations = Vec::new();
892        for fragment in self.database_info.fragment_infos() {
893            let downstream_fragment_id = fragment.fragment_id;
894            visit_stream_node_cont(&fragment.nodes, |node| {
895                if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
896                    && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
897                {
898                    no_shuffle_relations.push((merge.upstream_fragment_id, downstream_fragment_id));
899                }
900                true
901            });
902        }
903
904        for job in self.independent_checkpoint_job_controls.values() {
905            if let Some(fragment_infos) = job.fragment_infos() {
906                for fragment_info in fragment_infos.values() {
907                    let downstream_fragment_id = fragment_info.fragment_id;
908                    visit_stream_node_cont(&fragment_info.nodes, |node| {
909                        if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
910                            && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
911                        {
912                            no_shuffle_relations
913                                .push((merge.upstream_fragment_id, downstream_fragment_id));
914                        }
915                        true
916                    });
917                }
918            }
919        }
920        no_shuffle_relations
921    }
922
923    fn collect_reschedule_blocked_jobs_for_independent_jobs_inflight(
924        &self,
925    ) -> MetaResult<HashSet<JobId>> {
926        let mut initial_blocked_fragment_ids = HashSet::new();
927        for job in self.independent_checkpoint_job_controls.values() {
928            if let Some(fragment_infos) = job.fragment_infos() {
929                for fragment_info in fragment_infos.values() {
930                    if fragment_has_online_unreschedulable_scan(fragment_info) {
931                        initial_blocked_fragment_ids.insert(fragment_info.fragment_id);
932                        collect_fragment_upstream_fragment_ids(
933                            fragment_info,
934                            &mut initial_blocked_fragment_ids,
935                        );
936                    }
937                }
938            }
939        }
940
941        let mut blocked_fragment_ids = initial_blocked_fragment_ids.clone();
942        if !initial_blocked_fragment_ids.is_empty() {
943            let no_shuffle_relations =
944                self.collect_no_shuffle_fragment_relations_for_reschedule_check();
945            let (forward_edges, backward_edges) =
946                build_no_shuffle_fragment_graph_edges(no_shuffle_relations);
947            let initial_blocked_fragment_ids: Vec<_> =
948                initial_blocked_fragment_ids.iter().copied().collect();
949            for ensemble in find_no_shuffle_graphs(
950                &initial_blocked_fragment_ids,
951                &forward_edges,
952                &backward_edges,
953            )? {
954                blocked_fragment_ids.extend(ensemble.fragments());
955            }
956        }
957
958        let mut blocked_job_ids = HashSet::new();
959        blocked_job_ids.extend(
960            blocked_fragment_ids
961                .into_iter()
962                .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id)),
963        );
964        Ok(blocked_job_ids)
965    }
966
967    fn collect_reschedule_blocked_job_ids(
968        &self,
969        reschedules: &HashMap<FragmentId, Reschedule>,
970        fragment_actors: &HashMap<FragmentId, HashSet<ActorId>>,
971        blocked_job_ids: &HashSet<JobId>,
972    ) -> HashSet<JobId> {
973        let mut affected_fragment_ids: HashSet<FragmentId> = reschedules.keys().copied().collect();
974        affected_fragment_ids.extend(fragment_actors.keys().copied());
975        for reschedule in reschedules.values() {
976            affected_fragment_ids.extend(reschedule.downstream_fragment_ids.iter().copied());
977            affected_fragment_ids.extend(
978                reschedule
979                    .upstream_fragment_dispatcher_ids
980                    .iter()
981                    .map(|(fragment_id, _)| *fragment_id),
982            );
983        }
984
985        affected_fragment_ids
986            .into_iter()
987            .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id))
988            .filter(|job_id| blocked_job_ids.contains(job_id))
989            .collect()
990    }
991
992    fn next_complete_barrier_task(
993        &mut self,
994        periodic_barriers: &mut PeriodicBarriers,
995        partial_graph_manager: &mut PartialGraphManager,
996        task: &mut Option<CompleteBarrierTask>,
997        hummock_version_stats: &HummockVersionStats,
998    ) {
999        // `Vec::new` is a const fn, and do not have memory allocation, and therefore is lightweight enough
1000        let mut independent_jobs_task = vec![];
1001        let mut finished_jobs = Vec::new();
1002        let min_upstream_inflight_barrier = partial_graph_manager
1003            .first_inflight_barrier(self.partial_graph_id)
1004            .map(|epoch| epoch.prev);
1005        for (job_id, job) in &mut self.independent_checkpoint_job_controls {
1006            let Some(job) = job.ready_mut() else {
1007                continue;
1008            };
1009            match job {
1010                IndependentCheckpointJob::CreatingStreamingJob(creating_job) => {
1011                    if let Some((epoch, resps, info, is_finish_epoch)) = creating_job
1012                        .start_completing(partial_graph_manager, min_upstream_inflight_barrier)
1013                    {
1014                        let resps = resps.into_values().collect_vec();
1015                        if is_finish_epoch {
1016                            assert!(info.notifier.is_none());
1017                            finished_jobs.push((*job_id, epoch, resps));
1018                            continue;
1019                        };
1020                        independent_jobs_task.push((*job_id, epoch, resps, info));
1021                    }
1022                }
1023                IndependentCheckpointJob::BatchRefresh(batch_refresh_job) => {
1024                    if let Some((epoch, resps, info, tracking_job)) =
1025                        batch_refresh_job.start_completing(partial_graph_manager)
1026                    {
1027                        let resps = resps.into_values().collect_vec();
1028                        if let Some(tracking_job) = tracking_job {
1029                            let task = task.get_or_insert_default();
1030                            task.finished_jobs.push(tracking_job);
1031                        }
1032                        independent_jobs_task.push((*job_id, epoch, resps, info));
1033                    }
1034                }
1035                IndependentCheckpointJob::IcebergV3(iceberg_job) => {
1036                    if let Some((epoch, resps, info, tracking_job)) = iceberg_job
1037                        .start_completing(partial_graph_manager, min_upstream_inflight_barrier)
1038                    {
1039                        if let Some(tracking_job) = tracking_job {
1040                            let task = task.get_or_insert_default();
1041                            task.finished_jobs.push(tracking_job);
1042                        }
1043                        independent_jobs_task.push((
1044                            *job_id,
1045                            epoch,
1046                            resps.into_values().collect_vec(),
1047                            info,
1048                        ));
1049                    }
1050                }
1051            }
1052        }
1053        if !finished_jobs.is_empty() {
1054            partial_graph_manager.remove_partial_graphs(
1055                finished_jobs
1056                    .iter()
1057                    .map(|(job_id, ..)| to_partial_graph_id(self.database_id, Some(*job_id)))
1058                    .collect(),
1059            );
1060        }
1061        for (job_id, epoch, resps) in finished_jobs {
1062            debug!(epoch, %job_id, "finish creating job");
1063            // It's safe to remove the creating job, because on CompleteJobType::Finished,
1064            // all previous barriers have been collected and completed.
1065            // `finished_jobs` was populated above only from a Running creating job, and the
1066            // map cannot change between that scan and this removal.
1067            let Some(IndependentCheckpointJobControl::Running {
1068                job: IndependentCheckpointJob::CreatingStreamingJob(creating_streaming_job),
1069                ..
1070            }) = self.independent_checkpoint_job_controls.remove(&job_id)
1071            else {
1072                panic!("finished job {job_id} should be a creating streaming job");
1073            };
1074            let tracking_job = creating_streaming_job.into_tracking_job();
1075            self.finishing_jobs_collector
1076                .collect(epoch, job_id, (resps, tracking_job));
1077        }
1078        let mut observed_non_checkpoint = false;
1079        self.finishing_jobs_collector.advance_collected();
1080        let epoch_end_bound = self
1081            .finishing_jobs_collector
1082            .first_inflight_epoch()
1083            .map_or(Unbounded, |epoch| Excluded(epoch.prev));
1084        if let Some((epoch, resps, info)) = partial_graph_manager.start_completing(
1085            self.partial_graph_id,
1086            epoch_end_bound,
1087            |_, resps, post_collect_command| {
1088                observed_non_checkpoint = true;
1089                self.handle_refresh_table_info(task, &resps);
1090                self.database_info.apply_collected_command(
1091                    &post_collect_command,
1092                    &resps,
1093                    hummock_version_stats,
1094                );
1095            },
1096        ) {
1097            self.handle_refresh_table_info(task, &resps);
1098            self.database_info.apply_collected_command(
1099                &info.post_collect_command,
1100                &resps,
1101                hummock_version_stats,
1102            );
1103            let mut resps_to_commit = resps.into_values().collect_vec();
1104            let mut staging_commit_info = self.database_info.take_staging_commit_info();
1105            if let Some((_, finished_jobs, _)) =
1106                self.finishing_jobs_collector
1107                    .take_collected_if(|collected_epoch| {
1108                        assert!(epoch <= collected_epoch.prev);
1109                        epoch == collected_epoch.prev
1110                    })
1111            {
1112                finished_jobs
1113                    .into_iter()
1114                    .for_each(|(_, (resps, tracking_job))| {
1115                        resps_to_commit.extend(resps);
1116                        staging_commit_info.finished_jobs.push(tracking_job);
1117                    });
1118            }
1119            {
1120                let task = task.get_or_insert_default();
1121                Command::collect_commit_epoch_info(
1122                    &self.database_info,
1123                    &info,
1124                    task,
1125                    resps_to_commit,
1126                    self.collect_backfill_pinned_upstream_tables(),
1127                );
1128                self.completing_barrier = Some(info.barrier_info.epoch());
1129                task.finished_jobs.extend(staging_commit_info.finished_jobs);
1130                task.finished_cdc_table_backfill
1131                    .extend(staging_commit_info.finished_cdc_table_backfill);
1132                task.epoch_infos
1133                    .try_insert(self.partial_graph_id, info)
1134                    .expect("non duplicate");
1135                task.commit_info
1136                    .truncate_tables
1137                    .extend(staging_commit_info.table_ids_to_truncate);
1138            }
1139        } else if observed_non_checkpoint
1140            && self.database_info.has_pending_finished_jobs()
1141            && !partial_graph_manager.has_pending_checkpoint_barrier(self.partial_graph_id)
1142        {
1143            periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
1144        }
1145        if !independent_jobs_task.is_empty() {
1146            let task = task.get_or_insert_default();
1147            for (job_id, epoch, resps, info) in independent_jobs_task {
1148                collect_independent_job_commit_epoch_info(task, epoch, resps, &info);
1149                task.epoch_infos
1150                    .try_insert(to_partial_graph_id(self.database_id, Some(job_id)), info)
1151                    .expect("non duplicate");
1152            }
1153        }
1154    }
1155
1156    fn ack_completed(
1157        &mut self,
1158        partial_graph_manager: &mut PartialGraphManager,
1159        command_prev_epoch: Option<u64>,
1160        independent_job_epochs: Vec<(JobId, u64)>,
1161    ) {
1162        {
1163            if let Some(epoch) = self.completing_barrier.take() {
1164                assert_eq!(command_prev_epoch, Some(epoch.prev));
1165                self.committed_epoch = Some(epoch.prev);
1166                partial_graph_manager.ack_completed(self.partial_graph_id, epoch.prev);
1167                for job in self.independent_checkpoint_job_controls.values_mut() {
1168                    job.on_upstream_database_ack_completed(epoch.prev);
1169                }
1170                self.last_committed_barrier_time
1171                    .get_or_insert_with(|| {
1172                        GLOBAL_META_METRICS
1173                            .last_committed_barrier_time
1174                            .with_guarded_label_values(&[&self.database_id.to_string()])
1175                    })
1176                    .set(Epoch(epoch.curr).as_unix_secs() as i64);
1177            } else {
1178                assert_eq!(command_prev_epoch, None);
1179            };
1180            for (job_id, epoch) in independent_job_epochs {
1181                if let Some(job) = self.independent_checkpoint_job_controls.get_mut(&job_id) {
1182                    job.ack_completed(partial_graph_manager, epoch);
1183                }
1184                // If the job is not found, it was dropped and already removed
1185                // by `on_partial_graph_reset` while the completing task was running.
1186            }
1187        }
1188    }
1189
1190    fn handle_refresh_table_info(
1191        &self,
1192        task: &mut Option<CompleteBarrierTask>,
1193        resps: &HashMap<WorkerId, BarrierCompleteResponse>,
1194    ) {
1195        let list_finished_info = resps
1196            .values()
1197            .flat_map(|resp| resp.list_finished_sources.clone())
1198            .collect::<Vec<_>>();
1199        if !list_finished_info.is_empty() {
1200            let task = task.get_or_insert_default();
1201            task.list_finished_source_ids.extend(list_finished_info);
1202        }
1203
1204        let load_finished_info = resps
1205            .values()
1206            .flat_map(|resp| resp.load_finished_sources.clone())
1207            .collect::<Vec<_>>();
1208        if !load_finished_info.is_empty() {
1209            let task = task.get_or_insert_default();
1210            task.load_finished_source_ids.extend(load_finished_info);
1211        }
1212
1213        let refresh_finished_actors = resps
1214            .values()
1215            .flat_map(|resp| resp.refresh_finished_actors.clone())
1216            .collect::<Vec<_>>();
1217        if !refresh_finished_actors.is_empty() {
1218            let task = task.get_or_insert_default();
1219            task.refresh_finished_actors.extend(refresh_finished_actors);
1220        }
1221    }
1222}
1223
1224impl DatabaseCheckpointControl {
1225    /// Handle the new barrier from the scheduled queue and inject it.
1226    fn handle_new_barrier(
1227        &mut self,
1228        command: Option<(Command, Notifier)>,
1229        checkpoint: bool,
1230        span: tracing::Span,
1231        partial_graph_manager: &mut PartialGraphManager,
1232        hummock_version_stats: &HummockVersionStats,
1233        worker_nodes: &HashMap<WorkerId, WorkerNode>,
1234    ) -> MetaResult<()> {
1235        let curr_epoch = self.state.in_flight_prev_epoch().next();
1236
1237        let (mut command, notifier) = if let Some((command, notifier)) = command {
1238            (Some(command), Some(notifier))
1239        } else {
1240            (None, None)
1241        };
1242
1243        debug_assert!(
1244            !matches!(
1245                &command,
1246                Some(Command::RescheduleIntent {
1247                    reschedule_plan: None,
1248                    ..
1249                })
1250            ),
1251            "reschedule intent should be resolved before injection"
1252        );
1253
1254        let mut notifier_start = notifier.map(Notifier::start);
1255        if let Some(Command::DropStreamingJobs {
1256            streaming_job_ids, ..
1257        }) = &mut command
1258        {
1259            streaming_job_ids.retain(|job_id| {
1260                let Some(job) = self.independent_checkpoint_job_controls.get_mut(job_id) else {
1261                    return true;
1262                };
1263                !job.drop(notifier_start.as_mut(), partial_graph_manager)
1264            });
1265            if streaming_job_ids.is_empty() {
1266                if let Some(notifier) = notifier_start {
1267                    notifier.started();
1268                }
1269                return Ok(());
1270            }
1271        }
1272
1273        if let Some(Command::RescheduleIntent {
1274            reschedule_plan: Some(reschedule_plan),
1275            ..
1276        }) = &command
1277            && !self.independent_checkpoint_job_controls.is_empty()
1278        {
1279            let blocked_job_ids =
1280                self.collect_reschedule_blocked_jobs_for_independent_jobs_inflight()?;
1281            let blocked_reschedule_job_ids = self.collect_reschedule_blocked_job_ids(
1282                &reschedule_plan.reschedules,
1283                &reschedule_plan.fragment_actors,
1284                &blocked_job_ids,
1285            );
1286            if !blocked_reschedule_job_ids.is_empty() {
1287                warn!(
1288                    blocked_reschedule_job_ids = ?blocked_reschedule_job_ids,
1289                    "reject reschedule fragments related to creating unreschedulable backfill jobs"
1290                );
1291                if let Some(notifier) = notifier_start {
1292                    notifier.notify_start_failed(
1293                        anyhow!(
1294                            "cannot reschedule jobs {:?} when creating jobs with unreschedulable backfill fragments",
1295                            blocked_reschedule_job_ids
1296                        )
1297                            .into(),
1298                    );
1299                }
1300                return Ok(());
1301            }
1302        }
1303
1304        if !matches!(&command, Some(Command::CreateStreamingJob { .. }))
1305            && self.database_info.is_empty()
1306            && self
1307                .independent_checkpoint_job_controls
1308                .values()
1309                .all(|job| job.running().is_none())
1310        {
1311            // Drop the guard to remove the metric series of this database.
1312            self.last_committed_barrier_time = None;
1313            // skip the command when there is nothing to do with the barrier
1314            if let Some(notifier) = notifier_start {
1315                notifier.started();
1316            }
1317            return Ok(());
1318        };
1319
1320        if let Some(Command::CreateStreamingJob {
1321            job_type: CreateStreamingJobType::Independent { .. },
1322            ..
1323        }) = &command
1324            && self.state.is_paused()
1325        {
1326            warn!("cannot create streaming job with snapshot backfill when paused");
1327            if let Some(notifier) = notifier_start {
1328                notifier.notify_start_failed(
1329                    anyhow!("cannot create streaming job with snapshot backfill when paused",)
1330                        .into(),
1331                );
1332            }
1333            return Ok(());
1334        }
1335
1336        let barrier_info = self.state.next_barrier_info(checkpoint, curr_epoch);
1337        // Tracing related stuff
1338        barrier_info.prev_epoch.span().in_scope(|| {
1339            tracing::info!(target: "rw_tracing", epoch = barrier_info.curr_epoch(), "new barrier enqueued");
1340        });
1341        span.record("epoch", barrier_info.curr_epoch());
1342
1343        let epoch = barrier_info.epoch();
1344        let ApplyCommandInfo { jobs_to_wait } = match self.apply_command(
1345            command,
1346            &mut notifier_start,
1347            barrier_info,
1348            partial_graph_manager,
1349            hummock_version_stats,
1350            worker_nodes,
1351        ) {
1352            Ok(info) => {
1353                assert!(notifier_start.is_none());
1354                info
1355            }
1356            Err(err) => {
1357                if let Some(notifier) = notifier_start {
1358                    notifier.notify_start_failed(err.clone());
1359                }
1360                fail_point!("inject_barrier_err_success");
1361                return Err(err);
1362            }
1363        };
1364
1365        // Record the in-flight barrier.
1366        self.enqueue_command(epoch, jobs_to_wait);
1367
1368        Ok(())
1369    }
1370
1371    // ── Batch refresh trigger helpers ────────────────────────────────────────
1372
1373    /// Get the last committed epoch for a batch refresh job.
1374    pub(crate) fn get_batch_refresh_trigger_info(&self, job_id: JobId) -> u64 {
1375        // `CheckpointControl::next_event` emits a trigger only after matching this job as a
1376        // Running batch-refresh job. The worker reads this value before awaiting any other work.
1377        let job = self
1378            .independent_checkpoint_job_controls
1379            .get(&job_id)
1380            .expect("batch refresh job should exist")
1381            .running()
1382            .expect("batch refresh job should be running");
1383        match job {
1384            IndependentCheckpointJob::BatchRefresh(br_job) => br_job
1385                .last_committed_epoch()
1386                .expect("idle job must have a last_committed_epoch"),
1387            _ => panic!("job {} should be a batch refresh job", job_id),
1388        }
1389    }
1390
1391    /// Whether the batch refresh job already has its cached context populated.
1392    /// Start a batch refresh logstore consumption run.
1393    /// Returns true if a run was started, false if no log epochs to consume.
1394    pub(crate) fn start_batch_refresh_run(
1395        &mut self,
1396        job_id: JobId,
1397        context: &BatchRefreshJobTriggerContext,
1398        worker_nodes: &HashMap<WorkerId, WorkerNode>,
1399        actor_id_counter: &AtomicU32,
1400        partial_graph_manager: &mut PartialGraphManager,
1401    ) -> MetaResult<bool> {
1402        // The global barrier worker handles the trigger serially. Although metadata loading is
1403        // asynchronous, it does not process a drop/reset event before reaching this call.
1404        let term_id = self.term_id.as_str();
1405        let job = self
1406            .independent_checkpoint_job_controls
1407            .get_mut(&job_id)
1408            .expect("batch refresh job should exist")
1409            .running_mut()
1410            .expect("batch refresh job should be running");
1411        match job {
1412            IndependentCheckpointJob::BatchRefresh(br_job) => br_job.start_refresh_run(
1413                context,
1414                worker_nodes,
1415                actor_id_counter,
1416                term_id,
1417                partial_graph_manager,
1418            ),
1419            _ => panic!("job {} should be a batch refresh job", job_id),
1420        }
1421    }
1422}