Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion core/connectors/sources/http_source/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand Down Expand Up @@ -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` |
Expand Down
3 changes: 2 additions & 1 deletion core/connectors/sources/http_source/src/auth.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
11 changes: 7 additions & 4 deletions core/connectors/sources/http_source/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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=";
Expand Down Expand Up @@ -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
Comment thread
hubcio marked this conversation as resolved.
/// accept, and `Schema::Raw` cannot fail to decode.
fn on_nack(&self) -> Result<(), Error> {
let mut staged = self.lock_staged();
Expand Down Expand Up @@ -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}"#,
Comment thread
hubcio marked this conversation as resolved.
),
] {
assert!(
Expand Down
40 changes: 40 additions & 0 deletions core/connectors/sources/http_source/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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]
Comment thread
hubcio marked this conversation as resolved.
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() {
Comment thread
hubcio marked this conversation as resolved.
let mut config = config(free_port(), free_port(), &[ENDPOINT_ONE]);
Expand Down
42 changes: 42 additions & 0 deletions core/integration/tests/connectors/http/http_source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading