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