Skip to main content

risingwave_meta/stream/stream_graph/
actor.rs

1// Copyright 2023 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::{BTreeMap, HashMap};
16
17use assert_matches::assert_matches;
18use risingwave_common::bail;
19use risingwave_common::hash::{IsSingleton, VnodeCount, VnodeCountCompat};
20use risingwave_common::util::iter_util::ZipEqFast;
21use risingwave_common::util::stream_graph_visitor::visit_tables;
22use risingwave_pb::stream_plan::stream_node::NodeBody;
23use risingwave_pb::stream_plan::{
24    DispatchStrategy, DispatcherType, MergeNode, StreamNode, StreamScanType,
25};
26
27use crate::MetaResult;
28use crate::model::{Fragment, FragmentDownstreamRelation, FragmentId, FragmentReplaceUpstream};
29use crate::stream::stream_graph::fragment::{
30    CompleteStreamFragmentGraph, DownstreamExternalEdgeId, EdgeId, EitherFragment,
31    StreamFragmentEdge,
32};
33use crate::stream::stream_graph::id::GlobalFragmentId;
34use crate::stream::stream_graph::schedule;
35use crate::stream::stream_graph::schedule::Distribution;
36
37impl FragmentActorBuilder {
38    /// Rewrite the actor body.
39    ///
40    /// During this process, the following things will be done:
41    /// 1. Replace the logical `Exchange` in node's input with `Merge`, which can be executed on the
42    ///    compute nodes.
43    fn rewrite(&self) -> MetaResult<StreamNode> {
44        self.rewrite_inner(&self.node, 0)
45    }
46
47    fn rewrite_inner(&self, stream_node: &StreamNode, depth: usize) -> MetaResult<StreamNode> {
48        match stream_node.get_node_body()? {
49            // Leaf node `Exchange`.
50            NodeBody::Exchange(exchange) => {
51                // The exchange node should always be the bottom of the plan node. If we find one
52                // when the depth is 0, it means that the plan node is not well-formed.
53                if depth == 0 {
54                    bail!(
55                        "there should be no ExchangeNode on the top of the plan node: {:#?}",
56                        stream_node
57                    )
58                }
59                assert!(!stream_node.get_fields().is_empty());
60                assert!(stream_node.input.is_empty());
61
62                // Index the upstreams by the an internal edge ID.
63                let (upstream_fragment_id, _) = &self.upstreams[&EdgeId::Internal {
64                    link_id: stream_node.get_operator_id().as_raw_id(),
65                }];
66
67                let upstream_fragment_id = upstream_fragment_id.as_global_id();
68
69                Ok(StreamNode {
70                    node_body: Some(NodeBody::Merge(Box::new({
71                        MergeNode {
72                            upstream_fragment_id,
73                            upstream_dispatcher_type: exchange.get_strategy()?.r#type,
74                            ..Default::default()
75                        }
76                    }))),
77                    identity: "MergeExecutor".to_owned(),
78                    ..stream_node.clone()
79                })
80            }
81
82            // "Leaf" node `StreamScan`.
83            NodeBody::StreamScan(stream_scan) => {
84                let input = stream_node.get_input();
85                if stream_scan.stream_scan_type() == StreamScanType::CrossDbSnapshotBackfill {
86                    // CrossDbSnapshotBackfill is a special case, which doesn't have any upstream actor
87                    // and always reads from log store.
88                    assert!(input.is_empty());
89                    return Ok(stream_node.clone());
90                }
91                assert_eq!(input.len(), 2);
92
93                let merge_node = &input[0];
94                assert_matches!(merge_node.node_body, Some(NodeBody::Merge(_)));
95                let batch_plan_node = &input[1];
96                assert_matches!(batch_plan_node.node_body, Some(NodeBody::BatchPlan(_)));
97
98                // Index the upstreams by the an external edge ID.
99                let (upstream_fragment_id, upstream_no_shuffle_actor) = &self.upstreams
100                    [&EdgeId::UpstreamExternal {
101                        upstream_job_id: stream_scan.table_id.as_job_id(),
102                        downstream_fragment_id: self.fragment_id,
103                    }];
104
105                let is_shuffled_backfill = stream_scan.stream_scan_type
106                    == StreamScanType::ArrangementBackfill as i32
107                    || stream_scan.stream_scan_type == StreamScanType::SnapshotBackfill as i32;
108                if !is_shuffled_backfill {
109                    assert!(*upstream_no_shuffle_actor);
110                }
111
112                let upstream_dispatcher_type = if is_shuffled_backfill {
113                    // FIXME(kwannoel): Should the upstream dispatcher type depends on the upstream distribution?
114                    // If singleton, use `Simple` dispatcher, otherwise use `Hash` dispatcher.
115                    DispatcherType::Hash as _
116                } else {
117                    DispatcherType::NoShuffle as _
118                };
119
120                let upstream_fragment_id = upstream_fragment_id.as_global_id();
121
122                let input = vec![
123                    // Fill the merge node body with correct upstream info.
124                    StreamNode {
125                        node_body: Some(NodeBody::Merge(Box::new({
126                            MergeNode {
127                                upstream_fragment_id,
128                                upstream_dispatcher_type,
129                                ..Default::default()
130                            }
131                        }))),
132                        ..merge_node.clone()
133                    },
134                    batch_plan_node.clone(),
135                ];
136
137                Ok(StreamNode {
138                    input,
139                    ..stream_node.clone()
140                })
141            }
142
143            // "Leaf" node `CdcFilter` and `SourceBackfill`. They both `Merge` an upstream `Source`
144            // cdc_filter -> backfill -> mview
145            // source_backfill -> mview
146            NodeBody::CdcFilter(_) | NodeBody::SourceBackfill(_) => {
147                let input = stream_node.get_input();
148                assert_eq!(input.len(), 1);
149
150                let merge_node = &input[0];
151                assert_matches!(merge_node.node_body, Some(NodeBody::Merge(_)));
152
153                let upstream_source_id = match stream_node.get_node_body()? {
154                    NodeBody::CdcFilter(node) => node.upstream_source_id,
155                    NodeBody::SourceBackfill(node) => node.upstream_source_id,
156                    _ => unreachable!(),
157                };
158
159                // Index the upstreams by the an external edge ID.
160                let (upstream_fragment_id, upstream_is_no_shuffle) = &self.upstreams
161                    [&EdgeId::UpstreamExternal {
162                        upstream_job_id: upstream_source_id.as_share_source_job_id(),
163                        downstream_fragment_id: self.fragment_id,
164                    }];
165
166                assert!(
167                    *upstream_is_no_shuffle,
168                    "Upstream Cdc Source should be singleton. \
169                    SourceBackfill is NoShuffle 1-1 correspondence. \
170                    So they both should have only one upstream actor."
171                );
172
173                let upstream_fragment_id = upstream_fragment_id.as_global_id();
174
175                // rewrite the input
176                let input = vec![
177                    // Fill the merge node body with correct upstream info.
178                    StreamNode {
179                        node_body: Some(NodeBody::Merge(Box::new({
180                            MergeNode {
181                                upstream_fragment_id,
182                                upstream_dispatcher_type: DispatcherType::NoShuffle as _,
183                                ..Default::default()
184                            }
185                        }))),
186                        ..merge_node.clone()
187                    },
188                ];
189                Ok(StreamNode {
190                    input,
191                    ..stream_node.clone()
192                })
193            }
194
195            NodeBody::IcebergWithPkIndexWriter(_) => {
196                let mut new_stream_node = stream_node.clone();
197                // Compaction alternates between the normal and resolver inputs, so either merge
198                // must be able to remain alive while temporarily disconnected.
199                for (input, new_input) in stream_node
200                    .input
201                    .iter()
202                    .zip_eq_fast(&mut new_stream_node.input)
203                {
204                    *new_input = self.rewrite_inner(input, depth + 1)?;
205                    let Some(NodeBody::Merge(merge)) = new_input.node_body.as_mut() else {
206                        bail!("iceberg pk-index writer input must be a merge after actor rewrite");
207                    };
208                    merge.allow_empty_upstream = true;
209                }
210                Ok(new_stream_node)
211            }
212
213            // For other nodes, visit the children recursively.
214            _ => {
215                let mut new_stream_node = stream_node.clone();
216                for (input, new_input) in stream_node
217                    .input
218                    .iter()
219                    .zip_eq_fast(&mut new_stream_node.input)
220                {
221                    *new_input = self.rewrite_inner(input, depth + 1)?;
222                }
223                Ok(new_stream_node)
224            }
225        }
226    }
227}
228
229/// The required changes to an existing external actor to build the graph of a streaming job.
230///
231/// For example, when we're creating an mview on an existing mview, we need to add new downstreams
232/// to the upstream actors, by adding new dispatchers.
233#[derive(Default)]
234struct UpstreamFragmentChange {
235    /// The new downstreams to be added.
236    new_downstreams: HashMap<GlobalFragmentId, DispatchStrategy>,
237}
238
239#[derive(Default)]
240struct DownstreamFragmentChange {
241    /// The new upstreams to be added (replaced), indexed by the edge id to upstream fragment.
242    /// `edge_id` -> new upstream fragment id
243    new_upstreams: HashMap<DownstreamExternalEdgeId, GlobalFragmentId>,
244}
245
246impl UpstreamFragmentChange {
247    /// Add a dispatcher to the external actor.
248    fn add_dispatcher(
249        &mut self,
250        downstream_fragment_id: GlobalFragmentId,
251        dispatch: DispatchStrategy,
252    ) {
253        self.new_downstreams
254            .try_insert(downstream_fragment_id, dispatch)
255            .unwrap();
256    }
257}
258
259impl DownstreamFragmentChange {
260    /// Add an upstream to the external actor.
261    fn add_upstream(
262        &mut self,
263        edge_id: DownstreamExternalEdgeId,
264        new_upstream_fragment_id: GlobalFragmentId,
265    ) {
266        self.new_upstreams
267            .try_insert(edge_id, new_upstream_fragment_id)
268            .unwrap();
269    }
270}
271
272#[derive(Debug)]
273struct FragmentActorBuilder {
274    fragment_id: GlobalFragmentId,
275    node: StreamNode,
276    downstreams: HashMap<GlobalFragmentId, DispatchStrategy>,
277    // edge_id -> (upstream fragment_id, whether the edge is no shuffle)
278    upstreams: HashMap<EdgeId, (GlobalFragmentId, bool)>,
279}
280
281impl FragmentActorBuilder {
282    fn new(fragment_id: GlobalFragmentId, node: StreamNode) -> Self {
283        Self {
284            fragment_id,
285            node,
286            downstreams: Default::default(),
287            upstreams: Default::default(),
288        }
289    }
290}
291
292/// The actual mutable state of building an actor graph.
293///
294/// When the fragments are visited in a topological order, actor builders will be added to this
295/// state and the scheduled locations will be added. As the building process is run on the
296/// **complete graph** which also contains the info of the existing (external) fragments, the info
297/// of them will be also recorded.
298#[derive(Default)]
299struct ActorGraphBuildStateInner {
300    /// The builders of the actors to be built.
301    fragment_actor_builders: BTreeMap<GlobalFragmentId, FragmentActorBuilder>,
302
303    /// The required changes to the external downstream fragment. See [`DownstreamFragmentChange`].
304    /// Indexed by the `fragment_id` of fragments that have updates on its downstream.
305    downstream_fragment_changes: BTreeMap<GlobalFragmentId, DownstreamFragmentChange>,
306
307    /// The required changes to the external upstream fragment. See [`UpstreamFragmentChange`].
308    /// /// Indexed by the `fragment_id` of fragments that have updates on its upstream.
309    upstream_fragment_changes: BTreeMap<GlobalFragmentId, UpstreamFragmentChange>,
310}
311
312/// The information of a fragment, used for parameter passing for `Inner::add_link`.
313struct FragmentLinkNode {
314    fragment_id: GlobalFragmentId,
315}
316
317impl ActorGraphBuildStateInner {
318    /// Add the new downstream fragment relation to a fragment.
319    ///
320    /// - If the fragment is to be built, the fragment relation will be added to the fragment actor builder.
321    /// - If the fragment is an external existing fragment, the fragment relation will be added to the external changes.
322    fn add_dispatcher(
323        &mut self,
324        fragment_id: GlobalFragmentId,
325        downstream_fragment_id: GlobalFragmentId,
326        dispatch: DispatchStrategy,
327    ) {
328        if let Some(builder) = self.fragment_actor_builders.get_mut(&fragment_id) {
329            builder
330                .downstreams
331                .try_insert(downstream_fragment_id, dispatch)
332                .unwrap();
333        } else {
334            self.upstream_fragment_changes
335                .entry(fragment_id)
336                .or_default()
337                .add_dispatcher(downstream_fragment_id, dispatch);
338        }
339    }
340
341    /// Add the new upstream for an actor.
342    ///
343    /// - If the actor is to be built, the upstream will be added to the actor builder.
344    /// - If the actor is an external actor, the upstream will be added to the external changes.
345    fn add_upstream(
346        &mut self,
347        fragment_id: GlobalFragmentId,
348        edge_id: EdgeId,
349        upstream_fragment_id: GlobalFragmentId,
350        is_no_shuffle: bool,
351    ) {
352        if let Some(builder) = self.fragment_actor_builders.get_mut(&fragment_id) {
353            builder
354                .upstreams
355                .try_insert(edge_id, (upstream_fragment_id, is_no_shuffle))
356                .unwrap();
357        } else {
358            let EdgeId::DownstreamExternal(edge_id) = edge_id else {
359                unreachable!("edge from internal to external must be `DownstreamExternal`")
360            };
361            self.downstream_fragment_changes
362                .entry(fragment_id)
363                .or_default()
364                .add_upstream(edge_id, upstream_fragment_id);
365        }
366    }
367
368    /// Add a "link" between two fragments in the graph.
369    ///
370    /// The `edge` will be transformed into the fragment relation (downstream - upstream) pair between two fragments,
371    /// based on the distribution and the dispatch strategy. They will be
372    /// finally transformed to `Dispatcher` and `Merge` nodes when building the actors.
373    ///
374    /// If there're existing (external) fragments, the info will be recorded in `upstream_fragment_changes` and `downstream_fragment_changes`,
375    /// instead of the actor builders.
376    fn add_link(
377        &mut self,
378        upstream: FragmentLinkNode,
379        downstream: FragmentLinkNode,
380        edge: &StreamFragmentEdge,
381    ) {
382        let dt = edge.dispatch_strategy.r#type();
383
384        match dt {
385            // For `NoShuffle`, make n "1-1" links between the actors.
386            DispatcherType::NoShuffle => {
387                // Create a new dispatcher just between these two actors.
388                self.add_dispatcher(
389                    upstream.fragment_id,
390                    downstream.fragment_id,
391                    edge.dispatch_strategy.clone(),
392                );
393
394                // Also record the upstream for the downstream actor.
395                self.add_upstream(downstream.fragment_id, edge.id, upstream.fragment_id, true);
396            }
397
398            // Otherwise, make m * n links between the actors.
399            DispatcherType::Hash | DispatcherType::Broadcast | DispatcherType::Simple => {
400                self.add_dispatcher(
401                    upstream.fragment_id,
402                    downstream.fragment_id,
403                    edge.dispatch_strategy.clone(),
404                );
405                self.add_upstream(downstream.fragment_id, edge.id, upstream.fragment_id, false);
406            }
407
408            DispatcherType::Unspecified => unreachable!(),
409        }
410    }
411}
412
413/// The mutable state of building an actor graph. See [`ActorGraphBuildStateInner`].
414struct ActorGraphBuildState {
415    /// The actual state.
416    inner: ActorGraphBuildStateInner,
417}
418
419impl ActorGraphBuildState {
420    /// Create an empty state with the given id generator.
421    fn new() -> Self {
422        Self {
423            inner: Default::default(),
424        }
425    }
426
427    /// Finish the build and return the inner state.
428    fn finish(self) -> ActorGraphBuildStateInner {
429        self.inner
430    }
431}
432
433/// The result of a built actor graph. Will be further embedded into the `Context` for building
434/// actors on the compute nodes.
435pub struct ActorGraphBuildResult {
436    /// The graph of sealed fragments, including all actors.
437    pub graph: BTreeMap<FragmentId, Fragment>,
438    /// The downstream fragments of the fragments from the new graph to be created.
439    /// Including the fragment relation to external downstream fragment.
440    pub downstream_fragment_relations: FragmentDownstreamRelation,
441
442    /// The new dispatchers to be added to the upstream mview actors. Used for MV on MV.
443    pub upstream_fragment_downstreams: FragmentDownstreamRelation,
444
445    /// The updates to be applied to the downstream fragment merge node. Used for schema change (replace
446    /// table plan).
447    pub replace_upstream: FragmentReplaceUpstream,
448}
449
450/// [`ActorGraphBuilder`] builds the actor graph for the given complete fragment graph, based on the
451/// current cluster info and the required parallelism.
452#[derive(Debug)]
453pub struct ActorGraphBuilder {
454    /// The pre-scheduled distribution for each building fragment.
455    distributions: HashMap<GlobalFragmentId, Distribution>,
456
457    /// The complete fragment graph.
458    fragment_graph: CompleteStreamFragmentGraph,
459}
460
461impl ActorGraphBuilder {
462    /// Create a new actor graph builder with the given "complete" graph. Returns an error if the
463    /// graph is failed to be scheduled.
464    pub fn new(fragment_graph: CompleteStreamFragmentGraph) -> MetaResult<Self> {
465        let expected_vnode_count = fragment_graph.max_parallelism();
466        let scheduler = schedule::Scheduler::new(expected_vnode_count)?;
467
468        let distributions = scheduler.schedule(&fragment_graph)?;
469
470        // Fill the vnode count for each internal table, based on schedule result.
471        let mut fragment_graph = fragment_graph;
472        for (id, fragment) in fragment_graph.building_fragments_mut() {
473            let mut error = None;
474            let fragment_vnode_count = distributions[id].vnode_count();
475            visit_tables(fragment, |table, _| {
476                if error.is_some() {
477                    return;
478                }
479                // There are special cases where a hash-distributed fragment contains singleton
480                // internal tables, e.g., the state table of `Source` executors.
481                let vnode_count = if table.is_singleton() {
482                    if fragment_vnode_count > 1 {
483                        tracing::info!(
484                            table.name,
485                            "found singleton table in hash-distributed fragment"
486                        );
487                    }
488                    1
489                } else {
490                    fragment_vnode_count
491                };
492                match table.vnode_count_inner().value_opt() {
493                    // Vnode count of this table is not set to placeholder, meaning that we are replacing
494                    // a streaming job, and the existing state table requires a specific vnode count.
495                    // Check if it's the same with what we derived from the schedule result.
496                    //
497                    // Typically, inconsistency should not happen as we force to align the vnode count
498                    // when planning the new streaming job in the frontend.
499                    Some(required_vnode_count) if required_vnode_count != vnode_count => {
500                        error = Some(format!(
501                            "failed to align vnode count for table {}({}): required {}, but got {}",
502                            table.id, table.name, required_vnode_count, vnode_count
503                        ));
504                    }
505                    // Normal cases.
506                    _ => table.maybe_vnode_count = VnodeCount::set(vnode_count).to_protobuf(),
507                }
508            });
509            if let Some(error) = error {
510                bail!(error);
511            }
512        }
513
514        Ok(Self {
515            distributions,
516            fragment_graph,
517        })
518    }
519
520    /// Build a stream graph by duplicating each fragment as parallel actors. Returns
521    /// [`ActorGraphBuildResult`] that will be further used to build actors on the compute nodes.
522    pub fn generate_graph(self) -> MetaResult<ActorGraphBuildResult> {
523        // Build the actor graph and get the final state.
524        let ActorGraphBuildStateInner {
525            fragment_actor_builders,
526            downstream_fragment_changes,
527            upstream_fragment_changes,
528        } = self.build_actor_graph()?;
529
530        let mut downstream_fragment_relations: FragmentDownstreamRelation = HashMap::new();
531        // Serialize the graph into sealed fragments.
532        let graph = {
533            // As all fragments are processed, we can now `rewrite` the stream nodes where the
534            // `Exchange` and `Chain` are rewritten.
535            let mut fragment_nodes: HashMap<GlobalFragmentId, StreamNode> = HashMap::new();
536
537            for (fragment_id, builder) in fragment_actor_builders {
538                let global_fragment_id = fragment_id.as_global_id();
539                let node = builder.rewrite()?;
540                downstream_fragment_relations
541                    .try_insert(
542                        global_fragment_id,
543                        builder
544                            .downstreams
545                            .into_iter()
546                            .map(|(id, dispatch)| (id.as_global_id(), dispatch).into())
547                            .collect(),
548                    )
549                    .expect("non-duplicate");
550                fragment_nodes
551                    .try_insert(fragment_id, node)
552                    .expect("non-duplicate");
553            }
554
555            let mut graph = BTreeMap::new();
556            for (fragment_id, stream_node) in fragment_nodes {
557                let distribution = self.distributions[&fragment_id].clone();
558                let fragment =
559                    self.fragment_graph
560                        .seal_fragment(fragment_id, distribution, stream_node);
561                let fragment_id = fragment_id.as_global_id();
562                graph.insert(fragment_id, fragment);
563            }
564            graph
565        };
566
567        // Extract the new fragment relation from the external changes.
568        let upstream_fragment_downstreams = upstream_fragment_changes
569            .into_iter()
570            .map(|(fragment_id, changes)| {
571                (
572                    fragment_id.as_global_id(),
573                    changes
574                        .new_downstreams
575                        .into_iter()
576                        .map(|(downstream_fragment_id, new_dispatch)| {
577                            (downstream_fragment_id.as_global_id(), new_dispatch).into()
578                        })
579                        .collect(),
580                )
581            })
582            .collect();
583
584        // Extract the updates for merge executors from the external changes.
585        let replace_upstream = downstream_fragment_changes
586            .into_iter()
587            .map(|(fragment_id, changes)| {
588                let fragment_id = fragment_id.as_global_id();
589                (
590                    fragment_id,
591                    changes
592                        .new_upstreams
593                        .into_iter()
594                        .map(move |(edge_id, upstream_fragment_id)| {
595                            let upstream_fragment_id = upstream_fragment_id.as_global_id();
596                            let DownstreamExternalEdgeId {
597                                original_upstream_fragment_id,
598                                ..
599                            } = edge_id;
600                            (
601                                original_upstream_fragment_id.as_global_id(),
602                                upstream_fragment_id,
603                            )
604                        })
605                        .collect(),
606                )
607            })
608            .filter(|(_, fragment_changes): &(_, HashMap<_, _>)| !fragment_changes.is_empty())
609            .collect();
610
611        Ok(ActorGraphBuildResult {
612            graph,
613            downstream_fragment_relations,
614            upstream_fragment_downstreams,
615            replace_upstream,
616        })
617    }
618
619    /// Build actor graph for each fragment, using topological order.
620    fn build_actor_graph(&self) -> MetaResult<ActorGraphBuildStateInner> {
621        let mut state = ActorGraphBuildState::new();
622
623        // Use topological sort to build the graph from downstream to upstream. (The first fragment
624        // popped out from the heap will be the top-most node in plan, or the sink in stream graph.)
625        for fragment_id in self.fragment_graph.topo_order()? {
626            self.build_actor_graph_fragment(fragment_id, &mut state)?;
627        }
628
629        Ok(state.finish())
630    }
631
632    /// Build actor graph for a specific fragment.
633    fn build_actor_graph_fragment(
634        &self,
635        fragment_id: GlobalFragmentId,
636        state: &mut ActorGraphBuildState,
637    ) -> MetaResult<()> {
638        let current_fragment = self.fragment_graph.get_fragment(fragment_id);
639
640        // First, add or record the actors for the current fragment into the state.
641        match current_fragment {
642            // For building fragments, we need to generate the actor builders.
643            EitherFragment::Building(current_fragment) => {
644                let node = current_fragment.node.clone().unwrap();
645                state
646                    .inner
647                    .fragment_actor_builders
648                    .try_insert(fragment_id, FragmentActorBuilder::new(fragment_id, node))
649                    .expect("non-duplicate");
650            }
651
652            // For existing fragments, we only need to record the actor locations.
653            EitherFragment::Existing => {}
654        };
655
656        // Then, add links between the current fragment and its downstream fragments.
657        for (downstream_fragment_id, edge) in self.fragment_graph.get_downstreams(fragment_id) {
658            state.inner.add_link(
659                FragmentLinkNode { fragment_id },
660                FragmentLinkNode {
661                    fragment_id: downstream_fragment_id,
662                },
663                edge,
664            );
665        }
666
667        Ok(())
668    }
669}