mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-29 11:15:34 +00:00
volume server: run the vacuum compaction copy without the store lock (#11482)
* volume server: run the vacuum compaction copy without the store lock VacuumVolumeCompact held the store write lock for the whole live-needle copy, including every progress blocking_send on the 16-deep stream. On a large volume that is minutes with every read, write and heartbeat on the node parked behind it, long enough for the master to unregister the node. Split compaction the way Go's CompactByIndex runs it. A short locked step claims the volume's compacting flag, records the makeup_diff watermark (index size and compaction revision) and opens fresh .dat/.idx handles. The copy then replays .idx up to the watermark and copies from those handles with the store lock released; writes that land meanwhile are replayed by makeup_diff at commit, as before. The flag is an Arc<AtomicBool> released when the job is dropped, so every exit path clears it. Because the flag is now visible to other callers, the operations that would pull the files out from under the copy refuse while it is set: unmount (and VolumeConfigure, which unmounts and remounts), delete (checked before the volume is removed from the map, which a refused destroy used to leave unmounted), cleanup, and index relocation. A second compact and a commit stay no-ops, as in Go. The pre-copy fsync is dropped: the copy reads its own handles through the page cache and .cpd/.cpx are fsynced before commit. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume: keep a read-only in-memory index's size for the compaction copy The unlocked copy replays .idx up to index_file_size(). A read-only volume whose .sdx could not be built loads its index into memory without a writer, so that size stayed 0: the copy came out empty and the commit replaced the volume with it. CompactNeedleMap::load_from_idx now records the rows it loaded, which is also what Go's IndexFileSize reports for a read-only index. The copy's index replay now stops reading at the recorded size instead of walking rows appended since, which makeup_diff replays anyway. Adds tests for compacting a read-only volume on both the sorted index and the in-memory fallback, and for VolumeConfigure stopping when the unmount is refused during a copy. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * volume server: stop a vacuum copy as soon as its client is gone The progress callback only noticed a closed response stream when a report was due, every 128 MiB. With the copy now running outside the store lock, a copy nobody waits for keeps the volume marked compacting and so keeps refusing unmount, delete and cleanup until that next report. Check the stream on every callback. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
5c9c424a84
commit
00310f6588
@@ -1153,47 +1153,7 @@ impl VolumeServer for VolumeGrpcService {
|
||||
let (tx, rx) = tokio::sync::mpsc::channel(16);
|
||||
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let compact_start = std::time::Instant::now();
|
||||
let report_interval: i64 = 128 * 1024 * 1024;
|
||||
let next_report = std::sync::atomic::AtomicI64::new(report_interval);
|
||||
|
||||
let tx_clone = tx.clone();
|
||||
let result = {
|
||||
let mut store = state.store.write().unwrap();
|
||||
store.compact_volume(vid, preallocate, 0, |processed| {
|
||||
let target = next_report.load(std::sync::atomic::Ordering::Relaxed);
|
||||
if processed > target {
|
||||
let resp = volume_server_pb::VacuumVolumeCompactResponse {
|
||||
processed_bytes: processed,
|
||||
load_avg_1m: 0.0,
|
||||
};
|
||||
// If send fails (client disconnected), stop compaction
|
||||
if tx_clone.blocking_send(Ok(resp)).is_err() {
|
||||
return false;
|
||||
}
|
||||
next_report.store(
|
||||
processed + report_interval,
|
||||
std::sync::atomic::Ordering::Relaxed,
|
||||
);
|
||||
}
|
||||
true
|
||||
})
|
||||
};
|
||||
|
||||
let success = result.is_ok();
|
||||
crate::metrics::VACUUMING_HISTOGRAM
|
||||
.with_label_values(&["compact"])
|
||||
.observe(compact_start.elapsed().as_secs_f64());
|
||||
crate::metrics::VACUUMING_COMPACT_COUNTER
|
||||
.with_label_values(&[if success { "true" } else { "false" }])
|
||||
.inc();
|
||||
|
||||
if let Err(e) = result {
|
||||
let _ = tx.blocking_send(Err(crate::server::status_with_context(
|
||||
&format!("compact volume {vid}"),
|
||||
e,
|
||||
)));
|
||||
}
|
||||
run_vacuum_compact(&state, vid, preallocate, COMPACT_REPORT_INTERVAL, &tx);
|
||||
});
|
||||
|
||||
let stream = tokio_stream::wrappers::ReceiverStream::new(rx);
|
||||
@@ -1564,7 +1524,10 @@ impl VolumeServer for VolumeGrpcService {
|
||||
let vid = VolumeId(request.into_inner().volume_id);
|
||||
let mut store = self.state.store.write().unwrap();
|
||||
// Go returns nil when volume is not found (idempotent unmount)
|
||||
if store.unmount_volume(vid) {
|
||||
let unmounted = store
|
||||
.unmount_volume(vid)
|
||||
.map_err(|e| crate::server::status_with_context(&format!("unmount volume {vid}"), e))?;
|
||||
if unmounted {
|
||||
self.state.volume_state_notify.notify_one();
|
||||
}
|
||||
Ok(Response::new(volume_server_pb::VolumeUnmountResponse {}))
|
||||
@@ -1710,8 +1673,12 @@ impl VolumeServer for VolumeGrpcService {
|
||||
let mut store = self.state.store.write().unwrap();
|
||||
|
||||
// Unmount the volume (Go propagates unmount errors via resp.Error;
|
||||
// Rust unmount_volume returns bool, so not-found falls through to configure_volume)
|
||||
store.unmount_volume(vid);
|
||||
// not-found falls through to configure_volume)
|
||||
if let Err(e) = store.unmount_volume(vid) {
|
||||
return Ok(Response::new(volume_server_pb::VolumeConfigureResponse {
|
||||
error: format!("volume configure unmount {}: {}", vid, e),
|
||||
}));
|
||||
}
|
||||
|
||||
// Modify the super block on disk (replica_placement byte)
|
||||
if let Err(e) = store.configure_volume(vid, rp) {
|
||||
@@ -6161,6 +6128,70 @@ fn get_disk_usage(path: &str) -> (u64, u64) {
|
||||
}
|
||||
}
|
||||
|
||||
/// Bytes compacted between two `VacuumVolumeCompact` progress reports.
|
||||
const COMPACT_REPORT_INTERVAL: i64 = 128 * 1024 * 1024;
|
||||
|
||||
/// The blocking body of `vacuum_volume_compact`. The store lock is held only
|
||||
/// to start the job; the copy and its progress sends run without it.
|
||||
fn run_vacuum_compact(
|
||||
state: &VolumeServerState,
|
||||
vid: VolumeId,
|
||||
preallocate: u64,
|
||||
report_interval: i64,
|
||||
tx: &tokio::sync::mpsc::Sender<Result<volume_server_pb::VacuumVolumeCompactResponse, Status>>,
|
||||
) {
|
||||
let compact_start = std::time::Instant::now();
|
||||
let next_report = std::sync::atomic::AtomicI64::new(report_interval);
|
||||
let progress = |processed: i64| {
|
||||
// A copy nobody waits for would keep refusing unmount/delete/cleanup.
|
||||
if tx.is_closed() {
|
||||
return false;
|
||||
}
|
||||
let target = next_report.load(std::sync::atomic::Ordering::Relaxed);
|
||||
if processed > target {
|
||||
let resp = volume_server_pb::VacuumVolumeCompactResponse {
|
||||
processed_bytes: processed,
|
||||
load_avg_1m: 0.0,
|
||||
};
|
||||
// If send fails (client disconnected), stop compaction
|
||||
if tx.blocking_send(Ok(resp)).is_err() {
|
||||
return false;
|
||||
}
|
||||
next_report.store(
|
||||
processed + report_interval,
|
||||
std::sync::atomic::Ordering::Relaxed,
|
||||
);
|
||||
}
|
||||
true
|
||||
};
|
||||
|
||||
let job = state
|
||||
.store
|
||||
.write()
|
||||
.unwrap()
|
||||
.begin_compact_volume(vid, preallocate);
|
||||
let result = match job {
|
||||
Ok(Some(job)) => job.run(progress),
|
||||
Ok(None) => Ok(()), // already compacting
|
||||
Err(e) => Err(e),
|
||||
};
|
||||
|
||||
let success = result.is_ok();
|
||||
crate::metrics::VACUUMING_HISTOGRAM
|
||||
.with_label_values(&["compact"])
|
||||
.observe(compact_start.elapsed().as_secs_f64());
|
||||
crate::metrics::VACUUMING_COMPACT_COUNTER
|
||||
.with_label_values(&[if success { "true" } else { "false" }])
|
||||
.inc();
|
||||
|
||||
if let Err(e) = result {
|
||||
let _ = tx.blocking_send(Err(crate::server::status_with_context(
|
||||
&format!("compact volume {vid}"),
|
||||
e,
|
||||
)));
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -8077,6 +8108,157 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
// The compaction copy must not hold the store lock: a progress send parked
|
||||
// on a client that stops reading would otherwise stall every read, write
|
||||
// and heartbeat on the node for the rest of the copy.
|
||||
#[test]
|
||||
fn test_vacuum_compact_releases_store_lock_while_progress_send_parks() {
|
||||
let (service, _tmp) = make_local_service_with_volume("", None);
|
||||
let vid = VolumeId(1);
|
||||
let write = |store: &mut crate::storage::store::Store, id: u64, data: &[u8]| {
|
||||
let mut n = Needle {
|
||||
id: NeedleId(id),
|
||||
cookie: Cookie(id as u32),
|
||||
data_size: data.len() as u32,
|
||||
data: data.to_vec(),
|
||||
..Needle::default()
|
||||
};
|
||||
store.write_volume_needle(vid, &mut n, true).unwrap();
|
||||
};
|
||||
{
|
||||
let mut store = service.state.store.write().unwrap();
|
||||
for i in 100..140u64 {
|
||||
write(&mut store, i, format!("data-{i}").as_bytes());
|
||||
}
|
||||
}
|
||||
|
||||
// Report every needle into a one-slot channel nobody reads, so the
|
||||
// copy parks in `blocking_send` after the first report.
|
||||
let (tx, mut rx) = tokio::sync::mpsc::channel(1);
|
||||
let state = service.state.clone();
|
||||
let copy = std::thread::spawn(move || run_vacuum_compact(&state, vid, 0, 0, &tx));
|
||||
|
||||
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
|
||||
loop {
|
||||
if let Ok(store) = service.state.store.try_read()
|
||||
&& store.find_volume(vid).unwrap().1.is_compacting()
|
||||
{
|
||||
break;
|
||||
}
|
||||
if std::time::Instant::now() > deadline {
|
||||
drop(rx);
|
||||
panic!("the store lock stayed held while the compaction copy was parked");
|
||||
}
|
||||
std::thread::sleep(std::time::Duration::from_millis(10));
|
||||
}
|
||||
{
|
||||
let mut store = service
|
||||
.state
|
||||
.store
|
||||
.try_write()
|
||||
.expect("the store write lock must be free while the copy is parked");
|
||||
assert!(store.find_volume(vid).unwrap().1.is_compacting());
|
||||
write(&mut store, 999, b"late-write");
|
||||
}
|
||||
|
||||
let mut reports = 0;
|
||||
while let Some(msg) = rx.blocking_recv() {
|
||||
msg.expect("the compaction must succeed");
|
||||
reports += 1;
|
||||
}
|
||||
copy.join().unwrap();
|
||||
assert!(reports > 1, "every needle should have been reported");
|
||||
|
||||
let mut store = service.state.store.write().unwrap();
|
||||
store.commit_compact_volume(vid).unwrap();
|
||||
let mut n = Needle {
|
||||
id: NeedleId(999),
|
||||
cookie: Cookie(999),
|
||||
..Needle::default()
|
||||
};
|
||||
store.read_volume_needle(vid, &mut n).unwrap();
|
||||
assert_eq!(n.data, b"late-write");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_vacuum_compact_stops_when_the_client_is_gone_before_a_report() {
|
||||
let (service, _tmp) = make_local_service_with_volume("", None);
|
||||
let vid = VolumeId(1);
|
||||
let cpd = {
|
||||
let mut store = service.state.store.write().unwrap();
|
||||
for i in 100..110u64 {
|
||||
let data = format!("data-{i}");
|
||||
let mut n = Needle {
|
||||
id: NeedleId(i),
|
||||
cookie: Cookie(i as u32),
|
||||
data_size: data.len() as u32,
|
||||
data: data.into_bytes(),
|
||||
..Needle::default()
|
||||
};
|
||||
store.write_volume_needle(vid, &mut n, true).unwrap();
|
||||
}
|
||||
store.find_volume(vid).unwrap().1.file_name(".cpd")
|
||||
};
|
||||
|
||||
// The client is gone and no report is due for 128 MiB.
|
||||
let (tx, rx) = tokio::sync::mpsc::channel(1);
|
||||
drop(rx);
|
||||
run_vacuum_compact(&service.state, vid, 0, COMPACT_REPORT_INTERVAL, &tx);
|
||||
|
||||
assert!(
|
||||
!std::path::Path::new(&cpd).exists(),
|
||||
"the copy ran to the end after its client disconnected"
|
||||
);
|
||||
let store = service.state.store.read().unwrap();
|
||||
assert!(!store.find_volume(vid).unwrap().1.is_compacting());
|
||||
}
|
||||
|
||||
// VolumeConfigure unmounts, rewrites the super block and remounts. While
|
||||
// a copy is in flight the unmount is refused, and configure must stop
|
||||
// there rather than rewrite the .dat the copy is reading.
|
||||
#[tokio::test]
|
||||
async fn test_volume_configure_refused_while_compacting() {
|
||||
let (service, _tmp) = make_local_service_with_volume("", None);
|
||||
let vid = VolumeId(1);
|
||||
let job = service
|
||||
.state
|
||||
.store
|
||||
.write()
|
||||
.unwrap()
|
||||
.begin_compact_volume(vid, 0)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
|
||||
let resp = service
|
||||
.volume_configure(Request::new(volume_server_pb::VolumeConfigureRequest {
|
||||
volume_id: vid.0,
|
||||
replication: "001".to_string(),
|
||||
}))
|
||||
.await
|
||||
.unwrap()
|
||||
.into_inner();
|
||||
assert!(resp.error.contains("is compacting"), "{}", resp.error);
|
||||
{
|
||||
let store = service.state.store.read().unwrap();
|
||||
let (_, v) = store.find_volume(vid).expect("still mounted");
|
||||
assert_eq!(v.super_block.replica_placement.to_string(), "000");
|
||||
}
|
||||
|
||||
job.run(|_| true).unwrap();
|
||||
let resp = service
|
||||
.volume_configure(Request::new(volume_server_pb::VolumeConfigureRequest {
|
||||
volume_id: vid.0,
|
||||
replication: "001".to_string(),
|
||||
}))
|
||||
.await
|
||||
.unwrap()
|
||||
.into_inner();
|
||||
assert_eq!(resp.error, "");
|
||||
let store = service.state.store.read().unwrap();
|
||||
let (_, v) = store.find_volume(vid).unwrap();
|
||||
assert_eq!(v.super_block.replica_placement.to_string(), "001");
|
||||
}
|
||||
|
||||
// Regression test for comparing the wrong compaction-revision field.
|
||||
// last_compact_revision() is bookkeeping recorded just before a compaction
|
||||
// starts (for makeup-diff catch-up) and is intentionally left behind
|
||||
@@ -8561,7 +8743,8 @@ mod tests {
|
||||
.store
|
||||
.write()
|
||||
.unwrap()
|
||||
.unmount_volume(VolumeId(1)),
|
||||
.unmount_volume(VolumeId(1))
|
||||
.unwrap(),
|
||||
"the volume must still be mounted when the lookup ran"
|
||||
);
|
||||
drop(park);
|
||||
|
||||
@@ -1906,7 +1906,7 @@ mod tests {
|
||||
1.0
|
||||
);
|
||||
|
||||
assert!(store.unmount_volume(VolumeId(21)));
|
||||
assert!(store.unmount_volume(VolumeId(21)).unwrap());
|
||||
build_heartbeat(&test_config(), &mut store);
|
||||
|
||||
assert_eq!(
|
||||
|
||||
@@ -601,6 +601,12 @@ impl DiskLocation {
|
||||
only_garbage: bool,
|
||||
keep_remote_data: bool,
|
||||
) -> Result<(), VolumeError> {
|
||||
// Refuse before removing: a refused destroy must leave it mounted.
|
||||
if let Some(v) = self.volumes.get(&vid)
|
||||
&& v.is_compacting()
|
||||
{
|
||||
return Err(v.compacting_error());
|
||||
}
|
||||
if let Some(mut v) = self.volumes.remove(&vid) {
|
||||
crate::metrics::VOLUME_GAUGE
|
||||
.with_label_values(&[&v.collection, "volume"])
|
||||
|
||||
@@ -247,6 +247,8 @@ impl CompactNeedleMap {
|
||||
pub fn load_from_idx<R: Read + Seek>(reader: &mut R, version: Version) -> io::Result<Self> {
|
||||
let mut nm = CompactNeedleMap::new();
|
||||
idx::walk_index_file(reader, 0, |key, offset, size| {
|
||||
// A read-only load attaches no writer, so this is its only size.
|
||||
nm.idx_file_offset += NEEDLE_MAP_ENTRY_SIZE as u64;
|
||||
nm.metric.maybe_set_max_needle_end(offset, size, version);
|
||||
if offset.is_zero() || size.is_deleted() {
|
||||
nm.delete_from_map(key);
|
||||
|
||||
@@ -18,7 +18,7 @@ use crate::storage::needle::needle::Needle;
|
||||
use crate::storage::needle_map::NeedleMapKind;
|
||||
use crate::storage::super_block::ReplicaPlacement;
|
||||
use crate::storage::types::*;
|
||||
use crate::storage::volume::{VifVolumeInfo, VolumeError, VolumeSpec};
|
||||
use crate::storage::volume::{CompactionJob, VifVolumeInfo, VolumeError, VolumeSpec};
|
||||
|
||||
/// Top-level storage manager containing all disk locations and their volumes.
|
||||
pub struct Store {
|
||||
@@ -406,14 +406,20 @@ impl Store {
|
||||
Err(VolumeError::NotFound)
|
||||
}
|
||||
|
||||
/// Unload (unmount) a volume without deleting its files.
|
||||
pub fn unmount_volume(&mut self, vid: VolumeId) -> bool {
|
||||
/// Unload (unmount) a volume without deleting its files. Refused while
|
||||
/// compacting, since a remount could start a second copy into .cpd.
|
||||
pub fn unmount_volume(&mut self, vid: VolumeId) -> Result<bool, VolumeError> {
|
||||
if let Some((_, v)) = self.find_volume(vid)
|
||||
&& v.is_compacting()
|
||||
{
|
||||
return Err(v.compacting_error());
|
||||
}
|
||||
for loc in &mut self.locations {
|
||||
if loc.unload_volume(vid).is_some() {
|
||||
return true;
|
||||
return Ok(true);
|
||||
}
|
||||
}
|
||||
false
|
||||
Ok(false)
|
||||
}
|
||||
|
||||
/// Reports whether any local volume or EC shard is currently quarantined
|
||||
@@ -1472,12 +1478,25 @@ impl Store {
|
||||
&mut self,
|
||||
vid: VolumeId,
|
||||
preallocate: u64,
|
||||
max_bytes_per_second: i64,
|
||||
_max_bytes_per_second: i64,
|
||||
progress_fn: F,
|
||||
) -> Result<(), VolumeError>
|
||||
where
|
||||
F: Fn(i64) -> bool,
|
||||
{
|
||||
match self.begin_compact_volume(vid, preallocate)? {
|
||||
Some(job) => job.run(progress_fn),
|
||||
None => Ok(()),
|
||||
}
|
||||
}
|
||||
|
||||
/// The part of `compact_volume` that needs the store: check free space and
|
||||
/// start the compaction. The returned job runs without the store lock.
|
||||
pub(crate) fn begin_compact_volume(
|
||||
&mut self,
|
||||
vid: VolumeId,
|
||||
preallocate: u64,
|
||||
) -> Result<Option<CompactionJob>, VolumeError> {
|
||||
// Required space matches Go's CompactVolume check: the larger of the
|
||||
// requested preallocation and the estimated volume size.
|
||||
let (loc_idx, space_needed) = {
|
||||
@@ -1501,7 +1520,7 @@ impl Store {
|
||||
let (_, v) = self
|
||||
.find_volume_mut(vid)
|
||||
.ok_or(VolumeError::VolumeNotFound(vid))?;
|
||||
v.compact_by_index(preallocate, max_bytes_per_second, progress_fn)
|
||||
v.begin_compact_by_index()
|
||||
}
|
||||
|
||||
/// Commit a completed compaction: swap files and reload.
|
||||
@@ -1801,7 +1820,7 @@ mod tests {
|
||||
store
|
||||
.write_volume_needle(VolumeId(7), &mut n, false)
|
||||
.unwrap();
|
||||
assert!(store.unmount_volume(VolumeId(7)));
|
||||
assert!(store.unmount_volume(VolumeId(7)).unwrap());
|
||||
|
||||
store.mount_volume_by_id(VolumeId(7), Some("coll")).unwrap();
|
||||
assert!(store.find_volume(VolumeId(7)).is_some());
|
||||
@@ -1843,7 +1862,7 @@ mod tests {
|
||||
store
|
||||
.write_volume_needle(VolumeId(9), &mut n, false)
|
||||
.unwrap();
|
||||
assert!(store.unmount_volume(VolumeId(9)));
|
||||
assert!(store.unmount_volume(VolumeId(9)).unwrap());
|
||||
|
||||
// The hint is accepted and mounts the volume.
|
||||
store
|
||||
@@ -1888,7 +1907,7 @@ mod tests {
|
||||
store
|
||||
.write_volume_needle(VolumeId(11), &mut n, false)
|
||||
.unwrap();
|
||||
assert!(store.unmount_volume(VolumeId(11)));
|
||||
assert!(store.unmount_volume(VolumeId(11)).unwrap());
|
||||
|
||||
// Simulate an interrupted copy: drop a .note marker.
|
||||
let base = volume_file_name(dir, "coll", VolumeId(11));
|
||||
@@ -1945,7 +1964,7 @@ mod tests {
|
||||
store
|
||||
.write_volume_needle(VolumeId(13), &mut n, false)
|
||||
.unwrap();
|
||||
assert!(store.unmount_volume(VolumeId(13)));
|
||||
assert!(store.unmount_volume(VolumeId(13)).unwrap());
|
||||
|
||||
// No hint: the fallback scan finds the sidecar on disk 0 first (skip,
|
||||
// no .dat), then the real .dat on disk 1 (mount).
|
||||
@@ -2000,7 +2019,7 @@ mod tests {
|
||||
store
|
||||
.write_volume_needle(VolumeId(15), &mut n, false)
|
||||
.unwrap();
|
||||
assert!(store.unmount_volume(VolumeId(15)));
|
||||
assert!(store.unmount_volume(VolumeId(15)).unwrap());
|
||||
// Clear the low-space flag so mount_volume_by_id considers disk 0.
|
||||
store.locations[0]
|
||||
.is_disk_space_low
|
||||
@@ -2508,6 +2527,162 @@ mod tests {
|
||||
assert_eq!(v.dat_file_size().unwrap(), volume_size);
|
||||
}
|
||||
|
||||
fn write_test_needle(store: &mut Store, vid: VolumeId, id: u64, data: &[u8]) {
|
||||
let mut n = Needle {
|
||||
id: NeedleId(id),
|
||||
cookie: Cookie(id as u32),
|
||||
data_size: data.len() as u32,
|
||||
data: data.to_vec(),
|
||||
..Needle::default()
|
||||
};
|
||||
store.write_volume_needle(vid, &mut n, true).unwrap();
|
||||
}
|
||||
|
||||
fn read_test_needle(store: &Store, vid: VolumeId, id: u64) -> Result<Vec<u8>, VolumeError> {
|
||||
let mut n = Needle {
|
||||
id: NeedleId(id),
|
||||
cookie: Cookie(id as u32),
|
||||
..Needle::default()
|
||||
};
|
||||
store.read_volume_needle(vid, &mut n)?;
|
||||
Ok(n.data)
|
||||
}
|
||||
|
||||
/// While a compaction copy runs off the store lock, nothing may pull the
|
||||
/// volume's files out from under it or start a second copy into .cpd.
|
||||
#[test]
|
||||
fn test_compaction_in_flight_guards_the_volume() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let dir = tmp.path().to_str().unwrap();
|
||||
let other = TempDir::new().unwrap();
|
||||
let mut store = make_test_store(&[dir]);
|
||||
let vid = VolumeId(1);
|
||||
store
|
||||
.add_volume(vid, DiskType::HardDrive, &VolumeSpec::default())
|
||||
.unwrap();
|
||||
for i in 1..=3u64 {
|
||||
write_test_needle(&mut store, vid, i, format!("data-{i}").as_bytes());
|
||||
}
|
||||
let revision = {
|
||||
let (_, v) = store.find_volume(vid).unwrap();
|
||||
v.super_block.compaction_revision
|
||||
};
|
||||
|
||||
let job = store
|
||||
.begin_compact_volume(vid, 0)
|
||||
.unwrap()
|
||||
.expect("the first compaction claims the volume");
|
||||
|
||||
assert!(store.begin_compact_volume(vid, 0).unwrap().is_none());
|
||||
assert!(store.unmount_volume(vid).is_err());
|
||||
assert!(store.delete_volume(vid, false, false, false).is_err());
|
||||
store.delete_collection("").unwrap();
|
||||
assert!(store.cleanup_compact_volume(vid).is_err());
|
||||
let (_, v) = store.find_volume_mut(vid).unwrap();
|
||||
assert!(v.relocate_index_to(other.path().to_str().unwrap()).is_err());
|
||||
// Go parity: a commit that finds the volume compacting is a no-op.
|
||||
store.commit_compact_volume(vid).unwrap();
|
||||
|
||||
let (_, v) = store.find_volume(vid).expect("still mounted");
|
||||
assert!(v.is_compacting());
|
||||
assert_eq!(v.super_block.compaction_revision, revision);
|
||||
|
||||
job.run(|_| true).unwrap();
|
||||
assert!(!store.find_volume(vid).unwrap().1.is_compacting());
|
||||
store.commit_compact_volume(vid).unwrap();
|
||||
let (_, v) = store.find_volume(vid).unwrap();
|
||||
assert_eq!(v.super_block.compaction_revision, revision + 1);
|
||||
assert!(store.unmount_volume(vid).unwrap());
|
||||
}
|
||||
|
||||
/// Writes, overwrites and deletes that land while the copy is parked
|
||||
/// off the store lock must all survive the commit via makeup_diff.
|
||||
fn check_writes_during_compaction_survive_commit(kind: NeedleMapKind) {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let dir = tmp.path().to_str().unwrap();
|
||||
let mut store = Store::new(kind);
|
||||
store
|
||||
.add_location(
|
||||
dir,
|
||||
dir,
|
||||
10,
|
||||
DiskType::HardDrive,
|
||||
MinFreeSpace::Percent(1.0),
|
||||
Vec::new(),
|
||||
)
|
||||
.unwrap();
|
||||
let vid = VolumeId(1);
|
||||
store
|
||||
.add_volume(vid, DiskType::HardDrive, &VolumeSpec::default())
|
||||
.unwrap();
|
||||
let (_, v) = store.find_volume(vid).unwrap();
|
||||
assert_eq!(
|
||||
v.live_meta_idx_size_for_test().is_some(),
|
||||
kind == NeedleMapKind::Redb
|
||||
);
|
||||
for i in 1..=6u64 {
|
||||
write_test_needle(&mut store, vid, i, format!("data-{i}").as_bytes());
|
||||
}
|
||||
let mut del = Needle {
|
||||
id: NeedleId(2),
|
||||
cookie: Cookie(2),
|
||||
..Needle::default()
|
||||
};
|
||||
store.delete_volume_needle(vid, &mut del).unwrap();
|
||||
let revision = {
|
||||
let (_, v) = store.find_volume(vid).unwrap();
|
||||
v.super_block.compaction_revision
|
||||
};
|
||||
|
||||
let job = store.begin_compact_volume(vid, 0).unwrap().unwrap();
|
||||
let (entered_tx, entered_rx) = std::sync::mpsc::channel();
|
||||
let (release_tx, release_rx) = std::sync::mpsc::channel::<()>();
|
||||
let copy = std::thread::spawn(move || {
|
||||
job.run(move |_| {
|
||||
let _ = entered_tx.send(());
|
||||
let _ = release_rx.recv();
|
||||
true
|
||||
})
|
||||
});
|
||||
entered_rx.recv().unwrap();
|
||||
|
||||
write_test_needle(&mut store, vid, 99, b"late-write");
|
||||
write_test_needle(&mut store, vid, 3, b"overwritten");
|
||||
let mut del = Needle {
|
||||
id: NeedleId(4),
|
||||
cookie: Cookie(4),
|
||||
..Needle::default()
|
||||
};
|
||||
store.delete_volume_needle(vid, &mut del).unwrap();
|
||||
|
||||
drop(release_tx);
|
||||
copy.join().unwrap().unwrap();
|
||||
store.commit_compact_volume(vid).unwrap();
|
||||
|
||||
let (_, v) = store.find_volume(vid).unwrap();
|
||||
assert_eq!(v.super_block.compaction_revision, revision + 1);
|
||||
assert_eq!(read_test_needle(&store, vid, 99).unwrap(), b"late-write");
|
||||
assert_eq!(read_test_needle(&store, vid, 3).unwrap(), b"overwritten");
|
||||
assert!(read_test_needle(&store, vid, 4).is_err());
|
||||
assert!(read_test_needle(&store, vid, 2).is_err());
|
||||
for i in [1u64, 5, 6] {
|
||||
assert_eq!(
|
||||
read_test_needle(&store, vid, i).unwrap(),
|
||||
format!("data-{i}").as_bytes()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_writes_during_compaction_survive_commit_in_memory() {
|
||||
check_writes_during_compaction_survive_commit(NeedleMapKind::InMemory);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_writes_during_compaction_survive_commit_redb() {
|
||||
check_writes_during_compaction_survive_commit(NeedleMapKind::Redb);
|
||||
}
|
||||
|
||||
/// Build a Store with N HDD disk locations under a single TempDir.
|
||||
/// Returns the store and the TempDir guard so callers keep the dirs
|
||||
/// alive for the test's lifetime.
|
||||
|
||||
@@ -116,6 +116,59 @@ fn exceeds_expected_compacted_size(expected_live_bytes: u64, dst_dat_size: u64)
|
||||
expected_live_bytes > dst_dat_size
|
||||
}
|
||||
|
||||
/// Read and parse the needle at `offset` through `read_at`.
|
||||
fn read_needle_with(
|
||||
read_at: impl Fn(&mut [u8], u64) -> Result<(), VolumeError>,
|
||||
n: &mut Needle,
|
||||
offset: i64,
|
||||
size: Size,
|
||||
version: Version,
|
||||
) -> Result<(), VolumeError> {
|
||||
match parse_needle_at(&read_at, n, offset, size, version) {
|
||||
Ok(()) => Ok(()),
|
||||
#[cfg(not(feature = "5bytes"))]
|
||||
Err(VolumeError::Needle(NeedleError::SizeMismatch { offset: o, .. }))
|
||||
if o < MAX_POSSIBLE_VOLUME_SIZE as i64 =>
|
||||
{
|
||||
// Double-read: in 4-byte offset mode, the actual data may be
|
||||
// beyond 32GB due to offset wrapping. Retry at offset + 32GB.
|
||||
parse_needle_at(
|
||||
&read_at,
|
||||
n,
|
||||
offset + MAX_POSSIBLE_VOLUME_SIZE as i64,
|
||||
size,
|
||||
version,
|
||||
)
|
||||
}
|
||||
Err(e) => Err(e),
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_needle_at(
|
||||
read_at: &impl Fn(&mut [u8], u64) -> Result<(), VolumeError>,
|
||||
n: &mut Needle,
|
||||
offset: i64,
|
||||
size: Size,
|
||||
version: Version,
|
||||
) -> Result<(), VolumeError> {
|
||||
// Storage guard: negativity-only (Go parity — storage allocates what the
|
||||
// index says). Size(0) and >1GiB map sizes must still read; only
|
||||
// negative wraps/panics. Transport cap lives in RPC handlers only.
|
||||
if size.0 < 0 {
|
||||
return Err(VolumeError::Io(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidData,
|
||||
format!("invalid needle size {}", size.0),
|
||||
)));
|
||||
}
|
||||
let actual_size = get_actual_size(size, version);
|
||||
|
||||
let mut buf = vec![0u8; actual_size as usize];
|
||||
read_at(&mut buf, offset as u64)?;
|
||||
|
||||
n.read_bytes(&buf, offset, size, version)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// VolumeInfo (.vif persistence)
|
||||
// ============================================================================
|
||||
@@ -624,6 +677,231 @@ impl DatScanPlan {
|
||||
}
|
||||
}
|
||||
|
||||
/// Holds `Volume::is_compacting` set; dropping it clears the flag.
|
||||
struct CompactionClaim(Arc<AtomicBool>);
|
||||
|
||||
impl CompactionClaim {
|
||||
fn try_claim(flag: &Arc<AtomicBool>) -> Option<Self> {
|
||||
flag.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
|
||||
.ok()
|
||||
.map(|_| CompactionClaim(flag.clone()))
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for CompactionClaim {
|
||||
fn drop(&mut self) {
|
||||
self.0.store(false, Ordering::Release);
|
||||
}
|
||||
}
|
||||
|
||||
/// A compaction copy that runs off the store lock: `.dat`/`.idx` are
|
||||
/// append-only below `idx_size`, and `makeup_diff` replays what lands after.
|
||||
pub(crate) struct CompactionJob {
|
||||
id: VolumeId,
|
||||
src_dat: NeedleStreamSource,
|
||||
src_idx: File,
|
||||
idx_size: u64,
|
||||
cpd_path: String,
|
||||
cpx_path: String,
|
||||
version: Version,
|
||||
super_block: SuperBlock,
|
||||
io_errors: Arc<IoErrorTracker>,
|
||||
_claim: CompactionClaim,
|
||||
}
|
||||
|
||||
impl CompactionJob {
|
||||
pub(crate) fn run<F>(self, progress_fn: F) -> Result<(), VolumeError>
|
||||
where
|
||||
F: Fn(i64) -> bool,
|
||||
{
|
||||
let cpd_path = &self.cpd_path;
|
||||
let cpx_path = &self.cpx_path;
|
||||
let version = self.version;
|
||||
|
||||
// Write new super block with incremented compaction revision
|
||||
let mut new_sb = self.super_block.clone();
|
||||
new_sb.compaction_revision += 1;
|
||||
let sb_bytes = new_sb.to_bytes();
|
||||
|
||||
let mut dst = OpenOptions::new()
|
||||
.read(true)
|
||||
.write(true)
|
||||
.create(true)
|
||||
.truncate(true)
|
||||
.open(cpd_path)?;
|
||||
dst.write_all(&sb_bytes)?;
|
||||
let mut new_offset = sb_bytes.len() as i64;
|
||||
|
||||
// Build new index in memory
|
||||
let mut new_nm = CompactNeedleMap::new();
|
||||
let now = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
|
||||
let entries = self.live_entries()?;
|
||||
|
||||
let mut skipped_needles: u64 = 0;
|
||||
let mut skipped_data_bytes: u64 = 0;
|
||||
let mut expected_live_bytes: u64 = 0;
|
||||
for (id, offset, size) in entries {
|
||||
// Progress callback
|
||||
if !progress_fn(offset.to_actual_offset()) {
|
||||
// Interrupted
|
||||
let _ = fs::remove_file(cpd_path);
|
||||
return Err(VolumeError::Io(io::Error::new(
|
||||
io::ErrorKind::Interrupted,
|
||||
"compaction interrupted",
|
||||
)));
|
||||
}
|
||||
|
||||
// Read needle from source
|
||||
let mut n = Needle {
|
||||
id,
|
||||
..Needle::default()
|
||||
};
|
||||
let read = read_needle_with(
|
||||
|buf, at| self.src_dat.read_exact_at(buf, at).map_err(VolumeError::Io),
|
||||
&mut n,
|
||||
offset.to_actual_offset(),
|
||||
size,
|
||||
version,
|
||||
);
|
||||
if let Err(e) = read {
|
||||
// Record EIO for health monitoring (parity with Go's checkReadWriteError).
|
||||
if let VolumeError::Io(ref io_err) = e {
|
||||
self.io_errors.check_read_write_error(Some(io_err));
|
||||
}
|
||||
// Only drop the entry when the failure is one of the well-
|
||||
// known permanent-corruption shapes. A transient disk fault,
|
||||
// a tiered-read timeout, or a Windows hardware error (which
|
||||
// surfaces as a generic Io rather than UnexpectedEof) must
|
||||
// abort so an operator notices, rather than silently
|
||||
// compacting away data that might come back on retry.
|
||||
if !is_skippable_needle_read_error(&e) {
|
||||
return Err(VolumeError::Io(io::Error::other(format!(
|
||||
"cannot hydrate needle from file: {}",
|
||||
e
|
||||
))));
|
||||
}
|
||||
skipped_needles += 1;
|
||||
if size.is_valid() {
|
||||
skipped_data_bytes += size.0 as u64;
|
||||
}
|
||||
warn!(
|
||||
volume_id = self.id.0,
|
||||
key = id.0,
|
||||
offset = offset.to_actual_offset(),
|
||||
size = size.0,
|
||||
error = %e,
|
||||
"vacuum: dropping unreadable needle"
|
||||
);
|
||||
continue;
|
||||
}
|
||||
|
||||
// Skip TTL-expired needles using the volume's TTL (matches Go's volume_vacuum.go)
|
||||
if n.has_ttl() {
|
||||
let ttl_minutes = self.super_block.ttl.minutes();
|
||||
if ttl_minutes > 0 && n.last_modified > 0 {
|
||||
let expire_at = n.last_modified + (ttl_minutes as u64) * 60;
|
||||
if now >= expire_at {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Tally the live bytes from the frozen snapshot this loop copied
|
||||
// from, not the live needle map. Unreadable needles return before
|
||||
// this point, so no further skipped-byte adjustment is needed.
|
||||
expected_live_bytes += size.0 as u64;
|
||||
|
||||
// Write needle to destination
|
||||
let bytes = n.write_bytes(version);
|
||||
dst.write_all(&bytes)?;
|
||||
|
||||
// Update new index
|
||||
new_nm.put(id, Offset::from_actual_offset(new_offset), n.size)?;
|
||||
new_offset += bytes.len() as i64;
|
||||
}
|
||||
|
||||
if skipped_needles > 0 {
|
||||
warn!(
|
||||
volume_id = self.id.0,
|
||||
skipped_needles,
|
||||
skipped_data_bytes,
|
||||
"vacuum: dropped unreadable index entries during compaction"
|
||||
);
|
||||
}
|
||||
|
||||
dst.sync_all()?;
|
||||
|
||||
if self.super_block.ttl.is_empty() {
|
||||
let dst_dat_size = dst.metadata()?.len();
|
||||
if exceeds_expected_compacted_size(expected_live_bytes, dst_dat_size) {
|
||||
let _ = fs::remove_file(cpd_path);
|
||||
let _ = fs::remove_file(cpx_path);
|
||||
return Err(VolumeError::Io(io::Error::new(
|
||||
io::ErrorKind::UnexpectedEof,
|
||||
format!(
|
||||
"volume {} unexpected new data size: {} does not match expected live content size {} from the pre-compaction snapshot",
|
||||
self.id.0, dst_dat_size, expected_live_bytes
|
||||
),
|
||||
)));
|
||||
}
|
||||
}
|
||||
|
||||
// Save new index
|
||||
new_nm.save_to_idx(cpx_path)?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Replay `.idx` up to the recorded size, as Go's `LoadFromIdx` does, and
|
||||
/// return the live entries in `.dat` order. A short read fails: treating
|
||||
/// it as the whole index would compact away every needle past it.
|
||||
fn live_entries(&self) -> Result<Vec<(NeedleId, Offset, Size)>, VolumeError> {
|
||||
let rows = self.idx_size / NEEDLE_MAP_ENTRY_SIZE as u64;
|
||||
let mut snapshot = CompactNeedleMap::new();
|
||||
let mut seen = 0u64;
|
||||
// Rows past the recorded size are makeup_diff's; do not read them.
|
||||
let mut reader = io::BufReader::new((&self.src_idx).take(self.idx_size));
|
||||
let mut buf = [0u8; NEEDLE_MAP_ENTRY_SIZE];
|
||||
while seen < rows {
|
||||
match reader.read_exact(&mut buf) {
|
||||
Ok(()) => {}
|
||||
Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
|
||||
Err(e) => return Err(e.into()),
|
||||
}
|
||||
let (key, offset, size) = idx_entry_from_bytes(&buf);
|
||||
if offset.is_zero() || size.is_deleted() {
|
||||
snapshot.delete(key, offset)?;
|
||||
} else {
|
||||
snapshot.put(key, offset, size)?;
|
||||
}
|
||||
seen += 1;
|
||||
}
|
||||
if seen < rows {
|
||||
return Err(VolumeError::Io(io::Error::new(
|
||||
io::ErrorKind::UnexpectedEof,
|
||||
format!(
|
||||
"volume {} index has {} entries, expected {}",
|
||||
self.id.0, seen, rows
|
||||
),
|
||||
)));
|
||||
}
|
||||
|
||||
let mut entries: Vec<(NeedleId, Offset, Size)> = Vec::new();
|
||||
let _ = snapshot.ascending_visit(|id, nv| {
|
||||
if !nv.offset.is_zero() && !nv.size.is_deleted() {
|
||||
entries.push((id, nv.offset, nv.size));
|
||||
}
|
||||
Ok::<(), std::convert::Infallible>(())
|
||||
});
|
||||
entries.sort_by_key(|(_, offset, _)| *offset);
|
||||
Ok(entries)
|
||||
}
|
||||
}
|
||||
|
||||
pub struct NeedleStreamInfo {
|
||||
/// Stream source for the dat file, local or remote.
|
||||
pub(crate) source: NeedleStreamSource,
|
||||
@@ -734,14 +1012,15 @@ pub struct Volume {
|
||||
last_compact_index_offset: u64,
|
||||
last_compact_revision: u16,
|
||||
|
||||
is_compacting: bool,
|
||||
/// Shared with an in-flight `CompactionJob`, which runs off the store lock.
|
||||
is_compacting: Arc<AtomicBool>,
|
||||
|
||||
/// Compaction speed limit in bytes per second (0 = unlimited).
|
||||
pub compaction_byte_per_second: i64,
|
||||
|
||||
/// Consecutive storage-media errors and the quarantine they lead to,
|
||||
/// for volume health monitoring.
|
||||
io_errors: IoErrorTracker,
|
||||
io_errors: Arc<IoErrorTracker>,
|
||||
|
||||
/// Protobuf VolumeInfo for tiered storage (.vif file).
|
||||
///
|
||||
@@ -824,9 +1103,9 @@ impl Volume {
|
||||
last_disk_check_ns: Arc::new(std::sync::atomic::AtomicI64::new(0)),
|
||||
last_compact_index_offset: 0,
|
||||
last_compact_revision: 0,
|
||||
is_compacting: false,
|
||||
is_compacting: Arc::new(AtomicBool::new(false)),
|
||||
compaction_byte_per_second: 0,
|
||||
io_errors: IoErrorTracker::default(),
|
||||
io_errors: Arc::default(),
|
||||
volume_info: PbVolumeInfo::default(),
|
||||
};
|
||||
|
||||
@@ -864,16 +1143,23 @@ impl Volume {
|
||||
last_disk_check_ns: Arc::new(std::sync::atomic::AtomicI64::new(0)),
|
||||
last_compact_index_offset: 0,
|
||||
last_compact_revision: 0,
|
||||
is_compacting: false,
|
||||
is_compacting: Arc::new(AtomicBool::new(false)),
|
||||
compaction_byte_per_second: 0,
|
||||
io_errors: IoErrorTracker::default(),
|
||||
io_errors: Arc::default(),
|
||||
volume_info: PbVolumeInfo::default(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns true if the volume is currently being compacted.
|
||||
pub fn is_compacting(&self) -> bool {
|
||||
self.is_compacting
|
||||
self.is_compacting.load(Ordering::Acquire)
|
||||
}
|
||||
|
||||
pub(crate) fn compacting_error(&self) -> VolumeError {
|
||||
VolumeError::Io(io::Error::other(format!(
|
||||
"volume {} is compacting",
|
||||
self.id
|
||||
)))
|
||||
}
|
||||
|
||||
// ---- File naming (matching Go) ----
|
||||
@@ -1562,43 +1848,13 @@ impl Volume {
|
||||
size: Size,
|
||||
_read_option: &mut ReadOption,
|
||||
) -> Result<(), VolumeError> {
|
||||
match self.read_needle_blob_and_parse(n, offset, size) {
|
||||
Ok(()) => Ok(()),
|
||||
#[cfg(not(feature = "5bytes"))]
|
||||
Err(VolumeError::Needle(NeedleError::SizeMismatch { offset: o, .. }))
|
||||
if o < MAX_POSSIBLE_VOLUME_SIZE as i64 =>
|
||||
{
|
||||
// Double-read: in 4-byte offset mode, the actual data may be
|
||||
// beyond 32GB due to offset wrapping. Retry at offset + 32GB.
|
||||
self.read_needle_blob_and_parse(n, offset + MAX_POSSIBLE_VOLUME_SIZE as i64, size)
|
||||
}
|
||||
Err(e) => Err(e),
|
||||
}
|
||||
}
|
||||
|
||||
fn read_needle_blob_and_parse(
|
||||
&self,
|
||||
n: &mut Needle,
|
||||
offset: i64,
|
||||
size: Size,
|
||||
) -> Result<(), VolumeError> {
|
||||
let version = self.version();
|
||||
// Storage guard: negativity-only (Go parity — storage allocates what the
|
||||
// index says). Size(0) and >1GiB map sizes must still read; only
|
||||
// negative wraps/panics. Transport cap lives in RPC handlers only.
|
||||
if size.0 < 0 {
|
||||
return Err(VolumeError::Io(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidData,
|
||||
format!("invalid needle size {}", size.0),
|
||||
)));
|
||||
}
|
||||
let actual_size = get_actual_size(size, version);
|
||||
|
||||
let mut buf = vec![0u8; actual_size as usize];
|
||||
self.read_exact_at_backend(&mut buf, offset as u64)?;
|
||||
|
||||
n.read_bytes(&buf, offset, size, version)?;
|
||||
Ok(())
|
||||
read_needle_with(
|
||||
|buf, at| self.read_exact_at_backend(buf, at),
|
||||
n,
|
||||
offset,
|
||||
size,
|
||||
self.version(),
|
||||
)
|
||||
}
|
||||
|
||||
/// Read raw needle blob at a specific offset.
|
||||
@@ -1612,7 +1868,7 @@ impl Volume {
|
||||
|
||||
fn read_needle_blob_unlocked(&self, offset: i64, size: Size) -> Result<Vec<u8>, VolumeError> {
|
||||
let version = self.version();
|
||||
// Storage guard: negativity-only (Go parity). See read_needle_blob_and_parse.
|
||||
// Storage guard: negativity-only (Go parity). See parse_needle_at.
|
||||
if size.0 < 0 {
|
||||
return Err(VolumeError::Io(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidData,
|
||||
@@ -3429,7 +3685,7 @@ impl Volume {
|
||||
|
||||
/// Throttle IO during compaction to avoid saturating disk.
|
||||
pub fn maybe_throttle_compaction(&self, bytes_written: u64) {
|
||||
if self.compaction_byte_per_second <= 0 || !self.is_compacting {
|
||||
if self.compaction_byte_per_second <= 0 || !self.is_compacting() {
|
||||
return;
|
||||
}
|
||||
// Simple throttle: sleep based on bytes written vs allowed rate
|
||||
@@ -3628,7 +3884,7 @@ impl Volume {
|
||||
if self.is_read_only() {
|
||||
return Err(VolumeError::ReadOnly);
|
||||
}
|
||||
// Storage guard: negativity-only (Go parity). See read_needle_blob_and_parse.
|
||||
// Storage guard: negativity-only (Go parity). See parse_needle_at.
|
||||
if size.0 < 0 {
|
||||
return Err(VolumeError::Io(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidData,
|
||||
@@ -3833,195 +4089,70 @@ impl Volume {
|
||||
where
|
||||
F: Fn(i64) -> bool,
|
||||
{
|
||||
if self.is_compacting {
|
||||
return Ok(()); // already compacting
|
||||
match self.begin_compact_by_index()? {
|
||||
Some(job) => job.run(progress_fn),
|
||||
None => Ok(()), // already compacting
|
||||
}
|
||||
self.is_compacting = true;
|
||||
|
||||
let result = self.do_compact_by_index(progress_fn);
|
||||
|
||||
self.is_compacting = false;
|
||||
result
|
||||
}
|
||||
|
||||
fn do_compact_by_index<F>(&mut self, progress_fn: F) -> Result<(), VolumeError>
|
||||
where
|
||||
F: Fn(i64) -> bool,
|
||||
{
|
||||
// Guard against nil needle map (matches Go's nil check before compaction sync)
|
||||
if self.nm.is_none() {
|
||||
/// Claim the compaction, record the `makeup_diff` watermark, and open the
|
||||
/// sources the copy reads, so the caller can run the copy after dropping
|
||||
/// its store guard. `None` means a compaction is already running.
|
||||
pub(crate) fn begin_compact_by_index(&mut self) -> Result<Option<CompactionJob>, VolumeError> {
|
||||
let Some(claim) = CompactionClaim::try_claim(&self.is_compacting) else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
// Guard against nil needle map (matches Go's nil check)
|
||||
let Some(nm) = self.nm.as_ref() else {
|
||||
return Err(VolumeError::Io(io::Error::other(format!(
|
||||
"volume {} needle map is nil",
|
||||
self.id
|
||||
))));
|
||||
};
|
||||
if let Some(e) = self.unavailable_error() {
|
||||
return Err(e);
|
||||
}
|
||||
let idx_size = nm.index_file_size();
|
||||
|
||||
// Fresh opens, not `try_clone`: see `dat_scan_plan`.
|
||||
let src_dat = if self.dat_file.is_some() {
|
||||
NeedleStreamSource::Local(open_volume_file(
|
||||
OpenOptions::new().read(true),
|
||||
self.file_name(".dat"),
|
||||
)?)
|
||||
} else if let Some(remote) = self.remote_dat_file() {
|
||||
NeedleStreamSource::Remote(remote)
|
||||
} else {
|
||||
return Err(VolumeError::Io(io::Error::other("dat file not open")));
|
||||
};
|
||||
let src_idx = open_volume_file(OpenOptions::new().read(true), self.file_name(".idx"))?;
|
||||
|
||||
// Record state before compaction for makeupDiff
|
||||
self.last_compact_index_offset = self.nm.as_ref().map_or(0, |nm| nm.index_file_size());
|
||||
self.last_compact_index_offset = idx_size;
|
||||
self.last_compact_revision = self.super_block.compaction_revision;
|
||||
|
||||
// Sync current data
|
||||
self.sync_to_disk()?;
|
||||
|
||||
let cpd_path = self.file_name(".cpd");
|
||||
let cpx_path = self.file_name(".cpx");
|
||||
let version = self.version();
|
||||
|
||||
// Write new super block with incremented compaction revision
|
||||
let mut new_sb = self.super_block.clone();
|
||||
new_sb.compaction_revision += 1;
|
||||
let sb_bytes = new_sb.to_bytes();
|
||||
|
||||
let mut dst = OpenOptions::new()
|
||||
.read(true)
|
||||
.write(true)
|
||||
.create(true)
|
||||
.truncate(true)
|
||||
.open(&cpd_path)?;
|
||||
dst.write_all(&sb_bytes)?;
|
||||
let mut new_offset = sb_bytes.len() as i64;
|
||||
|
||||
// Build new index in memory
|
||||
let mut new_nm = CompactNeedleMap::new();
|
||||
let now = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
|
||||
// Collect live entries from needle map (sorted ascending)
|
||||
let nm = self.nm.as_ref().ok_or(VolumeError::NotInitialized)?;
|
||||
let mut entries: Vec<(NeedleId, Offset, Size)> = Vec::new();
|
||||
for (id, nv) in nm.iter_entries().map_err(VolumeError::Io)? {
|
||||
if nv.offset.is_zero() || nv.size.is_deleted() {
|
||||
continue;
|
||||
}
|
||||
entries.push((id, nv.offset, nv.size));
|
||||
}
|
||||
entries.sort_by_key(|(_, offset, _)| *offset);
|
||||
|
||||
let mut skipped_needles: u64 = 0;
|
||||
let mut skipped_data_bytes: u64 = 0;
|
||||
let mut expected_live_bytes: u64 = 0;
|
||||
for (id, offset, size) in entries {
|
||||
// Progress callback
|
||||
if !progress_fn(offset.to_actual_offset()) {
|
||||
// Interrupted
|
||||
let _ = fs::remove_file(&cpd_path);
|
||||
return Err(VolumeError::Io(io::Error::new(
|
||||
io::ErrorKind::Interrupted,
|
||||
"compaction interrupted",
|
||||
)));
|
||||
}
|
||||
|
||||
// Read needle from source
|
||||
let mut n = Needle {
|
||||
id,
|
||||
..Needle::default()
|
||||
};
|
||||
match self.read_needle_data_at(&mut n, offset.to_actual_offset(), size) {
|
||||
Ok(()) => {}
|
||||
Err(e) => {
|
||||
// Record EIO for health monitoring (parity with Go's checkReadWriteError).
|
||||
if let VolumeError::Io(ref io_err) = e {
|
||||
self.check_read_write_error(Some(io_err));
|
||||
}
|
||||
// Only drop the entry when the failure is one of the well-
|
||||
// known permanent-corruption shapes. A transient disk fault,
|
||||
// a tiered-read timeout, or a Windows hardware error (which
|
||||
// surfaces as a generic Io rather than UnexpectedEof) must
|
||||
// abort so an operator notices, rather than silently
|
||||
// compacting away data that might come back on retry.
|
||||
if !is_skippable_needle_read_error(&e) {
|
||||
return Err(VolumeError::Io(io::Error::other(format!(
|
||||
"cannot hydrate needle from file: {}",
|
||||
e
|
||||
))));
|
||||
}
|
||||
skipped_needles += 1;
|
||||
if size.is_valid() {
|
||||
skipped_data_bytes += size.0 as u64;
|
||||
}
|
||||
warn!(
|
||||
volume_id = self.id.0,
|
||||
key = id.0,
|
||||
offset = offset.to_actual_offset(),
|
||||
size = size.0,
|
||||
error = %e,
|
||||
"vacuum: dropping unreadable needle"
|
||||
);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
// Skip TTL-expired needles using the volume's TTL (matches Go's volume_vacuum.go)
|
||||
if n.has_ttl() {
|
||||
let ttl_minutes = self.super_block.ttl.minutes();
|
||||
if ttl_minutes > 0 && n.last_modified > 0 {
|
||||
let expire_at = n.last_modified + (ttl_minutes as u64) * 60;
|
||||
if now >= expire_at {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Tally the live bytes from the frozen snapshot this loop copied
|
||||
// from, not the live needle map. Unreadable needles return before
|
||||
// this point, so no further skipped-byte adjustment is needed.
|
||||
expected_live_bytes += size.0 as u64;
|
||||
|
||||
// Write needle to destination
|
||||
let bytes = n.write_bytes(version);
|
||||
dst.write_all(&bytes)?;
|
||||
|
||||
// Update new index
|
||||
new_nm.put(id, Offset::from_actual_offset(new_offset), n.size)?;
|
||||
new_offset += bytes.len() as i64;
|
||||
}
|
||||
|
||||
if skipped_needles > 0 {
|
||||
warn!(
|
||||
volume_id = self.id.0,
|
||||
skipped_needles,
|
||||
skipped_data_bytes,
|
||||
"vacuum: dropped unreadable index entries during compaction"
|
||||
);
|
||||
}
|
||||
|
||||
dst.sync_all()?;
|
||||
|
||||
if self.super_block.ttl.is_empty() {
|
||||
let dst_dat_size = dst.metadata()?.len();
|
||||
if exceeds_expected_compacted_size(expected_live_bytes, dst_dat_size) {
|
||||
let _ = fs::remove_file(&cpd_path);
|
||||
let _ = fs::remove_file(&cpx_path);
|
||||
return Err(VolumeError::Io(io::Error::new(
|
||||
io::ErrorKind::UnexpectedEof,
|
||||
format!(
|
||||
"volume {} unexpected new data size: {} does not match expected live content size {} from the pre-compaction snapshot",
|
||||
self.id.0, dst_dat_size, expected_live_bytes
|
||||
),
|
||||
)));
|
||||
}
|
||||
}
|
||||
|
||||
// Save new index
|
||||
new_nm.save_to_idx(&cpx_path)?;
|
||||
|
||||
Ok(())
|
||||
Ok(Some(CompactionJob {
|
||||
id: self.id,
|
||||
src_dat,
|
||||
src_idx,
|
||||
idx_size,
|
||||
cpd_path: self.file_name(".cpd"),
|
||||
cpx_path: self.file_name(".cpx"),
|
||||
version: self.version(),
|
||||
super_block: self.super_block.clone(),
|
||||
io_errors: self.io_errors.clone(),
|
||||
_claim: claim,
|
||||
}))
|
||||
}
|
||||
|
||||
/// Commit a previously completed compaction: swap .cpd/.cpx to .dat/.idx and reload.
|
||||
/// Matches Go's isCompactionInProgress CompareAndSwap guard.
|
||||
pub fn commit_compact(&mut self) -> Result<(), VolumeError> {
|
||||
if self.is_compacting {
|
||||
let Some(_claim) = CompactionClaim::try_claim(&self.is_compacting) else {
|
||||
return Ok(()); // already compacting, silently skip (matches Go)
|
||||
}
|
||||
self.is_compacting = true;
|
||||
|
||||
let result = self.do_commit_compact();
|
||||
|
||||
self.is_compacting = false;
|
||||
result
|
||||
};
|
||||
self.do_commit_compact()
|
||||
}
|
||||
|
||||
fn do_commit_compact(&mut self) -> Result<(), VolumeError> {
|
||||
@@ -4191,6 +4322,10 @@ impl Volume {
|
||||
|
||||
/// Clean up leftover compaction files (.cpd, .cpx).
|
||||
pub fn cleanup_compact(&self) -> Result<(), VolumeError> {
|
||||
// An in-flight copy is still writing .cpd/.cpx.
|
||||
if self.is_compacting() {
|
||||
return Err(self.compacting_error());
|
||||
}
|
||||
// Refuse to unlink .cpd/.cpx while a .cpc marker exists: those temp files
|
||||
// are the only inputs reconcile can roll forward to, so removing them
|
||||
// mid-commit would strand a decided swap.
|
||||
@@ -4395,6 +4530,10 @@ impl Volume {
|
||||
/// or nothing is co-located to move, and is used to pull an index a decode
|
||||
/// or reconstruct left beside the data back into the configured `-dir.idx`.
|
||||
pub fn relocate_index_to(&mut self, new_idx_dir: &str) -> Result<(), VolumeError> {
|
||||
// Moving .idx would strand the .cpx an in-flight copy writes beside it.
|
||||
if self.is_compacting() {
|
||||
return Err(self.compacting_error());
|
||||
}
|
||||
let _guard = self.data_file_access_control.write_lock();
|
||||
|
||||
if self.dir_idx == new_idx_dir {
|
||||
@@ -4464,11 +4603,8 @@ impl Volume {
|
||||
if (only_empty || only_garbage) && !empty_ok && !garbage_ok {
|
||||
return Err(VolumeError::NotEmpty);
|
||||
}
|
||||
if self.is_compacting {
|
||||
return Err(VolumeError::Io(io::Error::other(format!(
|
||||
"volume {} is compacting",
|
||||
self.id
|
||||
))));
|
||||
if self.is_compacting() {
|
||||
return Err(self.compacting_error());
|
||||
}
|
||||
|
||||
let (storage_name, storage_key) = self.remote_storage_name_key();
|
||||
@@ -7847,23 +7983,24 @@ mod tests {
|
||||
.remove("s3.vif_tierdown_test");
|
||||
}
|
||||
|
||||
// A .sdx that cannot be read end to end must abort compaction. Treating the
|
||||
// short scan as the complete live set would commit a volume missing every
|
||||
// needle past the truncation.
|
||||
// An .idx that cannot be read up to the recorded size must abort
|
||||
// compaction. Treating the short scan as the complete live set would
|
||||
// commit a volume missing every needle past the truncation.
|
||||
#[test]
|
||||
fn test_compaction_aborts_on_unreadable_sorted_index() {
|
||||
fn test_compaction_aborts_on_truncated_index() {
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let dir = tmp.path().to_str().unwrap();
|
||||
let mut v = reload_as_tiered(dir, "vif_compact_test", 4);
|
||||
let Some(NeedleMap::SortedFile(ref nm)) = v.nm else {
|
||||
let Some(NeedleMap::SortedFile(_)) = v.nm else {
|
||||
panic!("tiered volume should search the on-disk .sdx");
|
||||
};
|
||||
let sdx_path = nm.db_file_name().to_string();
|
||||
|
||||
let sdx = OpenOptions::new().write(true).open(&sdx_path).unwrap();
|
||||
sdx.set_len(NEEDLE_MAP_ENTRY_SIZE as u64).unwrap();
|
||||
drop(sdx);
|
||||
crate::storage::needle_map::file_pool::pooled_index_files().discard(&sdx_path);
|
||||
let idx = OpenOptions::new()
|
||||
.write(true)
|
||||
.open(v.file_name(".idx"))
|
||||
.unwrap();
|
||||
idx.set_len(NEEDLE_MAP_ENTRY_SIZE as u64).unwrap();
|
||||
drop(idx);
|
||||
|
||||
let err = v
|
||||
.compact_by_index(0, 0, |_| true)
|
||||
@@ -7938,6 +8075,79 @@ mod tests {
|
||||
assert_eq!(probe.data, b"still-readable");
|
||||
}
|
||||
|
||||
// Compaction copies the .idx prefix index_file_size() reports. A read-only
|
||||
// volume's map has no writer, so that must still be the loaded size, or the
|
||||
// copy is empty and the commit drops every needle.
|
||||
#[cfg(unix)]
|
||||
fn check_read_only_compaction_keeps_needles(index_dir_writable: bool) {
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
|
||||
let tmp = TempDir::new().unwrap();
|
||||
let dir = tmp.path().to_str().unwrap().to_string();
|
||||
{
|
||||
let mut v = make_test_volume(&dir);
|
||||
for i in 1..=3u64 {
|
||||
write_test_needle(&mut v, i, format!("data-{i}").as_bytes());
|
||||
}
|
||||
v.set_read_only_persist(false, true).unwrap();
|
||||
v.sync_to_disk().unwrap();
|
||||
}
|
||||
let idx_len = fs::metadata(format!("{dir}/1.idx")).unwrap().len();
|
||||
|
||||
let set_mode = |mode: u32| {
|
||||
let mut perms = std::fs::metadata(tmp.path()).unwrap().permissions();
|
||||
perms.set_mode(mode);
|
||||
std::fs::set_permissions(tmp.path(), perms).unwrap();
|
||||
};
|
||||
if !index_dir_writable {
|
||||
set_mode(0o555);
|
||||
}
|
||||
let loaded = Volume::new(
|
||||
&dir,
|
||||
&dir,
|
||||
VolumeId(1),
|
||||
NeedleMapKind::InMemory,
|
||||
&VolumeSpec::default(),
|
||||
);
|
||||
set_mode(0o755);
|
||||
|
||||
let mut v = loaded.unwrap();
|
||||
assert_eq!(
|
||||
matches!(v.nm, Some(NeedleMap::SortedFile(_))),
|
||||
index_dir_writable
|
||||
);
|
||||
assert_eq!(v.idx_file_size(), idx_len);
|
||||
|
||||
v.compact_by_index(0, 0, |_| true).unwrap();
|
||||
v.commit_compact().unwrap();
|
||||
for i in 1..=3u64 {
|
||||
let mut n = Needle {
|
||||
id: NeedleId(i),
|
||||
..Needle::default()
|
||||
};
|
||||
v.read_needle(&mut n).unwrap();
|
||||
assert_eq!(n.data, format!("data-{i}").as_bytes());
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
#[cfg(unix)]
|
||||
fn test_compacting_a_read_only_volume_keeps_its_needles_sorted_index() {
|
||||
check_read_only_compaction_keeps_needles(true);
|
||||
}
|
||||
|
||||
// The .sdx could not be built, so the index fell back to memory.
|
||||
#[test]
|
||||
#[cfg(unix)]
|
||||
fn test_compacting_a_read_only_volume_keeps_its_needles_in_memory_fallback() {
|
||||
// root ignores the directory mode, so there is nothing to simulate.
|
||||
// SAFETY: `geteuid` takes no arguments, reads no memory and cannot fail.
|
||||
if unsafe { libc::geteuid() } == 0 {
|
||||
return;
|
||||
}
|
||||
check_read_only_compaction_keeps_needles(false);
|
||||
}
|
||||
|
||||
// set_writable clears the read-only flags before it can know the .idx writer
|
||||
// will attach. If attaching fails the flags have to go back: a volume that
|
||||
// advertises writable while its map has no writer takes puts into memory and
|
||||
|
||||
Reference in New Issue
Block a user