From 924a683397dc7650ead3c0de6c3fbf69d5cb24c9 Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Sat, 10 Oct 2026 05:05:00 +0300 Subject: [PATCH] volume server: answer chunk-manifest reads from the sum of chunk sizes (#11628) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * volume: answer chunk-manifest HEAD from the chunk-size sum * volume: size chunk-manifest reads by the sum of chunk sizes * volume: saturate the sum of chunk-manifest Range lengths * volume: fetch remote chunk-manifest ranges with a bounded Range * volume: fmt chunk-manifest handlers * volume: cap whole-window chunk reads at the visible window * volume: reject an over-cap chunk span and an illegal manifest MIME * volume: take the manifest transform extension from the URL first Go's tryHandleChunkedFile keeps the extension parsed from the URL — including a fid suffix like /vid,fid.png — and only falls back to the resolved filename's extension when the URL has none. The manifest path now does the same, so a transform requested on a fid-suffixed URL runs as it does in Go. Also: compute the chunk-size sum once, drop an unused remote_read_for parameter, share the crop/resize predicate with parse_read_request, and settle a couple of clippy lints. --------- Co-authored-by: Chris Lu --- seaweed-volume/src/server/handlers.rs | 1530 +++++++++++++++++++++---- 1 file changed, 1325 insertions(+), 205 deletions(-) diff --git a/seaweed-volume/src/server/handlers.rs b/seaweed-volume/src/server/handlers.rs index b85332dae..8441c9430 100644 --- a/seaweed-volume/src/server/handlers.rs +++ b/seaweed-volume/src/server/handlers.rs @@ -161,7 +161,11 @@ impl SourceReadRequest { /// 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)) + // HEAD is meta-only unless the payload is needed to expand a manifest. + if self.is_head { + return !(n.is_chunk_manifest() && !self.bypass_cm); + } + self.direct(n) && (self.has_range || n.data_size > STREAMING_THRESHOLD) } /// Meta first only if the needle could be served from source; its data @@ -1101,8 +1105,7 @@ async fn get_or_head_handler_inner( return resp; } - // Chunk manifest expansion needs the full data; ETag is passed so expanded - // responses keep it (Go sets it before tryHandleChunkedFile runs). + // Chunk manifest expansion needs the full data. if n.is_chunk_manifest() && !request_kind.bypass_cm && let Some(expanded) = try_expand_chunk_manifest( @@ -1110,9 +1113,8 @@ async fn get_or_head_handler_inner( &n, &path, &query, - &etag, - &last_modified_str, range.as_deref().filter(|_| method == Method::GET), + &method, ) .await { @@ -1319,21 +1321,10 @@ fn parse_read_request( ) -> (String, SourceReadRequest) { 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) - && (query.width.unwrap_or(0) > 0 || query.height.unwrap_or(0) > 0); - // shouldCropImages requires x2 > x1 && y2 > y1; only a real crop disables streaming. - let has_crop_ops = is_image_crop_ext(&ext) && { - let x1 = query.crop_x1.unwrap_or(0); - let y1 = query.crop_y1.unwrap_or(0); - let x2 = query.crop_x2.unwrap_or(0); - let y2 = query.crop_y2.unwrap_or(0); - x2 > x1 && y2 > y1 - }; - let has_image_ops = has_resize_ops || has_crop_ops; let request_kind = SourceReadRequest { is_head: method == Method::HEAD, has_range, - has_image_ops, + has_image_ops: has_image_ops(&ext, query), bypass_cm: query.cm.as_deref() == Some("false"), }; (ext, request_kind) @@ -1454,7 +1445,8 @@ async fn read_volume_needle( let can_stream = can_direct_source_read && n.data_size > STREAMING_THRESHOLD && !has_range && !is_head; - // Go uses meta-only reads for all HEAD requests, regardless of compression/chunked files. + // HEAD is meta-only when served_from_source loaded no payload. A chunk + // manifest without cm=false is read in full so expansion can run. let can_handle_head_from_meta = stream_info.is_some() && is_head; let can_handle_range_from_source = can_direct_source_read && has_range; @@ -1961,7 +1953,13 @@ fn parse_range_header(s: &str, size: i64) -> Result, &'static str } fn sum_ranges_size(ranges: &[HttpRange]) -> i64 { - ranges.iter().map(|r| r.length).sum() + ranges.iter().fold(0i64, |acc, r| { + if r.length > 0 { + acc.saturating_add(r.length) + } else { + acc + } + }) } fn range_content_range(r: HttpRange, total: i64) -> String { @@ -1975,10 +1973,7 @@ fn range_error_response(mut headers: HeaderMap, msg: &str) -> Response { header::CONTENT_TYPE, "text/plain; charset=utf-8".parse().unwrap(), ); - headers.insert( - header::X_CONTENT_TYPE_OPTIONS, - "nosniff".parse().unwrap(), - ); + headers.insert(header::X_CONTENT_TYPE_OPTIONS, "nosniff".parse().unwrap()); let mut response = Response::new(Body::from(msg.to_string())); *response.status_mut() = StatusCode::RANGE_NOT_SATISFIABLE; *response.headers_mut() = headers; @@ -2208,6 +2203,19 @@ fn extract_filename_from_path(path: &str) -> String { } } +fn has_image_ops(ext: &str, query: &ReadQueryParams) -> bool { + let resize = + is_image_resize_ext(ext) && (query.width.unwrap_or(0) > 0 || query.height.unwrap_or(0) > 0); + let crop = is_image_crop_ext(ext) && { + let x1 = query.crop_x1.unwrap_or(0); + let y1 = query.crop_y1.unwrap_or(0); + let x2 = query.crop_x2.unwrap_or(0); + let y2 = query.crop_y2.unwrap_or(0); + x2 > x1 && y2 > y1 + }; + resize || crop +} + fn path_base(path: &str) -> String { let trimmed = path.trim_end_matches('/'); trimmed @@ -3602,93 +3610,23 @@ struct ChunkInfo { size: i64, } -/// 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, - n: &Needle, - path: &str, - query: &ReadQueryParams, - etag: &str, - last_modified_str: &Option, - get_range: Option<&str>, -) -> Option, HeaderMap)>> { - let data = if n.is_compressed() { - match maybe_decompress_gzip(&n.data) { - Ok(d) => d, - Err(GunzipError::TooLarge) => { - return Some(ControlFlow::Break( - ( - StatusCode::PAYLOAD_TOO_LARGE, - "compressed manifest exceeds decompression limit", - ) - .into_response(), - )); - } - Err(GunzipError::Decode) => return None, - } - } else { - n.data.clone() - }; - - let manifest: ChunkManifest = match serde_json::from_slice(&data) { - Ok(m) => m, - Err(_) => return None, - }; - - // Guard the attacker-controlled manifest size before allocating: a negative - // value would wrap to a huge usize (capacity-overflow panic) and an oversized - // one would OOM-kill the server. - if manifest.size < 0 || manifest.size as u64 > MAX_EXPANSION_BYTES { - return None; - } - - // Determine filename: URL path filename, then manifest name - // (Go's tryHandleChunkedFile does NOT fall back to needle name) +fn manifest_filename(path: &str, manifest_name: &str) -> String { let mut filename = extract_filename_from_path(path); - if filename.is_empty() && !manifest.name.is_empty() { - filename = manifest.name.clone(); + if filename.is_empty() && !manifest_name.is_empty() { + filename = manifest_name.to_string(); } - 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)); + filename +} - 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)>> = HashMap::new(); +fn extension_of(filename: &str) -> String { + filename + .rfind('.') + .map(|dot| filename[dot..].to_lowercase()) + .unwrap_or_default() +} - // 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. +fn reject_bad_chunks(chunks: &[ChunkInfo]) -> Option, HeaderMap)>> { + for chunk in chunks { if chunk.offset < 0 || chunk.size < 0 { return Some(ControlFlow::Break( ( @@ -3698,78 +3636,59 @@ async fn try_expand_chunk_manifest( .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]); - } } + None +} - // Determine MIME type: manifest mime, but fall back to extension detection - // if empty or application/octet-stream (matching Go behavior) +/// `None` is a u64 overflow or a sum that does not fit in i64. +fn sum_chunk_sizes(chunks: &[ChunkInfo]) -> Option { + let mut sum: u64 = 0; + for chunk in chunks { + if chunk.size < 0 { + return None; + } + sum = sum.checked_add(chunk.size as u64)?; + } + if sum > i64::MAX as u64 { + return None; + } + Some(sum as i64) +} + +fn exceeds_expansion_limit(buffered: u64) -> bool { + buffered > MAX_EXPANSION_BYTES +} + +fn manifest_too_large() -> ControlFlow, HeaderMap)> { + ControlFlow::Break( + ( + StatusCode::PAYLOAD_TOO_LARGE, + "chunk manifest exceeds expansion limit", + ) + .into_response(), + ) +} + +fn chunk_manifest_response_headers( + n: &Needle, + query: &ReadQueryParams, + etag: &str, + last_modified_str: &Option, + filename: &str, + manifest_mime: &str, +) -> HeaderMap { let content_type = { - let mime_str = &manifest.mime; - if !mime_str.is_empty() && !mime_str.starts_with("application/octet-stream") { - mime_str.clone() + if !manifest_mime.is_empty() && !manifest_mime.starts_with("application/octet-stream") { + manifest_mime.to_string() } else { - // Try to detect from filename extension - let ext = if !filename.is_empty() { - if let Some(dot_pos) = filename.rfind('.') { - filename[dot_pos..].to_lowercase() - } else { - String::new() - } - } else { - String::new() - }; + let ext = extension_of(filename); if !ext.is_empty() { mime_guess::from_ext(ext.trim_start_matches('.')) .first() .map(|m| m.to_string()) .unwrap_or_else(|| "application/octet-stream".to_string()) - } else if !mime_str.is_empty() { - mime_str.clone() + } else if !manifest_mime.is_empty() { + manifest_mime.to_string() } else { "application/octet-stream".to_string() } @@ -3777,21 +3696,22 @@ async fn try_expand_chunk_manifest( }; let mut response_headers = HeaderMap::new(); - // Preserve ETag from the needle (matches Go: ETag is set before tryHandleChunkedFile) if let Ok(etag_val) = etag.parse() { response_headers.insert(header::ETAG, etag_val); } - response_headers.insert(header::CONTENT_TYPE, content_type.parse().unwrap()); + // `mime` is stored JSON; fall back when it is not a legal header value. + let content_type_value = content_type + .parse() + .unwrap_or_else(|_| header::HeaderValue::from_static("application/octet-stream")); + response_headers.insert(header::CONTENT_TYPE, content_type_value); response_headers.insert("X-File-Store", "chunked".parse().unwrap()); - // Last-Modified — Go sets this on the response writer before tryHandleChunkedFile if let Some(lm) = last_modified_str && let Ok(hval) = lm.parse() { response_headers.insert(header::LAST_MODIFIED, hval); } - // Pairs — Go sets needle pairs on the response writer before tryHandleChunkedFile if n.has_pairs() && !n.pairs.is_empty() && let Ok(pair_map) = @@ -3807,7 +3727,6 @@ async fn try_expand_chunk_manifest( } } - // S3 response passthrough headers — Go sets these via AdjustPassthroughHeaders if let Some(ref cc) = query.response_cache_control && let Ok(hval) = cc.parse() { @@ -3834,7 +3753,6 @@ async fn try_expand_chunk_manifest( response_headers.insert(header::CONTENT_DISPOSITION, hval); } - // Content-Disposition if !filename.is_empty() { let disposition_type = if let Some(ref dl_val) = query.dl { if parse_go_bool(dl_val).unwrap_or(false) { @@ -3845,15 +3763,183 @@ async fn try_expand_chunk_manifest( } else { "inline" }; - let disposition = format_content_disposition(disposition_type, &filename); + let disposition = format_content_disposition(disposition_type, filename); if let Ok(hval) = disposition.parse() { response_headers.insert(header::CONTENT_DISPOSITION, hval); } } + response_headers +} + +/// 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, + n: &Needle, + path: &str, + query: &ReadQueryParams, + get_range: Option<&str>, + method: &Method, +) -> Option, HeaderMap)>> { + let data = if n.is_compressed() { + match maybe_decompress_gzip(&n.data) { + Ok(d) => d, + Err(GunzipError::TooLarge) => { + return Some(ControlFlow::Break( + ( + StatusCode::PAYLOAD_TOO_LARGE, + "compressed manifest exceeds decompression limit", + ) + .into_response(), + )); + } + Err(GunzipError::Decode) => return None, + } + } else { + n.data.clone() + }; + + let manifest: ChunkManifest = match serde_json::from_slice(&data) { + Ok(m) => m, + Err(_) => return None, + }; + + if let Some(flow) = reject_bad_chunks(&manifest.chunks) { + return Some(flow); + } + + let Some(sum) = sum_chunk_sizes(&manifest.chunks) else { + return Some(manifest_too_large()); + }; + + // Go's tryHandleChunkedFile: ext comes from the URL (including a fid + // suffix like /vid,fid.png) and falls back to the resolved filename's. + let filename = manifest_filename(path, &manifest.name); + let url_ext = extract_extension_from_path(path); + let cm_ext = if url_ext.is_empty() { + extension_of(&filename) + } else { + url_ext + }; + let transforms_image = has_image_ops(&cm_ext, query); + let (etag, last_modified_str) = etag_and_last_modified(n); + let mut response_headers = chunk_manifest_response_headers( + n, + query, + &etag, + &last_modified_str, + &filename, + &manifest.mime, + ); + if *method == Method::HEAD && !transforms_image { + response_headers.insert(header::ACCEPT_RANGES, "bytes".parse().unwrap()); + response_headers.insert(header::CONTENT_LENGTH, sum.to_string().parse().unwrap()); + return Some(ControlFlow::Break( + (StatusCode::OK, response_headers).into_response(), + )); + } + + let mut wanted: Option<(&str, Vec)> = None; + if transforms_image { + if exceeds_expansion_limit(sum as u64) { + return Some(manifest_too_large()); + } + } else if let Some(range) = get_range { + let parsed = parse_range_header(range, sum); + let fetch = match &parsed { + Err(_) => false, + Ok(ranges) if ranges.is_empty() || sum_ranges_size(ranges) > sum => false, + Ok(ranges) => { + let covered = sum_ranges_size(ranges); + if covered > 0 && exceeds_expansion_limit(covered as u64) { + return Some(manifest_too_large()); + } + if buffered_span_exceeds_limit(&manifest.chunks, sum, ranges) { + return Some(manifest_too_large()); + } + true + } + }; + if !fetch { + response_headers.insert(header::ACCEPT_RANGES, "bytes".parse().unwrap()); + return Some(ControlFlow::Break(handle_range_request_with( + range, + sum, + |_buf: &mut Vec, _, _| {}, + response_headers, + None, + ))); + } + wanted = Some((range, parsed.unwrap())); + } else if exceeds_expansion_limit(sum as u64) { + return Some(manifest_too_large()); + } + + let size = sum as usize; + type Parts = HashMap<(usize, usize), Vec<(usize, Vec)>>; + let mut parts: Parts = HashMap::new(); + let mut result = if wanted.is_some() { + Vec::new() + } else { + vec![0u8; size] + }; + for chunk in &manifest.chunks { + let Some(window) = chunk_window(chunk.offset, chunk.size, sum) else { + continue; + }; + let ranges = wanted.as_ref().map(|(_, ranges)| ranges.as_slice()); + let Some(how) = remote_read_for(chunk.offset, window, ranges) else { + continue; + }; + let fetched = match read_chunk_needle(state, &chunk.fid, how).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; + let file_lo = offset + fetched.origin; + if file_lo >= size { + continue; + } + let window_end = offset + window as usize; + let copy_len = fetched + .data + .len() + .min(window_end.saturating_sub(file_lo)) + .min(size - file_lo); + if let Some((_, ranges)) = &wanted { + let file_hi = file_lo + copy_len; + let key = |r: &HttpRange| (r.start as usize, (r.start + r.length) as usize); + for r in ranges { + if r.length <= 0 { + continue; + } + let rs = r.start as usize; + let re = rs + r.length as usize; + let lo = rs.max(file_lo); + let hi = re.min(file_hi); + if lo < hi { + parts + .entry(key(r)) + .or_default() + .push((lo, fetched.data[lo - file_lo..hi - file_lo].to_vec())); + } + } + } else { + result[file_lo..file_lo + copy_len].copy_from_slice(&fetched.data[..copy_len]); + } + } 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, start: usize, end: usize| { let base = buf.len(); buf.resize(base + end - start, 0); @@ -3869,31 +3955,140 @@ async fn try_expand_chunk_manifest( }; return Some(ControlFlow::Break(handle_range_request_with( range, - manifest.size, + sum, read, response_headers, None, ))); } - // Go's tryHandleChunkedFile applies crop then resize to expanded chunk data - // (L344-345: conditionallyCropImages, conditionallyResizeImages). if is_image_crop_ext(&cm_ext) { result = maybe_crop_image(&result, &cm_ext, query); } if is_image_resize_ext(&cm_ext) { result = maybe_resize_image(&result, &cm_ext, query); } - Some(ControlFlow::Continue((result, response_headers))) } +#[derive(Clone, Copy)] +enum RemoteRead { + /// GET with no Range. Keep at most `cap` bytes. Origin 0. + Whole { cap: u64 }, + /// `Range: bytes={lo}-{hi-1}`. The returned bytes start at chunk offset `lo`. + Range { lo: u64, hi: u64 }, +} + +struct ChunkBytes { + /// Chunk offset of `data[0]`. + origin: usize, + data: Vec, +} + +fn chunk_window(offset: i64, size: i64, sum: i64) -> Option { + if size <= 0 || offset >= sum { + return None; + } + Some(size.min(sum - offset)) +} + +/// Overlap of `[offset, offset+window)` with `ranges`, as chunk-relative `[lo, hi)`. +fn overlap_span(offset: i64, window: i64, ranges: &[HttpRange]) -> Option<(i64, i64)> { + let file_hi = offset + window; + let mut lo = i64::MAX; + let mut hi = i64::MIN; + for r in ranges { + if r.length <= 0 { + continue; + } + let start = r.start.max(offset); + let end = (r.start + r.length).min(file_hi); + if start < end { + lo = lo.min(start); + hi = hi.max(end); + } + } + if lo < hi { + Some((lo - offset, hi - offset)) + } else { + None + } +} + +/// One fetched span over `MAX_EXPANSION_BYTES`. The covered-length check does +/// not see the gap between two ranges on the same chunk. +fn buffered_span_exceeds_limit(chunks: &[ChunkInfo], sum: i64, ranges: &[HttpRange]) -> bool { + chunks.iter().any(|chunk| { + let Some(window) = chunk_window(chunk.offset, chunk.size, sum) else { + return false; + }; + let Some((lo, hi)) = overlap_span(chunk.offset, window, ranges) else { + return false; + }; + exceeds_expansion_limit((hi - lo) as u64) + }) +} + +fn remote_read_for(offset: i64, window: i64, ranges: Option<&[HttpRange]>) -> Option { + let (lo, hi) = match ranges { + None => (0, window), + Some(ranges) => overlap_span(offset, window, ranges)?, + }; + if lo == 0 && hi == window { + Some(RemoteRead::Whole { cap: window as u64 }) + } else { + Some(RemoteRead::Range { + lo: lo as u64, + hi: hi as u64, + }) + } +} + +fn content_range_first_byte(headers: &HeaderMap) -> Option { + let values: Vec<_> = headers.get_all(header::CONTENT_RANGE).iter().collect(); + if values.len() != 1 { + return None; + } + let value = values[0].to_str().ok()?; + let rest = value.trim().strip_prefix("bytes ")?; + let (start_end, _) = rest.split_once('/')?; + let (start, end) = start_end.split_once('-')?; + if start.trim().is_empty() || start.trim() == "*" || end.trim().is_empty() { + return None; + } + start.trim().parse().ok() +} + +async fn read_limited_body(resp: &mut reqwest::Response, cap: usize) -> Result, String> { + let mut out = Vec::new(); + while out.len() < cap { + match resp.chunk().await { + Ok(Some(chunk)) => { + let room = cap - out.len(); + if chunk.len() <= room { + out.extend_from_slice(&chunk); + } else { + out.extend_from_slice(&chunk[..room]); + break; + } + } + Ok(None) => break, + Err(e) => return Err(format!("read body: {e}")), + } + } + Ok(out) +} + /// Read one chunk-manifest chunk's final (decompressed) content bytes from /// wherever it lives: a local regular volume, a local EC volume /// (reconstruct-on-read from surviving shards), or a peer resolved via the /// master. Mirrors Go's ChunkedFileReader, which looks every chunk up through /// the master instead of assuming a local regular needle. -async fn read_chunk_needle(state: &Arc, fid: &str) -> Result, String> { +async fn read_chunk_needle( + state: &Arc, + fid: &str, + how: RemoteRead, +) -> Result { let (vid, nid, cookie) = parse_url_path(fid).ok_or_else(|| format!("invalid chunk fid: {}", fid))?; @@ -3914,7 +4109,8 @@ async fn read_chunk_needle(state: &Arc, fid: &str) -> Result< return store .read_volume_needle(vid, &mut n) .map_err(|e| format!("{}", e)) - .and_then(|_| cookie_checked_chunk(n, cookie)); + .and_then(|_| cookie_checked_chunk(n, cookie)) + .map(|data| ChunkBytes { origin: 0, data }); } else if store.find_ec_volume(vid).is_some() { Placement::Ec } else { @@ -3925,13 +4121,15 @@ async fn read_chunk_needle(state: &Arc, fid: &str) -> Result< match placement { Placement::Ec => { match crate::server::store_ec::read_ec_shard_needle_distributed(state, vid, nid).await { - Ok(Some(n)) => cookie_checked_chunk(n, cookie), + Ok(Some(n)) => { + cookie_checked_chunk(n, cookie).map(|data| ChunkBytes { origin: 0, data }) + } Ok(None) => Err("not found".to_string()), Err(e) => Err(format!("{}", e)), } } // The peer serves through its own GET handler, which validates the cookie. - Placement::Remote => read_remote_chunk_needle(state, vid, fid).await, + Placement::Remote => read_remote_chunk_needle(state, vid, fid, how).await, } } @@ -3970,7 +4168,8 @@ async fn read_remote_chunk_needle( state: &Arc, vid: VolumeId, fid: &str, -) -> Result, String> { + how: RemoteRead, +) -> Result { let locations = lookup_volume( &state.http_client, &state.outgoing_http_scheme, @@ -3999,13 +4198,39 @@ async fn read_remote_chunk_needle( continue; } }; - match state.http_client.get(&url).send().await { - Ok(resp) if resp.status().is_success() => match resp.bytes().await { - Ok(b) => return Ok(b.to_vec()), - Err(e) => last_err = format!("read body from {}: {}", url, e), - }, - Ok(resp) => last_err = format!("{} returned {}", url, resp.status()), - Err(e) => last_err = format!("request to {} failed: {}", url, e), + let mut req = state.http_client.get(&url); + let (origin, cap) = match how { + RemoteRead::Whole { cap } => (0usize, cap as usize), + RemoteRead::Range { lo, hi } => { + req = req.header(header::RANGE, format!("bytes={lo}-{}", hi - 1)); + (lo as usize, (hi - lo) as usize) + } + }; + match req.send().await { + Ok(mut resp) => { + let status = resp.status(); + let err = match how { + RemoteRead::Whole { .. } if status.is_success() => None, + RemoteRead::Range { lo, .. } + if status.as_u16() == StatusCode::PARTIAL_CONTENT.as_u16() => + { + match content_range_first_byte(resp.headers()) { + Some(start) if start == lo => None, + _ => Some(format!("{url} Content-Range does not start at {lo}")), + } + } + _ => Some(format!("{url} returned {status}")), + }; + if let Some(err) = err { + last_err = err; + continue; + } + match read_limited_body(&mut resp, cap).await { + Ok(data) => return Ok(ChunkBytes { origin, data }), + Err(e) => last_err = format!("read body from {url}: {e}"), + } + } + Err(e) => last_err = format!("request to {url} failed: {e}"), } } Err(if last_err.is_empty() { @@ -5174,6 +5399,103 @@ mod tests { format!("/1,{:x}{:08x}", id, TEST_COOKIE) } + fn watch_reads( + ids: &[u64], + ) -> ( + Vec, + std::sync::Arc<[std::sync::atomic::AtomicUsize]>, + ) { + use crate::storage::volume::needle_read_hook; + use std::sync::atomic::AtomicUsize; + let reads: std::sync::Arc<[AtomicUsize]> = std::sync::Arc::from( + std::iter::repeat_with(|| AtomicUsize::new(0)) + .take(ids.len()) + .collect::>() + .into_boxed_slice(), + ); + let regs = ids + .iter() + .enumerate() + .map(|(i, &id)| { + let reads = reads.clone(); + needle_read_hook::register(NeedleId(id), move |_| { + reads[i].fetch_add(1, std::sync::atomic::Ordering::SeqCst); + }) + }) + .collect(); + (regs, reads) + } + + fn taken(reads: &std::sync::Arc<[std::sync::atomic::AtomicUsize]>) -> Vec { + use std::sync::atomic::Ordering; + reads.iter().map(|c| c.swap(0, Ordering::SeqCst)).collect() + } + + fn put_manifest(state: &VolumeServerState, id: u64, json: &str) -> String { + put_test_needle_with(state, id, json.as_bytes(), |n| n.set_is_chunk_manifest()) + } + + fn put_gzip_manifest(state: &VolumeServerState, id: u64, json: &str) -> String { + let gz = try_gzip_data(json.as_bytes()).unwrap(); + put_test_needle_with(state, id, &gz, |n| { + n.set_is_chunk_manifest(); + n.set_is_compressed(); + }) + } + + /// Write `data` as a chunk manifest on a fresh EC volume 2. Chunks stay on volume 1. + fn put_ec_manifest( + state: &VolumeServerState, + tmp: &tempfile::TempDir, + id: u64, + data: &[u8], + compressed: bool, + ) -> String { + 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}; + + let dir = tmp.path().to_str().unwrap(); + 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.to_vec(), + data_size: data.len() as u32, + ..Needle::default() + }; + n.set_is_chunk_manifest(); + if compressed { + n.set_is_compressed(); + } + 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(); + format!("/2,{:x}{:08x}", id, TEST_COOKIE) + } + + fn slash_name(stored: &str, name: &str) -> String { + let fid = stored.trim_start_matches('/'); + let (vid, needle) = fid.split_once(',').unwrap(); + format!("/{vid}/{needle}/{name}") + } + async fn send_read( state: &Arc, method: Method, @@ -5436,8 +5758,8 @@ mod tests { } } - /// A chunk manifest applies Range to the assembled object, as Go's - /// writeResponseContent does; its HEAD is still read from the meta. + /// A chunk manifest applies Range to the assembled object. HEAD reports the + /// sum of chunk sizes, including when a Range header is present. #[tokio::test] async fn test_chunk_manifest_range() { let tmp = tempfile::TempDir::new().unwrap(); @@ -5487,15 +5809,18 @@ mod tests { assert_eq!(status, StatusCode::RANGE_NOT_SATISFIABLE); assert_eq!(body, b"invalid range"); - let (status, headers, _) = send_read(&state, Method::HEAD, &path, None).await; + let (status, headers, body) = send_read(&state, Method::HEAD, &path, None).await; assert_eq!(status, StatusCode::OK); - let (range_status, range_headers, _) = + assert!(body.is_empty()); + assert_eq!(headers[header::CONTENT_LENGTH], "16"); + assert_eq!(headers["X-File-Store"], "chunked"); + assert_eq!(headers[header::ACCEPT_RANGES], "bytes"); + let (range_status, range_headers, range_body) = 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] - ); + assert_eq!(range_status, StatusCode::OK); + assert!(range_body.is_empty()); + assert_eq!(range_headers[header::CONTENT_LENGTH], "16"); + assert_eq!(range_headers["X-File-Store"], "chunked"); } /// A ranged GET of a chunk manifest fetches only the chunks its ranges @@ -5598,9 +5923,9 @@ mod tests { 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 { + assert_eq!(full, b"AAAABBBB\0\0CCCC"); + for start in 0..14 { + for end in start..14 { let range = format!("bytes={start}-{end}"); let (status, _, body) = send_read(&state, Method::GET, &path, Some(range.as_bytes())).await; @@ -5610,6 +5935,801 @@ mod tests { } } + /// Regular and EC HEAD of a chunk manifest, plain or gzipped, report the + /// sum and do not read chunks. The sum is not the stored needle length. + #[tokio::test] + async fn test_chunk_manifest_head_sum() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let a = put_test_needle(&state, 0x6e7a_0a01, b"AAAAAAAA"); + let b = put_test_needle(&state, 0x6e7a_0a02, b"BBBBBBBB"); + let json = format!( + r#"{{"name":"obj.txt","size":100,"chunks":[{{"fid":"{}","offset":0,"size":8}},{{"fid":"{}","offset":8,"size":8}}]}}"#, + &a[1..], + &b[1..] + ); + assert_ne!(json.len(), 16); + let path = put_manifest(&state, 0x6e7a_0a03, &json); + let (_regs, reads) = watch_reads(&[0x6e7a_0a01, 0x6e7a_0a02]); + + 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], "16"); + assert_eq!(headers["X-File-Store"], "chunked"); + assert_eq!(headers[header::ACCEPT_RANGES], "bytes"); + assert_eq!(taken(&reads), vec![0, 0]); + + let gz_path = put_gzip_manifest(&state, 0x6e7a_0a04, &json); + let (status, headers, _) = send_read(&state, Method::HEAD, &gz_path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(headers[header::CONTENT_LENGTH], "16"); + assert_eq!(headers["X-File-Store"], "chunked"); + assert_eq!(taken(&reads), vec![0, 0]); + + let cached = format!("{path}?response-cache-control=max-age%3D9"); + let (status, headers, _) = send_read(&state, Method::HEAD, &cached, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(headers[header::CACHE_CONTROL], "max-age=9"); + assert_eq!(headers[header::CONTENT_LENGTH], "16"); + assert_eq!(taken(&reads), vec![0, 0]); + + // Regression pins: expansion does not claim these. + let raw = put_manifest(&state, 0x6e7a_0a05, "not-json"); + let (status, headers, body) = send_read(&state, Method::HEAD, &raw, None).await; + assert_eq!(status, StatusCode::OK); + assert!(body.is_empty()); + assert_eq!(headers[header::CONTENT_LENGTH], "8"); + assert!(!headers.contains_key("X-File-Store")); + + let bypass = format!("{path}?cm=false"); + let (status, headers, _) = send_read(&state, Method::HEAD, &bypass, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(headers[header::CONTENT_LENGTH], json.len().to_string()); + assert!(!headers.contains_key("X-File-Store")); + } + + /// An EC manifest used to fetch every chunk on HEAD. It now reports the sum. + #[tokio::test] + async fn test_chunk_manifest_head_ec() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let a = put_test_needle(&state, 0x6e7a_0a11, b"AAAAAAAA"); + let b = put_test_needle(&state, 0x6e7a_0a12, b"BBBBBBBB"); + let json = format!( + r#"{{"name":"obj.txt","size":100,"chunks":[{{"fid":"{}","offset":0,"size":8}},{{"fid":"{}","offset":8,"size":8}}]}}"#, + &a[1..], + &b[1..] + ); + let path = put_ec_manifest(&state, &tmp, 0x6e7a_0a13, json.as_bytes(), false); + let (_regs, reads) = watch_reads(&[0x6e7a_0a11, 0x6e7a_0a12]); + 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], "16"); + assert_eq!(headers["X-File-Store"], "chunked"); + assert_eq!(taken(&reads), vec![0, 0]); + drop(_regs); + + // A second mount of volume 2 on the same store fails. Gzip gets its own. + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let a = put_test_needle(&state, 0x6e7a_0a15, b"AAAAAAAA"); + let b = put_test_needle(&state, 0x6e7a_0a16, b"BBBBBBBB"); + let json = format!( + r#"{{"name":"obj.txt","size":100,"chunks":[{{"fid":"{}","offset":0,"size":8}},{{"fid":"{}","offset":8,"size":8}}]}}"#, + &a[1..], + &b[1..] + ); + let gz = try_gzip_data(json.as_bytes()).unwrap(); + assert_ne!(gz.len(), 16); + let path = put_ec_manifest(&state, &tmp, 0x6e7a_0a17, &gz, true); + let (_regs, reads) = watch_reads(&[0x6e7a_0a15, 0x6e7a_0a16]); + 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], "16"); + assert_eq!(headers["X-File-Store"], "chunked"); + assert_eq!(taken(&reads), vec![0, 0]); + } + + fn solid_png(width: u32, height: u32) -> Vec { + let mut img = image::RgbImage::new(width, height); + for px in img.pixels_mut() { + *px = image::Rgb([10, 20, 30]); + } + let mut buf = Vec::new(); + img.write_to(&mut std::io::Cursor::new(&mut buf), image::ImageFormat::Png) + .unwrap(); + buf + } + + /// A zero crop is not a transform. A real crop of one PNG still fetches, + /// and Range applies to the transformed bytes. + #[tokio::test] + async fn test_chunk_manifest_head_image() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let first = put_test_needle(&state, 0x6e7a_0a21, b"AAAA"); + let second = put_test_needle(&state, 0x6e7a_0a22, b"BBBB"); + let raw = format!( + r#"{{"name":"obj.txt","size":8,"chunks":[{{"fid":"{}","offset":0,"size":4}},{{"fid":"{}","offset":4,"size":4}}]}}"#, + &first[1..], + &second[1..] + ); + let stored = put_manifest(&state, 0x6e7a_0a23, &raw); + let (_regs, reads) = watch_reads(&[0x6e7a_0a21, 0x6e7a_0a22]); + for query in ["crop_x2=0&crop_y2=0", "width=0&height=0"] { + let url = format!("{}?{query}", slash_name(&stored, "obj.jpg")); + let (status, headers, _) = send_read(&state, Method::HEAD, &url, None).await; + assert_eq!(status, StatusCode::OK, "{query}"); + assert_eq!(headers[header::CONTENT_LENGTH], "8", "{query}"); + assert_eq!(headers["X-File-Store"], "chunked", "{query}"); + assert_eq!(taken(&reads), vec![0, 0], "{query}"); + } + let partial = format!("{}?crop_x2=0&crop_y2=0", slash_name(&stored, "obj.jpg")); + let (status, _, body) = send_read(&state, Method::GET, &partial, Some(b"bytes=0-3")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + assert_eq!(body, b"AAAA"); + assert_eq!(taken(&reads), vec![1, 0]); + + let png = solid_png(8, 8); + let chunk = put_test_needle(&state, 0x6e7a_0a24, &png); + let image = format!( + r#"{{"name":"obj.txt","size":{},"chunks":[{{"fid":"{}","offset":0,"size":{}}}]}}"#, + png.len(), + &chunk[1..], + png.len() + ); + let stored = put_manifest(&state, 0x6e7a_0a25, &image); + let (_regs, reads) = watch_reads(&[0x6e7a_0a24]); + let crop = format!( + "{}?crop_x1=0&crop_y1=0&crop_x2=2&crop_y2=2", + slash_name(&stored, "obj.jpg") + ); + let (get_status, _, cropped) = send_read(&state, Method::GET, &crop, None).await; + assert_eq!(get_status, StatusCode::OK); + assert_ne!(cropped.len(), png.len(), "the crop must change the length"); + assert_eq!(taken(&reads), vec![1]); + + let (status, headers, body) = + send_read(&state, Method::HEAD, &crop, Some(b"bytes=0-1")).await; + assert_eq!(status, StatusCode::OK); + assert!(body.is_empty()); + assert_eq!(headers[header::CONTENT_LENGTH], cropped.len().to_string()); + assert_eq!(taken(&reads), vec![1]); + + let past = format!("bytes={}-", cropped.len()); + let (status, _, _) = send_read(&state, Method::GET, &crop, Some(past.as_bytes())).await; + assert_eq!(status, StatusCode::RANGE_NOT_SATISFIABLE); + assert_eq!(taken(&reads), vec![1]); + } + + /// Go takes the transform extension from the URL first — including a fid + /// suffix like /vid,fid.png — before the resolved filename's. + #[tokio::test] + async fn test_chunk_manifest_fid_extension_drives_transform() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let png = solid_png(8, 8); + let chunk = put_test_needle(&state, 0x6e7a_0f01, &png); + let (_regs, reads) = watch_reads(&[0x6e7a_0f01]); + let json = format!( + r#"{{"name":"obj","size":{},"chunks":[{{"fid":"{}","offset":0,"size":{}}}]}}"#, + png.len(), + &chunk[1..], + png.len() + ); + let stored = put_manifest(&state, 0x6e7a_0f02, &json); + let url = format!("{stored}.png?width=2&height=2"); + let (status, _, body) = send_read(&state, Method::GET, &url, None).await; + assert_eq!(status, StatusCode::OK); + assert_ne!(body.len(), png.len(), "resize must run on the fid suffix"); + assert_eq!(taken(&reads), vec![1]); + + let (status, headers, head_body) = send_read(&state, Method::HEAD, &url, None).await; + assert_eq!(status, StatusCode::OK); + assert!(head_body.is_empty()); + assert_eq!(headers[header::CONTENT_LENGTH], body.len().to_string()); + assert_eq!(taken(&reads), vec![1]); + } + + #[test] + fn test_expansion_limit_allows_exactly_two_gib() { + assert!(!exceeds_expansion_limit(MAX_EXPANSION_BYTES)); + assert!(exceeds_expansion_limit(MAX_EXPANSION_BYTES + 1)); + assert!(!exceeds_expansion_limit(0)); + } + + #[test] + fn test_remote_whole_window_caps_at_visible_window() { + let ranges = [HttpRange { + start: 100, + length: 4, + }]; + let how = remote_read_for(100, 4, Some(&ranges)).unwrap(); + assert!(matches!(how, RemoteRead::Whole { cap: 4 })); + let full = remote_read_for(0, 8, None).unwrap(); + assert!(matches!(full, RemoteRead::Whole { cap: 8 })); + let prefix = remote_read_for( + 0, + 8, + Some(&[HttpRange { + start: 0, + length: 4, + }]), + ) + .unwrap(); + assert!(matches!(prefix, RemoteRead::Range { lo: 0, hi: 4 })); + } + + #[test] + fn test_sum_ranges_size_ignores_non_positive() { + let ranges = [ + HttpRange { + start: 0, + length: 8, + }, + HttpRange { + start: 0, + length: 8, + }, + HttpRange { + start: 0, + length: -1_073_741_824, + }, + ]; + assert_eq!(sum_ranges_size(&ranges), 16); + let over = [ + HttpRange { + start: 0, + length: MAX_EXPANSION_BYTES as i64 + 1, + }, + HttpRange { + start: 0, + length: -1, + }, + ]; + assert_eq!(sum_ranges_size(&over), MAX_EXPANSION_BYTES as i64 + 1); + assert!(exceeds_expansion_limit(sum_ranges_size(&over) as u64)); + } + + /// `"size"` is not the object length. Holes stay zero. Chunks at or past + /// the sum are not read. + #[tokio::test] + async fn test_chunk_manifest_length_is_sum() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + + let a = put_test_needle(&state, 0x6e7a_0c01, b"AAAA"); + let b = put_test_needle(&state, 0x6e7a_0c02, b"BBBB"); + let omitted = format!( + r#"{{"chunks":[{{"fid":"{}","offset":0,"size":4}},{{"fid":"{}","offset":4,"size":4}}]}}"#, + &a[1..], + &b[1..] + ); + let path = put_manifest(&state, 0x6e7a_0c03, &omitted); + let (status, headers, body) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, b"AAAABBBB"); + assert_eq!(headers[header::CONTENT_LENGTH], "8"); + + let padded = format!( + r#"{{"size":100,"chunks":[{{"fid":"{}","offset":0,"size":4}}]}}"#, + &a[1..] + ); + let path = put_manifest(&state, 0x6e7a_0c04, &padded); + let (status, headers, body) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, b"AAAA"); + assert_eq!(headers[header::CONTENT_LENGTH], "4"); + + let gap_b = put_test_needle(&state, 0x6e7a_0c05, b"BBBB"); + let (_regs, reads) = watch_reads(&[0x6e7a_0c01, 0x6e7a_0c05]); + let gap = format!( + r#"{{"size":100,"chunks":[{{"fid":"{}","offset":0,"size":4}},{{"fid":"{}","offset":8,"size":4}}]}}"#, + &a[1..], + &gap_b[1..] + ); + let path = put_manifest(&state, 0x6e7a_0c06, &gap); + let (status, _, body) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, b"AAAA\0\0\0\0"); + assert_eq!(taken(&reads), vec![1, 0]); + + let zero = format!( + r#"{{"chunks":[{{"fid":"9,0000000a12345678","offset":0,"size":0}},{{"fid":"{}","offset":0,"size":4}}]}}"#, + &a[1..] + ); + let path = put_manifest(&state, 0x6e7a_0c07, &zero); + let (status, _, body) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, b"AAAA"); + + let empty = put_manifest(&state, 0x6e7a_0c08, r#"{"size":5,"chunks":[]}"#); + let (status, headers, body) = send_read(&state, Method::HEAD, &empty, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(headers[header::CONTENT_LENGTH], "0"); + assert!(body.is_empty()); + let (status, headers, body) = send_read(&state, Method::GET, &empty, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(headers[header::CONTENT_LENGTH], "0"); + assert!(body.is_empty()); + let (status, headers, _) = send_read(&state, Method::GET, &empty, Some(b"bytes=0-0")).await; + assert_eq!(status, StatusCode::RANGE_NOT_SATISFIABLE); + assert_eq!(headers["Content-Range"], "bytes */0"); + + let ten = put_test_needle(&state, 0x6e7a_0c09, b"0123456789"); + let (_regs, reads) = watch_reads(&[0x6e7a_0c09]); + let dup = format!( + r#"{{"chunks":[{{"fid":"{}","offset":0,"size":10}}]}}"#, + &ten[1..] + ); + let path = put_manifest(&state, 0x6e7a_0c0a, &dup); + let (status, headers, body) = + send_read(&state, Method::GET, &path, Some(b"bytes=0-7,0-7")).await; + assert_eq!(status, StatusCode::OK); + assert!(body.is_empty()); + assert_eq!(headers[header::ACCEPT_RANGES], "bytes"); + assert_eq!(taken(&reads), vec![0]); + } + + /// A negative suffix must not shrink the covered length under the sum. + /// Duplicate ranges stay an empty 200 and read nothing. + #[tokio::test] + async fn test_chunk_manifest_negative_suffix_stays_empty() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let ten = put_test_needle(&state, 0x6e7a_0c41, b"0123456789"); + let (_regs, reads) = watch_reads(&[0x6e7a_0c41]); + let json = format!( + r#"{{"chunks":[{{"fid":"{}","offset":0,"size":10}}]}}"#, + &ten[1..] + ); + let path = put_manifest(&state, 0x6e7a_0c42, &json); + let (status, headers, body) = send_read( + &state, + Method::GET, + &path, + Some(b"bytes=0-7,0-7,--1073741824"), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert!(body.is_empty()); + assert_eq!(headers[header::ACCEPT_RANGES], "bytes"); + assert_eq!(taken(&reads), vec![0]); + } + + /// A negative field is rejected before any read, including a chunk the + /// range does not overlap. + #[tokio::test] + async fn test_chunk_manifest_negative_reads_nothing() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let inside = put_test_needle(&state, 0x6e7a_0c11, b"AAAA"); + let outside = put_test_needle(&state, 0x6e7a_0c12, b"BBBB"); + let (_regs, reads) = watch_reads(&[0x6e7a_0c11, 0x6e7a_0c12]); + let json = format!( + r#"{{"chunks":[{{"fid":"{}","offset":0,"size":4}},{{"fid":"{}","offset":100,"size":4}},{{"fid":"9,0000000b12345678","offset":0,"size":-1}}]}}"#, + &inside[1..], + &outside[1..] + ); + let path = put_manifest(&state, 0x6e7a_0c13, &json); + let (status, _, body) = send_read(&state, Method::GET, &path, Some(b"bytes=0-1")).await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); + let body = String::from_utf8(body).unwrap(); + assert!( + body.contains("invalid negative chunk offset/size"), + "{body}" + ); + assert_eq!(taken(&reads), vec![0, 0]); + + let (status, _, _) = send_read(&state, Method::HEAD, &path, None).await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); + assert_eq!(taken(&reads), vec![0, 0]); + } + + /// The cap is the bytes that would be buffered. Exactly 2 GiB is allowed + /// by the predicate. This test does not allocate that buffer. + #[tokio::test] + async fn test_chunk_manifest_cap_is_buffered_bytes() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let one = put_test_needle(&state, 0x6e7a_0c21, b"Z"); + let (_regs, reads) = watch_reads(&[0x6e7a_0c21]); + let big = MAX_EXPANSION_BYTES + 1; + let json = format!( + r#"{{"size":1,"chunks":[{{"fid":"{}","offset":0,"size":{big}}}]}}"#, + &one[1..] + ); + let path = put_manifest(&state, 0x6e7a_0c22, &json); + + let (status, headers, body) = send_read(&state, Method::HEAD, &path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(headers[header::CONTENT_LENGTH], big.to_string()); + assert!(body.is_empty()); + assert_eq!(taken(&reads), vec![0]); + + let (status, _, body) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::PAYLOAD_TOO_LARGE); + assert_eq!(body, b"chunk manifest exceeds expansion limit"); + assert_eq!(taken(&reads), vec![0]); + + let (status, _, body) = send_read(&state, Method::GET, &path, Some(b"bytes=0-3")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + assert_eq!(body, b"Z\0\0\0"); + assert_eq!(taken(&reads), vec![1]); + + let (status, _, body) = send_read(&state, Method::GET, &path, Some(b"bytes=0-")).await; + assert_eq!(status, StatusCode::PAYLOAD_TOO_LARGE); + assert_eq!(body, b"chunk manifest exceeds expansion limit"); + assert_eq!(taken(&reads), vec![0]); + + let overflow = format!( + r#"{{"chunks":[{{"fid":"{}","offset":0,"size":1}},{{"fid":"9,0000000c12345678","offset":0,"size":9223372036854775807}},{{"fid":"9,0000000d12345678","offset":0,"size":9223372036854775807}},{{"fid":"9,0000000e12345678","offset":0,"size":9223372036854775807}}]}}"#, + &one[1..] + ); + let path = put_manifest(&state, 0x6e7a_0c23, &overflow); + let (status, _, body) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::PAYLOAD_TOO_LARGE); + assert_eq!(body, b"chunk manifest exceeds expansion limit"); + assert_eq!(taken(&reads), vec![0]); + } + + /// Two open ranges on an `i64::MAX` object must not wrap the covered + /// length. The saturated sum equals the total, which is over the buffer + /// cap, so this is 413 and reads nothing. + #[tokio::test] + async fn test_chunk_manifest_open_ranges_do_not_wrap() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let one = put_test_needle(&state, 0x6e7a_0c31, b"Z"); + let (_regs, reads) = watch_reads(&[0x6e7a_0c31]); + let json = format!( + r#"{{"chunks":[{{"fid":"{}","offset":0,"size":{}}}]}}"#, + &one[1..], + i64::MAX + ); + let path = put_manifest(&state, 0x6e7a_0c32, &json); + let (status, _, body) = send_read(&state, Method::GET, &path, Some(b"bytes=0-,0-")).await; + assert_eq!(status, StatusCode::PAYLOAD_TOO_LARGE); + assert_eq!(body, b"chunk manifest exceeds expansion limit"); + assert_eq!(taken(&reads), vec![0]); + } + + #[derive(Clone)] + enum PeerMode { + /// Honor a single `bytes=lo-hi` range. No Range header returns `chunk`. + Honor { + chunk: Vec, + }, + Always200 { + body: Vec, + }, + Status(u16), + /// 206 whose Content-Range starts at 0, whatever Range was asked. + WrongStart { + body: Vec, + }, + /// 206, correct start, body shorter than the asked span. + Short { + body: Vec, + }, + /// 206, correct start, body longer than the asked span. + Long { + body: Vec, + }, + /// 206 with no Content-Range header. + NoContentRange { + body: Vec, + }, + } + + struct PeerHit { + range: Option, + } + + async fn spawn_chunk_peer( + mode: std::sync::Arc>, + hits: std::sync::Arc>>, + ) -> (std::net::SocketAddr, tokio::task::JoinHandle<()>) { + use axum::routing::get; + let app = axum::Router::new().route( + "/{*rest}", + get(move |headers: HeaderMap| { + let mode = mode.clone(); + let hits = hits.clone(); + async move { + let range = headers + .get(header::RANGE) + .and_then(|v| v.to_str().ok()) + .map(|s| s.to_string()); + hits.lock().unwrap().push(PeerHit { + range: range.clone(), + }); + let mode = mode.lock().unwrap().clone(); + match mode { + PeerMode::Honor { chunk } => match range { + None => (StatusCode::OK, chunk).into_response(), + Some(spec) => { + let nums = spec.trim_start_matches("bytes="); + let (lo, hi) = nums.split_once('-').unwrap(); + let lo: usize = lo.parse().unwrap(); + let hi: usize = hi.parse().unwrap(); + let body = chunk[lo..=hi].to_vec(); + let cr = format!("bytes {lo}-{hi}/{}", chunk.len()); + ( + StatusCode::PARTIAL_CONTENT, + [(header::CONTENT_RANGE, cr)], + body, + ) + .into_response() + } + }, + PeerMode::Always200 { body } => (StatusCode::OK, body).into_response(), + PeerMode::Status(code) => { + StatusCode::from_u16(code).unwrap().into_response() + } + PeerMode::WrongStart { body } => ( + StatusCode::PARTIAL_CONTENT, + [( + header::CONTENT_RANGE, + format!("bytes 0-{}/8", body.len().saturating_sub(1)), + )], + body, + ) + .into_response(), + PeerMode::Short { body } => { + let spec = range.unwrap(); + let nums = spec.trim_start_matches("bytes="); + let (lo, hi) = nums.split_once('-').unwrap(); + ( + StatusCode::PARTIAL_CONTENT, + [(header::CONTENT_RANGE, format!("bytes {lo}-{hi}/8"))], + body, + ) + .into_response() + } + PeerMode::Long { body } => { + let spec = range.unwrap(); + let nums = spec.trim_start_matches("bytes="); + let (lo, hi) = nums.split_once('-').unwrap(); + ( + StatusCode::PARTIAL_CONTENT, + [(header::CONTENT_RANGE, format!("bytes {lo}-{hi}/8"))], + body, + ) + .into_response() + } + PeerMode::NoContentRange { body } => { + (StatusCode::PARTIAL_CONTENT, body).into_response() + } + } + } + }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let handle = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + (addr, handle) + } + + async fn spawn_lookup(peer: std::net::SocketAddr) -> String { + use axum::routing::get; + let app = axum::Router::new().route( + "/dir/lookup", + get(move || async move { + axum::Json(serde_json::json!({ + "volumeOrFileId": "9", + "locations": [{ "url": peer.to_string() }] + })) + }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + format!("127.0.0.1:{}", addr.port()) + } + + fn remote_manifest_state( + tmp: &tempfile::TempDir, + master_url: String, + ) -> Arc { + let mut state = volume_test_state(tmp); + let s = Arc::get_mut(&mut state).unwrap(); + // Empty self_url makes `url.contains(self_url)` true for every peer. + s.self_url = "198.51.100.8:8080".to_string(); + s.master_url = master_url; + state + } + + /// Remote overlap is one bounded Range, indexed from `lo`. Local short + /// needles still zero-fill. + #[tokio::test] + async fn test_chunk_manifest_remote_range() { + let tmp = tempfile::TempDir::new().unwrap(); + let mode = std::sync::Arc::new(std::sync::Mutex::new(PeerMode::Honor { + chunk: b"ABCDEFGH".to_vec(), + })); + let hits = std::sync::Arc::new(std::sync::Mutex::new(Vec::new())); + let (peer, peer_handle) = spawn_chunk_peer(mode.clone(), hits.clone()).await; + let master = spawn_lookup(peer).await; + let state = remote_manifest_state(&tmp, master); + let local = put_test_needle(&state, 0x6e7a_0d01, b"0123456789"); + let remote_fid = "9,0000000a12345678"; + let json = format!( + r#"{{"chunks":[{{"fid":"{}","offset":0,"size":10}},{{"fid":"{remote_fid}","offset":10,"size":8}}]}}"#, + &local[1..] + ); + let path = put_manifest(&state, 0x6e7a_0d02, &json); + + hits.lock().unwrap().clear(); + let (status, headers, body) = + send_read(&state, Method::GET, &path, Some(b"bytes=12-15")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + assert_eq!(body, b"CDEF"); + assert_eq!(headers["Content-Range"], "bytes 12-15/18"); + let seen = hits.lock().unwrap().drain(..).collect::>(); + assert_eq!(seen.len(), 1); + assert_eq!(seen[0].range.as_deref(), Some("bytes=2-5")); + + hits.lock().unwrap().clear(); + let (status, _, body) = + send_read(&state, Method::GET, &path, Some(b"bytes=11-12,15-16")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + let text = String::from_utf8(body).unwrap(); + assert!( + text.contains("Content-Range: bytes 11-12/18\r\n\r\nBC"), + "{text}" + ); + assert!( + text.contains("Content-Range: bytes 15-16/18\r\n\r\nFG"), + "{text}" + ); + let seen = hits.lock().unwrap().drain(..).collect::>(); + assert_eq!(seen.len(), 1, "disjoint ranges share one span"); + assert_eq!(seen[0].range.as_deref(), Some("bytes=1-6")); + + hits.lock().unwrap().clear(); + let whole = put_manifest( + &state, + 0x6e7a_0d03, + &format!(r#"{{"chunks":[{{"fid":"{remote_fid}","offset":0,"size":8}}]}}"#), + ); + let (status, _, body) = send_read(&state, Method::GET, &whole, Some(b"bytes=0-7")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + assert_eq!(body, b"ABCDEFGH"); + let seen = hits.lock().unwrap().drain(..).collect::>(); + assert_eq!(seen.len(), 1); + assert_eq!(seen[0].range, None); + + *mode.lock().unwrap() = PeerMode::Short { + body: b"AB".to_vec(), + }; + hits.lock().unwrap().clear(); + let (status, _, body) = send_read(&state, Method::GET, &whole, Some(b"bytes=0-3")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + assert_eq!(body, b"AB\0\0"); + + *mode.lock().unwrap() = PeerMode::Long { + body: b"ABCDEFGH".to_vec(), + }; + let (status, _, body) = send_read(&state, Method::GET, &whole, Some(b"bytes=0-3")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + assert_eq!(body, b"ABCD"); + + *mode.lock().unwrap() = PeerMode::WrongStart { + body: b"ABCD".to_vec(), + }; + let (status, _, body) = send_read(&state, Method::GET, &whole, Some(b"bytes=2-5")).await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); + assert!(String::from_utf8_lossy(&body).contains("read chunk")); + + *mode.lock().unwrap() = PeerMode::NoContentRange { + body: b"CDEF".to_vec(), + }; + let (status, _, body) = send_read(&state, Method::GET, &whole, Some(b"bytes=2-5")).await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); + assert!(String::from_utf8_lossy(&body).contains("read chunk")); + + *mode.lock().unwrap() = PeerMode::Always200 { + body: b"ABCDEFGH".to_vec(), + }; + let (status, _, body) = send_read(&state, Method::GET, &whole, Some(b"bytes=2-5")).await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); + assert!(String::from_utf8_lossy(&body).contains("read chunk")); + + *mode.lock().unwrap() = PeerMode::Status(416); + let (status, _, body) = send_read(&state, Method::GET, &whole, Some(b"bytes=1-2")).await; + assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); + assert!(String::from_utf8_lossy(&body).contains("read chunk")); + + peer_handle.abort(); + } + + /// A local needle shorter than declared zero-fills. The remote 416 above + /// is the intentional divergence. + #[tokio::test] + async fn test_chunk_manifest_local_short_needle_zero_fills() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let short = put_test_needle(&state, 0x6e7a_0d11, b"AB"); + let json = format!( + r#"{{"chunks":[{{"fid":"{}","offset":0,"size":4}}]}}"#, + &short[1..] + ); + let path = put_manifest(&state, 0x6e7a_0d12, &json); + let (status, _, body) = send_read(&state, Method::GET, &path, Some(b"bytes=1-2")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + assert_eq!(body, b"B\0"); + let (status, _, body) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, b"AB\0\0"); + } + + /// A short range of a huge remote chunk is fetched. The span between two + /// distant ranges is not, once it would buffer more than 2 GiB. + #[tokio::test] + async fn test_chunk_manifest_remote_span_over_cap_is_rejected() { + let tmp = tempfile::TempDir::new().unwrap(); + let mode = std::sync::Arc::new(std::sync::Mutex::new(PeerMode::Honor { + chunk: b"ABCDEFGH".to_vec(), + })); + let hits = std::sync::Arc::new(std::sync::Mutex::new(Vec::new())); + let (peer, peer_handle) = spawn_chunk_peer(mode, hits.clone()).await; + let master = spawn_lookup(peer).await; + let state = remote_manifest_state(&tmp, master); + let big = MAX_EXPANSION_BYTES + 1; + let path = put_manifest( + &state, + 0x6e7a_0e01, + &format!(r#"{{"chunks":[{{"fid":"9,0000000a12345678","offset":0,"size":{big}}}]}}"#), + ); + + hits.lock().unwrap().clear(); + let (status, _, body) = send_read(&state, Method::GET, &path, Some(b"bytes=0-3")).await; + assert_eq!(status, StatusCode::PARTIAL_CONTENT); + assert_eq!(body, b"ABCD"); + assert_eq!(hits.lock().unwrap().len(), 1); + + hits.lock().unwrap().clear(); + let last = big - 1; + let wide = format!("bytes=0-0,{last}-{last}"); + let (status, _, body) = send_read(&state, Method::GET, &path, Some(wide.as_bytes())).await; + assert_eq!(status, StatusCode::PAYLOAD_TOO_LARGE); + assert_eq!(body, b"chunk manifest exceeds expansion limit"); + assert!(hits.lock().unwrap().is_empty()); + peer_handle.abort(); + } + + /// An illegal manifest MIME is not a legal header value. HEAD reaches the + /// header builder, so both HEAD and GET fall back instead of panicking. + #[tokio::test] + async fn test_chunk_manifest_illegal_mime_falls_back() { + let tmp = tempfile::TempDir::new().unwrap(); + let state = volume_test_state(&tmp); + let chunk = put_test_needle(&state, 0x6e7a_0e11, b"AAAA"); + let (_regs, reads) = watch_reads(&[0x6e7a_0e11]); + let json = format!( + r#"{{"name":"obj.txt","mime":"text/plain\n","size":4,"chunks":[{{"fid":"{}","offset":0,"size":4}}]}}"#, + &chunk[1..] + ); + let path = put_manifest(&state, 0x6e7a_0e12, &json); + + 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_TYPE], "application/octet-stream"); + assert_eq!(headers[header::CONTENT_LENGTH], "4"); + assert_eq!(taken(&reads), vec![0]); + + let (status, headers, body) = send_read(&state, Method::GET, &path, None).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body, b"AAAA"); + assert_eq!(headers[header::CONTENT_TYPE], "application/octet-stream"); + assert_eq!(taken(&reads), vec![1]); + } + /// 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]