Skip to main content

akar_storage/
column.rs

1//! Columnar storage — page-based column with BufferManager-backed I/O.
2//!
3//! Each `Column` stores values of a single `LogicalType` across multiple
4//! fixed-size pages managed by the `BufferManager`. Values are serialized
5//! in a compact binary format and packed sequentially within pages.
6//!
7//! # Page Layout
8//!
9//! Each page (default 8 KiB) stores a sequence of serialized values.
10//! A `PageHeader` records how many values are in the page and the byte
11//! offset of each value, enabling direct lookup without scanning.
12//!
13//! # Value Format
14//!
15//! Every value is stored as:
16//!   1. Tag byte (the Value discriminant, 0x00–0x1B)
17//!   2. Type-specific payload (primitives as fixed-size LE bytes,
18//!      variable-length types with a u32 length prefix)
19
20use 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
28// ---------------------------------------------------------------------------
29// Tag bytes for Value discriminant (must match Value::DISCRIMINANT ordering)
30// ---------------------------------------------------------------------------
31pub(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
62/// Maximum values stored per page (keeps the offset array fixed-size).
63const MAX_VALS_PER_PAGE: usize = 256;
64
65/// Fixed-size header: [num_values: u32][end_offsets: 256×u32]
66/// Header size = 4 + 256*4 = 1028 bytes.
67const PAGE_HEADER_SIZE: usize = 4 + MAX_VALS_PER_PAGE * 4;
68
69/// Metadata stored at the start of each page.
70#[derive(Debug, Clone, Copy)]
71struct PageHeader {
72    /// Number of values stored in this page.
73    num_values: u32,
74    /// End byte offsets of each value from the start of the data area.
75    /// `end_offsets[i]` = cumulative bytes of data after value i.
76    /// So value i spans bytes `[prev_end..end_offsets[i]]` in the data area.
77    offsets: [u32; MAX_VALS_PER_PAGE],
78}
79
80/// A column stores values of a single type across multiple pages.
81///
82/// The column owns a dedicated file (named `col_{table_id}_{col_idx}`) and
83/// uses the `BufferManager` for all page-level I/O with automatic caching.
84#[derive(Debug)]
85pub struct Column {
86    /// The logical type of values stored in this column.
87    pub logical_type: LogicalTypeID,
88    /// Physical type derived from `logical_type`.
89    pub physical_type: PhysicalTypeID,
90    /// The owning table ID (used to construct the file name).
91    pub table_id: u64,
92    /// The column index within the table.
93    pub col_idx: u32,
94    /// File name used in the BufferManager's file registry.
95    pub file_name: String,
96    /// File handle for low-level page I/O.
97    pub file_handle: FileHandle,
98    /// Shared buffer manager for page caching and eviction.
99    pub buffer_manager: Arc<Mutex<BufferManager>>,
100    /// Compression algorithm applied to serialized values.
101    pub compression_type: CompressionType,
102    /// Byte size of the serialized primitive value (0 for variable-length types).
103    pub value_size: usize,
104    /// Total number of values stored.
105    pub num_values: u64,
106    /// Number of pages allocated.
107    pub num_pages: u64,
108    /// Cumulative value count per page (for binary-search lookup).
109    /// `page_row_offsets[i]` = the global row index of the first value in page i.
110    pub page_row_offsets: Vec<u64>,
111}
112
113impl Column {
114    /// Create a new column backed by a file in `db_path`.
115    ///
116    /// The file is named `col_{table_id}_{col_idx}` and registered with
117    /// the `BufferManager` automatically.
118    ///
119    /// `compression_type` determines the compression algorithm applied to
120    /// serialized values. Use `CompressionType::Uncompressed` for no compression.
121    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    /// Create a new column with a specific compression algorithm.
141    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        // Register the file with the buffer manager so it can manage its pages.
154        {
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    // ------------------------------------------------------------------
181    // Public API
182    // ------------------------------------------------------------------
183
184    /// Append a single value to the column.
185    ///
186    /// Append a single value to the column.
187    ///
188    /// Automatically allocates a new page when the current one is full or
189    /// when the current page has reached the maximum values per page (256).
190    ///
191    /// Serialized values are compressed according to `self.compression_type`.
192    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        // Check if the current page is full (too many values or no space) before trying.
196        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                // Page has hit the max values limit; allocate a new one.
202                let new_page = self.allocate_new_page()?;
203                return self.write_value_to_page(new_page, &serialized);
204            }
205            // Check if the value fits in the remaining space.
206            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    /// Get a single value by row index.
223    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    /// Scan a range of values (inclusive of `start`, exclusive of `start + count`).
239    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    /// Read the raw serialised bytes of a single value (useful for compression).
253    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    /// Flush the column's dirty pages to disk.
262    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    /// Save column metadata to a `.meta` sidecar file so it can be
274    /// reconstructed after a crash or restart without scanning every page.
275    ///
276    /// Layout (all little-endian):
277    ///   magic: 4 bytes b"CMET"
278    ///   version: u32
279    ///   logical_type: u32
280    ///   table_id: u64
281    ///   col_idx: u32
282    ///   num_values: u64
283    ///   num_pages: u64
284    ///   page_row_offsets: [u64; num_pages]
285    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()); // version
291        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    /// Load column metadata from a `.meta` sidecar file.
305    ///
306    /// Returns `Ok(true)` if metadata was loaded successfully,
307    /// `Ok(false)` if no metadata file exists (fresh column).
308    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    // ------------------------------------------------------------------
362    // Internal helpers
363    // ------------------------------------------------------------------
364
365    /// Serialise `value` into a compact byte sequence.
366    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    /// Deserialize a single value from its serialised byte sequence.
373    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    /// Deserialise a Value from a byte slice starting at the tag byte.
519    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    /// Write a serialized value to a page (append to the data area).
701    ///
702    /// Page layout (fixed-size header, never shifts):
703    ///   ```text
704    ///   [0..4)            = num_values (u32 LE)
705    ///   [4..PAGE_HEADER_SIZE) = end_offsets[256] (u32 × 256)
706    ///   [PAGE_HEADER_SIZE..)   = data area (serialised values packed sequentially)
707    ///   ```
708    /// Header size is always `PAGE_HEADER_SIZE` (= 1028 bytes), so the data
709    /// area never moves regardless of how many values are in the page.
710    ///
711    /// Value i occupies bytes `[start_off..end_off)` in the data area where:
712    ///   start_off = 0 if i == 0 else end_offsets[i-1]
713    ///   end_off   = end_offsets[i]
714    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        // Data area starts at PAGE_HEADER_SIZE (fixed, never shifts).
725        let data_area_start = PAGE_HEADER_SIZE;
726
727        // Compute the byte offset into the data area where the new value begins.
728        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        // Ensure the value fits.
738        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        // Write the end offset for this value (= prev_end + serialized.len()).
751        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        // Write value data.
756        frame.data[data_write_pos..data_write_pos + serialized.len()].copy_from_slice(serialized);
757
758        // Update num_values.
759        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    /// Determine which page a row lives in via binary search.
770    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    /// Ensure at least one page exists, allocating a new one if needed.
784    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    /// Allocate a new, empty page and update tracking metadata.
794    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    /// Read raw page data from the buffer manager.
807    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    /// Parse the page header from raw page bytes.
819    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    /// Extract raw bytes of a single value from a page.
837    ///
838    /// Data area starts at `PAGE_HEADER_SIZE` (fixed). Value i occupies
839    /// bytes `[PAGE_HEADER_SIZE + start_off .. PAGE_HEADER_SIZE + end_off)`
840    /// where `start_off = 0 if i==0 else offsets[i-1]`, `end_off = offsets[i]`.
841    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    /// Deserialize a value from a page at the given local row index.
864    ///
865    /// Decompresses the stored bytes according to `self.compression_type`
866    /// before deserializing the Value.
867    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// ---------------------------------------------------------------------------
881// Tests
882// ---------------------------------------------------------------------------
883
884#[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, // table_id
905            0, // col_idx
906            &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    // ==================== Compression integration ====================
1090
1091    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, // table_id
1104            0, // col_idx
1105            &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        // Small integers benefit from bitpacking
1119        for i in 0i64..50 {
1120            col.append_value(&Value::Int64(i)).unwrap();
1121        }
1122        assert_eq!(col.num_values, 50);
1123
1124        // Verify roundtrip
1125        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        // Large integers should still roundtrip correctly
1170        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        // String values should pass through correctly with IntegerBitpacking
1186        // (value_size=0 for strings, so compression is pass-through)
1187        let (mut col, _dir) = setup_compressed_column(CompressionType::IntegerBitpacking);
1188        // Actually, for String type, value_size=0, so it behaves like pass-through
1189        // This test uses Int64 to verify compression doesn't break data
1190        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}