volume server: refuse a tier move while compacting, and a commit once tiered (#11539)

* volume server: refuse a tier move while compacting, and a commit once tiered

A tier move to remote and a vacuum compaction of the same volume could
interleave and leave the volume unreadable:

- A compaction committing while the upload ran swapped .dat/.idx under
  the transfer, which reopens the .dat by path per part. The move then
  published an object holding the old (or a mixed) layout against the
  compacted .idx, and with keep_local_dat_file=false deleted the only
  compacted .dat.
- A tier move finishing while the compaction copy ran (or between the
  copy and the commit) let the commit swap in the compacted .idx while
  the reload served the pre-compaction remote object through it.

The tier move now refuses to start while the volume is compacting, and
re-checks the compaction revision under the store write lock before it
records the remote file; on a mismatch it deletes the uploaded object
and fails with FailedPrecondition, leaving the volume local. Committing
a compaction on a volume that has a remote file is refused and its
.cpd/.cpx removed, since the reload would read the remote object
through the compacted index.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* volume server: abort a tier move whose volume was replaced or removed

The tier-up bookkeeping looked the volume up by id only and compared the
compaction revision. A delete and re-create of the same id during the upload
yields a fresh volume at the same revision, so the move recorded the old
volume's object on the new one and, without keep_local_dat_file, removed the
new .dat. An unmounted volume was skipped and the move reported success,
leaving the uploaded object referenced by nothing.

Capture the volume instance (its data-file access control Arc, as the scan
and read plans do) with the revision, and require both under the store write
lock. A replaced volume fails with FailedPrecondition, a missing one with
NotFound; either way nothing is recorded and the object is deleted after the
lock is released. Go fails in both cases because deleting or unmounting closes
the descriptor its copy reads.

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:
Eliah RusinandClaude Opus 5.5 authored and GitHub committed 2026-10-01 20:20:16 +08:00
1 parent f2498e122a
commit fea14c01a7
2 files changed
+487 -10

No files matched your search

+468 -10
View File
@@ -4045,7 +4045,7 @@ impl VolumeServer for VolumeGrpcService {
let vid = VolumeId(req.volume_id);
// Validate volume exists and collection matches
let dat_path = {
let (dat_path, instance, compaction_revision) = {
let store = self.state.store.read().unwrap();
let (_, vol) = store
.find_volume(vid)
@@ -4070,6 +4070,13 @@ impl VolumeServer for VolumeGrpcService {
));
}
if vol.is_compacting() {
return Err(Status::failed_precondition(format!(
"volume {} is compacting",
vid
)));
}
// Check if the destination backend already exists in volume info
let (backend_type, backend_id) =
crate::remote_storage::s3_tier::backend_name_to_type_id(
@@ -4084,7 +4091,11 @@ impl VolumeServer for VolumeGrpcService {
}
}
dat_path
(
dat_path,
vol.instance(),
vol.super_block.compaction_revision,
)
};
// Store the source .dat mtime, not the upload time, so a reload computes
@@ -4182,15 +4193,32 @@ impl VolumeServer for VolumeGrpcService {
// state consistent; stopping now would leave the object paid
// for and referenced by nothing. Go does not gate here either:
// its progress callback only runs during the transfer.
// Update volume info with remote file reference
{
// Update volume info with remote file reference, unless the
// volume was compacted, replaced or removed mid-upload: the
// object would not match what now holds this id.
let abort = {
let mut store = state.store.write().unwrap();
if let Some((_, vol)) = store.find_volume_mut(vid) {
let abort = match store.find_volume(vid) {
None => Some(Status::not_found(format!("volume {} not found", vid))),
Some((_, v))
if !v.is_instance(&instance)
|| v.super_block.compaction_revision != compaction_revision =>
{
Some(Status::failed_precondition(format!(
"volume {} was compacted or replaced during tier move to remote",
vid
)))
}
Some(_) => None,
};
if abort.is_none()
&& let Some((_, vol)) = store.find_volume_mut(vid)
{
vol.update_remote_files(|files| {
files.push(volume_server_pb::RemoteFile {
backend_type: backend_type.clone(),
backend_id: backend_id.clone(),
key,
key: key.clone(),
offset: 0,
file_size: size,
modified_time: dat_modified_secs,
@@ -4227,6 +4255,18 @@ impl VolumeServer for VolumeGrpcService {
let _ = std::fs::remove_file(&dat);
}
}
abort
};
if let Some(status) = abort {
if let Err(e) = backend.delete_file(&key).await {
tracing::warn!(
"volume {} could not delete stale tier object {}: {}",
vid,
key,
e
);
}
return Err(status);
}
// Go does NOT send a final 100% progress message after upload completion
@@ -7459,10 +7499,8 @@ mod tests {
.remove("s3.tier_down_keep");
}
// The tier-up handler has no end-to-end test — exercising it needs a fake
// S3 that accepts multipart uploads — so this probes only the part that
// changed: the destination is resolved from the process-wide registry, now
// the only one. A backend registered nowhere else has to get past that
// The destination is resolved from the process-wide registry, now the
// only one. A backend registered nowhere else has to get past that
// lookup. The stream stays open until the spawned transfer reports its
// terminal error, so the task cannot race a dropped receiver or outlive
// the test.
@@ -7517,6 +7555,426 @@ mod tests {
}
}
fn register_tier_backend(name: &str, endpoint: String) {
global_s3_tier_registry().write().unwrap().register(
name.to_string(),
S3TierBackend::new(&S3TierConfig {
access_key: "access".to_string(),
secret_key: "secret".to_string(),
region: "us-east-1".to_string(),
bucket: "bucket-a".to_string(),
endpoint,
storage_class: "STANDARD".to_string(),
force_path_style: true,
}),
);
}
fn tier_up_request(
backend: &str,
keep_local_dat_file: bool,
) -> Request<volume_server_pb::VolumeTierMoveDatToRemoteRequest> {
Request::new(volume_server_pb::VolumeTierMoveDatToRemoteRequest {
volume_id: 1,
collection: String::new(),
destination_backend_name: backend.to_string(),
keep_local_dat_file,
})
}
/// Starts a tier move of volume 1 and waits until `s3` holds its part.
async fn start_parked_tier_move(
service: &VolumeGrpcService,
s3: &mut FakeMultipartS3,
backend: &str,
keep_local_dat_file: bool,
) -> BoxStream<volume_server_pb::VolumeTierMoveDatToRemoteResponse> {
let stream = service
.volume_tier_move_dat_to_remote(tier_up_request(backend, keep_local_dat_file))
.await
.unwrap()
.into_inner();
s3.parked
.recv()
.await
.expect("the upload must reach its part");
stream
}
async fn tier_move_error(
mut stream: BoxStream<volume_server_pb::VolumeTierMoveDatToRemoteResponse>,
) -> Option<Status> {
while let Some(message) = stream.next().await {
if let Err(e) = message {
return Some(e);
}
}
None
}
/// Writes needles 1 and 2 and deletes 1, so a compaction moves needle 2.
fn seed_compactable_volume(service: &VolumeGrpcService) {
let mut store = service.state.store.write().unwrap();
for id in [1u64, 2] {
let data = format!("needle-{id}").into_bytes();
let mut n = Needle {
id: NeedleId(id),
cookie: Cookie(id as u32),
data_size: data.len() as u32,
data,
..Needle::default()
};
store
.write_volume_needle(VolumeId(1), &mut n, true)
.unwrap();
}
let mut n = Needle {
id: NeedleId(1),
cookie: Cookie(1),
..Needle::default()
};
store.delete_volume_needle(VolumeId(1), &mut n).unwrap();
}
fn read_surviving_needle(
service: &VolumeGrpcService,
) -> Result<Vec<u8>, crate::storage::volume::VolumeError> {
let mut n = Needle {
id: NeedleId(2),
cookie: Cookie(2),
..Needle::default()
};
let store = service.state.store.read().unwrap();
store.read_volume_needle(VolumeId(1), &mut n)?;
Ok(n.data)
}
fn compact_volume(service: &VolumeGrpcService) {
let job = service
.state
.store
.write()
.unwrap()
.begin_compact_volume(VolumeId(1), 0)
.unwrap()
.expect("no compaction should be running");
job.run(|_| true).unwrap();
}
struct FakeMultipartS3 {
endpoint: String,
shutdown: tokio::sync::oneshot::Sender<()>,
/// Receives one message per part as it arrives, before it is answered.
parked: tokio::sync::mpsc::UnboundedReceiver<()>,
/// Each permit answers one parked part.
release: Arc<tokio::sync::Semaphore>,
delete_count: Arc<std::sync::atomic::AtomicUsize>,
}
/// An S3 endpoint that accepts multipart uploads, holding every part
/// until the test releases it.
fn spawn_multipart_s3_server() -> FakeMultipartS3 {
use axum::http::{Method, StatusCode, Uri, header};
use std::sync::atomic::{AtomicUsize, Ordering};
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
listener.set_nonblocking(true).unwrap();
let (shutdown, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let (parked_tx, parked) = tokio::sync::mpsc::unbounded_channel();
let release = Arc::new(tokio::sync::Semaphore::new(0));
let delete_count = Arc::new(AtomicUsize::new(0));
let (ready_tx, ready_rx) = std::sync::mpsc::channel::<()>();
let handler_release = release.clone();
let handler_deletes = delete_count.clone();
std::thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
runtime.block_on(async move {
let app = axum::Router::new().fallback(axum::routing::any(
move |method: Method, uri: Uri| {
let parked_tx = parked_tx.clone();
let release = handler_release.clone();
let deletes = handler_deletes.clone();
async move {
let query = uri.query().unwrap_or_default();
let xml =
|body: String| (StatusCode::OK, [(header::ETAG, "\"etag\"")], body);
match method {
Method::DELETE => {
deletes.fetch_add(1, Ordering::SeqCst);
(
StatusCode::NO_CONTENT,
[(header::ETAG, "\"etag\"")],
String::new(),
)
}
Method::POST if query.contains("uploads") => {
xml("<InitiateMultipartUploadResult><Bucket>bucket-a</Bucket>\
<Key>k</Key><UploadId>upload-1</UploadId>\
</InitiateMultipartUploadResult>"
.to_string())
}
Method::PUT if query.contains("partNumber") => {
let _ = parked_tx.send(());
release.acquire().await.unwrap().forget();
xml(String::new())
}
Method::POST => {
xml("<CompleteMultipartUploadResult><Bucket>bucket-a</Bucket>\
<Key>k</Key><ETag>\"etag\"</ETag>\
</CompleteMultipartUploadResult>"
.to_string())
}
_ => (
StatusCode::NOT_FOUND,
[(header::ETAG, "\"etag\"")],
String::new(),
),
}
}
},
));
let listener = tokio::net::TcpListener::from_std(listener).unwrap();
let _ = ready_tx.send(());
axum::serve(listener, app)
.with_graceful_shutdown(async move {
let _ = shutdown_rx.await;
})
.await
.unwrap();
});
});
ready_rx.recv().unwrap();
FakeMultipartS3 {
endpoint: format!("http://{}", addr),
shutdown,
parked,
release,
delete_count,
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_tier_move_to_remote_refused_while_compacting() {
let (service, _tmp) = make_local_service_with_volume("", None);
// Nothing listens here, so a move that does start fails fast.
register_tier_backend("s3.tier_up_compacting", "http://127.0.0.1:1".to_string());
let job = service
.state
.store
.write()
.unwrap()
.begin_compact_volume(VolumeId(1), 0)
.unwrap()
.unwrap();
let result = service
.volume_tier_move_dat_to_remote(tier_up_request("s3.tier_up_compacting", true))
.await;
global_s3_tier_registry()
.write()
.unwrap()
.remove("s3.tier_up_compacting");
let Err(err) = result else {
panic!("a tier move must not start while the volume is compacting");
};
assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}");
assert_eq!(err.message(), "volume 1 is compacting");
drop(job);
}
// A commit that lands while the upload runs changes the .dat under it;
// the object must not be published against the compacted .idx.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_tier_move_to_remote_aborts_when_compaction_commits_mid_upload() {
let (service, tmp) = make_local_service_with_volume("", None);
seed_compactable_volume(&service);
let mut s3 = spawn_multipart_s3_server();
register_tier_backend("s3.tier_up_mid_commit", s3.endpoint.clone());
let stream = start_parked_tier_move(&service, &mut s3, "s3.tier_up_mid_commit", true).await;
compact_volume(&service);
service
.state
.store
.write()
.unwrap()
.commit_compact_volume(VolumeId(1))
.unwrap();
s3.release.add_permits(1);
let terminal = tier_move_error(stream).await;
global_s3_tier_registry()
.write()
.unwrap()
.remove("s3.tier_up_mid_commit");
let err = terminal.expect("the tier move must fail once the volume was compacted");
assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}");
assert_eq!(
s3.delete_count.load(std::sync::atomic::Ordering::SeqCst),
1,
"the stale object must be deleted"
);
{
let store = service.state.store.read().unwrap();
let (_, vol) = store.find_volume(VolumeId(1)).unwrap();
assert!(!vol.has_remote_file());
assert!(vol.volume_info().files.is_empty());
}
assert!(tmp.path().join("1.dat").exists());
assert_eq!(read_surviving_needle(&service).unwrap(), b"needle-2");
let _ = s3.shutdown.send(());
}
// A volume deleted and re-created under the same id mid-upload is back at
// the same compaction revision; the object holds the old volume's bytes.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_tier_move_to_remote_aborts_when_volume_is_recreated_mid_upload() {
let (service, tmp) = make_local_service_with_volume("", None);
let mut s3 = spawn_multipart_s3_server();
register_tier_backend("s3.tier_up_recreated", s3.endpoint.clone());
let stream = start_parked_tier_move(&service, &mut s3, "s3.tier_up_recreated", false).await;
{
let mut store = service.state.store.write().unwrap();
store
.delete_volume(VolumeId(1), false, false, false)
.unwrap();
store
.add_volume(VolumeId(1), DiskType::HardDrive, &VolumeSpec::default())
.unwrap();
let (_, vol) = store.find_volume(VolumeId(1)).unwrap();
assert_eq!(vol.super_block.compaction_revision, 0);
}
s3.release.add_permits(1);
let terminal = tier_move_error(stream).await;
global_s3_tier_registry()
.write()
.unwrap()
.remove("s3.tier_up_recreated");
{
let store = service.state.store.read().unwrap();
let (_, vol) = store.find_volume(VolumeId(1)).unwrap();
assert!(!vol.has_remote_file(), "the new volume must not be tiered");
assert!(vol.volume_info().files.is_empty());
}
assert!(tmp.path().join("1.dat").exists());
if let Ok(vif) = std::fs::read_to_string(tmp.path().join("1.vif")) {
assert!(!vif.contains("tier_up_recreated"), "{vif}");
}
let err = terminal.expect("the tier move must fail once the volume was replaced");
assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}");
assert_eq!(
s3.delete_count.load(std::sync::atomic::Ordering::SeqCst),
1,
"the stale object must be deleted"
);
let _ = s3.shutdown.send(());
}
// Go fails here too: unmounting closes the descriptor its copy reads.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_tier_move_to_remote_fails_when_volume_is_unmounted_mid_upload() {
let (service, tmp) = make_local_service_with_volume("", None);
let mut s3 = spawn_multipart_s3_server();
register_tier_backend("s3.tier_up_unmounted", s3.endpoint.clone());
let stream = start_parked_tier_move(&service, &mut s3, "s3.tier_up_unmounted", false).await;
assert!(
service
.state
.store
.write()
.unwrap()
.unmount_volume(VolumeId(1))
.unwrap()
);
s3.release.add_permits(1);
let terminal = tier_move_error(stream).await;
global_s3_tier_registry()
.write()
.unwrap()
.remove("s3.tier_up_unmounted");
let err = terminal.expect("the tier move must fail once the volume was unmounted");
assert_eq!(err.code(), tonic::Code::NotFound, "{err:?}");
assert_eq!(
s3.delete_count.load(std::sync::atomic::Ordering::SeqCst),
1,
"the unreferenced object must be deleted"
);
assert!(tmp.path().join("1.dat").exists());
if let Ok(vif) = std::fs::read_to_string(tmp.path().join("1.vif")) {
assert!(!vif.contains("tier_up_unmounted"), "{vif}");
}
let _ = s3.shutdown.send(());
}
// A tier move that lands while the copy runs leaves a compaction that
// must not be committed: the reload would read the remote object, which
// keeps the pre-compaction layout, through the compacted .idx.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_compaction_commit_refused_once_the_volume_is_tiered() {
let (service, tmp) = make_local_service_with_volume("", None);
seed_compactable_volume(&service);
let job = service
.state
.store
.write()
.unwrap()
.begin_compact_volume(VolumeId(1), 0)
.unwrap()
.unwrap();
// Tier it the way the tier-up bookkeeping does, keeping the local .dat.
let dat_bytes = std::fs::read(tmp.path().join("1.dat")).unwrap();
let (endpoint, shutdown_tx, _deletes) = spawn_fake_s3_server(dat_bytes.clone());
register_tier_backend("s3.tier_up_then_commit", endpoint);
{
let mut store = service.state.store.write().unwrap();
let (_, vol) = store.find_volume_mut(VolumeId(1)).unwrap();
vol.update_remote_files(|files| {
files.push(volume_server_pb::RemoteFile {
backend_type: "s3".to_string(),
backend_id: "tier_up_then_commit".to_string(),
key: "remote-key".to_string(),
offset: 0,
file_size: dat_bytes.len() as u64,
modified_time: 0,
extension: ".dat".to_string(),
})
})
.unwrap();
vol.save_volume_info().unwrap();
vol.load_remote_dat_file().unwrap();
}
job.run(|_| true).unwrap();
let result = service
.state
.store
.write()
.unwrap()
.commit_compact_volume(VolumeId(1));
let read = read_surviving_needle(&service);
global_s3_tier_registry()
.write()
.unwrap()
.remove("s3.tier_up_then_commit");
let _ = shutdown_tx.send(());
let err = result.expect_err("a tiered volume must refuse the commit");
assert!(err.to_string().contains("tiered"), "{err}");
assert!(!tmp.path().join("1.cpd").exists());
assert!(!tmp.path().join("1.cpx").exists());
assert_eq!(read.unwrap(), b"needle-2");
}
/// Build a local service whose volume has a `.dat` large enough to span
/// several 2MB copy chunks, so the streaming copy paths are exercised
/// across multiple messages rather than a single buffer.
+19
View File
@@ -1308,6 +1308,16 @@ impl Volume {
}
}
/// Identifies this volume instance: a re-created or remounted volume gets
/// a new one, even at the same compaction revision.
pub(crate) fn instance(&self) -> Arc<DataFileAccessControl> {
self.data_file_access_control.clone()
}
pub(crate) fn is_instance(&self, instance: &Arc<DataFileAccessControl>) -> bool {
Arc::ptr_eq(instance, &self.data_file_access_control)
}
/// Returns true if the volume is currently being compacted.
pub fn is_compacting(&self) -> bool {
self.is_compacting.load(Ordering::Acquire)
@@ -4248,6 +4258,15 @@ impl Volume {
let Some(_claim) = CompactionClaim::try_claim(&self.is_compacting) else {
return Ok(()); // already compacting, silently skip (matches Go)
};
// The reload would read the remote object through the compacted .idx.
if self.has_remote_file() {
let _ = fs::remove_file(self.file_name(".cpd"));
let _ = fs::remove_file(self.file_name(".cpx"));
return Err(VolumeError::Io(io::Error::other(format!(
"volume {} is tiered to remote storage, cannot commit compaction",
self.id
))));
}
self.do_commit_compact()
}