Skip to content

Commit 2e30e53

Browse files
apollo_http_server: accept compressed HTTP requests (gzip, zstd, brotli) (#13366)
1 parent 6bf0fde commit 2e30e53

5 files changed

Lines changed: 91 additions & 0 deletions

File tree

Cargo.lock

Lines changed: 3 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

crates/apollo_http_server/Cargo.toml

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ serde_json.workspace = true
4242
starknet_api.workspace = true
4343
thiserror.workspace = true
4444
tokio = { workspace = true, features = ["rt"] }
45+
tower-http = { workspace = true, features = ["decompression-full", "limit"] }
4546
tracing.workspace = true
4647

4748

@@ -62,3 +63,4 @@ rstest.workspace = true
6263
serde_json.workspace = true
6364
starknet-types-core.workspace = true
6465
tracing-test.workspace = true
66+
zstd.workspace = true

crates/apollo_http_server/src/http_server.rs

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,8 @@ use starknet_api::transaction::fields::ValidResourceBounds;
3535
use tokio::net::TcpListener;
3636
use tokio::sync::watch::{channel, Receiver, Sender};
3737
use tokio::time;
38+
use tower_http::decompression::RequestDecompressionLayer;
39+
use tower_http::limit::RequestBodyLimitLayer;
3840
use tracing::{debug, info, instrument, warn};
3941

4042
use crate::deprecated_gateway_transaction::DeprecatedGatewayTransactionV3;
@@ -142,6 +144,13 @@ impl HttpServer {
142144
get(|| futures::future::ready("Gateway is ready".to_owned()))
143145
)
144146
.layer(Extension(self.app_state.clone()))
147+
// Hard streaming limit on decompressed bytes — wraps the body in
148+
// http_body_util::Limited which errors during poll_frame() once the
149+
// limit is exceeded, preventing zip bombs from expanding in memory.
150+
.layer(RequestBodyLimitLayer::new(self.config.static_config.max_request_body_size))
151+
.layer(RequestDecompressionLayer::new())
152+
// Cap compressed wire bytes to bound network I/O.
153+
.layer(RequestBodyLimitLayer::new(self.config.static_config.max_request_body_size))
145154
}
146155

147156
fn post_method_router<H, T, S>(&self, handler: H) -> MethodRouter<S>

crates/apollo_http_server/src/http_server_test.rs

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
use std::io::Write;
2+
13
use apollo_gateway_types::communication::{GatewayClientError, MockGatewayClient};
24
use apollo_gateway_types::deprecated_gateway_error::{
35
KnownStarknetErrorCode,
@@ -385,3 +387,67 @@ async fn request_body_size_limit_enforced(
385387
let response = http_client.add_tx(tx).await;
386388
assert_eq!(response.status(), expected_status, "Unexpected status: {}", response.status());
387389
}
390+
391+
#[tokio::test]
392+
async fn zstd_compressed_request_decompression() {
393+
let mut mock_gateway_client = MockGatewayClient::new();
394+
mock_gateway_client.expect_add_tx().times(1).return_const(Ok(default_gateway_output()));
395+
396+
let http_client = HttpClientServerSetupBuilder::new(unique_u16!())
397+
.with_mock_gateway_client(mock_gateway_client)
398+
.build()
399+
.await;
400+
401+
let tx_json = serde_json::to_string(&rpc_invoke_tx()).unwrap();
402+
let mut encoder = zstd::stream::write::Encoder::new(Vec::new(), 0).unwrap();
403+
encoder.write_all(tx_json.as_bytes()).unwrap();
404+
let compressed_body = encoder.finish().unwrap();
405+
406+
let response = http_client.add_rpc_tx_with_zstd(compressed_body).await;
407+
408+
assert_eq!(response.status(), StatusCode::OK, "Request should be decompressed and handled");
409+
let response_body = response.text().await.unwrap();
410+
let gateway_output: GatewayOutput =
411+
serde_json::from_str(&response_body).expect("Response should be valid GatewayOutput");
412+
assert_eq!(gateway_output.transaction_hash(), EXPECTED_TX_HASH);
413+
}
414+
415+
#[tokio::test]
416+
async fn zstd_compressed_request_too_large() {
417+
let tx_json = serde_json::to_string(&rpc_invoke_tx()).unwrap();
418+
let mut encoder = zstd::stream::write::Encoder::new(Vec::new(), 0).unwrap();
419+
encoder.write_all(tx_json.as_bytes()).unwrap();
420+
let compressed_body = encoder.finish().unwrap();
421+
422+
let http_client = HttpClientServerSetupBuilder::new(unique_u16!())
423+
.with_max_request_body_size(compressed_body.len() - 1)
424+
.build()
425+
.await;
426+
427+
let response = http_client.add_rpc_tx_with_zstd(compressed_body).await;
428+
429+
assert_eq!(response.status(), StatusCode::PAYLOAD_TOO_LARGE);
430+
}
431+
432+
#[tokio::test]
433+
async fn zstd_decompressed_request_too_large() {
434+
// 10 KB of repeated bytes — compresses to ~50 bytes with zstd.
435+
let large_body = vec![b'a'; 10 * 1024];
436+
let mut encoder = zstd::stream::write::Encoder::new(Vec::new(), 0).unwrap();
437+
encoder.write_all(&large_body).unwrap();
438+
let compressed_body = encoder.finish().unwrap();
439+
440+
// Limit between compressed size and decompressed size.
441+
// compressed_body is ~50 bytes; decompressed is 10240 bytes.
442+
let max_request_body_size = large_body.len() - 1;
443+
assert!(compressed_body.len() < max_request_body_size);
444+
445+
let http_client = HttpClientServerSetupBuilder::new(unique_u16!())
446+
.with_max_request_body_size(max_request_body_size)
447+
.build()
448+
.await;
449+
450+
let response = http_client.add_rpc_tx_with_zstd(compressed_body).await;
451+
452+
assert_eq!(response.status(), StatusCode::PAYLOAD_TOO_LARGE);
453+
}

crates/apollo_http_server/src/test_utils.rs

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,17 @@ impl HttpTestClient {
9494
.await
9595
.unwrap()
9696
}
97+
98+
pub async fn add_rpc_tx_with_zstd(&self, compressed_body: Vec<u8>) -> Response {
99+
self.client
100+
.post(format!("http://{}/gateway/add_rpc_transaction", self.socket))
101+
.header("content-type", "application/json")
102+
.header("content-encoding", "zstd")
103+
.body(Body::from(compressed_body))
104+
.send()
105+
.await
106+
.unwrap()
107+
}
97108
}
98109

99110
pub fn create_http_server_config(socket: SocketAddr) -> HttpServerConfig {

0 commit comments

Comments
 (0)