Compare commits

...
5 Commits
Author SHA1 Message Date
Chris Lu 32e89b8774 rust: trim the new comments in the volume loader further 2026-08-06 15:37:40 -07:00
Chris Lu a67079dba5 rust: trim the new comments in the volume loader 2026-08-06 15:30:20 -07:00
Chris Lu 12c8df3405 rust: let a failed volume open fall back to the next candidate
Two collections can name the same volume id on one disk. Deduping the
load queue by id claimed the id for whichever candidate the scan saw
first, so a corrupt one shadowed a good one behind it; the serial loader
this replaced only claimed an id once a volume had actually opened.
Carry every claiming collection per id and try them in scan order until
one loads.
2026-08-06 15:22:36 -07:00
Chris Lu 680378e7ea rust: load a disk's volumes concurrently
Opening a volume is dominated by reading its .idx into the needle map,
and the loader did them one at a time, so a disk holding thousands of
volumes needed thousands of serial index reads before the server came
up. Go's concurrentLoadingVolumes spreads the same work over
max(cores, 10) workers; do the same, keeping the directory pre-pass and
the insert serial so only the open is parallel.
2026-08-06 13:58:55 -07:00
Chris Lu 22d492c9c7 rust: only compare the .dat tail on v3 volumes
Go's verifyNeedleIntegrity does the "does .dat end exactly at the last
indexed needle" comparison inside its v3 branch -- it rides along with
the v3 append-timestamp read -- so a v1/v2 volume carrying an unindexed
trailing record loads read-write and silent. The Rust check ran it at
every version, so booting the Rust server on a legacy cluster warned on
and quarantined volumes the Go server had been serving happily.
2026-08-06 13:57:18 -07:00
2 changed files with 186 additions and 36 deletions
+123 -22
View File
@@ -7,8 +7,8 @@
use std::collections::{HashMap, HashSet};
use std::fs;
use std::io;
use std::sync::atomic::{AtomicBool, AtomicI32, AtomicU64, Ordering};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicI32, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use tracing::warn;
@@ -122,6 +122,8 @@ impl DiskLocation {
// Scan for .dat files
let entries = fs::read_dir(&self.directory)?;
let mut dat_files: Vec<(String, VolumeId)> = Vec::new();
let mut to_load: Vec<(VolumeId, Vec<String>)> = Vec::new();
let mut queued: HashMap<VolumeId, usize> = HashMap::new();
let mut seen = HashSet::new();
for entry in entries {
@@ -201,6 +203,7 @@ impl DiskLocation {
continue;
}
// Load existing data only; never create a phantom `.dat`. A lone
// `.vif`/`.idx` (e.g. an EC sidecar whose `.ecx` is on a sibling
// disk) would otherwise have Volume::new write an 8-byte stub that
@@ -221,30 +224,22 @@ impl DiskLocation {
continue;
}
match Volume::new(
&self.directory,
&self.idx_directory,
&collection,
vid,
needle_map_kind,
None, // replica placement read from superblock
None, // TTL read from superblock
0, // no preallocate on load
Version::current(),
) {
Ok(mut v) => {
v.location_disk_space_low = self.is_disk_space_low.clone();
crate::metrics::VOLUME_GAUGE
.with_label_values(&[&collection, "volume"])
.inc();
self.volumes.insert(vid, v);
}
Err(e) => {
warn!(volume_id = vid.0, error = %e, "failed to load volume");
match queued.get(&vid) {
Some(&i) => to_load[i].1.push(collection),
None => {
queued.insert(vid, to_load.len());
to_load.push((vid, vec![collection]));
}
}
}
for (collection, vid, v) in self.open_volumes(to_load, needle_map_kind) {
crate::metrics::VOLUME_GAUGE
.with_label_values(&[&collection, "volume"])
.inc();
self.volumes.insert(vid, v);
}
// After regular volumes, auto-discover EC shards on disk so a
// fresh restart picks up shards without an explicit
// VolumeEcShardsMount RPC. Mirrors Go's loadExistingVolumes
@@ -257,6 +252,60 @@ impl DiskLocation {
Ok(())
}
/// Opening a volume is dominated by reading its .idx, so they load on
/// max(cores, 10) workers, as Go's concurrentLoadingVolumes does. An id is
/// only spoken for once one of its candidates opens.
fn open_volumes(
&self,
to_load: Vec<(VolumeId, Vec<String>)>,
needle_map_kind: NeedleMapKind,
) -> Vec<(String, VolumeId, Volume)> {
if to_load.is_empty() {
return Vec::new();
}
let workers = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1)
.max(10)
.min(to_load.len());
let next = AtomicUsize::new(0);
let opened = Mutex::new(Vec::with_capacity(to_load.len()));
std::thread::scope(|scope| {
for _ in 0..workers {
scope.spawn(|| loop {
let i = next.fetch_add(1, Ordering::Relaxed);
let Some((vid, collections)) = to_load.get(i) else {
return;
};
for collection in collections {
match Volume::new(
&self.directory,
&self.idx_directory,
collection,
*vid,
needle_map_kind,
None, // replica placement read from superblock
None, // TTL read from superblock
0, // no preallocate on load
Version::current(),
) {
Ok(mut v) => {
v.location_disk_space_low = self.is_disk_space_low.clone();
opened.lock().unwrap().push((collection.clone(), *vid, v));
break;
}
Err(e) => {
warn!(volume_id = vid.0, error = %e, "failed to load volume");
}
}
}
});
}
});
opened.into_inner().unwrap_or_else(|e| e.into_inner())
}
/// Directory pre-pass that recovers interrupted compaction commits. Collects
/// every volume id that still has a .cpc commit marker or a leftover
/// .cpd/.cpx temp file across the data and idx directories, then rolls the
@@ -1493,6 +1542,58 @@ mod tests {
assert!(ids.contains(&VolumeId(2)));
}
// Two collections can name the same volume id on one disk.
#[test]
fn test_open_volumes_falls_back_past_a_corrupt_candidate() {
let tmp = TempDir::new().unwrap();
let dir = tmp.path().to_str().unwrap();
{
let mut loc = DiskLocation::new(
dir,
dir,
10,
DiskType::HardDrive,
MinFreeSpace::Percent(1.0),
Vec::new(),
)
.unwrap();
loc.create_volume(
VolumeId(9),
"good",
NeedleMapKind::InMemory,
None,
None,
0,
Version::current(),
)
.unwrap();
loc.close();
}
let mut bad = vec![0u8; 16];
bad[0] = 9; // unsupported version
std::fs::write(format!("{}/bad_9.dat", dir), &bad).unwrap();
let loc = DiskLocation::new(
dir,
dir,
10,
DiskType::HardDrive,
MinFreeSpace::Percent(1.0),
Vec::new(),
)
.unwrap();
let opened = loc.open_volumes(
vec![(VolumeId(9), vec!["bad".to_string(), "good".to_string()])],
NeedleMapKind::InMemory,
);
assert_eq!(opened.len(), 1, "the good candidate should still open");
assert_eq!(opened[0].0, "good");
assert_eq!(opened[0].1, VolumeId(9));
}
#[test]
fn test_disk_location_delete_volume() {
let tmp = TempDir::new().unwrap();
+63 -14
View File
@@ -1970,11 +1970,12 @@ impl Volume {
}
/// Verify the live needle at the .dat tail: its header must match the
/// .idx entry and the .dat must end exactly at its on-disk end. Extra
/// bytes mean an unindexed trailing record (a torn append); appending
/// after one would place the next needle at an offset the 8-byte .idx
/// encoding may not represent, so the volume is quarantined read-only
/// instead. Mirrors Go's verifyNeedleIntegrity.
/// .idx entry and, on v3, the .dat must end exactly at its on-disk end.
/// Extra bytes mean an unindexed trailing record (a torn append);
/// appending after one would place the next needle at an offset the
/// 8-byte .idx encoding may not represent, so the volume is quarantined
/// read-only instead. Mirrors Go's verifyNeedleIntegrity, including its
/// v3-only tail check.
fn verify_needle_integrity(
&mut self,
actual_offset: i64,
@@ -2013,15 +2014,16 @@ impl Volume {
)));
}
if version == VERSION_3 {
let ts_offset =
checked_offset as u64 + NEEDLE_HEADER_SIZE as u64 + size.0 as u64 + 4; // skip checksum
let mut ts_buf = [0u8; 8];
self.read_exact_at_backend(&mut ts_buf, ts_offset)?;
let ts = u64::from_be_bytes(ts_buf);
if ts > 0 {
self.last_append_at_ns = ts;
}
if version != VERSION_3 {
return Ok(());
}
let ts_offset = checked_offset as u64 + NEEDLE_HEADER_SIZE as u64 + size.0 as u64 + 4; // skip checksum
let mut ts_buf = [0u8; 8];
self.read_exact_at_backend(&mut ts_buf, ts_offset)?;
let ts = u64::from_be_bytes(ts_buf);
if ts > 0 {
self.last_append_at_ns = ts;
}
self.verify_dat_ends_at(checked_offset + get_actual_size(size, version))
@@ -4115,6 +4117,53 @@ mod tests {
);
}
// Go's tail check is v3-only, so it serves a v1/v2 volume with an unindexed
// tail read-write; checking it here flipped whole legacy disks read-only.
#[test]
fn test_integrity_skips_dat_tail_check_before_v3() {
let tmp = TempDir::new().unwrap();
let dir = tmp.path().to_str().unwrap();
{
let mut v = Volume::new(
dir,
dir,
"",
VolumeId(1),
NeedleMapKind::InMemory,
None,
None,
0,
VERSION_2,
)
.unwrap();
for i in 1..=3u64 {
let data = format!("data {}", i);
let mut n = Needle {
id: NeedleId(i),
cookie: Cookie(i as u32),
data: data.as_bytes().to_vec(),
data_size: data.len() as u32,
..Needle::default()
};
v.write_needle(&mut n, true).unwrap();
}
v.sync_to_disk().unwrap();
}
let dat_path = volume_file_name(dir, "", VolumeId(1)) + ".dat";
let mut f = OpenOptions::new().append(true).open(&dat_path).unwrap();
f.write_all(&[0xFFu8; 96]).unwrap();
f.sync_all().unwrap();
drop(f);
let v = reload_volume(dir);
assert_eq!(v.version(), VERSION_2, "volume should reload as v2");
assert!(
!v.is_no_write_or_delete(),
"v2 volume with unindexed .dat tail must stay writable, as in Go"
);
}
// Same torn tail, but with a deletion tombstone as the last indexed
// record — the tombstone path must also verify the .dat ends at it.
#[test]