reifydb_sub_flow/operator/sink/
series_view.rs1use reifydb_abi::operator::capabilities::OperatorCapability;
5use reifydb_codec::{
6 encoded::{row::EncodedRow, shape::RowShape},
7 key::encoded::EncodedKey,
8};
9use reifydb_core::{
10 interface::{
11 catalog::{flow::FlowNodeId, id::SeriesId, series::SeriesKey, shape::ShapeId, view::View},
12 change::{Change, ChangeOrigin, Diff},
13 resolved::ResolvedView,
14 },
15 key::row::RowKey,
16 row::row_shape_from_columns,
17 value::column::columns::Columns,
18};
19use reifydb_value::{
20 Result,
21 value::{datetime::DateTime, row_number::RowNumber},
22};
23use smallvec::smallvec;
24
25use super::{coerce_columns, encode_row_at_index, shape_field_columns, view::dictionary_encode_view_columns};
26use crate::{Operator, operator::OperatorCell, transaction::FlowTransaction};
27
28pub struct SinkSeriesViewOperator {
29 #[allow(dead_code)]
30 parent: OperatorCell,
31 node: FlowNodeId,
32 view: ResolvedView,
33 series_id: SeriesId,
34 #[allow(dead_code)]
35 key: SeriesKey,
36}
37
38impl SinkSeriesViewOperator {
39 pub fn new(
40 parent: OperatorCell,
41 node: FlowNodeId,
42 view: ResolvedView,
43 series_id: SeriesId,
44 key: SeriesKey,
45 ) -> Self {
46 Self {
47 parent,
48 node,
49 view,
50 series_id,
51 key,
52 }
53 }
54}
55
56impl Operator for SinkSeriesViewOperator {
57 fn id(&self) -> FlowNodeId {
58 self.node
59 }
60
61 fn capabilities(&self) -> &[OperatorCapability] {
62 OperatorCapability::STANDARD
63 }
64
65 fn apply(&self, txn: &mut FlowTransaction, change: Change) -> Result<Change> {
66 let view = self.view.def().clone();
67 let shape = row_shape_from_columns(view.columns());
68 let object_id = ShapeId::series(self.series_id);
69
70 for diff in change.diffs.iter() {
71 match diff {
72 Diff::Insert {
73 post,
74 ..
75 } => self.apply_series_view_insert(txn, &view, &shape, object_id, post)?,
76 Diff::Update {
77 pre,
78 post,
79 ..
80 } => self.apply_series_view_update(txn, &view, &shape, object_id, pre, post)?,
81 Diff::Remove {
82 pre,
83 ..
84 } => self.apply_series_view_remove(txn, &view, object_id, pre)?,
85 }
86 }
87
88 Ok(Change::from_flow(self.node, change.version, Vec::new(), change.changed_at))
89 }
90}
91
92impl SinkSeriesViewOperator {
93 #[inline]
94 fn apply_series_view_insert(
95 &self,
96 txn: &mut FlowTransaction,
97 view: &View,
98 shape: &RowShape,
99 object_id: ShapeId,
100 post: &Columns,
101 ) -> Result<()> {
102 let coerced = coerce_columns(post, view.columns())?;
103 let dict_encoded = dictionary_encode_view_columns(txn, view, &coerced)?;
104 let source = dict_encoded.as_ref().unwrap_or(&coerced);
105 let row_count = source.row_count();
106 let field_columns = shape_field_columns(source, shape);
107 let mut ids: Vec<RowNumber> = Vec::with_capacity(row_count);
108 let mut encoded_rows: Vec<EncodedRow> = Vec::with_capacity(row_count);
109 for row_idx in 0..row_count {
110 let row_number = source.row_numbers[row_idx];
111 let (_, encoded) = encode_row_at_index(source, row_idx, shape, row_number, &field_columns)?;
112 ids.push(row_number);
113 encoded_rows.push(encoded);
114 }
115 for (row_number, encoded) in ids.iter().zip(encoded_rows.iter()) {
116 let key = RowKey::encoded(object_id, *row_number);
117 txn.set(&key, encoded.clone())?;
118 }
119 emit_view_change(txn, view, Diff::insert(coerced));
120 Ok(())
121 }
122
123 #[inline]
124 fn apply_series_view_update(
125 &self,
126 txn: &mut FlowTransaction,
127 view: &View,
128 shape: &RowShape,
129 object_id: ShapeId,
130 pre: &Columns,
131 post: &Columns,
132 ) -> Result<()> {
133 let coerced_pre = coerce_columns(pre, view.columns())?;
134 let coerced_post = coerce_columns(post, view.columns())?;
135 let dict_pre = dictionary_encode_view_columns(txn, view, &coerced_pre)?;
136 let dict_post = dictionary_encode_view_columns(txn, view, &coerced_post)?;
137 let source_pre = dict_pre.as_ref().unwrap_or(&coerced_pre);
138 let source_post = dict_post.as_ref().unwrap_or(&coerced_post);
139 let row_count = source_post.row_count();
140 let field_columns = shape_field_columns(source_post, shape);
141 let mut pre_keys: Vec<EncodedKey> = Vec::with_capacity(row_count);
142 let mut post_keys: Vec<EncodedKey> = Vec::with_capacity(row_count);
143 let mut post_encoded_rows: Vec<EncodedRow> = Vec::with_capacity(row_count);
144 for row_idx in 0..row_count {
145 let pre_row_number = source_pre.row_numbers[row_idx];
146 let post_row_number = source_post.row_numbers[row_idx];
147 let (_, post_encoded) =
148 encode_row_at_index(source_post, row_idx, shape, post_row_number, &field_columns)?;
149
150 pre_keys.push(RowKey::encoded(object_id, pre_row_number));
151 post_keys.push(RowKey::encoded(object_id, post_row_number));
152 post_encoded_rows.push(post_encoded);
153 }
154 for ((pre_key, post_key), post_encoded) in
155 pre_keys.iter().zip(post_keys.iter()).zip(post_encoded_rows.iter())
156 {
157 txn.remove(pre_key)?;
158 txn.set(post_key, post_encoded.clone())?;
159 }
160 emit_view_change(txn, view, Diff::update(coerced_pre, coerced_post));
161 Ok(())
162 }
163
164 #[inline]
165 fn apply_series_view_remove(
166 &self,
167 txn: &mut FlowTransaction,
168 view: &View,
169 object_id: ShapeId,
170 pre: &Columns,
171 ) -> Result<()> {
172 let coerced = coerce_columns(pre, view.columns())?;
173 let row_count = coerced.row_count();
174 let mut ids: Vec<RowNumber> = Vec::with_capacity(row_count);
175 for row_idx in 0..row_count {
176 let row_number = coerced.row_numbers[row_idx];
177 ids.push(row_number);
178 }
179 for row_number in ids.iter() {
180 let key = RowKey::encoded(object_id, *row_number);
181 txn.remove(&key)?;
182 }
183 emit_view_change(txn, view, Diff::remove(coerced));
184 Ok(())
185 }
186}
187
188#[inline]
189fn emit_view_change(txn: &mut FlowTransaction, view: &View, diff: Diff) {
190 let version = txn.version();
191 let changed_at = DateTime::from_nanos(txn.clock().now_nanos());
192 txn.track_flow_change(Change {
193 origin: ChangeOrigin::Shape(ShapeId::view(view.id())),
194 version,
195 diffs: smallvec![diff],
196 changed_at,
197 });
198}