Skip to main content

risingwave_meta/barrier/checkpoint/
control.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::future::{Future, poll_fn};
18use std::ops::Bound::{Excluded, Unbounded};
19use std::sync::atomic::AtomicU32;
20use std::task::Poll;
21
22use anyhow::anyhow;
23use fail::fail_point;
24use itertools::Itertools;
25use risingwave_common::catalog::{DatabaseId, TableId};
26use risingwave_common::id::JobId;
27use risingwave_common::metrics::{LabelGuardedHistogram, LabelGuardedIntGauge};
28use risingwave_common::util::epoch::{Epoch, EpochPair};
29use risingwave_common::util::stream_graph_visitor::visit_stream_node_cont;
30use risingwave_meta_model::WorkerId;
31use risingwave_pb::common::WorkerNode;
32use risingwave_pb::hummock::HummockVersionStats;
33use risingwave_pb::id::{FragmentId, PartialGraphId};
34use risingwave_pb::stream_plan::stream_node::NodeBody;
35use risingwave_pb::stream_plan::{DispatcherType as PbDispatcherType, PbSubscriptionUpstreamInfo};
36use risingwave_pb::stream_service::BarrierCompleteResponse;
37use risingwave_pb::stream_service::streaming_control_stream_response::ResetPartialGraphResponse;
38use tracing::{debug, warn};
39use uuid::Uuid;
40
41use crate::barrier::cdc_progress::CdcProgress;
42use crate::barrier::checkpoint::independent_job::{
43    BatchRefreshJobTriggerContext, IndependentCheckpointJob, IndependentCheckpointJobControl,
44};
45use crate::barrier::checkpoint::recovery::{
46    DatabaseRecoveringState, DatabaseStatusAction, EnterInitializing, EnterRunning,
47    RecoveringStateAction,
48};
49use crate::barrier::checkpoint::state::{ApplyCommandInfo, BarrierWorkerState};
50use crate::barrier::complete_task::{BarrierCompleteOutput, CompleteBarrierTask};
51use crate::barrier::info::{InflightDatabaseInfo, SharedActorInfos};
52use crate::barrier::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        // This is called synchronously after `start_batch_refresh_run` returns `true`, with no
516        // barrier-worker event in between that could transition the job to Resetting.
517        let database = self
518            .databases
519            .get_mut(&database_id)
520            .and_then(|s| s.running_state_mut())
521            .expect("database should be running");
522        let br_job = match database
523            .independent_checkpoint_job_controls
524            .get(&job_id)
525            .expect("job should exist")
526            .running()
527            .expect("job should be running")
528        {
529            IndependentCheckpointJob::BatchRefresh(job) => job,
530            _ => panic!("expected batch refresh job"),
531        };
532        if let Some(fragment_infos) = br_job.fragment_infos() {
533            database
534                .database_info
535                .shared_actor_infos
536                .upsert(database_id, fragment_infos.values().map(|f| (f, job_id)));
537        }
538    }
539}
540
541pub(crate) enum CheckpointControlEvent<'a> {
542    EnteringInitializing(DatabaseStatusAction<'a, EnterInitializing>),
543    EnteringRunning(DatabaseStatusAction<'a, EnterRunning>),
544    /// A batch refresh job is idle and its upstream has advanced past the refresh interval.
545    /// Carries owned values so the async handler can call into context without borrowing self.
546    BatchRefreshTrigger {
547        database_id: DatabaseId,
548        job_id: JobId,
549    },
550}
551
552impl CheckpointControl {
553    pub(crate) fn on_partial_graph_reset(
554        &mut self,
555        partial_graph_id: PartialGraphId,
556        reset_resps: HashMap<WorkerId, ResetPartialGraphResponse>,
557    ) {
558        let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
559        match self.databases.get_mut(&database_id).expect("should exist") {
560            DatabaseCheckpointControlStatus::Running(database) => {
561                if let Some(independent_job_id) = independent_job_id {
562                    match database
563                        .independent_checkpoint_job_controls
564                        .remove(&independent_job_id)
565                    {
566                        Some(independent_job) => {
567                            database
568                                .pending_independent_job_subscriptions_to_drop
569                                .extend(independent_job.on_partial_graph_reset());
570                        }
571                        None => {
572                            if cfg!(debug_assertions) {
573                                panic!(
574                                    "receive reset partial graph resp on non-existing independent job {independent_job_id} in database {database_id}"
575                                )
576                            }
577                            warn!(
578                                %database_id,
579                                %independent_job_id,
580                                "ignore reset partial graph resp on non-existing independent job on running database"
581                            );
582                        }
583                    }
584                } else {
585                    unreachable!("should not receive reset database resp when database running")
586                }
587            }
588            DatabaseCheckpointControlStatus::Recovering(state) => {
589                state.on_partial_graph_reset(partial_graph_id, reset_resps);
590            }
591        }
592    }
593
594    pub(crate) fn on_partial_graph_initialized(
595        &mut self,
596        partial_graph_id: PartialGraphId,
597        partial_graph_manager: &mut PartialGraphManager,
598    ) -> MetaResult<()> {
599        let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
600        match self.databases.get_mut(&database_id).expect("should exist") {
601            DatabaseCheckpointControlStatus::Running(database) => {
602                let Some(independent_job_id) = independent_job_id else {
603                    unreachable!("database partial graph should not initialize when running")
604                };
605                let job = database
606                    .independent_checkpoint_job_controls
607                    .get_mut(&independent_job_id)
608                    .expect("independent job should exist");
609                match job.running_mut().expect("job should be running") {
610                    IndependentCheckpointJob::BatchRefresh(job) => {
611                        job.on_log_store_initialized(partial_graph_manager)
612                    }
613                    IndependentCheckpointJob::CreatingStreamingJob(_) => {
614                        unreachable!("creating streaming job should not initialize when running")
615                    }
616                }
617            }
618            DatabaseCheckpointControlStatus::Recovering(state) => {
619                state.partial_graph_initialized(partial_graph_id);
620                Ok(())
621            }
622        }
623    }
624
625    pub(crate) fn next_event(
626        &mut self,
627    ) -> impl Future<Output = CheckpointControlEvent<'_>> + Send + '_ {
628        let mut this = Some(self);
629        poll_fn(move |cx| {
630            let Some(this_mut) = this.as_mut() else {
631                unreachable!("should not be polled after poll ready")
632            };
633            for (&database_id, database_status) in &mut this_mut.databases {
634                match database_status {
635                    DatabaseCheckpointControlStatus::Running(database) => {
636                        // Check if any idle batch refresh job should start a refresh run.
637                        if let Some(committed_epoch) = database.committed_epoch {
638                            for (job_id, job) in &database.independent_checkpoint_job_controls {
639                                if let Some(IndependentCheckpointJob::BatchRefresh(br_job)) =
640                                    job.running()
641                                    && br_job.should_start_refresh(committed_epoch)
642                                {
643                                    let job_id = *job_id;
644                                    let _ = this.take().expect("checked Some");
645                                    return Poll::Ready(
646                                        CheckpointControlEvent::BatchRefreshTrigger {
647                                            database_id,
648                                            job_id,
649                                        },
650                                    );
651                                }
652                            }
653                        }
654                    }
655                    DatabaseCheckpointControlStatus::Recovering(state) => {
656                        let poll_result = state.poll_next_event(cx);
657                        if let Poll::Ready(action) = poll_result {
658                            let this = this.take().expect("checked Some");
659                            return Poll::Ready(match action {
660                                RecoveringStateAction::EnterInitializing(reset_workers) => {
661                                    CheckpointControlEvent::EnteringInitializing(
662                                        this.new_database_status_action(
663                                            database_id,
664                                            EnterInitializing(reset_workers),
665                                        ),
666                                    )
667                                }
668                                RecoveringStateAction::EnterRunning => {
669                                    CheckpointControlEvent::EnteringRunning(
670                                        this.new_database_status_action(database_id, EnterRunning),
671                                    )
672                                }
673                            });
674                        }
675                    }
676                }
677            }
678            Poll::Pending
679        })
680    }
681}
682
683pub(crate) enum DatabaseCheckpointControlStatus {
684    Running(DatabaseCheckpointControl),
685    Recovering(DatabaseRecoveringState),
686}
687
688impl DatabaseCheckpointControlStatus {
689    fn running_state(&self) -> Option<&DatabaseCheckpointControl> {
690        match self {
691            DatabaseCheckpointControlStatus::Running(state) => Some(state),
692            DatabaseCheckpointControlStatus::Recovering(_) => None,
693        }
694    }
695
696    fn running_state_mut(&mut self) -> Option<&mut DatabaseCheckpointControl> {
697        match self {
698            DatabaseCheckpointControlStatus::Running(state) => Some(state),
699            DatabaseCheckpointControlStatus::Recovering(_) => None,
700        }
701    }
702
703    fn expect_running(&mut self, reason: &'static str) -> &mut DatabaseCheckpointControl {
704        match self {
705            DatabaseCheckpointControlStatus::Running(state) => state,
706            DatabaseCheckpointControlStatus::Recovering(_) => {
707                panic!("should be at running: {}", reason)
708            }
709        }
710    }
711}
712
713pub(in crate::barrier) struct DatabaseCheckpointControlMetrics {
714    barrier_latency: LabelGuardedHistogram,
715    in_flight_barrier_nums: LabelGuardedIntGauge,
716    all_barrier_nums: LabelGuardedIntGauge,
717}
718
719impl DatabaseCheckpointControlMetrics {
720    pub(in crate::barrier) fn new(database_id: DatabaseId) -> Self {
721        let database_id_str = database_id.to_string();
722        let barrier_latency = GLOBAL_META_METRICS
723            .barrier_latency
724            .with_guarded_label_values(&[&database_id_str]);
725        let in_flight_barrier_nums = GLOBAL_META_METRICS
726            .in_flight_barrier_nums
727            .with_guarded_label_values(&[&database_id_str]);
728        let all_barrier_nums = GLOBAL_META_METRICS
729            .all_barrier_nums
730            .with_guarded_label_values(&[&database_id_str]);
731        Self {
732            barrier_latency,
733            in_flight_barrier_nums,
734            all_barrier_nums,
735        }
736    }
737}
738
739impl PartialGraphStat for DatabaseCheckpointControlMetrics {
740    fn observe_barrier_latency(&self, _epoch: EpochPair, barrier_latency_secs: f64) {
741        self.barrier_latency.observe(barrier_latency_secs);
742    }
743
744    fn observe_barrier_num(&self, inflight_barrier_num: usize, collected_barrier_num: usize) {
745        self.in_flight_barrier_nums.set(inflight_barrier_num as _);
746        self.all_barrier_nums
747            .set((inflight_barrier_num + collected_barrier_num) as _);
748    }
749}
750
751/// Controls the concurrent execution of commands.
752pub(in crate::barrier) struct DatabaseCheckpointControl {
753    pub(super) database_id: DatabaseId,
754    /// Identifies the current recovery incarnation of this database.
755    pub(super) term_id: String,
756    partial_graph_id: PartialGraphId,
757    pub(super) state: BarrierWorkerState,
758
759    finishing_jobs_collector:
760        BarrierItemCollector<JobId, (Vec<BarrierCompleteResponse>, TrackingJob), ()>,
761    /// The barrier that are completing.
762    completing_barrier: Option<EpochPair>,
763
764    committed_epoch: Option<u64>,
765
766    /// `None` while the database has no streaming job, so that a frozen timestamp does not
767    /// render as an ever-growing barrier pending time.
768    last_committed_barrier_time: Option<LabelGuardedIntGauge>,
769
770    pub(super) database_info: InflightDatabaseInfo,
771    pub independent_checkpoint_job_controls: HashMap<JobId, IndependentCheckpointJobControl>,
772    pub(super) pending_independent_job_subscriptions_to_drop: Vec<PbSubscriptionUpstreamInfo>,
773}
774
775impl DatabaseCheckpointControl {
776    fn new(database_id: DatabaseId, shared_actor_infos: SharedActorInfos) -> Self {
777        Self {
778            database_id,
779            term_id: Uuid::new_v4().to_string(),
780            partial_graph_id: to_partial_graph_id(database_id, None),
781            state: BarrierWorkerState::new(),
782            finishing_jobs_collector: BarrierItemCollector::new(false),
783            completing_barrier: None,
784            committed_epoch: None,
785            last_committed_barrier_time: None,
786            database_info: InflightDatabaseInfo::empty(database_id, shared_actor_infos),
787            independent_checkpoint_job_controls: Default::default(),
788            pending_independent_job_subscriptions_to_drop: Default::default(),
789        }
790    }
791
792    pub(crate) fn recovery(
793        database_id: DatabaseId,
794        term_id: String,
795        state: BarrierWorkerState,
796        committed_epoch: u64,
797        database_info: InflightDatabaseInfo,
798        independent_checkpoint_job_controls: HashMap<JobId, IndependentCheckpointJobControl>,
799    ) -> Self {
800        Self {
801            database_id,
802            term_id,
803            partial_graph_id: to_partial_graph_id(database_id, None),
804            state,
805            finishing_jobs_collector: BarrierItemCollector::new(false),
806            completing_barrier: None,
807            committed_epoch: Some(committed_epoch),
808            last_committed_barrier_time: None,
809            database_info,
810            independent_checkpoint_job_controls,
811            pending_independent_job_subscriptions_to_drop: Default::default(),
812        }
813    }
814
815    pub(super) fn term_id(&self) -> &str {
816        &self.term_id
817    }
818
819    pub(crate) fn is_valid_after_worker_err(&self, worker_id: WorkerId) -> bool {
820        !self.database_info.contains_worker(worker_id as _)
821            && self
822                .independent_checkpoint_job_controls
823                .values()
824                .all(|job| {
825                    job.fragment_infos()
826                        .map(|fragment_infos| {
827                            !InflightFragmentInfo::contains_worker(
828                                fragment_infos.values(),
829                                worker_id,
830                            )
831                        })
832                        .unwrap_or(true)
833                })
834    }
835
836    /// Enqueue a barrier command
837    fn enqueue_command(&mut self, epoch: EpochPair, independent_jobs_to_wait: HashSet<JobId>) {
838        let prev_epoch = epoch.prev;
839        tracing::trace!(prev_epoch, ?independent_jobs_to_wait, "enqueue command");
840        if !independent_jobs_to_wait.is_empty() {
841            self.finishing_jobs_collector
842                .enqueue(epoch, independent_jobs_to_wait, ());
843        }
844    }
845
846    /// Change the state of this `prev_epoch` to `Completed`. Return continuous nodes
847    /// with `Completed` starting from first node [`Completed`..`InFlight`) and remove them.
848    fn barrier_collected(
849        &mut self,
850        partial_graph_id: PartialGraphId,
851        collected_barrier: CollectedBarrier<'_>,
852        periodic_barriers: &mut PeriodicBarriers,
853    ) -> MetaResult<()> {
854        let prev_epoch = collected_barrier.epoch.prev;
855        tracing::trace!(
856            prev_epoch,
857            partial_graph_id = %partial_graph_id,
858            "barrier collected"
859        );
860        let (database_id, independent_job_id) = from_partial_graph_id(partial_graph_id);
861        assert_eq!(self.database_id, database_id);
862        if let Some(independent_job_id) = independent_job_id {
863            let job = self
864                .independent_checkpoint_job_controls
865                .get_mut(&independent_job_id)
866                .expect("should exist");
867            let should_force_checkpoint = job.collect(collected_barrier);
868            if should_force_checkpoint {
869                periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
870            }
871        }
872        Ok(())
873    }
874}
875
876impl DatabaseCheckpointControl {
877    fn collect_backfill_pinned_upstream_tables(&self) -> HashSet<TableId> {
878        self.independent_checkpoint_job_controls
879            .values()
880            .flat_map(|job| job.pinned_upstream_tables().iter().copied())
881            .collect()
882    }
883
884    fn collect_no_shuffle_fragment_relations_for_reschedule_check(
885        &self,
886    ) -> Vec<(FragmentId, FragmentId)> {
887        let mut no_shuffle_relations = Vec::new();
888        for fragment in self.database_info.fragment_infos() {
889            let downstream_fragment_id = fragment.fragment_id;
890            visit_stream_node_cont(&fragment.nodes, |node| {
891                if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
892                    && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
893                {
894                    no_shuffle_relations.push((merge.upstream_fragment_id, downstream_fragment_id));
895                }
896                true
897            });
898        }
899
900        for job in self.independent_checkpoint_job_controls.values() {
901            if let Some(fragment_infos) = job.fragment_infos() {
902                for fragment_info in fragment_infos.values() {
903                    let downstream_fragment_id = fragment_info.fragment_id;
904                    visit_stream_node_cont(&fragment_info.nodes, |node| {
905                        if let Some(NodeBody::Merge(merge)) = node.node_body.as_ref()
906                            && merge.upstream_dispatcher_type == PbDispatcherType::NoShuffle as i32
907                        {
908                            no_shuffle_relations
909                                .push((merge.upstream_fragment_id, downstream_fragment_id));
910                        }
911                        true
912                    });
913                }
914            }
915        }
916        no_shuffle_relations
917    }
918
919    fn collect_reschedule_blocked_jobs_for_independent_jobs_inflight(
920        &self,
921    ) -> MetaResult<HashSet<JobId>> {
922        let mut initial_blocked_fragment_ids = HashSet::new();
923        for job in self.independent_checkpoint_job_controls.values() {
924            if let Some(fragment_infos) = job.fragment_infos() {
925                for fragment_info in fragment_infos.values() {
926                    if fragment_has_online_unreschedulable_scan(fragment_info) {
927                        initial_blocked_fragment_ids.insert(fragment_info.fragment_id);
928                        collect_fragment_upstream_fragment_ids(
929                            fragment_info,
930                            &mut initial_blocked_fragment_ids,
931                        );
932                    }
933                }
934            }
935        }
936
937        let mut blocked_fragment_ids = initial_blocked_fragment_ids.clone();
938        if !initial_blocked_fragment_ids.is_empty() {
939            let no_shuffle_relations =
940                self.collect_no_shuffle_fragment_relations_for_reschedule_check();
941            let (forward_edges, backward_edges) =
942                build_no_shuffle_fragment_graph_edges(no_shuffle_relations);
943            let initial_blocked_fragment_ids: Vec<_> =
944                initial_blocked_fragment_ids.iter().copied().collect();
945            for ensemble in find_no_shuffle_graphs(
946                &initial_blocked_fragment_ids,
947                &forward_edges,
948                &backward_edges,
949            )? {
950                blocked_fragment_ids.extend(ensemble.fragments());
951            }
952        }
953
954        let mut blocked_job_ids = HashSet::new();
955        blocked_job_ids.extend(
956            blocked_fragment_ids
957                .into_iter()
958                .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id)),
959        );
960        Ok(blocked_job_ids)
961    }
962
963    fn collect_reschedule_blocked_job_ids(
964        &self,
965        reschedules: &HashMap<FragmentId, Reschedule>,
966        fragment_actors: &HashMap<FragmentId, HashSet<ActorId>>,
967        blocked_job_ids: &HashSet<JobId>,
968    ) -> HashSet<JobId> {
969        let mut affected_fragment_ids: HashSet<FragmentId> = reschedules.keys().copied().collect();
970        affected_fragment_ids.extend(fragment_actors.keys().copied());
971        for reschedule in reschedules.values() {
972            affected_fragment_ids.extend(reschedule.downstream_fragment_ids.iter().copied());
973            affected_fragment_ids.extend(
974                reschedule
975                    .upstream_fragment_dispatcher_ids
976                    .iter()
977                    .map(|(fragment_id, _)| *fragment_id),
978            );
979        }
980
981        affected_fragment_ids
982            .into_iter()
983            .filter_map(|fragment_id| self.database_info.job_id_by_fragment(fragment_id))
984            .filter(|job_id| blocked_job_ids.contains(job_id))
985            .collect()
986    }
987
988    fn next_complete_barrier_task(
989        &mut self,
990        periodic_barriers: &mut PeriodicBarriers,
991        partial_graph_manager: &mut PartialGraphManager,
992        task: &mut Option<CompleteBarrierTask>,
993        hummock_version_stats: &HummockVersionStats,
994    ) {
995        // `Vec::new` is a const fn, and do not have memory allocation, and therefore is lightweight enough
996        let mut independent_jobs_task = vec![];
997        let mut finished_jobs = Vec::new();
998        let min_upstream_inflight_barrier = partial_graph_manager
999            .first_inflight_barrier(self.partial_graph_id)
1000            .map(|epoch| epoch.prev);
1001        for (job_id, job) in &mut self.independent_checkpoint_job_controls {
1002            let Some(job) = job.ready_mut() else {
1003                continue;
1004            };
1005            match job {
1006                IndependentCheckpointJob::CreatingStreamingJob(creating_job) => {
1007                    if let Some((epoch, resps, info, is_finish_epoch)) = creating_job
1008                        .start_completing(partial_graph_manager, min_upstream_inflight_barrier)
1009                    {
1010                        let resps = resps.into_values().collect_vec();
1011                        if is_finish_epoch {
1012                            assert!(info.notifier.is_none());
1013                            finished_jobs.push((*job_id, epoch, resps));
1014                            continue;
1015                        };
1016                        independent_jobs_task.push((*job_id, epoch, resps, info));
1017                    }
1018                }
1019                IndependentCheckpointJob::BatchRefresh(batch_refresh_job) => {
1020                    if let Some((epoch, resps, info, tracking_job)) =
1021                        batch_refresh_job.start_completing(partial_graph_manager)
1022                    {
1023                        let resps = resps.into_values().collect_vec();
1024                        if let Some(tracking_job) = tracking_job {
1025                            let task = task.get_or_insert_default();
1026                            task.finished_jobs.push(tracking_job);
1027                        }
1028                        independent_jobs_task.push((*job_id, epoch, resps, info));
1029                    }
1030                }
1031            }
1032        }
1033        if !finished_jobs.is_empty() {
1034            partial_graph_manager.remove_partial_graphs(
1035                finished_jobs
1036                    .iter()
1037                    .map(|(job_id, ..)| to_partial_graph_id(self.database_id, Some(*job_id)))
1038                    .collect(),
1039            );
1040        }
1041        for (job_id, epoch, resps) in finished_jobs {
1042            debug!(epoch, %job_id, "finish creating job");
1043            // It's safe to remove the creating job, because on CompleteJobType::Finished,
1044            // all previous barriers have been collected and completed.
1045            // `finished_jobs` was populated above only from a Running creating job, and the
1046            // map cannot change between that scan and this removal.
1047            let Some(IndependentCheckpointJobControl::Running {
1048                job: IndependentCheckpointJob::CreatingStreamingJob(creating_streaming_job),
1049                ..
1050            }) = self.independent_checkpoint_job_controls.remove(&job_id)
1051            else {
1052                panic!("finished job {job_id} should be a creating streaming job");
1053            };
1054            let tracking_job = creating_streaming_job.into_tracking_job();
1055            self.finishing_jobs_collector
1056                .collect(epoch, job_id, (resps, tracking_job));
1057        }
1058        let mut observed_non_checkpoint = false;
1059        self.finishing_jobs_collector.advance_collected();
1060        let epoch_end_bound = self
1061            .finishing_jobs_collector
1062            .first_inflight_epoch()
1063            .map_or(Unbounded, |epoch| Excluded(epoch.prev));
1064        if let Some((epoch, resps, info)) = partial_graph_manager.start_completing(
1065            self.partial_graph_id,
1066            epoch_end_bound,
1067            |_, resps, post_collect_command| {
1068                observed_non_checkpoint = true;
1069                self.handle_refresh_table_info(task, &resps);
1070                self.database_info.apply_collected_command(
1071                    &post_collect_command,
1072                    &resps,
1073                    hummock_version_stats,
1074                );
1075            },
1076        ) {
1077            self.handle_refresh_table_info(task, &resps);
1078            self.database_info.apply_collected_command(
1079                &info.post_collect_command,
1080                &resps,
1081                hummock_version_stats,
1082            );
1083            let mut resps_to_commit = resps.into_values().collect_vec();
1084            let mut staging_commit_info = self.database_info.take_staging_commit_info();
1085            if let Some((_, finished_jobs, _)) =
1086                self.finishing_jobs_collector
1087                    .take_collected_if(|collected_epoch| {
1088                        assert!(epoch <= collected_epoch.prev);
1089                        epoch == collected_epoch.prev
1090                    })
1091            {
1092                finished_jobs
1093                    .into_iter()
1094                    .for_each(|(_, (resps, tracking_job))| {
1095                        resps_to_commit.extend(resps);
1096                        staging_commit_info.finished_jobs.push(tracking_job);
1097                    });
1098            }
1099            {
1100                let task = task.get_or_insert_default();
1101                Command::collect_commit_epoch_info(
1102                    &self.database_info,
1103                    &info,
1104                    task,
1105                    resps_to_commit,
1106                    self.collect_backfill_pinned_upstream_tables(),
1107                );
1108                self.completing_barrier = Some(info.barrier_info.epoch());
1109                task.finished_jobs.extend(staging_commit_info.finished_jobs);
1110                task.finished_cdc_table_backfill
1111                    .extend(staging_commit_info.finished_cdc_table_backfill);
1112                task.epoch_infos
1113                    .try_insert(self.partial_graph_id, info)
1114                    .expect("non duplicate");
1115                task.commit_info
1116                    .truncate_tables
1117                    .extend(staging_commit_info.table_ids_to_truncate);
1118            }
1119        } else if observed_non_checkpoint
1120            && self.database_info.has_pending_finished_jobs()
1121            && !partial_graph_manager.has_pending_checkpoint_barrier(self.partial_graph_id)
1122        {
1123            periodic_barriers.force_checkpoint_in_next_barrier(self.database_id);
1124        }
1125        if !independent_jobs_task.is_empty() {
1126            let task = task.get_or_insert_default();
1127            for (job_id, epoch, resps, info) in independent_jobs_task {
1128                collect_independent_job_commit_epoch_info(task, epoch, resps, &info);
1129                task.epoch_infos
1130                    .try_insert(to_partial_graph_id(self.database_id, Some(job_id)), info)
1131                    .expect("non duplicate");
1132            }
1133        }
1134    }
1135
1136    fn ack_completed(
1137        &mut self,
1138        partial_graph_manager: &mut PartialGraphManager,
1139        command_prev_epoch: Option<u64>,
1140        independent_job_epochs: Vec<(JobId, u64)>,
1141    ) {
1142        {
1143            if let Some(epoch) = self.completing_barrier.take() {
1144                assert_eq!(command_prev_epoch, Some(epoch.prev));
1145                self.committed_epoch = Some(epoch.prev);
1146                partial_graph_manager.ack_completed(self.partial_graph_id, epoch.prev);
1147                for job in self.independent_checkpoint_job_controls.values_mut() {
1148                    job.on_upstream_database_ack_completed(epoch.prev);
1149                }
1150                self.last_committed_barrier_time
1151                    .get_or_insert_with(|| {
1152                        GLOBAL_META_METRICS
1153                            .last_committed_barrier_time
1154                            .with_guarded_label_values(&[&self.database_id.to_string()])
1155                    })
1156                    .set(Epoch(epoch.curr).as_unix_secs() as i64);
1157            } else {
1158                assert_eq!(command_prev_epoch, None);
1159            };
1160            for (job_id, epoch) in independent_job_epochs {
1161                if let Some(job) = self.independent_checkpoint_job_controls.get_mut(&job_id) {
1162                    job.ack_completed(partial_graph_manager, epoch);
1163                }
1164                // If the job is not found, it was dropped and already removed
1165                // by `on_partial_graph_reset` while the completing task was running.
1166            }
1167        }
1168    }
1169
1170    fn handle_refresh_table_info(
1171        &self,
1172        task: &mut Option<CompleteBarrierTask>,
1173        resps: &HashMap<WorkerId, BarrierCompleteResponse>,
1174    ) {
1175        let list_finished_info = resps
1176            .values()
1177            .flat_map(|resp| resp.list_finished_sources.clone())
1178            .collect::<Vec<_>>();
1179        if !list_finished_info.is_empty() {
1180            let task = task.get_or_insert_default();
1181            task.list_finished_source_ids.extend(list_finished_info);
1182        }
1183
1184        let load_finished_info = resps
1185            .values()
1186            .flat_map(|resp| resp.load_finished_sources.clone())
1187            .collect::<Vec<_>>();
1188        if !load_finished_info.is_empty() {
1189            let task = task.get_or_insert_default();
1190            task.load_finished_source_ids.extend(load_finished_info);
1191        }
1192
1193        let refresh_finished_table_ids: Vec<JobId> = resps
1194            .values()
1195            .flat_map(|resp| {
1196                resp.refresh_finished_tables
1197                    .iter()
1198                    .map(|table_id| table_id.as_job_id())
1199            })
1200            .collect::<Vec<_>>();
1201        if !refresh_finished_table_ids.is_empty() {
1202            let task = task.get_or_insert_default();
1203            task.refresh_finished_table_job_ids
1204                .extend(refresh_finished_table_ids);
1205        }
1206    }
1207}
1208
1209impl DatabaseCheckpointControl {
1210    /// Handle the new barrier from the scheduled queue and inject it.
1211    fn handle_new_barrier(
1212        &mut self,
1213        command: Option<(Command, Notifier)>,
1214        checkpoint: bool,
1215        span: tracing::Span,
1216        partial_graph_manager: &mut PartialGraphManager,
1217        hummock_version_stats: &HummockVersionStats,
1218        worker_nodes: &HashMap<WorkerId, WorkerNode>,
1219    ) -> MetaResult<()> {
1220        let curr_epoch = self.state.in_flight_prev_epoch().next();
1221
1222        let (mut command, notifier) = if let Some((command, notifier)) = command {
1223            (Some(command), Some(notifier))
1224        } else {
1225            (None, None)
1226        };
1227
1228        debug_assert!(
1229            !matches!(
1230                &command,
1231                Some(Command::RescheduleIntent {
1232                    reschedule_plan: None,
1233                    ..
1234                })
1235            ),
1236            "reschedule intent should be resolved before injection"
1237        );
1238
1239        let mut notifier_start = notifier.map(Notifier::start);
1240        if let Some(Command::DropStreamingJobs {
1241            streaming_job_ids, ..
1242        }) = &mut command
1243        {
1244            streaming_job_ids.retain(|job_id| {
1245                let Some(job) = self.independent_checkpoint_job_controls.get_mut(job_id) else {
1246                    return true;
1247                };
1248                !job.drop(notifier_start.as_mut(), partial_graph_manager)
1249            });
1250            if streaming_job_ids.is_empty() {
1251                if let Some(notifier) = notifier_start {
1252                    notifier.started();
1253                }
1254                return Ok(());
1255            }
1256        }
1257
1258        if let Some(Command::RescheduleIntent {
1259            reschedule_plan: Some(reschedule_plan),
1260            ..
1261        }) = &command
1262            && !self.independent_checkpoint_job_controls.is_empty()
1263        {
1264            let blocked_job_ids =
1265                self.collect_reschedule_blocked_jobs_for_independent_jobs_inflight()?;
1266            let blocked_reschedule_job_ids = self.collect_reschedule_blocked_job_ids(
1267                &reschedule_plan.reschedules,
1268                &reschedule_plan.fragment_actors,
1269                &blocked_job_ids,
1270            );
1271            if !blocked_reschedule_job_ids.is_empty() {
1272                warn!(
1273                    blocked_reschedule_job_ids = ?blocked_reschedule_job_ids,
1274                    "reject reschedule fragments related to creating unreschedulable backfill jobs"
1275                );
1276                if let Some(notifier) = notifier_start {
1277                    notifier.notify_start_failed(
1278                        anyhow!(
1279                            "cannot reschedule jobs {:?} when creating jobs with unreschedulable backfill fragments",
1280                            blocked_reschedule_job_ids
1281                        )
1282                            .into(),
1283                    );
1284                }
1285                return Ok(());
1286            }
1287        }
1288
1289        if !matches!(&command, Some(Command::CreateStreamingJob { .. }))
1290            && self.database_info.is_empty()
1291        {
1292            assert!(
1293                self.independent_checkpoint_job_controls.is_empty(),
1294                "should not have snapshot backfill job when there is no normal job in database"
1295            );
1296            // Drop the guard to remove the metric series of this database.
1297            self.last_committed_barrier_time = None;
1298            // skip the command when there is nothing to do with the barrier
1299            if let Some(notifier) = notifier_start {
1300                notifier.started();
1301            }
1302            return Ok(());
1303        };
1304
1305        if let Some(Command::CreateStreamingJob {
1306            job_type:
1307                CreateStreamingJobType::SnapshotBackfill { .. }
1308                | CreateStreamingJobType::BatchRefresh(_),
1309            ..
1310        }) = &command
1311            && self.state.is_paused()
1312        {
1313            warn!("cannot create streaming job with snapshot backfill when paused");
1314            if let Some(notifier) = notifier_start {
1315                notifier.notify_start_failed(
1316                    anyhow!("cannot create streaming job with snapshot backfill when paused",)
1317                        .into(),
1318                );
1319            }
1320            return Ok(());
1321        }
1322
1323        let barrier_info = self.state.next_barrier_info(checkpoint, curr_epoch);
1324        // Tracing related stuff
1325        barrier_info.prev_epoch.span().in_scope(|| {
1326            tracing::info!(target: "rw_tracing", epoch = barrier_info.curr_epoch(), "new barrier enqueued");
1327        });
1328        span.record("epoch", barrier_info.curr_epoch());
1329
1330        let epoch = barrier_info.epoch();
1331        let ApplyCommandInfo { jobs_to_wait } = match self.apply_command(
1332            command,
1333            &mut notifier_start,
1334            barrier_info,
1335            partial_graph_manager,
1336            hummock_version_stats,
1337            worker_nodes,
1338        ) {
1339            Ok(info) => {
1340                assert!(notifier_start.is_none());
1341                info
1342            }
1343            Err(err) => {
1344                if let Some(notifier) = notifier_start {
1345                    notifier.notify_start_failed(err.clone());
1346                }
1347                fail_point!("inject_barrier_err_success");
1348                return Err(err);
1349            }
1350        };
1351
1352        // Record the in-flight barrier.
1353        self.enqueue_command(epoch, jobs_to_wait);
1354
1355        Ok(())
1356    }
1357
1358    // ── Batch refresh trigger helpers ────────────────────────────────────────
1359
1360    /// Get the last committed epoch for a batch refresh job.
1361    pub(crate) fn get_batch_refresh_trigger_info(&self, job_id: JobId) -> u64 {
1362        // `CheckpointControl::next_event` emits a trigger only after matching this job as a
1363        // Running batch-refresh job. The worker reads this value before awaiting any other work.
1364        let job = self
1365            .independent_checkpoint_job_controls
1366            .get(&job_id)
1367            .expect("batch refresh job should exist")
1368            .running()
1369            .expect("batch refresh job should be running");
1370        match job {
1371            IndependentCheckpointJob::BatchRefresh(br_job) => br_job
1372                .last_committed_epoch()
1373                .expect("idle job must have a last_committed_epoch"),
1374            _ => panic!("job {} should be a batch refresh job", job_id),
1375        }
1376    }
1377
1378    /// Whether the batch refresh job already has its cached context populated.
1379    /// Start a batch refresh logstore consumption run.
1380    /// Returns true if a run was started, false if no log epochs to consume.
1381    pub(crate) fn start_batch_refresh_run(
1382        &mut self,
1383        job_id: JobId,
1384        context: &BatchRefreshJobTriggerContext,
1385        worker_nodes: &HashMap<WorkerId, WorkerNode>,
1386        actor_id_counter: &AtomicU32,
1387        partial_graph_manager: &mut PartialGraphManager,
1388    ) -> MetaResult<bool> {
1389        // The global barrier worker handles the trigger serially. Although metadata loading is
1390        // asynchronous, it does not process a drop/reset event before reaching this call.
1391        let term_id = self.term_id.as_str();
1392        let job = self
1393            .independent_checkpoint_job_controls
1394            .get_mut(&job_id)
1395            .expect("batch refresh job should exist")
1396            .running_mut()
1397            .expect("batch refresh job should be running");
1398        match job {
1399            IndependentCheckpointJob::BatchRefresh(br_job) => br_job.start_refresh_run(
1400                context,
1401                worker_nodes,
1402                actor_id_counter,
1403                term_id,
1404                partial_graph_manager,
1405            ),
1406            _ => panic!("job {} should be a batch refresh job", job_id),
1407        }
1408    }
1409}