1use 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
92unsafe 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}