Skip to main content

reifydb_sub_flow/operator/join/
operator.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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		// Latest-mode dispatch iterates `order` to issue one strategy call per key. The order must be
93		// first-occurrence so that, combined with latest's left-row-number reuse, output identity is
94		// stable; and every input index must land in exactly one group with none dropped or duplicated.
95		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		// None-key rows take the per-row undefined path (they never probe); they must not create a
111		// group keyed by some sentinel hash, or an undefined row would be mis-joined.
112		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		// Seed the version epoch so eviction cutoffs resolve to commit version 1 - the version every
849		// write in these single-transaction tests carries. Selectivity across versions (old evicted /
850		// young kept) needs multiple commits and is covered by the integration flow tests; the
851		// conservative cold-start (empty epoch -> no eviction) is covered by the append unit tests.
852		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		// A join mints one row-number mapping per (left,right) output pair. If those mappings are
866		// never evicted once the left row ages past the left TTL, the join's internal state grows
867		// without bound (observed: 430M mapping rows / 66GB on a live ingestor). evict_rownumbers
868		// must drop the aged mappings and keep the fresh ones.
869		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		// evict_left must drop left-store rows older than the left TTL.
905		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		// The snapshot right store accumulates one row per churned upstream RowNumber (observed
933		// ~2875 rows per hot mint, since upstream TTL drops never emit a Remove). evict_right must
934		// drop right-store rows past the right TTL so the probe fan-out and storage stay bounded.
935		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		// With neither side's TTL configured the join must not evict anything (mappings retained,
963		// exactly as before this change; the central GC still bounds the data stores).
964		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		// Evicting every mapping must NOT reset the monotonic counter; a fresh mapping after a
991		// full eviction must get a strictly larger number, or a recycled id would corrupt any
992		// downstream consumer that tracks rows by number.
993		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		// The engine calls enforce_tick_capability before tick() and aborts the process if Tick is
1021		// absent; capabilities must include Tick unconditionally, even when no TTL is set.
1022		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}