Skip to main content

reifydb_sub_flow/execution/
batch.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}