Skip to main content

reifydb_sub_flow/operator/
append.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::{cell::RefCell, ops::Bound};
5
6use reifydb_abi::operator::capabilities::OperatorCapability;
7use reifydb_codec::{
8	encoded::shape::RowShape,
9	key::{
10		encoded::{EncodedKey, EncodedKeyRange},
11		serializer::KeySerializer,
12	},
13};
14use reifydb_core::{
15	common::CommitVersion,
16	interface::{
17		catalog::flow::FlowNodeId,
18		change::{Change, ChangeOrigin, Diff},
19	},
20	value::column::columns::Columns,
21};
22use reifydb_runtime::version_epoch::VersionEpoch;
23use reifydb_sdk::operator::Tick;
24use reifydb_value::{
25	Result,
26	error::Error,
27	reifydb_assertions,
28	value::{duration::Duration, row_number::RowNumber},
29};
30
31use crate::{
32	error::FlowGraphError,
33	operator::{
34		Operator, OperatorCell,
35		stateful::{
36			row::RowNumberProvider,
37			utils::{internal_state_drop, internal_state_range_versioned, internal_state_set},
38		},
39	},
40	transaction::FlowTransaction,
41};
42
43const TIMESTAMP_PREFIX: u8 = b'T';
44
45pub struct AppendOperator {
46	node: FlowNodeId,
47
48	parents: Vec<OperatorCell>,
49
50	input_nodes: Vec<FlowNodeId>,
51
52	row_number_provider: RowNumberProvider,
53
54	ttl_nanos: Option<u64>,
55
56	version_epoch: VersionEpoch,
57
58	evict_cursor: RefCell<Option<EncodedKey>>,
59}
60
61impl AppendOperator {
62	pub fn new(
63		node: FlowNodeId,
64		parents: Vec<OperatorCell>,
65		input_nodes: Vec<FlowNodeId>,
66		ttl_nanos: Option<u64>,
67		version_epoch: VersionEpoch,
68	) -> Self {
69		reifydb_assertions! {
70			assert_eq!(parents.len(), input_nodes.len());
71			assert!(parents.len() >= 2, "Append requires at least 2 inputs");
72		}
73
74		Self {
75			node,
76			parents,
77			input_nodes,
78			row_number_provider: RowNumberProvider::new(node),
79			ttl_nanos,
80			version_epoch,
81			evict_cursor: RefCell::new(None),
82		}
83	}
84
85	#[cfg(test)]
86	pub(crate) fn new_for_state_tests(node: FlowNodeId, ttl_nanos: Option<u64>) -> Self {
87		Self {
88			node,
89			parents: Vec::new(),
90			input_nodes: Vec::new(),
91			row_number_provider: RowNumberProvider::new(node),
92			ttl_nanos,
93			version_epoch: VersionEpoch::new(),
94			evict_cursor: RefCell::new(None),
95		}
96	}
97
98	pub(crate) fn output_schema(&self) -> Option<Columns> {
99		self.parents[0].output_schema()
100	}
101
102	fn parent_index_for_origin(&self, origin: &ChangeOrigin) -> Option<usize> {
103		match origin {
104			ChangeOrigin::Flow(from_node) => self.input_nodes.iter().position(|n| n == from_node),
105			ChangeOrigin::Shape(_) => None,
106		}
107	}
108
109	fn make_composite_key(parent_index: u8, source_row: RowNumber) -> EncodedKey {
110		let mut serializer = KeySerializer::new();
111		serializer.extend_u8(parent_index);
112		serializer.extend_u64(source_row.0);
113		serializer.finish()
114	}
115
116	fn make_timestamp_key(composite_key: &EncodedKey) -> EncodedKey {
117		let mut bytes = Vec::with_capacity(1 + composite_key.len());
118		bytes.push(TIMESTAMP_PREFIX);
119		bytes.extend_from_slice(composite_key.as_ref());
120		EncodedKey::new(bytes)
121	}
122
123	fn touch(&self, txn: &mut FlowTransaction, composite_key: &EncodedKey) -> Result<()> {
124		if self.ttl_nanos.is_none() {
125			return Ok(());
126		}
127		let key = Self::make_timestamp_key(composite_key);
128		let row = RowShape::operator_state().allocate();
129		internal_state_set(self.node, txn, &key, row)
130	}
131
132	fn forget_mapping(&self, txn: &mut FlowTransaction, composite_key: &EncodedKey) -> Result<()> {
133		self.row_number_provider.remove_for_key(txn, composite_key)?;
134		let ts_key = Self::make_timestamp_key(composite_key);
135		internal_state_drop(self.node, txn, &ts_key)
136	}
137}
138
139impl Operator for AppendOperator {
140	fn id(&self) -> FlowNodeId {
141		self.node
142	}
143
144	fn capabilities(&self) -> &[OperatorCapability] {
145		OperatorCapability::STANDARD_WITH_TICK
146	}
147
148	fn ticks(&self) -> Option<Duration> {
149		if self.ttl_nanos.is_some() {
150			Some(Duration::from_seconds(1).unwrap())
151		} else {
152			None
153		}
154	}
155
156	fn apply(&self, txn: &mut FlowTransaction, change: Change) -> Result<Change> {
157		let parent_origin = change.origin.clone();
158		let mut result_diffs = Vec::with_capacity(change.diffs.len());
159
160		for diff in change.diffs {
161			let diff_origin = diff.origin().cloned().unwrap_or_else(|| parent_origin.clone());
162			let parent_index = self.parent_index_for_origin(&diff_origin).ok_or_else(|| {
163				Error::from(FlowGraphError::UnknownDiffOrigin {
164					operator: "Append",
165					origin: Some(format!("{:?}", diff_origin)),
166				})
167			})?;
168			match diff {
169				Diff::Insert {
170					post,
171					..
172				} => {
173					if let Some(d) = self.translate_append_insert(txn, parent_index, post)? {
174						result_diffs.push(d);
175					}
176				}
177				Diff::Update {
178					pre,
179					post,
180					..
181				} => {
182					if let Some(d) = self.translate_append_update(txn, parent_index, pre, post)? {
183						result_diffs.push(d);
184					}
185				}
186				Diff::Remove {
187					pre,
188					..
189				} => {
190					if let Some(d) = self.translate_append_remove(txn, parent_index, pre)? {
191						result_diffs.push(d);
192					}
193				}
194			}
195		}
196
197		Ok(Change::from_flow(self.node, change.version, result_diffs, change.changed_at))
198	}
199
200	fn tick(&self, txn: &mut FlowTransaction, tick: Tick) -> Result<Option<Change>> {
201		let Some(ttl_nanos) = self.ttl_nanos else {
202			return Ok(None);
203		};
204
205		let now_nanos = tick.now.to_nanos();
206		let Some(cutoff_nanos) = now_nanos.checked_sub(ttl_nanos) else {
207			return Ok(None);
208		};
209		let Some(cutoff_version) = self.version_epoch.floor_version_at(cutoff_nanos).map(CommitVersion) else {
210			return Ok(None);
211		};
212
213		const EVICT_BATCH: usize = 4096;
214		let prefix = [TIMESTAMP_PREFIX];
215		let base = EncodedKeyRange::prefix(&prefix);
216		let start = match self.evict_cursor.borrow().clone() {
217			Some(cursor) => Bound::Excluded(cursor),
218			None => base.start.clone(),
219		};
220		let range = EncodedKeyRange::new(start, base.end.clone());
221		let batch = internal_state_range_versioned(self.node, txn, range)
222			.take(EVICT_BATCH)
223			.collect::<Result<Vec<_>>>()?;
224		let reached_end = batch.len() < EVICT_BATCH;
225		let last_key = batch.last().map(|(key, _, _)| key.clone());
226
227		for (storage_key, version, _row) in batch {
228			if version > cutoff_version {
229				continue;
230			}
231
232			let bytes = storage_key.as_ref();
233			if bytes.is_empty() || bytes[0] != TIMESTAMP_PREFIX {
234				continue;
235			}
236			let composite_key = EncodedKey::new(bytes[1..].to_vec());
237			self.forget_mapping(txn, &composite_key)?;
238		}
239
240		*self.evict_cursor.borrow_mut() = if reached_end {
241			None
242		} else {
243			last_key
244		};
245		Ok(None)
246	}
247}
248
249impl AppendOperator {
250	#[inline]
251	fn translate_create_row_numbers(
252		&self,
253		txn: &mut FlowTransaction,
254		parent_index: usize,
255		source: &Columns,
256	) -> Result<Vec<RowNumber>> {
257		let row_count = source.row_count();
258		let mut output_row_numbers = Vec::with_capacity(row_count);
259		for row_idx in 0..row_count {
260			let source_row_number = source.row_numbers[row_idx];
261			let composite_key = Self::make_composite_key(parent_index as u8, source_row_number);
262			let (output_row_number, _) =
263				self.row_number_provider.get_or_create_row_number(txn, &composite_key)?;
264			self.touch(txn, &composite_key)?;
265			output_row_numbers.push(output_row_number);
266		}
267		Ok(output_row_numbers)
268	}
269
270	#[inline]
271	fn lookup_row_numbers(
272		&self,
273		txn: &mut FlowTransaction,
274		parent_index: usize,
275		source: &Columns,
276	) -> Result<Option<(Vec<RowNumber>, Vec<EncodedKey>)>> {
277		let row_count = source.row_count();
278		let mut output_row_numbers = Vec::with_capacity(row_count);
279		let mut composite_keys = Vec::with_capacity(row_count);
280		for row_idx in 0..row_count {
281			let source_row_number = source.row_numbers[row_idx];
282			let composite_key = Self::make_composite_key(parent_index as u8, source_row_number);
283			let Some(row_number) = self.row_number_provider.get_row_number(txn, &composite_key)? else {
284				return Ok(None);
285			};
286			output_row_numbers.push(row_number);
287			composite_keys.push(composite_key);
288		}
289		Ok(Some((output_row_numbers, composite_keys)))
290	}
291
292	#[inline]
293	fn translate_append_insert(
294		&self,
295		txn: &mut FlowTransaction,
296		parent_index: usize,
297		post: Columns,
298	) -> Result<Option<Diff>> {
299		if post.row_count() == 0 {
300			return Ok(None);
301		}
302		let output_row_numbers = self.translate_create_row_numbers(txn, parent_index, &post)?;
303		let output = post.with_row_numbers(output_row_numbers);
304		Ok(Some(Diff::insert(output)))
305	}
306
307	#[inline]
308	fn translate_append_update(
309		&self,
310		txn: &mut FlowTransaction,
311		parent_index: usize,
312		pre: Columns,
313		post: Columns,
314	) -> Result<Option<Diff>> {
315		if post.row_count() == 0 {
316			return Ok(None);
317		}
318		let Some((output_row_numbers, composite_keys)) = self.lookup_row_numbers(txn, parent_index, &pre)?
319		else {
320			return Ok(None);
321		};
322		for composite_key in &composite_keys {
323			self.touch(txn, composite_key)?;
324		}
325		let pre_output = pre.with_row_numbers(output_row_numbers.clone());
326		let post_output = post.with_row_numbers(output_row_numbers);
327		Ok(Some(Diff::update(pre_output, post_output)))
328	}
329
330	#[inline]
331	fn translate_append_remove(
332		&self,
333		txn: &mut FlowTransaction,
334		parent_index: usize,
335		pre: Columns,
336	) -> Result<Option<Diff>> {
337		if pre.row_count() == 0 {
338			return Ok(None);
339		}
340		let Some((output_row_numbers, composite_keys)) = self.lookup_row_numbers(txn, parent_index, &pre)?
341		else {
342			return Ok(None);
343		};
344		for composite_key in &composite_keys {
345			self.forget_mapping(txn, composite_key)?;
346		}
347		let output = pre.with_row_numbers(output_row_numbers);
348		Ok(Some(Diff::remove(output)))
349	}
350}
351
352#[cfg(test)]
353mod tests {
354	use reifydb_catalog::catalog::Catalog;
355	use reifydb_core::common::CommitVersion;
356	use reifydb_engine::test_harness::TestEngine;
357	use reifydb_runtime::context::clock::Clock;
358	use reifydb_sdk::operator::Tick;
359	use reifydb_transaction::interceptor::interceptors::Interceptors;
360	use reifydb_value::value::{datetime::DateTime, identity::IdentityId};
361
362	use super::*;
363	use crate::operator::stateful::utils::internal_state_get;
364
365	fn make_tick(clock: &Clock) -> Tick {
366		Tick {
367			now: DateTime::from_nanos(clock.now_nanos()),
368		}
369	}
370
371	fn composite(parent: u8, source_row: u64) -> EncodedKey {
372		AppendOperator::make_composite_key(parent, RowNumber(source_row))
373	}
374
375	#[test]
376	fn translate_create_assigns_and_persists_mapping() {
377		let engine = TestEngine::new();
378		let admin = engine.begin_admin(IdentityId::system()).unwrap();
379		let mut txn = FlowTransaction::deferred(
380			&admin,
381			CommitVersion(1),
382			Catalog::testing(),
383			Interceptors::new(),
384			engine.clock().clone(),
385		);
386		let op = AppendOperator::new_for_state_tests(FlowNodeId(1), None);
387
388		let key = composite(0, 42);
389		assert_eq!(op.row_number_provider.get_row_number(&mut txn, &key).unwrap(), None);
390
391		let (assigned, was_new) = op.row_number_provider.get_or_create_row_number(&mut txn, &key).unwrap();
392		assert!(was_new);
393		assert_eq!(op.row_number_provider.get_row_number(&mut txn, &key).unwrap(), Some(assigned));
394	}
395
396	#[test]
397	fn forget_mapping_removes_forward_and_touch_entries() {
398		let engine = TestEngine::new();
399		let admin = engine.begin_admin(IdentityId::system()).unwrap();
400		let mut txn = FlowTransaction::deferred(
401			&admin,
402			CommitVersion(1),
403			Catalog::testing(),
404			Interceptors::new(),
405			engine.clock().clone(),
406		);
407		let op = AppendOperator::new_for_state_tests(FlowNodeId(2), Some(1_000));
408
409		let key = composite(1, 7);
410		let (_assigned, _) = op.row_number_provider.get_or_create_row_number(&mut txn, &key).unwrap();
411		op.touch(&mut txn, &key).unwrap();
412
413		assert!(op.row_number_provider.get_row_number(&mut txn, &key).unwrap().is_some());
414		let ts_key = AppendOperator::make_timestamp_key(&key);
415		assert!(internal_state_get(op.node, &mut txn, &ts_key).unwrap().is_some());
416
417		op.forget_mapping(&mut txn, &key).unwrap();
418
419		assert!(op.row_number_provider.get_row_number(&mut txn, &key).unwrap().is_none());
420		assert!(internal_state_get(op.node, &mut txn, &ts_key).unwrap().is_none());
421	}
422
423	#[test]
424	fn touch_is_noop_when_ttl_disabled() {
425		// without ttl we must not waste storage on touch entries, since they would never be consulted
426		let engine = TestEngine::new();
427		let admin = engine.begin_admin(IdentityId::system()).unwrap();
428		let mut txn = FlowTransaction::deferred(
429			&admin,
430			CommitVersion(1),
431			Catalog::testing(),
432			Interceptors::new(),
433			engine.clock().clone(),
434		);
435		let op = AppendOperator::new_for_state_tests(FlowNodeId(3), None);
436
437		let key = composite(0, 5);
438		op.touch(&mut txn, &key).unwrap();
439
440		let ts_key = AppendOperator::make_timestamp_key(&key);
441		assert!(
442			internal_state_get(op.node, &mut txn, &ts_key).unwrap().is_none(),
443			"touch must not be written when ttl is disabled"
444		);
445	}
446
447	#[test]
448	fn touch_writes_touch_key_when_ttl_enabled() {
449		// the touch key carries no header timestamp now - its commit version is the last-touch marker.
450		let engine = TestEngine::new();
451		let admin = engine.begin_admin(IdentityId::system()).unwrap();
452		let mut txn = FlowTransaction::deferred(
453			&admin,
454			CommitVersion(1),
455			Catalog::testing(),
456			Interceptors::new(),
457			engine.clock().clone(),
458		);
459		let op = AppendOperator::new_for_state_tests(FlowNodeId(4), Some(60_000_000_000));
460
461		let key = composite(0, 1);
462		op.touch(&mut txn, &key).unwrap();
463
464		let ts_key = AppendOperator::make_timestamp_key(&key);
465		assert!(
466			internal_state_get(op.node, &mut txn, &ts_key).unwrap().is_some(),
467			"touch must write the touch key so its version marks the last access"
468		);
469	}
470
471	#[test]
472	fn tick_evicts_mappings_at_or_below_cutoff_version() {
473		let engine = TestEngine::new();
474		let mock_clock = engine.mock_clock();
475		let admin = engine.begin_admin(IdentityId::system()).unwrap();
476		let mut txn = FlowTransaction::deferred(
477			&admin,
478			CommitVersion(1),
479			Catalog::testing(),
480			Interceptors::new(),
481			engine.clock().clone(),
482		);
483		let ttl_nanos = 50_000_000; // 50ms
484		let op = AppendOperator::new_for_state_tests(FlowNodeId(5), Some(ttl_nanos));
485		// Seed the epoch so any cutoff time maps to commit version 1 - the version every write in
486		// this deferred transaction carries.
487		op.version_epoch.record(0, 1);
488
489		let key = composite(0, 100);
490		op.row_number_provider.get_or_create_row_number(&mut txn, &key).unwrap();
491		op.touch(&mut txn, &key).unwrap();
492		assert!(op.row_number_provider.get_row_number(&mut txn, &key).unwrap().is_some());
493
494		// Advance past the TTL: cutoff = floor_version_at(now - ttl) = 1, at/above the entry's version.
495		mock_clock.advance_millis(100);
496		let result = op.tick(&mut txn, make_tick(&engine.clock())).unwrap();
497		assert!(result.is_none(), "append tick never produces a downstream change");
498
499		assert!(
500			op.row_number_provider.get_row_number(&mut txn, &key).unwrap().is_none(),
501			"a mapping whose touch version is at or below the cutoff must be evicted"
502		);
503	}
504
505	#[test]
506	fn tick_is_conservative_when_epoch_has_no_sample() {
507		let engine = TestEngine::new();
508		let mock_clock = engine.mock_clock();
509		let admin = engine.begin_admin(IdentityId::system()).unwrap();
510		let mut txn = FlowTransaction::deferred(
511			&admin,
512			CommitVersion(1),
513			Catalog::testing(),
514			Interceptors::new(),
515			engine.clock().clone(),
516		);
517		let ttl_nanos = 50_000_000; // 50ms
518		let op = AppendOperator::new_for_state_tests(FlowNodeId(6), Some(ttl_nanos));
519
520		let key = composite(0, 1);
521		op.row_number_provider.get_or_create_row_number(&mut txn, &key).unwrap();
522		op.touch(&mut txn, &key).unwrap();
523
524		// No epoch sample: floor_version_at returns None, so nothing may be evicted (cold-start
525		// conservative contract - never delete when a version cannot be dated).
526		mock_clock.advance_millis(100);
527		op.tick(&mut txn, make_tick(&engine.clock())).unwrap();
528		assert!(
529			op.row_number_provider.get_row_number(&mut txn, &key).unwrap().is_some(),
530			"with no epoch sample the cutoff is None and nothing may be evicted"
531		);
532	}
533
534	#[test]
535	fn tick_is_noop_when_ttl_disabled() {
536		let engine = TestEngine::new();
537		let admin = engine.begin_admin(IdentityId::system()).unwrap();
538		let mut txn = FlowTransaction::deferred(
539			&admin,
540			CommitVersion(1),
541			Catalog::testing(),
542			Interceptors::new(),
543			engine.clock().clone(),
544		);
545		let op = AppendOperator::new_for_state_tests(FlowNodeId(7), None);
546
547		let key = composite(0, 1);
548		op.row_number_provider.get_or_create_row_number(&mut txn, &key).unwrap();
549
550		let result = op.tick(&mut txn, make_tick(&engine.clock())).unwrap();
551		assert!(result.is_none());
552		assert!(op.row_number_provider.get_row_number(&mut txn, &key).unwrap().is_some());
553	}
554
555	#[test]
556	fn capabilities_always_include_tick() {
557		// Mirrors join/distinct: the operator always declares the Tick capability so the
558		// engine can route per-flow ticks (set via `with { tick: ... }` on the view) here
559		// even when TTL is disabled. Tick is a no-op in that case, but the capability is
560		// required to avoid the engine's enforce_tick_capability abort.
561		let with_ttl = AppendOperator::new_for_state_tests(FlowNodeId(8), Some(100));
562		assert!(with_ttl.capabilities().contains(&OperatorCapability::Tick));
563		let without_ttl = AppendOperator::new_for_state_tests(FlowNodeId(9), None);
564		assert!(without_ttl.capabilities().contains(&OperatorCapability::Tick));
565	}
566}