1use std::{cell::RefCell, collections::HashMap, sync::LazyLock};
5
6use postcard::to_extend;
7use reifydb_abi::operator::capabilities::OperatorCapability;
8use reifydb_codec::{
9 encoded::shape::RowShape,
10 key::{encoded::EncodedKey, serializer::KeySerializer},
11};
12use reifydb_core::{
13 common::{CommitVersion, JoinType},
14 interface::{
15 catalog::flow::FlowNodeId,
16 change::{Change, ChangeOrigin, Diff},
17 },
18 value::column::{ColumnWithName, columns::Columns},
19};
20use reifydb_engine::{
21 expression::{
22 compile::{CompiledExpr, compile_expression},
23 context::{CompileContext, EvalContext},
24 },
25 vm::{executor::Executor, stack::SymbolTable},
26};
27use reifydb_routine::routine::registry::Routines;
28use reifydb_rql::expression::Expression;
29use reifydb_runtime::context::RuntimeContext;
30use reifydb_sdk::operator::Tick;
31use reifydb_value::{
32 Result,
33 error::Error,
34 params::Params,
35 util::hash::{Hash128, xxh3_128},
36 value::{
37 Value, datetime::DateTime, duration::Duration, identity::IdentityId, row_number::RowNumber,
38 value_type::ValueType,
39 },
40};
41
42use super::{
43 column::JoinedColumnsBuilder,
44 state::{JoinSide, JoinState},
45 store::Store,
46 strategy::{JoinContext, JoinStrategy, UpdateKeys},
47};
48use crate::{
49 error::{FlowGraphError, FlowStateError},
50 operator::{
51 Operator,
52 stateful::{raw::RawStatefulOperator, row::RowNumberProvider, single::SingleStateful},
53 },
54 transaction::FlowTransaction,
55};
56
57static EMPTY_PARAMS: Params = Params::None;
58static EMPTY_SYMBOL_TABLE: LazyLock<SymbolTable> = LazyLock::new(SymbolTable::new);
59
60pub(crate) const EVICT_BATCH: usize = 4096;
61
62fn group_by_key(keys: &[Option<Hash128>]) -> (Vec<Hash128>, HashMap<Hash128, Vec<usize>>, Vec<usize>) {
63 let mut order: Vec<Hash128> = Vec::new();
64 let mut groups: HashMap<Hash128, Vec<usize>> = HashMap::new();
65 let mut undefined: Vec<usize> = Vec::new();
66 for (row_idx, key) in keys.iter().enumerate() {
67 match key {
68 Some(key_hash) => {
69 groups.entry(*key_hash)
70 .or_insert_with(|| {
71 order.push(*key_hash);
72 Vec::new()
73 })
74 .push(row_idx);
75 }
76 None => undefined.push(row_idx),
77 }
78 }
79 (order, groups, undefined)
80}
81
82#[cfg(test)]
83mod group_by_key_tests {
84 use super::*;
85
86 fn h(v: u128) -> Hash128 {
87 Hash128(v)
88 }
89
90 #[test]
91 fn groups_duplicate_keys_and_preserves_first_occurrence_order() {
92 let keys = vec![Some(h(0xA)), Some(h(0xB)), Some(h(0xA)), Some(h(0xC)), Some(h(0xB))];
96 let (order, groups, undefined) = group_by_key(&keys);
97
98 assert_eq!(order, vec![h(0xA), h(0xB), h(0xC)], "keys must appear in first-occurrence order");
99 assert_eq!(groups[&h(0xA)], vec![0, 2], "indices within a group keep input order");
100 assert_eq!(groups[&h(0xB)], vec![1, 4]);
101 assert_eq!(groups[&h(0xC)], vec![3]);
102 assert!(undefined.is_empty());
103
104 let regrouped: usize = order.iter().map(|k| groups[k].len()).sum();
105 assert_eq!(regrouped, 5, "every defined row is grouped exactly once");
106 }
107
108 #[test]
109 fn routes_none_keys_to_undefined_without_grouping_them() {
110 let keys = vec![None, Some(h(0xA)), None, Some(h(0xA))];
113 let (order, groups, undefined) = group_by_key(&keys);
114
115 assert_eq!(order, vec![h(0xA)]);
116 assert_eq!(groups[&h(0xA)], vec![1, 3]);
117 assert_eq!(undefined, vec![0, 2], "none-key rows are collected in input order, separately");
118 }
119
120 #[test]
121 fn empty_input_yields_empty_partitions() {
122 let (order, groups, undefined) = group_by_key(&[]);
123 assert!(order.is_empty());
124 assert!(groups.is_empty());
125 assert!(undefined.is_empty());
126 }
127}
128
129pub struct JoinSideConfig {
130 pub node: FlowNodeId,
131 pub exprs: Vec<Expression>,
132 pub schema: Columns,
133}
134
135pub struct JoinOperator {
136 node: FlowNodeId,
137 strategy: JoinStrategy,
138 left_node: FlowNodeId,
139 right_node: FlowNodeId,
140 compiled_left_exprs: Vec<CompiledExpr>,
141 compiled_right_exprs: Vec<CompiledExpr>,
142 alias: Option<String>,
143 shape: RowShape,
144 right_schema: Columns,
145 row_number_provider: RowNumberProvider,
146 routines: Routines,
147 runtime_context: RuntimeContext,
148 pub(crate) snapshot: bool,
149 natural: bool,
150 pub(crate) latest: bool,
151 left_ttl: Option<Duration>,
152 right_ttl: Option<Duration>,
153 left_evict_cursor: RefCell<Option<EncodedKey>>,
154 right_evict_cursor: RefCell<Option<EncodedKey>>,
155 rownumber_evict_cursor: RefCell<Option<EncodedKey>>,
156}
157
158impl JoinOperator {
159 #[allow(clippy::too_many_arguments)]
160 pub fn new(
161 left: JoinSideConfig,
162 right: JoinSideConfig,
163 node: FlowNodeId,
164 join_type: JoinType,
165 alias: Option<String>,
166 executor: Executor,
167 snapshot: bool,
168 natural: bool,
169 latest: bool,
170 left_ttl: Option<Duration>,
171 right_ttl: Option<Duration>,
172 ) -> Self {
173 let left_node = left.node;
174 let right_node = right.node;
175 let left_exprs = left.exprs;
176 let right_exprs = right.exprs;
177 let right_schema = right.schema;
178 let strategy = JoinStrategy::from(join_type, latest);
179 let shape = Self::state_shape();
180 let row_number_provider = RowNumberProvider::new(node);
181
182 let compile_ctx = CompileContext {
183 symbols: &EMPTY_SYMBOL_TABLE,
184 };
185
186 let compiled_left_exprs: Vec<CompiledExpr> = left_exprs
187 .iter()
188 .map(|e| compile_expression(&compile_ctx, e))
189 .collect::<Result<Vec<_>>>()
190 .expect("Failed to compile left expressions");
191
192 let compiled_right_exprs: Vec<CompiledExpr> = right_exprs
193 .iter()
194 .map(|e| compile_expression(&compile_ctx, e))
195 .collect::<Result<Vec<_>>>()
196 .expect("Failed to compile right expressions");
197
198 let routines = executor.routines.clone();
199 let runtime_context = executor.runtime_context.clone();
200
201 Self {
202 node,
203 strategy,
204 left_node,
205 right_node,
206 compiled_left_exprs,
207 compiled_right_exprs,
208 alias,
209 shape,
210 right_schema,
211 row_number_provider,
212 routines,
213 runtime_context,
214 snapshot,
215 natural,
216 latest,
217 left_ttl,
218 right_ttl,
219 left_evict_cursor: RefCell::new(None),
220 right_evict_cursor: RefCell::new(None),
221 rownumber_evict_cursor: RefCell::new(None),
222 }
223 }
224
225 fn state_shape() -> RowShape {
226 RowShape::operator_state()
227 }
228
229 #[cfg(test)]
230 #[allow(clippy::too_many_arguments)]
231 pub(crate) fn new_for_state_tests(
232 node: FlowNodeId,
233 left_ttl: Option<Duration>,
234 right_ttl: Option<Duration>,
235 routines: Routines,
236 runtime_context: RuntimeContext,
237 ) -> Self {
238 Self {
239 node,
240 strategy: JoinStrategy::from(JoinType::Inner, false),
241 left_node: FlowNodeId(0),
242 right_node: FlowNodeId(0),
243 compiled_left_exprs: Vec::new(),
244 compiled_right_exprs: Vec::new(),
245 alias: None,
246 shape: Self::state_shape(),
247 right_schema: Columns::empty(),
248 row_number_provider: RowNumberProvider::new(node),
249 routines,
250 runtime_context,
251 snapshot: true,
252 natural: false,
253 latest: false,
254 left_ttl,
255 right_ttl,
256 left_evict_cursor: RefCell::new(None),
257 right_evict_cursor: RefCell::new(None),
258 rownumber_evict_cursor: RefCell::new(None),
259 }
260 }
261
262 fn evict_left(&self, txn: &mut FlowTransaction, now: DateTime) -> Result<()> {
263 let Some(ttl) = self.left_ttl else {
264 return Ok(());
265 };
266 let Some(cutoff_nanos) = now.to_nanos().checked_sub(ttl.get_nanos() as u64) else {
267 return Ok(());
268 };
269 let Some(cutoff_version) =
270 self.runtime_context.version_epoch.floor_version_at(cutoff_nanos).map(CommitVersion)
271 else {
272 return Ok(());
273 };
274 let left = Store::new(self.node, JoinSide::Left);
275 let mut cursor = self.left_evict_cursor.borrow_mut().take();
276 left.evict_expired(txn, cutoff_version, &mut cursor, EVICT_BATCH)?;
277 *self.left_evict_cursor.borrow_mut() = cursor;
278 Ok(())
279 }
280
281 fn evict_right(&self, txn: &mut FlowTransaction, now: DateTime) -> Result<()> {
282 let Some(ttl) = self.right_ttl else {
283 return Ok(());
284 };
285 let Some(cutoff_nanos) = now.to_nanos().checked_sub(ttl.get_nanos() as u64) else {
286 return Ok(());
287 };
288 let Some(cutoff_version) =
289 self.runtime_context.version_epoch.floor_version_at(cutoff_nanos).map(CommitVersion)
290 else {
291 return Ok(());
292 };
293 let right = Store::new(self.node, JoinSide::Right);
294 let mut cursor = self.right_evict_cursor.borrow_mut().take();
295 right.evict_expired(txn, cutoff_version, &mut cursor, EVICT_BATCH)?;
296 *self.right_evict_cursor.borrow_mut() = cursor;
297 Ok(())
298 }
299
300 fn evict_rownumbers(&self, txn: &mut FlowTransaction, now: DateTime) -> Result<()> {
301 let Some(ttl) = self.left_ttl else {
302 return Ok(());
303 };
304 let Some(cutoff_nanos) = now.to_nanos().checked_sub(ttl.get_nanos() as u64) else {
305 return Ok(());
306 };
307 let Some(cutoff_version) =
308 self.runtime_context.version_epoch.floor_version_at(cutoff_nanos).map(CommitVersion)
309 else {
310 return Ok(());
311 };
312 let mut cursor = self.rownumber_evict_cursor.borrow_mut().take();
313 self.row_number_provider.evict_expired(txn, cutoff_version, &mut cursor, EVICT_BATCH)?;
314 *self.rownumber_evict_cursor.borrow_mut() = cursor;
315 Ok(())
316 }
317
318 pub(crate) fn compute_join_keys(
319 &self,
320 columns: &Columns,
321 compiled_exprs: &[CompiledExpr],
322 ) -> Result<Vec<Option<Hash128>>> {
323 let row_count = columns.row_count();
324 if row_count == 0 {
325 return Ok(Vec::new());
326 }
327
328 let session = EvalContext {
329 params: &EMPTY_PARAMS,
330 symbols: &EMPTY_SYMBOL_TABLE,
331 routines: &self.routines,
332 runtime_context: &self.runtime_context,
333 arena: None,
334 identity: IdentityId::root(),
335 is_aggregate_context: false,
336 columns: Columns::empty(),
337 row_count: 1,
338 target: None,
339 take: None,
340 };
341 let exec_ctx = session.with_eval(columns.clone(), row_count);
342
343 let mut expr_columns = Vec::with_capacity(compiled_exprs.len());
344 for compiled_expr in compiled_exprs.iter() {
345 let col: ColumnWithName = if let Some(col_name) = compiled_expr.access_column_name() {
346 columns.column(col_name)
347 .map(|c| ColumnWithName::new(c.name().clone(), c.data().clone()))
348 .unwrap_or_else(|| {
349 ColumnWithName::undefined_typed(col_name, ValueType::Boolean, row_count)
350 })
351 } else {
352 compiled_expr.execute(&exec_ctx)?
353 };
354 expr_columns.push(col);
355 }
356
357 let mut hashes = Vec::with_capacity(row_count);
358 let mut buf: Vec<u8> = Vec::with_capacity(256);
359 for row_idx in 0..row_count {
360 buf.clear();
361 let mut has_undefined = false;
362
363 for col in &expr_columns {
364 let value = col.data().get_value(row_idx);
365
366 if matches!(value, Value::None { .. }) {
367 has_undefined = true;
368 break;
369 }
370
371 buf = to_extend(&value, buf).map_err(|e| {
372 Error::from(FlowStateError::Encode {
373 state: "value for hash",
374 cause: e.to_string(),
375 })
376 })?;
377 }
378
379 if has_undefined {
380 hashes.push(None);
381 } else {
382 hashes.push(Some(xxh3_128(&buf)));
383 }
384 }
385
386 Ok(hashes)
387 }
388
389 pub(crate) fn unmatched_left_columns(
390 &self,
391 txn: &mut FlowTransaction,
392 left: &Columns,
393 left_idx: usize,
394 ) -> Result<Columns> {
395 let left_row_number = left.row_numbers[left_idx];
396
397 let mut serializer = KeySerializer::new();
398 serializer.extend_u8(b'L');
399 serializer.extend_u64(left_row_number.0);
400 let composite_key = serializer.finish();
401
402 let (result_row_number, _is_new) =
403 self.row_number_provider.get_or_create_row_number(txn, &composite_key)?;
404
405 let builder = JoinedColumnsBuilder::new(left, &self.right_schema, &self.alias, self.natural);
406 Ok(builder.unmatched_left(result_row_number, left, left_idx, &self.right_schema))
407 }
408
409 pub(crate) fn unmatched_left_columns_batch(
410 &self,
411 txn: &mut FlowTransaction,
412 left: &Columns,
413 left_indices: &[usize],
414 ) -> Result<Columns> {
415 if left_indices.is_empty() {
416 return Ok(Columns::empty());
417 }
418
419 let composite_keys: Vec<EncodedKey> = left_indices
420 .iter()
421 .map(|&idx| {
422 let left_row_number = left.row_numbers[idx];
423 let mut serializer = KeySerializer::new();
424 serializer.extend_u8(b'L');
425 serializer.extend_u64(left_row_number.0);
426 serializer.finish()
427 })
428 .collect();
429
430 let row_numbers_with_flags =
431 self.row_number_provider.get_or_create_row_numbers(txn, composite_keys.iter())?;
432 let row_numbers: Vec<RowNumber> = row_numbers_with_flags.iter().map(|(rn, _)| *rn).collect();
433
434 let builder = JoinedColumnsBuilder::new(left, &self.right_schema, &self.alias, self.natural);
435 Ok(builder.unmatched_left_batch(&row_numbers, left, left_indices, &self.right_schema))
436 }
437
438 pub(crate) fn cleanup_left_row_joins(&self, txn: &mut FlowTransaction, left_number: u64) -> Result<()> {
439 let mut serializer = KeySerializer::new();
440 serializer.extend_u8(b'L');
441 serializer.extend_u64(left_number);
442 let prefix = serializer.finish();
443
444 self.row_number_provider.remove_by_prefix(txn, &prefix)
445 }
446
447 fn make_composite_key(left_num: RowNumber, right_num: RowNumber) -> EncodedKey {
448 let mut serializer = KeySerializer::new();
449 serializer.extend_u8(b'L');
450 serializer.extend_u64(left_num.0);
451 serializer.extend_u64(right_num.0);
452 serializer.finish()
453 }
454
455 pub(crate) fn join_columns_one_to_many(
456 &self,
457 txn: &mut FlowTransaction,
458 left: &Columns,
459 left_idx: usize,
460 right: &Columns,
461 ) -> Result<Columns> {
462 let right_count = right.row_count();
463 if right_count == 0 {
464 return Ok(Columns::empty());
465 }
466
467 let left_row_number = left.row_numbers[left_idx];
468
469 let composite_keys: Vec<EncodedKey> = (0..right_count)
470 .map(|right_idx| {
471 let right_row_number = right.row_numbers[right_idx];
472 Self::make_composite_key(left_row_number, right_row_number)
473 })
474 .collect();
475
476 let row_numbers_with_flags =
477 self.row_number_provider.get_or_create_row_numbers(txn, composite_keys.iter())?;
478 let row_numbers: Vec<RowNumber> = row_numbers_with_flags.iter().map(|(rn, _)| *rn).collect();
479
480 let builder = JoinedColumnsBuilder::new(left, right, &self.alias, self.natural);
481 Ok(builder.join_one_to_many(&row_numbers, left, left_idx, right))
482 }
483
484 pub(crate) fn join_columns_many_to_one(
485 &self,
486 txn: &mut FlowTransaction,
487 left: &Columns,
488 right: &Columns,
489 right_idx: usize,
490 ) -> Result<Columns> {
491 let left_count = left.row_count();
492 if left_count == 0 {
493 return Ok(Columns::empty());
494 }
495
496 let right_row_number = right.row_numbers[right_idx];
497
498 let composite_keys: Vec<EncodedKey> = (0..left_count)
499 .map(|left_idx| {
500 let left_row_number = left.row_numbers[left_idx];
501 Self::make_composite_key(left_row_number, right_row_number)
502 })
503 .collect();
504
505 let row_numbers_with_flags =
506 self.row_number_provider.get_or_create_row_numbers(txn, composite_keys.iter())?;
507 let row_numbers: Vec<RowNumber> = row_numbers_with_flags.iter().map(|(rn, _)| *rn).collect();
508
509 let builder = JoinedColumnsBuilder::new(left, right, &self.alias, self.natural);
510 Ok(builder.join_many_to_one(&row_numbers, left, right, right_idx))
511 }
512
513 pub(crate) fn join_columns_cartesian(
514 &self,
515 txn: &mut FlowTransaction,
516 left: &Columns,
517 left_indices: &[usize],
518 right: &Columns,
519 right_indices: &[usize],
520 ) -> Result<Columns> {
521 let left_count = left_indices.len();
522 let right_count = right_indices.len();
523 if left_count == 0 || right_count == 0 {
524 return Ok(Columns::empty());
525 }
526
527 let total_results = left_count * right_count;
528 let mut composite_keys = Vec::with_capacity(total_results);
529
530 for &left_idx in left_indices {
531 let left_row_number = left.row_numbers[left_idx];
532 for &right_idx in right_indices {
533 let right_row_number = right.row_numbers[right_idx];
534 composite_keys.push(Self::make_composite_key(left_row_number, right_row_number));
535 }
536 }
537
538 let row_numbers_with_flags =
539 self.row_number_provider.get_or_create_row_numbers(txn, composite_keys.iter())?;
540 let row_numbers: Vec<RowNumber> = row_numbers_with_flags.iter().map(|(rn, _)| *rn).collect();
541
542 let builder = JoinedColumnsBuilder::new(left, right, &self.alias, self.natural);
543 Ok(builder.join_cartesian(&row_numbers, left, left_indices, right, right_indices))
544 }
545
546 pub(crate) fn join_left_with_slot(&self, left: &Columns, left_indices: &[usize], slot: &Columns) -> Columns {
547 let row_numbers: Vec<RowNumber> = left_indices.iter().map(|&idx| left.row_numbers[idx]).collect();
548 let builder = JoinedColumnsBuilder::new(left, slot, &self.alias, self.natural);
549 builder.join_cartesian(&row_numbers, left, left_indices, slot, &[0])
550 }
551
552 pub(crate) fn unmatched_left_latest(&self, left: &Columns, left_indices: &[usize]) -> Columns {
553 let row_numbers: Vec<RowNumber> = left_indices.iter().map(|&idx| left.row_numbers[idx]).collect();
554 let builder = JoinedColumnsBuilder::new(left, &self.right_schema, &self.alias, self.natural);
555 builder.unmatched_left_batch(&row_numbers, left, left_indices, &self.right_schema)
556 }
557
558 fn determine_side_from_origin(&self, origin: &ChangeOrigin) -> Option<JoinSide> {
559 match origin {
560 ChangeOrigin::Flow(from_node) => {
561 if *from_node == self.left_node {
562 Some(JoinSide::Left)
563 } else if *from_node == self.right_node {
564 Some(JoinSide::Right)
565 } else {
566 None
567 }
568 }
569 _ => None,
570 }
571 }
572}
573
574impl RawStatefulOperator for JoinOperator {}
575
576impl SingleStateful for JoinOperator {
577 fn layout(&self) -> RowShape {
578 self.shape.clone()
579 }
580}
581
582impl Operator for JoinOperator {
583 fn id(&self) -> FlowNodeId {
584 self.node
585 }
586
587 fn capabilities(&self) -> &[OperatorCapability] {
588 OperatorCapability::STANDARD_WITH_TICK
589 }
590
591 fn ticks(&self) -> Option<Duration> {
592 if self.left_ttl.is_some() || self.right_ttl.is_some() {
593 Some(Duration::from_seconds(1).unwrap())
594 } else {
595 None
596 }
597 }
598
599 fn tick(&self, txn: &mut FlowTransaction, tick: Tick) -> Result<Option<Change>> {
600 self.evict_left(txn, tick.now)?;
601 if !self.latest {
602 self.evict_right(txn, tick.now)?;
603 self.evict_rownumbers(txn, tick.now)?;
604 }
605 Ok(None)
606 }
607
608 fn apply(&self, txn: &mut FlowTransaction, change: Change) -> Result<Change> {
609 if let ChangeOrigin::Flow(from_node) = &change.origin
610 && *from_node == self.node
611 {
612 return Ok(Change::from_flow(self.node, change.version, Vec::new(), DateTime::default()));
613 }
614
615 if self.natural && self.compiled_left_exprs.is_empty() {
616 return Ok(Change::from_flow(self.node, change.version, Vec::new(), change.changed_at));
617 }
618
619 let mut state = JoinState::new(self.node);
620 let mut result = Vec::with_capacity(change.diffs.len() * 2);
621
622 let version = change.version;
623 let parent_origin = change.origin.clone();
624 for diff in change.diffs {
625 let diff_origin = diff.origin().cloned().unwrap_or_else(|| parent_origin.clone());
626 let side = self.determine_side_from_origin(&diff_origin).ok_or_else(|| {
627 Error::from(FlowGraphError::UnknownDiffOrigin {
628 operator: "Join",
629 origin: None,
630 })
631 })?;
632 let compiled_exprs = match side {
633 JoinSide::Left => &self.compiled_left_exprs,
634 JoinSide::Right => &self.compiled_right_exprs,
635 };
636 match diff {
637 Diff::Insert {
638 post,
639 ..
640 } => self.apply_join_insert(txn, &post, compiled_exprs, side, &mut state, &mut result)?,
641 Diff::Remove {
642 pre,
643 ..
644 } => self.apply_join_remove(txn, &pre, compiled_exprs, side, &mut state, &mut result)?,
645 Diff::Update {
646 pre,
647 post,
648 ..
649 } => self.apply_join_update(
650 txn,
651 &pre,
652 &post,
653 compiled_exprs,
654 side,
655 &mut state,
656 &mut result,
657 )?,
658 }
659 }
660
661 Ok(Change::from_flow(self.node, version, result, change.changed_at))
662 }
663}
664
665impl JoinOperator {
666 #[inline]
667 #[allow(clippy::too_many_arguments)]
668 fn apply_join_insert(
669 &self,
670 txn: &mut FlowTransaction,
671 post: &Columns,
672 compiled_exprs: &[CompiledExpr],
673 side: JoinSide,
674 state: &mut JoinState,
675 result: &mut Vec<Diff>,
676 ) -> Result<()> {
677 let keys = self.compute_join_keys(post, compiled_exprs)?;
678
679 if !self.latest {
680 for (row_idx, key) in keys.iter().enumerate() {
681 let mut ctx = JoinContext {
682 side,
683 state,
684 operator: self,
685 };
686 let diffs = match key {
687 Some(key_hash) => self.strategy.handle_insert(
688 txn,
689 post,
690 &[row_idx],
691 key_hash,
692 &mut ctx,
693 )?,
694 None => self.strategy.handle_insert_undefined(txn, post, row_idx, &mut ctx)?,
695 };
696 result.extend(diffs);
697 }
698 return Ok(());
699 }
700
701 let (order, groups, undefined) = group_by_key(&keys);
702
703 for key_hash in &order {
704 let indices = &groups[key_hash];
705 let mut ctx = JoinContext {
706 side,
707 state,
708 operator: self,
709 };
710 result.extend(self.strategy.handle_insert(txn, post, indices, key_hash, &mut ctx)?);
711 }
712
713 for row_idx in undefined {
714 let mut ctx = JoinContext {
715 side,
716 state,
717 operator: self,
718 };
719 result.extend(self.strategy.handle_insert_undefined(txn, post, row_idx, &mut ctx)?);
720 }
721
722 Ok(())
723 }
724
725 #[inline]
726 #[allow(clippy::too_many_arguments)]
727 fn apply_join_remove(
728 &self,
729 txn: &mut FlowTransaction,
730 pre: &Columns,
731 compiled_exprs: &[CompiledExpr],
732 side: JoinSide,
733 state: &mut JoinState,
734 result: &mut Vec<Diff>,
735 ) -> Result<()> {
736 let keys = self.compute_join_keys(pre, compiled_exprs)?;
737
738 if !self.latest {
739 for (row_idx, key) in keys.iter().enumerate() {
740 let mut ctx = JoinContext {
741 side,
742 state,
743 operator: self,
744 };
745 let diffs = match key {
746 Some(key_hash) => {
747 self.strategy.handle_remove(txn, pre, &[row_idx], key_hash, &mut ctx)?
748 }
749 None => self.strategy.handle_remove_undefined(txn, pre, row_idx, &mut ctx)?,
750 };
751 result.extend(diffs);
752 }
753 return Ok(());
754 }
755
756 let (order, groups, undefined) = group_by_key(&keys);
757
758 for key_hash in &order {
759 let indices = &groups[key_hash];
760 let mut ctx = JoinContext {
761 side,
762 state,
763 operator: self,
764 };
765 result.extend(self.strategy.handle_remove(txn, pre, indices, key_hash, &mut ctx)?);
766 }
767
768 for row_idx in undefined {
769 let mut ctx = JoinContext {
770 side,
771 state,
772 operator: self,
773 };
774 result.extend(self.strategy.handle_remove_undefined(txn, pre, row_idx, &mut ctx)?);
775 }
776
777 Ok(())
778 }
779
780 #[inline]
781 #[allow(clippy::too_many_arguments)]
782 fn apply_join_update(
783 &self,
784 txn: &mut FlowTransaction,
785 pre: &Columns,
786 post: &Columns,
787 compiled_exprs: &[CompiledExpr],
788 side: JoinSide,
789 state: &mut JoinState,
790 result: &mut Vec<Diff>,
791 ) -> Result<()> {
792 let pre_keys = self.compute_join_keys(pre, compiled_exprs)?;
793 let post_keys = self.compute_join_keys(post, compiled_exprs)?;
794 let row_count = post.row_count();
795
796 for row_idx in 0..row_count {
797 let mut ctx = JoinContext {
798 side,
799 state,
800 operator: self,
801 };
802 let diffs = match (pre_keys[row_idx], post_keys[row_idx]) {
803 (Some(pre_key), Some(post_key)) => {
804 let keys = UpdateKeys {
805 pre: &pre_key,
806 post: &post_key,
807 };
808 self.strategy.handle_update(txn, pre, post, &[row_idx], keys, &mut ctx)?
809 }
810 _ => self.strategy.handle_update_undefined(txn, pre, post, row_idx, &mut ctx)?,
811 };
812 result.extend(diffs);
813 }
814
815 Ok(())
816 }
817}
818
819#[cfg(test)]
820mod tick_tests {
821 use reifydb_catalog::catalog::Catalog;
822 use reifydb_codec::encoded::row::EncodedRow;
823 use reifydb_core::common::CommitVersion;
824 use reifydb_engine::test_harness::TestEngine;
825 use reifydb_transaction::interceptor::interceptors::Interceptors;
826 use reifydb_value::value::blob::Blob;
827
828 use super::*;
829
830 fn ttl(millis: i64) -> Duration {
831 Duration::from_milliseconds_const(millis)
832 }
833
834 fn make_tick(engine: &TestEngine) -> Tick {
835 Tick {
836 now: DateTime::from_nanos(engine.clock().now_nanos()),
837 }
838 }
839
840 fn make_op(
841 node: u64,
842 left_ttl: Option<Duration>,
843 right_ttl: Option<Duration>,
844 engine: &TestEngine,
845 ) -> JoinOperator {
846 let routines = engine.executor().routines.clone();
847 let rc = RuntimeContext::with_clock(engine.clock().clone());
848 rc.version_epoch.record(0, 1);
853 JoinOperator::new_for_state_tests(FlowNodeId(node), left_ttl, right_ttl, routines, rc)
854 }
855
856 fn op_row(payload: u8) -> EncodedRow {
857 let shape = RowShape::operator_state();
858 let mut r = shape.allocate();
859 shape.set_blob(&mut r, 0, &Blob::from(vec![payload]));
860 r
861 }
862
863 #[test]
864 fn tick_evicts_rownumbers_past_ttl() {
865 let engine = TestEngine::new();
870 let mock_clock = engine.mock_clock();
871 let op = make_op(30, Some(ttl(50)), None, &engine);
872 let admin = engine.begin_admin(IdentityId::system()).unwrap();
873 let mut txn = FlowTransaction::deferred(
874 &admin,
875 CommitVersion(1),
876 Catalog::testing(),
877 Interceptors::new(),
878 engine.clock().clone(),
879 );
880
881 let old = JoinOperator::make_composite_key(RowNumber(1), RowNumber(1));
882 op.row_number_provider.get_or_create_row_number(&mut txn, &old).unwrap();
883
884 mock_clock.advance_millis(40);
885 let young = JoinOperator::make_composite_key(RowNumber(2), RowNumber(1));
886 op.row_number_provider.get_or_create_row_number(&mut txn, &young).unwrap();
887
888 mock_clock.advance_millis(20);
889 let emitted = op.tick(&mut txn, make_tick(&engine)).unwrap();
890 assert!(emitted.is_none(), "join tick must be silent (no downstream change)");
891
892 assert!(
893 op.row_number_provider.get_row_number(&mut txn, &old).unwrap().is_none(),
894 "a mapping whose touch version is at or below the cutoff must be evicted"
895 );
896 assert!(
897 op.row_number_provider.get_row_number(&mut txn, &young).unwrap().is_none(),
898 "every mapping at or below the cutoff version is evicted (cross-version selectivity is integration-tested)"
899 );
900 }
901
902 #[test]
903 fn tick_evicts_left_store_past_ttl() {
904 let engine = TestEngine::new();
906 let mock_clock = engine.mock_clock();
907 let op = make_op(30, Some(ttl(50)), None, &engine);
908 let admin = engine.begin_admin(IdentityId::system()).unwrap();
909 let mut txn = FlowTransaction::deferred(
910 &admin,
911 CommitVersion(1),
912 Catalog::testing(),
913 Interceptors::new(),
914 engine.clock().clone(),
915 );
916
917 let left = Store::new(FlowNodeId(30), JoinSide::Left);
918 let hash = Hash128(0xABC);
919 left.put_row(&mut txn, &hash, RowNumber(1), &op_row(0x10)).unwrap();
920 mock_clock.advance_millis(40);
921 left.put_row(&mut txn, &hash, RowNumber(2), &op_row(0x20)).unwrap();
922
923 mock_clock.advance_millis(20);
924 op.tick(&mut txn, make_tick(&engine)).unwrap();
925
926 let remaining = left.rows_for_key(&mut txn, &hash).unwrap();
927 assert!(remaining.is_empty(), "left-store rows at or below the cutoff version are evicted");
928 }
929
930 #[test]
931 fn tick_evicts_right_store_past_ttl() {
932 let engine = TestEngine::new();
936 let mock_clock = engine.mock_clock();
937 let op = make_op(30, None, Some(ttl(50)), &engine);
938 let admin = engine.begin_admin(IdentityId::system()).unwrap();
939 let mut txn = FlowTransaction::deferred(
940 &admin,
941 CommitVersion(1),
942 Catalog::testing(),
943 Interceptors::new(),
944 engine.clock().clone(),
945 );
946
947 let right = Store::new(FlowNodeId(30), JoinSide::Right);
948 let hash = Hash128(0xABC);
949 right.put_row(&mut txn, &hash, RowNumber(1), &op_row(0x10)).unwrap();
950 mock_clock.advance_millis(40);
951 right.put_row(&mut txn, &hash, RowNumber(2), &op_row(0x20)).unwrap();
952
953 mock_clock.advance_millis(20);
954 op.tick(&mut txn, make_tick(&engine)).unwrap();
955
956 let remaining = right.rows_for_key(&mut txn, &hash).unwrap();
957 assert!(remaining.is_empty(), "right-store rows at or below the cutoff version are evicted");
958 }
959
960 #[test]
961 fn tick_is_noop_when_no_ttl_set() {
962 let engine = TestEngine::new();
965 let mock_clock = engine.mock_clock();
966 let op = make_op(30, None, None, &engine);
967 let admin = engine.begin_admin(IdentityId::system()).unwrap();
968 let mut txn = FlowTransaction::deferred(
969 &admin,
970 CommitVersion(1),
971 Catalog::testing(),
972 Interceptors::new(),
973 engine.clock().clone(),
974 );
975
976 let key = JoinOperator::make_composite_key(RowNumber(1), RowNumber(1));
977 op.row_number_provider.get_or_create_row_number(&mut txn, &key).unwrap();
978
979 mock_clock.advance_millis(10_000);
980 let emitted = op.tick(&mut txn, make_tick(&engine)).unwrap();
981 assert!(emitted.is_none());
982 assert!(
983 op.row_number_provider.get_row_number(&mut txn, &key).unwrap().is_some(),
984 "with no TTL configured the tick must retain mappings"
985 );
986 }
987
988 #[test]
989 fn tick_preserves_row_number_counter() {
990 let engine = TestEngine::new();
994 let mock_clock = engine.mock_clock();
995 let op = make_op(30, Some(ttl(50)), None, &engine);
996 let admin = engine.begin_admin(IdentityId::system()).unwrap();
997 let mut txn = FlowTransaction::deferred(
998 &admin,
999 CommitVersion(1),
1000 Catalog::testing(),
1001 Interceptors::new(),
1002 engine.clock().clone(),
1003 );
1004
1005 let first = JoinOperator::make_composite_key(RowNumber(1), RowNumber(1));
1006 let (n1, _) = op.row_number_provider.get_or_create_row_number(&mut txn, &first).unwrap();
1007
1008 mock_clock.advance_millis(100);
1009 op.tick(&mut txn, make_tick(&engine)).unwrap();
1010 assert!(op.row_number_provider.get_row_number(&mut txn, &first).unwrap().is_none());
1011
1012 let second = JoinOperator::make_composite_key(RowNumber(7), RowNumber(7));
1013 let (n2, is_new) = op.row_number_provider.get_or_create_row_number(&mut txn, &second).unwrap();
1014 assert!(is_new);
1015 assert!(n2.0 > n1.0, "counter must keep advancing past evicted mappings, not recycle ids");
1016 }
1017
1018 #[test]
1019 fn capabilities_always_include_tick() {
1020 let engine = TestEngine::new();
1023 let with = make_op(1, Some(ttl(100)), None, &engine);
1024 assert!(with.capabilities().contains(&OperatorCapability::Tick));
1025 let without = make_op(2, None, None, &engine);
1026 assert!(without.capabilities().contains(&OperatorCapability::Tick));
1027 }
1028}