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