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_value::{
12	Result,
13	fragment::Fragment,
14	reifydb_assertions,
15	util::cowvec::CowVec,
16	value::{Value, constraint::Constraint, datetime::DateTime, row_number::RowNumber, value_type::ValueType},
17};
18use serde::{Deserialize, Serialize};
19
20use crate::{
21	encoded::{
22		row::EncodedRow,
23		shape::{RowShape, RowShapeField},
24	},
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			..
95		} => ColumnBuffer::none_typed(ValueType::Boolean, 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	/// Builds a one-column `Columns` from `buffer`, extracts `indices`, and asserts the extracted
966	/// column reports the right row count, keeps its value type, and reproduces the source value at
967	/// every requested index in order. This is the type-agnostic core check: it compares
968	/// `get_value` of the extraction against `get_value` of the source so it works for every
969	/// `ColumnBuffer` variant without hand-constructing each `Value`.
970	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	/// Regression: the change accumulator coalesces row-keyed inserts by calling `extract_row`
1267	/// per row, and a deferred view over a dictionary-encoded column then decodes using the
1268	/// buffer's `dictionary_id`. If extraction drops that metadata the view can no longer resolve
1269	/// the dictionary and inserts are silently lost. This pins that `extract_by_indices` carries
1270	/// the `dictionary_id` through.
1271	#[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]
1474	fn test_single_row_normal_column_names_work() {
1475		let columns = Columns::single_row([("normal_column", Value::Int4(42))]);
1476		assert_eq!(columns.len(), 1);
1477		assert_eq!(columns.column("normal_column").unwrap().data().get_value(0), Value::Int4(42));
1478	}
1479
1480	#[test]
1481	fn push_rows_matches_sequential_push_row_for_multiple_rows() {
1482		// push_rows pre-sizes the column buffers to the row count instead of growing them
1483		// one row at a time like push_row. The CDC producer relies on the two producing an
1484		// identical Columns; this pins that invariant for the multi-row case, which is the
1485		// only case where push_rows takes its own (pre-sizing) branch.
1486		let shape = RowShape::new(vec![
1487			RowShapeField::unconstrained("id".to_string(), ValueType::Int4),
1488			RowShapeField::unconstrained("label".to_string(), ValueType::Utf8),
1489		]);
1490
1491		let make = |number: u64, id: i32, label: &str| {
1492			let mut encoded = shape.allocate();
1493			shape.set_values(&mut encoded, &[Value::Int4(id), Value::Utf8(label.to_string())]);
1494			Row {
1495				number: RowNumber::from(number),
1496				encoded,
1497				shape: shape.clone(),
1498			}
1499		};
1500
1501		let rows = vec![make(1, 10, "a"), make(2, 20, "bb"), make(3, 30, "ccc")];
1502
1503		let mut sequential = Columns::empty();
1504		for row in &rows {
1505			sequential.push_row(row);
1506		}
1507
1508		let mut bulk = Columns::empty();
1509		bulk.push_rows(&rows);
1510
1511		assert_eq!(bulk.row_count(), 3);
1512		assert_eq!(bulk.len(), sequential.len());
1513		assert_eq!(bulk.row_count(), sequential.row_count());
1514
1515		for i in 0..sequential.len() {
1516			assert_eq!(bulk.name_at(i), sequential.name_at(i), "column {i} name diverged");
1517			assert_eq!(
1518				bulk.data_at(i).get_type(),
1519				sequential.data_at(i).get_type(),
1520				"column {i} type diverged"
1521			);
1522		}
1523		for r in 0..sequential.row_count() {
1524			assert_eq!(bulk.get_row(r), sequential.get_row(r), "row {r} values diverged");
1525		}
1526
1527		let bulk_numbers: Vec<RowNumber> = bulk.row_numbers.iter().cloned().collect();
1528		let seq_numbers: Vec<RowNumber> = sequential.row_numbers.iter().cloned().collect();
1529		assert_eq!(bulk_numbers, seq_numbers, "row numbers diverged");
1530	}
1531
1532	#[test]
1533	fn push_rows_on_empty_slice_is_a_noop() {
1534		let mut columns = Columns::empty();
1535		columns.push_rows(&[]);
1536		assert!(columns.is_empty());
1537		assert_eq!(columns.row_count(), 0);
1538	}
1539
1540	#[test]
1541	fn push_rows_preserves_values_when_first_row_is_none_on_option_field() {
1542		// Regression: push_rows used to pre-size Option columns to length=capacity (all
1543		// None), then push capacity more values on top, producing a 2x-long buffer whose
1544		// first half (the readable half) was all None. Manifested when projecting sumtype
1545		// variant fields through a deferred view: the row at index 1 of an INSERT batch
1546		// would lose its variant payload whenever row 0's value for that field was None
1547		// (because row 0 carried a different variant). The bulk and sequential paths must
1548		// produce identical Columns even when the first row's value at an Option field is
1549		// None, since field.constraint.get_type() is consulted in that case and would hit
1550		// the Option branch.
1551		let shape = RowShape::new(vec![
1552			RowShapeField::unconstrained("id".to_string(), ValueType::Int4),
1553			RowShapeField::unconstrained(
1554				"opt_val".to_string(),
1555				ValueType::Option(Box::new(ValueType::Float8)),
1556			),
1557		]);
1558
1559		let make = |number: u64, id: i32, opt: Value| {
1560			let mut encoded = shape.allocate();
1561			shape.set_values(&mut encoded, &[Value::Int4(id), opt]);
1562			Row {
1563				number: RowNumber::from(number),
1564				encoded,
1565				shape: shape.clone(),
1566			}
1567		};
1568
1569		let v = Value::Float8(OrderedF64::try_from(3.0).unwrap());
1570		let rows = vec![make(1, 1, Value::none()), make(2, 2, v.clone()), make(3, 3, Value::none())];
1571
1572		let mut sequential = Columns::empty();
1573		for row in &rows {
1574			sequential.push_row(row);
1575		}
1576
1577		let mut bulk = Columns::empty();
1578		bulk.push_rows(&rows);
1579
1580		assert_eq!(bulk.row_count(), 3);
1581		assert_eq!(bulk.row_count(), sequential.row_count());
1582		for r in 0..sequential.row_count() {
1583			assert_eq!(bulk.get_row(r), sequential.get_row(r), "row {r} values diverged");
1584		}
1585
1586		// The defining assertion: the real value at row 1 must survive the bulk path.
1587		let opt_col = bulk.column("opt_val").unwrap();
1588		assert_eq!(opt_col.data().get_value(0), Value::none());
1589		assert_eq!(opt_col.data().get_value(1), v);
1590		assert_eq!(opt_col.data().get_value(2), Value::none());
1591		assert_eq!(opt_col.data().len(), 3, "Option column has more entries than rows pushed");
1592	}
1593}