mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-06 06:22:05 +02:00
volume server: read a non-ASCII or empty Range header as Go does (#11534)
* 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> --------- 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:
1 file changed
+123
-25
@@ -1077,7 +1077,12 @@ async fn get_or_head_handler_inner(
|
||||
};
|
||||
|
||||
let read_deleted = query.read_deleted.as_deref() == Some("true");
|
||||
let (ext, request_kind) = parse_read_request(&path, &headers, &query, &method);
|
||||
// Go reads Range once and ignores an empty value; a non-UTF-8 byte fails its parse.
|
||||
let range = headers
|
||||
.get(header::RANGE)
|
||||
.map(|v| String::from_utf8_lossy(v.as_bytes()))
|
||||
.filter(|r| !r.is_empty());
|
||||
let (ext, request_kind) = parse_read_request(&path, range.is_some(), &query, &method);
|
||||
|
||||
// EC volumes always do a full read (no streaming/meta-only).
|
||||
let plan = if has_ec_volume && !has_volume {
|
||||
@@ -1128,19 +1133,14 @@ async fn get_or_head_handler_inner(
|
||||
return head_from_meta_response(&info, response_headers);
|
||||
}
|
||||
ReadStrategy::RangeFromSource(info) => {
|
||||
if let Some(range_header) = headers.get(header::RANGE)
|
||||
&& let Ok(range_str) = range_header.to_str()
|
||||
{
|
||||
return range_from_source_response(
|
||||
&state,
|
||||
range_str,
|
||||
info,
|
||||
response_headers,
|
||||
track_download,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
// An unreadable Range header falls through to the buffered path.
|
||||
return range_from_source_response(
|
||||
&state,
|
||||
range.as_deref().unwrap_or_default(),
|
||||
info,
|
||||
response_headers,
|
||||
track_download,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
ReadStrategy::Buffered => {}
|
||||
}
|
||||
@@ -1158,7 +1158,7 @@ async fn get_or_head_handler_inner(
|
||||
};
|
||||
buffered_response(
|
||||
&state,
|
||||
&headers,
|
||||
range.as_deref(),
|
||||
&method,
|
||||
data,
|
||||
response_headers,
|
||||
@@ -1304,11 +1304,10 @@ async fn check_download_limit(
|
||||
/// The URL extension, and how the reply may be served.
|
||||
fn parse_read_request(
|
||||
path: &str,
|
||||
headers: &HeaderMap,
|
||||
has_range: bool,
|
||||
query: &ReadQueryParams,
|
||||
method: &Method,
|
||||
) -> (String, SourceReadRequest) {
|
||||
let has_range = headers.contains_key(header::RANGE);
|
||||
let ext = extract_extension_from_path(path);
|
||||
// Go checks resize and crop extensions separately: resize supports .webp, crop does not.
|
||||
let has_resize_ops = is_image_resize_ext(&ext)
|
||||
@@ -1846,7 +1845,7 @@ fn buffered_payload(
|
||||
/// The buffered reply over the payload, whole or a range.
|
||||
fn buffered_response(
|
||||
state: &Arc<VolumeServerState>,
|
||||
headers: &HeaderMap,
|
||||
range: Option<&str>,
|
||||
method: &Method,
|
||||
data: Vec<u8>,
|
||||
mut response_headers: HeaderMap,
|
||||
@@ -1863,11 +1862,9 @@ fn buffered_response(
|
||||
return (StatusCode::OK, response_headers).into_response();
|
||||
}
|
||||
|
||||
if let Some(range_header) = headers.get(header::RANGE)
|
||||
&& let Ok(range_str) = range_header.to_str()
|
||||
{
|
||||
if let Some(range) = range {
|
||||
return handle_range_request(
|
||||
range_str,
|
||||
range,
|
||||
&data,
|
||||
response_headers,
|
||||
track_download.then(|| state.clone()),
|
||||
@@ -5091,12 +5088,15 @@ mod tests {
|
||||
state: &Arc<VolumeServerState>,
|
||||
method: Method,
|
||||
path: &str,
|
||||
range: Option<&str>,
|
||||
range: Option<&[u8]>,
|
||||
) -> (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);
|
||||
req = req.header(
|
||||
header::RANGE,
|
||||
header::HeaderValue::from_bytes(range).unwrap(),
|
||||
);
|
||||
}
|
||||
let resp = super::super::volume_server::build_public_router(state.clone())
|
||||
.oneshot(req.body(Body::empty()).unwrap())
|
||||
@@ -5232,7 +5232,7 @@ mod tests {
|
||||
assert!(headers.contains_key(header::ETAG));
|
||||
|
||||
let (status, headers, body) =
|
||||
send_read(&state, Method::GET, &path, Some("bytes=10-19")).await;
|
||||
send_read(&state, Method::GET, &path, Some(b"bytes=10-19")).await;
|
||||
assert_eq!(status, StatusCode::PARTIAL_CONTENT);
|
||||
assert_eq!(body, &data[10..20]);
|
||||
assert_eq!(
|
||||
@@ -5248,6 +5248,104 @@ mod tests {
|
||||
assert_eq!(status, StatusCode::NOT_FOUND);
|
||||
}
|
||||
|
||||
/// Range is read as Go reads it: a byte its parseRange cannot parse is a
|
||||
/// 416, Unicode whitespace around a range is trimmed, an empty value is no
|
||||
/// range, and HEAD ignores it.
|
||||
#[tokio::test]
|
||||
async fn test_get_range_header_is_read_as_go_reads_it() {
|
||||
const PLAIN: u64 = 0x6e7a_0501;
|
||||
const GZIPPED: u64 = 0x6e7a_0502;
|
||||
let tmp = tempfile::TempDir::new().unwrap();
|
||||
let state = volume_test_state(&tmp);
|
||||
let data = b"range over a needle".to_vec();
|
||||
let plain_path = put_test_needle(&state, PLAIN, &data);
|
||||
let gz = try_gzip_data(&data).unwrap();
|
||||
let gz_path = put_test_needle_with(&state, GZIPPED, &gz, |n| n.set_is_compressed());
|
||||
|
||||
for path in [&plain_path, &gz_path] {
|
||||
let (status, _, body) =
|
||||
send_read(&state, Method::GET, path, Some(b"bytes=0-1\xff")).await;
|
||||
assert_eq!(status, StatusCode::RANGE_NOT_SATISFIABLE, "{path}");
|
||||
assert_eq!(body, b"invalid range", "{path}");
|
||||
|
||||
let (status, _, _) =
|
||||
send_read(&state, Method::HEAD, path, Some(b"bytes=0-1\xff")).await;
|
||||
assert_eq!(status, StatusCode::OK, "{path}");
|
||||
|
||||
for range in ["bytes=0-1\u{a0}", "bytes=\u{85}0-1"] {
|
||||
let (status, headers, body) =
|
||||
send_read(&state, Method::GET, path, Some(range.as_bytes())).await;
|
||||
assert_eq!(status, StatusCode::PARTIAL_CONTENT, "{path} {range:?}");
|
||||
assert_eq!(body, &data[..2], "{path} {range:?}");
|
||||
assert_eq!(
|
||||
headers["Content-Range"],
|
||||
format!("bytes 0-1/{}", data.len())
|
||||
);
|
||||
}
|
||||
|
||||
let (status, _, body) = send_read(&state, Method::GET, path, Some(b"")).await;
|
||||
assert_eq!(status, StatusCode::OK, "{path}");
|
||||
assert_eq!(body, data, "{path}");
|
||||
}
|
||||
}
|
||||
|
||||
/// An EC needle is served from memory: its GET parses Range like any
|
||||
/// other, and its HEAD ignores it as Go's writeResponseContent does.
|
||||
#[tokio::test]
|
||||
async fn test_ec_needle_range_and_head() {
|
||||
use crate::storage::erasure_coding::ec_encoder::write_ec_files;
|
||||
use crate::storage::erasure_coding::ec_shard::ShardId;
|
||||
use crate::storage::needle_map::NeedleMapKind;
|
||||
use crate::storage::volume::{Volume, VolumeSpec};
|
||||
|
||||
const ID: u64 = 0x6e7a_0601;
|
||||
let tmp = tempfile::TempDir::new().unwrap();
|
||||
let state = volume_test_state(&tmp);
|
||||
let dir = tmp.path().to_str().unwrap();
|
||||
let data: Vec<u8> = (0..4096u32).map(|i| (i % 251) as u8).collect();
|
||||
let mut v = Volume::new(
|
||||
dir,
|
||||
dir,
|
||||
VolumeId(2),
|
||||
NeedleMapKind::InMemory,
|
||||
&VolumeSpec::default(),
|
||||
)
|
||||
.unwrap();
|
||||
let mut n = Needle {
|
||||
id: NeedleId(ID),
|
||||
cookie: Cookie(TEST_COOKIE),
|
||||
data: data.clone(),
|
||||
data_size: data.len() as u32,
|
||||
..Needle::default()
|
||||
};
|
||||
v.write_needle(&mut n, true, false).unwrap();
|
||||
v.sync_to_disk().unwrap();
|
||||
v.close();
|
||||
write_ec_files(dir, dir, "", VolumeId(2), 10, 4).unwrap();
|
||||
let shard_ids: Vec<ShardId> = (0..14).collect();
|
||||
state
|
||||
.store
|
||||
.write()
|
||||
.unwrap()
|
||||
.mount_ec_shards(VolumeId(2), "", &shard_ids)
|
||||
.unwrap();
|
||||
let path = format!("/2,{:x}{:08x}", ID, TEST_COOKIE);
|
||||
|
||||
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, _, body) = send_read(&state, Method::GET, &path, Some(b"bytes=10-19")).await;
|
||||
assert_eq!(status, StatusCode::PARTIAL_CONTENT);
|
||||
assert_eq!(body, &data[10..20]);
|
||||
|
||||
for range in [&b"bytes=10-19"[..], b"bytes=0-1\xff"] {
|
||||
let (status, headers, _) = send_read(&state, Method::HEAD, &path, Some(range)).await;
|
||||
assert_eq!(status, StatusCode::OK);
|
||||
assert_eq!(headers[header::CONTENT_LENGTH], data.len().to_string());
|
||||
}
|
||||
}
|
||||
|
||||
/// A large compressed needle cannot be streamed as stored: its meta is
|
||||
/// read first, then the payload exactly once.
|
||||
#[tokio::test]
|
||||
|
||||
Reference in new issue
Block a user