rust volume: stop a tier move whose caller has gone (#11192)

Both tier-move handlers run in a detached tokio::spawn and report
progress through a closure that returns (), with the send result
discarded. Nothing observes the caller leaving, so an abandoned move
uploads or downloads the whole .dat anyway and then commits the
transition.

Go aborts both. Its progress callback returns `stream.Send`'s error,
which surfaces out of the reader in s3_upload.go:99 and the writer in
s3_download.go:84 and fails the transfer, so the volume info is never
rewritten. The Rust port dropped that by typing the callback as
FnMut(i64, f32) with no result.

Give the callback Go's signature -- FnMut(i64, f32) -> Result<(), String>
-- and abort when the caller's channel is closed. Checked on every part
rather than only where progress is reported, since the report is
rate-limited to one a second and would miss a caller that left in
between. A merely full channel is a slow reader, not a departed one, so
only TrySendError::Closed counts as cancellation.

Two consequences of aborting mid-transfer that the old code never had to
handle:

- upload_file now aborts the multipart upload when the transfer fails.
  An abandoned multipart upload does not show up in an ordinary object
  listing but still accrues storage charges until a lifecycle rule reaps
  it, and cancellation makes that a routine path rather than a rare one.
- The tier-down handler removes the partial .dat. download_file
  pre-allocates the destination to the object's full size, so an aborted
  download leaves a file of the right length and the wrong content --
  and this handler refuses to run at all when a local .dat exists, so
  leaving one wedges every retry on "already on local disk" and a
  restart would load the sparse file as the volume's data.

There is deliberately no check between a finished transfer and the
bookkeeping that follows. Once the object is in S3, or the .dat is on
disk, that bookkeeping is what makes the state consistent; stopping
there would leave an object paid for and referenced by nothing, or a
complete local .dat the volume still calls remote. Go does not gate
there either -- its callback only runs during the transfer.


Claude-Session: https://claude.ai/code/session_0122W3eqt6gmLUMxmRoZdPAb

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Eliah Rusin
2026-09-06 12:18:47 -07:00
committed by GitHub
co-authored by Claude Opus 5
parent 4b41329e12
commit ade4bdf9e6
2 changed files with 276 additions and 42 deletions
+71 -30
View File
@@ -81,13 +81,17 @@ impl S3TierBackend {
/// Returns (s3_key, file_size) on success.
/// The progress callback receives (bytes_uploaded, percentage).
/// Uses 64MB part size and 5 concurrent uploads (matches Go s3manager).
/// `progress_fn` returning `Err` aborts the upload, mirroring Go's
/// `fn(progressed, percentage) error` in `s3_upload.go`, where the error
/// surfaces out of `ReadAt` and fails the transfer. It is what lets a
/// caller that has hung up stop the work it is no longer waiting for.
pub async fn upload_file<F>(
&self,
file_path: &str,
progress_fn: F,
) -> Result<(String, u64), String>
where
F: FnMut(i64, f32) + Send + Sync + 'static,
F: FnMut(i64, f32) -> Result<(), String> + Send + Sync + 'static,
{
let key = uuid::Uuid::new_v4().to_string();
@@ -184,8 +188,10 @@ impl S3TierBackend {
let e_tag = upload_part_resp.e_tag().unwrap_or_default().to_string();
// Report progress
{
// Report progress. The lock is released before the result is
// propagated so an aborting callback cannot poison the mutex
// for the other parts still in flight.
let progress_result = {
let mut guard = progress.lock().unwrap();
guard.0 += size as u64;
let uploaded = guard.0;
@@ -194,8 +200,9 @@ impl S3TierBackend {
} else {
100.0
};
(guard.1)(uploaded as i64, pct);
}
(guard.1)(uploaded as i64, pct)
};
progress_result?;
Ok::<_, String>(
CompletedPart::builder()
@@ -206,29 +213,58 @@ impl S3TierBackend {
}));
}
// Collect results, preserving part order
let mut completed_parts = Vec::with_capacity(handles.len());
for handle in handles {
let part = handle
let finish = async {
// Collect results, preserving part order
let mut completed_parts = Vec::with_capacity(handles.len());
for handle in handles {
let part = handle
.await
.map_err(|e| format!("upload task panicked: {}", e))??;
completed_parts.push(part);
}
// Complete multipart upload
let completed_upload = CompletedMultipartUpload::builder()
.set_parts(Some(completed_parts))
.build();
self.client
.complete_multipart_upload()
.bucket(&self.bucket)
.key(&key)
.upload_id(&upload_id)
.multipart_upload(completed_upload)
.send()
.await
.map_err(|e| format!("upload task panicked: {}", e))??;
completed_parts.push(part);
.map_err(|e| format!("failed to complete multipart upload: {}", e))?;
Ok::<(), String>(())
}
.await;
// Complete multipart upload
let completed_upload = CompletedMultipartUpload::builder()
.set_parts(Some(completed_parts))
.build();
self.client
.complete_multipart_upload()
.bucket(&self.bucket)
.key(&key)
.upload_id(&upload_id)
.multipart_upload(completed_upload)
.send()
.await
.map_err(|e| format!("failed to complete multipart upload: {}", e))?;
if let Err(e) = finish {
// An abandoned multipart upload does not appear in an ordinary
// object listing but still accrues storage charges until a
// lifecycle rule reaps it. Now that a departing caller aborts the
// transfer this is a routine path, not a rare one.
if let Err(abort_err) = self
.client
.abort_multipart_upload()
.bucket(&self.bucket)
.key(&key)
.upload_id(&upload_id)
.send()
.await
{
tracing::warn!(
"failed to abort multipart upload {} for key {}: {}",
upload_id,
key,
abort_err
);
}
return Err(e);
}
Ok((key, file_size))
}
@@ -238,6 +274,8 @@ impl S3TierBackend {
///
/// Returns the file size on success.
/// Uses 64MB part size and 5 concurrent downloads (matches Go s3manager).
/// `progress_fn` returning `Err` aborts the download, mirroring Go's
/// `fn(progressed, percentage) error` in `s3_download.go`.
pub async fn download_file<F>(
&self,
dest_path: &str,
@@ -245,7 +283,7 @@ impl S3TierBackend {
progress_fn: F,
) -> Result<u64, String>
where
F: FnMut(i64, f32) + Send + Sync + 'static,
F: FnMut(i64, f32) -> Result<(), String> + Send + Sync + 'static,
{
// Get file size first
let head_resp = self
@@ -340,8 +378,10 @@ impl S3TierBackend {
.await
.map_err(|e| format!("failed to write to {}: {}", dp, e))?;
// Report progress
{
// Report progress. The lock is released before the result is
// propagated so an aborting callback cannot poison the mutex
// for the other parts still in flight.
let progress_result = {
let mut guard = progress.lock().unwrap();
guard.0 += bytes.len() as u64;
let downloaded = guard.0;
@@ -350,8 +390,9 @@ impl S3TierBackend {
} else {
100.0
};
(guard.1)(downloaded as i64, pct);
}
(guard.1)(downloaded as i64, pct)
};
progress_result?;
Ok::<_, String>(())
}));
+205 -12
View File
@@ -3756,21 +3756,52 @@ impl VolumeServer for VolumeGrpcService {
tokio::spawn(async move {
let result: Result<(), Status> = async {
// Nothing below is worth doing for a caller that has already
// gone: the upload would spend bandwidth and S3 storage on a
// transition nobody is waiting to hear the result of.
if tx.is_closed() {
return Err(Status::cancelled(format!(
"volume {} tier move to remote cancelled by caller",
vid
)));
}
// Upload the .dat file to S3 with progress
let tx_progress = tx.clone();
let mut last_report = std::time::Instant::now();
let (key, size) = backend
.upload_file(&dat_path, move |processed, percentage| {
// Checked on every part, not only where progress is
// reported: the report is rate-limited to one a second,
// so its result alone would miss a caller that left in
// between. Go aborts here too -- the progress callback
// returns stream.Send's error (s3_upload.go).
if tx_progress.is_closed() {
return Err(format!(
"volume {} tier move to remote cancelled by caller",
vid
));
}
let now = std::time::Instant::now();
if now.duration_since(last_report) >= std::time::Duration::from_secs(1) {
last_report = now;
let _ = tx_progress.try_send(Ok(
volume_server_pb::VolumeTierMoveDatToRemoteResponse {
processed,
processed_percentage: percentage,
},
));
if let Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) =
tx_progress.try_send(Ok(
volume_server_pb::VolumeTierMoveDatToRemoteResponse {
processed,
processed_percentage: percentage,
},
))
{
return Err(format!(
"volume {} tier move to remote cancelled by caller",
vid
));
}
// A merely full channel is a slow reader, not a
// departed one; drop the report and keep going.
}
Ok(())
})
.await
.map_err(|e| {
@@ -3780,6 +3811,11 @@ impl VolumeServer for VolumeGrpcService {
))
})?;
// Deliberately no cancellation check here. Once the object is
// in S3 the cheap bookkeeping that follows is what makes the
// 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
{
let mut store = state.store.write().unwrap();
@@ -3831,6 +3867,16 @@ impl VolumeServer for VolumeGrpcService {
.await;
if let Err(e) = result {
// The error otherwise goes to a channel nobody is
// reading, leaving an abandoned tier move with no trace.
if e.code() == tonic::Code::Cancelled {
tracing::info!(
"volume {} tier move to remote abandoned by its caller",
vid
);
} else {
tracing::warn!("volume {} tier move to remote failed: {}", vid, e);
}
let _ = tx.send(Err(e)).await;
}
});
@@ -3913,25 +3959,73 @@ impl VolumeServer for VolumeGrpcService {
tokio::spawn(async move {
let result: Result<(), Status> = async {
// Nothing below is worth doing for a caller that has already
// gone, and the download would overwrite a .dat path the volume
// still reports as absent.
if tx.is_closed() {
return Err(Status::cancelled(format!(
"volume {} tier move from remote cancelled by caller",
vid
)));
}
// Download the .dat file from S3 with progress
let tx_progress = tx.clone();
let mut last_report = std::time::Instant::now();
let storage_name_clone = storage_name.clone();
let _size = backend
.download_file(&dat_path, &storage_key, move |processed, percentage| {
// Checked on every part, not only where progress is
// reported: the report is rate-limited to one a second,
// so its result alone would miss a caller that left in
// between. Go aborts here too -- the progress callback
// returns stream.Send's error (s3_download.go).
if tx_progress.is_closed() {
return Err(format!(
"volume {} tier move from remote cancelled by caller",
vid
));
}
let now = std::time::Instant::now();
if now.duration_since(last_report) >= std::time::Duration::from_secs(1) {
last_report = now;
let _ = tx_progress.try_send(Ok(
volume_server_pb::VolumeTierMoveDatFromRemoteResponse {
processed,
processed_percentage: percentage,
},
));
if let Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) =
tx_progress.try_send(Ok(
volume_server_pb::VolumeTierMoveDatFromRemoteResponse {
processed,
processed_percentage: percentage,
},
))
{
return Err(format!(
"volume {} tier move from remote cancelled by caller",
vid
));
}
// A merely full channel is a slow reader, not a
// departed one; drop the report and keep going.
}
Ok(())
})
.await
.map_err(|e| {
// download_file pre-allocates the destination to the
// object's full size, so an aborted download leaves a
// .dat that is the right length and the wrong content.
// It has to go: this handler refuses to run at all when
// a local .dat exists, so leaving one wedges every
// retry on "already on local disk", and a restart would
// load the sparse file as the volume's data.
if let Err(rm) = std::fs::remove_file(&dat_path) {
if rm.kind() != std::io::ErrorKind::NotFound {
tracing::warn!(
"volume {} could not remove the incomplete download {}: {}",
vid,
dat_path,
rm
);
}
}
Status::internal(format!(
"backend {} copy file {}: {}",
storage_name_clone, dat_path, e
@@ -4062,6 +4156,16 @@ impl VolumeServer for VolumeGrpcService {
.await;
if let Err(e) = result {
// The error otherwise goes to a channel nobody is
// reading, leaving an abandoned tier move with no trace.
if e.code() == tonic::Code::Cancelled {
tracing::info!(
"volume {} tier move from remote abandoned by its caller",
vid
);
} else {
tracing::warn!("volume {} tier move from remote failed: {}", vid, e);
}
let _ = tx.send(Err(e)).await;
}
});
@@ -5959,6 +6063,95 @@ mod tests {
.remove("s3.tier_down_delete");
}
// The progress callback is the only thing a tier transfer polls, so it is
// the only place a departing caller can be noticed. Go gives it an error
// return for exactly this (s3_upload.go:99, s3_download.go:84); this checks
// the Rust port now honours one.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_download_file_aborts_when_the_progress_callback_fails() {
let (service, tmp, shutdown_tx, _dat_bytes, _super_block_size, _delete_count) =
make_remote_only_service("tier_progress_abort");
let backend = global_s3_tier_registry()
.read()
.unwrap()
.get("s3.tier_progress_abort")
.expect("backend registered");
let dest = format!("{}/aborted.dat", tmp.path().to_str().unwrap());
let err = backend
.download_file(&dest, "remote-key", |_, _| {
Err("caller gone".to_string())
})
.await
.expect_err("a failing progress callback must abort the download");
assert!(err.contains("caller gone"), "unexpected error: {}", err);
drop(service);
let _ = shutdown_tx.send(());
global_s3_tier_registry()
.write()
.unwrap()
.remove("s3.tier_progress_abort");
}
// An abandoned tier-down must leave the volume exactly as it found it. The
// partial .dat matters most: download_file pre-allocates the destination to
// the object's full size, this handler refuses to run while a local .dat
// exists, and a restart would load that file as the volume's data -- so a
// leftover wedges every retry on "already on local disk" and corrupts the
// volume on the way.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_tier_move_from_remote_cancelled_leaves_the_volume_untouched() {
let (service, tmp, shutdown_tx, _dat_bytes, _super_block_size, delete_count) =
make_remote_only_service("tier_down_cancel");
let dat_path = format!("{}/1.dat", tmp.path().to_str().unwrap());
assert!(!std::path::Path::new(&dat_path).exists());
let response = service
.volume_tier_move_dat_from_remote(Request::new(
volume_server_pb::VolumeTierMoveDatFromRemoteRequest {
volume_id: 1,
collection: String::new(),
keep_remote_dat_file: false,
},
))
.await
.unwrap();
// The caller hangs up: dropping the response drops the receiving half
// of the channel, which is all a cancelled RPC amounts to here.
drop(response);
for _ in 0..60 {
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
assert!(
!std::path::Path::new(&dat_path).exists(),
"abandoned tier-down left a local .dat behind"
);
}
{
let store = service.state.store.read().unwrap();
let (_, vol) = store.find_volume(VolumeId(1)).unwrap();
assert!(
vol.has_remote_file,
"abandoned tier-down published the transition to local"
);
assert!(!vol.volume_info.files.is_empty());
}
assert_eq!(
delete_count.load(std::sync::atomic::Ordering::SeqCst),
0,
"abandoned tier-down deleted the shared remote object"
);
let _ = shutdown_tx.send(());
global_s3_tier_registry()
.write()
.unwrap()
.remove("s3.tier_down_cancel");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_tier_move_from_remote_keep_remote_still_swaps_to_local() {
let (service, tmp, shutdown_tx, dat_bytes, _super_block_size, delete_count) =