Skip to main content

reifydb_core/value/column/
columns.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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	/// Builds a one-column `Columns` from `buffer`, extracts `indices`, and asserts the extracted
998	/// column reports the right row count, keeps its value type, and reproduces the source value at
999	/// every requested index in order. This is the type-agnostic core check: it compares
1000	/// `get_value` of the extraction against `get_value` of the source so it works for every
1001	/// `ColumnBuffer` variant without hand-constructing each `Value`.
1002	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	/// Regression: the change accumulator coalesces row-keyed inserts by calling `extract_row`
1299	/// per row, and a deferred view over a dictionary-encoded column then decodes using the
1300	/// buffer's `dictionary_id`. If extraction drops that metadata the view can no longer resolve
1301	/// the dictionary and inserts are silently lost. This pins that `extract_by_indices` carries
1302	/// the `dictionary_id` through.
1303	#[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	// value_to_buffer must respect the actual `inner` type carried by a `Value::None`, not force
1506	// every None into a single hardcoded column type. `Value::PartialEq` now compares `inner`
1507	// (previously all `Value::None` compared equal regardless of type), so a wrong inner type here
1508	// would silently mistype every "all None" column.
1509
1510	#[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		// ColumnBuffer::none_typed already unwraps a nested Option(inner) type to its base type
1546		// (there is no separate column representation for Option<Option<T>>), so a value that is
1547		// itself Option<Option<Duration>>::None ends up as a Duration-typed None column.
1548		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		// Boolean is also value_to_buffer's old hardcoded default, so this case alone would not
1561		// have caught the bug; kept for symmetry with the other inner types above.
1562		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		// push_rows pre-sizes the column buffers to the row count instead of growing them
1581		// one row at a time like push_row. The CDC producer relies on the two producing an
1582		// identical Columns; this pins that invariant for the multi-row case, which is the
1583		// only case where push_rows takes its own (pre-sizing) branch.
1584		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		// Regression: push_rows used to pre-size Option columns to length=capacity (all
1641		// None), then push capacity more values on top, producing a 2x-long buffer whose
1642		// first half (the readable half) was all None. Manifested when projecting sumtype
1643		// variant fields through a deferred view: the row at index 1 of an INSERT batch
1644		// would lose its variant payload whenever row 0's value for that field was None
1645		// (because row 0 carried a different variant). The bulk and sequential paths must
1646		// produce identical Columns even when the first row's value at an Option field is
1647		// None, since field.constraint.get_type() is consulted in that case and would hit
1648		// the Option branch.
1649		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		// The defining assertion: the real value at row 1 must survive the bulk path.
1685		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}