1use std::{iter::once, ops::Bound};
4
5use reifydb_codec::key::{
6 encoded::{EncodedKey, EncodedKeyRange},
7 serializer::KeySerializer,
8};
9use reifydb_core::{
10 common::CommitVersion,
11 interface::catalog::flow::FlowNodeId,
12 key::{EncodableKey, flow_node_internal_state::FlowNodeInternalStateKey},
13};
14use reifydb_sdk::state::{decode_payload, encode_payload};
15use reifydb_transaction::multi::RangeScope;
16use reifydb_value::{Result, value::row_number::RowNumber};
17
18use crate::{
19 operator::stateful::utils::{
20 internal_state_drop, internal_state_get, internal_state_range_versioned, internal_state_set,
21 },
22 transaction::FlowTransaction,
23};
24
25pub fn allocate_row_numbers(txn: &mut FlowTransaction, node: FlowNodeId, count: u64) -> Result<u64> {
26 let registry = txn.row_allocators();
27 let counter_key = counter_key();
28 let seed = if registry.is_seeded(node) {
29 0
30 } else {
31 match internal_state_get(node, txn, &counter_key)? {
32 Some(row) => decode_payload::<u64>(&row)?,
33 None => 1,
34 }
35 };
36 let start = registry.allocate(node, count, seed);
37 let high_water = registry.high_water(node).expect("node seeded after allocate");
38 let now = txn.clock().now_nanos();
39 internal_state_set(node, txn, &counter_key, encode_payload(&high_water, now)?)?;
40 Ok(start)
41}
42
43fn counter_key() -> EncodedKey {
44 let mut serializer = KeySerializer::new();
45 serializer.extend_u8(FlowNodeInternalStateKey::ROW_NUMBER_COUNTER_TAG);
46 serializer.finish()
47}
48
49pub struct RowNumberProvider {
50 node: FlowNodeId,
51}
52
53impl RowNumberProvider {
54 pub fn new(node: FlowNodeId) -> Self {
55 Self {
56 node,
57 }
58 }
59
60 pub fn get_or_create_row_numbers<'a, I>(
61 &self,
62 txn: &mut FlowTransaction,
63 keys: I,
64 ) -> Result<Vec<(RowNumber, bool)>>
65 where
66 I: IntoIterator<Item = &'a EncodedKey>,
67 {
68 let now = txn.clock().now_nanos();
69 let keys: Vec<&EncodedKey> = keys.into_iter().collect();
70 let mut results: Vec<Option<(RowNumber, bool)>> = (0..keys.len()).map(|_| None).collect();
71 let mut new_positions: Vec<(usize, EncodedKey)> = Vec::new();
72
73 for (i, key) in keys.iter().enumerate() {
74 let map_key = self.make_map_key(key);
75 if let Some(existing_row) = internal_state_get(self.node, txn, &map_key)? {
76 results[i] = Some((RowNumber(decode_payload::<u64>(&existing_row)?), false));
77 } else {
78 new_positions.push((i, map_key));
79 }
80 }
81
82 if !new_positions.is_empty() {
83 let start = self.mint(txn, new_positions.len() as u64)?;
84 for (offset, (i, map_key)) in new_positions.iter().enumerate() {
85 let row_number = RowNumber(start + offset as u64);
86 internal_state_set(self.node, txn, map_key, encode_payload(&row_number.0, now)?)?;
87 results[*i] = Some((row_number, true));
88 }
89 }
90
91 Ok(results.into_iter().map(|r| r.expect("every position filled")).collect())
92 }
93
94 fn mint(&self, txn: &mut FlowTransaction, count: u64) -> Result<u64> {
95 allocate_row_numbers(txn, self.node, count)
96 }
97
98 pub fn get_or_create_row_number(
99 &self,
100 txn: &mut FlowTransaction,
101 key: &EncodedKey,
102 ) -> Result<(RowNumber, bool)> {
103 Ok(self.get_or_create_row_numbers(txn, once(key))?.into_iter().next().unwrap())
104 }
105
106 pub fn get_row_number(&self, txn: &mut FlowTransaction, key: &EncodedKey) -> Result<Option<RowNumber>> {
107 let map_key = self.make_map_key(key);
108 match internal_state_get(self.node, txn, &map_key)? {
109 Some(existing_row) => Ok(Some(RowNumber(decode_payload::<u64>(&existing_row)?))),
110 None => Ok(None),
111 }
112 }
113
114 pub fn remove_for_key(&self, txn: &mut FlowTransaction, key: &EncodedKey) -> Result<bool> {
115 let map_key = self.make_map_key(key);
116 if internal_state_get(self.node, txn, &map_key)?.is_none() {
117 return Ok(false);
118 }
119 internal_state_drop(self.node, txn, &map_key)?;
120 Ok(true)
121 }
122
123 fn make_map_key(&self, key: &EncodedKey) -> EncodedKey {
124 let mut serializer = KeySerializer::new();
125 serializer.extend_u8(FlowNodeInternalStateKey::ROW_NUMBER_MAPPING_TAG);
126 serializer.extend_bytes(key.as_ref());
127 serializer.finish()
128 }
129
130 pub fn remove_by_prefix(&self, txn: &mut FlowTransaction, key_prefix: &[u8]) -> Result<()> {
131 let mut prefix = Vec::new();
132 let mut serializer = KeySerializer::new();
133 serializer.extend_u8(FlowNodeInternalStateKey::ROW_NUMBER_MAPPING_TAG);
134 prefix.extend_from_slice(&serializer.finish());
135 prefix.extend_from_slice(key_prefix);
136
137 let state_prefix = FlowNodeInternalStateKey::new(self.node, prefix.clone());
138 let full_range = EncodedKeyRange::prefix(&state_prefix.encode());
139
140 let keys_to_remove = {
141 let stream = txn.range(full_range, RangeScope::All, 1024);
142 let mut keys = Vec::new();
143 for result in stream {
144 let multi = result?;
145 keys.push(multi.key);
146 }
147 keys
148 };
149
150 for key in keys_to_remove {
151 txn.remove(&key)?;
152 }
153
154 Ok(())
155 }
156
157 pub fn evict_expired(
158 &self,
159 txn: &mut FlowTransaction,
160 cutoff_version: CommitVersion,
161 cursor: &mut Option<EncodedKey>,
162 batch_size: usize,
163 ) -> Result<()> {
164 let prefix = {
165 let mut serializer = KeySerializer::new();
166 serializer.extend_u8(FlowNodeInternalStateKey::ROW_NUMBER_MAPPING_TAG);
167 serializer.finish()
168 };
169 let base = EncodedKeyRange::prefix(prefix.as_ref());
170 let start = match cursor.clone() {
171 Some(c) => Bound::Excluded(c),
172 None => base.start.clone(),
173 };
174 let range = EncodedKeyRange::new(start, base.end.clone());
175 let batch = internal_state_range_versioned(self.node, txn, range)
176 .take(batch_size)
177 .collect::<Result<Vec<_>>>()?;
178 let reached_end = batch.len() < batch_size;
179 let last_key = batch.last().map(|(key, _, _)| key.clone());
180
181 for (key, version, _row) in batch {
182 if version > cutoff_version {
183 continue;
184 }
185 internal_state_drop(self.node, txn, &key)?;
186 }
187
188 *cursor = if reached_end {
189 None
190 } else {
191 last_key
192 };
193 Ok(())
194 }
195}
196
197#[cfg(test)]
198pub mod tests {
199 use reifydb_catalog::catalog::Catalog;
200 use reifydb_core::common::CommitVersion;
201 use reifydb_runtime::context::clock::{Clock, MockClock};
202 use reifydb_transaction::interceptor::interceptors::Interceptors;
203
204 use super::*;
205 use crate::operator::stateful::test_utils::test::*;
206
207 #[test]
208 fn test_first_row_number() {
209 let mut txn = create_test_transaction();
210 let mut txn = FlowTransaction::deferred(
211 &mut txn,
212 CommitVersion(1),
213 Catalog::testing(),
214 Interceptors::new(),
215 Clock::Mock(MockClock::from_millis(1000)),
216 );
217 let provider = RowNumberProvider::new(FlowNodeId(1));
218
219 let key = test_key("first");
220 let (row_num, is_new) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
221
222 assert_eq!(row_num.0, 1);
223 assert!(is_new);
224 }
225
226 #[test]
227 fn test_duplicate_key_same_row_number() {
228 let mut txn = create_test_transaction();
229 let mut txn = FlowTransaction::deferred(
230 &mut txn,
231 CommitVersion(1),
232 Catalog::testing(),
233 Interceptors::new(),
234 Clock::Mock(MockClock::from_millis(1000)),
235 );
236 let provider = RowNumberProvider::new(FlowNodeId(1));
237
238 let key = test_key("duplicate");
239
240 let (row_num1, is_new1) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
242 assert_eq!(row_num1.0, 1);
243 assert!(is_new1);
244
245 let (row_num2, is_new2) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
247 assert_eq!(row_num2.0, 1);
248 assert!(!is_new2);
249
250 assert_eq!(row_num1, row_num2);
252 }
253
254 #[test]
255 fn test_sequential_row_numbers() {
256 let mut txn = create_test_transaction();
257 let mut txn = FlowTransaction::deferred(
258 &mut txn,
259 CommitVersion(1),
260 Catalog::testing(),
261 Interceptors::new(),
262 Clock::Mock(MockClock::from_millis(1000)),
263 );
264 let provider = RowNumberProvider::new(FlowNodeId(1));
265
266 for i in 1..=5 {
268 let key = test_key(&format!("key_{}", i));
269 let (row_num, is_new) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
270
271 assert_eq!(row_num.0, i as u64);
272 assert!(is_new);
273 }
274 }
275
276 #[test]
277 fn test_mixed_new_and_existing() {
278 let mut txn = create_test_transaction();
279 let mut txn = FlowTransaction::deferred(
280 &mut txn,
281 CommitVersion(1),
282 Catalog::testing(),
283 Interceptors::new(),
284 Clock::Mock(MockClock::from_millis(1000)),
285 );
286 let provider = RowNumberProvider::new(FlowNodeId(1));
287
288 let key1 = test_key("mixed_1");
290 let key2 = test_key("mixed_2");
291 let key3 = test_key("mixed_3");
292
293 let (rn1, new1) = provider.get_or_create_row_number(&mut txn, &key1).unwrap();
295 let (rn2, new2) = provider.get_or_create_row_number(&mut txn, &key2).unwrap();
296 let (rn3, new3) = provider.get_or_create_row_number(&mut txn, &key3).unwrap();
297
298 assert_eq!(rn1.0, 1);
299 assert!(new1);
300 assert_eq!(rn2.0, 2);
301 assert!(new2);
302 assert_eq!(rn3.0, 3);
303 assert!(new3);
304
305 let key4 = test_key("mixed_4");
307 let (rn2_again, new2_again) = provider.get_or_create_row_number(&mut txn, &key2).unwrap();
308 let (rn4, new4) = provider.get_or_create_row_number(&mut txn, &key4).unwrap();
309 let (rn1_again, new1_again) = provider.get_or_create_row_number(&mut txn, &key1).unwrap();
310
311 assert_eq!(rn2_again.0, 2);
312 assert!(!new2_again);
313 assert_eq!(rn4.0, 4); assert!(new4);
315 assert_eq!(rn1_again.0, 1);
316 assert!(!new1_again);
317 }
318
319 #[test]
320 fn test_multiple_providers_isolated() {
321 let mut txn = create_test_transaction();
322 let mut txn = FlowTransaction::deferred(
323 &mut txn,
324 CommitVersion(1),
325 Catalog::testing(),
326 Interceptors::new(),
327 Clock::Mock(MockClock::from_millis(1000)),
328 );
329 let provider1 = RowNumberProvider::new(FlowNodeId(1));
330 let provider2 = RowNumberProvider::new(FlowNodeId(2));
331
332 let key = test_key("shared_key");
333
334 let (rn1, _) = provider1.get_or_create_row_number(&mut txn, &key).unwrap();
336 let (rn2, _) = provider2.get_or_create_row_number(&mut txn, &key).unwrap();
337
338 assert_eq!(rn1.0, 1);
339 assert_eq!(rn2.0, 1);
340
341 let key2 = test_key("key2");
343 let (rn1_2, _) = provider1.get_or_create_row_number(&mut txn, &key2).unwrap();
344 assert_eq!(rn1_2.0, 2);
345
346 let (rn2_2, _) = provider2.get_or_create_row_number(&mut txn, &key2).unwrap();
348 assert_eq!(rn2_2.0, 2);
349 }
350
351 #[test]
352 fn test_counter_persistence() {
353 let mut txn = create_test_transaction();
354 let mut txn = FlowTransaction::deferred(
355 &mut txn,
356 CommitVersion(1),
357 Catalog::testing(),
358 Interceptors::new(),
359 Clock::Mock(MockClock::from_millis(1000)),
360 );
361 let provider = RowNumberProvider::new(FlowNodeId(1));
362
363 for i in 1..=3 {
365 let key = test_key(&format!("persist_{}", i));
366 let (rn, _) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
367 assert_eq!(rn.0, i as u64);
368 }
369
370 let new_key = test_key("persist_new");
372 let (rn, is_new) = provider.get_or_create_row_number(&mut txn, &new_key).unwrap();
373
374 assert_eq!(rn.0, 4);
376 assert!(is_new);
377 }
378
379 #[test]
380 fn test_large_row_numbers() {
381 let mut txn = create_test_transaction();
382 let mut txn = FlowTransaction::deferred(
383 &mut txn,
384 CommitVersion(1),
385 Catalog::testing(),
386 Interceptors::new(),
387 Clock::Mock(MockClock::from_millis(1000)),
388 );
389 let provider = RowNumberProvider::new(FlowNodeId(1));
390
391 for i in 1..=1000 {
393 let key = test_key(&format!("large_{}", i));
394 let (rn, is_new) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
395 assert_eq!(rn.0, i as u64);
396 assert!(is_new);
397 }
398
399 let key = test_key("large_1");
401 let (rn, is_new) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
402 assert_eq!(rn.0, 1);
403 assert!(!is_new);
404
405 let key = test_key("large_1001");
407 let (rn, is_new) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
408 assert_eq!(rn.0, 1001);
409 assert!(is_new);
410 }
411
412 #[test]
413 fn test_mixed_existing_and_new_keys() {
414 let mut txn = create_test_transaction();
415 let mut txn = FlowTransaction::deferred(
416 &mut txn,
417 CommitVersion(1),
418 Catalog::testing(),
419 Interceptors::new(),
420 Clock::Mock(MockClock::from_millis(1000)),
421 );
422 let provider = RowNumberProvider::new(FlowNodeId(1));
423
424 let key1 = test_key("key_1");
426 let key2 = test_key("key_2");
427 let key3 = test_key("key_3");
428
429 let (rn1, _) = provider.get_or_create_row_number(&mut txn, &key1).unwrap();
430 assert_eq!(rn1.0, 1);
431
432 let (rn2, _) = provider.get_or_create_row_number(&mut txn, &key2).unwrap();
433 assert_eq!(rn2.0, 2);
434
435 let (rn3, _) = provider.get_or_create_row_number(&mut txn, &key3).unwrap();
436 assert_eq!(rn3.0, 3);
437
438 let key4 = test_key("key_4");
440 let key5 = test_key("key_5");
441
442 let keys = vec![&key2, &key4, &key1, &key5, &key3];
444
445 let results = provider.get_or_create_row_numbers(&mut txn, keys.into_iter()).unwrap();
446
447 assert_eq!(results.len(), 5);
449
450 assert_eq!(results[0].0.0, 2);
452 assert!(!results[0].1);
453
454 assert_eq!(results[1].0.0, 4);
456 assert!(results[1].1);
457
458 assert_eq!(results[2].0.0, 1);
460 assert!(!results[2].1);
461
462 assert_eq!(results[3].0.0, 5);
464 assert!(results[3].1);
465
466 assert_eq!(results[4].0.0, 3);
468 assert!(!results[4].1);
469
470 let key6 = test_key("key_6");
473 let (rn6, is_new6) = provider.get_or_create_row_number(&mut txn, &key6).unwrap();
474 assert_eq!(rn6.0, 6);
475 assert!(is_new6);
476
477 let (check_rn4, is_new4) = provider.get_or_create_row_number(&mut txn, &key4).unwrap();
479 assert_eq!(check_rn4.0, 4);
480 assert!(!is_new4);
481
482 let (check_rn5, is_new5) = provider.get_or_create_row_number(&mut txn, &key5).unwrap();
483 assert_eq!(check_rn5.0, 5);
484 assert!(!is_new5);
485 }
486
487 #[test]
488 fn test_get_row_number_returns_none_for_unknown() {
489 let mut txn = create_test_transaction();
490 let mut txn = FlowTransaction::deferred(
491 &mut txn,
492 CommitVersion(1),
493 Catalog::testing(),
494 Interceptors::new(),
495 Clock::Mock(MockClock::from_millis(1000)),
496 );
497 let provider = RowNumberProvider::new(FlowNodeId(1));
498
499 let key = test_key("never_seen");
500 assert_eq!(provider.get_row_number(&mut txn, &key).unwrap(), None);
501 }
502
503 #[test]
504 fn test_get_row_number_returns_existing_without_creating() {
505 let mut txn = create_test_transaction();
506 let mut txn = FlowTransaction::deferred(
507 &mut txn,
508 CommitVersion(1),
509 Catalog::testing(),
510 Interceptors::new(),
511 Clock::Mock(MockClock::from_millis(1000)),
512 );
513 let provider = RowNumberProvider::new(FlowNodeId(1));
514
515 let key = test_key("lookup_hit");
516 let (created, was_new) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
517 assert!(was_new);
518
519 let looked_up = provider.get_row_number(&mut txn, &key).unwrap();
520 assert_eq!(looked_up, Some(created));
521
522 let another = test_key("another_missing");
523 assert_eq!(provider.get_row_number(&mut txn, &another).unwrap(), None);
524 let (after, was_new_after) = provider.get_or_create_row_number(&mut txn, &another).unwrap();
525 assert!(was_new_after);
526 assert_ne!(after, created);
527 }
528
529 #[test]
530 fn test_remove_for_key_clears_mapping() {
531 let mut txn = create_test_transaction();
532 let mut txn = FlowTransaction::deferred(
533 &mut txn,
534 CommitVersion(1),
535 Catalog::testing(),
536 Interceptors::new(),
537 Clock::Mock(MockClock::from_millis(1000)),
538 );
539 let provider = RowNumberProvider::new(FlowNodeId(1));
540
541 let key = test_key("to_remove");
542 let (_assigned, _) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
543 assert!(provider.get_row_number(&mut txn, &key).unwrap().is_some());
544
545 let removed = provider.remove_for_key(&mut txn, &key).unwrap();
546 assert!(removed);
547
548 assert_eq!(provider.get_row_number(&mut txn, &key).unwrap(), None);
549 }
550
551 #[test]
552 fn test_remove_for_key_is_idempotent() {
553 let mut txn = create_test_transaction();
554 let mut txn = FlowTransaction::deferred(
555 &mut txn,
556 CommitVersion(1),
557 Catalog::testing(),
558 Interceptors::new(),
559 Clock::Mock(MockClock::from_millis(1000)),
560 );
561 let provider = RowNumberProvider::new(FlowNodeId(1));
562
563 let key = test_key("absent");
564 assert!(!provider.remove_for_key(&mut txn, &key).unwrap());
565
566 let (_assigned, _) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
567 assert!(provider.remove_for_key(&mut txn, &key).unwrap());
568 assert!(!provider.remove_for_key(&mut txn, &key).unwrap());
569 }
570
571 #[test]
572 fn test_remove_for_key_then_recreate_assigns_new_number() {
573 let mut txn = create_test_transaction();
574 let mut txn = FlowTransaction::deferred(
575 &mut txn,
576 CommitVersion(1),
577 Catalog::testing(),
578 Interceptors::new(),
579 Clock::Mock(MockClock::from_millis(1000)),
580 );
581 let provider = RowNumberProvider::new(FlowNodeId(1));
582
583 let key = test_key("recycled");
584 let (first, _) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
585 assert!(provider.remove_for_key(&mut txn, &key).unwrap());
586
587 let (second, was_new) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
588 assert!(was_new, "after removal the next mapping should be created fresh");
589 assert_ne!(first, second, "counter must keep advancing, not recycle old row numbers");
590 }
591
592 #[test]
593 fn internal_state_tags_are_pairwise_distinct() {
594 let tags = [
600 FlowNodeInternalStateKey::ROW_NUMBER_COUNTER_TAG,
601 FlowNodeInternalStateKey::ROW_NUMBER_MAPPING_TAG,
602 FlowNodeInternalStateKey::WINDOW_META_TAG,
603 FlowNodeInternalStateKey::GATE_VISIBILITY_TAG,
604 ];
605 for i in 0..tags.len() {
606 for j in (i + 1)..tags.len() {
607 assert_ne!(tags[i], tags[j], "internal-state tag collision at {:#04x}", tags[i]);
608 }
609 }
610 }
611
612 #[test]
613 fn mapping_values_are_postcard_encoded() {
614 let mut txn = create_test_transaction();
618 let mut txn = FlowTransaction::deferred(
619 &mut txn,
620 CommitVersion(1),
621 Catalog::testing(),
622 Interceptors::new(),
623 Clock::Mock(MockClock::from_millis(1000)),
624 );
625 let provider = RowNumberProvider::new(FlowNodeId(1));
626
627 let key = test_key("encoded");
628 let (rn, _) = provider.get_or_create_row_number(&mut txn, &key).unwrap();
629
630 let forward =
631 internal_state_get(FlowNodeId(1), &mut txn, &provider.make_map_key(&key)).unwrap().unwrap();
632 assert_eq!(decode_payload::<u64>(&forward).unwrap(), rn.0);
633 }
634
635 #[test]
636 fn test_counter_survives_full_mapping_eviction() {
637 let mut txn = create_test_transaction();
642 let mut txn = FlowTransaction::deferred(
643 &mut txn,
644 CommitVersion(1),
645 Catalog::testing(),
646 Interceptors::new(),
647 Clock::Mock(MockClock::from_millis(1000)),
648 );
649 let provider = RowNumberProvider::new(FlowNodeId(1));
650
651 let keys = [test_key("a"), test_key("b"), test_key("c")];
652 let mut issued = Vec::new();
653 for key in &keys {
654 let (n, was_new) = provider.get_or_create_row_number(&mut txn, key).unwrap();
655 assert!(was_new);
656 issued.push(n);
657 }
658
659 for key in &keys {
660 assert!(provider.remove_for_key(&mut txn, key).unwrap());
661 }
662
663 let (fresh, was_new) = provider.get_or_create_row_number(&mut txn, &test_key("d")).unwrap();
664 assert!(was_new, "a brand-new key after full eviction must be assigned fresh");
665 for prev in &issued {
666 assert_ne!(&fresh, prev, "row number {:?} was reused after full eviction", prev);
667 }
668 assert!(
669 issued.iter().all(|prev| fresh.0 > prev.0),
670 "counter must keep advancing past every previously issued number, got {:?} after {:?}",
671 fresh,
672 issued
673 );
674 }
675}