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,