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