volume server: read GET/HEAD needles off the store lock, and only once (#11487)

The GET/HEAD handler read the needle synchronously on the tokio worker
while holding store.read(): first a stream-info read that loaded the
whole record just to parse its meta, then, for every needle that was not
streamed (small, compressed, chunk manifest, image ops), a second full
read. For a tiered volume each read is an S3 GET under the store lock,
and a writer queued behind it parks every other store reader.

The regular-volume read now runs in spawn_blocking. Under the store guard
it only resolves a NeedleReadPlan (index lookup, a freshly opened .dat
handle or the remote backend, offset, size); the guard is dropped before
any needle data I/O. No data-file lease is held across the read either,
since a writer waits for one while holding the store write lock. The
index size decides the read, as in Go's readNeedle: a HEAD, a ranged read
or a needle above the stream threshold reads only its header and meta
tail (ReadNeedleMeta) and hands off to StreamingBody or the range path;
everything else is read in full once, with its checksum verified. A
compressed or manifest needle found by the meta read is then read in
full once. The range-from-source read also moves to spawn_blocking.

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
This commit is contained in:
Eliah Rusin
2026-09-27 20:13:15 +08:00
committed by GitHub
co-authored by Claude Opus 5.5 Chris Lu
parent 67691a1eea
commit e57f8c4d87
3 changed files with 636 additions and 215 deletions
+404 -71
View File
@@ -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::NeedleStreamInfo>)>,
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<VolumeServerState> {
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<VolumeServerState> {
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<VolumeServerState> {
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<VolumeServerState>,
method: Method,
path: &str,
range: Option<&str>,
) -> (StatusCode, HeaderMap, Vec<u8>) {
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<u8> = (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<u8> = (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<u8> = (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);
}
}
+5 -5
View File
@@ -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<crate::storage::volume::NeedleStreamInfo, VolumeError> {
) -> Result<crate::storage::volume::NeedleReadPlan, VolumeError> {
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.
+227 -139
View File
@@ -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<dyn Fn(usize) + Send + Sync>;
static HOOKS: Mutex<Vec<(NeedleId, Hook)>> = 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<DataFileAccessControl>,
io_errors: Arc<IoErrorTracker>,
}
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<AtomicBool>);
@@ -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<NeedleStreamInfo, VolumeError> {
let _guard = self.data_file_access_control.read_lock();
) -> Result<NeedleReadPlan, VolumeError> {
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<NeedleStreamSource, VolumeError> {
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