volume server: fetch only the chunks a manifest Range needs, like Go (#11546)

* volume server: read GET/HEAD needles off the store lock, and only once

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>

* volume server: stream needle chunks without the store lock

StreamingBody::poll_frame took store.read() and find_volume for every
chunk to compare the volume's compaction revision, dup'd the source
handle, and allocated a fresh chunk buffer. With -hasSlowRead=false the
stream also holds a data-file read lease for its whole life, while a
writer waits for that lease under store.write(): the next chunk's
store.read() then waits for the writer and the writer for the stream.

The per-chunk re-lookup was also wrong. The stream reads a handle opened
at plan time, which pins the .dat inode the offset was resolved against;
a vacuum commit renames a new file over .dat and leaves that inode
untouched. The re-looked-up offset belongs to the new file but was read
from the old inode, so a stream whose needle a vacuum moved ended in a
checksum error. The pinned offset stays valid, so the check, and with
it every store access, is dropped, along with the now unused
re_lookup_needle_data_offset and the revision fields of the read plan.

The source is shared as an Arc instead of dup'd per chunk, and the chunk
buffer is a BytesMut that the blocking read hands back with its result,
so its allocation is reclaimed once the previous frame has been written.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* volume server: split get_or_head_handler_inner into phases

get_or_head_handler_inner was a ~650-line function. Its middle resolved
the needle and set five mutable flags (stream_info, can_stream,
can_handle_head_from_meta, can_handle_range_from_source, bypass_cm) that
three if-let reply paths then re-tested, each re-checking stream_info.

It is now a 126-line orchestrator over named phases: reject_read_jwt,
proxy_missing_volume, wait_for_download_slot, parse_read_request,
read_ec_needle / read_volume_needle, etag_and_last_modified,
not_modified_response, read_response_headers, and the reply phases
stream_response, head_from_meta_response, range_from_source_response,
buffered_payload and buffered_response. The read phases return a
ReadPlan whose ReadStrategy enum (Stream, HeadFromMeta, RangeFromSource,
Buffered) carries the NeedleStreamInfo only on the variants that use it,
so the reply is one match instead of three flag checks.

Pure refactor: every status code, header and header order, error text,
metric increment, lock and data-file lease scope, spawn_blocking
boundary and side-effect order is unchanged. Phases that can end the
request return ControlFlow<Response, T>. A Range header that is not
visible ASCII still falls through to the buffered path, as before.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* volume server: stop a needle stream once its volume becomes unavailable

Taking the store lock out of StreamingBody also dropped its per-chunk
unavailable_error() check. With -hasSlowRead a writer can take the
data-file lease between chunks, fail its fsync and its truncate, and mark
the volume unavailable; the stream then kept serving the rest of the
needle from its pinned handle.

The volume's io_unavailable reason is now an Arc-shared leaf mutex that
the read plan hands to the stream. Each chunk checks it under its
data-file lease, where the writer marks it, and fails with the same
"volume is unavailable: <reason>" error the old check returned.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* volume server: read a non-ASCII or empty Range header as Go does

A Range value with a byte >= 0x80 (obs-text, which hyper accepts) failed
HeaderValue::to_str at both range gates. For a needle served as stored the
handler had already chosen a meta-only read, so it fell through to the
buffered path with no payload and answered 200 with an empty body; a
compressed or EC needle answered 200 with the full body. Go's parseRange
fails on a byte it can neither trim nor parse and answers 416 "invalid
range", and trims Unicode whitespace such as NBSP into a normal 206. An
empty Range value was also a 200 with an empty body, where Go sends the
whole payload.

Read Range once with from_utf8_lossy, dropping an empty value, and hand
that one value to the read plan and to both range gates. A replaced byte
never parses, so it is a 416; str::trim trims the same Unicode whitespace
as strings.TrimSpace. A range read from the data file now always answers
itself instead of falling through with an empty needle. The buffered path
answers HEAD before it looks at Range, as Go's writeResponseContent does,
so an EC HEAD with a Range is a 200 with the full length.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* volume server: apply Range to chunk manifests and forward raw headers when proxying

A GET of a chunk manifest assembled the object and always answered 200 with
the whole body, ignoring Range. Go serves the expanded manifest through
writeResponseContent, which answers HEAD first and then hands Range to
ProcessRangeRequest: 206 for one range, multipart/byteranges for several,
416 for an unsatisfiable or unparsable one. try_expand_chunk_manifest now
returns the assembled body and headers, and the caller answers through
buffered_response, the same helper the buffered needle path uses.

A proxied read forwarded a request header only if HeaderValue::to_str
succeeded, so a Range with an obs-text byte was dropped and the target
answered 200 with the full body. Go copies every header value as is.
Forward the raw HeaderValue for every header.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* volume server: fetch only the chunks a manifest Range needs, like Go

A ranged GET of a chunk manifest fetched every chunk, assembled the whole
object and then sliced it, so reading a few bytes of a large object cost a
read of all of it, and a 416 still fetched everything. Go serves a manifest
through ChunkedFileReader, which seeks to each range and reads only the
chunks under it.

For a GET with a Range, try_expand_chunk_manifest now parses the ranges
against the manifest size, fetches only the chunks whose declared window
overlaps one of them (none when the reply carries no body), and answers
through handle_range_request_with, the reader-based core that
handle_range_request now wraps, so 206/416/multipart stay one code path.
The reader replays assembly: chunks clamped as before, later chunks over
earlier ones, zeros in gaps. HEAD, no-Range GETs and GETs that crop or
resize an image still assemble the whole object.

A missing chunk outside the requested ranges no longer turns a ranged GET
into a 500, as in Go; a missing chunk inside them still does, before any
headers are sent.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* volume server: reword a comment codespell flags

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* volume server: keep only range-covered bytes of fetched manifest chunks

A ranged GET retained every overlapping chunk's full contents; 1,000
overlapping 8 MiB chunks could pin ~8 GiB for a one-byte response. Clip
each fetched chunk to the bytes the requested ranges can actually read,
preserving the later-chunks-overwrite and zero-fill-gap semantics.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* volume server: bucket ranged manifest parts by range

Serving a multipart range scanned every retained part. Bucket the kept
intersections by their range so one range only reads its own parts.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
authored and GitHub committed 2026-10-05 21:17:19 +08:00
1 parent a7a590474d
commit 825c3dff8b
1 file changed
+379 -80
+379 -80
View File
@@ -900,11 +900,9 @@ async fn proxy_request(
// Build the proxy request
let mut req_builder = state.http_client.get(&target_url);
// Forward all original headers
// Forward all original headers as raw bytes, as Go does.
for (name, value) in &info.original_headers {
if let Ok(v) = value.to_str() {
req_builder = req_builder.header(name.as_str(), v);
}
req_builder = req_builder.header(name.clone(), value.clone());
}
let resp = match req_builder.send().await {
@@ -1107,18 +1105,29 @@ async fn get_or_head_handler_inner(
// responses keep it (Go sets it before tryHandleChunkedFile runs).
if n.is_chunk_manifest()
&& !request_kind.bypass_cm
&& let Some(resp) = try_expand_chunk_manifest(
&& let Some(expanded) = try_expand_chunk_manifest(
&state,
&n,
&method,
&path,
&query,
&etag,
&last_modified_str,
range.as_deref().filter(|_| method == Method::GET),
)
.await
{
return resp;
// HEAD, and a Range the manifest path did not answer, apply to the assembled object.
return match expanded {
ControlFlow::Continue((data, response_headers)) => buffered_response(
&state,
range.as_deref(),
&method,
data,
response_headers,
false,
),
ControlFlow::Break(resp) => resp,
};
}
// If manifest expansion fails (invalid JSON etc.), fall through to raw data
@@ -1975,10 +1984,27 @@ fn range_error_response(mut headers: HeaderMap, msg: &str) -> Response {
fn handle_range_request(
range_str: &str,
data: &[u8],
headers: HeaderMap,
state: Option<Arc<VolumeServerState>>,
) -> Response {
handle_range_request_with(
range_str,
data.len() as i64,
|buf, start, end| buf.extend_from_slice(&data[start..end]),
headers,
state,
)
}
/// Range reply over an object of `total` bytes; `read` appends `[start, end)`
/// of it and is called only for the ranges actually served.
fn handle_range_request_with(
range_str: &str,
total: i64,
read: impl Fn(&mut Vec<u8>, usize, usize),
mut headers: HeaderMap,
state: Option<Arc<VolumeServerState>>,
) -> Response {
let total = data.len() as i64;
let ranges = match parse_range_header(range_str, total) {
Ok(r) => r,
Err(msg) => {
@@ -2014,10 +2040,9 @@ fn handle_range_request(
if r.length <= 0 {
return (StatusCode::PARTIAL_CONTENT, headers).into_response();
}
let start = r.start as usize;
let end = (r.start + r.length) as usize;
let slice = &data[start..end];
finalize_bytes_response(StatusCode::PARTIAL_CONTENT, headers, slice.to_vec(), state)
let mut slice = Vec::with_capacity(r.length as usize);
read(&mut slice, r.start as usize, (r.start + r.length) as usize);
finalize_bytes_response(StatusCode::PARTIAL_CONTENT, headers, slice, state)
} else {
// Multi-range: build multipart/byteranges response
let boundary = "SeaweedFSBoundary";
@@ -2040,9 +2065,7 @@ fn handle_range_request(
format!("Content-Range: {}\r\n\r\n", range_content_range(*r, total)).as_bytes(),
);
if r.length > 0 {
let start = r.start as usize;
let end = (r.start + r.length) as usize;
body.extend_from_slice(&data[start..end]);
read(&mut body, r.start as usize, (r.start + r.length) as usize);
}
}
body.extend_from_slice(format!("\r\n--{}--\r\n", boundary).as_bytes());
@@ -3575,27 +3598,29 @@ struct ChunkInfo {
size: i64,
}
/// Try to expand a chunk manifest needle. Returns None if manifest can't be parsed.
/// Try to expand a chunk manifest needle into its assembled body and reply
/// headers. The Range of a GET (`get_range`) is answered here from only the
/// chunks it overlaps. Returns None if manifest can't be parsed.
async fn try_expand_chunk_manifest(
state: &Arc<VolumeServerState>,
n: &Needle,
method: &Method,
path: &str,
query: &ReadQueryParams,
etag: &str,
last_modified_str: &Option<String>,
) -> Option<Response> {
get_range: Option<&str>,
) -> Option<ControlFlow<Response, (Vec<u8>, HeaderMap)>> {
let data = if n.is_compressed() {
match maybe_decompress_gzip(&n.data) {
Ok(d) => d,
Err(GunzipError::TooLarge) => {
return Some(
return Some(ControlFlow::Break(
(
StatusCode::PAYLOAD_TOO_LARGE,
"compressed manifest exceeds decompression limit",
)
.into_response(),
);
));
}
Err(GunzipError::Decode) => return None,
}
@@ -3615,53 +3640,107 @@ async fn try_expand_chunk_manifest(
return None;
}
// Read and concatenate all chunks. Each chunk is resolved to wherever it
// lives — a local regular volume, a local EC volume (reconstruct-on-read),
// or a peer via master lookup — mirroring Go's ChunkedFileReader, which
// never assumes chunks are local regular needles.
let mut result = vec![0u8; manifest.size as usize];
for chunk in &manifest.chunks {
// Validate the attacker-controlled chunk offset before indexing: a
// negative value would wrap to a huge usize, and an out-of-range one has
// nowhere to land.
if chunk.offset < 0 || chunk.size < 0 {
return Some(
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("invalid negative chunk offset/size in {}", chunk.fid),
)
.into_response(),
);
}
let data = match read_chunk_needle(state, &chunk.fid).await {
Ok(d) => d,
Err(e) => {
return Some(
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("read chunk {}: {}", chunk.fid, e),
)
.into_response(),
);
}
};
let offset = chunk.offset as usize;
if offset >= result.len() {
continue;
}
// Clamp to the chunk's declared size so an over-long chunk can't bleed
// into the next chunk's window; also drop bytes past the buffer end.
let bound = (chunk.size as usize).min(result.len() - offset);
let copy_len = data.len().min(bound);
result[offset..offset + copy_len].copy_from_slice(&data[..copy_len]);
}
// Determine filename: URL path filename, then manifest name
// (Go's tryHandleChunkedFile does NOT fall back to needle name)
let mut filename = extract_filename_from_path(path);
if filename.is_empty() && !manifest.name.is_empty() {
filename = manifest.name.clone();
}
let cm_ext = if !filename.is_empty() {
if let Some(dot_pos) = filename.rfind('.') {
filename[dot_pos..].to_lowercase()
} else {
String::new()
}
} else {
String::new()
};
let transforms_image =
(is_image_crop_ext(&cm_ext) && query.crop_x2.is_some() && query.crop_y2.is_some())
|| (is_image_resize_ext(&cm_ext)
&& (query.width.unwrap_or(0) != 0 || query.height.unwrap_or(0) != 0));
let size = manifest.size as usize;
// A ranged GET needs only the chunks its ranges overlap; none for a
// range that gets no body.
let wanted = get_range.filter(|_| !transforms_image).map(|range| {
let ranges = parse_range_header(range, manifest.size)
.ok()
.filter(|ranges| sum_ranges_size(ranges) <= manifest.size)
.unwrap_or_default();
(range, ranges)
});
let mut parts: HashMap<(usize, usize), Vec<(usize, Vec<u8>)>> = HashMap::new();
// Read and concatenate all chunks. Each chunk is resolved to wherever it
// lives — a local regular volume, a local EC volume (reconstruct-on-read),
// or a peer via master lookup — mirroring Go's ChunkedFileReader, which
// never assumes chunks are local regular needles.
let mut result = if wanted.is_some() {
Vec::new()
} else {
vec![0u8; size]
};
for chunk in &manifest.chunks {
// Validate the attacker-controlled chunk offset before indexing: a
// negative value would wrap to a huge usize, and an out-of-range one has
// nowhere to land.
if chunk.offset < 0 || chunk.size < 0 {
return Some(ControlFlow::Break(
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("invalid negative chunk offset/size in {}", chunk.fid),
)
.into_response(),
));
}
if let Some((_, ranges)) = &wanted {
let end = chunk.offset.saturating_add(chunk.size);
if !ranges
.iter()
.any(|r| r.length > 0 && chunk.offset < r.start + r.length && r.start < end)
{
continue;
}
}
let data = match read_chunk_needle(state, &chunk.fid).await {
Ok(d) => d,
Err(e) => {
return Some(ControlFlow::Break(
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("read chunk {}: {}", chunk.fid, e),
)
.into_response(),
));
}
};
let offset = chunk.offset as usize;
if offset >= size {
continue;
}
// Clamp to the chunk's declared size so an over-long chunk can't bleed
// into the next chunk's window; also drop bytes past the buffer end.
let bound = (chunk.size as usize).min(size - offset);
let copy_len = data.len().min(bound);
if let Some((_, ranges)) = &wanted {
// Keep only the bytes each requested range can read, bucketed by
// range so serving one range never scans another's parts.
let key = |r: &HttpRange| (r.start as usize, (r.start + r.length) as usize);
for r in ranges {
let lo = (r.start as usize).clamp(offset, offset + copy_len);
let hi = ((r.start + r.length) as usize).clamp(offset, offset + copy_len);
if lo < hi {
parts
.entry(key(r))
.or_default()
.push((lo, data[lo - offset..hi - offset].to_vec()));
}
}
} else {
result[offset..offset + copy_len].copy_from_slice(&data[..copy_len]);
}
}
// Determine MIME type: manifest mime, but fall back to extension detection
// if empty or application/octet-stream (matching Go behavior)
@@ -3700,7 +3779,6 @@ async fn try_expand_chunk_manifest(
}
response_headers.insert(header::CONTENT_TYPE, content_type.parse().unwrap());
response_headers.insert("X-File-Store", "chunked".parse().unwrap());
response_headers.insert(header::ACCEPT_RANGES, "bytes".parse().unwrap());
// Last-Modified — Go sets this on the response writer before tryHandleChunkedFile
if let Some(lm) = last_modified_str
@@ -3769,17 +3847,33 @@ async fn try_expand_chunk_manifest(
}
}
if let Some((range, _)) = wanted {
response_headers.insert(header::ACCEPT_RANGES, "bytes".parse().unwrap());
// Later chunks overwrite earlier ones and gaps read as zeros, as in assembly.
let read = |buf: &mut Vec<u8>, start: usize, end: usize| {
let base = buf.len();
buf.resize(base + end - start, 0);
if let Some(range_parts) = parts.get(&(start, end)) {
for (offset, data) in range_parts {
let (lo, hi) = (start.max(*offset), end.min(offset + data.len()));
if lo < hi {
buf[base + lo - start..base + hi - start]
.copy_from_slice(&data[lo - offset..hi - offset]);
}
}
}
};
return Some(ControlFlow::Break(handle_range_request_with(
range,
manifest.size,
read,
response_headers,
None,
)));
}
// Go's tryHandleChunkedFile applies crop then resize to expanded chunk data
// (L344-345: conditionallyCropImages, conditionallyResizeImages).
let cm_ext = if !filename.is_empty() {
if let Some(dot_pos) = filename.rfind('.') {
filename[dot_pos..].to_lowercase()
} else {
String::new()
}
} else {
String::new()
};
if is_image_crop_ext(&cm_ext) {
result = maybe_crop_image(&result, &cm_ext, query);
}
@@ -3787,15 +3881,7 @@ async fn try_expand_chunk_manifest(
result = maybe_resize_image(&result, &cm_ext, query);
}
if *method == Method::HEAD {
response_headers.insert(
header::CONTENT_LENGTH,
result.len().to_string().parse().unwrap(),
);
return Some((StatusCode::OK, response_headers).into_response());
}
Some((StatusCode::OK, response_headers, result).into_response())
Some(ControlFlow::Continue((result, response_headers)))
}
/// Read one chunk-manifest chunk's final (decompressed) content bytes from
@@ -5346,6 +5432,219 @@ mod tests {
}
}
/// A chunk manifest applies Range to the assembled object, as Go's
/// writeResponseContent does; its HEAD is still read from the meta.
#[tokio::test]
async fn test_chunk_manifest_range() {
let tmp = tempfile::TempDir::new().unwrap();
let state = volume_test_state(&tmp);
let data = b"0123456789abcdef";
let a = put_test_needle(&state, 0x6e7a_0701, &data[..8]);
let b = put_test_needle(&state, 0x6e7a_0702, &data[8..]);
let manifest = format!(
r#"{{"name":"obj.txt","size":16,"chunks":[{{"fid":"{}","offset":0,"size":8}},{{"fid":"{}","offset":8,"size":8}}]}}"#,
&a[1..],
&b[1..]
);
let path = put_test_needle_with(&state, 0x6e7a_0703, manifest.as_bytes(), |n| {
n.set_is_chunk_manifest()
});
let (status, _, body) = send_read(&state, Method::GET, &path, None).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(body, data);
let (status, headers, body) =
send_read(&state, Method::GET, &path, Some(b"bytes=5-12")).await;
assert_eq!(status, StatusCode::PARTIAL_CONTENT);
assert_eq!(body, &data[5..13]);
assert_eq!(headers["Content-Range"], "bytes 5-12/16");
assert_eq!(headers["X-File-Store"], "chunked");
assert!(headers.contains_key(header::ETAG));
let (status, headers, body) =
send_read(&state, Method::GET, &path, Some(b"bytes=0-1,14-15")).await;
assert_eq!(status, StatusCode::PARTIAL_CONTENT);
assert!(
headers[header::CONTENT_TYPE]
.to_str()
.unwrap()
.starts_with("multipart/byteranges")
);
let body = String::from_utf8(body).unwrap();
assert!(body.contains("Content-Range: bytes 0-1/16\r\n\r\n01"));
assert!(body.contains("Content-Range: bytes 14-15/16\r\n\r\nef"));
let (status, headers, _) = send_read(&state, Method::GET, &path, Some(b"bytes=16-")).await;
assert_eq!(status, StatusCode::RANGE_NOT_SATISFIABLE);
assert_eq!(headers["Content-Range"], "bytes */16");
let (status, _, body) = send_read(&state, Method::GET, &path, Some(b"bytes=0-1\xff")).await;
assert_eq!(status, StatusCode::RANGE_NOT_SATISFIABLE);
assert_eq!(body, b"invalid range");
let (status, headers, _) = send_read(&state, Method::HEAD, &path, None).await;
assert_eq!(status, StatusCode::OK);
let (range_status, range_headers, _) =
send_read(&state, Method::HEAD, &path, Some(b"bytes=0-9")).await;
assert_eq!(range_status, status);
assert_eq!(
range_headers[header::CONTENT_LENGTH],
headers[header::CONTENT_LENGTH]
);
}
/// A ranged GET of a chunk manifest fetches only the chunks its ranges
/// overlap, as Go's ChunkedFileReader does.
#[tokio::test]
async fn test_chunk_manifest_range_fetches_only_overlapping_chunks() {
use crate::storage::volume::needle_read_hook;
use std::sync::atomic::{AtomicUsize, Ordering};
const IDS: [u64; 3] = [0x6e7a_0901, 0x6e7a_0902, 0x6e7a_0903];
let tmp = tempfile::TempDir::new().unwrap();
let state = volume_test_state(&tmp);
let data = b"chunk-1:chunk-2:chunk-3:";
let fids: Vec<String> = IDS
.iter()
.zip(data.chunks(8))
.map(|(&id, part)| put_test_needle(&state, id, part)[1..].to_string())
.collect();
let manifest = format!(
r#"{{"name":"obj.txt","size":24,"chunks":[{{"fid":"{}","offset":0,"size":8}},{{"fid":"{}","offset":8,"size":8}},{{"fid":"{}","offset":16,"size":8}}]}}"#,
fids[0], fids[1], fids[2]
);
let path = put_test_needle_with(&state, 0x6e7a_0904, manifest.as_bytes(), |n| {
n.set_is_chunk_manifest()
});
let reads: Arc<[AtomicUsize; 3]> = Arc::new(Default::default());
let _hooks: Vec<_> = IDS
.iter()
.enumerate()
.map(|(i, &id)| {
let reads = reads.clone();
needle_read_hook::register(NeedleId(id), move |_| {
reads[i].fetch_add(1, Ordering::SeqCst);
})
})
.collect();
let fetched =
|| -> [usize; 3] { std::array::from_fn(|i| reads[i].swap(0, Ordering::SeqCst)) };
let (status, _, body) = send_read(&state, Method::GET, &path, None).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(body, data);
assert_eq!(fetched(), [1, 1, 1]);
let (status, headers, body) =
send_read(&state, Method::GET, &path, Some(b"bytes=10-13")).await;
assert_eq!(status, StatusCode::PARTIAL_CONTENT);
assert_eq!(body, &data[10..14]);
assert_eq!(headers["Content-Range"], "bytes 10-13/24");
assert_eq!(headers[header::CONTENT_LENGTH], "4");
assert_eq!(headers[header::ACCEPT_RANGES], "bytes");
assert_eq!(headers["X-File-Store"], "chunked");
assert!(headers.contains_key(header::ETAG));
assert_eq!(fetched(), [0, 1, 0]);
let (status, _, body) = send_read(&state, Method::GET, &path, Some(b"bytes=4-11")).await;
assert_eq!(status, StatusCode::PARTIAL_CONTENT);
assert_eq!(body, &data[4..12]);
assert_eq!(fetched(), [1, 1, 0]);
let (status, headers, body) =
send_read(&state, Method::GET, &path, Some(b"bytes=0-1,20-23")).await;
assert_eq!(status, StatusCode::PARTIAL_CONTENT);
assert!(
headers[header::CONTENT_TYPE]
.to_str()
.unwrap()
.starts_with("multipart/byteranges")
);
let body = String::from_utf8(body).unwrap();
assert!(body.contains("Content-Range: bytes 0-1/24\r\n\r\nch"));
assert!(body.contains("Content-Range: bytes 20-23/24\r\n\r\nk-3:\r\n"));
assert_eq!(fetched(), [1, 0, 1]);
let (status, headers, _) = send_read(&state, Method::GET, &path, Some(b"bytes=24-")).await;
assert_eq!(status, StatusCode::RANGE_NOT_SATISFIABLE);
assert_eq!(headers["Content-Range"], "bytes */24");
assert_eq!(fetched(), [0, 0, 0]);
}
/// Ranges served from fetched chunks read what assembly would: clamped
/// chunks, later ones overwriting earlier ones, and zero-filled gaps.
#[tokio::test]
async fn test_chunk_manifest_range_matches_assembly() {
let tmp = tempfile::TempDir::new().unwrap();
let state = volume_test_state(&tmp);
let a = put_test_needle(&state, 0x6e7a_0905, b"AAAAAAAA");
let b = put_test_needle(&state, 0x6e7a_0906, b"BBBB");
let c = put_test_needle(&state, 0x6e7a_0907, b"CCCC");
let manifest = format!(
r#"{{"size":12,"chunks":[{{"fid":"{}","offset":0,"size":6}},{{"fid":"{}","offset":4,"size":4}},{{"fid":"{}","offset":10,"size":4}}]}}"#,
&a[1..],
&b[1..],
&c[1..]
);
let path = put_test_needle_with(&state, 0x6e7a_0908, manifest.as_bytes(), |n| {
n.set_is_chunk_manifest()
});
let (status, _, full) = send_read(&state, Method::GET, &path, None).await;
assert_eq!(status, StatusCode::OK);
assert_eq!(full, b"AAAABBBB\0\0CC");
for start in 0..12 {
for end in start..12 {
let range = format!("bytes={start}-{end}");
let (status, _, body) =
send_read(&state, Method::GET, &path, Some(range.as_bytes())).await;
assert_eq!(status, StatusCode::PARTIAL_CONTENT, "{range}");
assert_eq!(body, &full[start..=end], "{range}");
}
}
}
/// A proxied read forwards request headers as raw bytes, as Go does: a
/// Range the target cannot parse must reach it and come back as its 416.
#[tokio::test]
async fn test_proxy_forwards_non_ascii_range() {
let tmp = tempfile::TempDir::new().unwrap();
let target_state = volume_test_state(&tmp);
let path = put_test_needle(&target_state, 0x6e7a_0801, b"proxied payload");
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let app = super::super::volume_server::build_public_router(target_state);
let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
let proxy_state = test_state_with_store(crate::storage::store::Store::new(
crate::storage::needle_map::NeedleMapKind::InMemory,
));
let target = VolumeLocation {
url: addr.to_string(),
public_url: addr.to_string(),
grpc_port: 0,
read_only: false,
read_only_can_delete: false,
};
let mut headers = HeaderMap::new();
headers.insert(
header::RANGE,
header::HeaderValue::from_bytes(b"bytes=0-1\xff").unwrap(),
);
let info = build_proxy_request_info(&path, &headers, "").unwrap();
let resp = proxy_request(&proxy_state, &info, &target).await;
assert_eq!(resp.status(), StatusCode::RANGE_NOT_SATISFIABLE);
let body = axum::body::to_bytes(resp.into_body(), usize::MAX)
.await
.unwrap();
assert_eq!(&body[..], b"invalid range");
server.abort();
}
/// A large compressed needle cannot be streamed as stored: its meta is
/// read first, then the payload exactly once.
#[tokio::test]