Skip to main content

risingwave_meta/barrier/checkpoint/independent_job/creating_job/
mod.rs

1// Copyright 2026 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
15mod barrier_control;
16mod status;
17
18use std::cmp::max;
19use std::collections::{HashMap, HashSet, VecDeque, hash_map};
20use std::mem::take;
21use std::ops::Bound::{Excluded, Unbounded};
22use std::time::Duration;
23
24use risingwave_common::catalog::{DatabaseId, TableId};
25use risingwave_common::id::JobId;
26use risingwave_common::metrics::LabelGuardedIntGauge;
27use risingwave_common::util::epoch::Epoch;
28use risingwave_meta_model::WorkerId;
29use risingwave_pb::ddl_service::PbBackfillType;
30use risingwave_pb::hummock::HummockVersionStats;
31use risingwave_pb::id::{ActorId, FragmentId, PartialGraphId};
32use risingwave_pb::stream_plan::barrier::PbBarrierKind;
33use risingwave_pb::stream_plan::barrier_mutation::Mutation;
34use risingwave_pb::stream_plan::{AddMutation, StopMutation};
35use risingwave_pb::stream_service::BarrierCompleteResponse;
36use status::CreatingStreamingJobStatus;
37use tracing::{debug, info};
38
39use super::super::state::RenderResult;
40use super::{
41    IndependentCheckpointJob, IndependentCheckpointJobControl, IndependentCheckpointJobStatus,
42};
43use crate::MetaResult;
44use crate::barrier::backfill_order_control::get_nodes_with_backfill_dependencies;
45use crate::barrier::checkpoint::independent_job::creating_job::barrier_control::CreatingStreamingJobBarrierStats;
46use crate::barrier::checkpoint::independent_job::creating_job::status::CreateMviewLogStoreProgressTracker;
47use crate::barrier::command::{
48    PostCollectCommand, TableLogEpochs, ThrottleConfigMap, UpstreamTableLogEpochs,
49};
50use crate::barrier::context::CreateSnapshotBackfillJobCommandInfo;
51use crate::barrier::edge_builder::FragmentEdgeBuildResult;
52use crate::barrier::info::{BarrierInfo, InflightStreamingJobInfo};
53use crate::barrier::notifier::NotifierStarter;
54use crate::barrier::partial_graph::{
55    CollectedBarrier, PartialGraphBarrierInfo, PartialGraphManager, PartialGraphRecoverer,
56};
57use crate::barrier::progress::{CreateMviewProgressTracker, TrackingJob, collect_done_fragments};
58use crate::barrier::rpc::{build_locality_fragment_state_table_mapping, to_partial_graph_id};
59use crate::barrier::{
60    BackfillOrderState, BackfillProgress, BarrierKind, Command, FragmentBackfillProgress,
61    TracedEpoch,
62};
63use crate::controller::fragment::InflightFragmentInfo;
64use crate::manager::MetaOpts;
65use crate::model::{FragmentDownstreamRelation, StreamActor, StreamJobActorsToCreate};
66use crate::rpc::metrics::GLOBAL_META_METRICS;
67use crate::stream::source_manager::SplitAssignment;
68use crate::stream::{ExtendedFragmentBackfillOrder, build_actor_connector_splits};
69
70fn snapshot_backfill_max_pending_barrier_num(opts: &MetaOpts) -> usize {
71    opts.in_flight_barrier_nums
72        .saturating_mul(opts.snapshot_backfill_barrier_amplification_factor.max(1))
73}
74
75#[derive(Debug)]
76pub(crate) struct CreatingJobInfo {
77    pub fragment_infos: HashMap<FragmentId, InflightFragmentInfo>,
78    pub upstream_fragment_downstreams: FragmentDownstreamRelation,
79    pub downstreams: FragmentDownstreamRelation,
80    pub snapshot_backfill_upstream_tables: HashSet<TableId>,
81    pub stream_actors: HashMap<ActorId, StreamActor>,
82}
83
84#[derive(Debug)]
85pub(crate) struct CreatingStreamingJobControl {
86    job_id: JobId,
87    partial_graph_id: PartialGraphId,
88    snapshot_backfill_upstream_tables: HashSet<TableId>,
89    snapshot_epoch: u64,
90
91    node_actors: HashMap<WorkerId, HashSet<ActorId>>,
92    state_table_ids: HashSet<TableId>,
93
94    max_committed_epoch: Option<u64>,
95    status: CreatingStreamingJobStatus,
96    max_lagged_barrier_num: usize,
97    max_pending_barrier_num: usize,
98
99    upstream_lag: LabelGuardedIntGauge,
100}
101
102impl CreatingStreamingJobControl {
103    #[expect(clippy::too_many_arguments)]
104    pub(crate) fn new<'a>(
105        entry: hash_map::VacantEntry<'a, JobId, IndependentCheckpointJobControl>,
106        create_info: CreateSnapshotBackfillJobCommandInfo,
107        notifier: Option<&mut NotifierStarter>,
108        snapshot_backfill_upstream_tables: HashSet<TableId>,
109        snapshot_epoch: u64,
110        since_timestamp_upstream_log_epochs: Option<(&TableLogEpochs, PartialGraphId, u64)>,
111        version_stat: &HummockVersionStats,
112        term_id: &str,
113        partial_graph_manager: &mut PartialGraphManager,
114        edges: &mut FragmentEdgeBuildResult,
115        split_assignment: &SplitAssignment,
116        actors: &RenderResult,
117    ) -> MetaResult<&'a mut Self> {
118        let info = create_info.info.clone();
119        let job_id = info.stream_job_fragments.stream_job_id();
120        let database_id = info.streaming_job.database_id();
121        let is_since_timestamp = since_timestamp_upstream_log_epochs.is_some();
122        debug!(
123            %job_id,
124            definition = info.definition,
125            "new creating job"
126        );
127        let fragment_infos = info
128            .stream_job_fragments
129            .new_fragment_info(
130                &actors.stream_actors,
131                &actors.actor_location,
132                split_assignment,
133            )
134            .collect();
135        let snapshot_backfill_actors: HashSet<ActorId> =
136            InflightStreamingJobInfo::snapshot_backfill_actor_ids(&fragment_infos).collect();
137        let backfill_nodes_to_pause =
138            get_nodes_with_backfill_dependencies(&info.fragment_backfill_ordering)
139                .into_iter()
140                .collect();
141        let backfill_order_state = BackfillOrderState::new(
142            &info.fragment_backfill_ordering,
143            &fragment_infos,
144            info.locality_fragment_state_table_mapping.clone(),
145        );
146        let create_mview_tracker = CreateMviewProgressTracker::recover(
147            job_id,
148            &fragment_infos,
149            backfill_order_state,
150            version_stat,
151        );
152
153        let actors_to_create = Command::create_streaming_job_actors_to_create(
154            &info,
155            edges,
156            &actors.stream_actors,
157            &actors.actor_location,
158        );
159
160        let mut prev_epoch_fake_physical_time = 0;
161        let mut pending_non_checkpoint_barriers = vec![];
162
163        let (initial_barrier_info, log_store_barriers_to_inject) = if let Some((
164            upstream_log_epochs,
165            upstream_partial_graph_id,
166            new_upstream_barrier_prev_epoch,
167        )) =
168            since_timestamp_upstream_log_epochs
169        {
170            let (initial_barrier, barriers_to_inject) =
171                Self::resolve_since_timestamp_upstream_log_epochs(
172                    upstream_log_epochs,
173                    partial_graph_manager.pending_barrier_infos(upstream_partial_graph_id),
174                    snapshot_epoch,
175                    new_upstream_barrier_prev_epoch,
176                )?;
177            (initial_barrier, Some(barriers_to_inject))
178        } else {
179            (
180                CreatingStreamingJobStatus::new_fake_barrier(
181                    &mut prev_epoch_fake_physical_time,
182                    &mut pending_non_checkpoint_barriers,
183                    PbBarrierKind::Checkpoint,
184                ),
185                None,
186            )
187        };
188
189        let added_actors: Vec<ActorId> = actors
190            .stream_actors
191            .values()
192            .flatten()
193            .map(|actor| actor.actor_id)
194            .collect();
195        let actor_splits = split_assignment
196            .values()
197            .flat_map(build_actor_connector_splits)
198            .collect();
199
200        assert!(
201            info.cdc_table_snapshot_splits.is_none(),
202            "should not have cdc backfill for snapshot backfill job"
203        );
204
205        let initial_mutation = Mutation::Add(AddMutation {
206            // for mutation of snapshot backfill job, we won't include changes to dispatchers of upstream actors.
207            actor_dispatchers: Default::default(),
208            added_actors,
209            actor_splits,
210            // we assume that when handling snapshot backfill, the cluster must not be paused
211            pause: false,
212            subscriptions_to_add: Default::default(),
213            backfill_nodes_to_pause,
214            actor_cdc_table_snapshot_splits: None,
215            new_upstream_sinks: Default::default(),
216            dropped_actors: Default::default(),
217            sink_log_store_flush: Default::default(),
218        });
219
220        let node_actors = InflightFragmentInfo::actor_ids_to_collect(fragment_infos.values());
221        let state_table_ids =
222            InflightFragmentInfo::existing_table_ids(fragment_infos.values()).collect();
223
224        let partial_graph_id = to_partial_graph_id(database_id, Some(job_id));
225        let max_lagged_barrier_num = partial_graph_manager
226            .control_stream_manager()
227            .env
228            .opts
229            .snapshot_backfill_finish_max_lagged_barriers;
230        let opts = &partial_graph_manager.control_stream_manager().env.opts;
231        let max_pending_barrier_num = snapshot_backfill_max_pending_barrier_num(opts);
232
233        let mut job = Self {
234            partial_graph_id,
235            job_id,
236            snapshot_backfill_upstream_tables,
237            max_committed_epoch: None,
238            snapshot_epoch,
239            status: CreatingStreamingJobStatus::PlaceHolder, // filled in later code
240            max_lagged_barrier_num,
241            max_pending_barrier_num,
242            upstream_lag: GLOBAL_META_METRICS
243                .snapshot_backfill_lag
244                .with_guarded_label_values(&[&format!("{}", job_id)]),
245            node_actors,
246            state_table_ids,
247        };
248
249        let mut graph_adder = partial_graph_manager.add_partial_graph(
250            partial_graph_id,
251            term_id,
252            CreatingStreamingJobBarrierStats::new(job_id, snapshot_epoch),
253        );
254
255        if let Err(e) = Self::inject_barrier(
256            partial_graph_id,
257            graph_adder.manager(),
258            &job.node_actors,
259            &job.state_table_ids,
260            false,
261            initial_barrier_info,
262            Some(actors_to_create),
263            Some(initial_mutation),
264            notifier,
265            Some(create_info),
266        ) {
267            graph_adder.failed();
268            entry.insert(IndependentCheckpointJobControl::Resetting {
269                pinned_upstream_tables: job.snapshot_backfill_upstream_tables,
270                subscriptions_to_drop: vec![],
271                notifiers: vec![],
272            });
273            return Err(e);
274        }
275
276        graph_adder.added();
277        let job_info = CreatingJobInfo {
278            fragment_infos,
279            upstream_fragment_downstreams: info.upstream_fragment_downstreams.clone(),
280            downstreams: info.stream_job_fragments.downstreams,
281            snapshot_backfill_upstream_tables: job.snapshot_backfill_upstream_tables.clone(),
282            stream_actors: actors
283                .stream_actors
284                .values()
285                .flatten()
286                .map(|actor| (actor.actor_id, actor.clone()))
287                .collect(),
288        };
289        if let Some(log_store_barriers_to_inject) = log_store_barriers_to_inject {
290            let upstream_lag = log_store_barriers_to_inject
291                .last()
292                .map(|info| info.prev_epoch().saturating_sub(snapshot_epoch))
293                .unwrap_or(0);
294            job.status = CreatingStreamingJobStatus::ConsumingLogStore {
295                tracking_job: TrackingJob::recovered(job_id, &job_info.fragment_infos),
296                info: job_info,
297                log_store_progress_tracker: CreateMviewLogStoreProgressTracker::new(
298                    snapshot_backfill_actors.iter().cloned(),
299                    upstream_lag,
300                ),
301                pending_barriers: log_store_barriers_to_inject.into(),
302            };
303        } else {
304            assert!(pending_non_checkpoint_barriers.is_empty());
305            job.status = CreatingStreamingJobStatus::ConsumingSnapshot {
306                prev_epoch_fake_physical_time,
307                pending_upstream_barriers: vec![],
308                version_stats: version_stat.clone(),
309                create_mview_tracker,
310                snapshot_backfill_actors,
311                snapshot_epoch,
312                info: job_info,
313                pending_non_checkpoint_barriers,
314            };
315        }
316        let job_control = entry.insert(if is_since_timestamp {
317            // Since-timestamp resolution waits for current completion work and requires the
318            // resolved snapshot epoch to be older than the upstream committed epoch.
319            IndependentCheckpointJobControl::creating_streaming_job(
320                job_id,
321                partial_graph_id,
322                IndependentCheckpointJobStatus::Ready,
323                job,
324            )
325        } else {
326            IndependentCheckpointJobControl::creating_streaming_job(
327                job_id,
328                partial_graph_id,
329                IndependentCheckpointJobStatus::Initial { snapshot_epoch },
330                job,
331            )
332        });
333        let Some(IndependentCheckpointJob::CreatingStreamingJob(job)) = job_control.running_mut()
334        else {
335            unreachable!()
336        };
337        Ok(job)
338    }
339
340    pub(super) fn gen_fragment_backfill_progress(&self) -> Vec<FragmentBackfillProgress> {
341        match &self.status {
342            CreatingStreamingJobStatus::ConsumingSnapshot {
343                create_mview_tracker,
344                info,
345                ..
346            } => create_mview_tracker.collect_fragment_progress(&info.fragment_infos, true),
347            CreatingStreamingJobStatus::ConsumingLogStore { info, .. } => {
348                collect_done_fragments(self.job_id, &info.fragment_infos)
349            }
350            CreatingStreamingJobStatus::Finishing(_, _)
351            | CreatingStreamingJobStatus::PlaceHolder => vec![],
352        }
353    }
354
355    fn resolve_upstream_log_epochs(
356        snapshot_backfill_upstream_tables: &HashSet<TableId>,
357        upstream_table_log_epochs: &UpstreamTableLogEpochs,
358        exclusive_start_log_epoch: u64,
359        upstream_barrier_info: &BarrierInfo,
360    ) -> MetaResult<Vec<BarrierInfo>> {
361        let table_id = snapshot_backfill_upstream_tables
362            .iter()
363            .next()
364            .expect("snapshot backfill job should have upstream");
365        let epochs_iter = if let Some(epochs) = upstream_table_log_epochs.get(table_id) {
366            let mut epochs_iter = epochs.iter();
367            loop {
368                let (_, checkpoint_epoch) =
369                    epochs_iter.next().expect("not reach committed epoch yet");
370                if *checkpoint_epoch < exclusive_start_log_epoch {
371                    continue;
372                }
373                assert_eq!(*checkpoint_epoch, exclusive_start_log_epoch);
374                break;
375            }
376            epochs_iter
377        } else {
378            // snapshot backfill job has been marked as creating, but upstream table has not committed a new epoch yet, so no table change log
379            assert_eq!(
380                upstream_barrier_info.prev_epoch(),
381                exclusive_start_log_epoch
382            );
383            static EMPTY_VEC: Vec<(Vec<u64>, u64)> = Vec::new();
384            EMPTY_VEC.iter()
385        };
386
387        let mut ret = vec![];
388        let mut prev_epoch = exclusive_start_log_epoch;
389        let mut pending_non_checkpoint_barriers = vec![];
390        for (non_checkpoint_epochs, checkpoint_epoch) in epochs_iter {
391            for (i, epoch) in non_checkpoint_epochs
392                .iter()
393                .chain([checkpoint_epoch])
394                .enumerate()
395            {
396                assert!(*epoch > prev_epoch);
397                pending_non_checkpoint_barriers.push(prev_epoch);
398                ret.push(BarrierInfo {
399                    prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
400                    curr_epoch: TracedEpoch::new(Epoch(*epoch)),
401                    kind: if i == 0 {
402                        BarrierKind::Checkpoint(take(&mut pending_non_checkpoint_barriers))
403                    } else {
404                        BarrierKind::Barrier
405                    },
406                });
407                prev_epoch = *epoch;
408            }
409        }
410        ret.push(BarrierInfo {
411            prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
412            curr_epoch: TracedEpoch::new(Epoch(upstream_barrier_info.curr_epoch())),
413            kind: BarrierKind::Checkpoint(pending_non_checkpoint_barriers),
414        });
415        Ok(ret)
416    }
417
418    /// Resolves the log-store barriers that must be injected before the create barrier.
419    ///
420    /// Example with pending upstream barriers:
421    ///
422    /// ```text
423    /// snapshot epoch: 60
424    /// changelog after snapshot: [61, 62, 63, 64] + 65
425    /// pending upstream barriers: 66 -> 67 barrier, 67 -> 68 barrier,
426    ///                            68 -> 69 barrier, 69 -> 70 barrier
427    /// new create barrier: 70 -> 71
428    ///
429    /// injected: 60 -> 61 checkpoint, 61 -> 62 barrier, ..., 64 -> 65 barrier,
430    ///           65 -> 66 checkpoint, 66 -> 67 barrier, 67 -> 68 barrier, ...,
431    ///           69 -> 70 barrier
432    /// current create barrier later injects: 70 -> 71 checkpoint
433    /// ```
434    ///
435    /// Example without pending upstream barriers:
436    ///
437    /// ```text
438    /// snapshot epoch: 60
439    /// changelog after snapshot: [61, 62, 63, 64] + 65
440    /// new create barrier: 66 -> 67
441    ///
442    /// injected: 60 -> 61 checkpoint, 61 -> 62 barrier, ..., 64 -> 65 barrier,
443    ///           65 -> 66 checkpoint
444    /// current create barrier later injects: 66 -> 67 checkpoint
445    /// ```
446    fn resolve_since_timestamp_upstream_log_epochs(
447        upstream_log_epochs: &TableLogEpochs,
448        pending_upstream_barriers: impl Iterator<Item = &BarrierInfo>,
449        snapshot_epoch: u64,
450        new_upstream_barrier_prev_epoch: u64,
451    ) -> MetaResult<(BarrierInfo, Vec<BarrierInfo>)> {
452        let mut initial_barrier = None;
453        let mut barriers = vec![];
454        fn emit_barrier(
455            initial_barrier: &mut Option<BarrierInfo>,
456            barriers: &mut Vec<BarrierInfo>,
457            barrier: BarrierInfo,
458        ) {
459            if initial_barrier.is_none() {
460                *initial_barrier = Some(barrier);
461            } else {
462                barriers.push(barrier);
463            }
464        }
465
466        let mut prev_epoch = snapshot_epoch;
467        let mut pending_non_checkpoint_barriers = vec![];
468        for (non_checkpoint_epochs, checkpoint_epoch) in upstream_log_epochs {
469            for (i, epoch) in non_checkpoint_epochs
470                .iter()
471                .chain([checkpoint_epoch])
472                .enumerate()
473            {
474                assert!(
475                    *epoch > prev_epoch,
476                    "changelog epochs should be strictly increasing"
477                );
478                pending_non_checkpoint_barriers.push(prev_epoch);
479                emit_barrier(
480                    &mut initial_barrier,
481                    &mut barriers,
482                    BarrierInfo {
483                        prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
484                        curr_epoch: TracedEpoch::new(Epoch(*epoch)),
485                        kind: if i == 0 {
486                            BarrierKind::Checkpoint(take(&mut pending_non_checkpoint_barriers))
487                        } else {
488                            BarrierKind::Barrier
489                        },
490                    },
491                );
492                prev_epoch = *epoch;
493            }
494        }
495
496        let mut pending_upstream_barriers = pending_upstream_barriers.peekable();
497        pending_non_checkpoint_barriers.push(prev_epoch);
498        if pending_upstream_barriers.peek().is_none() {
499            assert!(
500                new_upstream_barrier_prev_epoch > prev_epoch,
501                "new upstream barrier prev epoch should be newer than the latest changelog epoch"
502            );
503            emit_barrier(
504                &mut initial_barrier,
505                &mut barriers,
506                BarrierInfo {
507                    prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
508                    curr_epoch: TracedEpoch::new(Epoch(new_upstream_barrier_prev_epoch)),
509                    kind: BarrierKind::Checkpoint(pending_non_checkpoint_barriers),
510                },
511            );
512        } else {
513            let first_pending_barrier = pending_upstream_barriers
514                .peek()
515                .expect("first pending upstream barrier should exist after peek");
516            assert!(
517                first_pending_barrier.prev_epoch() > prev_epoch,
518                "first pending upstream barrier should be newer than the latest resolved changelog epoch"
519            );
520            emit_barrier(
521                &mut initial_barrier,
522                &mut barriers,
523                BarrierInfo {
524                    prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
525                    curr_epoch: TracedEpoch::new(Epoch(first_pending_barrier.prev_epoch())),
526                    kind: BarrierKind::Checkpoint(take(&mut pending_non_checkpoint_barriers)),
527                },
528            );
529            prev_epoch = first_pending_barrier.prev_epoch();
530            for pending_barrier in pending_upstream_barriers {
531                assert_eq!(
532                    pending_barrier.prev_epoch(),
533                    prev_epoch,
534                    "pending upstream barriers should continue from resolved changelog epochs"
535                );
536                pending_non_checkpoint_barriers.push(prev_epoch);
537                emit_barrier(
538                    &mut initial_barrier,
539                    &mut barriers,
540                    BarrierInfo {
541                        prev_epoch: TracedEpoch::new(Epoch(prev_epoch)),
542                        curr_epoch: TracedEpoch::new(Epoch(pending_barrier.curr_epoch())),
543                        kind: if pending_barrier.kind.is_checkpoint() {
544                            BarrierKind::Checkpoint(take(&mut pending_non_checkpoint_barriers))
545                        } else {
546                            BarrierKind::Barrier
547                        },
548                    },
549                );
550                prev_epoch = pending_barrier.curr_epoch();
551            }
552            assert_eq!(
553                new_upstream_barrier_prev_epoch, prev_epoch,
554                "new upstream barrier prev epoch should match the latest pending log-store epoch"
555            );
556        }
557        let Some(initial_barrier) = initial_barrier else {
558            return Err(anyhow::anyhow!(
559                "missing lagging barriers for direct log-store start from snapshot epoch {}",
560                snapshot_epoch
561            )
562            .into());
563        };
564        assert!(initial_barrier.kind.is_checkpoint());
565        Ok((initial_barrier, barriers))
566    }
567
568    fn recover_consuming_snapshot(
569        job_id: JobId,
570        upstream_table_log_epochs: &UpstreamTableLogEpochs,
571        snapshot_epoch: u64,
572        committed_epoch: u64,
573        upstream_barrier_info: &BarrierInfo,
574        info: CreatingJobInfo,
575        backfill_order_state: BackfillOrderState,
576        version_stat: &HummockVersionStats,
577    ) -> MetaResult<(CreatingStreamingJobStatus, BarrierInfo)> {
578        let mut prev_epoch_fake_physical_time = Epoch(committed_epoch).physical_time();
579        let mut pending_non_checkpoint_barriers = vec![];
580        let create_mview_tracker = CreateMviewProgressTracker::recover(
581            job_id,
582            &info.fragment_infos,
583            backfill_order_state,
584            version_stat,
585        );
586        let barrier_info = CreatingStreamingJobStatus::new_fake_barrier(
587            &mut prev_epoch_fake_physical_time,
588            &mut pending_non_checkpoint_barriers,
589            PbBarrierKind::Initial,
590        );
591        Ok((
592            CreatingStreamingJobStatus::ConsumingSnapshot {
593                prev_epoch_fake_physical_time,
594                pending_upstream_barriers: Self::resolve_upstream_log_epochs(
595                    &info.snapshot_backfill_upstream_tables,
596                    upstream_table_log_epochs,
597                    snapshot_epoch,
598                    upstream_barrier_info,
599                )?,
600                version_stats: version_stat.clone(),
601                create_mview_tracker,
602                snapshot_backfill_actors: InflightStreamingJobInfo::snapshot_backfill_actor_ids(
603                    &info.fragment_infos,
604                )
605                .collect(),
606                info,
607                snapshot_epoch,
608                pending_non_checkpoint_barriers,
609            },
610            barrier_info,
611        ))
612    }
613
614    fn recover_consuming_log_store(
615        job_id: JobId,
616        upstream_table_log_epochs: &UpstreamTableLogEpochs,
617        committed_epoch: u64,
618        upstream_barrier_info: &BarrierInfo,
619        info: CreatingJobInfo,
620    ) -> MetaResult<(CreatingStreamingJobStatus, BarrierInfo)> {
621        let mut pending_barriers: VecDeque<_> = Self::resolve_upstream_log_epochs(
622            &info.snapshot_backfill_upstream_tables,
623            upstream_table_log_epochs,
624            committed_epoch,
625            upstream_barrier_info,
626        )?
627        .into();
628        let mut first_barrier = pending_barriers
629            .pop_front()
630            .expect("resolved upstream log epochs should not be empty");
631        assert!(first_barrier.kind.is_checkpoint());
632        first_barrier.kind = BarrierKind::Initial;
633
634        Ok((
635            CreatingStreamingJobStatus::ConsumingLogStore {
636                tracking_job: TrackingJob::recovered(job_id, &info.fragment_infos),
637                log_store_progress_tracker: CreateMviewLogStoreProgressTracker::new(
638                    InflightStreamingJobInfo::snapshot_backfill_actor_ids(&info.fragment_infos),
639                    pending_barriers
640                        .back()
641                        .map(|info| info.prev_epoch() - committed_epoch)
642                        .unwrap_or(0),
643                ),
644                pending_barriers,
645                info,
646            },
647            first_barrier,
648        ))
649    }
650
651    #[expect(clippy::too_many_arguments)]
652    pub(crate) fn recover(
653        database_id: DatabaseId,
654        job_id: JobId,
655        snapshot_backfill_upstream_tables: HashSet<TableId>,
656        upstream_table_log_epochs: &UpstreamTableLogEpochs,
657        snapshot_epoch: u64,
658        committed_epoch: u64,
659        upstream_barrier_info: &BarrierInfo,
660        fragment_infos: HashMap<FragmentId, InflightFragmentInfo>,
661        backfill_order: ExtendedFragmentBackfillOrder,
662        fragment_relations: &FragmentDownstreamRelation,
663        version_stat: &HummockVersionStats,
664        new_actors: StreamJobActorsToCreate,
665        initial_mutation: Mutation,
666        term_id: &str,
667        partial_graph_recoverer: &mut PartialGraphRecoverer<'_>,
668    ) -> MetaResult<Self> {
669        info!(
670            %job_id,
671            "recovered creating snapshot backfill job"
672        );
673
674        let node_actors = InflightFragmentInfo::actor_ids_to_collect(fragment_infos.values());
675        let state_table_ids: HashSet<_> =
676            InflightFragmentInfo::existing_table_ids(fragment_infos.values()).collect();
677
678        let mut upstream_fragment_downstreams: FragmentDownstreamRelation = Default::default();
679        for (upstream_fragment_id, downstreams) in fragment_relations {
680            if fragment_infos.contains_key(upstream_fragment_id) {
681                continue;
682            }
683            for downstream in downstreams {
684                if fragment_infos.contains_key(&downstream.downstream_fragment_id) {
685                    upstream_fragment_downstreams
686                        .entry(*upstream_fragment_id)
687                        .or_default()
688                        .push(downstream.clone());
689                }
690            }
691        }
692        let downstreams = fragment_infos
693            .keys()
694            .filter_map(|fragment_id| {
695                fragment_relations
696                    .get(fragment_id)
697                    .map(|relation| (*fragment_id, relation.clone()))
698            })
699            .collect();
700
701        let info = CreatingJobInfo {
702            fragment_infos,
703            upstream_fragment_downstreams,
704            downstreams,
705            snapshot_backfill_upstream_tables: snapshot_backfill_upstream_tables.clone(),
706            stream_actors: new_actors
707                .values()
708                .flat_map(|fragments| {
709                    fragments.values().flat_map(|(_, actors, _)| {
710                        actors
711                            .iter()
712                            .map(|(actor, _, _)| (actor.actor_id, actor.clone()))
713                    })
714                })
715                .collect(),
716        };
717
718        let (status, first_barrier_info) = if committed_epoch < snapshot_epoch {
719            let locality_fragment_state_table_mapping =
720                build_locality_fragment_state_table_mapping(&info.fragment_infos);
721            let backfill_order_state = BackfillOrderState::recover_from_fragment_infos(
722                &backfill_order,
723                &info.fragment_infos,
724                locality_fragment_state_table_mapping,
725            );
726            Self::recover_consuming_snapshot(
727                job_id,
728                upstream_table_log_epochs,
729                snapshot_epoch,
730                committed_epoch,
731                upstream_barrier_info,
732                info,
733                backfill_order_state,
734                version_stat,
735            )?
736        } else {
737            Self::recover_consuming_log_store(
738                job_id,
739                upstream_table_log_epochs,
740                committed_epoch,
741                upstream_barrier_info,
742                info,
743            )?
744        };
745
746        let partial_graph_id = to_partial_graph_id(database_id, Some(job_id));
747        let max_lagged_barrier_num = partial_graph_recoverer
748            .control_stream_manager()
749            .env
750            .opts
751            .snapshot_backfill_finish_max_lagged_barriers;
752        let opts = &partial_graph_recoverer.control_stream_manager().env.opts;
753        let max_pending_barrier_num = snapshot_backfill_max_pending_barrier_num(opts);
754
755        partial_graph_recoverer.recover_graph(
756            partial_graph_id,
757            term_id,
758            initial_mutation,
759            &first_barrier_info,
760            &node_actors,
761            state_table_ids.iter().copied(),
762            new_actors,
763            CreatingStreamingJobBarrierStats::new(job_id, snapshot_epoch),
764        )?;
765
766        Ok(Self {
767            job_id,
768            partial_graph_id,
769            snapshot_backfill_upstream_tables,
770            snapshot_epoch,
771            node_actors,
772            state_table_ids,
773            max_committed_epoch: Some(committed_epoch),
774            status,
775            max_lagged_barrier_num,
776            max_pending_barrier_num,
777            upstream_lag: GLOBAL_META_METRICS
778                .snapshot_backfill_lag
779                .with_guarded_label_values(&[&format!("{}", job_id)]),
780        })
781    }
782
783    pub(crate) fn gen_backfill_progress(&self) -> BackfillProgress {
784        let progress = match &self.status {
785            CreatingStreamingJobStatus::ConsumingSnapshot {
786                create_mview_tracker,
787                ..
788            } => {
789                if create_mview_tracker.is_finished() {
790                    "Snapshot finished".to_owned()
791                } else {
792                    let progress = create_mview_tracker.gen_backfill_progress();
793                    format!("Snapshot [{}]", progress)
794                }
795            }
796            CreatingStreamingJobStatus::ConsumingLogStore {
797                log_store_progress_tracker,
798                ..
799            } => {
800                format!(
801                    "LogStore [{}]",
802                    log_store_progress_tracker.gen_backfill_progress()
803                )
804            }
805            CreatingStreamingJobStatus::Finishing(finish_epoch, ..) => {
806                let committed_epoch = self.max_committed_epoch.expect("should have committed");
807                let lag = Duration::from_millis(
808                    Epoch(*finish_epoch).physical_time() - Epoch(committed_epoch).physical_time(),
809                );
810                format!("Finishing [epoch lag: {lag:?}]",)
811            }
812            CreatingStreamingJobStatus::PlaceHolder => {
813                unreachable!()
814            }
815        };
816        BackfillProgress {
817            progress,
818            backfill_type: PbBackfillType::SnapshotBackfill,
819        }
820    }
821
822    pub(super) fn pinned_upstream_tables(&self) -> &HashSet<TableId> {
823        &self.snapshot_backfill_upstream_tables
824    }
825
826    fn inject_barrier(
827        partial_graph_id: PartialGraphId,
828        partial_graph_manager: &mut PartialGraphManager,
829        node_actors: &HashMap<WorkerId, HashSet<ActorId>>,
830        state_table_ids: &HashSet<TableId>,
831        is_finishing: bool,
832        barrier_info: BarrierInfo,
833        new_actors: Option<StreamJobActorsToCreate>,
834        mutation: Option<Mutation>,
835        notifier: Option<&mut NotifierStarter>,
836        first_create_info: Option<CreateSnapshotBackfillJobCommandInfo>,
837    ) -> MetaResult<()> {
838        let (table_ids_to_sync, nodes_to_sync_table) = if !is_finishing {
839            (Some(state_table_ids), Some(node_actors.keys().copied()))
840        } else {
841            (None, None)
842        };
843        partial_graph_manager.inject_barrier(
844            partial_graph_id,
845            mutation,
846            node_actors,
847            table_ids_to_sync.into_iter().flatten().copied(),
848            nodes_to_sync_table.into_iter().flatten(),
849            new_actors,
850            PartialGraphBarrierInfo::new(
851                first_create_info.map_or_else(
852                    PostCollectCommand::barrier,
853                    CreateSnapshotBackfillJobCommandInfo::into_post_collect,
854                ),
855                barrier_info,
856                notifier,
857                state_table_ids.clone(),
858            ),
859        )?;
860        Ok(())
861    }
862
863    pub(crate) fn start_consume_upstream(
864        &mut self,
865        partial_graph_manager: &mut PartialGraphManager,
866        barrier_info: &BarrierInfo,
867    ) -> MetaResult<CreatingJobInfo> {
868        info!(
869            job_id = %self.job_id,
870            prev_epoch = barrier_info.prev_epoch(),
871            "start consuming upstream"
872        );
873        let info = self.status.start_consume_upstream(barrier_info);
874        Self::inject_barrier(
875            self.partial_graph_id,
876            partial_graph_manager,
877            &self.node_actors,
878            &self.state_table_ids,
879            true,
880            barrier_info.clone(),
881            None,
882            Some(Mutation::Stop(StopMutation {
883                // stop all actors
884                actors: info
885                    .fragment_infos
886                    .values()
887                    .flat_map(|info| info.actors.keys().copied())
888                    .collect(),
889                dropped_sink_fragments: vec![], // not related to sink-into-table
890            })),
891            None, // no notifier when start consuming upstream
892            None,
893        )?;
894        Ok(info)
895    }
896
897    pub(crate) fn on_new_upstream_barrier(
898        &mut self,
899        partial_graph_manager: &mut PartialGraphManager,
900        barrier_info: &BarrierInfo,
901        mutation: Option<(Mutation, Option<&mut NotifierStarter>)>,
902    ) -> MetaResult<()> {
903        let progress_epoch = if let Some(max_committed_epoch) = self.max_committed_epoch {
904            max(max_committed_epoch, self.snapshot_epoch)
905        } else {
906            self.snapshot_epoch
907        };
908        self.upstream_lag.set(
909            barrier_info
910                .prev_epoch
911                .value()
912                .0
913                .saturating_sub(progress_epoch) as _,
914        );
915        let (mut mutation, mut notifier) = match mutation {
916            Some((mutation, notifier)) => (Some(mutation), notifier),
917            None => (None, None),
918        };
919        for (barrier_to_inject, mutation) in self.status.on_new_upstream_epoch(
920            partial_graph_manager,
921            self.partial_graph_id,
922            self.max_pending_barrier_num,
923            barrier_info,
924            mutation.take(),
925        ) {
926            Self::inject_barrier(
927                self.partial_graph_id,
928                partial_graph_manager,
929                &self.node_actors,
930                &self.state_table_ids,
931                false,
932                barrier_to_inject,
933                None,
934                mutation,
935                notifier.take(),
936                None,
937            )?;
938        }
939        Ok(())
940    }
941
942    pub(crate) fn pre_apply_throttle(
943        &mut self,
944        config: &mut ThrottleConfigMap,
945    ) -> Option<Mutation> {
946        self.status.pre_apply_throttle(config)
947    }
948
949    /// Returns whether the next barrier should be forced to a checkpoint.
950    pub(crate) fn collect(&mut self, collected_barrier: CollectedBarrier<'_>) -> bool {
951        let pending_barrier_num = collected_barrier.pending_barrier_num;
952        self.status.update_progress(
953            collected_barrier
954                .resps
955                .values()
956                .flat_map(|resp| &resp.create_mview_progress),
957        );
958        self.is_ready_to_merge() && pending_barrier_num <= self.max_lagged_barrier_num
959    }
960
961    fn is_ready_to_merge(&self) -> bool {
962        if let CreatingStreamingJobStatus::ConsumingLogStore {
963            log_store_progress_tracker,
964            pending_barriers,
965            ..
966        } = &self.status
967            && pending_barriers.is_empty()
968            && log_store_progress_tracker.is_finished()
969        {
970            true
971        } else {
972            false
973        }
974    }
975
976    pub(crate) fn should_merge_to_upstream(
977        &self,
978        partial_graph_manager: &PartialGraphManager,
979    ) -> bool {
980        if !self.is_ready_to_merge() {
981            return false;
982        }
983
984        // A job that is ready to merge has finished initialization and is not resetting, so its
985        // partial graph must be running.
986        partial_graph_manager.pending_barrier_num(self.partial_graph_id)
987            <= self.max_lagged_barrier_num
988    }
989}
990
991impl CreatingStreamingJobControl {
992    pub(crate) fn start_completing(
993        &mut self,
994        partial_graph_manager: &mut PartialGraphManager,
995        min_upstream_inflight_epoch: Option<u64>,
996    ) -> Option<(
997        u64,
998        HashMap<WorkerId, BarrierCompleteResponse>,
999        PartialGraphBarrierInfo,
1000        bool,
1001    )> {
1002        let (finished_at_epoch, epoch_end_bound) = match &self.status {
1003            CreatingStreamingJobStatus::Finishing(finish_at_epoch, _) => {
1004                let epoch_end_bound = min_upstream_inflight_epoch
1005                    .map(|upstream_epoch| {
1006                        if upstream_epoch < *finish_at_epoch {
1007                            Excluded(upstream_epoch)
1008                        } else {
1009                            Unbounded
1010                        }
1011                    })
1012                    .unwrap_or(Unbounded);
1013                (Some(*finish_at_epoch), epoch_end_bound)
1014            }
1015            CreatingStreamingJobStatus::ConsumingSnapshot { .. }
1016            | CreatingStreamingJobStatus::ConsumingLogStore { .. } => (
1017                None,
1018                min_upstream_inflight_epoch
1019                    .map(Excluded)
1020                    .unwrap_or(Unbounded),
1021            ),
1022            CreatingStreamingJobStatus::PlaceHolder => {
1023                unreachable!()
1024            }
1025        };
1026        partial_graph_manager
1027            .start_completing(
1028                self.partial_graph_id,
1029                epoch_end_bound,
1030                |non_checkpoint_epoch, _, _| {
1031                    if let Some(finish_at_epoch) = finished_at_epoch {
1032                        assert!(non_checkpoint_epoch.prev < finish_at_epoch);
1033                    }
1034                },
1035            )
1036            .map(|(epoch, resps, info)| {
1037                let is_finish_epoch = if let Some(finish_at_epoch) = finished_at_epoch {
1038                    assert!(!info.post_collect_command.should_checkpoint());
1039                    if epoch == finish_at_epoch {
1040                        // TODO: can early remove partial graph here
1041                        self.ack_completed(partial_graph_manager, epoch);
1042                        true
1043                    } else {
1044                        false
1045                    }
1046                } else {
1047                    false
1048                };
1049                (epoch, resps, info, is_finish_epoch)
1050            })
1051    }
1052
1053    pub(super) fn ack_completed(
1054        &mut self,
1055        partial_graph_manager: &mut PartialGraphManager,
1056        completed_epoch: u64,
1057    ) {
1058        match &self.status {
1059            CreatingStreamingJobStatus::ConsumingSnapshot { .. }
1060            | CreatingStreamingJobStatus::ConsumingLogStore { .. }
1061            | CreatingStreamingJobStatus::Finishing(_, _) => {
1062                partial_graph_manager.ack_completed(self.partial_graph_id, completed_epoch);
1063                if let Some(prev_max_committed_epoch) =
1064                    self.max_committed_epoch.replace(completed_epoch)
1065                {
1066                    assert!(completed_epoch > prev_max_committed_epoch);
1067                }
1068            }
1069            CreatingStreamingJobStatus::PlaceHolder => {
1070                unreachable!()
1071            }
1072        }
1073    }
1074
1075    pub(crate) fn fragment_infos(&self) -> Option<&HashMap<FragmentId, InflightFragmentInfo>> {
1076        self.status.fragment_infos()
1077    }
1078
1079    pub fn into_tracking_job(self) -> TrackingJob {
1080        match self.status {
1081            CreatingStreamingJobStatus::ConsumingSnapshot { .. }
1082            | CreatingStreamingJobStatus::ConsumingLogStore { .. }
1083            | CreatingStreamingJobStatus::PlaceHolder => {
1084                unreachable!("expect finish")
1085            }
1086            CreatingStreamingJobStatus::Finishing(_, tracking_job) => tracking_job,
1087        }
1088    }
1089
1090    /// Whether the job can be dropped by resetting its independent partial graph.
1091    ///
1092    /// A finishing job has already been merged into the database graph, so it must be handled by
1093    /// the database-graph drop command instead.
1094    pub(super) fn can_drop_independently(&self) -> bool {
1095        match &self.status {
1096            CreatingStreamingJobStatus::ConsumingSnapshot { .. }
1097            | CreatingStreamingJobStatus::ConsumingLogStore { .. } => true,
1098            CreatingStreamingJobStatus::Finishing(_, _) => false,
1099            CreatingStreamingJobStatus::PlaceHolder => {
1100                unreachable!()
1101            }
1102        }
1103    }
1104}
1105
1106#[cfg(test)]
1107mod tests {
1108    use super::*;
1109
1110    #[test]
1111    fn test_snapshot_backfill_max_pending_barrier_num() {
1112        let mut opts = MetaOpts::test(false);
1113        opts.in_flight_barrier_nums = 10;
1114
1115        opts.snapshot_backfill_barrier_amplification_factor = 0;
1116        assert_eq!(snapshot_backfill_max_pending_barrier_num(&opts), 10);
1117
1118        opts.snapshot_backfill_barrier_amplification_factor = 1;
1119        assert_eq!(snapshot_backfill_max_pending_barrier_num(&opts), 10);
1120
1121        opts.snapshot_backfill_barrier_amplification_factor = 10;
1122        assert_eq!(snapshot_backfill_max_pending_barrier_num(&opts), 100);
1123
1124        opts.in_flight_barrier_nums = usize::MAX;
1125        assert_eq!(snapshot_backfill_max_pending_barrier_num(&opts), usize::MAX);
1126    }
1127
1128    #[test]
1129    fn test_resolve_since_timestamp_upstream_log_epochs() {
1130        let upstream_log_epochs = vec![(vec![45, 50], 55)];
1131
1132        let (initial_barrier, barriers) =
1133            CreatingStreamingJobControl::resolve_since_timestamp_upstream_log_epochs(
1134                &upstream_log_epochs,
1135                [].iter(),
1136                40,
1137                60,
1138            )
1139            .unwrap();
1140
1141        assert_eq!(
1142            (initial_barrier.prev_epoch(), initial_barrier.curr_epoch()),
1143            (40, 45)
1144        );
1145        assert!(initial_barrier.kind.is_checkpoint());
1146        assert_eq!(
1147            barriers
1148                .iter()
1149                .map(|barrier| (barrier.prev_epoch(), barrier.curr_epoch()))
1150                .collect::<Vec<_>>(),
1151            vec![(45, 50), (50, 55), (55, 60)]
1152        );
1153        assert_eq!(
1154            barriers
1155                .iter()
1156                .map(|barrier| match &barrier.kind {
1157                    BarrierKind::Checkpoint(epochs) => Some(epochs.clone()),
1158                    _ => None,
1159                })
1160                .collect::<Vec<_>>(),
1161            vec![None, None, Some(vec![45, 50, 55])]
1162        );
1163    }
1164
1165    #[test]
1166    fn test_resolve_since_timestamp_upstream_log_epochs_with_pending_barriers() {
1167        let upstream_log_epochs = vec![(vec![45, 50], 55)];
1168        let pending_upstream_barriers = [
1169            BarrierInfo {
1170                prev_epoch: TracedEpoch::new(Epoch(60)),
1171                curr_epoch: TracedEpoch::new(Epoch(65)),
1172                kind: BarrierKind::Barrier,
1173            },
1174            BarrierInfo {
1175                prev_epoch: TracedEpoch::new(Epoch(65)),
1176                curr_epoch: TracedEpoch::new(Epoch(70)),
1177                kind: BarrierKind::Checkpoint(vec![60, 65]),
1178            },
1179        ];
1180
1181        let (initial_barrier, barriers) =
1182            CreatingStreamingJobControl::resolve_since_timestamp_upstream_log_epochs(
1183                &upstream_log_epochs,
1184                pending_upstream_barriers.iter(),
1185                40,
1186                70,
1187            )
1188            .unwrap();
1189
1190        assert_eq!(
1191            (initial_barrier.prev_epoch(), initial_barrier.curr_epoch()),
1192            (40, 45)
1193        );
1194        assert!(initial_barrier.kind.is_checkpoint());
1195        assert_eq!(
1196            barriers
1197                .iter()
1198                .map(|barrier| (barrier.prev_epoch(), barrier.curr_epoch()))
1199                .collect::<Vec<_>>(),
1200            vec![(45, 50), (50, 55), (55, 60), (60, 65), (65, 70)]
1201        );
1202        assert_eq!(
1203            barriers
1204                .iter()
1205                .map(|barrier| match &barrier.kind {
1206                    BarrierKind::Checkpoint(epochs) => Some(epochs.clone()),
1207                    _ => None,
1208                })
1209                .collect::<Vec<_>>(),
1210            vec![None, None, Some(vec![45, 50, 55]), None, Some(vec![60, 65])]
1211        );
1212    }
1213
1214    #[test]
1215    fn test_resolve_since_timestamp_upstream_log_epochs_with_gap_before_pending_barriers() {
1216        let upstream_log_epochs = vec![(vec![61, 62, 63, 64], 65)];
1217        let pending_upstream_barriers = [
1218            BarrierInfo {
1219                prev_epoch: TracedEpoch::new(Epoch(66)),
1220                curr_epoch: TracedEpoch::new(Epoch(67)),
1221                kind: BarrierKind::Barrier,
1222            },
1223            BarrierInfo {
1224                prev_epoch: TracedEpoch::new(Epoch(67)),
1225                curr_epoch: TracedEpoch::new(Epoch(68)),
1226                kind: BarrierKind::Barrier,
1227            },
1228            BarrierInfo {
1229                prev_epoch: TracedEpoch::new(Epoch(68)),
1230                curr_epoch: TracedEpoch::new(Epoch(69)),
1231                kind: BarrierKind::Barrier,
1232            },
1233            BarrierInfo {
1234                prev_epoch: TracedEpoch::new(Epoch(69)),
1235                curr_epoch: TracedEpoch::new(Epoch(70)),
1236                kind: BarrierKind::Barrier,
1237            },
1238        ];
1239
1240        let (initial_barrier, barriers) =
1241            CreatingStreamingJobControl::resolve_since_timestamp_upstream_log_epochs(
1242                &upstream_log_epochs,
1243                pending_upstream_barriers.iter(),
1244                60,
1245                70,
1246            )
1247            .unwrap();
1248
1249        assert_eq!(
1250            (initial_barrier.prev_epoch(), initial_barrier.curr_epoch()),
1251            (60, 61)
1252        );
1253        assert!(initial_barrier.kind.is_checkpoint());
1254        assert_eq!(
1255            barriers
1256                .iter()
1257                .map(|barrier| (barrier.prev_epoch(), barrier.curr_epoch()))
1258                .collect::<Vec<_>>(),
1259            vec![
1260                (61, 62),
1261                (62, 63),
1262                (63, 64),
1263                (64, 65),
1264                (65, 66),
1265                (66, 67),
1266                (67, 68),
1267                (68, 69),
1268                (69, 70)
1269            ]
1270        );
1271        assert_eq!(
1272            barriers
1273                .iter()
1274                .map(|barrier| match &barrier.kind {
1275                    BarrierKind::Checkpoint(epochs) => Some(epochs.clone()),
1276                    _ => None,
1277                })
1278                .collect::<Vec<_>>(),
1279            vec![
1280                None,
1281                None,
1282                None,
1283                None,
1284                Some(vec![61, 62, 63, 64, 65]),
1285                None,
1286                None,
1287                None,
1288                None
1289            ]
1290        );
1291    }
1292
1293    #[test]
1294    fn test_resolve_since_timestamp_upstream_log_epochs_without_pending_barriers() {
1295        let upstream_log_epochs = vec![(vec![61, 62, 63, 64], 65)];
1296
1297        let (initial_barrier, barriers) =
1298            CreatingStreamingJobControl::resolve_since_timestamp_upstream_log_epochs(
1299                &upstream_log_epochs,
1300                [].iter(),
1301                60,
1302                66,
1303            )
1304            .unwrap();
1305
1306        assert_eq!(
1307            (initial_barrier.prev_epoch(), initial_barrier.curr_epoch()),
1308            (60, 61)
1309        );
1310        assert!(initial_barrier.kind.is_checkpoint());
1311        assert_eq!(
1312            barriers
1313                .iter()
1314                .map(|barrier| (barrier.prev_epoch(), barrier.curr_epoch()))
1315                .collect::<Vec<_>>(),
1316            vec![(61, 62), (62, 63), (63, 64), (64, 65), (65, 66)]
1317        );
1318        assert_eq!(
1319            barriers
1320                .iter()
1321                .map(|barrier| match &barrier.kind {
1322                    BarrierKind::Checkpoint(epochs) => Some(epochs.clone()),
1323                    _ => None,
1324                })
1325                .collect::<Vec<_>>(),
1326            vec![None, None, None, None, Some(vec![61, 62, 63, 64, 65])]
1327        );
1328    }
1329}