Skip to content

Commit 5635bee

Browse files
fix(storage): complete metadata I/O and surface backend errors
1 parent a6c3e18 commit 5635bee

15 files changed

Lines changed: 1388 additions & 109 deletions

File tree

src/compaction/full.rs

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -84,10 +84,12 @@ async fn compact_now_inner<V: Vfs + Clone>(db: &Db<V>) -> Result<CompactStats> {
8484
let all_segments = list_all_segments_inner(&db.pager, db.realm_id, &state).await?;
8585
for meta in all_segments {
8686
let live = crate::segment::writer::live_path(&meta.segment_id);
87-
let file_size = match db.vfs.open(&live, crate::vfs::types::OpenMode::Read).await {
88-
Ok(f) => f.len().await.unwrap_or(meta.total_bytes),
89-
Err(_) => continue,
90-
};
87+
let file_size = db
88+
.vfs
89+
.open(&live, crate::vfs::types::OpenMode::Read)
90+
.await?
91+
.len()
92+
.await?;
9193
// Skip segments with < 5% garbage.
9294
let threshold = meta.total_bytes.saturating_add(meta.total_bytes / 20);
9395
if file_size <= threshold {

src/pager/core.rs

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ use crate::pager::format::data_page::{
2323
use crate::pager::format::page_kind::PageKind;
2424
use crate::txn::db::rekey::EpochKeyring;
2525
use crate::vfs::types::{OpenMode, WriteReq};
26-
use crate::vfs::{Vfs, VfsFile};
26+
use crate::vfs::{Vfs, VfsFile, read_exact_at};
2727
use crate::{RealmId, Result};
2828
use rayon::prelude::*;
2929

@@ -741,13 +741,8 @@ impl<V: Vfs> Pager<V> {
741741
}
742742
let mut buf = vec![0u8; page_size];
743743
{
744-
let f = file_handle.lock().await;
745-
let n = f.read_at(page_offset, &mut buf).await?;
746-
if n < page_size {
747-
for b in &mut buf[n..] {
748-
*b = 0;
749-
}
750-
}
744+
let mut f = file_handle.lock().await;
745+
read_exact_at(&mut *f, page_offset, &mut buf).await?;
751746
}
752747

753748
// Extract the cipher_id and mk_epoch recorded in this specific page's

src/pager/header.rs

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ use crate::pager::format::structural_header::{
1111
MainDbHeaderFields, decode_main_db_header, encode_main_db_header,
1212
};
1313
use crate::vfs::types::OpenMode;
14-
use crate::vfs::{Vfs, VfsFile};
14+
use crate::vfs::{Vfs, VfsFile, read_exact_at, write_all_at};
1515

1616
/// Which header slot is the authoritative current header.
1717
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -57,15 +57,15 @@ pub async fn bootstrap_header<V: Vfs>(
5757
let bytes = encode_main_db_header(initial, hk, page_size)?;
5858
let mut f = vfs.open(path, OpenMode::CreateNew).await?;
5959
// Slot A at offset 0.
60-
f.write_at(0, &bytes).await?;
60+
write_all_at(&mut f, 0, &bytes).await?;
6161
// Slot B at offset page_size — write a zero-filled page so the slot is
6262
// materialised on disk. Decode of a zero buffer fails the magic check
6363
// and returns `Corruption(HeaderUnverifiable)`, which `open_header`
6464
// treats as "this slot is unverifiable; skip it."
6565
let zero = vec![0u8; page_size];
6666
let page_size_u64 = u64::try_from(page_size)
6767
.map_err(|_| PagedbError::Io(std::io::Error::other("page_size > u64")))?;
68-
f.write_at(page_size_u64, &zero).await?;
68+
write_all_at(&mut f, page_size_u64, &zero).await?;
6969
f.sync().await?;
7070
// Make the directory entry for the newly created file durable so a
7171
// power loss immediately after creation does not lose the file.
@@ -92,13 +92,13 @@ pub async fn open_header<V: Vfs>(
9292
hk: &DerivedKey,
9393
page_size: usize,
9494
) -> Result<(MainDbHeaderFields, ActiveSlot)> {
95-
let f = vfs.open(path, OpenMode::ReadWrite).await?;
95+
let mut f = vfs.open(path, OpenMode::ReadWrite).await?;
9696
let mut buf_a = vec![0u8; page_size];
9797
let mut buf_b = vec![0u8; page_size];
98-
let _ = f.read_at(0, &mut buf_a).await?;
98+
read_exact_at(&mut f, 0, &mut buf_a).await?;
9999
let page_size_u64 = u64::try_from(page_size)
100100
.map_err(|_| PagedbError::Io(std::io::Error::other("page_size > u64")))?;
101-
let _ = f.read_at(page_size_u64, &mut buf_b).await?;
101+
read_exact_at(&mut f, page_size_u64, &mut buf_b).await?;
102102
let a = decode_main_db_header(&buf_a, hk, page_size).ok();
103103
let b = decode_main_db_header(&buf_b, hk, page_size).ok();
104104
match (a, b) {
@@ -139,7 +139,7 @@ pub async fn commit_header<V: Vfs>(
139139
.ok()
140140
.map(|s| next.page_id().saturating_mul(s))
141141
.ok_or_else(|| PagedbError::Io(std::io::Error::other("offset arithmetic overflow")))?;
142-
f.write_at(offset, &bytes).await?;
142+
write_all_at(&mut f, offset, &bytes).await?;
143143
f.sync().await?;
144144
// No `sync_dir` here: a header rewrite is a data write to an existing,
145145
// already-durable inode (main.db). Architecture §883 requires `sync_dir`

src/segment/authenticated_metadata.rs

Lines changed: 19 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -11,8 +11,7 @@ use crate::pager::format::data_page::{body, extract_page_header_ids, open_data_p
1111
use crate::pager::format::page_kind::PageKind;
1212
use crate::pager::format::segment_footer::{SegmentFooterFields, decode_segment_footer};
1313
use crate::pager::format::structural_header::{SegmentHeaderFields, decode_segment_header};
14-
use crate::vfs::Vfs;
15-
use crate::vfs::VfsFile;
14+
use crate::vfs::{Vfs, VfsFile, checked_read_progress};
1615

1716
use super::types::{EXTENT_INDEX_ENTRY_LEN, ExtentIndexEntry};
1817
use super::writer::{live_path, staging_path};
@@ -85,17 +84,23 @@ pub(crate) async fn authenticate_segment_metadata<V: Vfs + Clone>(
8584
let master_key = pager.mk_for(meta.mk_epoch, cipher_id)?;
8685
let hk = derive_hk(&master_key)?;
8786
let mut header_bytes = vec![0u8; page_size];
88-
let header_read = file.read_at(0, &mut header_bytes).await?;
89-
if header_read != page_size {
90-
return Err(PagedbError::segment_geometry_invalid("header_read"));
87+
let mut header_offset = 0;
88+
let mut header_remaining = &mut header_bytes[..];
89+
while !header_remaining.is_empty() {
90+
let read = file.read_at(header_offset, header_remaining).await?;
91+
checked_read_progress(&mut header_offset, read, header_remaining.len())?;
92+
header_remaining = header_remaining.split_at_mut(read).1;
9193
}
9294
let header = decode_segment_header(&header_bytes, &hk, page_size)?;
9395
validate_header(&header, meta, parent_file_id, page_size)?;
9496

9597
let mut footer_bytes = vec![0u8; page_size];
96-
let footer_read = file.read_at(footer_offset, &mut footer_bytes).await?;
97-
if footer_read != page_size {
98-
return Err(PagedbError::segment_geometry_invalid("footer_read"));
98+
let mut next_footer_offset = footer_offset;
99+
let mut footer_remaining = &mut footer_bytes[..];
100+
while !footer_remaining.is_empty() {
101+
let read = file.read_at(next_footer_offset, footer_remaining).await?;
102+
checked_read_progress(&mut next_footer_offset, read, footer_remaining.len())?;
103+
footer_remaining = footer_remaining.split_at_mut(read).1;
99104
}
100105
let (footer, manifest) = {
101106
let mut lru = pager.dek_lru().lock();
@@ -332,8 +337,12 @@ async fn collect_and_decode_index_page<V: Vfs + Clone>(
332337
.checked_mul(page_size_u64)
333338
.ok_or_else(|| PagedbError::segment_geometry_invalid("index.offset"))?;
334339
let mut page = vec![0u8; context.page_size];
335-
if context.file.read_at(offset, &mut page).await? != context.page_size {
336-
return Err(PagedbError::segment_geometry_invalid("index.read"));
340+
let mut read_offset = offset;
341+
let mut remaining = &mut page[..];
342+
while !remaining.is_empty() {
343+
let read = context.file.read_at(read_offset, remaining).await?;
344+
checked_read_progress(&mut read_offset, read, remaining.len())?;
345+
remaining = remaining.split_at_mut(read).1;
337346
}
338347
let (cipher_id, mk_epoch) = extract_page_header_ids(&page)?;
339348
if cipher_id != context.cipher_id {

src/segment/reader.rs

Lines changed: 15 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ use crate::pager::Pager;
1414
use crate::pager::format::data_page::{body, extract_page_header_ids, open_data_page};
1515
use crate::pager::format::page_kind::PageKind;
1616
use crate::vfs::types::OpenMode;
17-
use crate::vfs::{Vfs, VfsFile};
17+
use crate::vfs::{Vfs, VfsFile, checked_read_progress};
1818

1919
use super::authenticated_metadata::{
2020
ExtentIndexDecodeContext, authenticate_segment_metadata, decode_extent_index,
@@ -67,11 +67,14 @@ impl<V: Vfs + Clone> SegmentReader<V> {
6767
) -> Result<Self> {
6868
let page_size = pager.page_size();
6969
let live = live_path(&catalog_meta.segment_id);
70-
let file = pager
71-
.vfs()
72-
.open(&live, OpenMode::Read)
73-
.await
74-
.map_err(|_| PagedbError::NotFound)?;
70+
let file = match pager.vfs().open(&live, OpenMode::Read).await {
71+
Ok(file) => file,
72+
Err(PagedbError::Io(error)) if error.kind() == std::io::ErrorKind::NotFound => {
73+
return Err(PagedbError::NotFound);
74+
}
75+
Err(PagedbError::NotFound) => return Err(PagedbError::NotFound),
76+
Err(error) => return Err(error),
77+
};
7578
Self::finish_open(
7679
pager,
7780
catalog_meta,
@@ -202,9 +205,12 @@ impl<V: Vfs + Clone> SegmentReader<V> {
202205
.checked_mul(page_size)
203206
.ok_or_else(|| PagedbError::arithmetic_overflow("segment page offset"))?;
204207
let mut buf = vec![0u8; self.page_size];
205-
let n = self.file.read_at(offset, &mut buf).await?;
206-
if n < self.page_size {
207-
return Err(PagedbError::NotFound);
208+
let mut read_offset = offset;
209+
let mut remaining = &mut buf[..];
210+
while !remaining.is_empty() {
211+
let read = self.file.read_at(read_offset, remaining).await?;
212+
checked_read_progress(&mut read_offset, read, remaining.len())?;
213+
remaining = remaining.split_at_mut(read).1;
208214
}
209215

210216
// Try each segment page kind; AAD binding rejects wrong ones.

src/segment/writer.rs

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ use crate::pager::format::segment_footer::{
1616
};
1717
use crate::pager::format::structural_header::{SegmentHeaderFields, encode_segment_header};
1818
use crate::vfs::types::OpenMode;
19-
use crate::vfs::{Vfs, VfsFile};
19+
use crate::vfs::{Vfs, VfsFile, write_all_at};
2020
use crate::{RealmId, Result};
2121
use tracing;
2222

@@ -89,7 +89,7 @@ impl<V: Vfs + Clone> SegmentWriter<V> {
8989
flags: 0,
9090
};
9191
let header_bytes = encode_segment_header(&header_fields, &hk, page_size)?;
92-
file.write_at(0, &header_bytes).await?;
92+
write_all_at(&mut file, 0, &header_bytes).await?;
9393

9494
let total_bytes = u64::try_from(page_size)
9595
.map_err(|_| PagedbError::Io(std::io::Error::other("page_size > u64")))?;
@@ -185,7 +185,7 @@ impl<V: Vfs + Clone> SegmentWriter<V> {
185185
let offset = page_id
186186
.checked_mul(self.page_size as u64)
187187
.ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?;
188-
self.file.write_at(offset, &buf).await?;
188+
write_all_at(&mut self.file, offset, &buf).await?;
189189
self.next_page_id += 1;
190190
self.total_bytes = self.total_bytes.saturating_add(self.page_size as u64);
191191
Ok(page_id)
@@ -327,7 +327,7 @@ impl<V: Vfs + Clone> SegmentWriter<V> {
327327
let offset = page_id
328328
.checked_mul(self.page_size as u64)
329329
.ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?;
330-
self.file.write_at(offset, &buf).await?;
330+
write_all_at(&mut self.file, offset, &buf).await?;
331331
self.next_page_id += 1;
332332
self.total_bytes = self.total_bytes.saturating_add(self.page_size as u64);
333333
}
@@ -362,7 +362,7 @@ impl<V: Vfs + Clone> SegmentWriter<V> {
362362
let offset = footer_page_id
363363
.checked_mul(self.page_size as u64)
364364
.ok_or_else(|| PagedbError::Io(std::io::Error::other("offset overflow")))?;
365-
self.file.write_at(offset, &footer_bytes).await?;
365+
write_all_at(&mut self.file, offset, &footer_bytes).await?;
366366
self.file.sync().await?;
367367
self.pager.vfs().sync_dir("seg/.staging").await?;
368368

src/txn/db/misc.rs

Lines changed: 8 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -112,20 +112,17 @@ impl<V: Vfs + Clone> Db<V> {
112112
let latest_commit_id = snapshot.commit_id;
113113

114114
// Durable free-list depth (chain rooted at the header's free_list_root).
115-
let free_list_pending_entries =
116-
crate::pager::freelist::read_chain(&self.pager, self.realm_id, free_list_root)
117-
.await
118-
.map_or(0, |(entries, _)| entries.len() as u64);
115+
let (free_list_entries, _) =
116+
crate::pager::freelist::read_chain(&self.pager, self.realm_id, free_list_root).await?;
117+
let free_list_pending_entries = free_list_entries.len() as u64;
119118

120119
// Main database file size.
121-
let main_db_bytes = match self
120+
let main_db_bytes = self
122121
.vfs
123122
.open(&self.main_db_path, crate::vfs::types::OpenMode::Read)
124-
.await
125-
{
126-
Ok(f) => f.len().await.unwrap_or(0),
127-
Err(_) => 0,
128-
};
123+
.await?
124+
.len()
125+
.await?;
129126

130127
// Buffer pool stats from cache.
131128
let buffer_pool_pages = { self.pager.inner.buffer_pool.lock().len() as u64 };
@@ -161,7 +158,7 @@ impl<V: Vfs + Clone> Db<V> {
161158
);
162159

163160
let seg_start = vec![CatalogRowKind::Segment as u8];
164-
let seg_rows = tree.scan_prefix(&seg_start).await.unwrap_or_default();
161+
let seg_rows = tree.scan_prefix(&seg_start).await?;
165162
let seg_count = u32::try_from(seg_rows.len()).unwrap_or(u32::MAX);
166163
let seg_bytes: u64 = seg_rows
167164
.iter()

src/txn/db/open/existing.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ use crate::options::OpenOptions;
1212
use crate::pager::header::ActiveSlot;
1313
use crate::pager::structural_header::MainDbHeaderFields;
1414
use crate::pager::{Pager, PagerConfig};
15-
use crate::vfs::{Vfs, VfsFile};
15+
use crate::vfs::{Vfs, read_exact_at};
1616
use crate::{RealmId, Result};
1717

1818
use super::super::super::mode::DbMode;
@@ -104,13 +104,13 @@ impl<V: Vfs + Clone> Db<V> {
104104
let capabilities = mode.open_capabilities();
105105
let file_mode = capabilities.main_db_open_mode();
106106
let read_only = capabilities.read_only_file_access();
107-
let f = vfs.open(&main_db_path, file_mode).await?;
107+
let mut f = vfs.open(&main_db_path, file_mode).await?;
108108
let mut buf_a = vec![0u8; page_size];
109109
let mut buf_b = vec![0u8; page_size];
110-
let _ = f.read_at(0, &mut buf_a).await?;
110+
read_exact_at(&mut f, 0, &mut buf_a).await?;
111111
let page_size_u64 = u64::try_from(page_size)
112112
.map_err(|_| PagedbError::Io(std::io::Error::other("page_size > u64")))?;
113-
let _ = f.read_at(page_size_u64, &mut buf_b).await?;
113+
read_exact_at(&mut f, page_size_u64, &mut buf_b).await?;
114114
drop(f);
115115

116116
let try_decode = |buf: &[u8]| -> Option<(MainDbHeaderFields, bool)> {

0 commit comments

Comments
 (0)