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