reifydb_sub_flow/execution/
batch.rs1use std::collections::{BTreeMap, HashMap};
5
6use reifydb_core::{
7 common::CommitVersion,
8 interface::{
9 catalog::flow::{FlowId, FlowNodeId},
10 change::Change,
11 },
12};
13use reifydb_rql::flow::flow::FlowDag;
14use reifydb_value::Result;
15use tracing::{Span, field, instrument};
16
17use crate::{engine::FlowEngineInner, transaction::FlowTransaction};
18
19impl FlowEngineInner {
20 #[instrument(name = "flow::engine::process", level = "debug", skip(self, txn, change), fields(
21 flow_id = ?flow_id,
22 origin = ?change.origin,
23 version = change.version.0,
24 diff_count = change.diffs.len(),
25 row_count = change.row_count(),
26 nodes_processed = field::Empty
27 ))]
28 pub fn process(&self, txn: &mut FlowTransaction, change: Change, flow_id: FlowId) -> Result<()> {
29 self.process_batch(txn, vec![change], flow_id)
30 }
31
32 #[instrument(name = "flow::engine::process_batch", level = "debug", skip(self, txn, changes), fields(
33 flow_id = ?flow_id,
34 batch_change_count = changes.len(),
35 batch_row_count = changes.iter().map(Change::row_count).sum::<usize>(),
36 version_count = field::Empty,
37 nodes_processed = field::Empty
38 ))]
39 pub fn process_batch(&self, txn: &mut FlowTransaction, changes: Vec<Change>, flow_id: FlowId) -> Result<()> {
40 let flow = match self.flows.get(&flow_id) {
41 Some(f) => f.clone(),
42 None => return Ok(()),
43 };
44
45 let mut by_version: BTreeMap<CommitVersion, Vec<Change>> = BTreeMap::new();
46 for change in changes {
47 by_version.entry(change.version).or_default().push(change);
48 }
49 Span::current().record("version_count", by_version.len());
50
51 let topo = flow.topological_order()?;
52 let mut nodes_processed = 0u32;
53
54 for (_, version_changes) in by_version {
55 nodes_processed += self.process_version(txn, &flow, flow_id, version_changes, &topo)?;
56 }
57
58 Span::current().record("nodes_processed", nodes_processed);
59 Ok(())
60 }
61
62 #[inline]
63 fn process_version(
64 &self,
65 txn: &mut FlowTransaction,
66 flow: &FlowDag,
67 flow_id: FlowId,
68 version_changes: Vec<Change>,
69 topo: &[FlowNodeId],
70 ) -> Result<u32> {
71 let mut pending: HashMap<FlowNodeId, Vec<Change>> = HashMap::new();
72 for change in version_changes {
73 self.seed_entry_nodes(flow, flow_id, change, &mut pending);
74 }
75
76 let mut nodes_processed = 0u32;
77 for node_id in topo {
78 let inbox = match pending.remove(node_id) {
79 Some(v) if !v.is_empty() => v,
80 _ => continue,
81 };
82
83 let node = match flow.get_node(node_id) {
84 Some(n) => n.clone(),
85 None => continue,
86 };
87
88 let combined_output = self.dispatch_node(txn, &node, inbox)?;
89 nodes_processed += 1;
90 if combined_output.diffs.is_empty() {
91 continue;
92 }
93
94 let child_count = node.outputs.len();
95 for (child_idx, child_id) in node.outputs.iter().enumerate() {
96 if child_idx + 1 == child_count {
97 pending.entry(*child_id).or_default().push(combined_output);
98 break;
99 }
100 pending.entry(*child_id).or_default().push(combined_output.clone());
101 }
102 }
103 Ok(nodes_processed)
104 }
105}