1use std::{
5 hash::Hash,
6 mem,
7 ops::{Index, IndexMut},
8};
9
10use indexmap::IndexMap;
11use reifydb_codec::encoded::{
12 row::EncodedRow,
13 shape::{RowShape, RowShapeField},
14};
15use reifydb_value::{
16 Result,
17 fragment::Fragment,
18 reifydb_assertions,
19 util::cowvec::CowVec,
20 value::{
21 Value, constraint::Constraint, datetime::DateTime, partition::Partition, row_number::RowNumber,
22 value_type::ValueType,
23 },
24};
25use serde::{Deserialize, Serialize};
26
27use crate::{
28 interface::catalog::column::Column as CatalogColumn,
29 return_internal_error,
30 row::Row,
31 value::column::{
32 ColumnBuffer, ColumnWithName, buffer::pool::ColumnBufferPool, data::Column, headers::ColumnHeaders,
33 },
34};
35
36#[derive(Debug, Clone, Serialize, Deserialize)]
37pub struct Columns {
38 pub row_numbers: CowVec<RowNumber>,
39 pub partitions: CowVec<Partition>,
40 pub created_at: CowVec<DateTime>,
41 pub updated_at: CowVec<DateTime>,
42 pub columns: CowVec<ColumnBuffer>,
43 pub names: CowVec<Fragment>,
44}
45
46#[derive(Debug, Clone, Copy)]
47pub struct ColumnRef<'a> {
48 name: &'a Fragment,
49 data: &'a ColumnBuffer,
50}
51
52impl Index<usize> for Columns {
53 type Output = ColumnBuffer;
54
55 fn index(&self, index: usize) -> &Self::Output {
56 &self.columns[index]
57 }
58}
59
60impl IndexMut<usize> for Columns {
61 fn index_mut(&mut self, index: usize) -> &mut Self::Output {
62 &mut self.columns.make_mut()[index]
63 }
64}
65
66impl<'a> ColumnRef<'a> {
67 pub fn new(name: &'a Fragment, data: &'a ColumnBuffer) -> Self {
68 Self {
69 name,
70 data,
71 }
72 }
73
74 pub fn name(&self) -> &'a Fragment {
75 self.name
76 }
77
78 pub fn data(&self) -> &'a ColumnBuffer {
79 self.data
80 }
81
82 pub fn get_type(&self) -> ValueType {
83 self.data.get_type()
84 }
85
86 pub fn column(&self) -> Column {
87 Column::from_column_buffer(self.data.clone())
88 }
89
90 pub fn with_new_data(&self, data: ColumnBuffer) -> ColumnWithName {
91 ColumnWithName::new(self.name.clone(), data)
92 }
93}
94
95fn value_to_buffer(value: Value) -> ColumnBuffer {
96 match value {
97 Value::None {
98 inner,
99 } => ColumnBuffer::none_typed(inner, 1),
100 Value::Boolean(v) => ColumnBuffer::bool([v]),
101 Value::Float4(v) => ColumnBuffer::float4([v.into()]),
102 Value::Float8(v) => ColumnBuffer::float8([v.into()]),
103 Value::Int1(v) => ColumnBuffer::int1([v]),
104 Value::Int2(v) => ColumnBuffer::int2([v]),
105 Value::Int4(v) => ColumnBuffer::int4([v]),
106 Value::Int8(v) => ColumnBuffer::int8([v]),
107 Value::Int16(v) => ColumnBuffer::int16([v]),
108 Value::Utf8(v) => ColumnBuffer::utf8([v]),
109 Value::Uint1(v) => ColumnBuffer::uint1([v]),
110 Value::Uint2(v) => ColumnBuffer::uint2([v]),
111 Value::Uint4(v) => ColumnBuffer::uint4([v]),
112 Value::Uint8(v) => ColumnBuffer::uint8([v]),
113 Value::Uint16(v) => ColumnBuffer::uint16([v]),
114 Value::Date(v) => ColumnBuffer::date([v]),
115 Value::DateTime(v) => ColumnBuffer::datetime([v]),
116 Value::Time(v) => ColumnBuffer::time([v]),
117 Value::Duration(v) => ColumnBuffer::duration([v]),
118 Value::IdentityId(v) => ColumnBuffer::identity_id([v]),
119 Value::Uuid4(v) => ColumnBuffer::uuid4([v]),
120 Value::Uuid7(v) => ColumnBuffer::uuid7([v]),
121 Value::Blob(v) => ColumnBuffer::blob([v]),
122 Value::Int(v) => ColumnBuffer::int(vec![v]),
123 Value::Uint(v) => ColumnBuffer::uint(vec![v]),
124 Value::Decimal(v) => ColumnBuffer::decimal(vec![v]),
125 Value::DictionaryId(v) => ColumnBuffer::dictionary_id(vec![v]),
126 Value::Any(v) => ColumnBuffer::any(vec![v]),
127 Value::Type(v) => ColumnBuffer::any(vec![Box::new(Value::Type(v))]),
128 Value::List(v) => ColumnBuffer::any(vec![Box::new(Value::List(v))]),
129 Value::Record(v) => ColumnBuffer::any(vec![Box::new(Value::Record(v))]),
130 Value::Tuple(v) => ColumnBuffer::any(vec![Box::new(Value::Tuple(v))]),
131 }
132}
133
134impl Columns {
135 pub fn scalar_value(&self) -> Value {
136 reifydb_assertions! {
137 assert_eq!(self.len(), 1, "scalar_value() requires exactly 1 column, got {}", self.len());
138 assert_eq!(
139 self.row_count(),
140 1,
141 "scalar_value() requires exactly 1 row, got {}",
142 self.row_count()
143 );
144 }
145 self.columns[0].get_value(0)
146 }
147
148 pub fn new(columns: Vec<ColumnWithName>) -> Self {
149 let n = columns.first().map_or(0, |c| c.data.len());
150 assert!(columns.iter().all(|c| c.data.len() == n));
151
152 let mut names = Vec::with_capacity(columns.len());
153 let mut buffers = Vec::with_capacity(columns.len());
154 for c in columns {
155 names.push(c.name);
156 buffers.push(c.data);
157 }
158
159 Self {
160 row_numbers: CowVec::new(Vec::new()),
161 partitions: CowVec::new(Vec::new()),
162 created_at: CowVec::new(Vec::new()),
163 updated_at: CowVec::new(Vec::new()),
164 columns: CowVec::new(buffers),
165 names: CowVec::new(names),
166 }
167 }
168
169 pub fn with_system_columns(
170 columns: Vec<ColumnWithName>,
171 row_numbers: Vec<RowNumber>,
172 created_at: Vec<DateTime>,
173 updated_at: Vec<DateTime>,
174 ) -> Self {
175 let n = columns.first().map_or(0, |c| c.data.len());
176 assert!(columns.iter().all(|c| c.data.len() == n));
177 assert_eq!(row_numbers.len(), n, "row_numbers length must match column data length");
178 assert_eq!(created_at.len(), n, "created_at length must match column data length");
179 assert_eq!(updated_at.len(), n, "updated_at length must match column data length");
180
181 let mut names = Vec::with_capacity(columns.len());
182 let mut buffers = Vec::with_capacity(columns.len());
183 for c in columns {
184 names.push(c.name);
185 buffers.push(c.data);
186 }
187
188 Self {
189 row_numbers: CowVec::new(row_numbers),
190 partitions: CowVec::new(Vec::new()),
191 created_at: CowVec::new(created_at),
192 updated_at: CowVec::new(updated_at),
193 columns: CowVec::new(buffers),
194 names: CowVec::new(names),
195 }
196 }
197
198 pub fn single_row<'b>(rows: impl IntoIterator<Item = (&'b str, Value)>) -> Columns {
199 let mut names = Vec::new();
200 let mut buffers = Vec::new();
201 for (name, value) in rows {
202 names.push(Fragment::internal(name));
203 buffers.push(value_to_buffer(value));
204 }
205 Self {
206 row_numbers: CowVec::new(Vec::new()),
207 partitions: CowVec::new(Vec::new()),
208 created_at: CowVec::new(Vec::new()),
209 updated_at: CowVec::new(Vec::new()),
210 columns: CowVec::new(buffers),
211 names: CowVec::new(names),
212 }
213 }
214
215 pub fn with_row_numbers(mut self, row_numbers: Vec<RowNumber>) -> Self {
216 let n = row_numbers.len();
217 self.row_numbers = CowVec::new(row_numbers);
218 if self.created_at.len() != n {
219 let now = DateTime::default();
220 self.created_at = CowVec::new(vec![now; n]);
221 self.updated_at = CowVec::new(vec![now; n]);
222 }
223 self
224 }
225
226 pub fn from_catalog_columns(cols: &[CatalogColumn]) -> Self {
227 let mut names = Vec::with_capacity(cols.len());
228 let mut buffers = Vec::with_capacity(cols.len());
229 for col in cols {
230 names.push(Fragment::internal(&col.name));
231 buffers.push(ColumnBuffer::with_capacity(col.constraint.get_type(), 0));
232 }
233 Self {
234 row_numbers: CowVec::new(Vec::new()),
235 partitions: CowVec::new(Vec::new()),
236 created_at: CowVec::new(Vec::new()),
237 updated_at: CowVec::new(Vec::new()),
238 columns: CowVec::new(buffers),
239 names: CowVec::new(names),
240 }
241 }
242
243 pub fn apply_headers(&mut self, headers: &ColumnHeaders) {
244 let n = self.len();
245 let names = self.names.make_mut();
246 for (i, name) in headers.columns.iter().enumerate() {
247 if i < n {
248 names[i] = name.clone();
249 }
250 }
251 }
252}
253
254impl Columns {
255 pub fn number(&self) -> RowNumber {
256 assert_eq!(self.row_count(), 1, "number() requires exactly 1 row, got {}", self.row_count());
257 if self.row_numbers.is_empty() {
258 RowNumber(0)
259 } else {
260 self.row_numbers[0]
261 }
262 }
263
264 pub fn shape(&self) -> (usize, usize) {
265 let row_count = if !self.row_numbers.is_empty() {
266 self.row_numbers.len()
267 } else {
268 self.columns.first().map(|c| c.len()).unwrap_or(0)
269 };
270 (row_count, self.len())
271 }
272
273 pub fn len(&self) -> usize {
274 self.columns.len()
275 }
276
277 pub fn is_empty(&self) -> bool {
278 self.columns.is_empty()
279 }
280
281 pub fn iter(&self) -> impl Iterator<Item = ColumnRef<'_>> + '_ {
282 self.names.iter().zip(self.columns.iter()).map(|(n, d)| ColumnRef::new(n, d))
283 }
284
285 pub fn first(&self) -> Option<ColumnRef<'_>> {
286 self.get(0)
287 }
288
289 pub fn last(&self) -> Option<ColumnRef<'_>> {
290 let n = self.len();
291 if n == 0 {
292 None
293 } else {
294 self.get(n - 1)
295 }
296 }
297
298 pub fn get(&self, index: usize) -> Option<ColumnRef<'_>> {
299 if index < self.len() {
300 Some(ColumnRef::new(&self.names[index], &self.columns[index]))
301 } else {
302 None
303 }
304 }
305
306 pub fn name_at(&self, index: usize) -> &Fragment {
307 &self.names[index]
308 }
309
310 pub fn data_at(&self, index: usize) -> &ColumnBuffer {
311 &self.columns[index]
312 }
313
314 pub fn data_at_mut(&mut self, index: usize) -> &mut ColumnBuffer {
315 &mut self.columns.make_mut()[index]
316 }
317
318 pub fn row(&self, i: usize) -> Vec<Value> {
319 self.columns.iter().map(|c| c.get_value(i)).collect()
320 }
321
322 pub fn column(&self, name: &str) -> Option<ColumnRef<'_>> {
323 self.names.iter().position(|n| n.text() == name).and_then(|i| self.get(i))
324 }
325
326 pub fn row_count(&self) -> usize {
327 if !self.row_numbers.is_empty() {
328 self.row_numbers.len()
329 } else {
330 self.columns.first().map_or(0, |col| col.len())
331 }
332 }
333
334 pub fn has_rows(&self) -> bool {
335 self.row_count() > 0
336 }
337
338 pub fn is_scalar(&self) -> bool {
339 self.len() == 1 && self.row_count() == 1
340 }
341
342 pub fn get_row(&self, index: usize) -> Vec<Value> {
343 self.columns.iter().map(|col| col.get_value(index)).collect()
344 }
345
346 #[track_caller]
347 pub fn assert_invariants(&self, ctx: &str) {
348 let n = self.columns.first().map_or(0, |c| c.len());
349 for (i, col) in self.columns.iter().enumerate() {
350 assert_eq!(
351 col.len(),
352 n,
353 "{ctx}: Columns column[{i}] has length {} but columns[0] has length {n}",
354 col.len(),
355 );
356 }
357 assert!(
358 self.row_numbers.is_empty() || self.row_numbers.len() == n,
359 "{ctx}: Columns.row_numbers.len() = {} but columns[0].len() = {n}",
360 self.row_numbers.len(),
361 );
362 assert!(
363 self.created_at.is_empty() || self.created_at.len() == n,
364 "{ctx}: Columns.created_at.len() = {} but columns[0].len() = {n}",
365 self.created_at.len(),
366 );
367 assert!(
368 self.updated_at.is_empty() || self.updated_at.len() == n,
369 "{ctx}: Columns.updated_at.len() = {} but columns[0].len() = {n}",
370 self.updated_at.len(),
371 );
372 }
373}
374
375impl Columns {
376 pub fn from_rows(names: &[&str], result_rows: &[Vec<Value>]) -> Self {
377 let column_count = names.len();
378
379 let mut name_vec: Vec<Fragment> = names.iter().map(Fragment::internal).collect();
380 let mut buffers: Vec<ColumnBuffer> =
381 (0..column_count).map(|_| ColumnBuffer::none_typed(ValueType::Boolean, 0)).collect();
382
383 for row in result_rows {
384 assert_eq!(row.len(), column_count, "row length does not match column count");
385 for (i, value) in row.iter().enumerate() {
386 buffers[i].push_value(value.clone());
387 }
388 }
389
390 let _ = &mut name_vec;
391 Self {
392 row_numbers: CowVec::new(Vec::new()),
393 partitions: CowVec::new(Vec::new()),
394 created_at: CowVec::new(Vec::new()),
395 updated_at: CowVec::new(Vec::new()),
396 columns: CowVec::new(buffers),
397 names: CowVec::new(name_vec),
398 }
399 }
400
401 pub fn from_encoded_rows(shape: &RowShape, ids: &[RowNumber], rows: &[EncodedRow]) -> Self {
402 assert_eq!(ids.len(), rows.len(), "ids length must match rows length");
403 let fields = shape.fields();
404 let row_count = rows.len();
405
406 let mut columns_vec: Vec<ColumnWithName> = Vec::with_capacity(fields.len());
407 for field in fields.iter() {
408 let mut data = ColumnBuffer::with_capacity(field.constraint.get_type(), row_count);
409 if field.constraint.get_type() == ValueType::DictionaryId
410 && let ColumnBuffer::DictionaryId(container) = &mut data
411 && let Some(Constraint::Dictionary(dict_id, _)) = field.constraint.constraint()
412 {
413 container.set_dictionary_id(*dict_id);
414 }
415 columns_vec.push(ColumnWithName {
416 name: Fragment::internal(&field.name),
417 data,
418 });
419 }
420
421 for encoded in rows {
422 for (i, _) in fields.iter().enumerate() {
423 columns_vec[i].data.push_value(shape.get_value(encoded, i));
424 }
425 }
426
427 let row_numbers: Vec<RowNumber> = ids.to_vec();
428 let created_at: Vec<DateTime> =
429 rows.iter().map(|r| DateTime::from_nanos(r.created_at_nanos())).collect();
430 let updated_at: Vec<DateTime> =
431 rows.iter().map(|r| DateTime::from_nanos(r.updated_at_nanos())).collect();
432
433 Self::with_system_columns(columns_vec, row_numbers, created_at, updated_at)
434 }
435}
436
437impl Columns {
438 pub fn empty() -> Self {
439 Self {
440 row_numbers: CowVec::new(Vec::new()),
441 partitions: CowVec::new(Vec::new()),
442 created_at: CowVec::new(Vec::new()),
443 updated_at: CowVec::new(Vec::new()),
444 columns: CowVec::new(Vec::new()),
445 names: CowVec::new(Vec::new()),
446 }
447 }
448}
449
450impl Default for Columns {
451 fn default() -> Self {
452 Self::empty()
453 }
454}
455
456impl Columns {
457 pub fn extract_by_indices(&self, indices: &[usize]) -> Columns {
458 if indices.is_empty() {
459 return Columns::empty();
460 }
461
462 let mut new_buffers: Vec<ColumnBuffer> = Vec::with_capacity(self.columns.len());
463 for col in self.columns.iter() {
464 let mut new_data = col.empty_like(indices.len());
465 for &idx in indices {
466 new_data.push_value(col.get_value(idx));
467 }
468 new_buffers.push(new_data);
469 }
470
471 let new_row_numbers: Vec<RowNumber> = if self.row_numbers.is_empty() {
472 Vec::new()
473 } else {
474 indices.iter().map(|&i| self.row_numbers[i]).collect()
475 };
476 let new_partitions: Vec<Partition> = if self.partitions.is_empty() {
477 Vec::new()
478 } else {
479 indices.iter().map(|&i| self.partitions[i]).collect()
480 };
481 let new_created_at: Vec<DateTime> = if self.created_at.is_empty() {
482 Vec::new()
483 } else {
484 indices.iter().map(|&i| self.created_at[i]).collect()
485 };
486 let new_updated_at: Vec<DateTime> = if self.updated_at.is_empty() {
487 Vec::new()
488 } else {
489 indices.iter().map(|&i| self.updated_at[i]).collect()
490 };
491 Columns {
492 row_numbers: CowVec::new(new_row_numbers),
493 partitions: CowVec::new(new_partitions),
494 created_at: CowVec::new(new_created_at),
495 updated_at: CowVec::new(new_updated_at),
496 columns: CowVec::new(new_buffers),
497 names: self.names.clone(),
498 }
499 }
500
501 pub fn extract_row(&self, index: usize) -> Columns {
502 self.extract_by_indices(&[index])
503 }
504
505 pub fn append_rows_by_indices(&mut self, source: &Columns, indices: &[usize]) {
506 if indices.is_empty() {
507 return;
508 }
509
510 if self.columns.is_empty() {
511 *self = source.extract_by_indices(indices);
512 return;
513 }
514
515 assert_eq!(
516 self.columns.len(),
517 source.columns.len(),
518 "append_rows: column count mismatch (self={}, source={})",
519 self.columns.len(),
520 source.columns.len(),
521 );
522
523 let self_cols = self.columns.make_mut();
524 for (i, src_col) in source.columns.iter().enumerate() {
525 for &idx in indices {
526 self_cols[i].push_value(src_col.get_value(idx));
527 }
528 }
529
530 if !source.row_numbers.is_empty() {
531 let rns = self.row_numbers.make_mut();
532 for &idx in indices {
533 rns.push(source.row_numbers[idx]);
534 }
535 }
536 if !source.partitions.is_empty() {
537 let parts = self.partitions.make_mut();
538 for &idx in indices {
539 parts.push(source.partitions[idx]);
540 }
541 }
542 if !source.created_at.is_empty() {
543 let cr = self.created_at.make_mut();
544 for &idx in indices {
545 cr.push(source.created_at[idx]);
546 }
547 }
548 if !source.updated_at.is_empty() {
549 let up = self.updated_at.make_mut();
550 for &idx in indices {
551 up.push(source.updated_at[idx]);
552 }
553 }
554 }
555
556 pub fn append_all(&mut self, source: Columns) -> Result<()> {
557 if source.row_count() == 0 {
558 return Ok(());
559 }
560 if self.columns.is_empty() {
561 *self = source;
562 return Ok(());
563 }
564
565 self.validate_append_compatibility(&source)?;
566 self.extend_data_columns(source.columns)?;
567 self.extend_system_columns(
568 &source.row_numbers,
569 &source.partitions,
570 &source.created_at,
571 &source.updated_at,
572 );
573 Ok(())
574 }
575
576 #[inline]
577 fn validate_append_compatibility(&self, source: &Columns) -> Result<()> {
578 if self.columns.len() != source.columns.len() {
579 return_internal_error!(
580 "Columns::append_all: column count mismatch (self={}, source={})",
581 self.columns.len(),
582 source.columns.len()
583 );
584 }
585
586 if self.row_numbers.is_empty() != source.row_numbers.is_empty() {
587 return_internal_error!(
588 "Columns::append_all: row_numbers population mismatch (self_empty={}, source_empty={})",
589 self.row_numbers.is_empty(),
590 source.row_numbers.is_empty()
591 );
592 }
593 if self.created_at.is_empty() != source.created_at.is_empty() {
594 return_internal_error!(
595 "Columns::append_all: created_at population mismatch (self_empty={}, source_empty={})",
596 self.created_at.is_empty(),
597 source.created_at.is_empty()
598 );
599 }
600 if self.updated_at.is_empty() != source.updated_at.is_empty() {
601 return_internal_error!(
602 "Columns::append_all: updated_at population mismatch (self_empty={}, source_empty={})",
603 self.updated_at.is_empty(),
604 source.updated_at.is_empty()
605 );
606 }
607 Ok(())
608 }
609
610 #[inline]
611 fn extend_data_columns(&mut self, source_columns: CowVec<ColumnBuffer>) -> Result<()> {
612 let dest_cols = self.columns.make_mut();
613 let source_cols = source_columns.into_inner();
614 reifydb_assertions! {
615 let dest_len = dest_cols.len();
616 let src_len = source_cols.len();
617 assert!(
618 dest_len == src_len,
619 "append_all extends destination columns by source index, so a source with more columns than \
620 the destination would index dest_cols out of bounds and panic mid-append, leaving self \
621 partially extended (dest_len={dest_len}, src_len={src_len})"
622 );
623 }
624 for (i, src_col) in source_cols.into_iter().enumerate() {
625 dest_cols[i].extend(src_col)?;
626 }
627 Ok(())
628 }
629
630 #[inline]
631 fn extend_system_columns(
632 &mut self,
633 source_row_numbers: &CowVec<RowNumber>,
634 source_partitions: &CowVec<Partition>,
635 source_created_at: &CowVec<DateTime>,
636 source_updated_at: &CowVec<DateTime>,
637 ) {
638 if !source_row_numbers.is_empty() {
639 self.row_numbers.extend_from_slice(source_row_numbers.as_slice());
640 }
641 if !source_partitions.is_empty() {
642 self.partitions.extend_from_slice(source_partitions.as_slice());
643 }
644 if !source_created_at.is_empty() {
645 self.created_at.extend_from_slice(source_created_at.as_slice());
646 }
647 if !source_updated_at.is_empty() {
648 self.updated_at.extend_from_slice(source_updated_at.as_slice());
649 }
650 }
651
652 pub fn concat(batches: Vec<Columns>) -> Result<Option<Columns>> {
653 let mut iter = batches.into_iter();
654 let mut merged = match iter.next() {
655 Some(first) => first,
656 None => return Ok(None),
657 };
658 for cols in iter {
659 merged.append_all(cols)?;
660 }
661 if merged.row_count() == 0 {
662 return Ok(None);
663 }
664 Ok(Some(merged))
665 }
666
667 pub fn remove_row(&mut self, row_number: RowNumber) -> bool {
668 let pos = self.row_numbers.iter().position(|&r| r == row_number);
669 let Some(idx) = pos else {
670 return false;
671 };
672
673 let kept_indices: Vec<usize> = (0..self.row_count()).filter(|&i| i != idx).collect();
674 *self = self.extract_by_indices(&kept_indices);
675 true
676 }
677
678 pub fn project_by_names(&self, names: &[String]) -> Columns {
679 let mut new_names = Vec::new();
680 let mut new_buffers = Vec::new();
681
682 for name in names {
683 if let Some(pos) = self.names.iter().position(|n| n.text() == name.as_str()) {
684 new_names.push(self.names[pos].clone());
685 new_buffers.push(self.columns[pos].clone());
686 }
687 }
688
689 if new_buffers.is_empty() {
690 return Columns::empty();
691 }
692
693 Columns {
694 row_numbers: self.row_numbers.clone(),
695 partitions: self.partitions.clone(),
696 created_at: self.created_at.clone(),
697 updated_at: self.updated_at.clone(),
698 columns: CowVec::new(new_buffers),
699 names: CowVec::new(new_names),
700 }
701 }
702
703 pub fn partition_by_keys<K: Hash + Eq + Clone>(&self, keys: &[K]) -> IndexMap<K, Columns> {
704 assert_eq!(keys.len(), self.row_count(), "keys length must match row count");
705
706 let mut key_to_indices: IndexMap<K, Vec<usize>> = IndexMap::new();
707 for (idx, key) in keys.iter().enumerate() {
708 key_to_indices.entry(key.clone()).or_default().push(idx);
709 }
710
711 key_to_indices.into_iter().map(|(key, indices)| (key, self.extract_by_indices(&indices))).collect()
712 }
713
714 pub fn from_row(row: &Row) -> Self {
715 let mut out = Columns::empty();
716 out.reset_from_row(row);
717 out
718 }
719
720 pub fn reset_from_row(&mut self, row: &Row) {
721 let field_count = row.shape.fields().len();
722
723 self.row_numbers.clear();
724 self.created_at.clear();
725 self.updated_at.clear();
726 self.columns.clear();
727 self.names.clear();
728
729 self.columns.make_mut().reserve(field_count);
730 self.names.make_mut().reserve(field_count);
731
732 self.row_numbers.push(row.number);
733 self.created_at.push(DateTime::from_nanos(row.encoded.created_at_nanos()));
734 self.updated_at.push(DateTime::from_nanos(row.encoded.updated_at_nanos()));
735
736 for (idx, field) in row.shape.fields().iter().enumerate() {
737 let value = row.shape.get_value(&row.encoded, idx);
738
739 let column_type = if matches!(value, Value::None { .. }) {
740 field.constraint.get_type()
741 } else {
742 value.get_type()
743 };
744
745 let mut data = if column_type.is_option() {
746 ColumnBuffer::none_typed(column_type.clone(), 0)
747 } else {
748 ColumnBuffer::with_capacity(column_type.clone(), 1)
749 };
750 data.push_value(value);
751
752 if column_type == ValueType::DictionaryId
753 && let ColumnBuffer::DictionaryId(container) = &mut data
754 && let Some(Constraint::Dictionary(dict_id, _)) = field.constraint.constraint()
755 {
756 container.set_dictionary_id(*dict_id);
757 }
758
759 let name = row.shape.get_field_name(idx).expect("RowShape missing name for field");
760
761 self.names.push(Fragment::internal(name));
762 self.columns.push(data);
763 }
764 }
765
766 pub fn reset_from_row_with_pool(&mut self, row: &Row, pool: &ColumnBufferPool) {
767 let field_count = row.shape.fields().len();
768
769 self.row_numbers.clear();
770 self.created_at.clear();
771 self.updated_at.clear();
772 self.names.clear();
773
774 self.row_numbers.push(row.number);
775 self.created_at.push(DateTime::from_nanos(row.encoded.created_at_nanos()));
776 self.updated_at.push(DateTime::from_nanos(row.encoded.updated_at_nanos()));
777
778 let columns_vec = self.columns.make_mut();
779 let names_vec = self.names.make_mut();
780
781 while columns_vec.len() > field_count {
782 if let Some(buf) = columns_vec.pop() {
783 pool.release(buf);
784 }
785 }
786
787 columns_vec.reserve(field_count);
788 names_vec.reserve(field_count);
789
790 for (idx, field) in row.shape.fields().iter().enumerate() {
791 let value = row.shape.get_value(&row.encoded, idx);
792
793 let column_type = if matches!(value, Value::None { .. }) {
794 field.constraint.get_type()
795 } else {
796 value.get_type()
797 };
798
799 if idx < columns_vec.len() {
800 if columns_vec[idx].get_type() == column_type {
801 columns_vec[idx].clear();
802 } else {
803 let replacement = if column_type.is_option() {
804 ColumnBuffer::none_typed(column_type.clone(), 0)
805 } else {
806 pool.acquire(&column_type, 1)
807 };
808 let old = mem::replace(&mut columns_vec[idx], replacement);
809 pool.release(old);
810 }
811 } else {
812 let fresh = if column_type.is_option() {
813 ColumnBuffer::none_typed(column_type.clone(), 0)
814 } else {
815 pool.acquire(&column_type, 1)
816 };
817 columns_vec.push(fresh);
818 }
819
820 columns_vec[idx].push_value(value);
821
822 if column_type == ValueType::DictionaryId
823 && let ColumnBuffer::DictionaryId(container) = &mut columns_vec[idx]
824 && let Some(Constraint::Dictionary(dict_id, _)) = field.constraint.constraint()
825 {
826 container.set_dictionary_id(*dict_id);
827 }
828
829 let name = row.shape.get_field_name(idx).expect("RowShape missing name for field");
830 names_vec.push(Fragment::internal(name));
831 }
832 }
833
834 pub fn push_row(&mut self, row: &Row) {
835 let field_count = row.shape.fields().len();
836
837 if self.columns.is_empty() {
838 self.columns.make_mut().reserve(field_count);
839 self.names.make_mut().reserve(field_count);
840 self.row_numbers.push(row.number);
841 self.created_at.push(DateTime::from_nanos(row.encoded.created_at_nanos()));
842 self.updated_at.push(DateTime::from_nanos(row.encoded.updated_at_nanos()));
843
844 for (idx, field) in row.shape.fields().iter().enumerate() {
845 let value = row.shape.get_value(&row.encoded, idx);
846
847 let column_type = if matches!(value, Value::None { .. }) {
848 field.constraint.get_type()
849 } else {
850 value.get_type()
851 };
852
853 let mut data = if column_type.is_option() {
854 ColumnBuffer::none_typed(column_type.clone(), 0)
855 } else {
856 ColumnBuffer::with_capacity(column_type.clone(), 1)
857 };
858 data.push_value(value);
859
860 if column_type == ValueType::DictionaryId
861 && let ColumnBuffer::DictionaryId(container) = &mut data
862 && let Some(Constraint::Dictionary(dict_id, _)) = field.constraint.constraint()
863 {
864 container.set_dictionary_id(*dict_id);
865 }
866
867 let name = row.shape.get_field_name(idx).expect("RowShape missing name for field");
868 self.names.push(Fragment::internal(name));
869 self.columns.push(data);
870 }
871 } else if self.columns.len() == field_count {
872 let columns_vec = self.columns.make_mut();
873 for (idx, column) in columns_vec.iter_mut().enumerate() {
874 let value = row.shape.get_value(&row.encoded, idx);
875 column.push_value(value);
876 }
877 self.row_numbers.push(row.number);
878 self.created_at.push(DateTime::from_nanos(row.encoded.created_at_nanos()));
879 self.updated_at.push(DateTime::from_nanos(row.encoded.updated_at_nanos()));
880 }
881 }
882
883 pub fn push_rows(&mut self, rows: &[Row]) {
884 let Some(first) = rows.first() else {
885 return;
886 };
887 if !self.columns.is_empty() {
888 for row in rows {
889 self.push_row(row);
890 }
891 return;
892 }
893
894 let capacity = rows.len();
895 let field_count = first.shape.fields().len();
896 self.columns.make_mut().reserve(field_count);
897 self.names.make_mut().reserve(field_count);
898 self.row_numbers.make_mut().reserve(capacity);
899 self.created_at.make_mut().reserve(capacity);
900 self.updated_at.make_mut().reserve(capacity);
901
902 for (idx, field) in first.shape.fields().iter().enumerate() {
903 let value = first.shape.get_value(&first.encoded, idx);
904
905 let column_type = if matches!(value, Value::None { .. }) {
906 field.constraint.get_type()
907 } else {
908 value.get_type()
909 };
910
911 let mut data = ColumnBuffer::with_capacity(column_type.clone(), capacity);
912
913 if column_type == ValueType::DictionaryId
914 && let ColumnBuffer::DictionaryId(container) = &mut data
915 && let Some(Constraint::Dictionary(dict_id, _)) = field.constraint.constraint()
916 {
917 container.set_dictionary_id(*dict_id);
918 }
919
920 let name = first.shape.get_field_name(idx).expect("RowShape missing name for field");
921 self.names.push(Fragment::internal(name));
922 self.columns.push(data);
923 }
924
925 let columns_vec = self.columns.make_mut();
926 for row in rows {
927 for (idx, column) in columns_vec.iter_mut().enumerate() {
928 column.push_value(row.shape.get_value(&row.encoded, idx));
929 }
930 }
931 for row in rows {
932 self.row_numbers.push(row.number);
933 self.created_at.push(DateTime::from_nanos(row.encoded.created_at_nanos()));
934 self.updated_at.push(DateTime::from_nanos(row.encoded.updated_at_nanos()));
935 }
936 }
937
938 pub fn to_single_row(&self) -> Row {
939 assert_eq!(self.row_count(), 1, "to_row() requires exactly 1 row, got {}", self.row_count());
940 assert_eq!(
941 self.row_numbers.len(),
942 1,
943 "to_row() requires exactly 1 row number, got {}",
944 self.row_numbers.len()
945 );
946
947 let row_number = *self.row_numbers.first().unwrap();
948
949 let fields: Vec<RowShapeField> = self
950 .names
951 .iter()
952 .zip(self.columns.iter())
953 .map(|(name, data)| RowShapeField::unconstrained(name.text().to_string(), data.get_type()))
954 .collect();
955
956 let layout = RowShape::new(fields);
957 let mut encoded = layout.allocate();
958
959 let values: Vec<Value> = self.columns.iter().map(|col| col.get_value(0)).collect();
960 layout.set_values(&mut encoded, &values);
961
962 Row {
963 number: row_number,
964 encoded,
965 shape: layout,
966 }
967 }
968}
969
970#[cfg(test)]
971pub mod tests {
972 use std::str::FromStr;
973
974 use reifydb_value::value::{
975 blob::Blob,
976 constraint::{bytes::MaxBytes, precision::Precision, scale::Scale},
977 date::Date,
978 datetime::DateTime,
979 decimal::Decimal,
980 dictionary::{DictionaryEntryId, DictionaryId},
981 duration::Duration,
982 identity::IdentityId,
983 int::Int,
984 ordered_f64::OrderedF64,
985 time::Time,
986 uint::Uint,
987 uuid::{Uuid4, Uuid7},
988 };
989 use uuid::{Timestamp, Uuid};
990
991 use super::*;
992
993 fn uuid7_at(a: u64, b: u16) -> Uuid7 {
994 Uuid7::from(Uuid::new_v7(Timestamp::from_gregorian_time(a, b)))
995 }
996
997 fn assert_extract_preserves_values(buffer: ColumnBuffer, indices: &[usize]) {
1003 let original = Columns::new(vec![ColumnWithName::new("c", buffer)]);
1004 let extracted = original.extract_by_indices(indices);
1005
1006 assert_eq!(extracted.len(), 1, "column count must be preserved");
1007 assert_eq!(extracted.row_count(), indices.len(), "row count must equal number of indices");
1008
1009 let src = original.data_at(0);
1010 let dst = extracted.data_at(0);
1011 assert_eq!(dst.get_type(), src.get_type(), "value type must be preserved");
1012 for (j, &idx) in indices.iter().enumerate() {
1013 assert_eq!(
1014 dst.get_value(j),
1015 src.get_value(idx),
1016 "value at extracted row {j} must equal source row {idx}"
1017 );
1018 }
1019 }
1020
1021 #[test]
1022 fn extract_by_indices_preserves_bool_values() {
1023 assert_extract_preserves_values(ColumnBuffer::bool([true, false, true, false]), &[3, 1, 2]);
1024 }
1025
1026 #[test]
1027 fn extract_by_indices_preserves_float4_values() {
1028 assert_extract_preserves_values(ColumnBuffer::float4([1.0f32, 2.5, -3.0, 4.25]), &[3, 1, 2]);
1029 }
1030
1031 #[test]
1032 fn extract_by_indices_preserves_float8_values() {
1033 assert_extract_preserves_values(ColumnBuffer::float8([1.0f64, 2.5, -3.0, 4.25]), &[3, 1, 2]);
1034 }
1035
1036 #[test]
1037 fn extract_by_indices_preserves_int1_values() {
1038 assert_extract_preserves_values(ColumnBuffer::int1([-1i8, 2, -3, 4]), &[3, 1, 2]);
1039 }
1040
1041 #[test]
1042 fn extract_by_indices_preserves_int2_values() {
1043 assert_extract_preserves_values(ColumnBuffer::int2([-1i16, 2, -3, 4]), &[3, 1, 2]);
1044 }
1045
1046 #[test]
1047 fn extract_by_indices_preserves_int4_values() {
1048 assert_extract_preserves_values(ColumnBuffer::int4([-1i32, 2, -3, 4]), &[3, 1, 2]);
1049 }
1050
1051 #[test]
1052 fn extract_by_indices_preserves_int8_values() {
1053 assert_extract_preserves_values(ColumnBuffer::int8([-1i64, 2, -3, 4]), &[3, 1, 2]);
1054 }
1055
1056 #[test]
1057 fn extract_by_indices_preserves_int16_values() {
1058 assert_extract_preserves_values(ColumnBuffer::int16([-1i128, 2, -3, 4]), &[3, 1, 2]);
1059 }
1060
1061 #[test]
1062 fn extract_by_indices_preserves_uint1_values() {
1063 assert_extract_preserves_values(ColumnBuffer::uint1([1u8, 2, 3, 4]), &[3, 1, 2]);
1064 }
1065
1066 #[test]
1067 fn extract_by_indices_preserves_uint2_values() {
1068 assert_extract_preserves_values(ColumnBuffer::uint2([1u16, 2, 3, 4]), &[3, 1, 2]);
1069 }
1070
1071 #[test]
1072 fn extract_by_indices_preserves_uint4_values() {
1073 assert_extract_preserves_values(ColumnBuffer::uint4([1u32, 2, 3, 4]), &[3, 1, 2]);
1074 }
1075
1076 #[test]
1077 fn extract_by_indices_preserves_uint8_values() {
1078 assert_extract_preserves_values(ColumnBuffer::uint8([1u64, 2, 3, 4]), &[3, 1, 2]);
1079 }
1080
1081 #[test]
1082 fn extract_by_indices_preserves_uint16_values() {
1083 assert_extract_preserves_values(ColumnBuffer::uint16([1u128, 2, 3, 4]), &[3, 1, 2]);
1084 }
1085
1086 #[test]
1087 fn extract_by_indices_preserves_utf8_values() {
1088 assert_extract_preserves_values(ColumnBuffer::utf8(["a", "bb", "ccc", "dddd"]), &[3, 1, 2]);
1089 }
1090
1091 #[test]
1092 fn extract_by_indices_preserves_date_values() {
1093 let data = [
1094 Date::from_ymd(2025, 1, 1).unwrap(),
1095 Date::from_ymd(2025, 6, 15).unwrap(),
1096 Date::from_ymd(2024, 12, 31).unwrap(),
1097 Date::from_ymd(2000, 2, 29).unwrap(),
1098 ];
1099 assert_extract_preserves_values(ColumnBuffer::date(data), &[3, 1, 2]);
1100 }
1101
1102 #[test]
1103 fn extract_by_indices_preserves_datetime_values() {
1104 let data = [
1105 DateTime::from_timestamp(1000).unwrap(),
1106 DateTime::from_timestamp(2000).unwrap(),
1107 DateTime::from_timestamp(3000).unwrap(),
1108 DateTime::from_timestamp(4000).unwrap(),
1109 ];
1110 assert_extract_preserves_values(ColumnBuffer::datetime(data), &[3, 1, 2]);
1111 }
1112
1113 #[test]
1114 fn extract_by_indices_preserves_time_values() {
1115 let data = [
1116 Time::from_hms(0, 0, 0).unwrap(),
1117 Time::from_hms(12, 30, 45).unwrap(),
1118 Time::from_hms(23, 59, 59).unwrap(),
1119 Time::from_hms(6, 15, 0).unwrap(),
1120 ];
1121 assert_extract_preserves_values(ColumnBuffer::time(data), &[3, 1, 2]);
1122 }
1123
1124 #[test]
1125 fn extract_by_indices_preserves_duration_values() {
1126 let data = [
1127 Duration::from_days(1).unwrap(),
1128 Duration::from_days(7).unwrap(),
1129 Duration::from_days(30).unwrap(),
1130 Duration::from_days(365).unwrap(),
1131 ];
1132 assert_extract_preserves_values(ColumnBuffer::duration(data), &[3, 1, 2]);
1133 }
1134
1135 #[test]
1136 fn extract_by_indices_preserves_identity_id_values() {
1137 let data = [IdentityId::root(), IdentityId::system(), IdentityId::anonymous(), IdentityId::root()];
1138 assert_extract_preserves_values(ColumnBuffer::identity_id(data), &[3, 1, 2]);
1139 }
1140
1141 #[test]
1142 fn extract_by_indices_preserves_uuid4_values() {
1143 let data = [Uuid4::generate(), Uuid4::generate(), Uuid4::generate(), Uuid4::generate()];
1144 assert_extract_preserves_values(ColumnBuffer::uuid4(data), &[3, 1, 2]);
1145 }
1146
1147 #[test]
1148 fn extract_by_indices_preserves_uuid7_values() {
1149 let data = [uuid7_at(1, 1), uuid7_at(1, 2), uuid7_at(2, 1), uuid7_at(2, 2)];
1150 assert_extract_preserves_values(ColumnBuffer::uuid7(data), &[3, 1, 2]);
1151 }
1152
1153 #[test]
1154 fn extract_by_indices_preserves_blob_values() {
1155 let data = [
1156 Blob::new(vec![1]),
1157 Blob::new(vec![2, 3]),
1158 Blob::new(vec![4, 5, 6]),
1159 Blob::new(vec![7, 8, 9, 10]),
1160 ];
1161 assert_extract_preserves_values(ColumnBuffer::blob(data), &[3, 1, 2]);
1162 }
1163
1164 #[test]
1165 fn extract_by_indices_preserves_int_values() {
1166 let data = [Int::from(-1i64), Int::from(2i64), Int::from(-3i64), Int::from(4i64)];
1167 assert_extract_preserves_values(ColumnBuffer::int(data), &[3, 1, 2]);
1168 }
1169
1170 #[test]
1171 fn extract_by_indices_preserves_uint_values() {
1172 let data = [Uint::from(1u64), Uint::from(2u64), Uint::from(3u64), Uint::from(4u64)];
1173 assert_extract_preserves_values(ColumnBuffer::uint(data), &[3, 1, 2]);
1174 }
1175
1176 #[test]
1177 fn extract_by_indices_preserves_decimal_values() {
1178 let data = [
1179 Decimal::from_str("1.50").unwrap(),
1180 Decimal::from_str("2.25").unwrap(),
1181 Decimal::from_str("-3.75").unwrap(),
1182 Decimal::from_str("4.00").unwrap(),
1183 ];
1184 assert_extract_preserves_values(ColumnBuffer::decimal(data), &[3, 1, 2]);
1185 }
1186
1187 #[test]
1188 fn extract_by_indices_preserves_any_values() {
1189 let data = [
1190 Box::new(Value::Int4(1)),
1191 Box::new(Value::Utf8("two".to_string())),
1192 Box::new(Value::Boolean(true)),
1193 Box::new(Value::none()),
1194 ];
1195 assert_extract_preserves_values(ColumnBuffer::any(data), &[3, 1, 2]);
1196 }
1197
1198 #[test]
1199 fn extract_by_indices_preserves_dictionary_id_values() {
1200 let data = [
1201 DictionaryEntryId::U2(10),
1202 DictionaryEntryId::U2(20),
1203 DictionaryEntryId::U2(30),
1204 DictionaryEntryId::U2(40),
1205 ];
1206 assert_extract_preserves_values(ColumnBuffer::dictionary_id(data), &[3, 1, 2]);
1207 }
1208
1209 #[test]
1210 fn extract_by_indices_preserves_option_values_including_none() {
1211 let mut buffer = ColumnBuffer::with_capacity(ValueType::Option(Box::new(ValueType::Int4)), 0);
1212 buffer.push_value(Value::Int4(1));
1213 buffer.push_value(Value::none());
1214 buffer.push_value(Value::Int4(3));
1215 buffer.push_value(Value::none());
1216 assert_extract_preserves_values(buffer, &[3, 1, 2, 0]);
1217 }
1218
1219 #[test]
1220 fn extract_by_indices_empty_indices_yields_empty_columns() {
1221 let original = Columns::new(vec![ColumnWithName::int4("c", [1, 2, 3])]);
1222 let extracted = original.extract_by_indices(&[]);
1223 assert_eq!(extracted.row_count(), 0);
1224 assert!(extracted.is_empty());
1225 }
1226
1227 #[test]
1228 fn extract_by_indices_full_identity_reproduces_all_rows() {
1229 assert_extract_preserves_values(ColumnBuffer::int4([10, 20, 30, 40]), &[0, 1, 2, 3]);
1230 }
1231
1232 #[test]
1233 fn extract_by_indices_duplicate_index_duplicates_row() {
1234 let original = Columns::new(vec![ColumnWithName::int4("c", [10, 20, 30])]);
1235 let extracted = original.extract_by_indices(&[1, 1, 1]);
1236 assert_eq!(extracted.row_count(), 3);
1237 assert_eq!(extracted.data_at(0).get_value(0), Value::Int4(20));
1238 assert_eq!(extracted.data_at(0).get_value(1), Value::Int4(20));
1239 assert_eq!(extracted.data_at(0).get_value(2), Value::Int4(20));
1240 }
1241
1242 #[test]
1243 fn extract_by_indices_extracts_multiple_columns_consistently() {
1244 let original = Columns::new(vec![
1245 ColumnWithName::int4("id", [1, 2, 3, 4]),
1246 ColumnWithName::utf8(
1247 "name",
1248 ["a".to_string(), "b".to_string(), "c".to_string(), "d".to_string()],
1249 ),
1250 ColumnWithName::bool("flag", [true, false, true, false]),
1251 ]);
1252 let extracted = original.extract_by_indices(&[2, 0]);
1253
1254 assert_eq!(extracted.len(), 3);
1255 assert_eq!(extracted.row_count(), 2);
1256 assert_eq!(extracted.column("id").unwrap().data().get_value(0), Value::Int4(3));
1257 assert_eq!(extracted.column("id").unwrap().data().get_value(1), Value::Int4(1));
1258 assert_eq!(extracted.column("name").unwrap().data().get_value(0), Value::Utf8("c".to_string()));
1259 assert_eq!(extracted.column("name").unwrap().data().get_value(1), Value::Utf8("a".to_string()));
1260 assert_eq!(extracted.column("flag").unwrap().data().get_value(0), Value::Boolean(true));
1261 assert_eq!(extracted.column("flag").unwrap().data().get_value(1), Value::Boolean(true));
1262 }
1263
1264 #[test]
1265 fn extract_by_indices_extracts_system_columns_in_order() {
1266 let columns = vec![ColumnWithName::int4("id", [10, 20, 30, 40])];
1267 let row_numbers = vec![RowNumber::from(1), RowNumber::from(2), RowNumber::from(3), RowNumber::from(4)];
1268 let created_at = vec![
1269 DateTime::from_timestamp(1000).unwrap(),
1270 DateTime::from_timestamp(2000).unwrap(),
1271 DateTime::from_timestamp(3000).unwrap(),
1272 DateTime::from_timestamp(4000).unwrap(),
1273 ];
1274 let updated_at = vec![
1275 DateTime::from_timestamp(1100).unwrap(),
1276 DateTime::from_timestamp(2200).unwrap(),
1277 DateTime::from_timestamp(3300).unwrap(),
1278 DateTime::from_timestamp(4400).unwrap(),
1279 ];
1280 let original = Columns::with_system_columns(columns, row_numbers, created_at, updated_at);
1281
1282 let extracted = original.extract_by_indices(&[3, 0]);
1283
1284 let rns: Vec<RowNumber> = extracted.row_numbers.iter().cloned().collect();
1285 assert_eq!(rns, vec![RowNumber::from(4), RowNumber::from(1)], "row_numbers must follow indices");
1286 assert_eq!(
1287 extracted.created_at.iter().cloned().collect::<Vec<_>>(),
1288 vec![DateTime::from_timestamp(4000).unwrap(), DateTime::from_timestamp(1000).unwrap()],
1289 "created_at must follow indices"
1290 );
1291 assert_eq!(
1292 extracted.updated_at.iter().cloned().collect::<Vec<_>>(),
1293 vec![DateTime::from_timestamp(4400).unwrap(), DateTime::from_timestamp(1100).unwrap()],
1294 "updated_at must follow indices"
1295 );
1296 }
1297
1298 #[test]
1304 fn extract_by_indices_preserves_dictionary_id_metadata() {
1305 let mut buffer = ColumnBuffer::dictionary_id([
1306 DictionaryEntryId::U2(10),
1307 DictionaryEntryId::U2(20),
1308 DictionaryEntryId::U2(30),
1309 ]);
1310 match &mut buffer {
1311 ColumnBuffer::DictionaryId(container) => container.set_dictionary_id(DictionaryId(42)),
1312 _ => unreachable!("dictionary_id factory must build a DictionaryId buffer"),
1313 }
1314
1315 let original = Columns::new(vec![ColumnWithName::new("token", buffer)]);
1316 let extracted = original.extract_by_indices(&[2, 0]);
1317
1318 match extracted.data_at(0) {
1319 ColumnBuffer::DictionaryId(container) => {
1320 assert_eq!(
1321 container.dictionary_id(),
1322 Some(DictionaryId(42)),
1323 "dictionary_id metadata must survive extraction"
1324 );
1325 }
1326 other => panic!("expected DictionaryId buffer, got {:?}", other.get_type()),
1327 }
1328 }
1329
1330 #[test]
1331 fn extract_by_indices_preserves_utf8_max_bytes_metadata() {
1332 let mut buffer = ColumnBuffer::utf8(["a", "bb", "ccc"]);
1333 match &mut buffer {
1334 ColumnBuffer::Utf8 {
1335 max_bytes,
1336 ..
1337 } => *max_bytes = MaxBytes::new(255),
1338 _ => unreachable!(),
1339 }
1340
1341 let original = Columns::new(vec![ColumnWithName::new("c", buffer)]);
1342 let extracted = original.extract_by_indices(&[2, 0]);
1343
1344 match extracted.data_at(0) {
1345 ColumnBuffer::Utf8 {
1346 max_bytes,
1347 ..
1348 } => assert_eq!(*max_bytes, MaxBytes::new(255), "Utf8 max_bytes must survive extraction"),
1349 other => panic!("expected Utf8 buffer, got {:?}", other.get_type()),
1350 }
1351 }
1352
1353 #[test]
1354 fn extract_by_indices_preserves_blob_max_bytes_metadata() {
1355 let mut buffer = ColumnBuffer::blob([Blob::new(vec![1]), Blob::new(vec![2, 3]), Blob::new(vec![4])]);
1356 match &mut buffer {
1357 ColumnBuffer::Blob {
1358 max_bytes,
1359 ..
1360 } => *max_bytes = MaxBytes::new(1024),
1361 _ => unreachable!(),
1362 }
1363
1364 let original = Columns::new(vec![ColumnWithName::new("c", buffer)]);
1365 let extracted = original.extract_by_indices(&[2, 0]);
1366
1367 match extracted.data_at(0) {
1368 ColumnBuffer::Blob {
1369 max_bytes,
1370 ..
1371 } => assert_eq!(*max_bytes, MaxBytes::new(1024), "Blob max_bytes must survive extraction"),
1372 other => panic!("expected Blob buffer, got {:?}", other.get_type()),
1373 }
1374 }
1375
1376 #[test]
1377 fn extract_by_indices_preserves_int_max_bytes_metadata() {
1378 let mut buffer = ColumnBuffer::int([Int::from(1i64), Int::from(2i64), Int::from(3i64)]);
1379 match &mut buffer {
1380 ColumnBuffer::Int {
1381 max_bytes,
1382 ..
1383 } => *max_bytes = MaxBytes::new(16),
1384 _ => unreachable!(),
1385 }
1386
1387 let original = Columns::new(vec![ColumnWithName::new("c", buffer)]);
1388 let extracted = original.extract_by_indices(&[2, 0]);
1389
1390 match extracted.data_at(0) {
1391 ColumnBuffer::Int {
1392 max_bytes,
1393 ..
1394 } => assert_eq!(*max_bytes, MaxBytes::new(16), "Int max_bytes must survive extraction"),
1395 other => panic!("expected Int buffer, got {:?}", other.get_type()),
1396 }
1397 }
1398
1399 #[test]
1400 fn extract_by_indices_preserves_uint_max_bytes_metadata() {
1401 let mut buffer = ColumnBuffer::uint([Uint::from(1u64), Uint::from(2u64), Uint::from(3u64)]);
1402 match &mut buffer {
1403 ColumnBuffer::Uint {
1404 max_bytes,
1405 ..
1406 } => *max_bytes = MaxBytes::new(8),
1407 _ => unreachable!(),
1408 }
1409
1410 let original = Columns::new(vec![ColumnWithName::new("c", buffer)]);
1411 let extracted = original.extract_by_indices(&[2, 0]);
1412
1413 match extracted.data_at(0) {
1414 ColumnBuffer::Uint {
1415 max_bytes,
1416 ..
1417 } => assert_eq!(*max_bytes, MaxBytes::new(8), "Uint max_bytes must survive extraction"),
1418 other => panic!("expected Uint buffer, got {:?}", other.get_type()),
1419 }
1420 }
1421
1422 #[test]
1423 fn extract_by_indices_preserves_decimal_precision_and_scale_metadata() {
1424 let mut buffer = ColumnBuffer::decimal([
1425 Decimal::from_str("1.50").unwrap(),
1426 Decimal::from_str("2.25").unwrap(),
1427 Decimal::from_str("3.75").unwrap(),
1428 ]);
1429 match &mut buffer {
1430 ColumnBuffer::Decimal {
1431 precision,
1432 scale,
1433 ..
1434 } => {
1435 *precision = Precision::new(10);
1436 *scale = Scale::new(2);
1437 }
1438 _ => unreachable!(),
1439 }
1440
1441 let original = Columns::new(vec![ColumnWithName::new("c", buffer)]);
1442 let extracted = original.extract_by_indices(&[2, 0]);
1443
1444 match extracted.data_at(0) {
1445 ColumnBuffer::Decimal {
1446 precision,
1447 scale,
1448 ..
1449 } => {
1450 assert_eq!(*precision, Precision::new(10), "Decimal precision must survive extraction");
1451 assert_eq!(*scale, Scale::new(2), "Decimal scale must survive extraction");
1452 }
1453 other => panic!("expected Decimal buffer, got {:?}", other.get_type()),
1454 }
1455 }
1456
1457 #[test]
1458 fn test_single_row_temporal_types() {
1459 let date = Date::from_ymd(2025, 1, 15).unwrap();
1460 let datetime = DateTime::from_timestamp(1642694400).unwrap();
1461 let time = Time::from_hms(14, 30, 45).unwrap();
1462 let duration = Duration::from_days(30).unwrap();
1463
1464 let columns = Columns::single_row([
1465 ("date_col", Value::Date(date.clone())),
1466 ("datetime_col", Value::DateTime(datetime.clone())),
1467 ("time_col", Value::Time(time.clone())),
1468 ("interval_col", Value::Duration(duration.clone())),
1469 ]);
1470
1471 assert_eq!(columns.len(), 4);
1472 assert_eq!(columns.shape(), (1, 4));
1473
1474 assert_eq!(columns.column("date_col").unwrap().data().get_value(0), Value::Date(date));
1475 assert_eq!(columns.column("datetime_col").unwrap().data().get_value(0), Value::DateTime(datetime));
1476 assert_eq!(columns.column("time_col").unwrap().data().get_value(0), Value::Time(time));
1477 assert_eq!(columns.column("interval_col").unwrap().data().get_value(0), Value::Duration(duration));
1478 }
1479
1480 #[test]
1481 fn test_single_row_mixed_types() {
1482 let date = Date::from_ymd(2025, 7, 15).unwrap();
1483 let time = Time::from_hms(9, 15, 30).unwrap();
1484
1485 let columns = Columns::single_row([
1486 ("bool_col", Value::Boolean(true)),
1487 ("int_col", Value::Int4(42)),
1488 ("str_col", Value::Utf8("hello".to_string())),
1489 ("date_col", Value::Date(date.clone())),
1490 ("time_col", Value::Time(time.clone())),
1491 ("none_col", Value::none()),
1492 ]);
1493
1494 assert_eq!(columns.len(), 6);
1495 assert_eq!(columns.shape(), (1, 6));
1496
1497 assert_eq!(columns.column("bool_col").unwrap().data().get_value(0), Value::Boolean(true));
1498 assert_eq!(columns.column("int_col").unwrap().data().get_value(0), Value::Int4(42));
1499 assert_eq!(columns.column("str_col").unwrap().data().get_value(0), Value::Utf8("hello".to_string()));
1500 assert_eq!(columns.column("date_col").unwrap().data().get_value(0), Value::Date(date));
1501 assert_eq!(columns.column("time_col").unwrap().data().get_value(0), Value::Time(time));
1502 assert_eq!(columns.column("none_col").unwrap().data().get_value(0), Value::none());
1503 }
1504
1505 #[test]
1511 fn test_single_row_none_of_int4_is_int4_typed() {
1512 let columns = Columns::single_row([("n", Value::none_of(ValueType::Int4))]);
1513 match columns.column("n").unwrap().data().get_value(0) {
1514 Value::None {
1515 inner,
1516 } => assert_eq!(inner, ValueType::Int4),
1517 other => panic!("expected Value::None, got {other:?}"),
1518 }
1519 }
1520
1521 #[test]
1522 fn test_single_row_none_of_utf8_is_utf8_typed() {
1523 let columns = Columns::single_row([("n", Value::none_of(ValueType::Utf8))]);
1524 match columns.column("n").unwrap().data().get_value(0) {
1525 Value::None {
1526 inner,
1527 } => assert_eq!(inner, ValueType::Utf8),
1528 other => panic!("expected Value::None, got {other:?}"),
1529 }
1530 }
1531
1532 #[test]
1533 fn test_single_row_bare_none_is_any_typed() {
1534 let columns = Columns::single_row([("n", Value::none())]);
1535 match columns.column("n").unwrap().data().get_value(0) {
1536 Value::None {
1537 inner,
1538 } => assert_eq!(inner, ValueType::Any),
1539 other => panic!("expected Value::None, got {other:?}"),
1540 }
1541 }
1542
1543 #[test]
1544 fn test_single_row_none_of_nested_option_collapses_to_base_type() {
1545 let inner_ty = ValueType::Option(Box::new(ValueType::Duration));
1549 let columns = Columns::single_row([("n", Value::none_of(inner_ty))]);
1550 match columns.column("n").unwrap().data().get_value(0) {
1551 Value::None {
1552 inner,
1553 } => assert_eq!(inner, ValueType::Duration),
1554 other => panic!("expected Value::None, got {other:?}"),
1555 }
1556 }
1557
1558 #[test]
1559 fn test_single_row_none_of_boolean_is_boolean_typed() {
1560 let columns = Columns::single_row([("n", Value::none_of(ValueType::Boolean))]);
1563 match columns.column("n").unwrap().data().get_value(0) {
1564 Value::None {
1565 inner,
1566 } => assert_eq!(inner, ValueType::Boolean),
1567 other => panic!("expected Value::None, got {other:?}"),
1568 }
1569 }
1570
1571 #[test]
1572 fn test_single_row_normal_column_names_work() {
1573 let columns = Columns::single_row([("normal_column", Value::Int4(42))]);
1574 assert_eq!(columns.len(), 1);
1575 assert_eq!(columns.column("normal_column").unwrap().data().get_value(0), Value::Int4(42));
1576 }
1577
1578 #[test]
1579 fn push_rows_matches_sequential_push_row_for_multiple_rows() {
1580 let shape = RowShape::new(vec![
1585 RowShapeField::unconstrained("id".to_string(), ValueType::Int4),
1586 RowShapeField::unconstrained("label".to_string(), ValueType::Utf8),
1587 ]);
1588
1589 let make = |number: u64, id: i32, label: &str| {
1590 let mut encoded = shape.allocate();
1591 shape.set_values(&mut encoded, &[Value::Int4(id), Value::Utf8(label.to_string())]);
1592 Row {
1593 number: RowNumber::from(number),
1594 encoded,
1595 shape: shape.clone(),
1596 }
1597 };
1598
1599 let rows = vec![make(1, 10, "a"), make(2, 20, "bb"), make(3, 30, "ccc")];
1600
1601 let mut sequential = Columns::empty();
1602 for row in &rows {
1603 sequential.push_row(row);
1604 }
1605
1606 let mut bulk = Columns::empty();
1607 bulk.push_rows(&rows);
1608
1609 assert_eq!(bulk.row_count(), 3);
1610 assert_eq!(bulk.len(), sequential.len());
1611 assert_eq!(bulk.row_count(), sequential.row_count());
1612
1613 for i in 0..sequential.len() {
1614 assert_eq!(bulk.name_at(i), sequential.name_at(i), "column {i} name diverged");
1615 assert_eq!(
1616 bulk.data_at(i).get_type(),
1617 sequential.data_at(i).get_type(),
1618 "column {i} type diverged"
1619 );
1620 }
1621 for r in 0..sequential.row_count() {
1622 assert_eq!(bulk.get_row(r), sequential.get_row(r), "row {r} values diverged");
1623 }
1624
1625 let bulk_numbers: Vec<RowNumber> = bulk.row_numbers.iter().cloned().collect();
1626 let seq_numbers: Vec<RowNumber> = sequential.row_numbers.iter().cloned().collect();
1627 assert_eq!(bulk_numbers, seq_numbers, "row numbers diverged");
1628 }
1629
1630 #[test]
1631 fn push_rows_on_empty_slice_is_a_noop() {
1632 let mut columns = Columns::empty();
1633 columns.push_rows(&[]);
1634 assert!(columns.is_empty());
1635 assert_eq!(columns.row_count(), 0);
1636 }
1637
1638 #[test]
1639 fn push_rows_preserves_values_when_first_row_is_none_on_option_field() {
1640 let shape = RowShape::new(vec![
1650 RowShapeField::unconstrained("id".to_string(), ValueType::Int4),
1651 RowShapeField::unconstrained(
1652 "opt_val".to_string(),
1653 ValueType::Option(Box::new(ValueType::Float8)),
1654 ),
1655 ]);
1656
1657 let make = |number: u64, id: i32, opt: Value| {
1658 let mut encoded = shape.allocate();
1659 shape.set_values(&mut encoded, &[Value::Int4(id), opt]);
1660 Row {
1661 number: RowNumber::from(number),
1662 encoded,
1663 shape: shape.clone(),
1664 }
1665 };
1666
1667 let v = Value::Float8(OrderedF64::try_from(3.0).unwrap());
1668 let rows = vec![make(1, 1, Value::none()), make(2, 2, v.clone()), make(3, 3, Value::none())];
1669
1670 let mut sequential = Columns::empty();
1671 for row in &rows {
1672 sequential.push_row(row);
1673 }
1674
1675 let mut bulk = Columns::empty();
1676 bulk.push_rows(&rows);
1677
1678 assert_eq!(bulk.row_count(), 3);
1679 assert_eq!(bulk.row_count(), sequential.row_count());
1680 for r in 0..sequential.row_count() {
1681 assert_eq!(bulk.get_row(r), sequential.get_row(r), "row {r} values diverged");
1682 }
1683
1684 let opt_col = bulk.column("opt_val").unwrap();
1686 assert_eq!(opt_col.data().get_value(0), Value::none_of(ValueType::Float8));
1687 assert_eq!(opt_col.data().get_value(1), v);
1688 assert_eq!(opt_col.data().get_value(2), Value::none_of(ValueType::Float8));
1689 assert_eq!(opt_col.data().len(), 3, "Option column has more entries than rows pushed");
1690 }
1691}