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