Skip to main content

reifydb_sub_flow/operator/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::{ops::Deref, sync::Arc};
5
6use reifydb_abi::operator::capabilities::OperatorCapability;
7use reifydb_core::{interface::catalog::flow::FlowNodeId, value::column::columns::Columns};
8use reifydb_sdk::operator::Tick;
9use reifydb_value::{Result, value::duration::Duration};
10
11use crate::transaction::FlowTransaction;
12
13pub mod append;
14pub mod apply;
15#[cfg(reifydb_target = "native")]
16pub mod context;
17pub mod distinct;
18pub mod extend;
19#[cfg(reifydb_target = "native")]
20pub mod ffi;
21pub mod filter;
22pub mod gate;
23pub mod guard;
24pub mod join;
25pub mod map;
26#[cfg(reifydb_target = "native")]
27pub mod native;
28pub mod scan;
29pub mod sink;
30pub mod sort;
31pub mod stateful;
32pub mod take;
33pub mod window;
34
35use append::AppendOperator;
36use apply::ApplyOperator;
37use distinct::operator::DistinctOperator;
38use extend::ExtendOperator;
39use filter::FilterOperator;
40use gate::GateOperator;
41use guard::{enforce_apply_capabilities, enforce_tick_capability};
42use join::operator::JoinOperator;
43use map::MapOperator;
44use reifydb_core::interface::change::Change;
45use scan::{
46	dictionary::PrimitiveDictionaryOperator, flow::PrimitiveFlowOperator, ringbuffer::PrimitiveRingBufferOperator,
47	series::PrimitiveSeriesOperator, table::PrimitiveTableOperator, view::PrimitiveViewOperator,
48};
49use sink::{
50	ringbuffer_view::SinkRingBufferViewOperator, series_view::SinkSeriesViewOperator, view::SinkTableViewOperator,
51};
52use sort::SortOperator;
53use take::TakeOperator;
54use window::{aggregate::AggregateOperator, operator::WindowOperator};
55
56pub trait Operator: Send {
57	fn id(&self) -> FlowNodeId;
58
59	fn capabilities(&self) -> &[OperatorCapability];
60
61	fn apply(&self, txn: &mut FlowTransaction, change: Change) -> Result<Change>;
62
63	fn tick(&self, _txn: &mut FlowTransaction, _tick: Tick) -> Result<Option<Change>> {
64		Ok(None)
65	}
66
67	fn ticks(&self) -> Option<Duration> {
68		None
69	}
70}
71
72pub type BoxedOperator = Box<dyn Operator + Send>;
73
74#[derive(Clone)]
75pub struct OperatorCell(Arc<Operators>);
76
77impl OperatorCell {
78	#[allow(clippy::arc_with_non_send_sync)]
79	pub fn new(operators: Operators) -> Self {
80		Self(Arc::new(operators))
81	}
82}
83
84impl Deref for OperatorCell {
85	type Target = Operators;
86
87	fn deref(&self) -> &Operators {
88		&self.0
89	}
90}
91
92// SAFETY: a flow and all of its operators are only ever accessed by a single thread at any one
93// time. Flows that execute in parallel on the rayon commit pool own disjoint operator sets
94// (operators are keyed by FlowNodeId and never shared between flows), so no Operators value is ever
95// reachable from two threads simultaneously. The inner Arc is only cloned and dereferenced from the
96// owning thread, so asserting Send and Sync over the !Sync Operators it holds is sound.
97unsafe impl Send for OperatorCell {}
98unsafe impl Sync for OperatorCell {}
99
100pub enum Operators {
101	SourceTable(PrimitiveTableOperator),
102	SourceView(PrimitiveViewOperator),
103	SourceFlow(PrimitiveFlowOperator),
104	SourceRingBuffer(PrimitiveRingBufferOperator),
105	SourceSeries(PrimitiveSeriesOperator),
106	SourceDictionary(PrimitiveDictionaryOperator),
107	Filter(FilterOperator),
108	Gate(GateOperator),
109	Map(MapOperator),
110	Extend(ExtendOperator),
111	Join(JoinOperator),
112	Sort(SortOperator),
113	Take(TakeOperator),
114	Distinct(DistinctOperator),
115	Append(AppendOperator),
116	Apply(ApplyOperator),
117	SinkTableView(SinkTableViewOperator),
118	SinkRingBufferView(SinkRingBufferViewOperator),
119	SinkSeriesView(SinkSeriesViewOperator),
120	Window(WindowOperator),
121	Aggregate(AggregateOperator),
122	Custom(BoxedOperator),
123}
124
125impl Operators {
126	pub fn id(&self) -> FlowNodeId {
127		match self {
128			Operators::Filter(op) => op.id(),
129			Operators::Gate(op) => op.id(),
130			Operators::Map(op) => op.id(),
131			Operators::Extend(op) => op.id(),
132			Operators::Join(op) => op.id(),
133			Operators::Sort(op) => op.id(),
134			Operators::Take(op) => op.id(),
135			Operators::Distinct(op) => op.id(),
136			Operators::Append(op) => op.id(),
137			Operators::Apply(op) => op.id(),
138			Operators::SinkTableView(op) => op.id(),
139			Operators::SinkRingBufferView(op) => op.id(),
140			Operators::SinkSeriesView(op) => op.id(),
141			Operators::Window(op) => op.id(),
142			Operators::Aggregate(op) => op.id(),
143			Operators::SourceTable(op) => op.id(),
144			Operators::SourceView(op) => op.id(),
145			Operators::SourceFlow(op) => op.id(),
146			Operators::SourceRingBuffer(op) => op.id(),
147			Operators::SourceSeries(op) => op.id(),
148			Operators::SourceDictionary(op) => op.id(),
149			Operators::Custom(op) => op.id(),
150		}
151	}
152
153	pub fn capabilities(&self) -> &[OperatorCapability] {
154		match self {
155			Operators::Filter(op) => op.capabilities(),
156			Operators::Gate(op) => op.capabilities(),
157			Operators::Map(op) => op.capabilities(),
158			Operators::Extend(op) => op.capabilities(),
159			Operators::Join(op) => op.capabilities(),
160			Operators::Sort(op) => op.capabilities(),
161			Operators::Take(op) => op.capabilities(),
162			Operators::Distinct(op) => op.capabilities(),
163			Operators::Append(op) => op.capabilities(),
164			Operators::Apply(op) => op.capabilities(),
165			Operators::SinkTableView(op) => op.capabilities(),
166			Operators::SinkRingBufferView(op) => op.capabilities(),
167			Operators::SinkSeriesView(op) => op.capabilities(),
168			Operators::Window(op) => op.capabilities(),
169			Operators::Aggregate(op) => op.capabilities(),
170			Operators::SourceTable(op) => op.capabilities(),
171			Operators::SourceView(op) => op.capabilities(),
172			Operators::SourceFlow(op) => op.capabilities(),
173			Operators::SourceRingBuffer(op) => op.capabilities(),
174			Operators::SourceSeries(op) => op.capabilities(),
175			Operators::SourceDictionary(op) => op.capabilities(),
176			Operators::Custom(op) => op.capabilities(),
177		}
178	}
179
180	pub fn ticks(&self) -> Option<Duration> {
181		match self {
182			Operators::Filter(op) => op.ticks(),
183			Operators::Gate(op) => op.ticks(),
184			Operators::Map(op) => op.ticks(),
185			Operators::Extend(op) => op.ticks(),
186			Operators::Join(op) => op.ticks(),
187			Operators::Sort(op) => op.ticks(),
188			Operators::Take(op) => op.ticks(),
189			Operators::Distinct(op) => op.ticks(),
190			Operators::Append(op) => op.ticks(),
191			Operators::Apply(op) => op.ticks(),
192			Operators::SinkTableView(op) => op.ticks(),
193			Operators::SinkRingBufferView(op) => op.ticks(),
194			Operators::SinkSeriesView(op) => op.ticks(),
195			Operators::Window(op) => op.ticks(),
196			Operators::Aggregate(op) => op.ticks(),
197			Operators::SourceTable(op) => op.ticks(),
198			Operators::SourceView(op) => op.ticks(),
199			Operators::SourceFlow(op) => op.ticks(),
200			Operators::SourceRingBuffer(op) => op.ticks(),
201			Operators::SourceSeries(op) => op.ticks(),
202			Operators::SourceDictionary(op) => op.ticks(),
203			Operators::Custom(op) => op.ticks(),
204		}
205	}
206
207	pub fn apply(&self, txn: &mut FlowTransaction, change: Change) -> Result<Change> {
208		enforce_apply_capabilities(self.id(), self.capabilities(), &change);
209		match self {
210			Operators::Filter(op) => op.apply(txn, change),
211			Operators::Gate(op) => op.apply(txn, change),
212			Operators::Map(op) => op.apply(txn, change),
213			Operators::Extend(op) => op.apply(txn, change),
214			Operators::Join(op) => op.apply(txn, change),
215			Operators::Sort(op) => op.apply(txn, change),
216			Operators::Take(op) => op.apply(txn, change),
217			Operators::Distinct(op) => op.apply(txn, change),
218			Operators::Append(op) => op.apply(txn, change),
219			Operators::Apply(op) => op.apply(txn, change),
220			Operators::SinkTableView(op) => op.apply(txn, change),
221			Operators::SinkRingBufferView(op) => op.apply(txn, change),
222			Operators::SinkSeriesView(op) => op.apply(txn, change),
223			Operators::Window(op) => op.apply(txn, change),
224			Operators::Aggregate(op) => op.apply(txn, change),
225			Operators::SourceTable(op) => op.apply(txn, change),
226			Operators::SourceView(op) => op.apply(txn, change),
227			Operators::SourceFlow(op) => op.apply(txn, change),
228			Operators::SourceRingBuffer(op) => op.apply(txn, change),
229			Operators::SourceSeries(op) => op.apply(txn, change),
230			Operators::SourceDictionary(op) => op.apply(txn, change),
231			Operators::Custom(op) => op.apply(txn, change),
232		}
233	}
234
235	pub fn tick(&self, txn: &mut FlowTransaction, tick: Tick) -> Result<Option<Change>> {
236		match self {
237			Operators::Window(op) => {
238				enforce_tick_capability(op.id(), op.capabilities());
239				op.tick(txn, tick)
240			}
241			Operators::Custom(op) => {
242				enforce_tick_capability(op.id(), op.capabilities());
243				op.tick(txn, tick)
244			}
245			Operators::Apply(op) => {
246				enforce_tick_capability(op.id(), op.capabilities());
247				op.tick(txn, tick)
248			}
249			Operators::Distinct(op) => {
250				enforce_tick_capability(op.id(), op.capabilities());
251				op.tick(txn, tick)
252			}
253			Operators::Join(op) => {
254				enforce_tick_capability(op.id(), op.capabilities());
255				op.tick(txn, tick)
256			}
257			Operators::Append(op) => {
258				enforce_tick_capability(op.id(), op.capabilities());
259				op.tick(txn, tick)
260			}
261			_ => Ok(None),
262		}
263	}
264
265	pub fn output_schema(&self) -> Option<Columns> {
266		match self {
267			Operators::SourceTable(op) => Some(op.output_schema()),
268			Operators::SourceView(op) => Some(op.output_schema()),
269			Operators::SourceRingBuffer(op) => Some(op.output_schema()),
270			Operators::SourceSeries(_) => Some(Columns::empty()),
271			Operators::SourceDictionary(_) => Some(Columns::empty()),
272			Operators::SourceFlow(_) => Some(Columns::empty()),
273			Operators::Filter(op) => op.output_schema(),
274			Operators::Gate(op) => op.output_schema(),
275			Operators::Map(op) => op.output_schema(),
276			Operators::Extend(op) => op.output_schema(),
277			Operators::Sort(op) => op.output_schema(),
278			Operators::Take(op) => op.output_schema(),
279			Operators::Distinct(op) => op.output_schema(),
280			Operators::Append(op) => op.output_schema(),
281			Operators::Window(op) => op.core.parent.output_schema(),
282			Operators::Aggregate(op) => op.output_schema(),
283			Operators::Apply(op) => op.output_schema(),
284			Operators::Join(_) => None,
285			Operators::SinkTableView(_) => None,
286			Operators::SinkRingBufferView(_) => None,
287			Operators::SinkSeriesView(_) => None,
288			Operators::Custom(_) => None,
289		}
290	}
291}