1use 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 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 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}