diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index c8d2783ad..a139055ed 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -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 { + 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 { + 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, + ) -> Option { + 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, 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, + delete_count: Arc, + } + + /// 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("bucket-a\ + kupload-1\ + " + .to_string()) + } + Method::PUT if query.contains("partNumber") => { + let _ = parked_tx.send(()); + release.acquire().await.unwrap().forget(); + xml(String::new()) + } + Method::POST => { + xml("bucket-a\ + k\"etag\"\ + " + .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. diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index 4e6ee064e..0f8322a8d 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -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 { + self.data_file_access_control.clone() + } + + pub(crate) fn is_instance(&self, instance: &Arc) -> 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() }