1use crate::buffer_manager::BufferManager;
21use crate::compression::{compress_serialized_value, decompress_serialized_value, serialized_value_size};
22use crate::page::FileHandle;
23use akar_common::enums::CompressionType;
24use akar_common::types::{LogicalTypeID, PhysicalTypeID, Value};
25
26use std::sync::{Arc, Mutex};
27
28pub(crate) const TAG_NULL: u8 = 0;
32pub(crate) const TAG_BOOL: u8 = 1;
33pub(crate) const TAG_INT64: u8 = 2;
34pub(crate) const TAG_INT32: u8 = 3;
35pub(crate) const TAG_INT16: u8 = 4;
36pub(crate) const TAG_INT8: u8 = 5;
37pub(crate) const TAG_UINT64: u8 = 6;
38pub(crate) const TAG_UINT32: u8 = 7;
39pub(crate) const TAG_UINT16: u8 = 8;
40pub(crate) const TAG_UINT8: u8 = 9;
41pub(crate) const TAG_INT128: u8 = 10;
42pub(crate) const TAG_DOUBLE: u8 = 11;
43pub(crate) const TAG_FLOAT: u8 = 12;
44pub(crate) const TAG_STRING: u8 = 13;
45pub(crate) const TAG_BLOB: u8 = 14;
46pub(crate) const TAG_DATE: u8 = 15;
47pub(crate) const TAG_TIMESTAMP: u8 = 16;
48pub(crate) const TAG_TIMESTAMP_TZ: u8 = 17;
49pub(crate) const TAG_TIMESTAMP_NS: u8 = 18;
50pub(crate) const TAG_TIMESTAMP_MS: u8 = 19;
51pub(crate) const TAG_TIMESTAMP_SEC: u8 = 20;
52pub(crate) const TAG_INTERVAL: u8 = 21;
53pub(crate) const TAG_INTERNAL_ID: u8 = 22;
54pub(crate) const TAG_LIST: u8 = 23;
55pub(crate) const TAG_MAP: u8 = 24;
56pub(crate) const TAG_STRUCT: u8 = 25;
57pub(crate) const TAG_UINT128: u8 = 26;
58pub(crate) const TAG_JSON: u8 = 27;
59pub(crate) const TAG_DTIME: u8 = 28;
60pub(crate) const TAG_UNION: u8 = 29;
61
62const MAX_VALS_PER_PAGE: usize = 256;
64
65const PAGE_HEADER_SIZE: usize = 4 + MAX_VALS_PER_PAGE * 4;
68
69#[derive(Debug, Clone, Copy)]
71struct PageHeader {
72 num_values: u32,
74 offsets: [u32; MAX_VALS_PER_PAGE],
78}
79
80#[derive(Debug)]
85pub struct Column {
86 pub logical_type: LogicalTypeID,
88 pub physical_type: PhysicalTypeID,
90 pub table_id: u64,
92 pub col_idx: u32,
94 pub file_name: String,
96 pub file_handle: FileHandle,
98 pub buffer_manager: Arc<Mutex<BufferManager>>,
100 pub compression_type: CompressionType,
102 pub value_size: usize,
104 pub num_values: u64,
106 pub num_pages: u64,
108 pub page_row_offsets: Vec<u64>,
111}
112
113impl Column {
114 pub fn new(
122 logical_type: LogicalTypeID,
123 table_id: u64,
124 col_idx: u32,
125 db_path: &std::path::Path,
126 buffer_manager: Arc<Mutex<BufferManager>>,
127 page_size: usize,
128 ) -> Self {
129 Self::with_compression(
130 logical_type,
131 table_id,
132 col_idx,
133 db_path,
134 buffer_manager,
135 page_size,
136 CompressionType::Uncompressed,
137 )
138 }
139
140 pub fn with_compression(
142 logical_type: LogicalTypeID,
143 table_id: u64,
144 col_idx: u32,
145 db_path: &std::path::Path,
146 buffer_manager: Arc<Mutex<BufferManager>>,
147 page_size: usize,
148 compression_type: CompressionType,
149 ) -> Self {
150 let file_name = format!("col_{}_{}", table_id, col_idx);
151 let col_file_path = db_path.join(&file_name);
152
153 {
155 let mut bm = buffer_manager.lock().unwrap();
156 bm.register_file(&file_name, col_file_path.clone());
157 }
158
159 let fh = FileHandle::new(col_file_path, page_size)
160 .with_free_space_manager(std::sync::Arc::new(crate::free_space_manager::FreeSpaceManager::new()));
161 let physical_type = akar_common::types::physical_type_from_logical(logical_type);
162 let value_size = serialized_value_size(physical_type);
163
164 Self {
165 logical_type,
166 physical_type,
167 table_id,
168 col_idx,
169 file_name,
170 file_handle: fh,
171 buffer_manager,
172 compression_type,
173 value_size,
174 num_values: 0,
175 num_pages: 0,
176 page_row_offsets: Vec::new(),
177 }
178 }
179
180 pub fn append_value(&mut self, value: &Value) -> std::io::Result<()> {
193 let raw = Self::serialize_value(value);
194 let serialized = compress_serialized_value(self.compression_type, &raw, self.value_size);
195 if self.num_pages > 0 {
197 let last_page = self.num_pages - 1;
198 let page_data = self.read_page_data(last_page as usize)?;
199 let num_vals = u32::from_le_bytes(page_data[..4].try_into().unwrap()) as usize;
200 if num_vals >= MAX_VALS_PER_PAGE {
201 let new_page = self.allocate_new_page()?;
203 return self.write_value_to_page(new_page, &serialized);
204 }
205 let data_end = if num_vals > 0 {
207 let last_off_pos = 4 + (num_vals - 1) * 4;
208 PAGE_HEADER_SIZE
209 + u32::from_le_bytes(page_data[last_off_pos..last_off_pos + 4].try_into().unwrap()) as usize
210 } else {
211 PAGE_HEADER_SIZE
212 };
213 if data_end + serialized.len() > self.file_handle.page_size {
214 let new_page = self.allocate_new_page()?;
215 return self.write_value_to_page(new_page, &serialized);
216 }
217 }
218 let page_idx = self.ensure_page_for_write()?;
219 self.write_value_to_page(page_idx, &serialized)
220 }
221
222 pub fn get_value(&self, row_idx: u64) -> std::io::Result<Value> {
224 if row_idx >= self.num_values {
225 return Err(std::io::Error::new(
226 std::io::ErrorKind::InvalidInput,
227 format!("row index {} out of range (num_values = {})", row_idx, self.num_values),
228 ));
229 }
230 let (page_idx, _) = self.locate_row(row_idx);
231 let serialized = self.read_page_data(page_idx)?;
232 let header = self.parse_page_header(&serialized)?;
233 let local_row = (row_idx - self.page_row_offsets[page_idx]) as usize;
234 let value = self.deserialize_value_from_page(&serialized, &header, local_row)?;
235 Ok(value)
236 }
237
238 pub fn scan_values(&self, start: u64, count: u64) -> std::io::Result<Vec<Value>> {
240 if count == 0 || start >= self.num_values {
241 return Ok(Vec::new());
242 }
243 let end = (start + count).min(self.num_values);
244 let mut results = Vec::with_capacity((end - start) as usize);
245
246 for row in start..end {
247 results.push(self.get_value(row)?);
248 }
249 Ok(results)
250 }
251
252 pub fn read_value_bytes(&self, row_idx: u64) -> std::io::Result<Vec<u8>> {
254 let (page_idx, _) = self.locate_row(row_idx);
255 let serialized = self.read_page_data(page_idx)?;
256 let header = self.parse_page_header(&serialized)?;
257 let local_row = (row_idx - self.page_row_offsets[page_idx]) as usize;
258 self.extract_value_bytes(&serialized, &header, local_row)
259 }
260
261 pub fn flush(&self) -> std::io::Result<()> {
263 for i in 0..self.num_pages {
264 let mut bm = self
265 .buffer_manager
266 .lock()
267 .map_err(|e| std::io::Error::other(format!("Lock poisoned: {e}")))?;
268 bm.flush(&self.file_name, i)?;
269 }
270 Ok(())
271 }
272
273 pub fn save_metadata(&self) -> std::io::Result<()> {
286 let meta_path = self.file_handle.path.with_extension("meta");
287 let mut buf = Vec::with_capacity(64 + self.num_pages as usize * 8);
288
289 buf.extend_from_slice(b"CMET");
290 buf.extend_from_slice(&1u32.to_le_bytes()); buf.extend_from_slice(&(self.logical_type as u32).to_le_bytes());
292 buf.extend_from_slice(&self.table_id.to_le_bytes());
293 buf.extend_from_slice(&self.col_idx.to_le_bytes());
294 buf.extend_from_slice(&self.num_values.to_le_bytes());
295 buf.extend_from_slice(&self.num_pages.to_le_bytes());
296 for offset in &self.page_row_offsets {
297 buf.extend_from_slice(&offset.to_le_bytes());
298 }
299
300 std::fs::write(&meta_path, &buf)?;
301 Ok(())
302 }
303
304 pub fn load_metadata(&mut self) -> std::io::Result<bool> {
309 let meta_path = self.file_handle.path.with_extension("meta");
310 if !meta_path.exists() {
311 return Ok(false);
312 }
313
314 let data = std::fs::read(&meta_path)?;
315 if data.len() < 36 {
316 return Err(std::io::Error::new(
317 std::io::ErrorKind::InvalidData,
318 "column metadata file too small",
319 ));
320 }
321
322 if &data[0..4] != b"CMET" {
323 return Err(std::io::Error::new(
324 std::io::ErrorKind::InvalidData,
325 "invalid column metadata magic bytes",
326 ));
327 }
328
329 let mut pos = 4;
330 let _version = u32::from_le_bytes(data[pos..pos + 4].try_into().unwrap());
331 pos += 4;
332 let _logical_type = u32::from_le_bytes(data[pos..pos + 4].try_into().unwrap());
333 pos += 4;
334 let _table_id = u64::from_le_bytes(data[pos..pos + 8].try_into().unwrap());
335 pos += 8;
336 let _col_idx = u32::from_le_bytes(data[pos..pos + 4].try_into().unwrap());
337 pos += 4;
338 self.num_values = u64::from_le_bytes(data[pos..pos + 8].try_into().unwrap());
339 pos += 8;
340 self.num_pages = u64::from_le_bytes(data[pos..pos + 8].try_into().unwrap());
341 pos += 8;
342
343 let num_pages = self.num_pages as usize;
344 if data.len() < pos + num_pages * 8 {
345 return Err(std::io::Error::new(
346 std::io::ErrorKind::InvalidData,
347 "column metadata file truncated (page_row_offsets)",
348 ));
349 }
350
351 self.page_row_offsets = Vec::with_capacity(num_pages);
352 for _ in 0..num_pages {
353 self.page_row_offsets
354 .push(u64::from_le_bytes(data[pos..pos + 8].try_into().unwrap()));
355 pos += 8;
356 }
357
358 Ok(true)
359 }
360
361 pub(crate) fn serialize_value(value: &Value) -> Vec<u8> {
367 let mut buf = Vec::with_capacity(16);
368 Self::serialize_into(&mut buf, value);
369 buf
370 }
371
372 pub(crate) fn deserialize_value_bytes(data: &[u8]) -> std::io::Result<Value> {
374 let mut pos = 0usize;
375 let value = Self::deserialize_value(data, &mut pos)?;
376 Ok(value)
377 }
378
379 fn serialize_into(buf: &mut Vec<u8>, value: &Value) {
380 match value {
381 Value::Null => buf.push(TAG_NULL),
382 Value::Bool(v) => {
383 buf.push(TAG_BOOL);
384 buf.push(if *v { 1 } else { 0 });
385 }
386 Value::Int64(v) => {
387 buf.push(TAG_INT64);
388 buf.extend_from_slice(&v.to_le_bytes());
389 }
390 Value::Int32(v) => {
391 buf.push(TAG_INT32);
392 buf.extend_from_slice(&v.to_le_bytes());
393 }
394 Value::Int16(v) => {
395 buf.push(TAG_INT16);
396 buf.extend_from_slice(&v.to_le_bytes());
397 }
398 Value::Int8(v) => {
399 buf.push(TAG_INT8);
400 buf.push(*v as u8);
401 }
402 Value::UInt64(v) => {
403 buf.push(TAG_UINT64);
404 buf.extend_from_slice(&v.to_le_bytes());
405 }
406 Value::UInt32(v) => {
407 buf.push(TAG_UINT32);
408 buf.extend_from_slice(&v.to_le_bytes());
409 }
410 Value::UInt16(v) => {
411 buf.push(TAG_UINT16);
412 buf.extend_from_slice(&v.to_le_bytes());
413 }
414 Value::UInt8(v) => {
415 buf.push(TAG_UINT8);
416 buf.push(*v);
417 }
418 Value::Int128(v) => {
419 buf.push(TAG_INT128);
420 buf.extend_from_slice(&v.to_le_bytes());
421 }
422 Value::Double(v) => {
423 buf.push(TAG_DOUBLE);
424 buf.extend_from_slice(&v.to_le_bytes());
425 }
426 Value::Float(v) => {
427 buf.push(TAG_FLOAT);
428 buf.extend_from_slice(&v.to_le_bytes());
429 }
430 Value::String(v) => {
431 buf.push(TAG_STRING);
432 buf.extend_from_slice(&(v.len() as u32).to_le_bytes());
433 buf.extend_from_slice(v.as_bytes());
434 }
435 Value::Blob(v) => {
436 buf.push(TAG_BLOB);
437 buf.extend_from_slice(&(v.len() as u32).to_le_bytes());
438 buf.extend_from_slice(v);
439 }
440 Value::Date(v) => {
441 buf.push(TAG_DATE);
442 buf.extend_from_slice(&v.0.to_le_bytes());
443 }
444 Value::Timestamp(v) | Value::TimestampNs(v) | Value::TimestampMs(v) | Value::TimestampSec(v) => {
445 let tag = match value {
446 Value::Timestamp(_) => TAG_TIMESTAMP,
447 Value::TimestampNs(_) => TAG_TIMESTAMP_NS,
448 Value::TimestampMs(_) => TAG_TIMESTAMP_MS,
449 Value::TimestampSec(_) => TAG_TIMESTAMP_SEC,
450 _ => unreachable!(),
451 };
452 buf.push(tag);
453 buf.extend_from_slice(&v.0.to_le_bytes());
454 }
455 Value::TimestampTz(v) => {
456 buf.push(TAG_TIMESTAMP_TZ);
457 buf.extend_from_slice(&v.0.to_le_bytes());
458 }
459 Value::UInt128(v) => {
460 buf.push(TAG_UINT128);
461 buf.extend_from_slice(&v.to_le_bytes());
462 }
463 Value::Json(v) => {
464 buf.push(TAG_JSON);
465 let json_str = v.to_string();
466 buf.extend_from_slice(&(json_str.len() as u32).to_le_bytes());
467 buf.extend_from_slice(json_str.as_bytes());
468 }
469 Value::DTime(v) => {
470 buf.push(TAG_DTIME);
471 buf.extend_from_slice(&v.to_le_bytes());
472 }
473 Value::Union(tag, val) => {
474 buf.push(TAG_UNION);
475 buf.extend_from_slice(&(tag.len() as u32).to_le_bytes());
476 buf.extend_from_slice(tag.as_bytes());
477 Self::serialize_into(buf, val);
478 }
479 Value::Interval(v) => {
480 buf.push(TAG_INTERVAL);
481 buf.extend_from_slice(&v.months.to_le_bytes());
482 buf.extend_from_slice(&v.days.to_le_bytes());
483 buf.extend_from_slice(&v.micros.to_le_bytes());
484 }
485 Value::InternalID(v) => {
486 buf.push(TAG_INTERNAL_ID);
487 buf.extend_from_slice(&v.table_id.to_le_bytes());
488 buf.extend_from_slice(&v.offset.to_le_bytes());
489 }
490 Value::List(v) => {
491 buf.push(TAG_LIST);
492 buf.extend_from_slice(&(v.len() as u32).to_le_bytes());
493 for elem in v {
494 Self::serialize_into(buf, elem);
495 }
496 }
497 Value::Map(v) => {
498 buf.push(TAG_MAP);
499 buf.extend_from_slice(&(v.len() as u32).to_le_bytes());
500 for (k, val) in v {
501 Self::serialize_into(buf, k);
502 Self::serialize_into(buf, val);
503 }
504 }
505 Value::Struct(v) => {
506 buf.push(TAG_STRUCT);
507 buf.extend_from_slice(&(v.len() as u32).to_le_bytes());
508 for (name, val) in v {
509 let name_bytes = name.as_bytes();
510 buf.extend_from_slice(&(name_bytes.len() as u32).to_le_bytes());
511 buf.extend_from_slice(name_bytes);
512 Self::serialize_into(buf, val);
513 }
514 }
515 }
516 }
517
518 pub(crate) fn deserialize_value(data: &[u8], pos: &mut usize) -> std::io::Result<Value> {
520 if *pos >= data.len() {
521 return Err(std::io::Error::new(
522 std::io::ErrorKind::UnexpectedEof,
523 "unexpected EOF reading value tag",
524 ));
525 }
526 let tag = data[*pos];
527 *pos += 1;
528
529 macro_rules! read_le {
530 ($ty:ty) => {{
531 let size = std::mem::size_of::<$ty>();
532 if *pos + size > data.len() {
533 return Err(std::io::Error::new(
534 std::io::ErrorKind::UnexpectedEof,
535 "unexpected EOF reading value",
536 ));
537 }
538 let mut arr = [0u8; std::mem::size_of::<$ty>()];
539 arr.copy_from_slice(&data[*pos..*pos + size]);
540 *pos += size;
541 <$ty>::from_le_bytes(arr)
542 }};
543 }
544
545 match tag {
546 TAG_NULL => Ok(Value::Null),
547 TAG_BOOL => {
548 if *pos >= data.len() {
549 return Err(std::io::Error::new(
550 std::io::ErrorKind::UnexpectedEof,
551 "eof reading bool",
552 ));
553 }
554 let v = data[*pos] != 0;
555 *pos += 1;
556 Ok(Value::Bool(v))
557 }
558 TAG_INT64 => Ok(Value::Int64(read_le!(i64))),
559 TAG_INT32 => Ok(Value::Int32(read_le!(i32))),
560 TAG_INT16 => Ok(Value::Int16(read_le!(i16))),
561 TAG_INT8 => {
562 if *pos >= data.len() {
563 return Err(std::io::Error::new(std::io::ErrorKind::UnexpectedEof, "eof reading i8"));
564 }
565 let v = data[*pos] as i8;
566 *pos += 1;
567 Ok(Value::Int8(v))
568 }
569 TAG_UINT64 => Ok(Value::UInt64(read_le!(u64))),
570 TAG_UINT32 => Ok(Value::UInt32(read_le!(u32))),
571 TAG_UINT16 => Ok(Value::UInt16(read_le!(u16))),
572 TAG_UINT8 => {
573 if *pos >= data.len() {
574 return Err(std::io::Error::new(std::io::ErrorKind::UnexpectedEof, "eof reading u8"));
575 }
576 let v = data[*pos];
577 *pos += 1;
578 Ok(Value::UInt8(v))
579 }
580 TAG_INT128 => Ok(Value::Int128(read_le!(i128))),
581 TAG_DOUBLE => Ok(Value::Double(read_le!(f64))),
582 TAG_FLOAT => Ok(Value::Float(read_le!(f32))),
583 TAG_STRING => {
584 let len = read_le!(u32) as usize;
585 if *pos + len > data.len() {
586 return Err(std::io::Error::new(
587 std::io::ErrorKind::UnexpectedEof,
588 "eof reading string data",
589 ));
590 }
591 let s = String::from_utf8_lossy(&data[*pos..*pos + len]).into_owned();
592 *pos += len;
593 Ok(Value::String(s))
594 }
595 TAG_BLOB => {
596 let len = read_le!(u32) as usize;
597 if *pos + len > data.len() {
598 return Err(std::io::Error::new(
599 std::io::ErrorKind::UnexpectedEof,
600 "eof reading blob data",
601 ));
602 }
603 let blob = data[*pos..*pos + len].to_vec();
604 *pos += len;
605 Ok(Value::Blob(blob))
606 }
607 TAG_DATE => Ok(Value::Date(akar_common::types::Date(read_le!(i32)))),
608 TAG_TIMESTAMP => Ok(Value::Timestamp(akar_common::types::Timestamp(read_le!(i64)))),
609 TAG_TIMESTAMP_TZ => Ok(Value::TimestampTz(akar_common::types::TimestampTZ(read_le!(i64)))),
610 TAG_TIMESTAMP_NS => Ok(Value::TimestampNs(akar_common::types::Timestamp(read_le!(i64)))),
611 TAG_TIMESTAMP_MS => Ok(Value::TimestampMs(akar_common::types::Timestamp(read_le!(i64)))),
612 TAG_TIMESTAMP_SEC => Ok(Value::TimestampSec(akar_common::types::Timestamp(read_le!(i64)))),
613 TAG_INTERVAL => {
614 let months = read_le!(i32);
615 let days = read_le!(i32);
616 let micros = read_le!(i64);
617 Ok(Value::Interval(akar_common::types::Interval { months, days, micros }))
618 }
619 TAG_INTERNAL_ID => {
620 let table_id = read_le!(u64);
621 let offset = read_le!(u64);
622 Ok(Value::InternalID(akar_common::types::InternalID { table_id, offset }))
623 }
624 TAG_LIST => {
625 let len = read_le!(u32) as usize;
626 let mut elems = Vec::with_capacity(len);
627 for _ in 0..len {
628 elems.push(Self::deserialize_value(data, pos)?);
629 }
630 Ok(Value::List(elems))
631 }
632 TAG_MAP => {
633 let len = read_le!(u32) as usize;
634 let mut elems = Vec::with_capacity(len);
635 for _ in 0..len {
636 let k = Self::deserialize_value(data, pos)?;
637 let v = Self::deserialize_value(data, pos)?;
638 elems.push((k, v));
639 }
640 Ok(Value::Map(elems))
641 }
642 TAG_STRUCT => {
643 let len = read_le!(u32) as usize;
644 let mut fields = Vec::with_capacity(len);
645 for _ in 0..len {
646 let name_len = read_le!(u32) as usize;
647 if *pos + name_len > data.len() {
648 return Err(std::io::Error::new(
649 std::io::ErrorKind::UnexpectedEof,
650 "eof reading struct field name",
651 ));
652 }
653 let name = String::from_utf8_lossy(&data[*pos..*pos + name_len]).into_owned();
654 *pos += name_len;
655 let val = Self::deserialize_value(data, pos)?;
656 fields.push((name, val));
657 }
658 Ok(Value::Struct(fields))
659 }
660 TAG_UINT128 => Ok(Value::UInt128(read_le!(u128))),
661 TAG_JSON => {
662 let len = read_le!(u32) as usize;
663 if *pos + len > data.len() {
664 return Err(std::io::Error::new(
665 std::io::ErrorKind::UnexpectedEof,
666 "eof reading json data",
667 ));
668 }
669 let s = String::from_utf8_lossy(&data[*pos..*pos + len]).into_owned();
670 *pos += len;
671 match serde_json::from_str(&s) {
672 Ok(v) => Ok(Value::Json(v)),
673 Err(e) => Err(std::io::Error::new(
674 std::io::ErrorKind::InvalidData,
675 format!("invalid json: {}", e),
676 )),
677 }
678 }
679 TAG_DTIME => Ok(Value::DTime(read_le!(i64))),
680 TAG_UNION => {
681 let len = read_le!(u32) as usize;
682 if *pos + len > data.len() {
683 return Err(std::io::Error::new(
684 std::io::ErrorKind::UnexpectedEof,
685 "eof reading union tag",
686 ));
687 }
688 let tag = String::from_utf8_lossy(&data[*pos..*pos + len]).into_owned();
689 *pos += len;
690 let val = Self::deserialize_value(data, pos)?;
691 Ok(Value::Union(tag, Box::new(val)))
692 }
693 _ => Err(std::io::Error::new(
694 std::io::ErrorKind::InvalidData,
695 format!("unknown value tag: 0x{:02x}", tag),
696 )),
697 }
698 }
699
700 fn write_value_to_page(&mut self, page_idx: u64, serialized: &[u8]) -> std::io::Result<()> {
715 let mut bm = self
716 .buffer_manager
717 .lock()
718 .map_err(|e| std::io::Error::other(format!("Lock poisoned: {e}")))?;
719 let frame = bm.pin_mut(&self.file_name, page_idx)?;
720 let page_size = self.file_handle.page_size;
721
722 let num_vals = u32::from_le_bytes(frame.data[..4].try_into().unwrap()) as usize;
723
724 let data_area_start = PAGE_HEADER_SIZE;
726
727 let prev_end = if num_vals > 0 {
729 let prev_off_pos = 4 + (num_vals - 1) * 4;
730 u32::from_le_bytes(frame.data[prev_off_pos..prev_off_pos + 4].try_into().unwrap()) as usize
731 } else {
732 0usize
733 };
734
735 let data_write_pos = data_area_start + prev_end;
736
737 if data_write_pos + serialized.len() > page_size {
739 return Err(std::io::Error::new(
740 std::io::ErrorKind::OutOfMemory,
741 format!(
742 "value of {} bytes does not fit in page (page_size={}, used={})",
743 serialized.len(),
744 page_size,
745 data_write_pos,
746 ),
747 ));
748 }
749
750 let new_end = (prev_end + serialized.len()) as u32;
752 let off_pos = 4 + num_vals * 4;
753 frame.data[off_pos..off_pos + 4].copy_from_slice(&new_end.to_le_bytes());
754
755 frame.data[data_write_pos..data_write_pos + serialized.len()].copy_from_slice(serialized);
757
758 frame.data[..4].copy_from_slice(&((num_vals + 1) as u32).to_le_bytes());
760 frame.mark_dirty();
761 bm.unpin(&self.file_name, page_idx);
762 drop(bm);
763
764 self.num_values += 1;
765
766 Ok(())
767 }
768
769 fn locate_row(&self, row_idx: u64) -> (usize, u64) {
771 match self.page_row_offsets.binary_search(&row_idx) {
772 Ok(i) => (i, row_idx - self.page_row_offsets[i]),
773 Err(i) => {
774 if i == 0 {
775 (0, row_idx)
776 } else {
777 (i - 1, row_idx - self.page_row_offsets[i - 1])
778 }
779 }
780 }
781 }
782
783 fn ensure_page_for_write(&mut self) -> std::io::Result<u64> {
785 if self.num_pages == 0 {
786 let page_num = self.allocate_new_page()?;
787 Ok(page_num)
788 } else {
789 Ok(self.num_pages - 1)
790 }
791 }
792
793 fn allocate_new_page(&mut self) -> std::io::Result<u64> {
795 let mut fh = self.file_handle.clone();
796 let page_num = fh.allocate_page();
797 let empty_header = vec![0u8; self.file_handle.page_size];
798 fh.write_page(page_num, &empty_header)?;
799 self.file_handle = fh;
800 self.page_row_offsets.push(self.num_values);
801 let page_idx = self.num_pages;
802 self.num_pages += 1;
803 Ok(page_idx)
804 }
805
806 fn read_page_data(&self, page_idx: usize) -> std::io::Result<Vec<u8>> {
808 let mut bm = self
809 .buffer_manager
810 .lock()
811 .map_err(|e| std::io::Error::other(format!("Lock poisoned: {e}")))?;
812 let frame = bm.pin(&self.file_name, page_idx as u64)?;
813 let data = frame.data.clone();
814 bm.unpin(&self.file_name, page_idx as u64);
815 Ok(data)
816 }
817
818 fn parse_page_header(&self, data: &[u8]) -> std::io::Result<PageHeader> {
820 if data.len() < PAGE_HEADER_SIZE {
821 return Err(std::io::Error::new(
822 std::io::ErrorKind::InvalidData,
823 format!("page too small for header: {} < {}", data.len(), PAGE_HEADER_SIZE),
824 ));
825 }
826 let num_values = u32::from_le_bytes(data[..4].try_into().unwrap());
827 let num_vals = num_values.min(MAX_VALS_PER_PAGE as u32) as usize;
828 let mut offsets = [0u32; MAX_VALS_PER_PAGE];
829 for (i, offset) in offsets.iter_mut().enumerate().take(num_vals) {
830 let off_pos = 4 + i * 4;
831 *offset = u32::from_le_bytes(data[off_pos..off_pos + 4].try_into().unwrap());
832 }
833 Ok(PageHeader { num_values, offsets })
834 }
835
836 fn extract_value_bytes(&self, data: &[u8], header: &PageHeader, local_row: usize) -> std::io::Result<Vec<u8>> {
842 if local_row >= header.num_values as usize {
843 return Err(std::io::Error::new(
844 std::io::ErrorKind::InvalidInput,
845 "local row out of range",
846 ));
847 }
848 let value_start = if local_row > 0 {
849 PAGE_HEADER_SIZE + header.offsets[local_row - 1] as usize
850 } else {
851 PAGE_HEADER_SIZE
852 };
853 let value_end = PAGE_HEADER_SIZE + header.offsets[local_row] as usize;
854 if value_start >= data.len() || value_end > data.len() {
855 return Err(std::io::Error::new(
856 std::io::ErrorKind::InvalidData,
857 "value offset out of bounds",
858 ));
859 }
860 Ok(data[value_start..value_end].to_vec())
861 }
862
863 fn deserialize_value_from_page(
868 &self,
869 data: &[u8],
870 header: &PageHeader,
871 local_row: usize,
872 ) -> std::io::Result<Value> {
873 let stored = self.extract_value_bytes(data, header, local_row)?;
874 let bytes = decompress_serialized_value(self.compression_type, &stored, self.value_size);
875 let mut pos = 0;
876 Self::deserialize_value(&bytes, &mut pos)
877 }
878}
879
880#[cfg(test)]
885mod tests {
886 use super::*;
887 use crate::page::DEFAULT_PAGE_SIZE;
888 use akar_common::enums::CompressionType;
889 use akar_common::memory::MemoryManager;
890 use akar_common::types::{Date, InternalID, Interval, Timestamp, TimestampTZ};
891
892 fn setup_column() -> (Column, tempfile::TempDir) {
893 let dir = tempfile::tempdir().unwrap();
894 let db_path = dir.path().to_path_buf();
895 let mm = Arc::new(MemoryManager::new(64 * 1024 * 1024));
896 let config = crate::buffer_manager::BufferManagerConfig::default();
897 let bm = Arc::new(Mutex::new(crate::buffer_manager::BufferManager::new(
898 db_path.clone(),
899 mm,
900 config,
901 )));
902 let col = Column::new(
903 LogicalTypeID::Int64,
904 0, 0, &db_path,
907 bm,
908 DEFAULT_PAGE_SIZE,
909 );
910 (col, dir)
911 }
912
913 #[test]
914 fn test_serialize_deserialize_primitives() {
915 let values = vec![
916 Value::Null,
917 Value::Bool(true),
918 Value::Bool(false),
919 Value::Int64(42),
920 Value::Int64(-1),
921 Value::Int32(12345),
922 Value::Int16(-32768),
923 Value::Int8(127),
924 Value::UInt64(u64::MAX),
925 Value::UInt32(99999),
926 Value::UInt16(65535),
927 Value::UInt8(255),
928 Value::Int128(i128::MAX),
929 Value::Double(3.15),
930 Value::Float(std::f32::consts::E),
931 ];
932
933 for v in &values {
934 let buf = Column::serialize_value(v);
935 let mut pos = 0;
936 let deserialized = Column::deserialize_value(&buf, &mut pos).unwrap();
937 assert_eq!(&deserialized, v, "roundtrip failed for value: {:?}", v);
938 }
939 }
940
941 #[test]
942 fn test_serialize_deserialize_string() {
943 let v = Value::String("hello world!".to_string());
944 let buf = Column::serialize_value(&v);
945 let mut pos = 0;
946 let deserialized = Column::deserialize_value(&buf, &mut pos).unwrap();
947 assert_eq!(deserialized, v);
948 }
949
950 #[test]
951 fn test_serialize_deserialize_blob() {
952 let v = Value::Blob(vec![0xDE, 0xAD, 0xBE, 0xEF]);
953 let buf = Column::serialize_value(&v);
954 let mut pos = 0;
955 let deserialized = Column::deserialize_value(&buf, &mut pos).unwrap();
956 assert_eq!(deserialized, v);
957 }
958
959 #[test]
960 fn test_serialize_deserialize_date_time() {
961 let vals = vec![
962 Value::Date(Date(12345)),
963 Value::Timestamp(Timestamp(1_700_000_000_000_000)),
964 Value::TimestampTz(TimestampTZ(1_700_000_000_000_000)),
965 Value::Interval(Interval {
966 months: 12,
967 days: 30,
968 micros: 1_000_000,
969 }),
970 Value::InternalID(InternalID {
971 table_id: 10,
972 offset: 42,
973 }),
974 ];
975 for v in &vals {
976 let buf = Column::serialize_value(v);
977 let mut pos = 0;
978 let deserialized = Column::deserialize_value(&buf, &mut pos).unwrap();
979 assert_eq!(&deserialized, v, "failed for {:?}", v);
980 }
981 }
982
983 #[test]
984 fn test_serialize_deserialize_list() {
985 let v = Value::List(vec![Value::Int64(1), Value::Int64(2), Value::Int64(3)]);
986 let buf = Column::serialize_value(&v);
987 let mut pos = 0;
988 let deserialized = Column::deserialize_value(&buf, &mut pos).unwrap();
989 assert_eq!(deserialized, v);
990 }
991
992 #[test]
993 fn test_serialize_deserialize_struct() {
994 let v = Value::Struct(vec![
995 ("name".to_string(), Value::String("Alice".to_string())),
996 ("age".to_string(), Value::Int64(30)),
997 ]);
998 let buf = Column::serialize_value(&v);
999 let mut pos = 0;
1000 let deserialized = Column::deserialize_value(&buf, &mut pos).unwrap();
1001 assert_eq!(deserialized, v);
1002 }
1003
1004 #[test]
1005 fn test_append_and_read() {
1006 let (mut col, _dir) = setup_column();
1007
1008 col.append_value(&Value::Int64(100)).unwrap();
1009 col.append_value(&Value::Int64(200)).unwrap();
1010 col.append_value(&Value::Int64(300)).unwrap();
1011
1012 assert_eq!(col.num_values, 3);
1013
1014 let v0 = col.get_value(0).unwrap();
1015 assert_eq!(v0, Value::Int64(100));
1016
1017 let v1 = col.get_value(1).unwrap();
1018 assert_eq!(v1, Value::Int64(200));
1019
1020 let v2 = col.get_value(2).unwrap();
1021 assert_eq!(v2, Value::Int64(300));
1022 }
1023
1024 #[test]
1025 fn test_scan_values() {
1026 let (mut col, _dir) = setup_column();
1027
1028 for i in 0..10 {
1029 col.append_value(&Value::Int64(i as i64)).unwrap();
1030 }
1031
1032 let scanned = col.scan_values(2, 5).unwrap();
1033 assert_eq!(scanned.len(), 5);
1034 assert_eq!(scanned[0], Value::Int64(2));
1035 assert_eq!(scanned[4], Value::Int64(6));
1036 }
1037
1038 #[test]
1039 fn test_out_of_range() {
1040 let (mut col, _dir) = setup_column();
1041 col.append_value(&Value::Int64(42)).unwrap();
1042
1043 let result = col.get_value(5);
1044 assert!(result.is_err());
1045 }
1046
1047 #[test]
1048 fn test_mixed_types() {
1049 let (mut col, _dir) = setup_column();
1050
1051 col.append_value(&Value::String("hello".to_string())).unwrap();
1052 col.append_value(&Value::Double(3.15)).unwrap();
1053 col.append_value(&Value::Bool(true)).unwrap();
1054 col.append_value(&Value::List(vec![Value::Int64(1), Value::Int64(2)]))
1055 .unwrap();
1056
1057 assert_eq!(col.get_value(0).unwrap(), Value::String("hello".to_string()));
1058 assert_eq!(col.get_value(1).unwrap(), Value::Double(3.15));
1059 assert_eq!(col.get_value(2).unwrap(), Value::Bool(true));
1060 assert_eq!(
1061 col.get_value(3).unwrap(),
1062 Value::List(vec![Value::Int64(1), Value::Int64(2)])
1063 );
1064 }
1065
1066 #[test]
1067 fn test_roundtrip_via_buffer_manager() {
1068 let (mut col, _dir) = setup_column();
1069
1070 for i in 0..50 {
1071 col.append_value(&Value::Int64(i as i64)).unwrap();
1072 }
1073
1074 col.flush().unwrap();
1075
1076 for i in 0..50 {
1077 let v = col.get_value(i as u64).unwrap();
1078 assert_eq!(v, Value::Int64(i as i64));
1079 }
1080 }
1081
1082 #[test]
1083 fn test_empty_column() {
1084 let (col, _dir) = setup_column();
1085 assert_eq!(col.num_values, 0);
1086 assert_eq!(col.num_pages, 0);
1087 }
1088
1089 fn setup_compressed_column(ctype: CompressionType) -> (Column, tempfile::TempDir) {
1092 let dir = tempfile::tempdir().unwrap();
1093 let db_path = dir.path().to_path_buf();
1094 let mm = Arc::new(MemoryManager::new(64 * 1024 * 1024));
1095 let config = crate::buffer_manager::BufferManagerConfig::default();
1096 let bm = Arc::new(Mutex::new(crate::buffer_manager::BufferManager::new(
1097 db_path.clone(),
1098 mm,
1099 config,
1100 )));
1101 let col = Column::with_compression(
1102 LogicalTypeID::Int64,
1103 0, 0, &db_path,
1106 bm,
1107 DEFAULT_PAGE_SIZE,
1108 ctype,
1109 );
1110 (col, dir)
1111 }
1112
1113 #[test]
1114 fn test_column_with_integer_bitpacking() {
1115 let (mut col, _dir) = setup_compressed_column(CompressionType::IntegerBitpacking);
1116 assert_eq!(col.compression_type, CompressionType::IntegerBitpacking);
1117
1118 for i in 0i64..50 {
1120 col.append_value(&Value::Int64(i)).unwrap();
1121 }
1122 assert_eq!(col.num_values, 50);
1123
1124 for i in 0i64..50 {
1126 let v = col.get_value(i as u64).unwrap();
1127 assert_eq!(v, Value::Int64(i));
1128 }
1129 }
1130
1131 #[test]
1132 fn test_column_with_float_compression() {
1133 let dir = tempfile::tempdir().unwrap();
1134 let db_path = dir.path().to_path_buf();
1135 let mm = Arc::new(MemoryManager::new(64 * 1024 * 1024));
1136 let config = crate::buffer_manager::BufferManagerConfig::default();
1137 let bm = Arc::new(Mutex::new(crate::buffer_manager::BufferManager::new(
1138 db_path.clone(),
1139 mm,
1140 config,
1141 )));
1142 let mut col = Column::with_compression(
1143 LogicalTypeID::Double,
1144 0,
1145 0,
1146 &db_path,
1147 bm,
1148 DEFAULT_PAGE_SIZE,
1149 CompressionType::Float,
1150 );
1151
1152 let vals = vec![1.0, 3.15, -2.5, 0.0, 1e10];
1153 for v in &vals {
1154 col.append_value(&Value::Double(*v)).unwrap();
1155 }
1156 assert_eq!(col.num_values, 5);
1157
1158 for (i, expected) in vals.iter().enumerate() {
1159 let v = col.get_value(i as u64).unwrap();
1160 match v {
1161 Value::Double(d) => assert!((d - expected).abs() < 1e-10, "mismatch at {}", i),
1162 _ => panic!("expected Double, got {:?}", v),
1163 }
1164 }
1165 }
1166
1167 #[test]
1168 fn test_column_compression_large_values() {
1169 let (mut col, _dir) = setup_compressed_column(CompressionType::IntegerBitpacking);
1171 let large: Vec<i64> = vec![i64::MAX, i64::MIN, 0, 1, -1, 1_000_000_000, -999_999_999];
1172
1173 for v in &large {
1174 col.append_value(&Value::Int64(*v)).unwrap();
1175 }
1176
1177 for (i, expected) in large.iter().enumerate() {
1178 let v = col.get_value(i as u64).unwrap();
1179 assert_eq!(v, Value::Int64(*expected), "mismatch at index {}", i);
1180 }
1181 }
1182
1183 #[test]
1184 fn test_column_compression_mixed_types() {
1185 let (mut col, _dir) = setup_compressed_column(CompressionType::IntegerBitpacking);
1188 col.append_value(&Value::Int64(42)).unwrap();
1191 col.append_value(&Value::Int64(0)).unwrap();
1192 col.append_value(&Value::Int64(-1)).unwrap();
1193 assert_eq!(col.get_value(0).unwrap(), Value::Int64(42));
1194 assert_eq!(col.get_value(1).unwrap(), Value::Int64(0));
1195 assert_eq!(col.get_value(2).unwrap(), Value::Int64(-1));
1196 }
1197}