Skip to main content

reddb_server/storage/engine/turboquant/
extent.rs

1//! Contiguous page extent storage for TurboQuant payloads.
2//!
3//! MIT notice: clean-room RedDB implementation for the turbovec-compatible
4//! TurboQuant surface; no upstream turbovec source is copied.
5
6use crate::storage::engine::pager::{ExtentId, PagerError};
7use crate::storage::engine::{Page, PageType, Pager, HEADER_SIZE, PAGE_SIZE};
8use std::sync::Arc;
9
10const PAYLOAD_BYTES_PER_PAGE: usize = PAGE_SIZE - HEADER_SIZE;
11
12pub struct TurboExtent {
13    pager: Arc<Pager>,
14    extents: Vec<ExtentId>,
15    write_offset: usize,
16}
17
18impl TurboExtent {
19    pub fn new(pager: Arc<Pager>) -> Result<Self, PagerError> {
20        let first = pager.reserve_contig_extent(64)?;
21        Ok(Self {
22            pager,
23            extents: vec![first],
24            write_offset: 0,
25        })
26    }
27
28    pub fn append(&mut self, bytes: &[u8]) -> Result<usize, PagerError> {
29        let offset = self.write_offset;
30        self.ensure_capacity(self.write_offset + bytes.len())?;
31        let mut written = 0;
32        while written < bytes.len() {
33            let absolute = self.write_offset;
34            let (page_id, page_offset) = self.locate(absolute).ok_or_else(|| {
35                PagerError::InvalidDatabase(
36                    "turbo extent offset outside reserved pages".to_string(),
37                )
38            })?;
39            let chunk_len = (PAYLOAD_BYTES_PER_PAGE - page_offset).min(bytes.len() - written);
40            let mut page = self
41                .pager
42                .read_page(page_id)
43                .unwrap_or_else(|_| Page::new(PageType::Vector, page_id));
44            page.content_mut()[page_offset..page_offset + chunk_len]
45                .copy_from_slice(&bytes[written..written + chunk_len]);
46            self.pager.write_page(page_id, page)?;
47            written += chunk_len;
48            self.write_offset += chunk_len;
49        }
50        Ok(offset)
51    }
52
53    pub fn read(&self, offset: usize, len: usize) -> Result<Vec<u8>, PagerError> {
54        let mut out = Vec::with_capacity(len);
55        while out.len() < len {
56            let absolute = offset + out.len();
57            let (page_id, page_offset) = self.locate(absolute).ok_or_else(|| {
58                PagerError::InvalidDatabase("turbo extent read outside reserved pages".to_string())
59            })?;
60            let page = self.pager.read_page(page_id)?;
61            let chunk_len = (PAYLOAD_BYTES_PER_PAGE - page_offset).min(len - out.len());
62            out.extend_from_slice(&page.content()[page_offset..page_offset + chunk_len]);
63        }
64        Ok(out)
65    }
66
67    fn ensure_capacity(&mut self, required_bytes: usize) -> Result<(), PagerError> {
68        while required_bytes > self.capacity_bytes() {
69            let next_pages = self
70                .extents
71                .last()
72                .map(|extent| extent.n_pages.saturating_mul(2))
73                .unwrap_or(64)
74                .max(1);
75            self.extents
76                .push(self.pager.reserve_contig_extent(next_pages)?);
77        }
78        Ok(())
79    }
80
81    fn capacity_bytes(&self) -> usize {
82        self.extents
83            .iter()
84            .map(|extent| extent.n_pages as usize * PAYLOAD_BYTES_PER_PAGE)
85            .sum()
86    }
87
88    fn locate(&self, mut offset: usize) -> Option<(u32, usize)> {
89        for extent in &self.extents {
90            let bytes = extent.n_pages as usize * PAYLOAD_BYTES_PER_PAGE;
91            if offset < bytes {
92                let page_delta = offset / PAYLOAD_BYTES_PER_PAGE;
93                let page_offset = offset % PAYLOAD_BYTES_PER_PAGE;
94                return Some((extent.start_page + page_delta as u32, page_offset));
95            }
96            offset -= bytes;
97        }
98        None
99    }
100}
101
102#[cfg(test)]
103mod tests {
104    use super::*;
105    use crate::storage::engine::PagerConfig;
106
107    struct TempPagerPath(std::path::PathBuf);
108
109    impl std::ops::Deref for TempPagerPath {
110        type Target = std::path::PathBuf;
111
112        fn deref(&self) -> &std::path::PathBuf {
113            &self.0
114        }
115    }
116
117    impl AsRef<std::path::Path> for TempPagerPath {
118        fn as_ref(&self) -> &std::path::Path {
119            &self.0
120        }
121    }
122
123    impl Drop for TempPagerPath {
124        fn drop(&mut self) {
125            let _ = std::fs::remove_file(&self.0);
126            for sidecar in reddb_file::layout::pager_shadow_sidecar_paths(&self.0) {
127                let _ = std::fs::remove_file(sidecar);
128            }
129        }
130    }
131
132    #[test]
133    fn turbo_extent_reads_across_page_boundaries() {
134        let path = TempPagerPath(
135            std::env::temp_dir().join(format!("reddb-turbo-extent-{}.db", std::process::id())),
136        );
137        let _ = std::fs::remove_file(&path);
138        let pager = Arc::new(Pager::open(&path, PagerConfig::default()).unwrap());
139        let mut extent = TurboExtent::new(pager).unwrap();
140        extent.write_offset = PAYLOAD_BYTES_PER_PAGE - 2;
141        extent.ensure_capacity(PAYLOAD_BYTES_PER_PAGE + 2).unwrap();
142        let offset = extent.append(&[1, 2, 3, 4]).unwrap();
143        assert_eq!(extent.read(offset, 4).unwrap(), vec![1, 2, 3, 4]);
144    }
145}