From db9ec0605eec8e99c9cc7cc5906b6f6d02f86b07 Mon Sep 17 00:00:00 2001 From: Hampus Date: Sat, 3 Oct 2026 02:10:55 +0200 Subject: [PATCH] fix(media-proxy): stop rejecting large storage transport chunks (#3144) --- .../src/server/external/fetch.rs | 34 +++++--------- .../src/server/external/tests/mod.rs | 41 ++++++++++++++++ .../src/storage/response_body.rs | 19 +++----- .../src/storage/tests/response_body.rs | 47 +++++++++++++++++++ 4 files changed, 107 insertions(+), 34 deletions(-) diff --git a/fluxer_media_proxy/src/server/external/fetch.rs b/fluxer_media_proxy/src/server/external/fetch.rs index 98ce3e1f8..972133315 100644 --- a/fluxer_media_proxy/src/server/external/fetch.rs +++ b/fluxer_media_proxy/src/server/external/fetch.rs @@ -253,8 +253,7 @@ async fn external_body_prefix( prefix .try_reserve_exact(EXTERNAL_SNIFF_PREFIX_BYTES) .map_err(|_| ExternalFetchError::BufferAllocationFailed)?; - let mut chunks_read = 0_u64; - let chunks_max = + let mut empty_chunks_remaining = response_body_limit::response_body_chunk_limit(constants::MAX_MEDIA_PROXY_BYTES as u64); while prefix.len() < EXTERNAL_SNIFF_PREFIX_BYTES { let Some(chunk) = response.chunk().await.map_err(|err| { @@ -265,17 +264,13 @@ async fn external_body_prefix( else { break; }; - if chunk.len() > response_body_limit::RESPONSE_BODY_TRANSPORT_CHUNK_BYTES_MAX { - warn!(url = %url, "external response transport chunk exceeded its byte bound"); - return Err(ExternalFetchError::PayloadTooLarge); - } - chunks_read = chunks_read - .checked_add(1) - .filter(|chunks| *chunks <= chunks_max) - .ok_or_else(|| { - warn!(url = %url, chunks_max, "external sniff prefix exceeded its chunk limit"); + if chunk.is_empty() { + empty_chunks_remaining = empty_chunks_remaining.checked_sub(1).ok_or_else(|| { + warn!(url = %url, "external sniff prefix exceeded its empty chunk limit"); ExternalFetchError::PayloadTooLarge })?; + continue; + } prefix.extend_from_slice(&chunk); } Ok(Bytes::from(prefix)) @@ -370,24 +365,19 @@ pub(super) async fn buffer_external_response( let mut reserved_bytes = grow_to_capacity(&mut reservation, metrics, buf.capacity(), initial_capacity)?; buf.extend_from_slice(&prefix); - let mut chunks_read = 0_u64; - let chunks_max = response_body_limit::response_body_chunk_limit(limit as u64); + let mut empty_chunks_remaining = response_body_limit::response_body_chunk_limit(limit as u64); while let Some(chunk) = response.chunk().await.map_err(|err| { warn!(url = %url, %err, "external body read failed"); metrics.record_fetch_failure(); ExternalFetchError::FetchFailed })? { - if chunk.len() > response_body_limit::RESPONSE_BODY_TRANSPORT_CHUNK_BYTES_MAX { - warn!(url = %url, "external response transport chunk exceeded its byte bound"); - return Err(ExternalFetchError::PayloadTooLarge); - } - chunks_read = chunks_read - .checked_add(1) - .filter(|chunks| *chunks <= chunks_max) - .ok_or_else(|| { - warn!(url = %url, chunks_max, "external payload exceeded its chunk limit"); + if chunk.is_empty() { + empty_chunks_remaining = empty_chunks_remaining.checked_sub(1).ok_or_else(|| { + warn!(url = %url, "external payload exceeded its empty chunk limit"); ExternalFetchError::PayloadTooLarge })?; + continue; + } let Some(next_len) = buf .len() .checked_add(chunk.len()) diff --git a/fluxer_media_proxy/src/server/external/tests/mod.rs b/fluxer_media_proxy/src/server/external/tests/mod.rs index b5b43a715..b5594382f 100644 --- a/fluxer_media_proxy/src/server/external/tests/mod.rs +++ b/fluxer_media_proxy/src/server/external/tests/mod.rs @@ -929,3 +929,44 @@ async fn an_upstream_that_ignores_the_client_range_still_streams_a_complete_body let body = to_bytes(response.into_body(), 64).await.unwrap(); assert_eq!(b"streamed bytes", body.as_ref()); } + +#[tokio::test] +async fn external_buffering_accepts_transport_chunks_of_any_size_within_the_limit() { + const LARGE_CHUNK_BYTES: usize = + crate::response_body_limit::RESPONSE_BODY_TRANSPORT_CHUNK_BYTES_MAX + 155_648; + const SMALL_CHUNK_BYTES: usize = 1448; + const SMALL_CHUNKS: usize = 512; + let budget = ByteBudget::new(constants::MAX_MEDIA_PROXY_BYTES * 4); + let metrics = ExternalMetrics::new(); + let buffer = |response: reqwest::Response, length: usize| { + buffer_external_response(ExternalBufferRequest { + response, + prefix: Bytes::new(), + url: "https://media.example.test/clip.webm", + budget: &budget, + metrics: &metrics, + content_length: Some(length as u64), + limit: constants::MAX_MEDIA_PROXY_BYTES, + }) + }; + + let large = buffer( + reqwest::Response::from(http::Response::new(vec![7u8; LARGE_CHUNK_BYTES])), + LARGE_CHUNK_BYTES, + ) + .await + .expect("a transport chunk larger than the reserved allowance"); + assert_eq!(LARGE_CHUNK_BYTES, large.as_bytes().len()); + + let small_body = reqwest::Body::wrap_stream(futures_util::stream::iter( + (0..SMALL_CHUNKS) + .map(|_| Ok::(Bytes::from(vec![9u8; SMALL_CHUNK_BYTES]))), + )); + let small = buffer( + reqwest::Response::from(http::Response::new(small_body)), + SMALL_CHUNK_BYTES * SMALL_CHUNKS, + ) + .await + .expect("packet-sized transport chunks"); + assert_eq!(SMALL_CHUNK_BYTES * SMALL_CHUNKS, small.as_bytes().len()); +} diff --git a/fluxer_media_proxy/src/storage/response_body.rs b/fluxer_media_proxy/src/storage/response_body.rs index a774b7f05..c49686294 100644 --- a/fluxer_media_proxy/src/storage/response_body.rs +++ b/fluxer_media_proxy/src/storage/response_body.rs @@ -80,22 +80,17 @@ pub(super) async fn read_response_bytes( body.try_reserve_exact(expected_length) .map_err(|_| StorageError::BufferAllocationFailed)?; buffer_budget.grow_to(body.capacity())?; - let mut chunks_read = 0_u64; - let chunks_max = response_body_limit::response_body_chunk_limit(expected_length as u64); + let mut empty_chunks_remaining = + response_body_limit::response_body_chunk_limit(expected_length as u64); while let Some(chunk) = response.chunk().await? { - if chunk.len() > response_body_limit::RESPONSE_BODY_TRANSPORT_CHUNK_BYTES_MAX { - return Err(StorageError::ObjectStorage(anyhow::anyhow!( - "object storage response transport chunk exceeded its byte bound" - ))); - } - chunks_read = chunks_read - .checked_add(1) - .filter(|chunks| *chunks <= chunks_max) - .ok_or_else(|| { + if chunk.is_empty() { + empty_chunks_remaining = empty_chunks_remaining.checked_sub(1).ok_or_else(|| { StorageError::ObjectStorage(anyhow::anyhow!( - "object storage response exceeded its chunk limit" + "object storage response exceeded its empty chunk limit" )) })?; + continue; + } let next_length = body .len() .checked_add(chunk.len()) diff --git a/fluxer_media_proxy/src/storage/tests/response_body.rs b/fluxer_media_proxy/src/storage/tests/response_body.rs index db0421414..19dc3b24d 100644 --- a/fluxer_media_proxy/src/storage/tests/response_body.rs +++ b/fluxer_media_proxy/src/storage/tests/response_body.rs @@ -282,3 +282,50 @@ async fn exact_stream_accepts_small_transport_chunks_and_bounds_empty_ones() { .expect_err("empty chunk flood"); assert_eq!(error.kind(), std::io::ErrorKind::InvalidData); } + +fn chunked_provider_response(chunks: Vec) -> reqwest::Response { + let body = reqwest::Body::wrap_stream(stream::iter( + chunks.into_iter().map(Ok::), + )); + reqwest::Response::from(http::Response::new(body)) +} + +#[tokio::test] +async fn response_reader_accepts_transport_chunks_of_any_size_within_the_length() { + const LARGE_CHUNK_BYTES: usize = + crate::response_body_limit::RESPONSE_BODY_TRANSPORT_CHUNK_BYTES_MAX + 155_648; + let budget = ByteBudget::new(4 * LARGE_CHUNK_BYTES); + let large = read_response_bytes( + provider_response(vec![7u8; LARGE_CHUNK_BYTES]), + LARGE_CHUNK_BYTES, + &budget, + ) + .await + .expect("a transport chunk larger than the reserved allowance"); + assert_eq!(large.as_ref().len(), LARGE_CHUNK_BYTES); + + const SMALL_CHUNK_BYTES: usize = 1448; + const SMALL_CHUNKS: usize = 512; + let small = read_response_bytes( + chunked_provider_response( + (0..SMALL_CHUNKS) + .map(|_| Bytes::from(vec![9u8; SMALL_CHUNK_BYTES])) + .collect(), + ), + SMALL_CHUNK_BYTES * SMALL_CHUNKS, + &budget, + ) + .await + .expect("packet-sized transport chunks"); + assert_eq!(small.as_ref().len(), SMALL_CHUNK_BYTES * SMALL_CHUNKS); + + assert!(matches!( + read_response_bytes( + chunked_provider_response((0..4096).map(|_| Bytes::new()).collect()), + 4, + &budget, + ) + .await, + Err(StorageError::ObjectStorage(_)) + )); +}