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 does not exist while handling the command");
344                            notifier.notify_start_failed(anyhow!("database {database_id} does not exist while 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    fn collect_backfill_pinned_upstream_tables(&self) -> HashSet<TableId> {
857        self.independent_checkpoint_job_controls
858            .values()
859            .flat_map(|job| job.pinned_upstream_tables())
860            .collect()
861    }
862
863    fn collect_no_shuffle_fragment_relations_for_reschedule_check(
864        &self,
865    ) -> Vec<(FragmentId, FragmentId)> {
866        let mut no_shuffle_relations = Vec::new();
867        for fragment in self.database_info.fragment_infos() {
868            let downstream_fragment_id = fragment.fragment_id;
869            visit_stream_node_cont(&fragment.nodes, |node| {
870                if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
871                    && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
872                {
873                    no_shuffle_relations.push((merge.upstream_fragment_id, downstream_fragment_id));
874                }
875                true
876            });
877        }
878
879        for job in self.independent_checkpoint_job_controls.values() {
880            if let Some(fragment_infos) = job.fragment_infos() {
881                for fragment_info in fragment_infos.values() {
882                    let downstream_fragment_id = fragment_info.fragment_id;
883                    visit_stream_node_cont(&fragment_info.nodes, |node| {
884                        if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
885                            && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
886                        {
887                            no_shuffle_relations
888                                .push((merge.upstream_fragment_id, downstream_fragment_id));
889                        }
890                        true
891                    });
892                }
893            }
894        }
895        no_shuffle_relations
896    }
897
898    fn collect_reschedule_blocked_jobs_for_independent_jobs_inflight(
899        &self,
900    ) -> MetaResult<HashSet<JobId>> {
901        let mut initial_blocked_fragment_ids = HashSet::new();
902        for job in self.independent_checkpoint_job_controls.values() {
903            if let Some(fragment_infos) = job.fragment_infos() {
904                for fragment_info in fragment_infos.values() {
905                    if fragment_has_online_unreschedulable_scan(fragment_info) {
906                        initial_blocked_fragment_ids.insert(fragment_info.fragment_id);
907                        collect_fragment_upstream_fragment_ids(
908                            fragment_info,
909                            &mut initial_blocked_fragment_ids,
910                        );
911                    }
912                }
913            }
914        }
915
916        let mut blocked_fragment_ids = initial_blocked_fragment_ids.clone();
917        if !initial_blocked_fragment_ids.is_empty() {
918            let no_shuffle_relations =
919                self.collect_no_shuffle_fragment_relations_for_reschedule_check();
920            let (forward_edges, backward_edges) =
921                build_no_shuffle_fragment_graph_edges(no_shuffle_relations);
922            let initial_blocked_fragment_ids: Vec<_> =
923                initial_blocked_fragment_ids.iter().copied().collect();
924            for ensemble in find_no_shuffle_graphs(
925                &initial_blocked_fragment_ids,
926                &forward_edges,
927                &backward_edges,
928            )? {
929                blocked_fragment_ids.extend(ensemble.fragments());
930            }
931        }
932
933        let mut blocked_job_ids = HashSet::new();
934        blocked_job_ids.extend(
935            blocked_fragment_ids
936                .into_iter()
937                .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id)),
938        );
939        Ok(blocked_job_ids)
940    }
941
942    fn collect_reschedule_blocked_job_ids(
943        &self,
944        reschedules: &HashMap<FragmentId, Reschedule>,
945        fragment_actors: &HashMap<FragmentId, HashSet<ActorId>>,
946        blocked_job_ids: &HashSet<JobId>,
947    ) -> HashSet<JobId> {
948        let mut affected_fragment_ids: HashSet<FragmentId> = reschedules.keys().copied().collect();
949        affected_fragment_ids.extend(fragment_actors.keys().copied());
950        for reschedule in reschedules.values() {
951            affected_fragment_ids.extend(reschedule.downstream_fragment_ids.iter().copied());
952            affected_fragment_ids.extend(
953                reschedule
954                    .upstream_fragment_dispatcher_ids
955                    .iter()
956                    .map(|(fragment_id, _)| *fragment_id),
957            );
958        }
959
960        affected_fragment_ids
961            .into_iter()
962            .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id))
963            .filter(|job_id| blocked_job_ids.contains(job_id))
964            .collect()
965    }
966
967    fn next_complete_barrier_task(
968        &mut self,
969        periodic_barriers: &mut PeriodicBarriers,
970        partial_graph_manager: &mut PartialGraphManager,
971        task: &mut Option<CompleteBarrierTask>,
972        hummock_version_stats: &HummockVersionStats,
973    ) {
974        // `Vec::new` is a const fn, and do not have memory allocation, and therefore is lightweight enough
975        let mut independent_jobs_task = vec![];
976        if let Some(committed_epoch) = self.committed_epoch {
977            // `Vec::new` is a const fn, and do not have memory allocation, and therefore is lightweight enough
978            let mut finished_jobs = Vec::new();
979            let min_upstream_inflight_barrier = partial_graph_manager
980                .first_inflight_barrier(self.partial_graph_id)
981                .map(|epoch| epoch.prev);
982            for (job_id, job) in &mut self.independent_checkpoint_job_controls {
983                match job {
984                    IndependentCheckpointJobControl::CreatingStreamingJob(creating_job) => {
985                        if let Some((epoch, resps, info, is_finish_epoch)) = creating_job
986                            .start_completing(
987                                partial_graph_manager,
988                                min_upstream_inflight_barrier,
989                                committed_epoch,
990                            )
991                        {
992                            let resps = resps.into_values().collect_vec();
993                            if is_finish_epoch {
994                                assert!(info.notifier.is_none());
995                                finished_jobs.push((*job_id, epoch, resps));
996                                continue;
997                            };
998                            independent_jobs_task.push((*job_id, epoch, resps, info));
999                        }
1000                    }
1001                    IndependentCheckpointJobControl::BatchRefresh(batch_refresh_job) => {
1002                        if let Some((epoch, resps, info, tracking_job)) = batch_refresh_job
1003                            .start_completing(partial_graph_manager, committed_epoch)
1004                        {
1005                            let resps = resps.into_values().collect_vec();
1006                            if let Some(tracking_job) = tracking_job {
1007                                let task = task.get_or_insert_default();
1008                                task.finished_jobs.push(tracking_job);
1009                            }
1010                            independent_jobs_task.push((*job_id, epoch, resps, info));
1011                        }
1012                    }
1013                }
1014            }
1015            if !finished_jobs.is_empty() {
1016                partial_graph_manager.remove_partial_graphs(
1017                    finished_jobs
1018                        .iter()
1019                        .map(|(job_id, ..)| to_partial_graph_id(self.database_id, Some(*job_id)))
1020                        .collect(),
1021                );
1022            }
1023            for (job_id, epoch, resps) in finished_jobs {
1024                debug!(epoch, %job_id, "finish creating job");
1025                // It's safe to remove the creating job, because on CompleteJobType::Finished,
1026                // all previous barriers have been collected and completed.
1027                let Some(IndependentCheckpointJobControl::CreatingStreamingJob(
1028                    creating_streaming_job,
1029                )) = self.independent_checkpoint_job_controls.remove(&job_id)
1030                else {
1031                    panic!("finished job {job_id} should be a creating streaming job");
1032                };
1033                let tracking_job = creating_streaming_job.into_tracking_job();
1034                self.finishing_jobs_collector
1035                    .collect(epoch, job_id, (resps, tracking_job));
1036            }
1037        }
1038        let mut observed_non_checkpoint = false;
1039        self.finishing_jobs_collector.advance_collected();
1040        let epoch_end_bound = self
1041            .finishing_jobs_collector
1042            .first_inflight_epoch()
1043            .map_or(Unbounded, |epoch| Excluded(epoch.prev));
1044        if let Some((epoch, resps, info)) = partial_graph_manager.start_completing(
1045            self.partial_graph_id,
1046            epoch_end_bound,
1047            |_, resps, post_collect_command| {
1048                observed_non_checkpoint = true;
1049                self.handle_refresh_table_info(task, &resps);
1050                self.database_info.apply_collected_command(
1051                    &post_collect_command,
1052                    &resps,
1053                    hummock_version_stats,
1054                );
1055            },
1056        ) {
1057            self.handle_refresh_table_info(task, &resps);
1058            self.database_info.apply_collected_command(
1059                &info.post_collect_command,
1060                &resps,
1061                hummock_version_stats,
1062            );
1063            let mut resps_to_commit = resps.into_values().collect_vec();
1064            let mut staging_commit_info = self.database_info.take_staging_commit_info();
1065            if let Some((_, finished_jobs, _)) =
1066                self.finishing_jobs_collector
1067                    .take_collected_if(|collected_epoch| {
1068                        assert!(epoch <= collected_epoch.prev);
1069                        epoch == collected_epoch.prev
1070                    })
1071            {
1072                finished_jobs
1073                    .into_iter()
1074                    .for_each(|(_, (resps, tracking_job))| {
1075                        resps_to_commit.extend(resps);
1076                        staging_commit_info.finished_jobs.push(tracking_job);
1077                    });
1078            }
1079            {
1080                let task = task.get_or_insert_default();
1081                Command::collect_commit_epoch_info(
1082                    &self.database_info,
1083                    &info,
1084                    task,
1085                    resps_to_commit,
1086                    self.collect_backfill_pinned_upstream_tables(),
1087                );
1088                self.completing_barrier = Some(info.barrier_info.epoch());
1089                task.finished_jobs.extend(staging_commit_info.finished_jobs);
1090                task.finished_cdc_table_backfill
1091                    .extend(staging_commit_info.finished_cdc_table_backfill);
1092                task.epoch_infos
1093                    .try_insert(self.partial_graph_id, info)
1094                    .expect("non duplicate");
1095                task.commit_info
1096                    .truncate_tables
1097                    .extend(staging_commit_info.table_ids_to_truncate);
1098            }
1099        } else if observed_non_checkpoint
1100            && self.database_info.has_pending_finished_jobs()
1101            && !partial_graph_manager.has_pending_checkpoint_barrier(self.partial_graph_id)
1102        {
1103            periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
1104        }
1105        if !independent_jobs_task.is_empty() {
1106            let task = task.get_or_insert_default();
1107            for (job_id, epoch, resps, info) in independent_jobs_task {
1108                collect_independent_job_commit_epoch_info(task, epoch, resps, &info);
1109                task.epoch_infos
1110                    .try_insert(to_partial_graph_id(self.database_id, Some(job_id)), info)
1111                    .expect("non duplicate");
1112            }
1113        }
1114    }
1115
1116    fn ack_completed(
1117        &mut self,
1118        partial_graph_manager: &mut PartialGraphManager,
1119        command_prev_epoch: Option<u64>,
1120        independent_job_epochs: Vec<(JobId, u64)>,
1121    ) {
1122        {
1123            if let Some(epoch) = self.completing_barrier.take() {
1124                assert_eq!(command_prev_epoch, Some(epoch.prev));
1125                self.committed_epoch = Some(epoch.prev);
1126                partial_graph_manager.ack_completed(self.partial_graph_id, epoch.prev);
1127                self.last_committed_barrier_time
1128                    .get_or_insert_with(|| {
1129                        GLOBAL_META_METRICS
1130                            .last_committed_barrier_time
1131                            .with_guarded_label_values(&[&self.database_id.to_string()])
1132                    })
1133                    .set(Epoch(epoch.curr).as_unix_secs() as i64);
1134            } else {
1135                assert_eq!(command_prev_epoch, None);
1136            };
1137            for (job_id, epoch) in independent_job_epochs {
1138                if let Some(job) = self.independent_checkpoint_job_controls.get_mut(&job_id) {
1139                    job.ack_completed(partial_graph_manager, epoch);
1140                }
1141                // If the job is not found, it was dropped and already removed
1142                // by `on_partial_graph_reset` while the completing task was running.
1143            }
1144        }
1145    }
1146
1147    fn handle_refresh_table_info(
1148        &self,
1149        task: &mut Option<CompleteBarrierTask>,
1150        resps: &HashMap<WorkerId, BarrierCompleteResponse>,
1151    ) {
1152        let list_finished_info = resps
1153            .values()
1154            .flat_map(|resp| resp.list_finished_sources.clone())
1155            .collect::<Vec<_>>();
1156        if !list_finished_info.is_empty() {
1157            let task = task.get_or_insert_default();
1158            task.list_finished_source_ids.extend(list_finished_info);
1159        }
1160
1161        let load_finished_info = resps
1162            .values()
1163            .flat_map(|resp| resp.load_finished_sources.clone())
1164            .collect::<Vec<_>>();
1165        if !load_finished_info.is_empty() {
1166            let task = task.get_or_insert_default();
1167            task.load_finished_source_ids.extend(load_finished_info);
1168        }
1169
1170        let refresh_finished_table_ids: Vec<JobId> = resps
1171            .values()
1172            .flat_map(|resp| {
1173                resp.refresh_finished_tables
1174                    .iter()
1175                    .map(|table_id| table_id.as_job_id())
1176            })
1177            .collect::<Vec<_>>();
1178        if !refresh_finished_table_ids.is_empty() {
1179            let task = task.get_or_insert_default();
1180            task.refresh_finished_table_job_ids
1181                .extend(refresh_finished_table_ids);
1182        }
1183    }
1184}
1185
1186impl DatabaseCheckpointControl {
1187    /// Handle the new barrier from the scheduled queue and inject it.
1188    fn handle_new_barrier(
1189        &mut self,
1190        command: Option<(Command, Notifier)>,
1191        checkpoint: bool,
1192        span: tracing::Span,
1193        partial_graph_manager: &mut PartialGraphManager,
1194        hummock_version_stats: &HummockVersionStats,
1195        worker_nodes: &HashMap<WorkerId, WorkerNode>,
1196    ) -> MetaResult<()> {
1197        let curr_epoch = self.state.in_flight_prev_epoch().next();
1198
1199        let (mut command, notifier) = if let Some((command, notifier)) = command {
1200            (Some(command), Some(notifier))
1201        } else {
1202            (None, None)
1203        };
1204
1205        debug_assert!(
1206            !matches!(
1207                &command,
1208                Some(Command::RescheduleIntent {
1209                    reschedule_plan: None,
1210                    ..
1211                })
1212            ),
1213            "reschedule intent should be resolved before injection"
1214        );
1215
1216        let mut notifier_start = notifier.map(Notifier::start);
1217        if let Some(Command::DropStreamingJobs {
1218            streaming_job_ids, ..
1219        }) = &mut command
1220        {
1221            streaming_job_ids.retain(|job_id| {
1222                let Some(job) = self.independent_checkpoint_job_controls.get_mut(job_id) else {
1223                    return true;
1224                };
1225                !job.drop(notifier_start.as_mut(), partial_graph_manager)
1226            });
1227            if streaming_job_ids.is_empty() {
1228                if let Some(notifier) = notifier_start {
1229                    notifier.started();
1230                }
1231                return Ok(());
1232            }
1233        }
1234
1235        if let Some(Command::RescheduleIntent {
1236            reschedule_plan: Some(reschedule_plan),
1237            ..
1238        }) = &command
1239            && !self.independent_checkpoint_job_controls.is_empty()
1240        {
1241            let blocked_job_ids =
1242                self.collect_reschedule_blocked_jobs_for_independent_jobs_inflight()?;
1243            let blocked_reschedule_job_ids = self.collect_reschedule_blocked_job_ids(
1244                &reschedule_plan.reschedules,
1245                &reschedule_plan.fragment_actors,
1246                &blocked_job_ids,
1247            );
1248            if !blocked_reschedule_job_ids.is_empty() {
1249                warn!(
1250                    blocked_reschedule_job_ids = ?blocked_reschedule_job_ids,
1251                    "reject reschedule fragments related to creating unreschedulable backfill jobs"
1252                );
1253                if let Some(notifier) = notifier_start {
1254                    notifier.notify_start_failed(
1255                        anyhow!(
1256                            "cannot reschedule jobs {:?} when creating jobs with unreschedulable backfill fragments",
1257                            blocked_reschedule_job_ids
1258                        )
1259                            .into(),
1260                    );
1261                }
1262                return Ok(());
1263            }
1264        }
1265
1266        if !matches!(&command, Some(Command::CreateStreamingJob { .. }))
1267            && self.database_info.is_empty()
1268        {
1269            assert!(
1270                self.independent_checkpoint_job_controls.is_empty(),
1271                "should not have snapshot backfill job when there is no normal job in database"
1272            );
1273            // Drop the guard to remove the metric series of this database.
1274            self.last_committed_barrier_time = None;
1275            // skip the command when there is nothing to do with the barrier
1276            if let Some(notifier) = notifier_start {
1277                notifier.started();
1278            }
1279            return Ok(());
1280        };
1281
1282        if let Some(Command::CreateStreamingJob {
1283            job_type:
1284                CreateStreamingJobType::SnapshotBackfill { .. }
1285                | CreateStreamingJobType::BatchRefresh(_),
1286            ..
1287        }) = &command
1288            && self.state.is_paused()
1289        {
1290            warn!("cannot create streaming job with snapshot backfill when paused");
1291            if let Some(notifier) = notifier_start {
1292                notifier.notify_start_failed(
1293                    anyhow!("cannot create streaming job with snapshot backfill when paused",)
1294                        .into(),
1295                );
1296            }
1297            return Ok(());
1298        }
1299
1300        let barrier_info = self.state.next_barrier_info(checkpoint, curr_epoch);
1301        // Tracing related stuff
1302        barrier_info.prev_epoch.span().in_scope(|| {
1303            tracing::info!(target: "rw_tracing", epoch = barrier_info.curr_epoch(), "new barrier enqueued");
1304        });
1305        span.record("epoch", barrier_info.curr_epoch());
1306
1307        let epoch = barrier_info.epoch();
1308        let ApplyCommandInfo { jobs_to_wait } = match self.apply_command(
1309            command,
1310            &mut notifier_start,
1311            barrier_info,
1312            partial_graph_manager,
1313            hummock_version_stats,
1314            worker_nodes,
1315        ) {
1316            Ok(info) => {
1317                assert!(notifier_start.is_none());
1318                info
1319            }
1320            Err(err) => {
1321                if let Some(notifier) = notifier_start {
1322                    notifier.notify_start_failed(err.clone());
1323                }
1324                fail_point!("inject_barrier_err_success");
1325                return Err(err);
1326            }
1327        };
1328
1329        // Record the in-flight barrier.
1330        self.enqueue_command(epoch, jobs_to_wait);
1331
1332        Ok(())
1333    }
1334
1335    // ── Batch refresh trigger helpers ────────────────────────────────────────
1336
1337    /// Get the last committed epoch for a batch refresh job.
1338    pub(crate) fn get_batch_refresh_trigger_info(&self, job_id: JobId) -> u64 {
1339        let job = self
1340            .independent_checkpoint_job_controls
1341            .get(&job_id)
1342            .expect("batch refresh job should exist");
1343        match job {
1344            IndependentCheckpointJobControl::BatchRefresh(br_job) => br_job
1345                .last_committed_epoch()
1346                .expect("idle job must have a last_committed_epoch"),
1347            _ => panic!("job {} should be a batch refresh job", job_id),
1348        }
1349    }
1350
1351    /// Whether the batch refresh job already has its cached context populated.
1352    /// Start a batch refresh logstore consumption run.
1353    /// Returns true if a run was started, false if no log epochs to consume.
1354    pub(crate) fn start_batch_refresh_run(
1355        &mut self,
1356        job_id: JobId,
1357        context: &BatchRefreshJobTriggerContext,
1358        worker_nodes: &HashMap<WorkerId, WorkerNode>,
1359        actor_id_counter: &AtomicU32,
1360        partial_graph_manager: &mut PartialGraphManager,
1361    ) -> MetaResult<bool> {
1362        let job = self
1363            .independent_checkpoint_job_controls
1364            .get_mut(&job_id)
1365            .expect("batch refresh job should exist");
1366        match job {
1367            IndependentCheckpointJobControl::BatchRefresh(br_job) => br_job.start_refresh_run(
1368                context,
1369                worker_nodes,
1370                actor_id_counter,
1371                partial_graph_manager,
1372            ),
1373            _ => panic!("job {} should be a batch refresh job", job_id),
1374        }
1375    }
1376}