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