From 4c50a643cf49f7447feda5601da3a7dad69fc857 Mon Sep 17 00:00:00 2001 From: mlevkov Date: Fri, 25 Sep 2026 18:11:17 -0700 Subject: [PATCH] fix(connectors): refuse http_source bodies Iggy cannot store Iggy refuses a message with an empty payload or one over MAX_PAYLOAD_SIZE. http_source answered 200 to both, and the runtime then could not build the message and NACKed the whole batch around it. The source replays a NACKed batch before anything newer, so nothing queued behind it was ever delivered, and the SDK stopped the poll task after five NACKs. An empty body now gets 400 before anything is queued. The ceiling on max_body_size_bytes is now MAX_PAYLOAD_SIZE (64,000,000 bytes) instead of 64 MiB, so every body the handlers accept can be stored, and a larger one gets the existing 413. A max_body_size_bytes above the new ceiling fails open(). --- core/connectors/sources/http_source/README.md | 3 +- .../sources/http_source/src/auth.rs | 3 +- .../connectors/sources/http_source/src/lib.rs | 11 +++-- .../sources/http_source/src/server.rs | 40 ++++++++++++++++++ .../tests/connectors/http/http_source.rs | 42 +++++++++++++++++++ 5 files changed, 93 insertions(+), 6 deletions(-) diff --git a/core/connectors/sources/http_source/README.md b/core/connectors/sources/http_source/README.md index b8c0029517..4ce82ea377 100644 --- a/core/connectors/sources/http_source/README.md +++ b/core/connectors/sources/http_source/README.md @@ -119,7 +119,7 @@ The batch is then replayed on every poll and the SDK stops the poll task after f | `topic_path` | string | none | Exposes `POST /topics/{topic_path}`. Unset leaves only secret-path endpoints. | | `auth_bearer_token` | string | none | Guards the named topic path. Unset leaves it unauthenticated, for deployments behind an authenticating gateway. A misspelled key reads as unset, so it opens the path rather than failing; see the note on unknown keys below. | | `management_token` | string | none | Enables `/admin/endpoints`. Unset means the management API does not exist. | -| `max_body_size_bytes` | usize | `1048576` | Request body limit, applied by the handlers rather than an extractor. Routing wins over it: an oversized POST to an unknown or revoked path answers 404 without the body being read. Must match across instances sharing a listener. **Max 67108864**; a larger value fails `open()`. | +| `max_body_size_bytes` | usize | `1048576` | Request body limit, applied by the handlers rather than an extractor. Routing wins over it: an oversized POST to an unknown or revoked path answers 404 without the body being read. Must match across instances sharing a listener. **Max 64000000**, Iggy's message payload cap, since a body becomes the payload unchanged; a larger value fails `open()`. | | `buffer_capacity` | usize | `10000` | Messages the instance bridge holds. A full bridge answers 429, which since #3855 signals either an arrival burst or a slow Iggy, since the poll loop stalls waiting for the previous batch to be acknowledged. **Max 1000000**; a larger value fails `open()`. | | `max_batch_size` | usize | `500` | Maximum messages a single `poll()` returns. **Max 100000**; a larger value fails `open()`. | | `include_http_metadata` | bool | `true` | Adds instance, peer address, and receive time as message headers. | @@ -184,6 +184,7 @@ Content-Type: application/json | 401 | Bearer or HMAC validation failed | `{"error":"unauthorized"}` | | 404 | Unknown path, or a revoked or expired endpoint | `{"error":"not found"}` | | 400 | Malformed request body, e.g. the client reset mid-send | `{"error":"bad request"}` | +| 400 | Empty body, which Iggy cannot store as a message | `{"error":"empty body"}` | | 413 | Body over `max_body_size_bytes` | `{"error":"payload too large"}` | | 405 | A known path with the wrong method | `{"error":"method not allowed"}` | | 429 | Bridge full | `{"error":"too many requests"}` plus `Retry-After: 1` | diff --git a/core/connectors/sources/http_source/src/auth.rs b/core/connectors/sources/http_source/src/auth.rs index 18100a2df0..0b887fe57c 100644 --- a/core/connectors/sources/http_source/src/auth.rs +++ b/core/connectors/sources/http_source/src/auth.rs @@ -54,7 +54,8 @@ fn strip_bearer(header_value: &str) -> Option<&str> { /// Unlike a revoke `reason`, which is written once onto a tombstone, these sit /// on active entries: every `mutate_registry` deep-clones them and every flush /// re-serializes the whole registry, so they are paid for repeatedly. Without a -/// cap the only limit was the body limit, up to 64 MiB, times `MAX_ENDPOINTS`. +/// cap the only limit was the body limit, up to `MAX_BODY_SIZE_BYTES_LIMIT`, +/// times `MAX_ENDPOINTS`. /// /// Generous on purpose. A real HMAC secret is tens of bytes and a header name /// is shorter still; these refuse abuse without refusing anything an operator diff --git a/core/connectors/sources/http_source/src/lib.rs b/core/connectors/sources/http_source/src/lib.rs index db633a0c5b..a0961ef63e 100644 --- a/core/connectors/sources/http_source/src/lib.rs +++ b/core/connectors/sources/http_source/src/lib.rs @@ -26,7 +26,7 @@ mod types; use arc_swap::{ArcSwap, Guard}; use async_trait::async_trait; use axum::http::{HeaderName, header}; -use iggy_common::{HeaderKey, HeaderValue}; +use iggy_common::{HeaderKey, HeaderValue, MAX_PAYLOAD_SIZE}; use iggy_connector_sdk::{ ConnectorState, Error, ProducedMessage, ProducedMessages, Schema, Source, source, source_connector, @@ -70,7 +70,10 @@ pub const BUFFER_CAPACITY_LIMIT: usize = 1_000_000; /// until it hits this many messages. Nothing is allocated up front, so an /// oversized value costs a longer drain rather than a large empty buffer. pub const MAX_BATCH_SIZE_LIMIT: usize = 100_000; -pub const MAX_BODY_SIZE_BYTES_LIMIT: usize = 64 * 1024 * 1024; +/// Iggy's own payload cap, because a body becomes the message payload byte +/// for byte. The runtime cannot send a larger one, and it NACKs the whole +/// batch around it on every replay. +pub const MAX_BODY_SIZE_BYTES_LIMIT: usize = MAX_PAYLOAD_SIZE as usize; pub const DEFAULT_HMAC_HEADER: &str = "X-Hub-Signature-256"; pub const DEFAULT_HMAC_PREFIX: &str = "sha256="; @@ -905,7 +908,7 @@ impl HttpSource { /// /// A batch that can never be delivered would replay forever, but this /// connector rejects malformed work at the door rather than mid-stream: - /// oversized bodies get 413 before a handler runs, headers are clamped on + /// an empty body gets 400 and an oversized one 413, headers are clamped on /// accept, and `Schema::Raw` cannot fail to decode. fn on_nack(&self) -> Result<(), Error> { let mut staged = self.lock_staged(); @@ -1352,7 +1355,7 @@ mod tests { ), ( "max_body_size_bytes", - r#"{"listen_addr": "0.0.0.0:9090", "max_body_size_bytes": 67108865}"#, + r#"{"listen_addr": "0.0.0.0:9090", "max_body_size_bytes": 64000001}"#, ), ] { assert!( diff --git a/core/connectors/sources/http_source/src/server.rs b/core/connectors/sources/http_source/src/server.rs index 5b98d1a85c..4f4d7ed6fc 100644 --- a/core/connectors/sources/http_source/src/server.rs +++ b/core/connectors/sources/http_source/src/server.rs @@ -1095,6 +1095,11 @@ fn enqueue( body: Bytes, metrics: &Metrics, ) -> Response { + // Iggy refuses an empty payload, and the runtime NACKs the whole batch + // around it on every replay. Checked first, because no retry can fix it. + if body.is_empty() { + return error_response(StatusCode::BAD_REQUEST, "empty body"); + } // Checked before the header map is built and the body copied, since a full // bridge throws both away. Both handlers check once more before they read // the body at all; this one catches a bridge that filled during the read. @@ -1577,6 +1582,41 @@ mod tests { close(&mut source).await; } + #[tokio::test] + async fn given_empty_body_when_posted_to_secret_path_should_answer_bad_request() { + let mut source = open(1, config(free_port(), free_port(), &[ENDPOINT_ONE])).await; + + let response = post_signed(&base_url(&source), ENDPOINT_ONE, "").await; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + let body: serde_json::Value = response + .json() + .await + .expect("the refusal must carry this API's json"); + assert_eq!(body["error"], "empty body"); + assert_eq!( + source.shared.sender.len(), + 0, + "an empty body queued is a batch the runtime NACKs on every replay" + ); + close(&mut source).await; + } + + #[tokio::test] + async fn given_empty_body_when_posted_to_named_path_should_answer_bad_request() { + let mut source = open(1, config(free_port(), free_port(), &[])).await; + + let response = client() + .post(format!("{}/topics/github", base_url(&source))) + .send() + .await + .expect("the request must reach the listener"); + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + assert_eq!(source.shared.sender.len(), 0); + close(&mut source).await; + } + #[tokio::test] async fn given_oversized_body_when_posted_should_answer_payload_too_large() { let mut config = config(free_port(), free_port(), &[ENDPOINT_ONE]); diff --git a/core/integration/tests/connectors/http/http_source.rs b/core/integration/tests/connectors/http/http_source.rs index e1720be777..5d457476e1 100644 --- a/core/integration/tests/connectors/http/http_source.rs +++ b/core/integration/tests/connectors/http/http_source.rs @@ -124,6 +124,48 @@ async fn named_path_post_produces_message_to_iggy( ); } +/// Iggy refuses to store an empty payload. Accepted with a 200, an empty body +/// made the runtime NACK its whole batch on every replay, so no webhook queued +/// behind it ever reached Iggy. +#[iggy_harness( + server(connectors_runtime(config_path = "tests/connectors/http/source.toml")), + seed = seeds::connector_multi_topic_stream +)] +async fn empty_body_post_is_refused_without_blocking_later_webhooks( + harness: &TestHarness, + fixture: HttpSourceFixture, +) { + let client = harness.root_client().await.unwrap(); + let http = webhook_client(); + wait_for_gateway(&http, &fixture).await; + + let empty = http + .post(fixture.named_url(seeds::names::TOPIC)) + .send() + .await + .expect("Failed to POST the empty body"); + let body = r#"{"event":"deployment","environment":"production"}"#; + let response = http + .post(fixture.named_url(seeds::names::TOPIC)) + .body(body) + .send() + .await + .expect("Failed to POST the webhook"); + assert_eq!(response.status(), StatusCode::OK); + + let messages = poll_payloads(&client, seeds::names::TOPIC, "http_source_cg_empty", 1).await; + assert_eq!( + String::from_utf8_lossy(&messages[0].0), + body, + "the webhook after an empty body must still be delivered" + ); + assert_eq!( + empty.status(), + StatusCode::BAD_REQUEST, + "an empty body must be refused, not queued" + ); +} + #[iggy_harness( server(connectors_runtime(config_path = "tests/connectors/http/source.toml")), seed = seeds::connector_multi_topic_stream