Skip to content

Commit 758afbd

Browse files
committed
fix(catalog): stream remaining segment-row scans in bounded batches
Compaction's segment listing, recovery's catalog reconciliation, rekey's segment/progress walks, and the reader-pin check each collected every catalog segment row into memory (or scanned it per operation) before this, so their resident cost and, for rekey's per-segment lookup, its running time scaled with how many segments an embedder had linked. Route them all through the existing bounded prefix-batch cursor instead, resuming each walk on the successor of the last key read.
1 parent 9eef246 commit 758afbd

5 files changed

Lines changed: 539 additions & 282 deletions

File tree

src/compaction/full.rs

Lines changed: 95 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -7,14 +7,15 @@
77
//! 3. Repacks segment files whose garbage ratio exceeds 5%.
88
99
use crate::Result;
10+
use crate::catalog::codec::{Catalog, CatalogRowKind, SegmentMeta};
1011
use crate::errors::PagedbError;
1112
use crate::segment::reader::SegmentReader;
1213
use crate::segment::types::SegmentPageKind;
1314
use crate::segment::writer::SegmentWriter;
14-
use crate::txn::db::Db;
15+
use crate::txn::db::{Db, WriterState};
1516
use crate::vfs::{Vfs, VfsFile};
1617

17-
use super::helpers::{find_segment_name_inner, list_all_segments_inner, replace_segment_compact};
18+
use super::helpers::{SEGMENT_ROW_BATCH, replace_segment_compact, segment_rows_from};
1819
use super::types::CompactStats;
1920

2021
/// Full online compaction. See module-level docs for the staged flow.
@@ -81,66 +82,104 @@ async fn compact_now_inner<V: Vfs + Clone>(db: &Db<V>) -> Result<CompactStats> {
8182
}
8283

8384
// ── 3. Repack segments ────────────────────────────────────────────────────
84-
let all_segments = list_all_segments_inner(&db.pager, db.realm_id, &state).await?;
85-
for meta in all_segments {
86-
let live = crate::segment::writer::live_path(&meta.segment_id);
87-
let file_size = db
88-
.vfs
89-
.open(&live, crate::vfs::types::OpenMode::Read)
90-
.await?
91-
.len()
92-
.await?;
93-
// Skip segments with < 5% garbage.
94-
let threshold = meta.total_bytes.saturating_add(meta.total_bytes / 20);
95-
if file_size <= threshold {
96-
continue;
85+
// The catalog is streamed in bounded batches rather than listed: a repack
86+
// must not size an allocation by how many segments the embedder has linked.
87+
// Each replacement rewrites the row under its own key, so the cursor stays
88+
// valid across the mutation and resumes on `key ‖ 0x00`.
89+
let segment_prefix = [CatalogRowKind::Segment as u8];
90+
let mut cursor: Vec<u8> = segment_prefix.to_vec();
91+
loop {
92+
let batch = segment_rows_from(&db.pager, db.realm_id, &state, &cursor).await?;
93+
let Some((last_key, _)) = batch.last() else {
94+
break;
95+
};
96+
cursor.clear();
97+
cursor.extend_from_slice(last_key);
98+
cursor.push(0);
99+
let exhausted = batch.len() < SEGMENT_ROW_BATCH;
100+
101+
for (key, meta) in batch {
102+
repack_one_segment(db, &mut state, &visibility_guard, &key, &meta, &mut result).await?;
97103
}
98104

99-
let mmap_limit = u64::try_from(db.options.mmap_view_scratch_bytes).unwrap_or(u64::MAX);
100-
let reader = SegmentReader::open_internal(
101-
db.pager.clone(),
102-
meta.clone(),
103-
db.mmap_bytes_in_use.clone(),
104-
mmap_limit,
105-
)
106-
.await?;
107-
db.vfs.mkdir_all("seg/.staging").await?;
108-
let new_segment_id = crate::crypto::random::segment_id()?;
109-
let mut writer = SegmentWriter::create_internal(
110-
db.pager.clone(),
111-
meta.realm_id,
112-
new_segment_id,
113-
db.file_id,
114-
meta.segment_kind,
115-
)
105+
if exhausted {
106+
break;
107+
}
108+
}
109+
110+
Ok(result)
111+
}
112+
113+
/// Repack one catalog-linked segment when its file carries more than 5% garbage,
114+
/// republishing the replacement under the row's own name.
115+
async fn repack_one_segment<V: Vfs + Clone>(
116+
db: &Db<V>,
117+
state: &mut WriterState,
118+
visibility_guard: &tokio::sync::RwLockWriteGuard<'_, ()>,
119+
key: &[u8],
120+
meta: &SegmentMeta,
121+
result: &mut CompactStats,
122+
) -> Result<()> {
123+
let live = crate::segment::writer::live_path(&meta.segment_id);
124+
let file_size = db
125+
.vfs
126+
.open(&live, crate::vfs::types::OpenMode::Read)
127+
.await?
128+
.len()
116129
.await?;
130+
// Skip segments with < 5% garbage.
131+
let threshold = meta.total_bytes.saturating_add(meta.total_bytes / 20);
132+
if file_size <= threshold {
133+
return Ok(());
134+
}
135+
136+
// The row key is `[kind] || realm_id || name`, so the name the replacement
137+
// must be relinked under comes straight from the row that carried this
138+
// meta. Validating it against the meta rejects a malformed key before the
139+
// repack republishes anything.
140+
let seg_name = String::from_utf8_lossy(Catalog::validate_segment_key(key, meta)?).into_owned();
141+
142+
let mmap_limit = u64::try_from(db.options.mmap_view_scratch_bytes).unwrap_or(u64::MAX);
143+
let reader = SegmentReader::open_internal(
144+
db.pager.clone(),
145+
meta.clone(),
146+
db.mmap_bytes_in_use.clone(),
147+
mmap_limit,
148+
)
149+
.await?;
150+
db.vfs.mkdir_all("seg/.staging").await?;
151+
let new_segment_id = crate::crypto::random::segment_id()?;
152+
let mut writer = SegmentWriter::create_internal(
153+
db.pager.clone(),
154+
meta.realm_id,
155+
new_segment_id,
156+
db.file_id,
157+
meta.segment_kind,
158+
)
159+
.await?;
117160

118-
for page_id in 1..meta.page_count.saturating_sub(1) {
119-
match reader.read_page(page_id).await {
120-
Ok(page_bytes) => {
121-
writer
122-
.append_page(SegmentPageKind::Data, &page_bytes)
123-
.await?;
124-
}
125-
Err(PagedbError::NotFound) => {}
126-
Err(e) => return Err(e),
161+
for page_id in 1..meta.page_count.saturating_sub(1) {
162+
match reader.read_page(page_id).await {
163+
Ok(page_bytes) => {
164+
writer
165+
.append_page(SegmentPageKind::Data, &page_bytes)
166+
.await?;
127167
}
168+
Err(PagedbError::NotFound) => {}
169+
Err(e) => return Err(e),
128170
}
129-
130-
let new_meta = writer.seal().await?;
131-
let seg_name =
132-
find_segment_name_inner(&db.pager, db.realm_id, &state, &meta.segment_id).await?;
133-
replace_segment_compact(
134-
db,
135-
&mut state,
136-
&visibility_guard,
137-
&seg_name,
138-
&meta.segment_id,
139-
&new_meta,
140-
)
141-
.await?;
142-
result.segments_repacked += 1;
143171
}
144172

145-
Ok(result)
173+
let new_meta = writer.seal().await?;
174+
replace_segment_compact(
175+
db,
176+
state,
177+
visibility_guard,
178+
&seg_name,
179+
&meta.segment_id,
180+
&new_meta,
181+
)
182+
.await?;
183+
result.segments_repacked += 1;
184+
Ok(())
146185
}

src/compaction/helpers.rs

Lines changed: 29 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -47,11 +47,31 @@ pub(super) async fn collect_catalog_split<V: Vfs + Clone>(
4747
collect_all_pairs(&tree).await
4848
}
4949

50-
pub(super) async fn list_all_segments_inner<V: Vfs + Clone>(
50+
/// Catalog segment rows read per batch while compaction walks them. A row is a
51+
/// fixed-width authenticated value plus a name capped at
52+
/// `MAX_SEGMENT_NAME_LEN`, so one batch is a few hundred KiB resident however
53+
/// many segments the catalog holds.
54+
pub(super) const SEGMENT_ROW_BATCH: usize = 256;
55+
56+
/// Read one bounded batch of catalog segment rows at or after `cursor`, as
57+
/// `(row key, decoded meta)` pairs.
58+
///
59+
/// Segment repack replaces a row's value under its own key, so key order is
60+
/// stable across the mutations and a caller resumes on `key ‖ 0x00` — the exact
61+
/// successor in the key ordering. The tree is opened per call from the live
62+
/// writer state because each replacement commits a new catalog root. Holding one
63+
/// batch is what keeps compaction's resident cost fixed instead of one
64+
/// `SegmentMeta` per linked segment.
65+
///
66+
/// The row key carries the embedder name, so a caller that needs the name to
67+
/// relink a replacement already has it here and never searches the catalog by
68+
/// segment identity.
69+
pub(super) async fn segment_rows_from<V: Vfs + Clone>(
5170
pager: &Arc<crate::pager::Pager<V>>,
5271
realm_id: crate::RealmId,
5372
state: &WriterState,
54-
) -> Result<Vec<SegmentMeta>> {
73+
cursor: &[u8],
74+
) -> Result<Vec<(Vec<u8>, SegmentMeta)>> {
5575
if state.catalog_root_page_id == 0 {
5676
return Ok(Vec::new());
5777
}
@@ -62,43 +82,18 @@ pub(super) async fn list_all_segments_inner<V: Vfs + Clone>(
6282
state.next_page_id,
6383
pager.page_size(),
6484
);
65-
let start = vec![CatalogRowKind::Segment as u8];
66-
let rows = tree.scan_prefix(&start).await?;
85+
let prefix = [CatalogRowKind::Segment as u8];
86+
let rows = tree
87+
.collect_prefix_batch_from(&prefix, cursor, SEGMENT_ROW_BATCH)
88+
.await?;
6789
let mut out = Vec::with_capacity(rows.len());
68-
for (_k, v) in rows {
69-
let meta = Catalog::decode_segment_meta(&v)?;
70-
out.push(meta);
90+
for (key, value) in rows {
91+
let meta = Catalog::decode_segment_meta(&value)?;
92+
out.push((key, meta));
7193
}
7294
Ok(out)
7395
}
7496

75-
pub(super) async fn find_segment_name_inner<V: Vfs + Clone>(
76-
pager: &Arc<crate::pager::Pager<V>>,
77-
realm_id: crate::RealmId,
78-
state: &WriterState,
79-
segment_id: &[u8; 16],
80-
) -> Result<String> {
81-
if state.catalog_root_page_id == 0 {
82-
return Err(PagedbError::NotFound);
83-
}
84-
let tree = BTree::open(
85-
pager.clone(),
86-
realm_id,
87-
state.catalog_root_page_id,
88-
state.next_page_id,
89-
pager.page_size(),
90-
);
91-
let start = vec![CatalogRowKind::Segment as u8];
92-
let rows = tree.scan_prefix(&start).await?;
93-
for (k, v) in rows {
94-
let meta = Catalog::decode_segment_meta(&v)?;
95-
if meta.segment_id == *segment_id && k.len() > 17 {
96-
return Ok(String::from_utf8_lossy(&k[17..]).into_owned());
97-
}
98-
}
99-
Err(PagedbError::NotFound)
100-
}
101-
10297
pub(super) async fn replace_segment_compact<V: Vfs + Clone>(
10398
db: &Db<V>,
10499
state: &mut WriterState,

0 commit comments

Comments
 (0)