Skip to main content

reifydb_sub_flow/operator/sink/
series_view.rs

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