diff --git a/seaweed-volume/src/server/handlers.rs b/seaweed-volume/src/server/handlers.rs index ce85156a2..5e96d49ef 100644 --- a/seaweed-volume/src/server/handlers.rs +++ b/seaweed-volume/src/server/handlers.rs @@ -143,6 +143,78 @@ const STREAMING_THRESHOLD: u32 = 1024 * 1024; // 1 MB /// Default chunk size for streaming reads from the dat file. const DEFAULT_STREAMING_CHUNK_SIZE: usize = 64 * 1024; // 64 KB +/// Whether a GET/HEAD can be answered from the needle meta and the data file. +#[derive(Clone, Copy)] +struct SourceReadRequest { + is_head: bool, + has_range: bool, + has_image_ops: bool, + bypass_cm: bool, +} + +impl SourceReadRequest { + /// The payload can be served as stored. + fn direct(&self, n: &Needle) -> bool { + !n.is_compressed() && !(n.is_chunk_manifest() && !self.bypass_cm) && !self.has_image_ops + } + + /// HEAD, a range, or a stream. + fn served_from_source(&self, n: &Needle) -> bool { + self.is_head || (self.direct(n) && (self.has_range || n.data_size > STREAMING_THRESHOLD)) + } + + /// Meta first only if the needle could be served from source; its data + /// size is below its index size. + fn try_meta_first(&self, size: Size) -> bool { + self.is_head + || (!self.has_image_ops && (self.has_range || size.0 > STREAMING_THRESHOLD as i32)) + } +} + +/// Blocking part of a GET/HEAD: the store guard is dropped before any needle +/// I/O, so a slow or tiered read cannot park store writers. `Ok(None)` is a +/// cookie mismatch. +fn read_needle_for_get( + state: &VolumeServerState, + vid: VolumeId, + needle_id: NeedleId, + cookie: Cookie, + read_deleted: bool, + request: SourceReadRequest, +) -> Result< + Option<(Needle, Option)>, + crate::storage::volume::VolumeError, +> { + let mut plan = state + .store + .read() + .unwrap() + .needle_read_plan(vid, needle_id, read_deleted)?; + let blank = || Needle { + id: needle_id, + cookie, + ..Needle::default() + }; + let mut n = blank(); + if request.try_meta_first(plan.size()) { + plan.read_meta(&mut n)?; + if n.cookie != cookie { + return Ok(None); + } + if request.served_from_source(&n) { + let info = plan.into_stream_info(&n); + return Ok(Some((n, Some(info)))); + } + // Compressed or a chunk manifest: the payload is needed after all. + n = blank(); + } + plan.read_full(&mut n)?; + if n.cookie != cookie { + return Ok(None); + } + Ok(Some((n, None))) +} + /// A body that streams needle data from the dat file in chunks using pread, /// avoiding loading the entire payload into memory at once. struct StreamingBody { @@ -1112,11 +1184,7 @@ async fn get_or_head_handler_inner( // Read needle — branching between regular volume and EC volume paths. // EC volumes always do a full read (no streaming/meta-only). - let mut n = Needle { - id: needle_id, - cookie, - ..Needle::default() - }; + let n: Needle; let read_deleted = query.read_deleted.as_deref() == Some("true"); let has_range = headers.contains_key(header::RANGE); @@ -1193,24 +1261,49 @@ async fn get_or_head_handler_inner( can_handle_range_from_source = false; } else { // ---- Regular volume read path (with streaming support) ---- + bypass_cm = query.cm.as_deref() == Some("false"); + track_download = download_guard.is_some(); + let request_kind = SourceReadRequest { + is_head: method == Method::HEAD, + has_range, + has_image_ops, + bypass_cm, + }; - // Try meta-only read first for potential streaming - let store = state.store.read().unwrap(); - let si_result = store.read_volume_needle_stream_info(vid, &mut n, read_deleted); - stream_info = match si_result { - Ok(info) => Some(info), - Err(crate::storage::volume::VolumeError::StreamingUnsupported) => None, - Err(crate::storage::volume::VolumeError::NotFound) => { + let read_state = state.clone(); + let read = tokio::task::spawn_blocking(move || { + read_needle_for_get( + &read_state, + vid, + needle_id, + cookie, + read_deleted, + request_kind, + ) + }) + .await; + (n, stream_info) = match read { + Ok(Ok(Some(found))) => found, + // Cookie mismatch + Ok(Ok(None)) => return StatusCode::NOT_FOUND.into_response(), + Ok(Err( + crate::storage::volume::VolumeError::NotFound + | crate::storage::volume::VolumeError::Deleted, + )) => { metrics::HANDLER_COUNTER .with_label_values(&[metrics::ERROR_GET_NOT_FOUND]) .inc(); return StatusCode::NOT_FOUND.into_response(); } - Err(crate::storage::volume::VolumeError::Deleted) => { + Ok(Err(e)) => { metrics::HANDLER_COUNTER - .with_label_values(&[metrics::ERROR_GET_NOT_FOUND]) + .with_label_values(&[metrics::ERROR_GET_INTERNAL]) .inc(); - return StatusCode::NOT_FOUND.into_response(); + return ( + StatusCode::INTERNAL_SERVER_ERROR, + format!("read error: {}", e), + ) + .into_response(); } Err(e) => { metrics::HANDLER_COUNTER @@ -1223,19 +1316,9 @@ async fn get_or_head_handler_inner( .into_response(); } }; - drop(store); - // Validate cookie - if n.cookie != cookie { - return StatusCode::NOT_FOUND.into_response(); - } - - bypass_cm = query.cm.as_deref() == Some("false"); - track_download = download_guard.is_some(); - let can_direct_source_read = stream_info.is_some() - && !n.is_compressed() - && !(n.is_chunk_manifest() && !bypass_cm) - && !has_image_ops; + // Stream info is only returned for a reply served from the data file. + let can_direct_source_read = stream_info.is_some() && request_kind.direct(&n); // Determine if we can stream (large, direct-source eligible, no range) can_stream = can_direct_source_read @@ -1246,41 +1329,6 @@ async fn get_or_head_handler_inner( // Go uses meta-only reads for all HEAD requests, regardless of compression/chunked files. can_handle_head_from_meta = stream_info.is_some() && method == Method::HEAD; can_handle_range_from_source = can_direct_source_read && has_range; - - // For chunk manifest or any non-streaming path, we need the full data. - // If we can't stream, do a full read now. - if !can_stream && !can_handle_head_from_meta && !can_handle_range_from_source { - // Re-read with full data - let mut n_full = Needle { - id: needle_id, - cookie, - ..Needle::default() - }; - let store = state.store.read().unwrap(); - match store.read_volume_needle_opt(vid, &mut n_full, read_deleted) { - Ok(count) => { - if count < 0 { - return StatusCode::NOT_FOUND.into_response(); - } - } - Err(crate::storage::volume::VolumeError::NotFound) => { - return StatusCode::NOT_FOUND.into_response(); - } - Err(crate::storage::volume::VolumeError::Deleted) => { - return StatusCode::NOT_FOUND.into_response(); - } - Err(e) => { - return ( - StatusCode::INTERNAL_SERVER_ERROR, - format!("read error: {}", e), - ) - .into_response(); - } - } - drop(store); - // Use the full needle from here (it has the same metadata + data) - n = n_full; - } } // Build ETag and Last-Modified BEFORE conditional checks and chunk manifest expansion @@ -1565,12 +1613,19 @@ async fn get_or_head_handler_inner( && let (Some(range_header), Some(info)) = (headers.get(header::RANGE), stream_info) && let Ok(range_str) = range_header.to_str() { - return handle_range_request_from_source( - range_str, - info, - response_headers, - track_download.then(|| state.clone()), - ); + let range_str = range_str.to_string(); + let tracking = track_download.then(|| state.clone()); + return tokio::task::spawn_blocking(move || { + handle_range_request_from_source(&range_str, info, response_headers, tracking) + }) + .await + .unwrap_or_else(|e| { + ( + StatusCode::INTERNAL_SERVER_ERROR, + format!("range read error: {}", e), + ) + .into_response() + }); } // ---- Buffered path: small files, compressed, images, range requests ---- @@ -4617,15 +4672,19 @@ mod tests { } fn streaming_test_state() -> Arc { - use crate::security::{Guard, SigningKey}; - use crate::server::volume_server::RuntimeMetricsConfig; use crate::storage::needle_map::NeedleMapKind; use crate::storage::store::Store; + test_state_with_store(Store::new(NeedleMapKind::InMemory)) + } + + fn test_state_with_store(store: crate::storage::store::Store) -> Arc { + use crate::security::{Guard, SigningKey}; + use crate::server::volume_server::RuntimeMetricsConfig; use std::sync::RwLock; use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU32}; Arc::new(VolumeServerState { - store: RwLock::new(Store::new(NeedleMapKind::InMemory)), + store: RwLock::new(store), guard: RwLock::new(Guard::new( &[], SigningKey(vec![]), @@ -4738,4 +4797,278 @@ mod tests { assert!(saw_err, "corrupted stream must fail"); assert!(got < data.len(), "delivered {got} of {} bytes", data.len()); } + + const TEST_COOKIE: u32 = 0x1234_5678; + + /// State whose store holds volume 1 in `tmp`. + fn volume_test_state(tmp: &tempfile::TempDir) -> Arc { + use crate::storage::needle_map::NeedleMapKind; + use crate::storage::store::Store; + use crate::storage::volume::VolumeSpec; + + let dir = tmp.path().to_str().unwrap(); + let mut store = Store::new(NeedleMapKind::InMemory); + store + .add_location( + dir, + dir, + 10, + DiskType::HardDrive, + crate::config::MinFreeSpace::Percent(1.0), + Vec::new(), + ) + .unwrap(); + store + .add_volume(VolumeId(1), DiskType::HardDrive, &VolumeSpec::default()) + .unwrap(); + test_state_with_store(store) + } + + /// Write a needle to volume 1 and return its URL path. + fn put_test_needle(state: &VolumeServerState, id: u64, data: &[u8]) -> String { + put_test_needle_with(state, id, data, |_| {}) + } + + fn put_test_needle_with( + state: &VolumeServerState, + id: u64, + data: &[u8], + prepare: impl FnOnce(&mut Needle), + ) -> String { + let mut n = Needle { + id: NeedleId(id), + cookie: Cookie(TEST_COOKIE), + data: data.to_vec(), + data_size: data.len() as u32, + ..Needle::default() + }; + prepare(&mut n); + state + .store + .write() + .unwrap() + .write_volume_needle(VolumeId(1), &mut n, false) + .unwrap(); + format!("/1,{:x}{:08x}", id, TEST_COOKIE) + } + + async fn send_read( + state: &Arc, + method: Method, + path: &str, + range: Option<&str>, + ) -> (StatusCode, HeaderMap, Vec) { + use tower::ServiceExt; + let mut req = Request::builder().method(method).uri(path); + if let Some(range) = range { + req = req.header(header::RANGE, range); + } + let resp = super::super::volume_server::build_public_router(state.clone()) + .oneshot(req.body(Body::empty()).unwrap()) + .await + .unwrap(); + let status = resp.status(); + let headers = resp.headers().clone(); + let body = axum::body::to_bytes(resp.into_body(), usize::MAX) + .await + .unwrap() + .to_vec(); + (status, headers, body) + } + + /// A GET whose needle read is parked must not hold the store lock: a + /// writer, here a real append to the same volume, must get through. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_get_needle_read_runs_off_the_store_lock() { + use crate::storage::volume::needle_read_hook; + use std::sync::Mutex; + use std::sync::mpsc; + use std::time::Duration; + + const ID: u64 = 0x6e7a_0001; + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let data = b"parked needle".to_vec(); + let path = put_test_needle(&state, ID, &data); + + let (entered_tx, entered_rx) = mpsc::channel::<()>(); + let (release_tx, release_rx) = mpsc::channel::<()>(); + let release_rx = Mutex::new(release_rx); + let _hook = needle_read_hook::register(NeedleId(ID), move |_| { + let _ = entered_tx.send(()); + // Returns once released, or at once when the sender is gone. + let _ = release_rx.lock().unwrap().recv(); + }); + + let get = tokio::spawn({ + let state = state.clone(); + async move { send_read(&state, Method::GET, &path, None).await } + }); + tokio::task::spawn_blocking(move || entered_rx.recv_timeout(Duration::from_secs(10))) + .await + .unwrap() + .expect("the GET never reached its needle read"); + + let try_write_ok = state.store.try_write().is_ok(); + let (wrote_tx, wrote_rx) = mpsc::channel(); + std::thread::spawn({ + let state = state.clone(); + move || { + put_test_needle(&state, ID + 1, b"written meanwhile"); + let _ = wrote_tx.send(()); + } + }); + let write_done = tokio::task::spawn_blocking(move || { + wrote_rx.recv_timeout(Duration::from_secs(5)).is_ok() + }) + .await + .unwrap(); + drop(release_tx); + + assert!( + try_write_ok, + "a parked GET read must not hold the store lock" + ); + assert!(write_done, "an append must not wait for a parked GET read"); + let (status, _, body) = get.await.unwrap(); + assert_eq!(status, StatusCode::OK); + assert_eq!(body, data); + } + + /// A small needle is read from disk once, not once for its meta and + /// again in full; a large one never reads its payload before streaming. + #[tokio::test] + async fn test_get_reads_small_needle_once_and_large_needle_meta_only() { + use crate::storage::volume::needle_read_hook; + use std::sync::Mutex; + + const SMALL: u64 = 0x6e7a_0101; + const LARGE: u64 = 0x6e7a_0102; + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let small = b"small needle".to_vec(); + let large: Vec = (0..(STREAMING_THRESHOLD as usize + 4096)) + .map(|i| (i % 251) as u8) + .collect(); + let small_path = put_test_needle(&state, SMALL, &small); + let large_path = put_test_needle(&state, LARGE, &large); + + let small_reads = Arc::new(Mutex::new(Vec::new())); + let large_reads = Arc::new(Mutex::new(Vec::new())); + let _small_hook = needle_read_hook::register(NeedleId(SMALL), { + let reads = small_reads.clone(); + move |len| reads.lock().unwrap().push(len) + }); + let _large_hook = needle_read_hook::register(NeedleId(LARGE), { + let reads = large_reads.clone(); + move |len| reads.lock().unwrap().push(len) + }); + + let (status, _, body) = send_read(&state, Method::GET, &small_path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, small); + let reads = small_reads.lock().unwrap().clone(); + assert_eq!(reads.len(), 1, "reads of the small needle: {reads:?}"); + + let (status, _, body) = send_read(&state, Method::GET, &large_path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, large); + let reads = large_reads.lock().unwrap().clone(); + assert!( + !reads.is_empty() && reads.iter().all(|&len| len < 64 * 1024), + "the large needle's payload must not be read before streaming: {reads:?}" + ); + } + + /// HEAD and ranged reads keep answering from the needle meta and the + /// data file, and a missing or wrong-cookie needle is still a 404. + #[tokio::test] + async fn test_get_head_range_and_not_found_after_single_read() { + const ID: u64 = 0x6e7a_0201; + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let data: Vec = (0..4096u32).map(|i| (i % 251) as u8).collect(); + let path = put_test_needle(&state, ID, &data); + + let (status, headers, body) = send_read(&state, Method::HEAD, &path, None).await; + assert_eq!(status, StatusCode::OK); + assert!(body.is_empty()); + assert_eq!(headers[header::CONTENT_LENGTH], data.len().to_string()); + assert!(headers.contains_key(header::ETAG)); + + let (status, headers, body) = + send_read(&state, Method::GET, &path, Some("bytes=10-19")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + assert_eq!(body, &data[10..20]); + assert_eq!( + headers["Content-Range"], + format!("bytes 10-19/{}", data.len()) + ); + + let wrong_cookie = format!("/1,{:x}{:08x}", ID, TEST_COOKIE ^ 1); + let (status, _, _) = send_read(&state, Method::GET, &wrong_cookie, None).await; + assert_eq!(status, StatusCode::NOT_FOUND); + let missing = format!("/1,{:x}{:08x}", ID + 1, TEST_COOKIE); + let (status, _, _) = send_read(&state, Method::GET, &missing, None).await; + assert_eq!(status, StatusCode::NOT_FOUND); + } + + /// A large compressed needle cannot be streamed as stored: its meta is + /// read first, then the payload exactly once. + #[tokio::test] + async fn test_get_large_compressed_needle_reads_payload_once() { + use crate::storage::volume::needle_read_hook; + use std::sync::Mutex; + + const ID: u64 = 0x6e7a_0301; + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + // Incompressible, so the stored gzip stays above the stream threshold. + let mut x = 0x2545_f491_u32; + let plain: Vec = (0..(STREAMING_THRESHOLD as usize + 64 * 1024)) + .map(|_| { + x ^= x << 13; + x ^= x >> 17; + x ^= x << 5; + x as u8 + }) + .collect(); + let gz = try_gzip_data(&plain).unwrap(); + assert!(gz.len() > STREAMING_THRESHOLD as usize); + let path = put_test_needle_with(&state, ID, &gz, |n| n.set_is_compressed()); + + let reads = Arc::new(Mutex::new(Vec::new())); + let _hook = needle_read_hook::register(NeedleId(ID), { + let reads = reads.clone(); + move |len| reads.lock().unwrap().push(len) + }); + let (status, _, body) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, plain); + let reads = reads.lock().unwrap().clone(); + let full = reads.iter().filter(|&&len| len >= gz.len()).count(); + assert_eq!(full, 1, "reads: {reads:?}"); + } + + /// The single full read still verifies the needle checksum. + #[tokio::test] + async fn test_get_small_needle_with_bad_checksum_fails() { + const ID: u64 = 0x6e7a_0401; + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let data = b"checksummed payload, to be corrupted".to_vec(); + let path = put_test_needle(&state, ID, &data); + + let dat = tmp.path().join("1.dat"); + let mut bytes = std::fs::read(&dat).unwrap(); + let at = bytes + .windows(data.len()) + .position(|w| w == data.as_slice()) + .unwrap(); + bytes[at] ^= 0xff; + std::fs::write(&dat, bytes).unwrap(); + + let (status, _, _) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); + } } diff --git a/seaweed-volume/src/storage/store.rs b/seaweed-volume/src/storage/store.rs index 260611d52..533dd8ab4 100644 --- a/seaweed-volume/src/storage/store.rs +++ b/seaweed-volume/src/storage/store.rs @@ -715,15 +715,15 @@ impl Store { vol.read_needle_opt(n, read_deleted) } - /// Read needle metadata and return streaming info for large file reads. - pub fn read_volume_needle_stream_info( + /// Resolve a needle read under this guard, to run after it is released. + pub(crate) fn needle_read_plan( &self, vid: VolumeId, - n: &mut Needle, + id: NeedleId, read_deleted: bool, - ) -> Result { + ) -> Result { let (_, vol) = self.find_volume(vid).ok_or(VolumeError::NotFound)?; - vol.read_needle_stream_info(n, read_deleted) + vol.needle_read_plan(id, read_deleted) } /// Re-lookup a needle's data-file offset after compaction may have moved it. diff --git a/seaweed-volume/src/storage/volume.rs b/seaweed-volume/src/storage/volume.rs index 4c1572aea..d36261d71 100644 --- a/seaweed-volume/src/storage/volume.rs +++ b/seaweed-volume/src/storage/volume.rs @@ -88,9 +88,6 @@ pub enum VolumeError { #[error("IO error: {0}")] Io(#[from] io::Error), - #[error("streaming from remote-backed volume requires buffered fallback")] - StreamingUnsupported, - #[error(transparent)] Tier(#[from] crate::remote_storage::s3_tier::TierError), } @@ -166,6 +163,8 @@ fn parse_needle_at( let actual_size = get_actual_size(size, version); let mut buf = vec![0u8; actual_size as usize]; + #[cfg(test)] + needle_read_hook::fire(n.id, buf.len()); read_at(&mut buf, offset as u64)?; n.read_bytes(&buf, offset, size, version)?; @@ -489,6 +488,45 @@ impl Drop for DataFileWriteLease { } } +/// Test seam: a hook sees (and may park) every data read of its needle. +#[cfg(test)] +pub(crate) mod needle_read_hook { + use super::NeedleId; + use std::sync::{Arc, Mutex}; + + type Hook = Arc; + static HOOKS: Mutex> = Mutex::new(Vec::new()); + + pub(crate) struct Registration(NeedleId); + + impl Drop for Registration { + fn drop(&mut self) { + HOOKS.lock().unwrap().retain(|(id, _)| *id != self.0); + } + } + + /// Needle ids are the key, so tests running in parallel must use distinct ones. + pub(crate) fn register( + id: NeedleId, + hook: impl Fn(usize) + Send + Sync + 'static, + ) -> Registration { + HOOKS.lock().unwrap().push((id, Arc::new(hook))); + Registration(id) + } + + pub(crate) fn fire(id: NeedleId, len: usize) { + let hook = HOOKS + .lock() + .unwrap() + .iter() + .find(|(hook_id, _)| *hook_id == id) + .map(|(_, hook)| hook.clone()); + if let Some(hook) = hook { + hook(len); + } + } +} + /// Information needed to stream needle data directly from the dat file /// without loading the entire payload into memory. pub(crate) enum NeedleStreamSource { @@ -680,6 +718,123 @@ impl DatScanPlan { } } +/// A needle read resolved under a store guard and run after it is released; +/// the handle pins the inode, as for `DatScanPlan`. It takes no data-file +/// lease: writers wait for one while holding the store write lock. +pub(crate) struct NeedleReadPlan { + source: NeedleStreamSource, + offset: i64, + size: Size, + version: Version, + volume_id: VolumeId, + needle_id: NeedleId, + compaction_revision: u16, + data_file_access_control: Arc, + io_errors: Arc, +} + +impl NeedleReadPlan { + /// The needle's size in the index. + pub(crate) fn size(&self) -> Size { + self.size + } + + /// Read the whole needle and verify its checksum. + pub(crate) fn read_full(&mut self, n: &mut Needle) -> Result<(), VolumeError> { + let (source, size, version) = (&self.source, self.size, self.version); + let read_at = |buf: &mut [u8], at: u64| Ok(source.read_exact_at(buf, at)?); + let result = with_offset_retry(&mut self.offset, |off| { + parse_needle_at(&read_at, n, off, size, version) + }); + match &result { + Ok(()) => self.io_errors.check_read_write_error(None), + Err(VolumeError::Io(e)) => self.io_errors.check_read_write_error(Some(e)), + Err(_) => {} + } + result?; + if needle_expired(n) { + return Err(VolumeError::NotFound); + } + Ok(()) + } + + /// Read only the needle's metadata, skipping the payload. + pub(crate) fn read_meta(&mut self, n: &mut Needle) -> Result<(), VolumeError> { + let (source, size, version) = (&self.source, self.size, self.version); + let read_at = |buf: &mut [u8], at: u64| Ok(source.read_exact_at(buf, at)?); + with_offset_retry(&mut self.offset, |off| { + if version == VERSION_1 { + // A V1 body is all data, so its "meta" is the whole record. + let mut buf = vec![0u8; get_actual_size(size, version) as usize]; + #[cfg(test)] + needle_read_hook::fire(n.id, buf.len()); + read_at(&mut buf, off as u64)?; + Ok(n.read_bytes_meta_only(&buf, off, size, version)?) + } else { + Volume::read_needle_meta_blob_and_parse(read_at, n, off, size, version) + } + })?; + if needle_expired(n) { + return Err(VolumeError::NotFound); + } + Ok(()) + } + + /// For a reply served from the data file, after `read_meta` filled `n`. + pub(crate) fn into_stream_info(self, n: &Needle) -> NeedleStreamInfo { + // V1 data starts right after the header, V2/V3 after the DataSize field. + let data_file_offset = if self.version == VERSION_1 { + self.offset as u64 + NEEDLE_HEADER_SIZE as u64 + } else { + self.offset as u64 + NEEDLE_HEADER_SIZE as u64 + DATA_SIZE_SIZE as u64 + }; + NeedleStreamInfo { + source: self.source, + data_file_offset, + data_size: n.data_size, + data_file_access_control: self.data_file_access_control, + volume_id: self.volume_id, + needle_id: self.needle_id, + compaction_revision: self.compaction_revision, + checksum: n.checksum.0, + } + } +} + +/// Run `read` at `offset`, retrying past the 32 GiB wrap of 4-byte offsets +/// on a size mismatch, and leave `offset` where the needle parsed. +fn with_offset_retry( + offset: &mut i64, + mut read: impl FnMut(i64) -> Result<(), VolumeError>, +) -> Result<(), VolumeError> { + match read(*offset) { + #[cfg(not(feature = "5bytes"))] + Err(VolumeError::Needle(NeedleError::SizeMismatch { offset: o, .. })) + if o < MAX_POSSIBLE_VOLUME_SIZE as i64 => + { + *offset += MAX_POSSIBLE_VOLUME_SIZE as i64; + read(*offset) + } + result => result, + } +} + +fn needle_expired(n: &Needle) -> bool { + let Some(ttl) = n.ttl.as_ref().filter(|_| n.has_ttl()) else { + return false; + }; + let ttl_minutes = ttl.minutes(); + if ttl_minutes == 0 || !n.has_last_modified_date() { + return false; + } + let expire_at_ns = n.append_at_ns + (ttl_minutes as u64) * 60 * 1_000_000_000; + let now_ns = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_nanos() as u64; + now_ns >= expire_at_ns +} + /// Holds `Volume::is_compacting` set; dropping it clears the flag. struct CompactionClaim(Arc); @@ -1809,21 +1964,8 @@ impl Volume { Err(e) => return Err(e), } - // TTL expiry check - if n.has_ttl() - && let Some(ref ttl) = n.ttl - { - let ttl_minutes = ttl.minutes(); - if ttl_minutes > 0 && n.has_last_modified_date() { - let expire_at_ns = n.append_at_ns + (ttl_minutes as u64) * 60 * 1_000_000_000; - let now_ns = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_nanos() as u64; - if now_ns >= expire_at_ns { - return Err(VolumeError::NotFound); - } - } + if needle_expired(n) { + return Err(VolumeError::NotFound); } Ok(n.data_size as i32) @@ -1910,35 +2052,28 @@ impl Volume { size: Size, ) -> Result<(), VolumeError> { let normalized_size = if size.is_deleted() { Size(0) } else { size }; - match self.read_needle_meta_blob_and_parse(n, offset, normalized_size) { - Ok(()) => Ok(()), - #[cfg(not(feature = "5bytes"))] - Err(VolumeError::Needle(NeedleError::SizeMismatch { offset: o, .. })) - if o < MAX_POSSIBLE_VOLUME_SIZE as i64 => - { - self.read_needle_meta_blob_and_parse( - n, - offset + MAX_POSSIBLE_VOLUME_SIZE as i64, - normalized_size, - ) - } - Err(e) => Err(e), - } + let read_at = |buf: &mut [u8], at: u64| self.read_exact_at_backend(buf, at); + let version = self.version(); + let mut offset = offset; + with_offset_retry(&mut offset, |off| { + Self::read_needle_meta_blob_and_parse(read_at, n, off, normalized_size, version) + }) } fn read_needle_meta_blob_and_parse( - &self, + read_at: impl Fn(&mut [u8], u64) -> Result<(), VolumeError>, n: &mut Needle, offset: i64, size: Size, + version: Version, ) -> Result<(), VolumeError> { - let version = self.version(); - // Step 1: Read only the first 20 bytes (header + DataSize). // Matches Go's ReadNeedleMeta which reads NeedleHeaderSize+DataSizeSize first. const HEADER_PREFIX: usize = NEEDLE_HEADER_SIZE + DATA_SIZE_SIZE; // 20 let mut header_buf = [0u8; HEADER_PREFIX]; - self.read_exact_at_backend(&mut header_buf, offset as u64)?; + #[cfg(test)] + needle_read_hook::fire(n.id, header_buf.len()); + read_at(&mut header_buf, offset as u64)?; // Parse header to get the needle's Size field for validation let (_, _, found_size) = Needle::parse_header(&header_buf); @@ -1967,7 +2102,9 @@ impl Volume { ))); } let mut meta_buf = vec![0u8; meta_size as usize]; - self.read_exact_at_backend(&mut meta_buf, (offset + NEEDLE_HEADER_SIZE as i64) as u64)?; + #[cfg(test)] + needle_read_hook::fire(n.id, meta_buf.len()); + read_at(&mut meta_buf, (offset + NEEDLE_HEADER_SIZE as i64) as u64)?; n.read_paged_meta(&header_buf, &meta_buf, offset, size, version)?; } else { // V2/V3: extract DataSize from bytes 16..20 @@ -1997,121 +2134,79 @@ impl Volume { // Step 3: Read only the meta tail (skip the data payload entirely) let mut meta_buf = vec![0u8; meta_size as usize]; - self.read_exact_at_backend(&mut meta_buf, start_offset as u64)?; + #[cfg(test)] + needle_read_hook::fire(n.id, meta_buf.len()); + read_at(&mut meta_buf, start_offset as u64)?; n.read_paged_meta(&header_buf, &meta_buf, offset, size, version)?; } Ok(()) } - /// Read needle metadata (header + flags/name/mime/etc) without loading the data payload, - /// and return a `NeedleStreamInfo` that can be used to stream data directly from the dat file. - /// - /// This is used for large needles to avoid loading the entire payload into memory. - pub fn read_needle_stream_info( + /// Resolve a needle read under the caller's store guard, without data I/O. + pub(crate) fn needle_read_plan( &self, - n: &mut Needle, + id: NeedleId, read_deleted: bool, - ) -> Result { - let _guard = self.data_file_access_control.read_lock(); + ) -> Result { if let Some(e) = self.unavailable_error() { return Err(e); } let nm = self.nm_or_not_found()?; - let nv = nm.get(n.id)?.ok_or(VolumeError::NotFound)?; + let nv = nm.get(id)?.ok_or(VolumeError::NotFound)?; if nv.offset.is_zero() { return Err(VolumeError::NotFound); } - let mut read_size = nv.size; - if read_size.is_deleted() { - if read_deleted && !read_size.is_tombstone() { - debug!("reading deleted {}", n.id); + let mut size = nv.size; + if size.is_deleted() { + if read_deleted && !size.is_tombstone() { + debug!("reading deleted {}", id); crate::metrics::HANDLER_COUNTER .with_label_values(&[crate::metrics::READ_DELETED_NEEDLE]) .inc(); - read_size = Size(-read_size.0); + size = Size(-size.0); } else { return Err(VolumeError::Deleted); } } - if read_size.0 == 0 { + if size.0 == 0 { return Err(VolumeError::NotFound); } - #[cfg_attr(feature = "5bytes", allow(unused_mut))] - let mut offset = nv.offset.to_actual_offset(); - let version = self.version(); - let actual_size = get_actual_size(read_size, version); - - // Read the full needle bytes (including data) for metadata parsing. - // We use read_bytes_meta_only which skips copying the data payload. - #[cfg_attr(feature = "5bytes", allow(unused_mut))] - let mut read_and_parse = |off: i64| -> Result<(), VolumeError> { - let mut buf = vec![0u8; actual_size as usize]; - self.read_exact_at_backend(&mut buf, off as u64)?; - n.read_bytes_meta_only(&buf, off, read_size, version)?; - Ok(()) - }; - - match read_and_parse(offset) { - Ok(()) => {} - #[cfg(not(feature = "5bytes"))] - Err(VolumeError::Needle(NeedleError::SizeMismatch { offset: o, .. })) - if o < MAX_POSSIBLE_VOLUME_SIZE as i64 => - { - offset += MAX_POSSIBLE_VOLUME_SIZE as i64; - read_and_parse(offset)?; - } - Err(e) => return Err(e), - } - - // TTL expiry check - if n.has_ttl() - && let Some(ref ttl) = n.ttl - { - let ttl_minutes = ttl.minutes(); - if ttl_minutes > 0 && n.has_last_modified_date() { - let expire_at_ns = n.append_at_ns + (ttl_minutes as u64) * 60 * 1_000_000_000; - let now_ns = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap_or_default() - .as_nanos() as u64; - if now_ns >= expire_at_ns { - return Err(VolumeError::NotFound); - } - } - } - - // For V1, data starts right after the header - // For V2/V3, data starts at header + 4 (DataSize field) - let data_file_offset = if version == VERSION_1 { - offset as u64 + NEEDLE_HEADER_SIZE as u64 - } else { - offset as u64 + NEEDLE_HEADER_SIZE as u64 + 4 // skip DataSize (4 bytes) - }; - - let source = match (self.dat_file.as_ref(), self.remote_dat_file.as_ref()) { - (Some(dat_file), _) => { - NeedleStreamSource::Local(dat_file.try_clone().map_err(VolumeError::Io)?) - } - (None, Some(remote_dat_file)) => NeedleStreamSource::Remote(remote_dat_file.clone()), - (None, None) => return Err(VolumeError::StreamingUnsupported), - }; - - Ok(NeedleStreamInfo { - source, - data_file_offset, - data_size: n.data_size, - data_file_access_control: self.data_file_access_control.clone(), + Ok(NeedleReadPlan { + source: self.read_source()?, + offset: nv.offset.to_actual_offset(), + size, + version: self.version(), volume_id: self.id, - needle_id: n.id, + needle_id: id, compaction_revision: self.super_block.compaction_revision, - checksum: n.checksum.0, + data_file_access_control: self.data_file_access_control.clone(), + io_errors: self.io_errors.clone(), }) } + /// A read handle on the data backend that stays valid without the guard. + fn read_source(&self) -> Result { + if self.dat_file.is_some() { + // A fresh open, not `try_clone`: a duplicated handle shares the + // file position, and on Windows `read_exact_at` goes through + // `seek_read`, which moves it under a concurrent append. Opened + // while the caller's guard excludes a vacuum swap, the path names + // the inode `dat_file` holds. + Ok(NeedleStreamSource::Local(open_volume_file( + OpenOptions::new().read(true), + self.file_name(".dat"), + )?)) + } else if let Some(remote) = self.remote_dat_file() { + Ok(NeedleStreamSource::Remote(remote)) + } else { + Err(VolumeError::Io(io::Error::other("dat file not open"))) + } + } + /// Re-lookup a needle's data-file offset after compaction may have moved it. /// /// Returns `(new_data_file_offset, current_compaction_revision)` or an error @@ -3247,23 +3342,8 @@ impl Volume { if let Some(e) = self.unavailable_error() { return Err(e); } - let source = if self.dat_file.is_some() { - // A fresh open, not `try_clone`: a duplicated handle shares the - // file position, and on Windows `read_exact_at` goes through - // `seek_read`, which moves it under a concurrent append. Opened - // while the caller's guard excludes a vacuum swap, the path names - // the inode `dat_file` holds. - NeedleStreamSource::Local(open_volume_file( - OpenOptions::new().read(true), - self.file_name(".dat"), - )?) - } else if let Some(remote) = self.remote_dat_file() { - NeedleStreamSource::Remote(remote) - } else { - return Err(VolumeError::Io(io::Error::other("dat file not open"))); - }; Ok(DatScanPlan { - source, + source: self.read_source()?, version: self.version(), from: from_offset, end: self.current_dat_file_size()?, @@ -6618,7 +6698,11 @@ mod tests { .unwrap(); let (_count, broken) = v.scrub().unwrap(); - assert!(broken.is_empty(), "healthy needle must pass, got {:?}", broken); + assert!( + broken.is_empty(), + "healthy needle must pass, got {:?}", + broken + ); let mut id_bytes = [0u8; NEEDLE_ID_SIZE]; NeedleId(99).to_bytes(&mut id_bytes); @@ -7650,7 +7734,9 @@ mod tests { cookie: Cookie(0xDEADBEEF), ..Needle::default() }; - let info = v.read_needle_stream_info(&mut read_n, false).unwrap(); + let mut plan = v.needle_read_plan(NeedleId(42), false).unwrap(); + plan.read_meta(&mut read_n).unwrap(); + let info = plan.into_stream_info(&read_n); assert_eq!(info.volume_id, VolumeId(1)); assert_eq!(info.needle_id, NeedleId(42)); @@ -8990,7 +9076,9 @@ mod tests { id: NeedleId(7), ..Needle::default() }; - let info = v.read_needle_stream_info(&mut meta, false).unwrap(); + let mut plan = v.needle_read_plan(NeedleId(7), false).unwrap(); + plan.read_meta(&mut meta).unwrap(); + let info = plan.into_stream_info(&meta); assert!(matches!(info.source, NeedleStreamSource::Remote(_))); let mut streamed = vec![0u8; info.data_size as usize]; info.source