reddb_server/storage/engine/turboquant/
extent.rs1use 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}