Skip to main content

risingwave_meta/barrier/
edge_builder.rs

1// Copyright 2025 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{HashMap, HashSet};
16use std::marker::PhantomData;
17use std::mem::take;
18
19use anyhow::anyhow;
20use risingwave_common::bitmap::Bitmap;
21use risingwave_meta_model::WorkerId;
22use risingwave_meta_model::fragment::DistributionType;
23use risingwave_pb::common::{ActorInfo, HostAddress};
24use risingwave_pb::id::{PartialGraphId, SubscriberId};
25use risingwave_pb::stream_plan::update_mutation::{DispatcherUpdate, MergeUpdate};
26use risingwave_pb::stream_plan::{AddMutation, PbDispatcher, StreamNode, UpdateMutation};
27use tracing::warn;
28
29use crate::MetaResult;
30use crate::barrier::rpc::ControlStreamManager;
31use crate::controller::fragment::InflightFragmentInfo;
32use crate::controller::utils::compose_dispatchers;
33use crate::model::{
34    ActorId, ActorNewNoShuffle, ActorUpstreams, DownstreamFragmentRelation, Fragment,
35    FragmentActorDispatchers, FragmentDownstreamRelation, FragmentId, StreamActor,
36    StreamJobActorsToCreate,
37};
38
39type ComposedEdge = (
40    HashMap<ActorId, PbDispatcher>,
41    HashMap<ActorId, ActorUpstreams>,
42    Option<HashMap<ActorId, ActorId>>,
43);
44
45/// Fragment information needed by [`FragmentEdgeBuilder`] to compute dispatchers and merge nodes.
46///
47/// Contains actor bitmaps and resolved host addresses.
48#[derive(Debug)]
49struct EdgeBuilderFragmentInfo {
50    distribution_type: DistributionType,
51    actors: HashMap<ActorId, Option<Bitmap>>,
52    actor_location: HashMap<ActorId, HostAddress>,
53    partial_graph_id: PartialGraphId,
54}
55
56#[derive(Debug)]
57enum FragmentStatus {
58    Existing(EdgeBuilderFragmentInfo),
59    New(EdgeBuilderFragmentInfo),
60    Changed {
61        before: EdgeBuilderFragmentInfo,
62        after: EdgeBuilderFragmentInfo,
63    },
64}
65
66impl EdgeBuilderFragmentInfo {
67    /// Build from an already-inflight fragment (actors already materialized).
68    fn from_inflight(
69        info: &InflightFragmentInfo,
70        partial_graph_id: PartialGraphId,
71        control_stream_manager: &ControlStreamManager,
72    ) -> Self {
73        let (actors, actor_location) = info
74            .actors
75            .iter()
76            .map(|(&actor_id, actor)| {
77                (
78                    (actor_id, actor.vnode_bitmap.clone()),
79                    (actor_id, control_stream_manager.host_addr(actor.worker_id)),
80                )
81            })
82            .unzip();
83        Self {
84            distribution_type: info.distribution_type,
85            actors,
86            actor_location,
87            partial_graph_id,
88        }
89    }
90
91    /// Build from a model `Fragment` with separately provided actors and locations.
92    fn from_fragment(
93        fragment: &Fragment,
94        stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
95        actor_worker: &HashMap<ActorId, WorkerId>,
96        partial_graph_id: PartialGraphId,
97        control_stream_manager: &ControlStreamManager,
98    ) -> Self {
99        let (actors, actor_location) = stream_actors
100            .get(&fragment.fragment_id)
101            .into_iter()
102            .flatten()
103            .map(|actor| {
104                (
105                    (actor.actor_id, actor.vnode_bitmap.clone()),
106                    (
107                        actor.actor_id,
108                        control_stream_manager.host_addr(actor_worker[&actor.actor_id]),
109                    ),
110                )
111            })
112            .unzip();
113        Self {
114            distribution_type: fragment.distribution_type.into(),
115            actors,
116            actor_location,
117            partial_graph_id,
118        }
119    }
120}
121
122#[derive(Debug)]
123pub(crate) struct FragmentEdgeBuildResult {
124    upstreams: HashMap<FragmentId, HashMap<ActorId, ActorUpstreams>>,
125    dispatchers: FragmentActorDispatchers,
126    merge_updates: HashMap<FragmentId, Vec<MergeUpdate>>,
127    dispatcher_updates: Vec<DispatcherUpdate>,
128}
129
130impl FragmentEdgeBuildResult {
131    fn validate_terminal_consumption(&self) {
132        let remaining_upstreams = self.upstreams.values().map(HashMap::len).sum::<usize>();
133        let remaining_dispatchers = self.dispatchers.values().map(HashMap::len).sum::<usize>();
134        let unapplied_dispatcher_updates = self.dispatcher_updates.len();
135        let unapplied_merge_updates = self.merge_updates.values().map(Vec::len).sum::<usize>();
136        let has_unconsumed_fields = remaining_upstreams != 0
137            || remaining_dispatchers != 0
138            || unapplied_dispatcher_updates != 0
139            || unapplied_merge_updates != 0;
140
141        debug_assert!(
142            !has_unconsumed_fields,
143            "edge result has unconsumed fields: remaining_upstreams={remaining_upstreams}, remaining_dispatchers={remaining_dispatchers}, unapplied_dispatcher_updates={unapplied_dispatcher_updates}, unapplied_merge_updates={unapplied_merge_updates}"
144        );
145        #[cfg(not(debug_assertions))]
146        if has_unconsumed_fields {
147            warn!(
148                remaining_upstreams,
149                remaining_dispatchers,
150                unapplied_dispatcher_updates,
151                unapplied_merge_updates,
152                "edge result has unconsumed fields"
153            );
154        }
155    }
156
157    /// Apply the remaining dispatchers after newly created actors have been collected.
158    pub(crate) fn apply_to_add_mutation(mut self, mutation: &mut AddMutation) {
159        let dispatchers = take(&mut self.dispatchers);
160        self.validate_terminal_consumption();
161        for (actor_id, dispatchers) in dispatchers.into_values().flatten() {
162            mutation
163                .actor_dispatchers
164                .entry(actor_id)
165                .or_default()
166                .dispatchers
167                .extend(dispatchers);
168        }
169    }
170
171    /// Apply the remaining edge updates after newly created actors have been collected.
172    pub(crate) fn apply_to_update_mutation(mut self, mutation: &mut UpdateMutation) {
173        let dispatchers = take(&mut self.dispatchers);
174        let dispatcher_updates = take(&mut self.dispatcher_updates);
175        let merge_updates = take(&mut self.merge_updates);
176        self.validate_terminal_consumption();
177        mutation.dispatcher_update.extend(dispatcher_updates);
178        mutation
179            .merge_update
180            .extend(merge_updates.into_values().flatten());
181        for (actor_id, dispatchers) in dispatchers.into_values().flatten() {
182            mutation
183                .actor_new_dispatchers
184                .entry(actor_id)
185                .or_default()
186                .dispatchers
187                .extend(dispatchers);
188        }
189    }
190
191    pub(crate) fn collect_actors_to_create(
192        &mut self,
193        actors: impl Iterator<
194            Item = (
195                FragmentId,
196                &StreamNode,
197                impl Iterator<Item = (&StreamActor, WorkerId)>,
198                impl IntoIterator<Item = SubscriberId>,
199            ),
200        >,
201    ) -> StreamJobActorsToCreate {
202        let mut actors_to_create = StreamJobActorsToCreate::default();
203        for (fragment_id, node, actors, subscriber_ids) in actors {
204            let subscriber_ids: HashSet<_> = subscriber_ids.into_iter().collect();
205            for (actor, worker_id) in actors {
206                let upstreams = self
207                    .upstreams
208                    .get_mut(&fragment_id)
209                    .and_then(|upstreams| upstreams.remove(&actor.actor_id))
210                    .unwrap_or_default();
211                let dispatchers = self
212                    .dispatchers
213                    .get_mut(&fragment_id)
214                    .and_then(|upstreams| upstreams.remove(&actor.actor_id))
215                    .unwrap_or_default();
216                actors_to_create
217                    .entry(worker_id)
218                    .or_default()
219                    .entry(fragment_id)
220                    .or_insert_with(|| (node.clone(), vec![], subscriber_ids.clone()))
221                    .1
222                    .push((actor.clone(), upstreams, dispatchers))
223            }
224        }
225        actors_to_create
226    }
227
228    /// Take the common upstream actor set for an edge to an existing fragment.
229    ///
230    /// This is used by sink-into-table, whose hash edge connects every downstream actor to the
231    /// same set of upstream sink actors. The downstream actors already exist, so their upstream
232    /// edge is installed through `new_upstream_sinks` rather than actor creation.
233    pub(crate) fn take_common_upstream_actors(
234        &mut self,
235        upstream_fragment_id: FragmentId,
236        downstream_fragment_id: FragmentId,
237    ) -> MetaResult<Vec<ActorInfo>> {
238        let downstream_actor_upstreams = self
239            .upstreams
240            .get_mut(&downstream_fragment_id)
241            .ok_or_else(|| {
242                anyhow!("cannot find upstreams for downstream fragment {downstream_fragment_id}")
243            })?;
244        let mut common_upstream_actors = None;
245        for (&downstream_actor_id, actor_upstreams) in downstream_actor_upstreams.iter_mut() {
246            let upstream_actors = actor_upstreams
247                .remove(&upstream_fragment_id)
248                .ok_or_else(|| {
249                    anyhow!(
250                        "cannot find upstream fragment {upstream_fragment_id} for downstream actor {downstream_actor_id}"
251                    )
252                })?;
253            if let Some(common) = &common_upstream_actors {
254                // `HashMap` equality is order-independent. Compare the complete `ActorInfo` values
255                // as well because matching actor IDs with different routing metadata do not form a
256                // common upstream layout.
257                if common != &upstream_actors {
258                    return Err(anyhow!(
259                        "upstream actors from fragment {upstream_fragment_id} are not common across downstream fragment {downstream_fragment_id}"
260                    )
261                    .into());
262                }
263            } else {
264                common_upstream_actors = Some(upstream_actors);
265            }
266        }
267        let common_upstream_actors = common_upstream_actors
268            .ok_or_else(|| anyhow!("downstream fragment {downstream_fragment_id} has no actors"))?;
269
270        downstream_actor_upstreams.retain(|_, actor_upstreams| !actor_upstreams.is_empty());
271        if downstream_actor_upstreams.is_empty() {
272            self.upstreams.remove(&downstream_fragment_id);
273        }
274
275        let mut common_upstream_actors: Vec<_> = common_upstream_actors.into_values().collect();
276        common_upstream_actors.sort_unstable_by_key(|actor| actor.actor_id);
277        Ok(common_upstream_actors)
278    }
279}
280
281pub(crate) struct RegisteringFragments;
282pub(crate) struct AddingRelations;
283
284pub(crate) struct FragmentEdgeBuilder<State> {
285    fragments: HashMap<FragmentId, FragmentStatus>,
286    result: FragmentEdgeBuildResult,
287    actor_new_no_shuffle: ActorNewNoShuffle,
288    _state: PhantomData<State>,
289}
290
291impl FragmentEdgeBuilder<RegisteringFragments> {
292    pub(crate) fn new() -> Self {
293        Self {
294            fragments: Default::default(),
295            result: FragmentEdgeBuildResult {
296                upstreams: Default::default(),
297                dispatchers: Default::default(),
298                merge_updates: Default::default(),
299                dispatcher_updates: Default::default(),
300            },
301            actor_new_no_shuffle: Default::default(),
302            _state: PhantomData,
303        }
304    }
305
306    fn add_existing_fragment_infos(
307        &mut self,
308        fragment_infos: impl IntoIterator<Item = (FragmentId, EdgeBuilderFragmentInfo)>,
309    ) {
310        for (fragment_id, info) in fragment_infos {
311            self.fragments
312                .try_insert(fragment_id, FragmentStatus::Existing(info))
313                .expect("non-duplicate");
314        }
315    }
316
317    pub(crate) fn add_existing_fragments<'a>(
318        mut self,
319        fragments: impl IntoIterator<Item = &'a InflightFragmentInfo>,
320        partial_graph_id: PartialGraphId,
321        control_stream_manager: &ControlStreamManager,
322    ) -> Self {
323        self.add_existing_fragment_infos(fragments.into_iter().map(|fragment| {
324            (
325                fragment.fragment_id,
326                EdgeBuilderFragmentInfo::from_inflight(
327                    fragment,
328                    partial_graph_id,
329                    control_stream_manager,
330                ),
331            )
332        }));
333        self
334    }
335
336    fn add_new_fragment_infos(
337        &mut self,
338        fragment_infos: impl IntoIterator<Item = (FragmentId, EdgeBuilderFragmentInfo)>,
339    ) {
340        for (fragment_id, info) in fragment_infos {
341            self.fragments
342                .try_insert(fragment_id, FragmentStatus::New(info))
343                .expect("new fragment must not already be registered");
344        }
345    }
346
347    pub(crate) fn add_new_fragments<'a>(
348        mut self,
349        fragments: impl IntoIterator<Item = &'a InflightFragmentInfo>,
350        partial_graph_id: PartialGraphId,
351        control_stream_manager: &ControlStreamManager,
352    ) -> Self {
353        self.add_new_fragment_infos(fragments.into_iter().map(|fragment| {
354            (
355                fragment.fragment_id,
356                EdgeBuilderFragmentInfo::from_inflight(
357                    fragment,
358                    partial_graph_id,
359                    control_stream_manager,
360                ),
361            )
362        }));
363        self
364    }
365
366    pub(crate) fn add_new_logical_fragments<'a>(
367        mut self,
368        fragments: impl IntoIterator<Item = (PartialGraphId, &'a Fragment)>,
369        stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
370        actor_worker: &HashMap<ActorId, WorkerId>,
371        control_stream_manager: &ControlStreamManager,
372    ) -> Self {
373        self.add_new_fragment_infos(fragments.into_iter().map(|(partial_graph_id, fragment)| {
374            (
375                fragment.fragment_id,
376                EdgeBuilderFragmentInfo::from_fragment(
377                    fragment,
378                    stream_actors,
379                    actor_worker,
380                    partial_graph_id,
381                    control_stream_manager,
382                ),
383            )
384        }));
385        self
386    }
387
388    fn replace_existing_fragment_actor_infos(
389        &mut self,
390        fragment_infos: impl IntoIterator<Item = (FragmentId, EdgeBuilderFragmentInfo)>,
391    ) {
392        for (fragment_id, after) in fragment_infos {
393            let status = self
394                .fragments
395                .remove(&fragment_id)
396                .expect("changed fragment must already be registered");
397            let FragmentStatus::Existing(before) = status else {
398                panic!("fragment {fragment_id} can only be changed once");
399            };
400            assert_eq!(
401                before.distribution_type, after.distribution_type,
402                "fragment distribution type cannot change"
403            );
404            for actor_id in before.actors.keys() {
405                assert!(
406                    !after.actors.contains_key(actor_id),
407                    "changed fragment {fragment_id} must use entirely fresh actor ids: retained {actor_id}"
408                );
409            }
410            self.fragments
411                .insert(fragment_id, FragmentStatus::Changed { before, after });
412        }
413    }
414
415    pub(crate) fn replace_existing_fragment_actors<'a>(
416        mut self,
417        fragments: impl IntoIterator<Item = &'a InflightFragmentInfo>,
418        partial_graph_id: PartialGraphId,
419        control_stream_manager: &ControlStreamManager,
420    ) -> Self {
421        self.replace_existing_fragment_actor_infos(fragments.into_iter().map(|fragment| {
422            (
423                fragment.fragment_id,
424                EdgeBuilderFragmentInfo::from_inflight(
425                    fragment,
426                    partial_graph_id,
427                    control_stream_manager,
428                ),
429            )
430        }));
431        self
432    }
433
434    pub(crate) fn finish_fragments(self) -> FragmentEdgeBuilder<AddingRelations> {
435        FragmentEdgeBuilder {
436            fragments: self.fragments,
437            result: self.result,
438            actor_new_no_shuffle: self.actor_new_no_shuffle,
439            _state: PhantomData,
440        }
441    }
442}
443
444impl FragmentEdgeBuilder<AddingRelations> {
445    pub(crate) fn add_relations(
446        mut self,
447        relations: &FragmentDownstreamRelation,
448    ) -> MetaResult<Self> {
449        for fragment_id in self.fragments.keys().copied().collect::<Vec<_>>() {
450            let Some(relations) = relations.get(&fragment_id) else {
451                continue;
452            };
453            for relation in relations {
454                if self
455                    .fragments
456                    .contains_key(&relation.downstream_fragment_id)
457                {
458                    self.add_edge_inner(fragment_id, relation)?;
459                }
460            }
461        }
462        Ok(self)
463    }
464
465    /// Add a relation that is not yet installed in the runtime graph.
466    pub(crate) fn add_edge(
467        mut self,
468        fragment_id: FragmentId,
469        downstream: &DownstreamFragmentRelation,
470    ) -> MetaResult<Self> {
471        self.add_edge_inner(fragment_id, downstream)?;
472        Ok(self)
473    }
474
475    fn add_edge_inner(
476        &mut self,
477        fragment_id: FragmentId,
478        downstream: &DownstreamFragmentRelation,
479    ) -> MetaResult<()> {
480        let Some(fragment_status) = self.fragments.get(&fragment_id) else {
481            if self
482                .fragments
483                .contains_key(&downstream.downstream_fragment_id)
484            {
485                return Err(anyhow!(
486                    "cannot find fragment {} with downstream {:?}",
487                    fragment_id,
488                    downstream
489                )
490                .into());
491            } else {
492                // ignore fragment relation with both upstream and downstream not in the set of fragments
493                return Ok(());
494            }
495        };
496        let Some(downstream_status) = self.fragments.get(&downstream.downstream_fragment_id) else {
497            // upstream is in the builder but downstream is not (e.g., edge to an independent job's fragment).
498            // Skip this edge.
499            return Ok(());
500        };
501        match (fragment_status, downstream_status) {
502            // `FragmentStatus` describes the lifecycle of the fragment's actor set, not whether
503            // the relation itself is new. Current mutation workflows only install relations where
504            // at least one endpoint is new or changed, so an existing-to-existing relation is
505            // already installed and needs no mutation.
506            (FragmentStatus::Existing(_), FragmentStatus::Existing(_)) => {}
507            (
508                FragmentStatus::Existing(fragment),
509                FragmentStatus::Changed {
510                    before,
511                    after: downstream_fragment,
512                },
513            ) => {
514                let (dispatchers, upstreams, no_shuffle_map) =
515                    Self::compose_edge(fragment_id, fragment, downstream, downstream_fragment);
516                Self::add_no_shuffle_mapping(
517                    &mut self.actor_new_no_shuffle,
518                    fragment_id,
519                    downstream.downstream_fragment_id,
520                    no_shuffle_map,
521                );
522                Self::add_upstreams(
523                    &mut self.result.upstreams,
524                    downstream.downstream_fragment_id,
525                    upstreams,
526                );
527                Self::add_dispatcher_updates(
528                    &mut self.result.dispatcher_updates,
529                    dispatchers,
530                    before.actors.keys().copied(),
531                );
532            }
533            (
534                FragmentStatus::Changed {
535                    before,
536                    after: fragment,
537                },
538                FragmentStatus::Existing(downstream_fragment),
539            ) => {
540                let (dispatchers, upstreams, no_shuffle_map) =
541                    Self::compose_edge(fragment_id, fragment, downstream, downstream_fragment);
542                Self::add_no_shuffle_mapping(
543                    &mut self.actor_new_no_shuffle,
544                    fragment_id,
545                    downstream.downstream_fragment_id,
546                    no_shuffle_map,
547                );
548                Self::add_dispatchers(&mut self.result.dispatchers, fragment_id, dispatchers);
549                Self::add_merge_updates(
550                    &mut self.result.merge_updates,
551                    fragment_id,
552                    downstream.downstream_fragment_id,
553                    downstream_fragment,
554                    upstreams,
555                    before.actors.keys().copied(),
556                );
557            }
558            (FragmentStatus::New(_), FragmentStatus::Changed { .. })
559            | (FragmentStatus::Changed { .. }, FragmentStatus::New(_)) => {
560                return Err(anyhow!(
561                    "an edge cannot connect new and changed fragments: {} -> {}",
562                    fragment_id,
563                    downstream.downstream_fragment_id,
564                )
565                .into());
566            }
567            (FragmentStatus::Existing(fragment), FragmentStatus::New(downstream_fragment))
568            | (FragmentStatus::New(fragment), FragmentStatus::Existing(downstream_fragment))
569            | (FragmentStatus::New(fragment), FragmentStatus::New(downstream_fragment))
570            | (
571                FragmentStatus::Changed {
572                    after: fragment, ..
573                },
574                FragmentStatus::Changed {
575                    after: downstream_fragment,
576                    ..
577                },
578            ) => {
579                let (dispatchers, upstreams, no_shuffle_map) =
580                    Self::compose_edge(fragment_id, fragment, downstream, downstream_fragment);
581                Self::add_no_shuffle_mapping(
582                    &mut self.actor_new_no_shuffle,
583                    fragment_id,
584                    downstream.downstream_fragment_id,
585                    no_shuffle_map,
586                );
587                Self::add_dispatchers(&mut self.result.dispatchers, fragment_id, dispatchers);
588                Self::add_upstreams(
589                    &mut self.result.upstreams,
590                    downstream.downstream_fragment_id,
591                    upstreams,
592                );
593            }
594        }
595        Ok(())
596    }
597
598    fn add_no_shuffle_mapping(
599        actor_new_no_shuffle: &mut ActorNewNoShuffle,
600        fragment_id: FragmentId,
601        downstream_fragment_id: FragmentId,
602        no_shuffle_map: Option<HashMap<ActorId, ActorId>>,
603    ) {
604        if let Some(no_shuffle_map) = no_shuffle_map {
605            actor_new_no_shuffle
606                .entry(fragment_id)
607                .or_default()
608                .insert(downstream_fragment_id, no_shuffle_map);
609        }
610    }
611
612    fn add_dispatchers(
613        fragment_dispatchers: &mut FragmentActorDispatchers,
614        fragment_id: FragmentId,
615        dispatchers: HashMap<ActorId, PbDispatcher>,
616    ) {
617        for (actor_id, dispatcher) in dispatchers {
618            fragment_dispatchers
619                .entry(fragment_id)
620                .or_default()
621                .entry(actor_id)
622                .or_default()
623                .push(dispatcher);
624        }
625    }
626
627    fn add_upstreams(
628        fragment_upstreams: &mut HashMap<FragmentId, HashMap<ActorId, ActorUpstreams>>,
629        fragment_id: FragmentId,
630        upstreams: HashMap<ActorId, ActorUpstreams>,
631    ) {
632        let target = fragment_upstreams.entry(fragment_id).or_default();
633        for (actor_id, actor_upstreams) in upstreams {
634            target.entry(actor_id).or_default().extend(actor_upstreams);
635        }
636    }
637
638    fn add_dispatcher_updates(
639        dispatcher_updates: &mut Vec<DispatcherUpdate>,
640        dispatchers: HashMap<ActorId, PbDispatcher>,
641        removed_downstream_actor_id: impl IntoIterator<Item = ActorId>,
642    ) {
643        let removed_downstream_actor_id: Vec<_> = removed_downstream_actor_id.into_iter().collect();
644        for (actor_id, dispatcher) in dispatchers {
645            dispatcher_updates.push(DispatcherUpdate {
646                actor_id,
647                dispatcher_id: dispatcher.dispatcher_id,
648                hash_mapping: dispatcher.hash_mapping,
649                added_downstream_actor_id: dispatcher.downstream_actor_id,
650                removed_downstream_actor_id: removed_downstream_actor_id.clone(),
651            });
652        }
653    }
654
655    fn add_merge_updates(
656        merge_updates: &mut HashMap<FragmentId, Vec<MergeUpdate>>,
657        fragment_id: FragmentId,
658        downstream_fragment_id: FragmentId,
659        downstream_fragment: &EdgeBuilderFragmentInfo,
660        upstreams: HashMap<ActorId, ActorUpstreams>,
661        removed_upstream_actor_id: impl IntoIterator<Item = ActorId>,
662    ) {
663        let removed_upstream_actor_id: Vec<_> = removed_upstream_actor_id.into_iter().collect();
664        let fragment_merge_updates = merge_updates.entry(downstream_fragment_id).or_default();
665        for actor_id in downstream_fragment.actors.keys() {
666            let added_upstream_actors = upstreams
667                .get(actor_id)
668                .and_then(|upstreams| upstreams.get(&fragment_id))
669                .into_iter()
670                .flat_map(|upstreams| upstreams.values().cloned())
671                .collect();
672            fragment_merge_updates.push(MergeUpdate {
673                actor_id: *actor_id,
674                upstream_fragment_id: fragment_id,
675                new_upstream_fragment_id: None,
676                added_upstream_actors,
677                removed_upstream_actor_id: removed_upstream_actor_id.clone(),
678            });
679        }
680    }
681
682    fn compose_edge_dispatchers(
683        fragment: &EdgeBuilderFragmentInfo,
684        downstream: &DownstreamFragmentRelation,
685        downstream_fragment: &EdgeBuilderFragmentInfo,
686    ) -> (
687        HashMap<ActorId, PbDispatcher>,
688        Option<HashMap<ActorId, ActorId>>,
689    ) {
690        compose_dispatchers(
691            fragment.distribution_type,
692            &fragment.actors,
693            downstream.downstream_fragment_id,
694            downstream_fragment.distribution_type,
695            &downstream_fragment.actors,
696            downstream.dispatcher_type,
697            downstream.dist_key_indices.clone(),
698            downstream.output_mapping.clone(),
699        )
700    }
701
702    fn compose_edge(
703        fragment_id: FragmentId,
704        fragment: &EdgeBuilderFragmentInfo,
705        downstream: &DownstreamFragmentRelation,
706        downstream_fragment: &EdgeBuilderFragmentInfo,
707    ) -> ComposedEdge {
708        let (dispatchers, no_shuffle_map) =
709            Self::compose_edge_dispatchers(fragment, downstream, downstream_fragment);
710        let mut upstreams: HashMap<ActorId, ActorUpstreams> = HashMap::new();
711        for (&actor_id, dispatcher) in &dispatchers {
712            let actor_location = &fragment.actor_location[&actor_id];
713            for &downstream_actor in &dispatcher.downstream_actor_id {
714                upstreams
715                    .entry(downstream_actor)
716                    .or_default()
717                    .entry(fragment_id)
718                    .or_default()
719                    .insert(
720                        actor_id,
721                        ActorInfo {
722                            actor_id,
723                            host: Some(actor_location.clone()),
724                            partial_graph_id: fragment.partial_graph_id,
725                        },
726                    );
727            }
728        }
729        (dispatchers, upstreams, no_shuffle_map)
730    }
731
732    pub(crate) fn replace_upstream(
733        mut self,
734        fragment_id: FragmentId,
735        original_upstream_fragment_id: FragmentId,
736        new_upstream_fragment_id: FragmentId,
737    ) -> Self {
738        let fragment_merge_updates = self.result.merge_updates.entry(fragment_id).or_default();
739        if let Some(fragment_upstreams) = self.result.upstreams.get_mut(&fragment_id) {
740            fragment_upstreams.retain(|&actor_id, actor_upstreams| {
741                if let Some(new_upstreams) = actor_upstreams.remove(&new_upstream_fragment_id) {
742                    fragment_merge_updates.push(MergeUpdate {
743                        actor_id,
744                        upstream_fragment_id: original_upstream_fragment_id,
745                        new_upstream_fragment_id: Some(new_upstream_fragment_id),
746                        added_upstream_actors: new_upstreams.into_values().collect(),
747                        removed_upstream_actor_id: vec![],
748                    })
749                } else if cfg!(debug_assertions) {
750                    panic!("cannot find new upstreams for actor {} in fragment {} to new_upstream {}. Current upstreams {:?}", actor_id, fragment_id, new_upstream_fragment_id, actor_upstreams);
751                } else {
752                    warn!(%actor_id, %fragment_id, %new_upstream_fragment_id, ?actor_upstreams, "cannot find new upstreams for actor");
753                }
754                !actor_upstreams.is_empty()
755            })
756        } else if cfg!(debug_assertions) {
757            panic!(
758                "cannot find new upstreams for fragment {} to new_upstream {} to replace {}. Current upstreams: {:?}",
759                fragment_id,
760                new_upstream_fragment_id,
761                original_upstream_fragment_id,
762                self.result.upstreams
763            );
764        } else {
765            warn!(%fragment_id, %new_upstream_fragment_id, %original_upstream_fragment_id, upstreams = ?self.result.upstreams, "cannot find new upstreams to replace");
766        }
767        self
768    }
769
770    /// Finalize the builder and return the generated edges and mutation-time deltas.
771    pub(crate) fn build(self) -> (FragmentEdgeBuildResult, ActorNewNoShuffle) {
772        (self.result, self.actor_new_no_shuffle)
773    }
774}
775
776#[cfg(test)]
777mod tests {
778    use risingwave_meta_model::DispatcherType;
779    use risingwave_pb::stream_plan::PbDispatchOutputMapping;
780
781    use super::*;
782
783    #[derive(Clone, Copy, Debug)]
784    enum TestStatus {
785        Existing,
786        New,
787        Changed,
788    }
789
790    fn actor(id: u32) -> ActorId {
791        ActorId::new(id)
792    }
793
794    fn fragment(id: u32) -> FragmentId {
795        FragmentId::new(id)
796    }
797
798    fn edge_info(
799        distribution_type: DistributionType,
800        actors: impl IntoIterator<Item = (u32, Option<Bitmap>)>,
801    ) -> EdgeBuilderFragmentInfo {
802        let actors: HashMap<_, _> = actors
803            .into_iter()
804            .map(|(id, bitmap)| (actor(id), bitmap))
805            .collect();
806        let actor_location = actors
807            .keys()
808            .map(|actor_id| {
809                (
810                    *actor_id,
811                    HostAddress {
812                        host: format!("actor-{actor_id}"),
813                        port: 1234,
814                    },
815                )
816            })
817            .collect();
818        EdgeBuilderFragmentInfo {
819            distribution_type,
820            actors,
821            actor_location,
822            partial_graph_id: PartialGraphId::new(1),
823        }
824    }
825
826    fn single_info(actor_id: u32) -> EdgeBuilderFragmentInfo {
827        edge_info(DistributionType::Single, [(actor_id, None)])
828    }
829
830    fn relation(target: FragmentId, dispatcher_type: DispatcherType) -> DownstreamFragmentRelation {
831        DownstreamFragmentRelation {
832            downstream_fragment_id: target,
833            dispatcher_type,
834            dist_key_indices: vec![],
835            output_mapping: PbDispatchOutputMapping::default(),
836        }
837    }
838
839    fn build_status_pair(
840        source: TestStatus,
841        target: TestStatus,
842    ) -> MetaResult<FragmentEdgeBuildResult> {
843        let source_fragment = fragment(1);
844        let target_fragment = fragment(2);
845        let existing = [(source_fragment, source, 1), (target_fragment, target, 11)]
846            .into_iter()
847            .filter(|(_, status, _)| !matches!(status, TestStatus::New))
848            .map(|(fragment_id, _, actor_id)| (fragment_id, single_info(actor_id)));
849        let mut builder = FragmentEdgeBuilder::new();
850        builder.add_existing_fragment_infos(existing);
851        for (fragment_id, status, actor_id) in
852            [(source_fragment, source, 2), (target_fragment, target, 12)]
853        {
854            match status {
855                TestStatus::Existing => {}
856                TestStatus::New => {
857                    builder.add_new_fragment_infos([(fragment_id, single_info(actor_id))]);
858                }
859                TestStatus::Changed => {
860                    builder.replace_existing_fragment_actor_infos([(
861                        fragment_id,
862                        single_info(actor_id),
863                    )]);
864                }
865            }
866        }
867        let (result, _) = builder
868            .finish_fragments()
869            .add_relations(&HashMap::from([(
870                source_fragment,
871                vec![relation(target_fragment, DispatcherType::Broadcast)],
872            )]))?
873            .build();
874        Ok(result)
875    }
876
877    #[test]
878    fn test_endpoint_status_matrix() {
879        for source in [TestStatus::Existing, TestStatus::New, TestStatus::Changed] {
880            for target in [TestStatus::Existing, TestStatus::New, TestStatus::Changed] {
881                let result = build_status_pair(source, target);
882                if matches!(
883                    (source, target),
884                    (TestStatus::New, TestStatus::Changed) | (TestStatus::Changed, TestStatus::New)
885                ) {
886                    assert!(result.is_err(), "source={source:?}, target={target:?}");
887                    continue;
888                }
889                let result = result.unwrap();
890                assert_eq!(
891                    result.dispatchers.contains_key(&fragment(1)),
892                    matches!(source, TestStatus::New | TestStatus::Changed)
893                        || matches!((source, target), (TestStatus::Existing, TestStatus::New)),
894                    "source={source:?}, target={target:?}"
895                );
896                assert_eq!(
897                    result.upstreams.contains_key(&fragment(2)),
898                    matches!(target, TestStatus::New | TestStatus::Changed)
899                        || matches!((source, target), (TestStatus::New, TestStatus::Existing)),
900                    "source={source:?}, target={target:?}"
901                );
902                assert_eq!(
903                    !result.dispatcher_updates.is_empty(),
904                    matches!(
905                        (source, target),
906                        (TestStatus::Existing, TestStatus::Changed)
907                    ),
908                    "source={source:?}, target={target:?}"
909                );
910                assert_eq!(
911                    result
912                        .merge_updates
913                        .values()
914                        .any(|updates| !updates.is_empty()),
915                    matches!(
916                        (source, target),
917                        (TestStatus::Changed, TestStatus::Existing)
918                    ),
919                    "source={source:?}, target={target:?}"
920                );
921            }
922        }
923    }
924
925    #[test]
926    fn test_hash_and_broadcast_dispatcher_updates() {
927        for dispatcher_type in [DispatcherType::Hash, DispatcherType::Broadcast] {
928            let source = fragment(1);
929            let target = fragment(2);
930            let bitmap = || Some(Bitmap::from_iter([true, true]));
931            let mut builder = FragmentEdgeBuilder::new();
932            builder.add_existing_fragment_infos([
933                (source, edge_info(DistributionType::Single, [(1, None)])),
934                (target, edge_info(DistributionType::Hash, [(11, bitmap())])),
935            ]);
936            builder.replace_existing_fragment_actor_infos([(
937                target,
938                edge_info(DistributionType::Hash, [(12, bitmap())]),
939            )]);
940            let (result, _) = builder
941                .finish_fragments()
942                .add_edge(source, &relation(target, dispatcher_type))
943                .unwrap()
944                .build();
945
946            let update = &result.dispatcher_updates[0];
947            assert_eq!(update.added_downstream_actor_id, vec![actor(12)]);
948            assert_eq!(update.removed_downstream_actor_id, vec![actor(11)]);
949            assert_eq!(
950                update.hash_mapping.is_some(),
951                dispatcher_type == DispatcherType::Hash
952            );
953        }
954    }
955
956    #[test]
957    fn test_no_shuffle_changed_fragments() {
958        let source = fragment(1);
959        let target = fragment(2);
960        let mut builder = FragmentEdgeBuilder::new();
961        builder.add_existing_fragment_infos([(source, single_info(1)), (target, single_info(11))]);
962        builder.replace_existing_fragment_actor_infos([
963            (source, single_info(2)),
964            (target, single_info(12)),
965        ]);
966        let (result, actor_new_no_shuffle) = builder
967            .finish_fragments()
968            .add_edge(source, &relation(target, DispatcherType::NoShuffle))
969            .unwrap()
970            .build();
971
972        assert!(result.dispatcher_updates.is_empty());
973        assert!(result.merge_updates.is_empty());
974        assert_eq!(actor_new_no_shuffle[&source][&target][&actor(2)], actor(12));
975    }
976
977    #[test]
978    fn test_no_shuffle_updates_remove_before_actor_superset() {
979        let source = fragment(1);
980        let target = fragment(2);
981        let actors = |left, right| {
982            [
983                (left, Some(Bitmap::from_iter([true, false]))),
984                (right, Some(Bitmap::from_iter([false, true]))),
985            ]
986        };
987        let relations =
988            HashMap::from([(source, vec![relation(target, DispatcherType::NoShuffle)])]);
989
990        let mut builder = FragmentEdgeBuilder::new();
991        builder.add_existing_fragment_infos([
992            (source, edge_info(DistributionType::Hash, actors(1, 2))),
993            (target, edge_info(DistributionType::Hash, actors(11, 12))),
994        ]);
995        builder.replace_existing_fragment_actor_infos([(
996            target,
997            edge_info(DistributionType::Hash, actors(13, 14)),
998        )]);
999        let (result, _) = builder
1000            .finish_fragments()
1001            .add_relations(&relations)
1002            .unwrap()
1003            .build();
1004
1005        assert_eq!(result.dispatcher_updates.len(), 2);
1006        for update in result.dispatcher_updates {
1007            assert_eq!(update.added_downstream_actor_id.len(), 1);
1008            let mut removed = update.removed_downstream_actor_id;
1009            removed.sort_unstable();
1010            assert_eq!(removed, vec![actor(11), actor(12)]);
1011        }
1012
1013        let mut builder = FragmentEdgeBuilder::new();
1014        builder.add_existing_fragment_infos([
1015            (source, edge_info(DistributionType::Hash, actors(1, 2))),
1016            (target, edge_info(DistributionType::Hash, actors(11, 12))),
1017        ]);
1018        builder.replace_existing_fragment_actor_infos([(
1019            source,
1020            edge_info(DistributionType::Hash, actors(3, 4)),
1021        )]);
1022        let (result, _) = builder
1023            .finish_fragments()
1024            .add_relations(&relations)
1025            .unwrap()
1026            .build();
1027
1028        let merge_updates = &result.merge_updates[&target];
1029        assert_eq!(merge_updates.len(), 2);
1030        for update in merge_updates {
1031            assert_eq!(update.added_upstream_actors.len(), 1);
1032            let mut removed = update.removed_upstream_actor_id.clone();
1033            removed.sort_unstable();
1034            assert_eq!(removed, vec![actor(1), actor(2)]);
1035        }
1036    }
1037
1038    #[test]
1039    fn test_replace_upstream() {
1040        let old_source = fragment(1);
1041        let new_source = fragment(2);
1042        let target = fragment(3);
1043        let mut builder = FragmentEdgeBuilder::new();
1044        builder.add_existing_fragment_infos([(target, single_info(11))]);
1045        builder.add_new_fragment_infos([(new_source, single_info(2))]);
1046        let (result, _) = builder
1047            .finish_fragments()
1048            .add_edge(new_source, &relation(target, DispatcherType::Broadcast))
1049            .unwrap()
1050            .replace_upstream(target, old_source, new_source)
1051            .build();
1052
1053        let update = &result.merge_updates[&target][0];
1054        assert_eq!(update.upstream_fragment_id, old_source);
1055        assert_eq!(update.new_upstream_fragment_id, Some(new_source));
1056        assert_eq!(update.added_upstream_actors[0].actor_id, actor(2));
1057    }
1058
1059    #[test]
1060    fn test_new_to_existing_upstreams_are_consumed_for_add_mutation() {
1061        let source = fragment(1);
1062        let target = fragment(2);
1063        let mut builder = FragmentEdgeBuilder::new();
1064        builder.add_existing_fragment_infos([(target, single_info(11))]);
1065        builder.add_new_fragment_infos([(source, single_info(1))]);
1066        let (mut result, _) = builder
1067            .finish_fragments()
1068            .add_edge(source, &relation(target, DispatcherType::Broadcast))
1069            .unwrap()
1070            .build();
1071        let source_actor = StreamActor {
1072            actor_id: actor(1),
1073            fragment_id: source,
1074            vnode_bitmap: None,
1075            mview_definition: Default::default(),
1076            expr_context: None,
1077            config_override: Default::default(),
1078        };
1079        let node = StreamNode::default();
1080        let worker_id: WorkerId = 1.into();
1081        let actors_to_create = result.collect_actors_to_create(std::iter::once((
1082            source,
1083            &node,
1084            std::iter::once((&source_actor, worker_id)),
1085            [],
1086        )));
1087
1088        assert_eq!(actors_to_create[&worker_id][&source].1[0].2.len(), 1);
1089        let upstream_actors = result.take_common_upstream_actors(source, target).unwrap();
1090        assert_eq!(upstream_actors.len(), 1);
1091        assert_eq!(upstream_actors[0].actor_id, actor(1));
1092        result.apply_to_add_mutation(&mut AddMutation::default());
1093    }
1094
1095    #[test]
1096    fn test_take_common_upstream_actors_from_hash_edge() {
1097        let source = fragment(1);
1098        let target = fragment(2);
1099        let actors = |left, right| {
1100            [
1101                (left, Some(Bitmap::from_iter([true, false]))),
1102                (right, Some(Bitmap::from_iter([false, true]))),
1103            ]
1104        };
1105        let mut builder = FragmentEdgeBuilder::new();
1106        builder.add_existing_fragment_infos([(
1107            target,
1108            edge_info(DistributionType::Hash, actors(11, 12)),
1109        )]);
1110        builder.add_new_fragment_infos([(source, edge_info(DistributionType::Hash, actors(1, 2)))]);
1111        let (mut result, _) = builder
1112            .finish_fragments()
1113            .add_edge(source, &relation(target, DispatcherType::Hash))
1114            .unwrap()
1115            .build();
1116
1117        let upstream_actors = result.take_common_upstream_actors(source, target).unwrap();
1118        assert_eq!(
1119            upstream_actors
1120                .into_iter()
1121                .map(|actor| actor.actor_id)
1122                .collect::<Vec<_>>(),
1123            vec![actor(1), actor(2)]
1124        );
1125        assert!(!result.upstreams.contains_key(&target));
1126    }
1127
1128    #[test]
1129    #[should_panic(expected = "remaining_upstreams=1")]
1130    fn test_add_mutation_rejects_uncollected_actor_upstreams() {
1131        let source = fragment(1);
1132        let target = fragment(2);
1133        let mut builder = FragmentEdgeBuilder::new();
1134        builder.add_existing_fragment_infos([(source, single_info(1))]);
1135        builder.add_new_fragment_infos([(target, single_info(11))]);
1136        let (result, _) = builder
1137            .finish_fragments()
1138            .add_edge(source, &relation(target, DispatcherType::Broadcast))
1139            .unwrap()
1140            .build();
1141
1142        result.apply_to_add_mutation(&mut AddMutation::default());
1143    }
1144
1145    #[test]
1146    fn test_update_mutation_consumes_edge_updates_before_validation() {
1147        let source = fragment(1);
1148        let mut result = build_status_pair(TestStatus::Changed, TestStatus::Existing).unwrap();
1149        let source_actor = StreamActor {
1150            actor_id: actor(2),
1151            fragment_id: source,
1152            vnode_bitmap: None,
1153            mview_definition: Default::default(),
1154            expr_context: None,
1155            config_override: Default::default(),
1156        };
1157        let node = StreamNode::default();
1158        result.collect_actors_to_create(std::iter::once((
1159            source,
1160            &node,
1161            std::iter::once((&source_actor, 1.into())),
1162            [],
1163        )));
1164        let mut mutation = UpdateMutation::default();
1165
1166        result.apply_to_update_mutation(&mut mutation);
1167
1168        assert_eq!(mutation.merge_update.len(), 1);
1169    }
1170
1171    #[test]
1172    fn test_attach_new_relation_dispatcher_to_add_mutation() {
1173        let source = fragment(1);
1174        let target = fragment(2);
1175        let mut builder = FragmentEdgeBuilder::new();
1176        builder.add_existing_fragment_infos([(source, single_info(1))]);
1177        builder.add_new_fragment_infos([(target, single_info(11))]);
1178        let (mut result, _) = builder
1179            .finish_fragments()
1180            .add_edge(source, &relation(target, DispatcherType::Broadcast))
1181            .unwrap()
1182            .build();
1183        let target_actor = StreamActor {
1184            actor_id: actor(11),
1185            fragment_id: target,
1186            vnode_bitmap: None,
1187            mview_definition: Default::default(),
1188            expr_context: None,
1189            config_override: Default::default(),
1190        };
1191        let node = StreamNode::default();
1192        result.collect_actors_to_create(std::iter::once((
1193            target,
1194            &node,
1195            std::iter::once((&target_actor, 1.into())),
1196            [],
1197        )));
1198        let mut mutation = AddMutation::default();
1199
1200        result.apply_to_add_mutation(&mut mutation);
1201
1202        assert_eq!(mutation.actor_dispatchers[&actor(1)].dispatchers.len(), 1);
1203        assert_eq!(
1204            mutation.actor_dispatchers[&actor(1)].dispatchers[0].downstream_actor_id,
1205            vec![actor(11)]
1206        );
1207    }
1208
1209    #[test]
1210    fn test_add_relations_only_uses_registered_fragments() {
1211        let external_source = fragment(1);
1212        let source = fragment(2);
1213        let target = fragment(3);
1214        let mut builder = FragmentEdgeBuilder::new();
1215        builder.add_new_fragment_infos([(source, single_info(2)), (target, single_info(3))]);
1216        let (result, _) = builder
1217            .finish_fragments()
1218            .add_relations(&HashMap::from([
1219                (
1220                    external_source,
1221                    vec![relation(source, DispatcherType::Broadcast)],
1222                ),
1223                (source, vec![relation(target, DispatcherType::Broadcast)]),
1224            ]))
1225            .unwrap()
1226            .build();
1227
1228        assert!(!result.dispatchers.contains_key(&external_source));
1229        assert!(!result.upstreams.contains_key(&source));
1230        assert!(result.dispatchers.contains_key(&source));
1231        assert!(result.upstreams.contains_key(&target));
1232    }
1233
1234    #[test]
1235    fn test_add_edge_rejects_unregistered_upstream() {
1236        let source = fragment(1);
1237        let target = fragment(2);
1238        let mut builder = FragmentEdgeBuilder::new();
1239        builder.add_existing_fragment_infos([(target, single_info(2))]);
1240
1241        let result = builder
1242            .finish_fragments()
1243            .add_edge(source, &relation(target, DispatcherType::Broadcast));
1244
1245        assert!(result.is_err());
1246    }
1247
1248    #[test]
1249    #[should_panic(expected = "must use entirely fresh actor ids")]
1250    fn test_changed_fragment_rejects_retained_actor_ids() {
1251        let fragment = fragment(1);
1252        let mut builder = FragmentEdgeBuilder::new();
1253        builder.add_existing_fragment_infos([(fragment, single_info(1))]);
1254        builder.replace_existing_fragment_actor_infos([(fragment, single_info(1))]);
1255    }
1256}