diff --git a/seaweed-volume/src/server/handlers.rs b/seaweed-volume/src/server/handlers.rs index 664204c29..7ed9c91f4 100644 --- a/seaweed-volume/src/server/handlers.rs +++ b/seaweed-volume/src/server/handlers.rs @@ -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, - headers: &HeaderMap, + range: Option<&str>, method: &Method, data: Vec, 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, method: Method, path: &str, - range: Option<&str>, + range: Option<&[u8]>, ) -> (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); + 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 = (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 = (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]