1use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
16use std::ops::{AddAssign, Deref};
17use std::sync::Arc;
18
19use itertools::Itertools;
20use risingwave_common::bitmap::Bitmap;
21use risingwave_common::catalog::{FragmentTypeFlag, FragmentTypeMask, TableId};
22use risingwave_common::hash::{IsSingleton, VirtualNode, VnodeCount, VnodeCountCompat};
23use risingwave_common::id::JobId;
24use risingwave_common::util::stream_graph_visitor::{self, visit_stream_node_body};
25use risingwave_meta_model::{DispatcherType, SourceId, StreamingParallelism, WorkerId, fragment};
26use risingwave_pb::catalog::Table;
27use risingwave_pb::common::ActorInfo;
28use risingwave_pb::id::SubscriberId;
29use risingwave_pb::meta::table_fragments::fragment::{
30 FragmentDistributionType, PbFragmentDistributionType,
31};
32use risingwave_pb::meta::table_fragments::{PbActorStatus, PbFragment, State};
33use risingwave_pb::meta::table_parallelism::{
34 FixedParallelism, Parallelism, PbAdaptiveParallelism, PbCustomParallelism, PbFixedParallelism,
35 PbParallelism,
36};
37use risingwave_pb::meta::{PbTableFragments, PbTableParallelism};
38use risingwave_pb::plan_common::PbExprContext;
39use risingwave_pb::stream_plan::stream_node::NodeBody;
40use risingwave_pb::stream_plan::{
41 DispatchStrategy, Dispatcher, PbDispatchOutputMapping, PbDispatcher, PbStreamActor,
42 PbStreamContext, StreamNode,
43};
44use strum::Display;
45
46use super::{ActorId, FragmentId};
47
48#[derive(Debug, Copy, Clone, Eq, PartialEq)]
50pub enum TableParallelism {
51 Adaptive,
53 Fixed(usize),
56 Custom,
63}
64
65impl From<PbTableParallelism> for TableParallelism {
66 fn from(value: PbTableParallelism) -> Self {
67 use Parallelism::*;
68 match &value.parallelism {
69 Some(Fixed(FixedParallelism { parallelism: n })) => Self::Fixed(*n as usize),
70 Some(Adaptive(_)) | Some(Auto(_)) => Self::Adaptive,
71 Some(Custom(_)) => Self::Custom,
72 _ => unreachable!(),
73 }
74 }
75}
76
77impl From<TableParallelism> for PbTableParallelism {
78 fn from(value: TableParallelism) -> Self {
79 use TableParallelism::*;
80
81 let parallelism = match value {
82 Adaptive => PbParallelism::Adaptive(PbAdaptiveParallelism {}),
83 Fixed(n) => PbParallelism::Fixed(PbFixedParallelism {
84 parallelism: n as u32,
85 }),
86 Custom => PbParallelism::Custom(PbCustomParallelism {}),
87 };
88
89 Self {
90 parallelism: Some(parallelism),
91 }
92 }
93}
94
95impl From<StreamingParallelism> for TableParallelism {
96 fn from(value: StreamingParallelism) -> Self {
97 match value {
98 StreamingParallelism::Adaptive => TableParallelism::Adaptive,
99 StreamingParallelism::Fixed(n) => TableParallelism::Fixed(n),
100 StreamingParallelism::Custom => TableParallelism::Custom,
101 }
102 }
103}
104
105impl From<TableParallelism> for StreamingParallelism {
106 fn from(value: TableParallelism) -> Self {
107 match value {
108 TableParallelism::Adaptive => StreamingParallelism::Adaptive,
109 TableParallelism::Fixed(n) => StreamingParallelism::Fixed(n),
110 TableParallelism::Custom => StreamingParallelism::Custom,
111 }
112 }
113}
114
115pub type ActorUpstreams = BTreeMap<FragmentId, HashMap<ActorId, ActorInfo>>;
116pub type StreamActorWithDispatchers = (StreamActor, Vec<PbDispatcher>);
117pub type StreamActorWithUpDownstreams = (StreamActor, ActorUpstreams, Vec<PbDispatcher>);
118pub type FragmentActorDispatchers = HashMap<FragmentId, HashMap<ActorId, Vec<PbDispatcher>>>;
119
120pub type FragmentDownstreamRelation = HashMap<FragmentId, Vec<DownstreamFragmentRelation>>;
121pub type FragmentReplaceUpstream = HashMap<FragmentId, HashMap<FragmentId, FragmentId>>;
123pub type ActorNewNoShuffle = HashMap<FragmentId, HashMap<FragmentId, HashMap<ActorId, ActorId>>>;
126
127#[derive(Debug, Clone)]
128pub struct DownstreamFragmentRelation {
129 pub downstream_fragment_id: FragmentId,
130 pub dispatcher_type: DispatcherType,
131 pub dist_key_indices: Vec<u32>,
132 pub output_mapping: PbDispatchOutputMapping,
133}
134
135impl From<(FragmentId, DispatchStrategy)> for DownstreamFragmentRelation {
136 fn from((fragment_id, dispatch): (FragmentId, DispatchStrategy)) -> Self {
137 Self {
138 downstream_fragment_id: fragment_id,
139 dispatcher_type: dispatch.get_type().unwrap().into(),
140 dist_key_indices: dispatch.dist_key_indices,
141 output_mapping: dispatch.output_mapping.unwrap(),
142 }
143 }
144}
145
146#[derive(Debug, Clone)]
147pub struct StreamJobFragmentsToCreate {
148 pub inner: StreamJobFragments,
149 pub downstreams: FragmentDownstreamRelation,
150}
151
152impl Deref for StreamJobFragmentsToCreate {
153 type Target = StreamJobFragments;
154
155 fn deref(&self) -> &Self::Target {
156 &self.inner
157 }
158}
159
160#[derive(Clone, Debug)]
161pub struct StreamActor {
162 pub actor_id: ActorId,
163 pub fragment_id: FragmentId,
164 pub vnode_bitmap: Option<Bitmap>,
165 pub mview_definition: String,
166 pub expr_context: Option<PbExprContext>,
167 pub config_override: Arc<str>,
169}
170
171impl StreamActor {
172 fn to_protobuf(&self, dispatchers: impl Iterator<Item = Dispatcher>) -> PbStreamActor {
173 PbStreamActor {
174 actor_id: self.actor_id,
175 fragment_id: self.fragment_id,
176 dispatcher: dispatchers.collect(),
177 vnode_bitmap: self
178 .vnode_bitmap
179 .as_ref()
180 .map(|bitmap| bitmap.to_protobuf()),
181 mview_definition: self.mview_definition.clone(),
182 expr_context: self.expr_context.clone(),
183 config_override: self.config_override.to_string(),
184 }
185 }
186}
187
188#[derive(Clone, Debug, Default)]
189pub struct Fragment {
190 pub fragment_id: FragmentId,
191 pub fragment_type_mask: FragmentTypeMask,
192 pub distribution_type: PbFragmentDistributionType,
193 pub state_table_ids: Vec<TableId>,
194 pub maybe_vnode_count: Option<u32>,
195 pub nodes: StreamNode,
196}
197
198impl Fragment {
199 pub fn to_protobuf(
200 &self,
201 actors: &[StreamActor],
202 upstream_fragments: impl Iterator<Item = FragmentId>,
203 dispatchers: Option<&HashMap<ActorId, Vec<Dispatcher>>>,
204 ) -> PbFragment {
205 PbFragment {
206 fragment_id: self.fragment_id,
207 fragment_type_mask: self.fragment_type_mask.into(),
208 distribution_type: self.distribution_type as _,
209 actors: actors
210 .iter()
211 .map(|actor| {
212 actor.to_protobuf(
213 dispatchers
214 .and_then(|dispatchers| dispatchers.get(&actor.actor_id))
215 .into_iter()
216 .flatten()
217 .cloned(),
218 )
219 })
220 .collect(),
221 state_table_ids: self.state_table_ids.clone(),
222 upstream_fragment_ids: upstream_fragments.collect(),
223 maybe_vnode_count: self.maybe_vnode_count,
224 nodes: Some(self.nodes.clone()),
225 }
226 }
227}
228
229impl VnodeCountCompat for Fragment {
230 fn vnode_count_inner(&self) -> VnodeCount {
231 VnodeCount::from_protobuf(self.maybe_vnode_count, || self.is_singleton())
232 }
233}
234
235impl IsSingleton for Fragment {
236 fn is_singleton(&self) -> bool {
237 matches!(self.distribution_type, FragmentDistributionType::Single)
238 }
239}
240
241impl From<fragment::Model> for Fragment {
242 fn from(model: fragment::Model) -> Self {
243 Self {
244 fragment_id: model.fragment_id,
245 fragment_type_mask: FragmentTypeMask::from(model.fragment_type_mask),
246 distribution_type: model.distribution_type.into(),
247 state_table_ids: model.state_table_ids.into_inner(),
248 maybe_vnode_count: VnodeCount::set(model.vnode_count).to_protobuf(),
249 nodes: model.stream_node.to_protobuf(),
250 }
251 }
252}
253
254#[derive(Debug, Clone)]
260pub struct StreamJobFragments {
261 pub stream_job_id: JobId,
263
264 pub state: State,
266
267 pub fragments: BTreeMap<FragmentId, Fragment>,
269
270 pub ctx: StreamContext,
272
273 pub max_parallelism: usize,
284}
285
286#[derive(Debug, Clone, Default)]
287pub struct StreamContext {
288 pub timezone: Option<String>,
290
291 pub config_override: Arc<str>,
293}
294
295impl StreamContext {
296 pub fn to_protobuf(&self) -> PbStreamContext {
297 PbStreamContext {
298 timezone: self.timezone.clone().unwrap_or("".into()),
299 config_override: self.config_override.to_string(),
300 }
301 }
302
303 pub fn to_expr_context(&self) -> PbExprContext {
304 PbExprContext {
305 time_zone: self.timezone.clone().unwrap_or("Empty Time Zone".into()),
307 strict_mode: false,
308 }
309 }
310
311 pub fn from_protobuf(prost: &PbStreamContext) -> Self {
312 Self {
313 timezone: if prost.get_timezone().is_empty() {
314 None
315 } else {
316 Some(prost.get_timezone().clone())
317 },
318 config_override: prost.get_config_override().as_str().into(),
319 }
320 }
321}
322
323#[easy_ext::ext(StreamingJobModelContextExt)]
324impl risingwave_meta_model::streaming_job::Model {
325 pub fn stream_context(&self) -> StreamContext {
326 StreamContext {
327 timezone: self.timezone.clone(),
328 config_override: self.config_override.clone().unwrap_or_default().into(),
329 }
330 }
331}
332
333impl StreamJobFragments {
334 pub fn to_protobuf(
335 &self,
336 fragment_actors: &HashMap<FragmentId, Vec<StreamActor>>,
337 fragment_upstreams: &HashMap<FragmentId, HashSet<FragmentId>>,
338 fragment_dispatchers: &FragmentActorDispatchers,
339 actor_status: HashMap<ActorId, PbActorStatus>,
340 ) -> PbTableFragments {
341 PbTableFragments {
342 table_id: self.stream_job_id,
343 state: self.state as _,
344 fragments: self
345 .fragments
346 .iter()
347 .map(|(id, fragment)| {
348 let actors = fragment_actors.get(id).map(|a| a.as_slice()).unwrap_or(&[]);
349 (
350 *id,
351 fragment.to_protobuf(
352 actors,
353 fragment_upstreams.get(id).into_iter().flatten().cloned(),
354 fragment_dispatchers.get(id),
355 ),
356 )
357 })
358 .collect(),
359 actor_status,
360 ctx: Some(self.ctx.to_protobuf()),
361 node_label: "".to_owned(),
362 backfill_done: true,
363 max_parallelism: Some(self.max_parallelism as _),
364 }
365 }
366}
367
368#[cfg(test)]
369mod tests {
370 use super::*;
371
372 #[test]
373 fn test_stream_context_protobuf_round_trip() {
374 let context = StreamContext {
375 timezone: Some("UTC".to_owned()),
376 config_override: "{\"a\":1}".into(),
377 };
378
379 let prost = context.to_protobuf();
380
381 let round_trip = StreamContext::from_protobuf(&prost);
382 assert_eq!(round_trip.timezone, Some("UTC".to_owned()));
383 assert_eq!(&*round_trip.config_override, "{\"a\":1}");
384 }
385
386 #[test]
387 fn test_stream_context_from_model() {
388 let model = risingwave_meta_model::streaming_job::Model {
389 job_id: 1.into(),
390 job_status: risingwave_meta_model::JobStatus::Created,
391 create_type: risingwave_meta_model::CreateType::Foreground,
392 timezone: Some("Asia/Shanghai".to_owned()),
393 config_override: Some("{\"parallelism\":2}".to_owned()),
394 adaptive_parallelism_strategy: Some("AUTO".to_owned()),
395 parallelism: StreamingParallelism::Adaptive,
396 backfill_parallelism: Some(StreamingParallelism::Adaptive),
397 backfill_adaptive_parallelism_strategy: Some("RATIO(0.25)".to_owned()),
398 backfill_orders: None,
399 max_parallelism: 32,
400 specific_resource_group: None,
401 is_serverless_backfill: false,
402 refresh_interval_sec: None,
403 };
404
405 let context = model.stream_context();
406 assert_eq!(context.timezone, Some("Asia/Shanghai".to_owned()));
407 assert_eq!(&*context.config_override, "{\"parallelism\":2}");
408 }
409
410 fn fragment_type_mask(flags: impl IntoIterator<Item = FragmentTypeFlag>) -> FragmentTypeMask {
411 let mut mask = FragmentTypeMask::empty();
412 for flag in flags {
413 mask.add(flag);
414 }
415 mask
416 }
417
418 #[test]
419 fn test_tracking_progress_skips_only_cdc_fragments() {
420 let nodes = StreamNode::default();
421 let fragments = [
422 (
423 fragment_type_mask([FragmentTypeFlag::CdcFilter]),
424 &nodes,
425 vec![ActorId::new(1)].into_iter(),
426 ),
427 (
428 fragment_type_mask([
429 FragmentTypeFlag::StreamCdcScan,
430 FragmentTypeFlag::StreamScan,
431 ]),
432 &nodes,
433 vec![ActorId::new(2)].into_iter(),
434 ),
435 (
436 fragment_type_mask([FragmentTypeFlag::StreamScan]),
437 &nodes,
438 vec![ActorId::new(3), ActorId::new(4)].into_iter(),
439 ),
440 (
441 fragment_type_mask([FragmentTypeFlag::SourceScan]),
442 &nodes,
443 vec![ActorId::new(5)].into_iter(),
444 ),
445 ];
446
447 assert_eq!(
448 StreamJobFragments::tracking_progress_actor_ids_impl(fragments),
449 vec![
450 (ActorId::new(3), BackfillUpstreamType::MView),
451 (ActorId::new(4), BackfillUpstreamType::MView),
452 (ActorId::new(5), BackfillUpstreamType::Source),
453 ]
454 );
455 }
456}
457
458pub type StreamJobActorsToCreate = HashMap<
459 WorkerId,
460 HashMap<
461 FragmentId,
462 (
463 StreamNode,
464 Vec<StreamActorWithUpDownstreams>,
465 HashSet<SubscriberId>,
466 ),
467 >,
468>;
469
470impl StreamJobFragments {
471 pub fn for_test(job_id: JobId, fragments: BTreeMap<FragmentId, Fragment>) -> Self {
473 Self::new(
474 job_id,
475 fragments,
476 StreamContext::default(),
477 VirtualNode::COUNT_FOR_TEST,
478 )
479 }
480
481 pub fn new(
483 stream_job_id: JobId,
484 fragments: BTreeMap<FragmentId, Fragment>,
485 ctx: StreamContext,
486 max_parallelism: usize,
487 ) -> Self {
488 Self {
489 stream_job_id,
490 state: State::Initial,
491 fragments,
492 ctx,
493 max_parallelism,
494 }
495 }
496
497 pub fn fragment_ids(&self) -> impl Iterator<Item = FragmentId> + '_ {
498 self.fragments.keys().cloned()
499 }
500
501 pub fn fragments(&self) -> impl Iterator<Item = &Fragment> {
502 self.fragments.values()
503 }
504
505 pub fn stream_job_id(&self) -> JobId {
507 self.stream_job_id
508 }
509
510 pub fn timezone(&self) -> Option<String> {
512 self.ctx.timezone.clone()
513 }
514
515 pub fn is_created(&self) -> bool {
517 self.state == State::Created
518 }
519
520 #[cfg(test)]
522 pub fn mview_fragment_ids(&self) -> Vec<FragmentId> {
523 self.fragments
524 .values()
525 .filter(move |fragment| {
526 fragment
527 .fragment_type_mask
528 .contains(FragmentTypeFlag::Mview)
529 })
530 .map(|fragment| fragment.fragment_id)
531 .collect()
532 }
533
534 pub fn tracking_progress_actor_ids_impl<'a>(
536 fragments: impl IntoIterator<
537 Item = (
538 FragmentTypeMask,
539 &'a StreamNode,
540 impl Iterator<Item = ActorId>,
541 ),
542 >,
543 ) -> Vec<(ActorId, BackfillUpstreamType)> {
544 let mut actor_ids = vec![];
545 for (fragment_type_mask, nodes, actors) in fragments {
546 if fragment_type_mask
547 .contains_any([FragmentTypeFlag::CdcFilter, FragmentTypeFlag::StreamCdcScan])
548 {
549 continue;
552 }
553 let mut has_cdc_scan = false;
557 if fragment_type_mask.contains(FragmentTypeFlag::StreamScan) {
558 stream_graph_visitor::visit_stream_node(nodes, |node| {
559 has_cdc_scan |=
560 matches!(node.node_body.as_ref(), Some(NodeBody::StreamCdcScan(_)));
561 });
562 }
563 if has_cdc_scan {
564 continue;
565 }
566 if fragment_type_mask.contains_any([
567 FragmentTypeFlag::Values,
568 FragmentTypeFlag::StreamScan,
569 FragmentTypeFlag::SourceScan,
570 FragmentTypeFlag::LocalityProvider,
571 ]) {
572 actor_ids.extend(actors.map(|actor_id| {
573 (
574 actor_id,
575 BackfillUpstreamType::from_fragment_type_mask(fragment_type_mask),
576 )
577 }));
578 }
579 }
580 actor_ids
581 }
582
583 pub fn root_fragment(&self) -> Option<Fragment> {
584 self.mview_fragment()
585 .or_else(|| self.sink_fragment())
586 .or_else(|| self.source_fragment())
587 }
588
589 pub fn mview_fragment(&self) -> Option<Fragment> {
591 self.fragments
592 .values()
593 .find(|fragment| {
594 fragment
595 .fragment_type_mask
596 .contains(FragmentTypeFlag::Mview)
597 })
598 .cloned()
599 }
600
601 pub fn source_fragment(&self) -> Option<Fragment> {
602 self.fragments
603 .values()
604 .find(|fragment| {
605 fragment
606 .fragment_type_mask
607 .contains(FragmentTypeFlag::Source)
608 })
609 .cloned()
610 }
611
612 pub fn sink_fragment(&self) -> Option<Fragment> {
613 self.fragments
614 .values()
615 .find(|fragment| fragment.fragment_type_mask.contains(FragmentTypeFlag::Sink))
616 .cloned()
617 }
618
619 pub fn stream_source_fragments(&self) -> HashMap<SourceId, BTreeSet<FragmentId>> {
622 let mut source_fragments = HashMap::new();
623
624 for fragment in self.fragments() {
625 {
626 if let Some(source_id) = fragment.nodes.find_stream_source() {
627 source_fragments
628 .entry(source_id)
629 .or_insert(BTreeSet::new())
630 .insert(fragment.fragment_id as FragmentId);
631 }
632 }
633 }
634 source_fragments
635 }
636
637 pub fn source_backfill_fragments(
638 &self,
639 ) -> HashMap<SourceId, BTreeSet<(FragmentId, FragmentId)>> {
640 Self::source_backfill_fragments_impl(
641 self.fragments
642 .iter()
643 .map(|(fragment_id, fragment)| (*fragment_id, &fragment.nodes)),
644 )
645 }
646
647 pub fn source_backfill_fragments_impl(
652 fragments: impl Iterator<Item = (FragmentId, &StreamNode)>,
653 ) -> HashMap<SourceId, BTreeSet<(FragmentId, FragmentId)>> {
654 let mut source_backfill_fragments = HashMap::new();
655
656 for (fragment_id, fragment_node) in fragments {
657 {
658 if let Some((source_id, upstream_source_fragment_id)) =
659 fragment_node.find_source_backfill()
660 {
661 source_backfill_fragments
662 .entry(source_id)
663 .or_insert(BTreeSet::new())
664 .insert((fragment_id, upstream_source_fragment_id));
665 }
666 }
667 }
668 source_backfill_fragments
669 }
670
671 pub fn union_fragment_for_table(&mut self) -> &mut Fragment {
674 let mut union_fragment_id = None;
675 for (fragment_id, fragment) in &self.fragments {
676 {
677 {
678 visit_stream_node_body(&fragment.nodes, |body| {
679 if let NodeBody::Union(_) = body {
680 if let Some(union_fragment_id) = union_fragment_id.as_mut() {
681 assert_eq!(*union_fragment_id, *fragment_id);
683 } else {
684 union_fragment_id = Some(*fragment_id);
685 }
686 }
687 })
688 }
689 }
690 }
691
692 let union_fragment_id =
693 union_fragment_id.expect("fragment of placeholder merger not found");
694
695 (self
696 .fragments
697 .get_mut(&union_fragment_id)
698 .unwrap_or_else(|| panic!("fragment {} not found", union_fragment_id))) as _
699 }
700
701 fn resolve_dependent_table(stream_node: &StreamNode, table_ids: &mut HashMap<TableId, usize>) {
703 let table_id = match stream_node.node_body.as_ref() {
704 Some(NodeBody::StreamScan(stream_scan)) => Some(stream_scan.table_id),
705 Some(NodeBody::StreamCdcScan(stream_scan)) => Some(stream_scan.table_id),
706 Some(NodeBody::LocalityProvider(state)) => {
707 Some(state.state_table.as_ref().expect("must have state").id)
708 }
709 _ => None,
710 };
711 if let Some(table_id) = table_id {
712 table_ids.entry(table_id).or_default().add_assign(1);
713 }
714
715 for child in &stream_node.input {
716 Self::resolve_dependent_table(child, table_ids);
717 }
718 }
719
720 pub fn upstream_table_counts(&self) -> HashMap<TableId, usize> {
721 Self::upstream_table_counts_impl(self.fragments.values().map(|fragment| &fragment.nodes))
722 }
723
724 pub fn upstream_table_counts_impl(
726 fragment_nodes: impl Iterator<Item = &StreamNode>,
727 ) -> HashMap<TableId, usize> {
728 let mut table_ids = HashMap::new();
729 fragment_nodes.for_each(|node| {
730 Self::resolve_dependent_table(node, &mut table_ids);
731 });
732
733 table_ids
734 }
735
736 pub fn mv_table_id(&self) -> Option<TableId> {
737 self.fragments
738 .values()
739 .flat_map(|f| f.state_table_ids.iter().copied())
740 .find(|table_id| self.stream_job_id.is_mv_table_id(*table_id))
741 }
742
743 pub fn collect_tables(fragments: impl Iterator<Item = &Fragment>) -> BTreeMap<TableId, Table> {
744 let mut tables = BTreeMap::new();
745 for fragment in fragments {
746 stream_graph_visitor::visit_stream_node_tables_inner(
747 &mut fragment.nodes.clone(),
748 false,
749 true,
750 |table, _| {
751 let table_id = table.id;
752 tables
753 .try_insert(table_id, table.clone())
754 .unwrap_or_else(|_| panic!("duplicated table id `{}`", table_id));
755 },
756 );
757 }
758 tables
759 }
760
761 pub fn internal_table_ids(&self) -> Vec<TableId> {
763 self.fragments
764 .values()
765 .flat_map(|f| f.state_table_ids.iter().copied())
766 .filter(|&t| !self.stream_job_id.is_mv_table_id(t))
767 .collect_vec()
768 }
769
770 pub fn all_table_ids(&self) -> impl Iterator<Item = TableId> + '_ {
772 self.fragments
773 .values()
774 .flat_map(|f| f.state_table_ids.clone())
775 }
776}
777
778#[derive(Debug, Display, Clone, Copy, PartialEq, Eq)]
779pub enum BackfillUpstreamType {
780 MView,
781 Values,
782 Source,
783 LocalityProvider,
784}
785
786impl BackfillUpstreamType {
787 pub fn from_fragment_type_mask(mask: FragmentTypeMask) -> Self {
788 let is_mview = mask.contains(FragmentTypeFlag::StreamScan);
789 let is_values = mask.contains(FragmentTypeFlag::Values);
790 let is_source = mask.contains(FragmentTypeFlag::SourceScan);
791 let is_locality_provider = mask.contains(FragmentTypeFlag::LocalityProvider);
792
793 debug_assert!(
796 is_mview as u8 + is_values as u8 + is_source as u8 + is_locality_provider as u8 == 1,
797 "a backfill fragment should either be mview, value, source, or locality provider, found {:?}",
798 mask
799 );
800
801 if is_mview {
802 BackfillUpstreamType::MView
803 } else if is_values {
804 BackfillUpstreamType::Values
805 } else if is_source {
806 BackfillUpstreamType::Source
807 } else if is_locality_provider {
808 BackfillUpstreamType::LocalityProvider
809 } else {
810 unreachable!("invalid fragment type mask: {:?}", mask);
811 }
812 }
813}