From 67691a1eeae9b87726ff09c8ceb317afb90ac3dd Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Sun, 27 Sep 2026 15:12:55 +0300 Subject: [PATCH] volume server: split volume_copy into phases and type the delete-after-status gate (#11485) volume_copy was one ~400-line handler, and the rule that an existing local replica is deleted only after the source's ReadVolumeFileStatus succeeded was held by statement order alone. The keep_remote_data=true that the pre-copy delete and the failed-copy rollback must share was kept in sync by a comment pointing from one to the other. The handler is now a ~60-line orchestrator over connect_to_copy_source, SourceVolumeStatus::fetch, delete_existing_replica, plan_copy_destination and a VolumeCopyJob whose run() drives preallocate_dat, transfer_files, finish_copied_files and mount_and_reply, with cleanup_failed_copy on error. delete_existing_replica takes a &SourceVolumeStatus, which only fetch can construct (private field in a child module), so the delete cannot be called before the status RPC. Both deletes go through delete_replica_keep_remote. Pure refactor: call order, status codes and messages, cancellation checks, throttling, progress reports and cleanup are unchanged. Co-authored-by: Claude Opus 5.5 (1M context) Co-authored-by: Chris Lu --- seaweed-volume/src/server/grpc_server.rs | 1079 +++++++++++++--------- 1 file changed, 631 insertions(+), 448 deletions(-) diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 561879bc8..1b82443bf 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -22,8 +22,8 @@ use crate::storage::types::*; use crate::storage::volume::VolumeSpec; use super::grpc_client::{ - GrpcDialOptions, connect_channel, connect_channel_guarded, filer_client, master_client, - volume_server_client, + GrpcDialOptions, VolumeServerGrpcClient, connect_channel, connect_channel_guarded, + filer_client, master_client, volume_server_client, }; use super::volume_server::VolumeServerState; @@ -1823,143 +1823,33 @@ impl VolumeServer for VolumeGrpcService { })?; } - // A pre-existing local replica is NOT deleted up front. Deleting before - // the source is confirmed reachable destroys a healthy copy on a - // transient source outage (and, on retry, can lose the volume - // entirely). The delete is deferred until read_volume_file_status below - // proves the source holds the volume; readability alone is the gate. + // A pre-existing local replica is NOT deleted up front; see + // delete_existing_replica for why and when. let had_existing_volume = { let store = self.state.store.read().unwrap(); store.find_volume(vid).is_some() }; - // Parse source_data_node address: "ip:port.grpcPort" or "ip:port" (grpc = port + 10000) - let source = &req.source_data_node; - let grpc_addr = parse_grpc_address(source).map_err(|e| { - Status::internal(format!( - "VolumeCopy volume {} invalid source_data_node {}: {}", - vid, source, e - )) - })?; - - let channel = connect_channel_guarded( - &grpc_addr, - source, - self.state.outgoing_grpc_tls.as_ref(), - GrpcDialOptions::stream(), - self.state.allow_untrusted_remote_endpoints, - ) - .await - .map_err(|e| { - Status::internal(format!( - "VolumeCopy volume {} connect to {}: {}", - vid, grpc_addr, e - )) - })?; - - let mut client = volume_server_client(channel); - - // The source's record counts before and after the copy decide whether - // the copied replica's counts can be checked. Without the "before" - // snapshot the check is skipped, not the copy. - let source_status_before = match client - .volume_status(volume_server_pb::VolumeStatusRequest { - volume_id: req.volume_id, - }) - .await - { - Ok(resp) => Some(resp.into_inner()), - Err(e) => { - tracing::warn!( - "failed to read source volume {} status before copy; skip record count validation: {}", - vid, - e - ); - None - } - }; - - // Get file status from source - let vol_info = client - .read_volume_file_status(volume_server_pb::ReadVolumeFileStatusRequest { - volume_id: req.volume_id, - }) - .await - .map_err(|e| Status::internal(format!("read volume file status failed, {}", e)))? - .into_inner(); - - let requested_disk_type = if !req.disk_type.is_empty() { - DiskType::from_string(&req.disk_type) - } else { - DiskType::from_string(&vol_info.disk_type) - }; - - let has_remote_dat = vol_info - .volume_info - .as_ref() - .map(|vi| !vi.files.is_empty()) - .unwrap_or(false); - // a remote-backed volume only lands its .idx/.vif locally; the .dat stays in the tier - let needed_space = if has_remote_dat { - vol_info.idx_file_size - } else { - vol_info.dat_file_size - }; - - // Find a free disk location using Go's Store.FindFreeLocation semantics. - // The slot of the replica being replaced counts as free. - let (data_base, idx_base, selected_disk_type) = { - let mut store = self.state.store.write().unwrap(); - let Some(loc_idx) = store.find_free_location_predicate( - |loc| { - loc.disk_type == requested_disk_type - && loc.available_space.load(Ordering::Relaxed) > needed_space - }, - Some(vid), - ) else { - return Err(Status::internal(format!( - "no space left {}", - requested_disk_type.readable_string() - ))); - }; - let loc = &store.locations[loc_idx]; - let selected = ( - loc.directory.clone(), - loc.idx_directory.clone(), - loc.disk_type.clone(), - ); - - // Source is reachable and a destination is reserved: only now is - // it safe to drop an existing local replica before overwriting its - // files. - if had_existing_volume { - // keep remote data: the inbound copy carries a .vif that may - // point at the same cloud-tier object the existing volume - // references. - store.delete_volume(vid, false, false, true).map_err(|e| { - Status::internal(format!("failed to delete existing volume {}: {}", vid, e)) - })?; - } - selected - }; - if had_existing_volume { - self.state.volume_state_notify.notify_one(); - } - - let data_base_name = - crate::storage::volume::volume_file_name(&data_base, &vol_info.collection, vid); - let idx_base_name = - crate::storage::volume::volume_file_name(&idx_base, &vol_info.collection, vid); - - // Write a .note file to indicate copy in progress. A leftover note - // fails the volume load on restart, so a write failure must abort. - let note_path = format!("{}.note", data_base_name); - std::fs::write(¬e_path, format!("copying from {}", source)) - .map_err(|e| Status::internal(format!("write .note for volume {}: {}", vid, e)))?; + let mut client = self + .connect_to_copy_source(vid, &req.source_data_node) + .await?; + let source_status_before = + fetch_source_status_before(&mut client, req.volume_id, vid).await; + let source = SourceVolumeStatus::fetch(&mut client, req.volume_id).await?; + let dest = self.plan_copy_destination(&req, vid, &source, had_existing_volume)?; + dest.write_note(vid, &req.source_data_node)?; let (tx, rx) = tokio::sync::mpsc::channel::>(16); - let state = self.state.clone(); + let job = VolumeCopyJob { + state: self.state.clone(), + req, + vid, + source, + source_status_before, + dest, + tx, + }; tokio::spawn(async move { // Tracks whether mount_volume succeeded so the error branch can @@ -1967,319 +1857,8 @@ impl VolumeServer for VolumeGrpcService { // races a departing caller is the one case where the existing // file-only cleanup leaves the volume loaded in memory. let mut mounted = false; - let result = async { - // Nothing below is worth doing for a caller that has already - // gone: the transfer would spend the source's bandwidth and the - // destination's disk on a volume nobody will take delivery of. - if tx.is_closed() { - return Err(Status::cancelled(format!( - "volume {} copy cancelled by caller", - vid - ))); - } - - let report_interval: i64 = 128 * 1024 * 1024; - let mut next_report_target: i64 = report_interval; - let io_byte_per_second = if req.io_byte_per_second > 0 { - req.io_byte_per_second - } else { - state.maintenance_byte_per_second - }; - let mut throttler = WriteThrottler::new(io_byte_per_second); - - // Query master for preallocation settings (matching Go VolumeCopy behavior). - let mut preallocate_size: i64 = 0; - if !has_remote_dat { - let grpc_addr = super::heartbeat::to_grpc_address(&state.master_url); - // Race the master-configuration RPC against the caller's - // response channel: a stalled master (or a slow leader - // election) would otherwise hold the task and its .note - // past a departing caller, since the per-chunk checks in - // copy_file_from_source are never reached. - let config = tokio::select! { - res = super::heartbeat::try_get_master_configuration( - &grpc_addr, - state.outgoing_grpc_tls.as_ref(), - ) => res, - _ = tx.closed() => { - return Err(Status::cancelled(format!( - "volume {} copy cancelled by caller", - vid - ))); - } - }; - match config { - Ok(resp) => { - if resp.volume_preallocate { - preallocate_size = resp.volume_size_limit_m_b as i64 * 1024 * 1024; - } - } - Err(e) => { - tracing::warn!("get master {} configuration: {}", state.master_url, e); - } - } - - if preallocate_size > 0 { - let dat_path = format!("{}.dat", data_base_name); - let file = std::fs::File::create(&dat_path).map_err(|e| { - Status::internal(format!( - "create preallocated volume file {}: {}", - dat_path, e - )) - })?; - file.set_len(preallocate_size as u64).map_err(|e| { - Status::internal(format!("preallocate volume file {}: {}", dat_path, e)) - })?; - } - } - - // Copy .dat file - let mut progress = CopyProgress { - tx: &tx, - next_report_target: &mut next_report_target, - report_interval, - throttler: &mut throttler, - }; - if !has_remote_dat { - let dat_path = format!("{}.dat", data_base_name); - let dat_modified_ts_ns = copy_file_from_source( - &mut client, - &CopyFileSpec { - is_ec_volume: false, - collection: &req.collection, - volume_id: req.volume_id, - compaction_revision: vol_info.compaction_revision, - stop_offset: vol_info.dat_file_size, - dest_path: &dat_path, - ext: ".dat", - is_append: false, - ignore_source_not_found: true, - report_progress: true, - }, - &mut progress, - ) - .await?; - if dat_modified_ts_ns > 0 { - let _ = set_file_mtime(&dat_path, dat_modified_ts_ns); - } - } - - // Copy .idx file - let idx_path = format!("{}.idx", idx_base_name); - let idx_modified_ts_ns = copy_file_from_source( - &mut client, - &CopyFileSpec { - is_ec_volume: false, - collection: &req.collection, - volume_id: req.volume_id, - compaction_revision: vol_info.compaction_revision, - stop_offset: vol_info.idx_file_size, - dest_path: &idx_path, - ext: ".idx", - is_append: false, - ignore_source_not_found: false, - report_progress: false, - }, - &mut progress, - ) - .await?; - if idx_modified_ts_ns > 0 { - let _ = set_file_mtime(&idx_path, idx_modified_ts_ns); - } - - // Copy .vif file (ignore if not found on source) - let vif_path = format!("{}.vif", data_base_name); - let vif_modified_ts_ns = copy_file_from_source( - &mut client, - &CopyFileSpec { - is_ec_volume: false, - collection: &req.collection, - volume_id: req.volume_id, - compaction_revision: vol_info.compaction_revision, - stop_offset: 1024 * 1024, - dest_path: &vif_path, - ext: ".vif", - is_append: false, - ignore_source_not_found: true, - report_progress: false, - }, - &mut progress, - ) - .await?; - if vif_modified_ts_ns > 0 { - let _ = set_file_mtime(&vif_path, vif_modified_ts_ns); - } - - // Remove the .note file. A leftover note fails the load on the - // next restart, so a removal failure must fail the copy. - if let Err(e) = std::fs::remove_file(¬e_path) - && e.kind() != std::io::ErrorKind::NotFound - { - return Err(Status::internal(format!( - "remove .note for volume {}: {}", - vid, e - ))); - } - - // Go passes stream.Context() here, so a departing caller - // cancels the call. Race it against the response channel the - // same way: a stalled source would otherwise hold the copied, - // unmounted files past the caller. - let source_status_after = tokio::select! { - res = client.volume_status(volume_server_pb::VolumeStatusRequest { - volume_id: req.volume_id, - }) => res.map_err(|e| { - Status::internal(format!( - "read source volume {} status after copy failed: {}", - vid, e - )) - })? - .into_inner(), - _ = tx.closed() => { - return Err(Status::cancelled(format!( - "volume {} copy cancelled by caller", - vid - ))); - } - }; - - // Verify file sizes - if !has_remote_dat { - let dat_path = format!("{}.dat", data_base_name); - check_copy_file_size(&dat_path, vol_info.dat_file_size)?; - } - if vol_info.idx_file_size > 0 { - check_copy_file_size(&idx_path, vol_info.idx_file_size)?; - } - - // Find last_append_at_ns from copied files - let last_append_at_ns = if !has_remote_dat { - find_last_append_at_ns( - &idx_path, - &format!("{}.dat", data_base_name), - vol_info.version, - ) - .unwrap_or(vol_info.dat_file_timestamp_seconds * 1_000_000_000) - } else { - vol_info.dat_file_timestamp_seconds * 1_000_000_000 - }; - - // An orphaned mount is how an abandoned copy does lasting - // damage: the destination carries that volume's index cache for - // the lifetime of the process, and under replication=000 the - // cluster is left holding one volume id on two servers, both - // writable, which concurrent writes can diverge. - if tx.is_closed() { - return Err(Status::cancelled(format!( - "volume {} copy cancelled by caller before mount", - vid - ))); - } - - let validate_counts = - copy_counts_stable(source_status_before.as_ref(), &source_status_after); - if !validate_counts { - tracing::debug!( - "source volume {} changed during copy; skip record count validation", - vid - ); - } - - // Mount and validate the volume under one store lock, so a - // replica that fails validation is unloaded before the next - // heartbeat can announce it. - { - let mut store = state.store.write().unwrap(); - store - .mount_volume(vid, &vol_info.collection, selected_disk_type) - .map_err(|e| { - Status::internal(format!( - "failed to mount or validate volume {}: {}", - vid, e - )) - })?; - if validate_counts - && let Some((_, v)) = store.find_volume(vid) - && let Err(e) = check_copy_counts( - &source_status_after, - v.file_count() as u64, - v.deleted_count() as u64, - ) - { - store.unmount_volume(vid); - return Err(Status::internal(format!( - "failed to mount or validate volume {}: {}", - vid, e - ))); - } - mounted = true; - } - state.volume_state_notify.notify_one(); - - // Send final response with last_append_at_ns. A failed send - // means the caller is gone: the mount above raced a departing - // receiver (the pre-mount is_closed() check cannot close that - // window), and leaving the volume mounted is exactly the orphan - // this PR prevents. Surface it as Cancelled so the error branch - // unmounts and deletes the replica it just created. - if tx - .send(Ok(volume_server_pb::VolumeCopyResponse { - last_append_at_ns, - processed_bytes: 0, - })) - .await - .is_err() - { - return Err(Status::cancelled(format!( - "volume {} copy cancelled by caller after mount", - vid - ))); - } - - Ok::<(), Status>(()) - } - .await; - - if let Err(e) = result { - // An abandoned copy is otherwise invisible here: the error goes - // to a channel nobody is reading. Logging it gives the operator - // the cause behind the balancer's "delete that copy, then re-run - // the move". - if e.code() == tonic::Code::Cancelled { - tracing::info!( - "volume {} copy from {} abandoned by its caller, discarding the partial copy", - vid, - req.source_data_node - ); - } else { - tracing::warn!( - "volume {} copy from {} failed: {}", - vid, - req.source_data_node, - e - ); - } - // Clean up on error. If the volume was mounted (the after-mount - // race), delete_volume unmounts it from the store AND removes - // the .dat/.idx/.vif in one step; the remove_file calls below - // cover the never-mounted partial-file case and are harmless - // no-ops when delete_volume already removed the files. - // - // keep_remote_data=true matches the pre-spawn delete_volume at - // the top of volume_copy: a remote-tier copy's .vif points at - // the same cloud object the source replica references, so - // destroying the abandoned destination with keep_remote_data= - // false would delete the source's remote data. - if mounted { - let mut store = state.store.write().unwrap(); - let _ = store.delete_volume(vid, false, false, true); - state.volume_state_notify.notify_one(); - } - let _ = std::fs::remove_file(format!("{}.dat", data_base_name)); - let _ = std::fs::remove_file(format!("{}.idx", idx_base_name)); - let _ = std::fs::remove_file(format!("{}.vif", data_base_name)); - let _ = std::fs::remove_file(¬e_path); - let _ = tx.send(Err(e)).await; + if let Err(e) = job.run(&mut client, &mut mounted).await { + job.cleanup_failed_copy(e, mounted).await; } }); @@ -5520,6 +5099,602 @@ impl VolumeServer for VolumeGrpcService { } } +/// The copy source's ReadVolumeFileStatus answer. `fetch` is the only +/// constructor, so holding one proves the source is reachable and holds the +/// volume. Readability alone is the gate: size/count comparisons invert after +/// divergent compaction and would block valid re-replication. +mod copy_source { + use tonic::Status; + + use crate::pb::volume_server_pb; + use crate::server::grpc_client::VolumeServerGrpcClient; + + pub(super) struct SourceVolumeStatus(volume_server_pb::ReadVolumeFileStatusResponse); + + impl SourceVolumeStatus { + pub(super) async fn fetch( + client: &mut VolumeServerGrpcClient, + volume_id: u32, + ) -> Result { + // Get file status from source + let vol_info = client + .read_volume_file_status(volume_server_pb::ReadVolumeFileStatusRequest { + volume_id, + }) + .await + .map_err(|e| Status::internal(format!("read volume file status failed, {}", e)))? + .into_inner(); + Ok(Self(vol_info)) + } + } + + impl std::ops::Deref for SourceVolumeStatus { + type Target = volume_server_pb::ReadVolumeFileStatusResponse; + + fn deref(&self) -> &Self::Target { + &self.0 + } + } +} + +use copy_source::SourceVolumeStatus; + +/// Deletes a VolumeCopy destination replica but keeps its remote data: its +/// .vif may point at the same cloud-tier object the source replica references, +/// so dropping that object would destroy the source's data. The pre-copy delete +/// and the failed-copy rollback both go through here. +fn delete_replica_keep_remote( + store: &mut crate::storage::store::Store, + vid: VolumeId, +) -> Result<(), crate::storage::volume::VolumeError> { + store.delete_volume(vid, false, false, true) +} + +/// Drops a pre-existing local replica before the copy overwrites its files. +/// Deleting before the source is confirmed reachable destroys a healthy copy +/// on a transient source outage (and, on retry, can lose the volume +/// entirely), so this takes the source's status as proof. The caller holds +/// the store write lock so the delete is atomic with the destination pick. +fn delete_existing_replica( + store: &mut crate::storage::store::Store, + vid: VolumeId, + _source: &SourceVolumeStatus, +) -> Result<(), Status> { + delete_replica_keep_remote(store, vid) + .map_err(|e| Status::internal(format!("failed to delete existing volume {}: {}", vid, e))) +} + +/// The source's record counts before and after the copy decide whether the +/// copied replica's counts can be checked. Without the "before" snapshot the +/// check is skipped, not the copy. +async fn fetch_source_status_before( + client: &mut VolumeServerGrpcClient, + volume_id: u32, + vid: VolumeId, +) -> Option { + match client + .volume_status(volume_server_pb::VolumeStatusRequest { volume_id }) + .await + { + Ok(resp) => Some(resp.into_inner()), + Err(e) => { + tracing::warn!( + "failed to read source volume {} status before copy; skip record count validation: {}", + vid, + e + ); + None + } + } +} + +/// Go passes stream.Context() here, so a departing caller cancels the call. +/// Race it against the response channel the same way: a stalled source would +/// otherwise hold the copied, unmounted files past the caller. +async fn fetch_source_status_after( + client: &mut VolumeServerGrpcClient, + volume_id: u32, + vid: VolumeId, + tx: &tokio::sync::mpsc::Sender>, +) -> Result { + tokio::select! { + res = client.volume_status(volume_server_pb::VolumeStatusRequest { + volume_id, + }) => res + .map(|r| r.into_inner()) + .map_err(|e| { + Status::internal(format!( + "read source volume {} status after copy failed: {}", + vid, e + )) + }), + _ = tx.closed() => Err(Status::cancelled(format!( + "volume {} copy cancelled by caller", + vid + ))), + } +} + +/// Where a VolumeCopy lands on this server. +struct CopyDestination { + data_base_name: String, + idx_base_name: String, + disk_type: DiskType, + /// A remote-backed volume only lands its .idx/.vif locally; the .dat stays in the tier. + has_remote_dat: bool, +} + +impl CopyDestination { + fn note_path(&self) -> String { + format!("{}.note", self.data_base_name) + } + + /// Write a .note file to indicate copy in progress. A leftover note + /// fails the volume load on restart, so a write failure must abort. + fn write_note(&self, vid: VolumeId, source: &str) -> Result<(), Status> { + let note_path = self.note_path(); + std::fs::write(¬e_path, format!("copying from {}", source)) + .map_err(|e| Status::internal(format!("write .note for volume {}: {}", vid, e)))?; + Ok(()) + } +} + +impl VolumeGrpcService { + async fn connect_to_copy_source( + &self, + vid: VolumeId, + source: &str, + ) -> Result { + // Parse source_data_node address: "ip:port.grpcPort" or "ip:port" (grpc = port + 10000) + let grpc_addr = parse_grpc_address(source).map_err(|e| { + Status::internal(format!( + "VolumeCopy volume {} invalid source_data_node {}: {}", + vid, source, e + )) + })?; + + let channel = connect_channel_guarded( + &grpc_addr, + source, + self.state.outgoing_grpc_tls.as_ref(), + GrpcDialOptions::stream(), + self.state.allow_untrusted_remote_endpoints, + ) + .await + .map_err(|e| { + Status::internal(format!( + "VolumeCopy volume {} connect to {}: {}", + vid, grpc_addr, e + )) + })?; + + Ok(volume_server_client(channel)) + } + + /// Picks the destination disk and, under the same write lock, drops a + /// pre-existing local replica so its slot counts as free for the pick. + fn plan_copy_destination( + &self, + req: &volume_server_pb::VolumeCopyRequest, + vid: VolumeId, + vol_info: &SourceVolumeStatus, + had_existing_volume: bool, + ) -> Result { + let requested_disk_type = if !req.disk_type.is_empty() { + DiskType::from_string(&req.disk_type) + } else { + DiskType::from_string(&vol_info.disk_type) + }; + + let has_remote_dat = vol_info + .volume_info + .as_ref() + .map(|vi| !vi.files.is_empty()) + .unwrap_or(false); + let needed_space = if has_remote_dat { + vol_info.idx_file_size + } else { + vol_info.dat_file_size + }; + + // Find a free disk location using Go's Store.FindFreeLocation + // semantics; the slot of the replica being replaced counts as free. + let (data_base, idx_base, selected_disk_type) = { + let mut store = self.state.store.write().unwrap(); + let Some(loc_idx) = store.find_free_location_predicate( + |loc| { + loc.disk_type == requested_disk_type + && loc.available_space.load(Ordering::Relaxed) > needed_space + }, + Some(vid), + ) else { + return Err(Status::internal(format!( + "no space left {}", + requested_disk_type.readable_string() + ))); + }; + let loc = &store.locations[loc_idx]; + let selected = ( + loc.directory.clone(), + loc.idx_directory.clone(), + loc.disk_type.clone(), + ); + if had_existing_volume { + delete_existing_replica(&mut store, vid, vol_info)?; + } + selected + }; + if had_existing_volume { + self.state.volume_state_notify.notify_one(); + } + + let data_base_name = + crate::storage::volume::volume_file_name(&data_base, &vol_info.collection, vid); + let idx_base_name = + crate::storage::volume::volume_file_name(&idx_base, &vol_info.collection, vid); + + Ok(CopyDestination { + data_base_name, + idx_base_name, + disk_type: selected_disk_type, + has_remote_dat, + }) + } +} + +/// The part of a VolumeCopy that runs after the handler has returned its stream. +struct VolumeCopyJob { + state: Arc, + req: volume_server_pb::VolumeCopyRequest, + vid: VolumeId, + source: SourceVolumeStatus, + source_status_before: Option, + dest: CopyDestination, + tx: tokio::sync::mpsc::Sender>, +} + +impl VolumeCopyJob { + async fn run( + &self, + client: &mut VolumeServerGrpcClient, + mounted: &mut bool, + ) -> Result<(), Status> { + let (state, req, vid, tx) = (&self.state, &self.req, self.vid, &self.tx); + // Nothing below is worth doing for a caller that has already + // gone: the transfer would spend the source's bandwidth and the + // destination's disk on a volume nobody will take delivery of. + if tx.is_closed() { + return Err(Status::cancelled(format!( + "volume {} copy cancelled by caller", + vid + ))); + } + + let report_interval: i64 = 128 * 1024 * 1024; + let mut next_report_target: i64 = report_interval; + let io_byte_per_second = if req.io_byte_per_second > 0 { + req.io_byte_per_second + } else { + state.maintenance_byte_per_second + }; + let mut throttler = WriteThrottler::new(io_byte_per_second); + + self.preallocate_dat().await?; + + let mut progress = CopyProgress { + tx, + next_report_target: &mut next_report_target, + report_interval, + throttler: &mut throttler, + }; + self.transfer_files(client, &mut progress).await?; + + let source_status_after = + fetch_source_status_after(client, self.req.volume_id, vid, tx).await?; + let last_append_at_ns = self.finish_copied_files()?; + self.mount_and_reply(last_append_at_ns, &source_status_after, mounted) + .await + } + + /// Query master for preallocation settings (matching Go VolumeCopy behavior). + async fn preallocate_dat(&self) -> Result<(), Status> { + let (state, vid, tx) = (&self.state, self.vid, &self.tx); + let (data_base_name, has_remote_dat) = + (&self.dest.data_base_name, self.dest.has_remote_dat); + let mut preallocate_size: i64 = 0; + if !has_remote_dat { + let grpc_addr = super::heartbeat::to_grpc_address(&state.master_url); + // Race the master-configuration RPC against the caller's + // response channel: a stalled master (or a slow leader + // election) would otherwise hold the task and its .note + // past a departing caller, since the per-chunk checks in + // copy_file_from_source are never reached. + let config = tokio::select! { + res = super::heartbeat::try_get_master_configuration( + &grpc_addr, + state.outgoing_grpc_tls.as_ref(), + ) => res, + _ = tx.closed() => { + return Err(Status::cancelled(format!( + "volume {} copy cancelled by caller", + vid + ))); + } + }; + match config { + Ok(resp) => { + if resp.volume_preallocate { + preallocate_size = resp.volume_size_limit_m_b as i64 * 1024 * 1024; + } + } + Err(e) => { + tracing::warn!("get master {} configuration: {}", state.master_url, e); + } + } + + if preallocate_size > 0 { + let dat_path = format!("{}.dat", data_base_name); + let file = std::fs::File::create(&dat_path).map_err(|e| { + Status::internal(format!( + "create preallocated volume file {}: {}", + dat_path, e + )) + })?; + file.set_len(preallocate_size as u64).map_err(|e| { + Status::internal(format!("preallocate volume file {}: {}", dat_path, e)) + })?; + } + } + Ok(()) + } + + /// Pull the .dat (unless it stays in the remote tier), .idx and .vif. + async fn transfer_files( + &self, + client: &mut VolumeServerGrpcClient, + progress: &mut CopyProgress<'_>, + ) -> Result<(), Status> { + let (req, vol_info) = (&self.req, &self.source); + let (data_base_name, idx_base_name, has_remote_dat) = ( + &self.dest.data_base_name, + &self.dest.idx_base_name, + self.dest.has_remote_dat, + ); + // Copy .dat file + if !has_remote_dat { + let dat_path = format!("{}.dat", data_base_name); + let dat_modified_ts_ns = copy_file_from_source( + client, + &CopyFileSpec { + is_ec_volume: false, + collection: &req.collection, + volume_id: req.volume_id, + compaction_revision: vol_info.compaction_revision, + stop_offset: vol_info.dat_file_size, + dest_path: &dat_path, + ext: ".dat", + is_append: false, + ignore_source_not_found: true, + report_progress: true, + }, + progress, + ) + .await?; + if dat_modified_ts_ns > 0 { + let _ = set_file_mtime(&dat_path, dat_modified_ts_ns); + } + } + + // Copy .idx file + let idx_path = format!("{}.idx", idx_base_name); + let idx_modified_ts_ns = copy_file_from_source( + client, + &CopyFileSpec { + is_ec_volume: false, + collection: &req.collection, + volume_id: req.volume_id, + compaction_revision: vol_info.compaction_revision, + stop_offset: vol_info.idx_file_size, + dest_path: &idx_path, + ext: ".idx", + is_append: false, + ignore_source_not_found: false, + report_progress: false, + }, + progress, + ) + .await?; + if idx_modified_ts_ns > 0 { + let _ = set_file_mtime(&idx_path, idx_modified_ts_ns); + } + + // Copy .vif file (ignore if not found on source) + let vif_path = format!("{}.vif", data_base_name); + let vif_modified_ts_ns = copy_file_from_source( + client, + &CopyFileSpec { + is_ec_volume: false, + collection: &req.collection, + volume_id: req.volume_id, + compaction_revision: vol_info.compaction_revision, + stop_offset: 1024 * 1024, + dest_path: &vif_path, + ext: ".vif", + is_append: false, + ignore_source_not_found: true, + report_progress: false, + }, + progress, + ) + .await?; + if vif_modified_ts_ns > 0 { + let _ = set_file_mtime(&vif_path, vif_modified_ts_ns); + } + Ok(()) + } + + /// Clear the .note, verify the copied sizes and return last_append_at_ns. + fn finish_copied_files(&self) -> Result { + let (vid, vol_info, note_path) = (self.vid, &self.source, self.dest.note_path()); + let (data_base_name, has_remote_dat) = + (&self.dest.data_base_name, self.dest.has_remote_dat); + let idx_path = format!("{}.idx", self.dest.idx_base_name); + // Remove the .note file. A leftover note fails the load on the + // next restart, so a removal failure must fail the copy. + if let Err(e) = std::fs::remove_file(¬e_path) + && e.kind() != std::io::ErrorKind::NotFound + { + return Err(Status::internal(format!( + "remove .note for volume {}: {}", + vid, e + ))); + } + + // Verify file sizes + if !has_remote_dat { + let dat_path = format!("{}.dat", data_base_name); + check_copy_file_size(&dat_path, vol_info.dat_file_size)?; + } + if vol_info.idx_file_size > 0 { + check_copy_file_size(&idx_path, vol_info.idx_file_size)?; + } + + // Find last_append_at_ns from copied files + let last_append_at_ns = if !has_remote_dat { + find_last_append_at_ns( + &idx_path, + &format!("{}.dat", data_base_name), + vol_info.version, + ) + .unwrap_or(vol_info.dat_file_timestamp_seconds * 1_000_000_000) + } else { + vol_info.dat_file_timestamp_seconds * 1_000_000_000 + }; + Ok(last_append_at_ns) + } + + async fn mount_and_reply( + &self, + last_append_at_ns: u64, + source_status_after: &volume_server_pb::VolumeStatusResponse, + mounted: &mut bool, + ) -> Result<(), Status> { + let (state, vid, vol_info, tx) = (&self.state, self.vid, &self.source, &self.tx); + let selected_disk_type = self.dest.disk_type.clone(); + // An orphaned mount is how an abandoned copy does lasting + // damage: the destination carries that volume's index cache for + // the lifetime of the process, and under replication=000 the + // cluster is left holding one volume id on two servers, both + // writable, which concurrent writes can diverge. + if tx.is_closed() { + return Err(Status::cancelled(format!( + "volume {} copy cancelled by caller before mount", + vid + ))); + } + + let validate_counts = + copy_counts_stable(self.source_status_before.as_ref(), source_status_after); + if !validate_counts { + tracing::debug!( + "source volume {} changed during copy; skip record count validation", + vid + ); + } + + // Mount and validate the volume under one store lock, so a replica that + // fails validation is unloaded before the next heartbeat can announce it. + { + let mut store = state.store.write().unwrap(); + store + .mount_volume(vid, &vol_info.collection, selected_disk_type) + .map_err(|e| { + Status::internal(format!("failed to mount or validate volume {}: {}", vid, e)) + })?; + if validate_counts + && let Some((_, v)) = store.find_volume(vid) + && let Err(e) = check_copy_counts( + source_status_after, + v.file_count() as u64, + v.deleted_count() as u64, + ) + { + store.unmount_volume(vid); + return Err(Status::internal(format!( + "failed to mount or validate volume {}: {}", + vid, e + ))); + } + *mounted = true; + } + state.volume_state_notify.notify_one(); + + // Send final response with last_append_at_ns. A failed send + // means the caller is gone: the mount above raced a departing + // receiver (the pre-mount is_closed() check cannot close that + // window), and leaving the volume mounted is exactly the orphan + // this PR prevents. Surface it as Cancelled so the error branch + // unmounts and deletes the replica it just created. + if tx + .send(Ok(volume_server_pb::VolumeCopyResponse { + last_append_at_ns, + processed_bytes: 0, + })) + .await + .is_err() + { + return Err(Status::cancelled(format!( + "volume {} copy cancelled by caller after mount", + vid + ))); + } + + Ok(()) + } + + async fn cleanup_failed_copy(&self, e: Status, mounted: bool) { + let (state, req, vid, tx) = (&self.state, &self.req, self.vid, &self.tx); + let (data_base_name, idx_base_name, note_path) = ( + &self.dest.data_base_name, + &self.dest.idx_base_name, + self.dest.note_path(), + ); + // An abandoned copy is otherwise invisible here: the error goes + // to a channel nobody is reading. Logging it gives the operator + // the cause behind the balancer's "delete that copy, then re-run + // the move". + if e.code() == tonic::Code::Cancelled { + tracing::info!( + "volume {} copy from {} abandoned by its caller, discarding the partial copy", + vid, + req.source_data_node + ); + } else { + tracing::warn!( + "volume {} copy from {} failed: {}", + vid, + req.source_data_node, + e + ); + } + // Clean up on error. If the volume was mounted (the after-mount + // race), delete_volume unmounts it from the store AND removes + // the .dat/.idx/.vif in one step; the remove_file calls below + // cover the never-mounted partial-file case and are harmless + // no-ops when delete_volume already removed the files. + if mounted { + let mut store = state.store.write().unwrap(); + let _ = delete_replica_keep_remote(&mut store, vid); + state.volume_state_notify.notify_one(); + } + let _ = std::fs::remove_file(format!("{}.dat", data_base_name)); + let _ = std::fs::remove_file(format!("{}.idx", idx_base_name)); + let _ = std::fs::remove_file(format!("{}.vif", data_base_name)); + let _ = std::fs::remove_file(¬e_path); + let _ = tx.send(Err(e)).await; + } +} + /// Dial a ping target, bounding the whole connect at 5s. /// /// The outer timeout is not redundant with `GrpcDialOptions`' connect timeout: @@ -7409,7 +7584,9 @@ mod tests { let (dest_service, dest_tmp) = make_local_service_with_volume("", None); { let mut store = dest_service.state.store.write().unwrap(); - store.delete_volume(VolumeId(1), false, false, false).unwrap(); + store + .delete_volume(VolumeId(1), false, false, false) + .unwrap(); // available_space is filled in by the periodic disk check, which // does not run in a unit test; without it VolumeCopy finds no // location with room and never gets as far as copying. @@ -7556,7 +7733,9 @@ mod tests { let (dest_service, dest_tmp) = make_local_service_with_volume("", None); { let mut store = dest_service.state.store.write().unwrap(); - store.delete_volume(VolumeId(1), false, false, false).unwrap(); + store + .delete_volume(VolumeId(1), false, false, false) + .unwrap(); for loc in &store.locations { loc.check_disk_space(); } @@ -7714,7 +7893,9 @@ mod tests { let (dest_service, dest_tmp) = make_local_service_with_volume("", None); { let mut store = dest_service.state.store.write().unwrap(); - store.delete_volume(VolumeId(1), false, false, false).unwrap(); + store + .delete_volume(VolumeId(1), false, false, false) + .unwrap(); for loc in &store.locations { loc.check_disk_space(); } @@ -7828,7 +8009,9 @@ mod tests { let (dest_service, dest_tmp) = make_local_service_with_volume("", None); { let mut store = dest_service.state.store.write().unwrap(); - store.delete_volume(VolumeId(1), false, false, false).unwrap(); + store + .delete_volume(VolumeId(1), false, false, false) + .unwrap(); for loc in &store.locations { loc.check_disk_space(); }