1use 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 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 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; let op = AppendOperator::new_for_state_tests(FlowNodeId(5), Some(ttl_nanos));
485 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 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; 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 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 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}