From 8c1ebbee32718c3173b301828dbfc93e24893e2e Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Sat, 3 Oct 2026 12:54:46 +0300 Subject: [PATCH] volume server: split get_or_head_handler_inner into phases (#11489) * 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) * 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) * 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. 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) * 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: " error the old check returned. Co-Authored-By: Claude Opus 5.5 (1M context) * volume server: mirror Go order in the buffered read path - check HEAD before Range in buffered_response (writeResponseContent order); an EC-volume HEAD with a Range header answered 206, Go answers 200 - treat the proxied flag as an exact query pair like Go's parsed lookup, not a substring - name the phases after their Go counterparts: check_download_limit and read_ec_shard_needle; reuse has_replication() - drop comments that restate the code or cite Go line numbers --------- Co-authored-by: Claude Opus 5.5 (1M context) Co-authored-by: Chris Lu Co-authored-by: Chris Lu --- seaweed-volume/src/server/handlers.rs | 959 +++++++++++++++----------- 1 file changed, 574 insertions(+), 385 deletions(-) diff --git a/seaweed-volume/src/server/handlers.rs b/seaweed-volume/src/server/handlers.rs index 4a8a223b9..664204c29 100644 --- a/seaweed-volume/src/server/handlers.rs +++ b/seaweed-volume/src/server/handlers.rs @@ -6,6 +6,7 @@ use std::collections::HashMap; use std::future::Future; +use std::ops::ControlFlow; use std::sync::Arc; use std::sync::atomic::Ordering; @@ -1050,27 +1051,10 @@ async fn get_or_head_handler_inner( request: Request, ) -> Response { let path = request.uri().path().to_string(); - let raw_query = request.uri().query().map(|q| q.to_string()); let method = request.method().clone(); - // JWT check for reads — must happen BEFORE path parsing to match Go behavior. - // Go's GetOrHeadHandler calls maybeCheckJwtAuthorization before NewVolumeId, - // so invalid paths with JWT enabled return 401, not 400. - let file_id = extract_file_id(&path); - let token = extract_jwt(&headers, request.uri()); - if state - .guard - .read() - .unwrap() - .check_jwt_for_file(token.as_deref(), &file_id, false) - .is_err() - { - let body = serde_json::json!({"error": "wrong jwt"}); - return Response::builder() - .status(StatusCode::UNAUTHORIZED) - .header(header::CONTENT_TYPE, "application/json") - .body(Body::from(serde_json::to_string(&body).unwrap())) - .unwrap(); + if let Some(resp) = reject_read_jwt(&state, &headers, request.uri(), &path) { + return resp; } let (vid, needle_id, cookie) = match parse_url_path(&path) { @@ -1078,284 +1062,46 @@ async fn get_or_head_handler_inner( None => return StatusCode::BAD_REQUEST.into_response(), }; - // Check if volume exists locally; if not, proxy/redirect based on read_mode. - // This mirrors Go's hasVolume + hasEcVolume check in GetOrHeadHandler. - // NOTE: The RwLockReadGuard must be dropped before any .await to keep the future Send. + // The RwLockReadGuard must drop before any .await to keep the future Send. let has_volume = state.store.read().unwrap().has_volume(vid); let has_ec_volume = state.store.read().unwrap().has_ec_volume(vid); if !has_volume && !has_ec_volume { - // Check if already proxied (loop prevention) - let query_string = request.uri().query().unwrap_or("").to_string(); - let is_proxied = query_string.contains("proxied=true"); - - if is_proxied || state.read_mode == ReadMode::Local || state.master_url.is_empty() { - return StatusCode::NOT_FOUND.into_response(); - } - - // For redirect, fid must be stripped of extension (Go parity: parseURLPath returns raw fid). - let info = match build_proxy_request_info(&path, request.headers(), &query_string) { - Some(info) => info, - None => return StatusCode::NOT_FOUND.into_response(), - }; - - return proxy_or_redirect_to_target(&state, info, vid, false).await; + return proxy_missing_volume(&state, request.uri(), request.headers(), &path, vid).await; } - // Download throttling — matches Go's checkDownloadLimit + waitForDownloadSlot - let download_guard = if state.concurrent_download_limit > 0 { - let timeout = state.inflight_download_data_timeout; - let deadline = tokio::time::Instant::now() + timeout; - let query_string = request.uri().query().unwrap_or("").to_string(); - - let current = state.inflight_download_bytes.load(Ordering::Relaxed); - if current > state.concurrent_download_limit { - metrics::HANDLER_COUNTER - .with_label_values(&[metrics::DOWNLOAD_LIMIT_COND]) - .inc(); - - // Go tries proxy to replica ONCE before entering the blocking wait - // loop (checkDownloadLimit L65). It does NOT retry on each wakeup. - let should_try_replica = - !query_string.contains("proxied=true") && !state.master_url.is_empty() && { - let store = state.store.read().unwrap(); - store.find_volume(vid).is_some_and(|(_, vol)| { - vol.super_block.replica_placement.get_copy_count() > 1 - }) - }; - if should_try_replica - && let Some(info) = - build_proxy_request_info(&path, request.headers(), &query_string) - { - return proxy_or_redirect_to_target(&state, info, vid, true).await; - } - - // Blocking wait loop (Go's waitForDownloadSlot) - loop { - if tokio::time::timeout_at(deadline, state.download_notify.notified()) - .await - .is_err() - { - return json_error_with_query( - StatusCode::TOO_MANY_REQUESTS, - "download limit exceeded", - raw_query.as_deref(), - ); - } - let current = state.inflight_download_bytes.load(Ordering::Relaxed); - if current <= state.concurrent_download_limit { - break; - } - } - } - // We'll set the actual bytes after reading the needle (once we know the size) - Some(state.clone()) - } else { - None - }; - - // Read needle — branching between regular volume and EC volume paths. - // EC volumes always do a full read (no streaming/meta-only). - let n: Needle; + let track_download = + match check_download_limit(&state, request.uri(), request.headers(), &path, vid).await { + ControlFlow::Continue(track_download) => track_download, + ControlFlow::Break(resp) => return resp, + }; let read_deleted = query.read_deleted.as_deref() == Some("true"); - 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) - && (query.width.unwrap_or(0) > 0 || query.height.unwrap_or(0) > 0); - // Go's shouldCropImages (L410) requires x2 > x1 && y2 > y1 (x1/y1 default 0). - // Only disable streaming when a real crop will actually happen. - 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 (ext, request_kind) = parse_read_request(&path, &headers, &query, &method); + + // EC volumes always do a full read (no streaming/meta-only). + let plan = if has_ec_volume && !has_volume { + read_ec_shard_needle(&state, vid, needle_id, cookie).await + } else { + read_volume_needle(&state, vid, needle_id, cookie, read_deleted, request_kind).await }; - let has_image_ops = has_resize_ops || has_crop_ops; - - // Stream info is only available for regular volumes, not EC volumes. - let stream_info; - let bypass_cm; - let track_download; - let can_stream; - let can_handle_head_from_meta; - let can_handle_range_from_source; - - if has_ec_volume && !has_volume { - // ---- EC volume read path (always full read, no streaming) ---- - // - // The distributed read path already does a local-first pass - // in its Snapshot phase under the same store read lock the - // legacy code would have taken — so calling it directly - // serves both the "all shards local" fast case and the - // "some intervals need peer fetch + reconstruct" general - // case without paying for the local interval reads twice. - match crate::server::store_ec::read_ec_shard_needle_distributed(&state, vid, needle_id) - .await - { - Ok(Some(ec_needle)) => { - n = ec_needle; - } - Ok(None) => { - metrics::HANDLER_COUNTER - .with_label_values(&[metrics::ERROR_GET_NOT_FOUND]) - .inc(); - return StatusCode::NOT_FOUND.into_response(); - } - Err(e) => { - let kind = if e.kind() == std::io::ErrorKind::NotFound { - metrics::ERROR_GET_NOT_FOUND - } else { - metrics::ERROR_GET_INTERNAL - }; - metrics::HANDLER_COUNTER.with_label_values(&[kind]).inc(); - if e.kind() == std::io::ErrorKind::NotFound { - return StatusCode::NOT_FOUND.into_response(); - } - return (StatusCode::INTERNAL_SERVER_ERROR, format!("ec read: {}", e)) - .into_response(); - } - } - - // Validate cookie (matches Go behavior after ReadEcShardNeedle) - if n.cookie != cookie { - return StatusCode::NOT_FOUND.into_response(); - } - - // EC volumes: no streaming support - stream_info = None; - bypass_cm = query.cm.as_deref() == Some("false"); - track_download = download_guard.is_some(); - can_stream = false; - can_handle_head_from_meta = false; - can_handle_range_from_source = false; - } else { - // ---- Regular volume read path (with streaming support) ---- - bypass_cm = query.cm.as_deref() == Some("false"); - track_download = download_guard.is_some(); - let request_kind = SourceReadRequest { - is_head: method == Method::HEAD, - has_range, - has_image_ops, - bypass_cm, - }; - - let read_state = state.clone(); - let read = tokio::task::spawn_blocking(move || { - read_needle_for_get( - &read_state, - vid, - needle_id, - cookie, - read_deleted, - request_kind, - ) - }) - .await; - (n, stream_info) = match read { - Ok(Ok(Some(found))) => found, - // Cookie mismatch - Ok(Ok(None)) => return StatusCode::NOT_FOUND.into_response(), - Ok(Err( - crate::storage::volume::VolumeError::NotFound - | crate::storage::volume::VolumeError::Deleted, - )) => { - metrics::HANDLER_COUNTER - .with_label_values(&[metrics::ERROR_GET_NOT_FOUND]) - .inc(); - return StatusCode::NOT_FOUND.into_response(); - } - Ok(Err(e)) => { - metrics::HANDLER_COUNTER - .with_label_values(&[metrics::ERROR_GET_INTERNAL]) - .inc(); - return ( - StatusCode::INTERNAL_SERVER_ERROR, - format!("read error: {}", e), - ) - .into_response(); - } - Err(e) => { - metrics::HANDLER_COUNTER - .with_label_values(&[metrics::ERROR_GET_INTERNAL]) - .inc(); - return ( - StatusCode::INTERNAL_SERVER_ERROR, - format!("read error: {}", e), - ) - .into_response(); - } - }; - - // Stream info is only returned for a reply served from the data file. - let can_direct_source_read = stream_info.is_some() && request_kind.direct(&n); - - // Determine if we can stream (large, direct-source eligible, no range) - can_stream = can_direct_source_read - && n.data_size > STREAMING_THRESHOLD - && !has_range - && method != Method::HEAD; - - // Go uses meta-only reads for all HEAD requests, regardless of compression/chunked files. - can_handle_head_from_meta = stream_info.is_some() && method == Method::HEAD; - can_handle_range_from_source = can_direct_source_read && has_range; - } - - // Build ETag and Last-Modified BEFORE conditional checks and chunk manifest expansion - // (matches Go order: conditional checks first, then chunk manifest) - let etag = format!("\"{}\"", n.etag()); - - // Build Last-Modified header (RFC 1123 format) — must be done before conditional checks - let last_modified_str = if n.last_modified > 0 { - use chrono::{TimeZone, Utc}; - Utc.timestamp_opt(n.last_modified as i64, 0) - .single() - .map(|dt| dt.format("%a, %d %b %Y %H:%M:%S GMT").to_string()) - } else { - None + let ReadPlan { + needle: n, + strategy, + } = match plan { + ControlFlow::Continue(plan) => plan, + ControlFlow::Break(resp) => return resp, }; - // Check If-Modified-Since FIRST (Go checks this before If-None-Match) - if n.last_modified > 0 - && let Some(ims_header) = headers.get(header::IF_MODIFIED_SINCE) - && let Ok(ims_str) = ims_header.to_str() - { - // Parse HTTP date format: "Mon, 02 Jan 2006 15:04:05 GMT" - if let Ok(ims_time) = - chrono::NaiveDateTime::parse_from_str(ims_str, "%a, %d %b %Y %H:%M:%S GMT") - && (n.last_modified as i64) <= ims_time.and_utc().timestamp() - { - let mut resp = StatusCode::NOT_MODIFIED.into_response(); - if let Some(ref lm) = last_modified_str { - resp.headers_mut() - .insert(header::LAST_MODIFIED, lm.parse().unwrap()); - } - // Go sets ETag AFTER the 304 return paths (L235), so 304 does NOT include ETag - return resp; - } - } - - // Check If-None-Match SECOND - if let Some(if_none_match) = headers.get(header::IF_NONE_MATCH) - && let Ok(inm) = if_none_match.to_str() - && inm == etag - { - let mut resp = StatusCode::NOT_MODIFIED.into_response(); - if let Some(ref lm) = last_modified_str { - resp.headers_mut() - .insert(header::LAST_MODIFIED, lm.parse().unwrap()); - } - // Go sets ETag AFTER the 304 return paths (L235), so 304 does NOT include ETag + let (etag, last_modified_str) = etag_and_last_modified(&n); + if let Some(resp) = not_modified_response(&n, &headers, &etag, &last_modified_str) { return resp; } - // Chunk manifest expansion (needs full data) — after conditional checks, before response - // Pass ETag so chunk manifest responses include it (matches Go: ETag is set on the - // response writer before tryHandleChunkedFile runs). + // Chunk manifest expansion needs the full data; ETag is passed so expanded + // responses keep it (Go sets it before tryHandleChunkedFile runs). if n.is_chunk_manifest() - && !bypass_cm + && !request_kind.bypass_cm && let Some(resp) = try_expand_chunk_manifest( &state, &n, @@ -1371,6 +1117,416 @@ async fn get_or_head_handler_inner( } // If manifest expansion fails (invalid JSON etc.), fall through to raw data + let (mut response_headers, ext) = + read_response_headers(&n, &path, &query, &etag, &last_modified_str, ext); + + match strategy { + ReadStrategy::Stream(info) => { + return stream_response(&state, info, response_headers, track_download); + } + ReadStrategy::HeadFromMeta(info) => { + 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. + } + ReadStrategy::Buffered => {} + } + + let data = match buffered_payload( + n, + &headers, + &query, + &ext, + request_kind.has_image_ops, + &mut response_headers, + ) { + ControlFlow::Continue(data) => data, + ControlFlow::Break(resp) => return resp, + }; + buffered_response( + &state, + &headers, + &method, + data, + response_headers, + track_download, + ) +} + +/// How a GET/HEAD reply is produced once the needle is read. +enum ReadStrategy { + /// Large, served as stored, no range: streamed from the data file. + Stream(crate::storage::volume::NeedleStreamInfo), + /// HEAD: answered from the needle meta. + HeadFromMeta(crate::storage::volume::NeedleStreamInfo), + /// A range over a payload served as stored: read from the data file. + RangeFromSource(crate::storage::volume::NeedleStreamInfo), + /// Everything else, from the needle data in memory. + Buffered, +} + +/// The needle a GET/HEAD resolved to, and how its reply is produced. +struct ReadPlan { + needle: Needle, + strategy: ReadStrategy, +} + +/// The 401 reply for a bad read JWT; Go checks it before NewVolumeId, so a +/// bad path with JWT enabled is a 401, not a 400. +fn reject_read_jwt( + state: &VolumeServerState, + headers: &HeaderMap, + uri: &axum::http::Uri, + path: &str, +) -> Option { + let file_id = extract_file_id(path); + let token = extract_jwt(headers, uri); + if state + .guard + .read() + .unwrap() + .check_jwt_for_file(token.as_deref(), &file_id, false) + .is_err() + { + let body = serde_json::json!({"error": "wrong jwt"}); + return Some( + Response::builder() + .status(StatusCode::UNAUTHORIZED) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(serde_json::to_string(&body).unwrap())) + .unwrap(), + ); + } + None +} + +/// The volume is not here: 404, or proxy/redirect to a holder per read_mode. +async fn proxy_missing_volume( + state: &Arc, + uri: &axum::http::Uri, + request_headers: &HeaderMap, + path: &str, + vid: VolumeId, +) -> Response { + let query_string = uri.query().unwrap_or("").to_string(); + if is_proxied_query(&query_string) + || state.read_mode == ReadMode::Local + || state.master_url.is_empty() + { + return StatusCode::NOT_FOUND.into_response(); + } + + // For redirect, fid must be stripped of extension (Go parity: parseURLPath returns raw fid). + let info = match build_proxy_request_info(path, request_headers, &query_string) { + Some(info) => info, + None => return StatusCode::NOT_FOUND.into_response(), + }; + + proxy_or_redirect_to_target(state, info, vid, false).await +} + +/// Go's `proxied` query param check (loop prevention). +fn is_proxied_query(query_string: &str) -> bool { + query_string.split('&').any(|kv| kv == "proxied=true") +} + +/// Download throttling — Go's checkDownloadLimit + waitForDownloadSlot. +/// `Continue(true)` when the reply must count toward the inflight download bytes. +async fn check_download_limit( + state: &Arc, + uri: &axum::http::Uri, + request_headers: &HeaderMap, + path: &str, + vid: VolumeId, +) -> ControlFlow { + if state.concurrent_download_limit <= 0 { + return ControlFlow::Continue(false); + } + let timeout = state.inflight_download_data_timeout; + let deadline = tokio::time::Instant::now() + timeout; + let query_string = uri.query().unwrap_or("").to_string(); + + let current = state.inflight_download_bytes.load(Ordering::Relaxed); + if current > state.concurrent_download_limit { + metrics::HANDLER_COUNTER + .with_label_values(&[metrics::DOWNLOAD_LIMIT_COND]) + .inc(); + + // Go tries proxy to replica once before the wait loop + // (checkDownloadLimit); it does not retry on each wakeup. + let should_try_replica = + !is_proxied_query(&query_string) && !state.master_url.is_empty() && { + let store = state.store.read().unwrap(); + store + .find_volume(vid) + .is_some_and(|(_, vol)| vol.super_block.replica_placement.has_replication()) + }; + if should_try_replica + && let Some(info) = build_proxy_request_info(path, request_headers, &query_string) + { + return ControlFlow::Break(proxy_or_redirect_to_target(state, info, vid, true).await); + } + + // Blocking wait loop (Go's waitForDownloadSlot) + loop { + if tokio::time::timeout_at(deadline, state.download_notify.notified()) + .await + .is_err() + { + return ControlFlow::Break(json_error_with_query( + StatusCode::TOO_MANY_REQUESTS, + "download limit exceeded", + uri.query(), + )); + } + let current = state.inflight_download_bytes.load(Ordering::Relaxed); + if current <= state.concurrent_download_limit { + break; + } + } + } + ControlFlow::Continue(true) +} + +/// The URL extension, and how the reply may be served. +fn parse_read_request( + path: &str, + headers: &HeaderMap, + 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) + && (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, + bypass_cm: query.cm.as_deref() == Some("false"), + }; + (ext, request_kind) +} + +/// Full read from an EC volume; there is no streaming for EC. +async fn read_ec_shard_needle( + state: &Arc, + vid: VolumeId, + needle_id: NeedleId, + cookie: Cookie, +) -> ControlFlow { + // The distributed read does a local-first pass in its Snapshot phase, so + // calling it directly covers both the all-local fast case and the + // peer-fetch + reconstruct case without reading local intervals twice. + let n = match crate::server::store_ec::read_ec_shard_needle_distributed(state, vid, needle_id) + .await + { + Ok(Some(ec_needle)) => ec_needle, + Ok(None) => { + metrics::HANDLER_COUNTER + .with_label_values(&[metrics::ERROR_GET_NOT_FOUND]) + .inc(); + return ControlFlow::Break(StatusCode::NOT_FOUND.into_response()); + } + Err(e) => { + let kind = if e.kind() == std::io::ErrorKind::NotFound { + metrics::ERROR_GET_NOT_FOUND + } else { + metrics::ERROR_GET_INTERNAL + }; + metrics::HANDLER_COUNTER.with_label_values(&[kind]).inc(); + if e.kind() == std::io::ErrorKind::NotFound { + return ControlFlow::Break(StatusCode::NOT_FOUND.into_response()); + } + return ControlFlow::Break( + (StatusCode::INTERNAL_SERVER_ERROR, format!("ec read: {}", e)).into_response(), + ); + } + }; + + // Validate cookie (matches Go behavior after ReadEcShardNeedle) + if n.cookie != cookie { + return ControlFlow::Break(StatusCode::NOT_FOUND.into_response()); + } + ControlFlow::Continue(ReadPlan { + needle: n, + strategy: ReadStrategy::Buffered, + }) +} + +/// Reads a regular volume's needle off the store lock, meta-only when the +/// reply can be served from the data file. +async fn read_volume_needle( + state: &Arc, + vid: VolumeId, + needle_id: NeedleId, + cookie: Cookie, + read_deleted: bool, + request_kind: SourceReadRequest, +) -> ControlFlow { + let has_range = request_kind.has_range; + let is_head = request_kind.is_head; + + let read_state = state.clone(); + let read = tokio::task::spawn_blocking(move || { + read_needle_for_get( + &read_state, + vid, + needle_id, + cookie, + read_deleted, + request_kind, + ) + }) + .await; + let (n, stream_info) = match read { + Ok(Ok(Some(found))) => found, + // Cookie mismatch + Ok(Ok(None)) => return ControlFlow::Break(StatusCode::NOT_FOUND.into_response()), + Ok(Err( + crate::storage::volume::VolumeError::NotFound + | crate::storage::volume::VolumeError::Deleted, + )) => { + metrics::HANDLER_COUNTER + .with_label_values(&[metrics::ERROR_GET_NOT_FOUND]) + .inc(); + return ControlFlow::Break(StatusCode::NOT_FOUND.into_response()); + } + Ok(Err(e)) => { + metrics::HANDLER_COUNTER + .with_label_values(&[metrics::ERROR_GET_INTERNAL]) + .inc(); + return ControlFlow::Break( + ( + StatusCode::INTERNAL_SERVER_ERROR, + format!("read error: {}", e), + ) + .into_response(), + ); + } + Err(e) => { + metrics::HANDLER_COUNTER + .with_label_values(&[metrics::ERROR_GET_INTERNAL]) + .inc(); + return ControlFlow::Break( + ( + StatusCode::INTERNAL_SERVER_ERROR, + format!("read error: {}", e), + ) + .into_response(), + ); + } + }; + + // Stream info is only returned for a reply served from the data file. + let can_direct_source_read = stream_info.is_some() && request_kind.direct(&n); + + 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. + let can_handle_head_from_meta = stream_info.is_some() && is_head; + let can_handle_range_from_source = can_direct_source_read && has_range; + + let strategy = match stream_info { + Some(info) if can_stream => ReadStrategy::Stream(info), + Some(info) if can_handle_head_from_meta => ReadStrategy::HeadFromMeta(info), + Some(info) if can_handle_range_from_source => ReadStrategy::RangeFromSource(info), + _ => ReadStrategy::Buffered, + }; + ControlFlow::Continue(ReadPlan { + needle: n, + strategy, + }) +} + +/// The ETag, and Last-Modified in RFC 1123 format. +fn etag_and_last_modified(n: &Needle) -> (String, Option) { + let etag = format!("\"{}\"", n.etag()); + let last_modified_str = if n.last_modified > 0 { + use chrono::{TimeZone, Utc}; + Utc.timestamp_opt(n.last_modified as i64, 0) + .single() + .map(|dt| dt.format("%a, %d %b %Y %H:%M:%S GMT").to_string()) + } else { + None + }; + (etag, last_modified_str) +} + +/// The 304 reply, if a conditional header matches. +fn not_modified_response( + n: &Needle, + headers: &HeaderMap, + etag: &str, + last_modified_str: &Option, +) -> Option { + // If-Modified-Since first (Go order), then If-None-Match. + if n.last_modified > 0 + && let Some(ims_header) = headers.get(header::IF_MODIFIED_SINCE) + && let Ok(ims_str) = ims_header.to_str() + { + // Parse HTTP date format: "Mon, 02 Jan 2006 15:04:05 GMT" + if let Ok(ims_time) = + chrono::NaiveDateTime::parse_from_str(ims_str, "%a, %d %b %Y %H:%M:%S GMT") + && (n.last_modified as i64) <= ims_time.and_utc().timestamp() + { + let mut resp = StatusCode::NOT_MODIFIED.into_response(); + if let Some(lm) = last_modified_str { + resp.headers_mut() + .insert(header::LAST_MODIFIED, lm.parse().unwrap()); + } + // Go sets ETag after the 304 paths, so a 304 has no ETag. + return Some(resp); + } + } + + if let Some(if_none_match) = headers.get(header::IF_NONE_MATCH) + && let Ok(inm) = if_none_match.to_str() + && inm == etag + { + let mut resp = StatusCode::NOT_MODIFIED.into_response(); + if let Some(lm) = last_modified_str { + resp.headers_mut() + .insert(header::LAST_MODIFIED, lm.parse().unwrap()); + } + // Go sets ETag after the 304 paths, so a 304 has no ETag. + return Some(resp); + } + None +} + +/// The 200 reply headers, and the extension, which falls back to the stored name's. +fn read_response_headers( + n: &Needle, + path: &str, + query: &ReadQueryParams, + etag: &str, + last_modified_str: &Option, + ext: String, +) -> (HeaderMap, String) { let mut response_headers = HeaderMap::new(); response_headers.insert(header::ETAG, etag.parse().unwrap()); @@ -1391,7 +1547,7 @@ async fn get_or_head_handler_inner( } // H8: Use needle stored name when URL path has no filename (only vid,fid) - let mut filename = extract_filename_from_path(&path); + let mut filename = extract_filename_from_path(path); let mut ext = ext; if n.name_size > 0 && filename.is_empty() { filename = String::from_utf8_lossy(&n.name).to_string(); @@ -1500,7 +1656,7 @@ async fn get_or_head_handler_inner( } // Last-Modified - if let Some(ref lm) = last_modified_str { + if let Some(lm) = last_modified_str { response_headers.insert(header::LAST_MODIFIED, lm.parse().unwrap()); } @@ -1521,105 +1677,126 @@ async fn get_or_head_handler_inner( response_headers.insert(header::CONTENT_DISPOSITION, hval); } } + (response_headers, ext) +} - // ---- Streaming path: large uncompressed files ---- - if can_stream && let Some(info) = stream_info { - response_headers.insert(header::ACCEPT_RANGES, "bytes".parse().unwrap()); - response_headers.insert( - header::CONTENT_LENGTH, - info.data_size.to_string().parse().unwrap(), - ); +/// Streaming path: large uncompressed files. +fn stream_response( + state: &Arc, + info: crate::storage::volume::NeedleStreamInfo, + mut response_headers: HeaderMap, + track_download: bool, +) -> Response { + response_headers.insert(header::ACCEPT_RANGES, "bytes".parse().unwrap()); + response_headers.insert( + header::CONTENT_LENGTH, + info.data_size.to_string().parse().unwrap(), + ); - let tracked_bytes = info.data_size as i64; - let tracking_state = if download_guard.is_some() { - let new_val = state - .inflight_download_bytes - .fetch_add(tracked_bytes, Ordering::Relaxed) - + tracked_bytes; - metrics::INFLIGHT_DOWNLOAD_SIZE.set(new_val); - Some(state.clone()) - } else { + let tracked_bytes = info.data_size as i64; + let tracking_state = if track_download { + let new_val = state + .inflight_download_bytes + .fetch_add(tracked_bytes, Ordering::Relaxed) + + tracked_bytes; + metrics::INFLIGHT_DOWNLOAD_SIZE.set(new_val); + Some(state.clone()) + } else { + None + }; + + let streaming = StreamingBody { + source: Arc::new(info.source), + data_offset: info.data_file_offset, + data_size: info.data_size, + pos: 0, + chunk_size: streaming_chunk_size(state.read_buffer_size_bytes, info.data_size as usize), + _held_read_lease: if state.has_slow_read { None - }; + } else { + Some(info.data_file_access_control.read_lock()) + }, + data_file_access_control: info.data_file_access_control, + hold_read_lock_for_stream: !state.has_slow_read, + io_unavailable: info.io_unavailable, + pending: None, + buf: bytes::BytesMut::new(), + state: tracking_state, + tracked_bytes, + needle_id: info.needle_id, + expected_checksum: info.checksum, + crc: CRC(0), + }; - let streaming = StreamingBody { - source: Arc::new(info.source), - data_offset: info.data_file_offset, - data_size: info.data_size, - pos: 0, - chunk_size: streaming_chunk_size(state.read_buffer_size_bytes, info.data_size as usize), - _held_read_lease: if state.has_slow_read { - None - } else { - Some(info.data_file_access_control.read_lock()) - }, - data_file_access_control: info.data_file_access_control, - hold_read_lock_for_stream: !state.has_slow_read, - io_unavailable: info.io_unavailable, - pending: None, - buf: bytes::BytesMut::new(), - state: tracking_state, - tracked_bytes, - needle_id: info.needle_id, - expected_checksum: info.checksum, - crc: CRC(0), - }; + let body = Body::new(streaming); + let mut resp = Response::new(body); + *resp.status_mut() = StatusCode::OK; + *resp.headers_mut() = response_headers; + resp +} - let body = Body::new(streaming); - let mut resp = Response::new(body); - *resp.status_mut() = StatusCode::OK; - *resp.headers_mut() = response_headers; - return resp; - } +/// HEAD from the needle meta: Content-Length is the stored data size. +fn head_from_meta_response( + info: &crate::storage::volume::NeedleStreamInfo, + mut response_headers: HeaderMap, +) -> Response { + response_headers.insert( + header::CONTENT_LENGTH, + info.data_size.to_string().parse().unwrap(), + ); + (StatusCode::OK, response_headers).into_response() +} - if can_handle_head_from_meta && let Some(info) = stream_info { - response_headers.insert( - header::CONTENT_LENGTH, - info.data_size.to_string().parse().unwrap(), - ); - return (StatusCode::OK, response_headers).into_response(); - } - - if can_handle_range_from_source - && let (Some(range_header), Some(info)) = (headers.get(header::RANGE), stream_info) - && let Ok(range_str) = range_header.to_str() - { - let range_str = range_str.to_string(); - let tracking = track_download.then(|| state.clone()); - return tokio::task::spawn_blocking(move || { - handle_range_request_from_source(&range_str, info, response_headers, tracking) - }) - .await - .unwrap_or_else(|e| { - ( - StatusCode::INTERNAL_SERVER_ERROR, - format!("range read error: {}", e), - ) - .into_response() - }); - } - - // ---- Buffered path: small files, compressed, images, range requests ---- +/// A range over a payload served as stored, read from the data file. +async fn range_from_source_response( + state: &Arc, + range_str: &str, + info: crate::storage::volume::NeedleStreamInfo, + response_headers: HeaderMap, + track_download: bool, +) -> Response { + let range_str = range_str.to_string(); + let tracking = track_download.then(|| state.clone()); + tokio::task::spawn_blocking(move || { + handle_range_request_from_source(&range_str, info, response_headers, tracking) + }) + .await + .unwrap_or_else(|e| { + ( + StatusCode::INTERNAL_SERVER_ERROR, + format!("range read error: {}", e), + ) + .into_response() + }) +} +/// The needle data, decompressed and image-processed as the request needs. +fn buffered_payload( + n: Needle, + headers: &HeaderMap, + query: &ReadQueryParams, + ext: &str, + needs_image_ops: bool, + response_headers: &mut HeaderMap, +) -> ControlFlow> { // Handle compressed data: if needle is compressed, either pass through or decompress let is_compressed = n.is_compressed(); let mut data = n.data; - // Check if image operations are needed — must decompress first regardless of Accept-Encoding - // Go checks resize (.webp OK) and crop (.webp NOT OK) separately. - let needs_image_ops = has_resize_ops || has_crop_ops; - + // Image operations must decompress first regardless of Accept-Encoding. if is_compressed { if needs_image_ops { // Always decompress for image operations (Go decompresses before resize/crop) match maybe_decompress_gzip(&data) { Ok(decompressed) => data = decompressed, Err(GunzipError::TooLarge) => { - return ( - StatusCode::PAYLOAD_TOO_LARGE, - "compressed object exceeds decompression limit", - ) - .into_response(); + return ControlFlow::Break( + ( + StatusCode::PAYLOAD_TOO_LARGE, + "compressed object exceeds decompression limit", + ) + .into_response(), + ); } Err(GunzipError::Decode) => {} // not valid gzip; keep raw bytes } @@ -1641,11 +1818,13 @@ async fn get_or_head_handler_inner( match maybe_decompress_gzip(&data) { Ok(decompressed) => data = decompressed, Err(GunzipError::TooLarge) => { - return ( - StatusCode::PAYLOAD_TOO_LARGE, - "compressed object exceeds decompression limit", - ) - .into_response(); + return ControlFlow::Break( + ( + StatusCode::PAYLOAD_TOO_LARGE, + "compressed object exceeds decompression limit", + ) + .into_response(), + ); } Err(GunzipError::Decode) => {} // not valid gzip; keep raw bytes } @@ -1655,17 +1834,35 @@ async fn get_or_head_handler_inner( // Image crop and resize — Go checks extensions separately per operation. // Crop: .png .jpg .jpeg .gif (no .webp). Resize: .png .jpg .jpeg .gif .webp. - if is_image_crop_ext(&ext) { - data = maybe_crop_image(&data, &ext, &query); + if is_image_crop_ext(ext) { + data = maybe_crop_image(&data, ext, query); } - if is_image_resize_ext(&ext) { - data = maybe_resize_image(&data, &ext, &query); + if is_image_resize_ext(ext) { + data = maybe_resize_image(&data, ext, query); } + ControlFlow::Continue(data) +} - // Accept-Ranges +/// The buffered reply over the payload, whole or a range. +fn buffered_response( + state: &Arc, + headers: &HeaderMap, + method: &Method, + data: Vec, + mut response_headers: HeaderMap, + track_download: bool, +) -> Response { response_headers.insert(header::ACCEPT_RANGES, "bytes".parse().unwrap()); - // Check Range header + // HEAD before Range (Go's writeResponseContent order). + if method == Method::HEAD { + response_headers.insert( + header::CONTENT_LENGTH, + data.len().to_string().parse().unwrap(), + ); + return (StatusCode::OK, response_headers).into_response(); + } + if let Some(range_header) = headers.get(header::RANGE) && let Ok(range_str) = range_header.to_str() { @@ -1677,14 +1874,6 @@ async fn get_or_head_handler_inner( ); } - if method == Method::HEAD { - response_headers.insert( - header::CONTENT_LENGTH, - data.len().to_string().parse().unwrap(), - ); - return (StatusCode::OK, response_headers).into_response(); - } - finalize_bytes_response( StatusCode::OK, response_headers,