diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 65fae7cab..4c4918fde 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -1409,7 +1409,7 @@ impl VolumeServer for VolumeGrpcService { Local(std::fs::File), Remote(crate::storage::volume::RemoteDatFile), } - let reader = if v.has_remote_file { + let reader = if v.has_remote_file() { match v.remote_dat_file() { Some(r) => DatReader::Remote(r), None => { @@ -2239,7 +2239,7 @@ impl VolumeServer for VolumeGrpcService { compaction_revision: vol.super_block.compaction_revision as u32, collection: vol.collection.clone(), disk_type: store.locations[loc_idx].disk_type.to_string(), - volume_info: Some(vol.volume_info.clone()), + volume_info: Some(vol.volume_info().clone()), version: vol.version().0 as u32, }, )) @@ -4304,7 +4304,7 @@ impl VolumeServer for VolumeGrpcService { // Match Go's DiskFile check: if the .dat file is still local, we can // keep tiering it even when remote file entries already exist. - if volume_is_remote_only(&dat_path, vol.has_remote_file) { + if volume_is_remote_only(&dat_path, vol.has_remote_file()) { // Already on remote -- return empty stream (matches Go: returns nil) let stream = tokio_stream::empty(); return Ok(Response::new( @@ -4317,7 +4317,7 @@ impl VolumeServer for VolumeGrpcService { crate::remote_storage::s3_tier::backend_name_to_type_id( &req.destination_backend_name, ); - for rf in &vol.volume_info.files { + for rf in &vol.volume_info().files { if rf.backend_type == backend_type && rf.backend_id == backend_id { return Err(Status::already_exists(format!( "destination {} already exists", @@ -4426,16 +4426,18 @@ impl VolumeServer for VolumeGrpcService { { let mut store = state.store.write().unwrap(); if let Some((_, vol)) = store.find_volume_mut(vid) { - vol.volume_info.files.push(volume_server_pb::RemoteFile { - backend_type: backend_type.clone(), - backend_id: backend_id.clone(), - key, - offset: 0, - file_size: size, - modified_time: dat_modified_secs, - extension: ".dat".to_string(), - }); - vol.refresh_remote_write_mode().map_err(|e| { + vol.update_remote_files(|files| { + files.push(volume_server_pb::RemoteFile { + backend_type: backend_type.clone(), + backend_id: backend_id.clone(), + key, + offset: 0, + file_size: size, + modified_time: dat_modified_secs, + extension: ".dat".to_string(), + }) + }) + .map_err(|e| { Status::internal(format!( "volume {} failed to refresh write mode: {}", vid, e @@ -4533,7 +4535,7 @@ impl VolumeServer for VolumeGrpcService { } let remote_modified_secs = vol - .volume_info + .volume_info() .files .first() .map(|f| f.modified_time) @@ -4672,7 +4674,7 @@ impl VolumeServer for VolumeGrpcService { // Trim the remote reference, persist the .vif, and swap to the local // .dat on BOTH paths BEFORE deleting the remote object. After this the - // volume serves from local disk (has_remote_file = false), so a crash + // volume serves from local disk (has_remote_file() is false), so a crash // before the delete only leaks the remote object; the .vif must never // reference an object that has already been deleted. { @@ -4693,28 +4695,31 @@ impl VolumeServer for VolumeGrpcService { } // Snapshot the remote reference before dropping it: the - // refresh below can fail, and a half-applied transition + // refresh it triggers can fail, and a half-applied transition // leaves the volume claiming local while the remote backend // is still attached and the on-disk .vif still says remote // — a state a retry reads as "already on local disk" and // refuses to finish. - let removed_remote = if vol.volume_info.files.is_empty() { - None - } else { - Some(vol.volume_info.files.remove(0)) - }; - // Swaps the read-only sorted map out before the volume is - // published as writable; without it the first write would - // append to the local .dat and then fail to index. - if let Err(e) = vol.refresh_remote_write_mode() { - if let Some(remote) = removed_remote { - vol.volume_info.files.insert(0, remote); + // + // update_remote_files also swaps the read-only sorted map + // out before the volume is published as writable; without + // it the first write would append to the local .dat and + // then fail to index. + let mut removed_remote = None; + if let Err(e) = vol.update_remote_files(|files| { + if !files.is_empty() { + removed_remote = Some(files.remove(0)); } - // Put the derived flags and the needle map back where - // the restored reference says they belong. Best effort: - // if even this fails the volume stays pinned read-only, - // which is the safe end of the transition. - if let Err(restore_err) = vol.refresh_remote_write_mode() { + }) { + // Put the reference, the derived flags and the needle + // map back where they belong. Best effort: if even this + // fails the volume stays pinned read-only, which is the + // safe end of the transition. + if let Err(restore_err) = vol.update_remote_files(|files| { + if let Some(remote) = removed_remote { + files.insert(0, remote); + } + }) { tracing::warn!( volume_id = vid.0, error = %restore_err, @@ -6630,8 +6635,8 @@ mod tests { { 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!(!vol.has_remote_file()); + assert!(vol.volume_info().files.is_empty()); assert!(vol.has_data_backend()); } @@ -6714,10 +6719,10 @@ mod tests { let store = service.state.store.read().unwrap(); let (_, vol) = store.find_volume(VolumeId(1)).unwrap(); assert!( - vol.has_remote_file, + vol.has_remote_file(), "abandoned tier-down published the transition to local" ); - assert!(!vol.volume_info.files.is_empty()); + assert!(!vol.volume_info().files.is_empty()); } assert_eq!( delete_count.load(std::sync::atomic::Ordering::SeqCst), @@ -6761,8 +6766,8 @@ mod tests { { 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!(!vol.has_remote_file()); + assert!(vol.volume_info().files.is_empty()); assert!(vol.has_data_backend()); } diff --git a/seaweed-volume/src/server/heartbeat.rs b/seaweed-volume/src/server/heartbeat.rs index b57240272..cce60ccac 100644 --- a/seaweed-volume/src/server/heartbeat.rs +++ b/seaweed-volume/src/server/heartbeat.rs @@ -973,7 +973,7 @@ fn build_heartbeat_with_ec_status( // whose .dat legitimately lives in cloud storage. Only a present .dat is // cached for 30s; a missing one is re-checked every heartbeat so the volume // stays suppressed until the file returns. See issues/10004 - if vol.file_count() > 0 && !vol.has_remote_file { + if vol.file_count() > 0 && !vol.has_remote_file() { const DISK_CHECK_INTERVAL_NS: i64 = 30 * 1_000_000_000; let now_ns = SystemTime::now() .duration_since(UNIX_EPOCH) @@ -1688,8 +1688,9 @@ mod tests { { let (_, volume) = store.find_volume_mut(VolumeId(17)).unwrap(); volume.set_read_only().unwrap(); - volume.volume_info.files.push(Default::default()); - volume.refresh_remote_write_mode().unwrap(); + volume + .update_remote_files(|files| files.push(Default::default())) + .unwrap(); } let heartbeat = build_heartbeat(&test_config(), &mut store); @@ -2054,15 +2055,15 @@ mod tests { .unwrap(); let (_, volume) = store.find_volume_mut(VolumeId(71)).unwrap(); volume - .volume_info - .files - .push(crate::storage::volume::PbRemoteFile { - backend_type: "s3".to_string(), - backend_id: "archive".to_string(), - key: "volumes/71.dat".to_string(), - ..Default::default() - }); - volume.refresh_remote_write_mode().unwrap(); + .update_remote_files(|files| { + files.push(crate::storage::volume::PbRemoteFile { + backend_type: "s3".to_string(), + backend_id: "archive".to_string(), + key: "volumes/71.dat".to_string(), + ..Default::default() + }) + }) + .unwrap(); let heartbeat = build_heartbeat(&test_config(), &mut store); diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index 0364382c1..c12181833 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -710,10 +710,12 @@ pub struct Volume { io_error_quarantined: std::sync::atomic::AtomicBool, /// Protobuf VolumeInfo for tiered storage (.vif file). - pub volume_info: PbVolumeInfo, - - /// Whether this volume has a remote file reference. - pub has_remote_file: bool, + /// + /// Private: `volume_info.files` and the write mode derived from it are the + /// same fact, so outside callers edit the list through + /// [`Volume::update_remote_files`] and read it through + /// [`Volume::volume_info`]. + volume_info: PbVolumeInfo, } /// What a volume is created with beyond its id, directories and index kind: @@ -791,7 +793,6 @@ impl Volume { io_error_count: std::sync::atomic::AtomicI32::new(0), io_error_quarantined: std::sync::atomic::AtomicBool::new(false), volume_info: PbVolumeInfo::default(), - has_remote_file: false, }; v.load(true, true, preallocate, version)?; @@ -831,7 +832,6 @@ impl Volume { io_error_count: std::sync::atomic::AtomicI32::new(0), io_error_quarantined: std::sync::atomic::AtomicBool::new(false), volume_info: PbVolumeInfo::default(), - has_remote_file: false, } } @@ -884,7 +884,7 @@ impl Volume { let has_volume_info_file = self.load_vif()?; - if self.volume_info.read_only && !self.has_remote_file { + if self.volume_info.read_only && !self.has_remote_file() { if self.volume_info.read_only_can_delete { self.no_write_can_delete = true; } else { @@ -892,7 +892,7 @@ impl Volume { } } - if self.has_remote_file { + if self.has_remote_file() { self.load_remote_dat_file()?; if let Some(remote_file) = self.volume_info.files.first() { if remote_file.modified_time > 0 { @@ -962,7 +962,7 @@ impl Volume { // Match Go: v.volumeInfo.Version = uint32(v.SuperBlock.Version) self.volume_info.version = self.super_block.version.0 as u32; } - Err(e) if self.has_remote_file => { + Err(e) if self.has_remote_file() => { warn!( volume_id = self.id.0, error = %e, @@ -988,7 +988,7 @@ impl Volume { // A changed --dir.idx leaves the new directory without an index. // The .dat still holds every row, so rebuild rather than mount the // volume with every needle invisible. - if !self.has_remote_file + if !self.has_remote_file() && !Path::new(&self.file_name(".idx")).exists() && self.current_dat_file_size()? > SUPER_BLOCK_SIZE as u64 { @@ -1021,7 +1021,7 @@ impl Volume { // Match Go: CheckVolumeDataIntegrity after loading index (volume_loading.go L154-159) // Only for non-remote volumes (remote storage may not have local .dat) - if !self.has_remote_file { + if !self.has_remote_file() { if let Err(e) = self.check_volume_data_integrity() { self.no_write_or_delete = true; warn!( @@ -1369,7 +1369,7 @@ impl Volume { /// Clone the remote backend handle (cheap: an `Arc` plus key/size) so a /// caller can stream a tiered `.dat` from S3 after dropping the store lock. - /// Returns `None` for local volumes. `has_remote_file` implies `dat_file` + /// Returns `None` for local volumes. `has_remote_file()` implies `dat_file` /// is `None`, so this selects the same backend `read_dat_slice` would. pub(crate) fn remote_dat_file(&self) -> Option { self.remote_dat_file.clone() @@ -2216,7 +2216,7 @@ impl Volume { n.data = vec![]; n.append_at_ns = get_append_at_ns(self.last_append_at_ns); - let offset = if !self.has_remote_file { + let offset = if !self.has_remote_file() { // Normal volume: append tombstone to .dat file let (offset, _, _) = self.append_needle(n)?; offset @@ -2961,14 +2961,14 @@ impl Volume { can_delete: bool, persist: bool, ) -> Result<(), VolumeError> { - if can_delete && !self.has_remote_file { + if can_delete && !self.has_remote_file() { // deletes append tombstones to .idx; a read-only boot attached no writer self.attach_idx_writer_if_missing()?; } self.no_write_or_delete = !can_delete; if can_delete { self.no_write_can_delete = true; - } else if !self.has_remote_file { + } else if !self.has_remote_file() { // downgrading a canDelete mark; remote volumes keep their derived flag self.no_write_can_delete = false; } @@ -3013,7 +3013,7 @@ impl Volume { let was_no_write_can_delete = self.no_write_can_delete; self.no_write_or_delete = false; // Remote-tiered volumes must stay no_write_can_delete regardless of marks. - if !self.has_remote_file { + if !self.has_remote_file() { self.no_write_can_delete = false; } @@ -3064,6 +3064,37 @@ impl Volume { self.load_index() } + /// The tiered-storage VolumeInfo backing the .vif file. + pub fn volume_info(&self) -> &PbVolumeInfo { + &self.volume_info + } + + /// Whether this volume is backed by a remote file, i.e. whether the .vif + /// holds a remote reference. Derived rather than mirrored, so it can never + /// disagree with the reference list it describes. + pub fn has_remote_file(&self) -> bool { + !self.volume_info.files.is_empty() + } + + /// Edit the remote reference list and recompute everything derived from it. + /// + /// This is the only way to change `volume_info.files` from outside the + /// module: the write/delete mode and the needle map follow from whether the + /// volume is remote, so an edit that skipped the refresh would leave the + /// volume claiming a mode its .vif contradicts. + /// + /// Not transactional: on error `f` has already been applied and the volume + /// is left pinned read-only, which is the safe end of a half-finished tier + /// transition. A caller that needs the old list back restores it with a + /// second call. + pub fn update_remote_files( + &mut self, + f: impl FnOnce(&mut Vec), + ) -> Result<(), VolumeError> { + f(&mut self.volume_info.files); + self.refresh_remote_write_mode() + } + /// Recompute the Go-style write/delete mode from the current remote tier /// state, and bring the needle map in line with it — a volume that stops /// being remote also stops using the read-only sorted index. @@ -3071,9 +3102,8 @@ impl Volume { /// If the map cannot be rebuilt the volume is pinned read-only rather than /// published as writable with an index that rejects every put; a restart /// recovers it. - pub fn refresh_remote_write_mode(&mut self) -> Result<(), VolumeError> { - self.has_remote_file = !self.volume_info.files.is_empty(); - if self.has_remote_file { + fn refresh_remote_write_mode(&mut self) -> Result<(), VolumeError> { + if self.has_remote_file() { self.no_write_can_delete = true; self.no_write_or_delete = false; } else if !self.volume_info.read_only_can_delete { @@ -3125,7 +3155,7 @@ impl Volume { if self.volume_info.version == 0 { self.volume_info.version = Version::current().0 as u32; } - if !self.has_remote_file && self.volume_info.bytes_offset == 0 { + if !self.has_remote_file() && self.volume_info.bytes_offset == 0 { self.volume_info.bytes_offset = OFFSET_SIZE as u32; } if self.volume_info.bytes_offset != 0 @@ -3154,7 +3184,7 @@ impl Volume { if self.volume_info.version == 0 { self.volume_info.version = Version::current().0 as u32; } - if !self.has_remote_file && self.volume_info.bytes_offset == 0 { + if !self.has_remote_file() && self.volume_info.bytes_offset == 0 { self.volume_info.bytes_offset = OFFSET_SIZE as u32; } if self.volume_info.bytes_offset != 0 @@ -3201,7 +3231,7 @@ impl Volume { // remoteness-derived no_write_can_delete is not an operator mark; sync // volume_info so refresh_remote_write_mode sees a real mark let marked_can_delete = - self.no_write_can_delete && !self.has_remote_file && !self.no_write_or_delete; + self.no_write_can_delete && !self.has_remote_file() && !self.no_write_or_delete; self.volume_info.read_only = self.no_write_or_delete || marked_can_delete; self.volume_info.read_only_can_delete = marked_can_delete; let mut vif = VifVolumeInfo::from_pb(&self.volume_info); @@ -3222,7 +3252,7 @@ impl Volume { /// Matches Go's SaveVolumeInfo which computes ExpireAtSec from TTL. pub fn save_volume_info(&mut self) -> Result<(), VolumeError> { let marked_can_delete = - self.no_write_can_delete && !self.has_remote_file && !self.no_write_or_delete; + self.no_write_can_delete && !self.has_remote_file() && !self.no_write_or_delete; self.volume_info.read_only = self.no_write_or_delete || marked_can_delete; self.volume_info.read_only_can_delete = marked_can_delete; @@ -4307,7 +4337,7 @@ impl Volume { let (storage_name, storage_key) = self.remote_storage_name_key(); if !keep_remote_data - && self.has_remote_file + && self.has_remote_file() && !storage_name.is_empty() && !storage_key.is_empty() { @@ -7429,6 +7459,53 @@ mod tests { .unwrap() } + /// The remote reference and the derived write mode are the same fact, so a + /// caller must not be able to move one without the other. Every edit goes + /// through update_remote_files, which refreshes the mode on the way out. + #[test] + fn test_update_remote_files_refreshes_the_derived_write_mode() { + let tmp = TempDir::new().unwrap(); + let dir = tmp.path().to_str().unwrap(); + let mut v = make_test_volume(dir); + let file_size = v.dat_file_size().unwrap(); + + assert!(!v.has_remote_file()); + assert!(!v.is_no_write_can_delete()); + assert!(!v.is_no_write_or_delete()); + + v.update_remote_files(|files| { + files.push(PbRemoteFile { + backend_type: "s3".to_string(), + backend_id: "default".to_string(), + key: "remote-key".to_string(), + offset: 0, + file_size, + modified_time: 123, + extension: ".dat".to_string(), + }) + }) + .unwrap(); + + assert!(v.has_remote_file()); + assert!( + v.is_no_write_can_delete(), + "a remote-backed volume only serves deletes" + ); + assert!(!v.is_no_write_or_delete()); + + v.update_remote_files(|files| { + files.remove(0); + }) + .unwrap(); + + assert!(!v.has_remote_file()); + assert!( + !v.is_no_write_can_delete(), + "dropping the last remote reference publishes the volume as writable" + ); + assert!(!v.is_no_write_or_delete()); + } + // Tier-down clears the remote mode and publishes the volume as writable. The // map it booted with is the read-only sorted one, whose put always fails, so // without a rebuild the first write appends to .dat and then cannot be @@ -7460,18 +7537,21 @@ mod tests { // What the tier-up handler does once the .dat is uploaded: record the // remote reference and reconcile the mode. - v.volume_info.files.push(PbRemoteFile { - backend_type: "s3".to_string(), - backend_id: "vif_tierup_test".to_string(), - key: "remote-key".to_string(), - offset: 0, - file_size: v.dat_file_size().unwrap(), - modified_time: 123, - extension: ".dat".to_string(), - }); - v.refresh_remote_write_mode().unwrap(); + let file_size = v.dat_file_size().unwrap(); + v.update_remote_files(|files| { + files.push(PbRemoteFile { + backend_type: "s3".to_string(), + backend_id: "vif_tierup_test".to_string(), + key: "remote-key".to_string(), + offset: 0, + file_size, + modified_time: 123, + extension: ".dat".to_string(), + }) + }) + .unwrap(); - assert!(v.has_remote_file); + assert!(v.has_remote_file()); assert!( matches!(v.nm, Some(NeedleMap::SortedFile(_))), "tier-up must install the sorted map without waiting for a restart" @@ -7494,12 +7574,14 @@ mod tests { assert!(matches!(v.nm, Some(NeedleMap::SortedFile(_)))); // The tail of the tier-down handler, once the .dat is back on disk. - v.volume_info.files.remove(0); - v.refresh_remote_write_mode().unwrap(); + v.update_remote_files(|files| { + files.remove(0); + }) + .unwrap(); v.save_volume_info().unwrap(); v.open_local_dat_backend().unwrap(); - assert!(!v.has_remote_file); + assert!(!v.has_remote_file()); assert!( !v.is_read_only(), "tier-down should publish a writable volume" @@ -7919,7 +8001,7 @@ mod tests { ) .unwrap(); - assert!(v.has_remote_file); + assert!(v.has_remote_file()); let Some(NeedleMap::SortedFile(ref nm)) = v.nm else { panic!( "tiered volume should search the on-disk .sdx, got {:?}", @@ -8003,16 +8085,19 @@ mod tests { let dir = tmp.path().to_str().unwrap(); let mut v = make_test_volume(dir); - v.volume_info.files.push(PbRemoteFile { - backend_type: "s3".to_string(), - backend_id: "default".to_string(), - key: "remote-key".to_string(), - offset: 0, - file_size: v.dat_file_size().unwrap(), - modified_time: 123, - extension: ".dat".to_string(), - }); - v.refresh_remote_write_mode().unwrap(); + let file_size = v.dat_file_size().unwrap(); + v.update_remote_files(|files| { + files.push(PbRemoteFile { + backend_type: "s3".to_string(), + backend_id: "default".to_string(), + key: "remote-key".to_string(), + offset: 0, + file_size, + modified_time: 123, + extension: ".dat".to_string(), + }) + }) + .unwrap(); v.set_writable().unwrap(); assert!(v.is_read_only()); @@ -8447,7 +8532,7 @@ mod tests { ) .unwrap(); - assert!(v.has_remote_file); + assert!(v.has_remote_file()); assert!(v.dat_file.is_none()); assert!(v.remote_dat_file.is_some());