From 4bb40732bbe8e1df76d14c7c575a12cbabb5edf6 Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Mon, 21 Sep 2026 09:26:37 +0300 Subject: [PATCH] volume server: ReceiveFile loses bytes and hides fsync failures (#11407) * volume server: ReceiveFile loses bytes and hides fsync failures Three defects in one handler, all on the path that receives a pushed .dat/.idx/.vif or EC shard: - `f.write(&content)` never compared the return to content.len(). A short write (ENOSPC, NFS) counted only the bytes that landed, so every later chunk was written at a shifted offset and the RPC answered error: "" with a byte count that looked right. Go's os.File.Write loops. Now write_all. - `let _ = f.sync_all();` discarded EIO and answered success with the full byte count. Go omits the check too, but ReceiveFileResponse carries an `error` field and the caller renames the staged file into place on success -- so a silent fsync failure publishes a file whose data never reached the platter. Flush and fsync failures are now reported. - Both the per-chunk write and the final fsync were blocking std::fs calls inside the async fn, on the runtime worker that is also driving the stream. Switched to tokio::fs + BufWriter, the shape `drain_copy_stream_to_file` in this same file already uses and documents. The partial-file cleanup on the error path moves to tokio::fs::remove_file for the same reason. The handler had no test at all, which is how the short-write bug survived. Added a round-trip over a real connection with ragged chunk boundaries, asserting the bytes on disk and not only the reported count -- a dropped or reordered chunk changes the file even when bytes_written still adds up. That test guards the rewrite; it does not reproduce the original faults. ENOSPC and EIO need fault injection that this suite has no harness for, so the short-write and fsync paths are argued from the code, not demonstrated by a failing test. Co-Authored-By: Claude Opus 5 (1M context) * volume server: remove the staged file on every ReceiveFile error reply Flush and fsync failures returned early and left the partial .copying or shard file behind, as did the pre-existing write-error path. Route all response-level errors through one cleanup block, matching Go's close-and-remove on a failed write. * volume server: tighten ReceiveFile comments --------- Co-authored-by: Claude Opus 5 (1M context) Co-authored-by: chrislusf Co-authored-by: Devin --- seaweed-volume/src/server/grpc_server.rs | 111 ++++++++++++++++++++--- 1 file changed, 100 insertions(+), 11 deletions(-) diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 07f918336..de651c611 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -2415,8 +2415,13 @@ impl VolumeServer for VolumeGrpcService { ) -> Result, Status> { self.state.check_maintenance()?; + use tokio::io::AsyncWriteExt; + let mut stream = request.into_inner(); - let mut target_file: Option = None; + // tokio::fs + BufWriter, as `drain_copy_stream_to_file` below already + // does: the chunk writes and the final fsync are disk I/O and must not + // run on the runtime worker that is also driving this stream. + let mut target_file: Option> = None; let mut file_path: Option = None; let mut bytes_written: u64 = 0; let mut resp_error: Option = None; @@ -2554,16 +2559,19 @@ impl VolumeServer for VolumeGrpcService { } }; - target_file = Some(std::fs::File::create(&path).map_err(|e| { + let f = tokio::fs::File::create(&path).await.map_err(|e| { Status::internal(format!("failed to create file: {}", e)) - })?); + })?; + target_file = Some(tokio::io::BufWriter::new(f)); file_path = Some(path); } Some(volume_server_pb::receive_file_request::Data::FileContent(content)) => { if let Some(ref mut f) = target_file { - use std::io::Write; - match f.write(&content) { - Ok(n) => bytes_written += n as u64, + // write_all, not write: a short write (ENOSPC, NFS) + // would be counted as success for however many + // bytes landed, silently shifting later chunks. + match f.write_all(&content).await { + Ok(()) => bytes_written += content.len() as u64, Err(e) => { // Match Go: write failures are response-level errors, not gRPC errors resp_error = Some(format!("failed to write file: {}", e)); @@ -2588,16 +2596,32 @@ impl VolumeServer for VolumeGrpcService { match result { Ok(()) => { - // Check for protocol-level errors (returned in response body, not gRPC status) + // Flush the BufWriter and fsync, reporting failure through the + // response `error` field: the caller renames the staged file + // into place on success, so a swallowed fsync error would + // publish data that never reached disk. + if resp_error.is_none() + && let Some(ref mut f) = target_file + { + if let Err(e) = f.flush().await { + resp_error = Some(format!("failed to flush file: {}", e)); + } else if let Err(e) = f.get_ref().sync_all().await { + resp_error = Some(format!("failed to sync file: {}", e)); + } + } + // Protocol-level errors are returned in the response body, not + // gRPC status. Any of them leaves a partial staged file behind; + // remove it as Go does on a failed write. if let Some(err_msg) = resp_error { + drop(target_file.take()); + if let Some(ref p) = file_path { + let _ = tokio::fs::remove_file(p).await; + } return Ok(Response::new(volume_server_pb::ReceiveFileResponse { error: err_msg, bytes_written: 0, })); } - if let Some(ref f) = target_file { - let _ = f.sync_all(); - } Ok(Response::new(volume_server_pb::ReceiveFileResponse { error: String::new(), bytes_written, @@ -2609,7 +2633,7 @@ impl VolumeServer for VolumeGrpcService { drop(f); } if let Some(ref p) = file_path { - let _ = std::fs::remove_file(p); + let _ = tokio::fs::remove_file(p).await; } Err(e) } @@ -7039,6 +7063,71 @@ mod tests { // delete_volume. Without the lock seam the task sees is_closed() at its // very first check and returns before mount_volume, exercising the wrong // path — the test would be green for the wrong reason. + /// ReceiveFile had no test at all, which is how a `write()` whose short + /// return was counted as success survived. This drives the real streaming + /// handler over a real connection with chunks that do not divide evenly, + /// and checks the bytes on disk rather than just the reported count -- a + /// dropped or reordered chunk changes the file even when `bytes_written` + /// still adds up. + #[tokio::test] + async fn receive_file_writes_every_chunk_and_reports_the_full_length() { + let (service, tmp) = make_local_service_with_volume("", None); + let dir = tmp.path().to_str().unwrap().to_string(); + let (port, _shutdown) = serve_source(service).await; + + let mut client = volume_server_pb::volume_server_client::VolumeServerClient::connect( + format!("http://127.0.0.1:{}", port), + ) + .await + .unwrap(); + + // Deliberately ragged chunk sizes, and a payload whose bytes are + // position-dependent so any shift is visible. + let payload: Vec = (0..70_001u32).map(|i| (i % 251) as u8).collect(); + let mut messages = vec![volume_server_pb::ReceiveFileRequest { + data: Some(volume_server_pb::receive_file_request::Data::Info( + volume_server_pb::ReceiveFileInfo { + volume_id: 1, + ext: ".recv_test".to_string(), + collection: String::new(), + is_ec_volume: false, + shard_id: 0, + file_size: payload.len() as u64, + disk_type: String::new(), + disk_id: 0, + }, + )), + }]; + for chunk in payload.chunks(7_777) { + messages.push(volume_server_pb::ReceiveFileRequest { + data: Some(volume_server_pb::receive_file_request::Data::FileContent( + chunk.to_vec(), + )), + }); + } + + let response = client + .receive_file(tokio_stream::iter(messages)) + .await + .unwrap() + .into_inner(); + + assert_eq!(response.error, "", "ReceiveFile reported an error"); + assert_eq!( + response.bytes_written, + payload.len() as u64, + "bytes_written must cover the whole payload" + ); + + let written = std::fs::read(format!("{}/1.recv_test", dir)).unwrap(); + assert_eq!( + written.len(), + payload.len(), + "file on disk is a different length than the payload" + ); + assert_eq!(written, payload, "file on disk does not match the payload"); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[expect( clippy::await_holding_lock,