diff --git a/seaweed-volume/src/server/grpc_server.rs b/seaweed-volume/src/server/grpc_server.rs index 036d1e2bc..1b704abbe 100644 --- a/seaweed-volume/src/server/grpc_server.rs +++ b/seaweed-volume/src/server/grpc_server.rs @@ -777,29 +777,114 @@ impl VolumeServer for VolumeGrpcService { } } }; - let mut results = Vec::new(); - let mut bytes_to_read = (dat_size - start_offset) as i64; - let buffer_size = 2 * 1024 * 1024; - let mut offset = start_offset; - - while bytes_to_read > 0 { - let chunk = std::cmp::min(bytes_to_read as usize, buffer_size); - match v.read_dat_slice(offset, chunk) { - Ok(buf) if buf.is_empty() => break, - Ok(buf) => { - let read_len = buf.len() as i64; - results.push(Ok(volume_server_pb::VolumeIncrementalCopyResponse { - file_content: buf, - })); - bytes_to_read -= read_len; - offset += read_len as u64; - } - Err(e) => return Err(Status::internal(e.to_string())), - } + // A corrupt index could place the located offset past the captured + // `.dat` size; nothing to send, and it would underflow `total` below. + if start_offset >= dat_size { + drop(store); + let stream = tokio_stream::iter(Vec::new()); + return Ok(Response::new(Box::pin(stream))); } - + // Capture a reader for the `.dat` while the volume lookup is still + // protected by the store lock, then drop the lock and stream the delta + // through a bounded channel fed by a blocking reader task (instead of + // buffering the whole delta into a Vec). Opening the local file -- or + // cloning the remote backend handle -- under the lock pins it to the + // volume that exists now, so a concurrent delete/recreate can't make us + // stream from the wrong file, and no store lock is held while a slow + // disk or S3 read runs. + enum DatReader { + Local(std::fs::File), + Remote(crate::storage::volume::RemoteDatFile), + } + let reader = if v.has_remote_file { + match v.remote_dat_file() { + Some(r) => DatReader::Remote(r), + None => { + return Err(Status::not_found(format!( + "remote dat not loaded for volume id {}", + vid + ))) + } + } + } else { + let path = v.file_name(".dat"); + let file = std::fs::File::open(&path) + .map_err(|e| Status::internal(format!("open {}: {}", path, e)))?; + DatReader::Local(file) + }; drop(store); - let stream = tokio_stream::iter(results); + + let total = dat_size - start_offset; + let (tx, rx) = tokio::sync::mpsc::channel::< + Result, + >(8); + + tokio::task::spawn_blocking(move || { + let buffer_size = 2 * 1024 * 1024u64; // 2MB chunks + let mut bytes_to_read = total; + let mut offset = start_offset; + + match reader { + DatReader::Local(file) => { + // Local volume: stream directly from the .dat, no lock held. + use std::io::{Read, Seek, SeekFrom}; + let mut reader = std::io::BufReader::new(file); + if let Err(e) = reader.seek(SeekFrom::Start(offset)) { + let _ = tx.blocking_send(Err(Status::internal(e.to_string()))); + return; + } + while bytes_to_read > 0 { + let chunk = std::cmp::min(bytes_to_read, buffer_size) as usize; + let mut buf = vec![0u8; chunk]; + match reader.read(&mut buf) { + Ok(0) => break, // EOF + Ok(n) => { + buf.truncate(n); + let msg = volume_server_pb::VolumeIncrementalCopyResponse { + file_content: buf, + }; + if tx.blocking_send(Ok(msg)).is_err() { + return; // receiver dropped + } + bytes_to_read -= n as u64; + } + Err(e) => { + let _ = tx.blocking_send(Err(Status::internal(e.to_string()))); + return; + } + } + } + } + DatReader::Remote(remote) => { + // Remote/tiered volume: read through the cloned backend + // handle. No store lock is held while the (potentially slow) + // S3 fetch runs, so it never blocks store writers. + while bytes_to_read > 0 { + let chunk = std::cmp::min(bytes_to_read, buffer_size) as usize; + match remote.read_slice(offset, chunk) { + Ok(buf) if buf.is_empty() => break, + Ok(buf) => { + let read_len = buf.len() as u64; + let msg = volume_server_pb::VolumeIncrementalCopyResponse { + file_content: buf, + }; + if tx.blocking_send(Ok(msg)).is_err() { + return; // receiver dropped + } + bytes_to_read -= read_len; + offset += read_len; + } + Err(e) => { + let _ = tx.blocking_send(Err(Status::internal(e.to_string()))); + return; + } + } + } + } + } + }); + + let stream = tokio_stream::wrappers::ReceiverStream::new(rx); Ok(Response::new(Box::pin(stream))) } @@ -1509,43 +1594,60 @@ impl VolumeServer for VolumeGrpcService { .map(|d| d.as_nanos() as i64) .unwrap_or(0); - let mut results: Vec> = Vec::new(); - let mut bytes_to_read = req.stop_offset as i64; - let mut reader = std::io::BufReader::new(file); - let buffer_size = 2 * 1024 * 1024; // 2MB chunks - let mut first = true; + // Stream the file in 2MB chunks from a blocking reader task instead of + // materializing the whole file into a Vec first. The previous code + // pushed every chunk into `results` and only then returned + // `tokio_stream::iter(results)`, so serving a large volume as a copy + // source (e.g. during `volume.fix.replication`) held the entire .dat + // (up to the volume size limit) resident at once and could OOM the + // process. A bounded channel caps memory at channel_capacity * 2MB and + // applies backpressure to the reader. + let stop_offset = req.stop_offset; + let (tx, rx) = + tokio::sync::mpsc::channel::>(8); - use std::io::Read; - while bytes_to_read > 0 { - let chunk_size = std::cmp::min(bytes_to_read as usize, buffer_size); - let mut buf = vec![0u8; chunk_size]; - match reader.read(&mut buf) { - Ok(0) => break, // EOF - Ok(n) => { - buf.truncate(n); - if n as i64 > bytes_to_read { - buf.truncate(bytes_to_read as usize); + tokio::task::spawn_blocking(move || { + use std::io::Read; + let buffer_size = 2 * 1024 * 1024u64; // 2MB chunks + let mut reader = std::io::BufReader::new(file); + let mut bytes_to_read = stop_offset; + let mut first = true; + + while bytes_to_read > 0 { + let chunk_size = std::cmp::min(bytes_to_read, buffer_size) as usize; + let mut buf = vec![0u8; chunk_size]; + match reader.read(&mut buf) { + Ok(0) => break, // EOF + Ok(n) => { + let n = std::cmp::min(n, bytes_to_read as usize); + buf.truncate(n); + let msg = volume_server_pb::CopyFileResponse { + file_content: buf, + modified_ts_ns: if first { mod_ts_ns } else { 0 }, + }; + if tx.blocking_send(Ok(msg)).is_err() { + return; // receiver dropped, stop reading + } + first = false; + bytes_to_read -= n as u64; + } + Err(e) => { + let _ = tx.blocking_send(Err(Status::internal(e.to_string()))); + return; } - results.push(Ok(volume_server_pb::CopyFileResponse { - file_content: buf, - modified_ts_ns: if first { mod_ts_ns } else { 0 }, - })); - first = false; - bytes_to_read -= n as i64; } - Err(e) => return Err(Status::internal(e.to_string())), } - } - // If no data was sent, still send ModifiedTsNs - if first && mod_ts_ns != 0 { - results.push(Ok(volume_server_pb::CopyFileResponse { - file_content: vec![], - modified_ts_ns: mod_ts_ns, - })); - } + // If no data was sent, still send ModifiedTsNs + if first && mod_ts_ns != 0 { + let _ = tx.blocking_send(Ok(volume_server_pb::CopyFileResponse { + file_content: vec![], + modified_ts_ns: mod_ts_ns, + })); + } + }); - let stream = tokio_stream::iter(results); + let stream = tokio_stream::wrappers::ReceiverStream::new(rx); Ok(Response::new(Box::pin(stream))) } @@ -4905,6 +5007,145 @@ mod tests { .remove("s3.incr_copy_test"); } + /// 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. + fn make_local_service_with_large_volume() -> (VolumeGrpcService, TempDir, Vec) { + let (service, tmp) = make_local_service_with_volume("", None); + { + let mut store = service.state.store.write().unwrap(); + let (_, v) = store.find_volume_mut(VolumeId(1)).unwrap(); + // Non-uniform payload so reassembly catches duplicated or reordered + // chunks that a repeated byte would hide. + let payload: Vec = (0..5 * 1024 * 1024usize) + .map(|i| (i as u8).wrapping_mul(31).wrapping_add(7)) + .collect(); + let mut needle = Needle { + id: NeedleId(99), + cookie: Cookie(0x1122), + data_size: payload.len() as u32, + data: payload, + ..Needle::default() + }; + v.write_needle(&mut needle, true).unwrap(); + v.sync_to_disk().unwrap(); + } + let dat_path = { + let store = service.state.store.read().unwrap(); + let (_, v) = store.find_volume(VolumeId(1)).unwrap(); + v.file_name(".dat") + }; + let dat_bytes = std::fs::read(&dat_path).unwrap(); + (service, tmp, dat_bytes) + } + + // copy_file must stream the whole .dat in 2MB chunks (not buffer it) and + // reassemble byte-for-byte, with the mtime carried only on the first message. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_copy_file_streams_full_dat_in_chunks() { + let (service, _tmp, dat_bytes) = make_local_service_with_large_volume(); + let stop_offset = dat_bytes.len() as u64; + assert!( + stop_offset > 2 * 1024 * 1024u64, + "fixture must exceed one 2MB chunk" + ); + + let response = service + .copy_file(Request::new(volume_server_pb::CopyFileRequest { + volume_id: 1, + ext: ".dat".to_string(), + compaction_revision: u32::MAX, + stop_offset, + collection: String::new(), + is_ec_volume: false, + ignore_source_file_not_found: false, + })) + .await + .unwrap(); + + let mut stream = response.into_inner(); + let mut messages = Vec::new(); + while let Some(m) = stream.next().await { + messages.push(m.unwrap()); + } + + assert!( + messages.len() >= 2, + "expected multiple streamed chunks, got {}", + messages.len() + ); + assert_ne!(messages[0].modified_ts_ns, 0); + assert!(messages[1..].iter().all(|m| m.modified_ts_ns == 0)); + + let copied: Vec = messages.iter().flat_map(|m| m.file_content.clone()).collect(); + assert_eq!(copied, dat_bytes); + } + + // copy_file must stop exactly at stop_offset, never streaming past it. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_copy_file_respects_stop_offset() { + let (service, _tmp, dat_bytes) = make_local_service_with_large_volume(); + let stop_offset = dat_bytes.len() as u64 - 100; + + let response = service + .copy_file(Request::new(volume_server_pb::CopyFileRequest { + volume_id: 1, + ext: ".dat".to_string(), + compaction_revision: u32::MAX, + stop_offset, + collection: String::new(), + is_ec_volume: false, + ignore_source_file_not_found: false, + })) + .await + .unwrap(); + + let mut stream = response.into_inner(); + let mut copied = Vec::new(); + while let Some(m) = stream.next().await { + copied.extend_from_slice(&m.unwrap().file_content); + } + + assert_eq!(copied.len() as u64, stop_offset); + assert_eq!(copied, dat_bytes[..stop_offset as usize]); + } + + // volume_incremental_copy must stream a local volume from the super block + // to EOF in chunks, reassembling byte-for-byte. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_volume_incremental_copy_streams_local_volume_in_chunks() { + let (service, _tmp, dat_bytes) = make_local_service_with_large_volume(); + let super_block_size = { + let store = service.state.store.read().unwrap(); + let (_, v) = store.find_volume(VolumeId(1)).unwrap(); + v.super_block.block_size() as u64 + }; + + let response = service + .volume_incremental_copy(Request::new( + volume_server_pb::VolumeIncrementalCopyRequest { + volume_id: 1, + since_ns: 0, + }, + )) + .await + .unwrap(); + + let mut stream = response.into_inner(); + let mut messages = Vec::new(); + while let Some(m) = stream.next().await { + messages.push(m.unwrap()); + } + + assert!( + messages.len() >= 2, + "expected multiple streamed chunks, got {}", + messages.len() + ); + let copied: Vec = messages.iter().flat_map(|m| m.file_content.clone()).collect(); + assert_eq!(copied, dat_bytes[super_block_size as usize..]); + } + /// Build a bare-bones service with no on-disk store but a configurable /// seed master set, for Ping admission tests. Matches the structure of /// `make_local_service_with_volume` minus the volume bits. diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index d7f39d397..a83c652e5 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -455,6 +455,19 @@ pub(crate) struct RemoteDatFile { } impl RemoteDatFile { + /// Read up to `size` bytes at `offset`, bounded by the remote file size. + /// Mirrors `Volume::read_dat_slice` for a remote-only backend, but holds no + /// volume/store lock so callers can stream off the store lock entirely. + pub(crate) fn read_slice(&self, offset: u64, size: usize) -> io::Result> { + if size == 0 || offset >= self.file_size { + return Ok(Vec::new()); + } + let read_len = std::cmp::min(size as u64, self.file_size - offset) as usize; + let mut buf = vec![0u8; read_len]; + self.read_exact_at(&mut buf, offset)?; + Ok(buf) + } + pub(crate) fn read_exact_at(&self, buf: &mut [u8], offset: u64) -> io::Result<()> { let data = self .backend @@ -1037,6 +1050,14 @@ impl Volume { Ok(buf) } + /// 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` + /// 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() + } + // ---- SuperBlock I/O ---- fn read_super_block(&mut self) -> Result<(), VolumeError> {