Skip to main content

reifydb_value/value/
system_columns.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::fmt::{self, Display, Formatter};
5
6use serde::{Deserialize, Serialize};
7
8use crate::{
9	util::bitvec::BitVec,
10	value::{datetime::DateTime, partition::Partition, row_number::RowNumber},
11};
12
13#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
14pub struct SystemColumns {
15	row_numbers: Vec<RowNumber>,
16	partitions: Vec<Partition>,
17	created_at: Vec<DateTime>,
18	updated_at: Vec<DateTime>,
19	time: Vec<DateTime>,
20}
21
22#[derive(Debug, Clone, Copy, PartialEq)]
23pub struct RowStamps {
24	pub row_number: Option<RowNumber>,
25	pub partition: Option<Partition>,
26	pub created_at: Option<DateTime>,
27	pub updated_at: Option<DateTime>,
28	pub time: Option<DateTime>,
29}
30
31#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub enum SystemColumn {
33	RowNumbers,
34	Partitions,
35	CreatedAt,
36	UpdatedAt,
37	Time,
38}
39
40impl SystemColumn {
41	pub const fn name(self) -> &'static str {
42		match self {
43			SystemColumn::RowNumbers => "#rownum",
44			SystemColumn::Partitions => "#partition",
45			SystemColumn::CreatedAt => "#created_at",
46			SystemColumn::UpdatedAt => "#updated_at",
47			SystemColumn::Time => "#time",
48		}
49	}
50}
51
52impl Display for SystemColumn {
53	fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
54		f.write_str(self.name())
55	}
56}
57
58#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
59pub enum SystemColumnsError {
60	#[error("cannot append rows: {column} is present on one side but not the other")]
61	PresenceMismatch {
62		column: SystemColumn,
63		target_present: bool,
64		source_present: bool,
65	},
66
67	#[error("{column} holds {len} entries but the batch has {row_count} rows")]
68	LengthMismatch {
69		column: SystemColumn,
70		len: usize,
71		row_count: usize,
72	},
73}
74
75#[inline]
76fn gather<T: Copy + PartialEq>(src: &[T], indices: &[usize]) -> Vec<T> {
77	if src.is_empty() {
78		return Vec::new();
79	}
80	indices.iter().map(|&i| src[i]).collect()
81}
82
83#[inline]
84fn retain<T: Copy + PartialEq>(src: &[T], mask: &BitVec) -> Vec<T> {
85	if src.is_empty() {
86		return Vec::new();
87	}
88	src.iter().enumerate().filter(|(i, _)| *i < mask.len() && mask.get(*i)).map(|(_, &v)| v).collect()
89}
90
91#[inline]
92fn head<T: Copy + PartialEq>(src: &[T], n: usize) -> Vec<T> {
93	if src.is_empty() {
94		return Vec::new();
95	}
96	src[..n.min(src.len())].to_vec()
97}
98
99#[inline]
100fn concat<T: Copy + PartialEq>(dst: &mut Vec<T>, src: &Vec<T>) {
101	if src.is_empty() {
102		return;
103	}
104	dst.extend_from_slice(src.as_slice());
105}
106
107impl SystemColumns {
108	pub fn empty() -> Self {
109		Self::default()
110	}
111
112	pub fn new(
113		row_numbers: Vec<RowNumber>,
114		partitions: Vec<Partition>,
115		created_at: Vec<DateTime>,
116		updated_at: Vec<DateTime>,
117		time: Vec<DateTime>,
118	) -> Self {
119		Self {
120			row_numbers,
121			partitions,
122			created_at,
123			updated_at,
124			time,
125		}
126	}
127
128	pub fn from_row_numbers(row_numbers: Vec<RowNumber>) -> Self {
129		let n = row_numbers.len();
130		let now = DateTime::default();
131		Self::new(row_numbers, Vec::new(), vec![now; n], vec![now; n], vec![now; n])
132	}
133
134	pub fn set_row_numbers(&mut self, row_numbers: Vec<RowNumber>) {
135		self.row_numbers = row_numbers;
136	}
137
138	pub fn set_partitions(&mut self, partitions: Vec<Partition>) {
139		self.partitions = partitions;
140	}
141
142	pub fn set_created_at(&mut self, created_at: Vec<DateTime>) {
143		self.created_at = created_at;
144	}
145
146	pub fn set_updated_at(&mut self, updated_at: Vec<DateTime>) {
147		self.updated_at = updated_at;
148	}
149
150	pub fn set_time(&mut self, time: Vec<DateTime>) {
151		self.time = time;
152	}
153}
154
155impl SystemColumns {
156	#[inline]
157	pub fn row_numbers(&self) -> &[RowNumber] {
158		self.row_numbers.as_slice()
159	}
160
161	#[inline]
162	pub fn partitions(&self) -> &[Partition] {
163		self.partitions.as_slice()
164	}
165
166	#[inline]
167	pub fn created_at(&self) -> &[DateTime] {
168		self.created_at.as_slice()
169	}
170
171	#[inline]
172	pub fn updated_at(&self) -> &[DateTime] {
173		self.updated_at.as_slice()
174	}
175
176	#[inline]
177	pub fn time(&self) -> &[DateTime] {
178		self.time.as_slice()
179	}
180
181	pub fn row_count(&self) -> Option<usize> {
182		let Self {
183			row_numbers,
184			partitions,
185			created_at,
186			updated_at,
187			time,
188		} = self;
189		[row_numbers.len(), partitions.len(), created_at.len(), updated_at.len(), time.len()]
190			.into_iter()
191			.find(|&len| len > 0)
192	}
193
194	pub fn is_empty(&self) -> bool {
195		self.row_count().is_none()
196	}
197
198	pub fn heap_size(&self) -> usize {
199		let Self {
200			row_numbers,
201			partitions,
202			created_at,
203			updated_at,
204			time,
205		} = self;
206		row_numbers.len() * size_of::<RowNumber>()
207			+ partitions.len() * size_of::<Partition>()
208			+ created_at.len() * size_of::<DateTime>()
209			+ updated_at.len() * size_of::<DateTime>()
210			+ time.len() * size_of::<DateTime>()
211	}
212}
213
214impl SystemColumns {
215	pub fn permute(&self, indices: &[usize]) -> Self {
216		let Self {
217			row_numbers,
218			partitions,
219			created_at,
220			updated_at,
221			time,
222		} = self;
223		Self {
224			row_numbers: gather(row_numbers, indices),
225			partitions: gather(partitions, indices),
226			created_at: gather(created_at, indices),
227			updated_at: gather(updated_at, indices),
228			time: gather(time, indices),
229		}
230	}
231
232	pub fn permute_in_place(&mut self, indices: &[usize]) {
233		*self = self.permute(indices);
234	}
235
236	pub fn filter(&mut self, mask: &BitVec) {
237		let Self {
238			row_numbers,
239			partitions,
240			created_at,
241			updated_at,
242			time,
243		} = self;
244		*row_numbers = retain(row_numbers, mask);
245		*partitions = retain(partitions, mask);
246		*created_at = retain(created_at, mask);
247		*updated_at = retain(updated_at, mask);
248		*time = retain(time, mask);
249	}
250
251	pub fn take(&mut self, n: usize) {
252		let Self {
253			row_numbers,
254			partitions,
255			created_at,
256			updated_at,
257			time,
258		} = self;
259		*row_numbers = head(row_numbers, n);
260		*partitions = head(partitions, n);
261		*created_at = head(created_at, n);
262		*updated_at = head(updated_at, n);
263		*time = head(time, n);
264	}
265
266	pub fn extend(&mut self, source: &Self) -> Result<(), SystemColumnsError> {
267		self.check_extendable(source)?;
268		let Self {
269			row_numbers,
270			partitions,
271			created_at,
272			updated_at,
273			time,
274		} = self;
275		concat(row_numbers, &source.row_numbers);
276		concat(partitions, &source.partitions);
277		concat(created_at, &source.created_at);
278		concat(updated_at, &source.updated_at);
279		concat(time, &source.time);
280		Ok(())
281	}
282
283	pub fn append_indices(&mut self, source: &Self, indices: &[usize]) {
284		let gathered = source.permute(indices);
285		let Self {
286			row_numbers,
287			partitions,
288			created_at,
289			updated_at,
290			time,
291		} = self;
292		concat(row_numbers, &gathered.row_numbers);
293		concat(partitions, &gathered.partitions);
294		concat(created_at, &gathered.created_at);
295		concat(updated_at, &gathered.updated_at);
296		concat(time, &gathered.time);
297	}
298
299	pub fn push(&mut self, stamps: RowStamps) {
300		let RowStamps {
301			row_number,
302			partition,
303			created_at,
304			updated_at,
305			time,
306		} = stamps;
307		if let Some(row_number) = row_number {
308			self.row_numbers.push(row_number);
309		}
310		if let Some(partition) = partition {
311			self.partitions.push(partition);
312		}
313		if let Some(created_at) = created_at {
314			self.created_at.push(created_at);
315		}
316		if let Some(updated_at) = updated_at {
317			self.updated_at.push(updated_at);
318		}
319		if let Some(time) = time {
320			self.time.push(time);
321		}
322	}
323
324	pub fn clear(&mut self) {
325		let Self {
326			row_numbers,
327			partitions,
328			created_at,
329			updated_at,
330			time,
331		} = self;
332		row_numbers.clear();
333		partitions.clear();
334		created_at.clear();
335		updated_at.clear();
336		time.clear();
337	}
338
339	fn check_extendable(&self, source: &Self) -> Result<(), SystemColumnsError> {
340		let Self {
341			row_numbers,
342			partitions,
343			created_at,
344			updated_at,
345			time,
346		} = self;
347		let pairs = [
348			(SystemColumn::RowNumbers, !row_numbers.is_empty(), !source.row_numbers.is_empty()),
349			(SystemColumn::Partitions, !partitions.is_empty(), !source.partitions.is_empty()),
350			(SystemColumn::CreatedAt, !created_at.is_empty(), !source.created_at.is_empty()),
351			(SystemColumn::UpdatedAt, !updated_at.is_empty(), !source.updated_at.is_empty()),
352			(SystemColumn::Time, !time.is_empty(), !source.time.is_empty()),
353		];
354		for (column, target_present, source_present) in pairs {
355			if target_present != source_present {
356				return Err(SystemColumnsError::PresenceMismatch {
357					column,
358					target_present,
359					source_present,
360				});
361			}
362		}
363		Ok(())
364	}
365
366	pub fn validate(&self, row_count: usize) -> Result<(), SystemColumnsError> {
367		let Self {
368			row_numbers,
369			partitions,
370			created_at,
371			updated_at,
372			time,
373		} = self;
374		let lengths = [
375			(SystemColumn::RowNumbers, row_numbers.len()),
376			(SystemColumn::Partitions, partitions.len()),
377			(SystemColumn::CreatedAt, created_at.len()),
378			(SystemColumn::UpdatedAt, updated_at.len()),
379			(SystemColumn::Time, time.len()),
380		];
381		for (column, len) in lengths {
382			if len != 0 && len != row_count {
383				return Err(SystemColumnsError::LengthMismatch {
384					column,
385					len,
386					row_count,
387				});
388			}
389		}
390		Ok(())
391	}
392
393	#[track_caller]
394	pub fn assert_invariants(&self, row_count: usize, ctx: &str) {
395		if let Err(err) = self.validate(row_count) {
396			panic!("{ctx}: {err}");
397		}
398	}
399}
400
401#[cfg(test)]
402mod tests {
403	use super::*;
404
405	fn dt(n: u64) -> DateTime {
406		DateTime::from_nanos(n)
407	}
408
409	fn partition(n: u128) -> Partition {
410		Partition::from(n)
411	}
412
413	fn populated() -> SystemColumns {
414		SystemColumns::new(
415			(1..5).map(RowNumber::from).collect(),
416			(0..4).map(|i| partition(i as u128)).collect(),
417			(0..4).map(|i| dt(1000 + i)).collect(),
418			(0..4).map(|i| dt(2000 + i)).collect(),
419			(0..4).map(|i| dt(3000 + i)).collect(),
420		)
421	}
422
423	#[track_caller]
424	fn assert_row_matches(actual: &SystemColumns, at: usize, source: &SystemColumns, from: usize) {
425		assert_eq!(actual.row_numbers()[at], source.row_numbers()[from], "row_numbers[{at}]");
426		assert_eq!(actual.partitions()[at], source.partitions()[from], "partitions[{at}]");
427		assert_eq!(actual.created_at()[at], source.created_at()[from], "created_at[{at}]");
428		assert_eq!(actual.updated_at()[at], source.updated_at()[from], "updated_at[{at}]");
429		assert_eq!(actual.time()[at], source.time()[from], "time[{at}]");
430	}
431
432	#[test]
433	fn permute_moves_every_sidecar_with_its_row() {
434		let source = populated();
435		let indices = [3, 0, 2, 1];
436		let permuted = source.permute(&indices);
437
438		assert_eq!(permuted.row_count(), Some(4));
439		for (at, &from) in indices.iter().enumerate() {
440			assert_row_matches(&permuted, at, &source, from);
441		}
442	}
443
444	#[test]
445	fn permute_trims_when_given_fewer_indices_than_rows() {
446		let source = populated();
447		let indices = [2, 0];
448		let permuted = source.permute(&indices);
449
450		assert_eq!(permuted.row_count(), Some(2));
451		for (at, &from) in indices.iter().enumerate() {
452			assert_row_matches(&permuted, at, &source, from);
453		}
454	}
455
456	#[test]
457	fn permute_duplicates_a_repeated_index() {
458		let source = populated();
459		let permuted = source.permute(&[1, 1, 1]);
460
461		assert_eq!(permuted.row_count(), Some(3));
462		for at in 0..3 {
463			assert_row_matches(&permuted, at, &source, 1);
464		}
465	}
466
467	#[test]
468	fn filter_keeps_masked_rows_intact() {
469		let source = populated();
470		let mut filtered = source.clone();
471		filtered.filter(&BitVec::from_slice(&[false, true, false, true]));
472
473		assert_eq!(filtered.row_count(), Some(2));
474		assert_row_matches(&filtered, 0, &source, 1);
475		assert_row_matches(&filtered, 1, &source, 3);
476	}
477
478	#[test]
479	fn take_trims_every_sidecar() {
480		let source = populated();
481		let mut taken = source.clone();
482		taken.take(2);
483
484		assert_eq!(taken.row_count(), Some(2));
485		assert_row_matches(&taken, 0, &source, 0);
486		assert_row_matches(&taken, 1, &source, 1);
487	}
488
489	#[test]
490	fn take_beyond_the_row_count_is_a_noop() {
491		let source = populated();
492		let mut taken = source.clone();
493		taken.take(99);
494		assert_eq!(taken, source);
495	}
496
497	#[test]
498	fn extend_concatenates_every_sidecar() {
499		let source = populated();
500		let mut acc = source.clone();
501		acc.extend(&source).unwrap();
502
503		assert_eq!(acc.row_count(), Some(8));
504		for i in 0..4 {
505			assert_row_matches(&acc, i, &source, i);
506			assert_row_matches(&acc, i + 4, &source, i);
507		}
508	}
509
510	#[test]
511	fn extend_rejects_a_presence_mismatch() {
512		let mut acc = populated();
513		let mut source = populated();
514		source.set_partitions(Vec::new());
515
516		assert_eq!(
517			acc.extend(&source).unwrap_err(),
518			SystemColumnsError::PresenceMismatch {
519				column: SystemColumn::Partitions,
520				target_present: true,
521				source_present: false,
522			}
523		);
524	}
525
526	#[test]
527	fn append_indices_appends_only_the_named_rows() {
528		let source = populated();
529		let mut acc = source.clone();
530		acc.append_indices(&source, &[3, 1]);
531
532		assert_eq!(acc.row_count(), Some(6));
533		assert_row_matches(&acc, 4, &source, 3);
534		assert_row_matches(&acc, 5, &source, 1);
535	}
536
537	#[test]
538	fn every_operation_leaves_an_absent_sidecar_absent() {
539		let mut source = populated();
540		source.set_partitions(Vec::new());
541
542		assert!(source.permute(&[1, 0]).partitions().is_empty(), "permute");
543
544		let mut filtered = source.clone();
545		filtered.filter(&BitVec::from_slice(&[true, false, true, false]));
546		assert!(filtered.partitions().is_empty(), "filter");
547
548		let mut taken = source.clone();
549		taken.take(2);
550		assert!(taken.partitions().is_empty(), "take");
551
552		let mut extended = source.clone();
553		extended.extend(&source).unwrap();
554		assert!(extended.partitions().is_empty(), "extend");
555
556		let mut appended = source.clone();
557		appended.append_indices(&source, &[0]);
558		assert!(appended.partitions().is_empty(), "append_indices");
559	}
560
561	#[test]
562	fn permuting_by_the_inverse_restores_the_original() {
563		let source = populated();
564		let forward = [2, 3, 1, 0];
565		let mut inverse = [0usize; 4];
566		for (at, &from) in forward.iter().enumerate() {
567			inverse[from] = at;
568		}
569		assert_eq!(source.permute(&forward).permute(&inverse), source);
570	}
571
572	#[test]
573	fn push_appends_one_row_to_every_sidecar() {
574		let mut acc = SystemColumns::empty();
575		acc.push(RowStamps {
576			row_number: Some(RowNumber::from(7)),
577			partition: Some(partition(2)),
578			created_at: Some(dt(10)),
579			updated_at: Some(dt(20)),
580			time: Some(dt(30)),
581		});
582
583		assert_eq!(acc.row_count(), Some(1));
584		assert_eq!(acc.row_numbers(), &[RowNumber::from(7)]);
585		assert_eq!(acc.partitions(), &[partition(2)]);
586		assert_eq!(acc.created_at(), &[dt(10)]);
587		assert_eq!(acc.updated_at(), &[dt(20)]);
588		assert_eq!(acc.time(), &[dt(30)]);
589	}
590
591	#[test]
592	fn pushing_rows_without_a_time_leaves_the_time_sidecar_absent() {
593		// A time-less object's rows must produce an empty #time vector, not a vector of sentinels.
594		// Absence is what downstream reads as "this source has no clock"; a filled vector of epochs
595		// would instead read as a source whose every row happened in 1970.
596		let mut acc = SystemColumns::empty();
597		for i in 0..3 {
598			acc.push(RowStamps {
599				row_number: Some(RowNumber::from(i + 1)),
600				partition: None,
601				created_at: Some(dt(10)),
602				updated_at: Some(dt(20)),
603				time: None,
604			});
605		}
606
607		assert_eq!(acc.row_count(), Some(3));
608		assert!(acc.time().is_empty(), "#time must stay absent rather than fill with sentinels");
609		acc.assert_invariants(3, "time-less push");
610	}
611
612	#[test]
613	fn a_time_less_batch_may_not_be_extended_by_a_timed_one() {
614		// Mixing the two would leave #time shorter than the row count, so every later positional read
615		// would silently attribute one row's time to a different row.
616		let mut untimed = SystemColumns::empty();
617		untimed.push(RowStamps {
618			row_number: Some(RowNumber::from(1)),
619			partition: None,
620			created_at: Some(dt(10)),
621			updated_at: Some(dt(20)),
622			time: None,
623		});
624
625		let mut timed = SystemColumns::empty();
626		timed.push(RowStamps {
627			row_number: Some(RowNumber::from(2)),
628			partition: None,
629			created_at: Some(dt(10)),
630			updated_at: Some(dt(20)),
631			time: Some(dt(30)),
632		});
633
634		assert_eq!(
635			untimed.extend(&timed).unwrap_err(),
636			SystemColumnsError::PresenceMismatch {
637				column: SystemColumn::Time,
638				target_present: false,
639				source_present: true,
640			}
641		);
642	}
643
644	#[test]
645	fn clear_empties_every_sidecar() {
646		let mut acc = populated();
647		acc.clear();
648		assert_eq!(acc.row_count(), None);
649		assert!(acc.is_empty());
650	}
651
652	#[test]
653	#[should_panic(expected = "time")]
654	fn assert_invariants_rejects_a_partial_sidecar() {
655		let mut partial = populated();
656		partial.time = vec![dt(1)];
657		partial.assert_invariants(4, "test");
658	}
659}