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