This is an automated email from the ASF dual-hosted git repository.
numinnex pushed a commit to branch kafka_proxy_auth
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/kafka_proxy_auth by this push:
new 927b492ec address review
927b492ec is described below
commit 927b492ecc1081ab3a6728c9b68b44d52256c91d
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Wed Sep 23 11:21:48 2026 +0200
address review
---
gateways/kafka/README.md | 15 +-
gateways/kafka/docs/AUTHENTICATION.md | 35 +--
gateways/kafka/src/auth.rs | 187 ++++++++++++++-
gateways/kafka/src/main.rs | 56 ++++-
gateways/kafka/src/protocol/api.rs | 7 +-
.../kafka/src/protocol/handlers/api_versions.rs | 3 +-
gateways/kafka/src/server.rs | 255 +++++++++++++++++----
gateways/kafka/tests/common/wire.rs | 1 -
gateways/kafka/tests/sasl_tests.rs | 127 +++++++++-
gateways/kafka/tests/version_firewall_tests.rs | 2 +-
10 files changed, 606 insertions(+), 82 deletions(-)
diff --git a/gateways/kafka/README.md b/gateways/kafka/README.md
index f890bf44b..ea776a38b 100644
--- a/gateways/kafka/README.md
+++ b/gateways/kafka/README.md
@@ -25,8 +25,8 @@ Default bind: `127.0.0.1:9093`. Environment variables:
| `IGGY_KAFKA_SHUTDOWN_DRAIN_TIMEOUT_SECS` | `25` | Seconds graceful shutdown
waits for in-flight connections before abandoning them |
| `IGGY_KAFKA_BRIDGE_ENABLED` | `false` | Connect the Iggy bridge at startup.
While false every API answers with its stub, and the `IGGY_KAFKA_IGGY_*`
variables below are read by nothing. A failed connection is fatal, not a
downgrade to stubs. |
| `IGGY_KAFKA_SASL_ENABLED` | `false` | Require SASL/PLAIN authentication
before serving any other API (`true` or `false`, nothing else) |
-| `IGGY_KAFKA_PRE_AUTH_TIMEOUT_SECS` | `15` | Seconds an unauthenticated
connection may sit between frames, and the ceiling on waiting for an
authentication slot plus the verification itself. Separate from the 10-minute
idle timeout that applies once authenticated |
-| `IGGY_KAFKA_MAX_CONCURRENT_AUTHENTICATIONS` | `4` | Credential verifications
allowed to run at once, across all connections. Each costs a password hash on
an Iggy shard thread, so this bounds what unauthenticated traffic can demand of
the server. Size it below the Iggy node's shard count |
+| `IGGY_KAFKA_PRE_AUTH_TIMEOUT_SECS` | `15` | Seconds an unauthenticated
connection may sit between frames. Waiting for an authentication slot and the
verification itself each get this budget, the verification's starting once it
holds a slot. Separate from the 10-minute idle timeout that applies once
authenticated |
+| `IGGY_KAFKA_MAX_CONCURRENT_AUTHENTICATIONS` | `4` | Credential verifications
the gateway runs at once, across all connections. Each costs a password hash on
an Iggy shard thread. This bounds the gateway's side only: a check that times
out frees its slot while its hash keeps running inside Iggy. Size it below the
Iggy node's shard count |
## Test
@@ -94,11 +94,13 @@ IGGY_KAFKA_SASL_ENABLED=true
IGGY_KAFKA_IGGY_ADDR=127.0.0.1:8090 cargo run -p ig
```
Transport security to Iggy is configured separately from the Kafka side,
because the two protect
-different hops:
+different hops. These variables cover only the connection the credential check
makes, and are
+refused while SASL is off. The bridge's own client connects without TLS at any
setting, and that
+connection carries `IGGY_KAFKA_IGGY_PASSWORD` in the clear:
| Variable | Default | Description |
| --- | --- | --- |
-| `IGGY_KAFKA_IGGY_TLS_ENABLED` | `false` | Encrypt the gateway's link to Iggy
(`true` or `false`, nothing else). Required if the Iggy server only accepts
TLS, otherwise every verification fails as unreachable |
+| `IGGY_KAFKA_IGGY_TLS_ENABLED` | `false` | Encrypt the credential check's
connection to Iggy (`true` or `false`, nothing else). The bridge connection is
not covered. Required if the Iggy server only accepts TLS, otherwise every
verification fails as unreachable |
| `IGGY_KAFKA_IGGY_TLS_DOMAIN` | derived from the address | Name checked
against the Iggy server certificate |
| `IGGY_KAFKA_IGGY_TLS_CA_FILE` | SDK bundled roots | PEM roots to trust. Note
the SDK does not use the system trust store |
@@ -119,8 +121,9 @@ Four things to know before switching it on:
`ApiVersions` advertisement only while it is on, and unauthenticated clients
are refused.
- **Every connection costs a login**, meaning one password hash on an Iggy
shard thread and one
replicated registration. Verification is deliberately not cached, since
caching it per username
- would let a second connection present any password. Connection churn is
therefore server load,
- bounded by `IGGY_KAFKA_MAX_CONCURRENT_AUTHENTICATIONS`.
+ would let a second connection present any password. Connection churn is
therefore server load.
+ `IGGY_KAFKA_MAX_CONCURRENT_AUTHENTICATIONS` bounds the checks the gateway
runs at once, and a peer
+ whose login was rejected is refused for a delay that doubles per rejection,
from 0.5s up to 30s.
- **Authentication only, for now.** The gateway verifies the credentials and
then drops the
session, because no handler consumes one yet. Iggy's permissions will decide
what a principal can
do once Produce and Fetch are wired to it
diff --git a/gateways/kafka/docs/AUTHENTICATION.md
b/gateways/kafka/docs/AUTHENTICATION.md
index 1c79f1b30..43aaacfa5 100644
--- a/gateways/kafka/docs/AUTHENTICATION.md
+++ b/gateways/kafka/docs/AUTHENTICATION.md
@@ -21,8 +21,10 @@ the only place a credential is defined or verified.
verified on every connection, against Iggy, and that verification is
deliberately not cached. See
[Verification is per
connection](#verification-is-per-connection-and-cannot-be-cached).
-**Authorization stays in Iggy.** The gateway does not implement its own ACL
model. It authenticates a Kafka
-client into an Iggy identity and lets the server's existing permission checks
decide what that identity can do.
+**Authorization is meant to stay in Iggy.** The gateway will not implement its
own ACL model. The intent is
+to carry the authenticated Iggy identity onto the data plane and let the
server's existing permission checks
+decide what it can do. That is later work: today the verified client is shut
down as soon as the credential
+clears, and nothing carries the principal past authentication.
## Why SCRAM is out
@@ -165,24 +167,29 @@ Notes that decide the implementation.
- **Produce with `acks=0` stays silent.** Answering an unauthenticated
fire-and-forget produce desyncs the
client's correlation stream, so that case closes without writing. The
existing rationale in
`protocol/api.rs` applies unchanged.
-- **Pre-authentication deadline.** An unauthenticated connection currently
holds a `max_connections` permit
- for up to the idle timeout, ten minutes by default. Authentication needs its
own, much shorter deadline.
+- **Pre-authentication deadline.** An unauthenticated connection holds a
`max_connections` permit, so it
+ gets `IGGY_KAFKA_PRE_AUTH_TIMEOUT_SECS` between frames rather than the
ten-minute idle timeout.
## Error mapping
At authentication time a rejected credential becomes
`SASL_AUTHENTICATION_FAILED` (58) with a generic
-message. An Iggy that cannot be reached, or an overloaded gateway, closes the
connection without a body
-instead: Kafka clients treat 58 as fatal and raise it to the application, so
borrowing it for a transient
+message. An Iggy that cannot be reached, an overloaded gateway, or a peer
still inside the delay a previous
+rejection earned, closes the connection without a body instead: Kafka clients
treat 58 as fatal and raise it to the application, so borrowing it for a
transient
condition turns a blip into a permanent failure for credentials that were
always correct. A close reads as
a transport failure, which is retriable, and still says nothing about whether
the account exists. Iggy's login path already runs a dummy hash for unknown
users to avoid a user-enumeration oracle,
so the gateway must not reintroduce one by distinguishing unknown user from
wrong password in the message
or by returning early.
+### Planned: authorization errors
+
+Not implemented. Today `bridge/error.rs` picks the Kafka code from the Iggy
error, and neither 30 nor 31
+is produced anywhere.
+
After authentication, Iggy reports exactly one permission-denied code,
`IggyError::Unauthorized` (41,
`core/common/src/error/iggy_error.rs:97`), alongside `Unauthenticated` (40).
Kafka distinguishes
`TOPIC_AUTHORIZATION_FAILED` (29), `GROUP_AUTHORIZATION_FAILED` (30) and
`CLUSTER_AUTHORIZATION_FAILED`
-(31). The gateway therefore picks the Kafka code from the operation it was
performing, not from the Iggy
-error, because the Iggy error cannot tell them apart.
+(31). Once handlers act as the authenticated principal, the gateway will have
to pick the Kafka code from
+the operation it was performing, not from the Iggy error, because the Iggy
error cannot tell them apart.
Iggy's data-plane permission checks read the local shard's view, so a
permission revocation is visible on
the control plane immediately and on the data plane only after that shard
applies it. Say so in the README
@@ -198,9 +205,10 @@ follow-up. The Iggy side already supports it on both ends
(client at
`core/configs/src/server_config/tcp.rs:31`), shipped disabled.
The two hops are configured independently. `IGGY_KAFKA_IGGY_TLS_ENABLED` and
its companions encrypt
-the gateway's link to Iggy and are what make the gateway usable at all against
a TLS-only Iggy
-server, where every verification would otherwise fail as unreachable. The
Kafka-side listener is the
-half that is still missing.
+only the connection the credential check makes to Iggy, and are what make SASL
usable at all against a
+TLS-only Iggy server, where every verification would otherwise fail as
unreachable. The bridge's own
+client has no TLS at any setting, so the bridge hop stays unencrypted and
carries
+`IGGY_KAFKA_IGGY_PASSWORD` in the clear. The Kafka-side listener is the other
half still missing.
Two limits worth recording. There is no mTLS anywhere in the tree, every
rustls config uses
`with_no_client_auth()`, so certificate-based Kafka client authentication
cannot map to an Iggy identity
@@ -224,8 +232,9 @@ has to resolve it first.
- SCRAM-SHA-256 and SCRAM-SHA-512, blocked on credential storage that does not
exist.
- SASL/OAUTHBEARER and GSSAPI.
- mTLS and certificate-based identity.
-- Mapping Kafka ACL administration APIs onto Iggy permissions. Authorization
is enforced, but the
- `DescribeAcls` and `CreateAcls` API keys stay unimplemented.
+- Authorization. No handler asks Iggy about permissions yet, since the
verified identity is not carried
+ onto the data plane; SASL is an admission gate only. The `DescribeAcls` and
`CreateAcls` API keys also
+ stay unimplemented.
- KIP-368 re-authentication.
## Open questions
diff --git a/gateways/kafka/src/auth.rs b/gateways/kafka/src/auth.rs
index 9d57c2ece..5a45c2995 100644
--- a/gateways/kafka/src/auth.rs
+++ b/gateways/kafka/src/auth.rs
@@ -21,7 +21,10 @@
//! password are an Iggy username and password, so authenticating is
forwarding them to an Iggy
//! login and seeing whether it succeeds. `docs/AUTHENTICATION.md` has the
reasoning.
-use std::time::Duration;
+use std::collections::HashMap;
+use std::net::{IpAddr, Ipv6Addr};
+use std::sync::{Mutex, PoisonError};
+use std::time::{Duration, Instant};
use async_trait::async_trait;
use iggy::prelude::{AutoLogin, Client, Credentials, IggyClientBuilder,
IggyError};
@@ -51,9 +54,11 @@ const VERIFY_TIMEOUT: Duration = Duration::from_secs(10);
/// Retries the dial makes before giving up, not the SDK's unlimited default.
///
-/// A Kafka client retries the whole authentication itself, so an unbounded
inner loop would only
-/// hide the failure underneath one the client cannot see.
-const VERIFY_RECONNECTION_RETRIES: u32 = 1;
+/// The SDK counts passes after the first, so zero still makes one full pass
over the endpoints.
+/// Any retry adds a second pass plus a `reconnection.interval` sleep while an
authentication slot
+/// is held, and a Kafka client retries the whole authentication itself
anyway, so an inner loop
+/// would only hide the failure underneath one the client cannot see.
+const VERIFY_RECONNECTION_RETRIES: u32 = 0;
/// Budget for tearing the verification client down again.
///
@@ -63,6 +68,113 @@ const VERIFY_RECONNECTION_RETRIES: u32 = 1;
/// Nothing is lost by cutting it short: `Drop` aborts the heartbeat task
regardless.
const TEARDOWN_TIMEOUT: Duration = Duration::from_secs(1);
+/// Delay imposed on a peer after its first rejected login, doubled on every
further rejection.
+const THROTTLE_BASE_DELAY: Duration = Duration::from_millis(500);
+
+/// Ceiling on the doubling, and how long a peer stays remembered once its
delay has run out.
+const THROTTLE_MAX_DELAY: Duration = Duration::from_secs(30);
+
+/// Peers remembered at once.
+///
+/// Bounds the table against a sweep of source addresses. Once it is full of
peers still inside
+/// their delay, new peers go untracked, which is no worse than having no
throttle at all.
+const THROTTLE_MAX_PEERS: usize = 4096;
+
+/// Refuses verifications from a peer whose last login was rejected, for an
escalating delay.
+///
+/// Every guess costs a full Argon2id hash on an Iggy shard thread, and
reconnecting is free, so
+/// without this a single peer can keep every authentication slot busy with
wrong passwords.
+/// Keyed on the IP rather than the socket address because the port changes on
every reconnect,
+/// and on the /64 for IPv6 because a single host is routinely handed a whole
/64 to draw from.
+/// Peers sharing a NAT share a delay, which is why it starts short.
+#[derive(Debug, Default)]
+pub struct FailedLoginThrottle {
+ peers: Mutex<HashMap<IpAddr, Strikes>>,
+}
+
+#[derive(Debug, Clone, Copy)]
+struct Strikes {
+ rejections: u32,
+ blocked_until: Instant,
+}
+
+impl Strikes {
+ fn is_forgotten(&self, now: Instant) -> bool {
+ now >= self.blocked_until + THROTTLE_MAX_DELAY
+ }
+}
+
+impl FailedLoginThrottle {
+ /// Whether `peer` is still inside the delay its last rejection earned.
+ #[must_use]
+ pub fn is_blocked(&self, peer: IpAddr) -> bool {
+ self.is_blocked_at(peer, Instant::now())
+ }
+
+ /// Records a rejected login from `peer` and extends its delay.
+ pub fn record_rejection(&self, peer: IpAddr) {
+ self.record_rejection_at(peer, Instant::now());
+ }
+
+ /// Forgets `peer` once it has proven a credential.
+ pub fn record_success(&self, peer: IpAddr) {
+ self.peers
+ .lock()
+ .unwrap_or_else(PoisonError::into_inner)
+ .remove(&throttle_key(peer));
+ }
+
+ fn is_blocked_at(&self, peer: IpAddr, now: Instant) -> bool {
+ let peers = self.peers.lock().unwrap_or_else(PoisonError::into_inner);
+ peers
+ .get(&throttle_key(peer))
+ .is_some_and(|strikes| now < strikes.blocked_until)
+ }
+
+ fn record_rejection_at(&self, peer: IpAddr, now: Instant) {
+ let peer = throttle_key(peer);
+ let mut peers =
self.peers.lock().unwrap_or_else(PoisonError::into_inner);
+ if !peers.contains_key(&peer) && peers.len() >= THROTTLE_MAX_PEERS {
+ peers.retain(|_, strikes| !strikes.is_forgotten(now));
+ if peers.len() >= THROTTLE_MAX_PEERS {
+ return;
+ }
+ }
+ let rejections = peers
+ .get(&peer)
+ .filter(|strikes| !strikes.is_forgotten(now))
+ .map_or(1, |strikes| strikes.rejections.saturating_add(1));
+ let doublings = (rejections - 1).min(16);
+ let delay = THROTTLE_BASE_DELAY
+ .saturating_mul(1 << doublings)
+ .min(THROTTLE_MAX_DELAY);
+ peers.insert(
+ peer,
+ Strikes {
+ rejections,
+ blocked_until: now + delay,
+ },
+ );
+ }
+}
+
+/// The address a peer is throttled under: IPv4 as is, IPv6 truncated to its
/64.
+///
+/// An IPv4-mapped IPv6 address is unwrapped first, so a dual-stack listener
throttles the same
+/// client under the same key whichever socket family it arrived on.
+fn throttle_key(peer: IpAddr) -> IpAddr {
+ match peer {
+ IpAddr::V4(_) => peer,
+ IpAddr::V6(v6) => v6.to_ipv4_mapped().map_or_else(
+ || {
+ let prefix = v6.to_bits() & !u128::from(u64::MAX);
+ IpAddr::V6(Ipv6Addr::from_bits(prefix))
+ },
+ IpAddr::V4,
+ ),
+ }
+}
+
/// Why a SASL exchange did not produce a verified identity.
///
/// Both variants reach the client as the same `SASL_AUTHENTICATION_FAILED`
with the same generic
@@ -413,6 +525,73 @@ mod tests {
assert_eq!(authenticator.to_string(), "iggy.internal:8090 over TLS");
}
+ #[test]
+ fn
given_a_rejected_peer_when_throttled_should_block_until_the_delay_runs_out() {
+ let throttle = FailedLoginThrottle::default();
+ let peer: IpAddr = [192, 0, 2, 1].into();
+ let other: IpAddr = [192, 0, 2, 2].into();
+ let start = Instant::now();
+
+ throttle.record_rejection_at(peer, start);
+ assert!(throttle.is_blocked_at(peer, start));
+ assert!(
+ !throttle.is_blocked_at(other, start),
+ "the delay is per peer"
+ );
+ assert!(!throttle.is_blocked_at(peer, start + THROTTLE_BASE_DELAY));
+
+ // A second rejection before the peer is forgotten doubles the delay.
+ let second = start + THROTTLE_BASE_DELAY;
+ throttle.record_rejection_at(peer, second);
+ assert!(throttle.is_blocked_at(peer, second + THROTTLE_BASE_DELAY));
+ assert!(!throttle.is_blocked_at(peer, second + THROTTLE_BASE_DELAY *
2));
+ }
+
+ #[test]
+ fn given_a_throttled_peer_when_it_authenticates_should_be_forgotten() {
+ let throttle = FailedLoginThrottle::default();
+ let peer: IpAddr = [192, 0, 2, 1].into();
+ let start = Instant::now();
+ throttle.record_rejection_at(peer, start);
+ throttle.record_rejection_at(peer, start);
+ throttle.record_success(peer);
+ assert!(!throttle.is_blocked_at(peer, start));
+
+ // The escalation restarts from the base delay rather than where it
left off.
+ throttle.record_rejection_at(peer, start);
+ assert!(!throttle.is_blocked_at(peer, start + THROTTLE_BASE_DELAY));
+ }
+
+ #[test]
+ fn given_an_ipv6_peer_should_share_a_delay_across_its_slash_64() {
+ let throttle = FailedLoginThrottle::default();
+ let start = Instant::now();
+ let first: IpAddr = "2001:db8:1:2::1".parse().expect("valid address");
+ let same_prefix: IpAddr = "2001:db8:1:2:ffff::9".parse().expect("valid
address");
+ let other_prefix: IpAddr = "2001:db8:1:3::1".parse().expect("valid
address");
+
+ throttle.record_rejection_at(first, start);
+ assert!(throttle.is_blocked_at(same_prefix, start));
+ assert!(!throttle.is_blocked_at(other_prefix, start));
+
+ let mapped: IpAddr = "::ffff:192.0.2.7".parse().expect("valid
address");
+ let plain: IpAddr = [192, 0, 2, 7].into();
+ throttle.record_rejection_at(mapped, start);
+ assert!(throttle.is_blocked_at(plain, start));
+ }
+
+ #[test]
+ fn given_repeated_rejections_should_cap_the_delay() {
+ let throttle = FailedLoginThrottle::default();
+ let peer: IpAddr = [192, 0, 2, 1].into();
+ let start = Instant::now();
+ for _ in 0..64 {
+ throttle.record_rejection_at(peer, start);
+ }
+ assert!(throttle.is_blocked_at(peer, start + THROTTLE_MAX_DELAY / 2));
+ assert!(!throttle.is_blocked_at(peer, start + THROTTLE_MAX_DELAY));
+ }
+
#[tokio::test]
async fn
given_an_unreachable_iggy_when_authenticating_should_report_unavailable() {
// Port 1 on loopback refuses immediately, so this exercises the
failure path without
diff --git a/gateways/kafka/src/main.rs b/gateways/kafka/src/main.rs
index 1b97b988e..3b92b4875 100644
--- a/gateways/kafka/src/main.rs
+++ b/gateways/kafka/src/main.rs
@@ -55,8 +55,9 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
);
if !authenticator.is_tls_enabled() {
warn!(
- "the link from this gateway to Iggy is also unencrypted; set \
- IGGY_KAFKA_IGGY_TLS_ENABLED=true to close the second hop"
+ "the credential check's link to Iggy is also unencrypted; set \
+ IGGY_KAFKA_IGGY_TLS_ENABLED=true to encrypt it. The bridge's
own connection has \
+ no TLS at any setting"
);
}
info!("SASL/PLAIN enabled; credentials verified against
{authenticator}");
@@ -243,10 +244,36 @@ fn load_config() -> Result<GatewayConfig, String> {
.map_err(|e| format!("invalid
IGGY_KAFKA_SHUTDOWN_DRAIN_TIMEOUT_SECS `{raw}`: {e}"))?;
config.shutdown_drain_timeout = Duration::from_secs(secs);
}
+ reject_iggy_tls_without_sasl(config.sasl_enabled)?;
Ok(config)
}
+/// Refuses `IGGY_KAFKA_IGGY_TLS_*` while SASL is off.
+///
+/// Only the credential verifier reads them, and it is built only with SASL
on. Accepting them
+/// otherwise lets an operator believe the link to Iggy is encrypted when
nothing reads the switch,
+/// the same mistake `IggyAuthenticator::from_env` refuses to start for.
+fn reject_iggy_tls_without_sasl(sasl_enabled: bool) -> Result<(), String> {
+ if sasl_enabled {
+ return Ok(());
+ }
+ let set: Vec<&str> = IggyAuthenticator::KNOWN_ENV_VARS
+ .iter()
+ .copied()
+ .filter(|var| var.starts_with("IGGY_KAFKA_IGGY_TLS_"))
+ .filter(|var| std::env::var(var).is_ok())
+ .collect();
+ if set.is_empty() {
+ return Ok(());
+ }
+ Err(format!(
+ "{} set but IGGY_KAFKA_SASL_ENABLED is not true; only the SASL
credential check reads them, \
+ so nothing would be encrypted",
+ set.join(", ")
+ ))
+}
+
fn env_var(key: &str) -> Option<String> {
std::env::var(key).ok()
}
@@ -296,11 +323,11 @@ async fn shutdown_signal() {
mod tests {
use serial_test::serial;
- use super::{parse_positive, reject_unknown_kafka_env_vars};
+ use super::{parse_positive, reject_iggy_tls_without_sasl,
reject_unknown_kafka_env_vars};
/// Sequential (not two separate `#[test]` fns), and `#[serial]` (unkeyed
- this binary's
- /// default group). This is the only `#[serial]` test compiled into *this*
binary
- /// (`main.rs` -> the `iggy-gateway-kafka` bin's own test harness) -
`auth`'s,
+ /// default group). The `#[serial]` tests in this module are the only ones
compiled into *this*
+ /// binary (`main.rs` -> the `iggy-gateway-kafka` bin's own test harness)
- `auth`'s,
/// `bridge::config`'s and `server`'s env-touching tests compile into the
separate lib test
/// binary, and `serial_test`'s mutex is process-local, so it does not
(and does not need to)
/// coordinate with any of those; `server.rs`'s own `#[serial]` test makes
the mirror-image
@@ -355,6 +382,25 @@ mod tests {
);
}
+ /// `#[serial]` for the same reason as the test above: it mutates
process-wide env state.
+ #[test]
+ #[serial]
+ fn given_iggy_tls_without_sasl_should_refuse_to_start() {
+ unsafe {
+ std::env::set_var("IGGY_KAFKA_IGGY_TLS_ENABLED", "true");
+ }
+ let without_sasl = reject_iggy_tls_without_sasl(false);
+ let with_sasl = reject_iggy_tls_without_sasl(true);
+ unsafe {
+ std::env::remove_var("IGGY_KAFKA_IGGY_TLS_ENABLED");
+ }
+ assert!(
+ without_sasl.is_err(),
+ "nothing reads the TLS switch with SASL off, so accepting it hides
that"
+ );
+ assert!(with_sasl.is_ok());
+ }
+
#[test]
fn parse_positive_rejects_zero() {
assert!(parse_positive::<usize>("KEY", "0").is_err());
diff --git a/gateways/kafka/src/protocol/api.rs
b/gateways/kafka/src/protocol/api.rs
index 907de22f9..e49a3e1a3 100644
--- a/gateways/kafka/src/protocol/api.rs
+++ b/gateways/kafka/src/protocol/api.rs
@@ -286,9 +286,10 @@ pub const fn advertised_min_version(api_key: i16,
firewall_min: i16) -> i16 {
/// Advertised only when SASL is switched on, and deliberately absent from
[`SUPPORTED_RANGES`].
///
/// These two keys never reach [`handle_request_bounded`]: the connection loop
routes them through
-/// the SASL state machine before dispatch. Keeping them out of the firewall
table means a gateway
-/// with SASL off treats them as any other unknown key and closes, which is
what stops enabling the
-/// feature later from silently widening what an unauthenticated client can
send today.
+/// the SASL state machine before dispatch, whether SASL is on or off. With it
off every connection
+/// starts authenticated, so both keys are answered `ILLEGAL_SASL_STATE` and
the connection stays
+/// open, the answer a real broker gives on a PLAINTEXT listener. Keeping them
out of the firewall
+/// table means dispatch never serves them on its own.
///
/// `SaslHandshake` is pinned to v1 on both ends. v0 selects the headerless
token framing (KIP-152)
/// that the frame reader cannot parse, so advertising it would invite exactly
the shape this
diff --git a/gateways/kafka/src/protocol/handlers/api_versions.rs
b/gateways/kafka/src/protocol/handlers/api_versions.rs
index f109ae341..4c3974665 100644
--- a/gateways/kafka/src/protocol/handlers/api_versions.rs
+++ b/gateways/kafka/src/protocol/handlers/api_versions.rs
@@ -73,7 +73,8 @@ pub async fn handle(state: &GatewayState, api_version: i16,
body: Bytes) -> Hand
/// Returns an error when `kafka_protocol` cannot encode the response at
`api_version`.
pub fn encode_response(api_version: i16, error_code: i16, sasl_enabled: bool)
-> Result<Bytes> {
// The SASL keys are advertised only while the feature is on, and are
deliberately kept out of
- // `SUPPORTED_RANGES` so a gateway with SASL off treats them as any other
unknown key.
+ // `SUPPORTED_RANGES` so dispatch never serves them. With SASL off the
state machine still
+ // answers them with `ILLEGAL_SASL_STATE` and keeps the connection open.
let sasl = if sasl_enabled {
sasl_advertised_ranges()
} else {
diff --git a/gateways/kafka/src/server.rs b/gateways/kafka/src/server.rs
index a2adb306b..8859f0f80 100644
--- a/gateways/kafka/src/server.rs
+++ b/gateways/kafka/src/server.rs
@@ -30,8 +30,10 @@ use tokio_util::sync::CancellationToken;
use tokio_util::task::TaskTracker;
use tracing::{debug, error, info, warn};
use tracing_appender::non_blocking::WorkerGuard;
+use tracing_subscriber::EnvFilter;
+use tracing_subscriber::filter::LevelFilter;
-use crate::auth::{AuthError, SaslAuthenticator};
+use crate::auth::{AuthError, FailedLoginThrottle, SaslAuthenticator};
use crate::bridge::IggyBridge;
use crate::error::{KafkaProtocolError, Result};
use crate::protocol::api::{
@@ -43,11 +45,13 @@ use crate::protocol::api::{
};
use crate::protocol::header::{request_header_version, response_header_version};
use crate::protocol::sasl::{
- SASL_AUTHENTICATE_MAX_VERSION, SASL_HANDSHAKE_VERSION, SaslAction,
SaslState, parse_plain,
+ PlainCredentials, SASL_AUTHENTICATE_MAX_VERSION, SASL_HANDSHAKE_VERSION,
SaslAction, SaslState,
+ parse_plain,
};
use std::io;
const READ_CHUNK: usize = 65536;
+const GATEWAY_LOG_TARGET: &str = env!("CARGO_CRATE_NAME");
/// Builds the log filter, forcing the Iggy SDK quiet unless the operator
asked otherwise.
///
@@ -61,17 +65,33 @@ const READ_CHUNK: usize = 65536;
/// and using it verbatim silently drops this the moment anyone sets it,
including on the run
/// commands this repository's own documentation gives. An explicit `iggy=`
directive still wins,
/// so raising it deliberately for debugging remains possible.
+///
+/// `EnvFilter` matches targets by prefix, so `iggy=warn` also catches this
crate. The gateway gets
+/// its own directive at the operator's default level, since the longer target
wins, and without it
+/// every authentication decision this gateway logs would be filtered out.
fn sdk_quieted_filter(rust_log: Option<&str>) -> String {
let base = rust_log.unwrap_or("info");
let base = if base.trim().is_empty() { "info" } else { base };
- let names_sdk = base
- .split(',')
- .any(|directive| directive.trim().starts_with("iggy="));
- if names_sdk {
- base.to_string()
- } else {
- format!("{base},iggy=warn")
+ let directives: Vec<&str> = base.split(',').map(str::trim).collect();
+ if directives
+ .iter()
+ .any(|directive| directive.starts_with("iggy="))
+ {
+ return base.to_string();
}
+ let names_gateway = directives
+ .iter()
+ .any(|directive| directive.split(['=', '[']).next() ==
Some(GATEWAY_LOG_TARGET));
+ if names_gateway {
+ return format!("{base},iggy=warn");
+ }
+ let default_level = directives
+ .iter()
+ .rev()
+ .copied()
+ .find(|directive| directive.parse::<LevelFilter>().is_ok())
+ .unwrap_or("info");
+ format!("{base},iggy=warn,{GATEWAY_LOG_TARGET}={default_level}")
}
#[derive(Debug, Clone)]
@@ -114,14 +134,18 @@ pub struct GatewayConfig {
/// A connection that has proven nothing still holds a `max_connections`
permit, so it gets a
/// budget measured in seconds instead.
pub pre_auth_timeout: Duration,
- /// Credential verifications allowed to run at once, across every
connection.
+ /// Credential verifications the gateway runs at once, across every
connection.
///
/// Each one costs an Argon2id verify on an Iggy shard thread that has no
blocking pool, so
- /// unauthenticated traffic can otherwise saturate the server's request
loop with nothing but
+ /// unauthenticated traffic can otherwise pile that work onto the server
with nothing but
/// shape-valid tokens. `max_connections` alone is not a bound on that: it
caps sockets, not
/// the work each one can ask Iggy to do. Verifications beyond this queue
rather than fail,
/// since a rejected login is indistinguishable from a wrong password to
the client.
///
+ /// This bounds the gateway's side only. A verification that times out
frees its slot, but the
+ /// login it started keeps hashing inside Iggy, so a slow server can
briefly carry more than
+ /// this many.
+ ///
/// Four by default, which stays under the shard count of any node with 16
or fewer physical
/// cores. A default above that is no bound at all on the deployments most
likely to run one
/// gateway in front of one node; a larger node raises it deliberately.
@@ -237,6 +261,17 @@ pub struct KafkaGateway {
struct SharedAuth {
authenticator: Option<Arc<dyn SaslAuthenticator>>,
slots: Semaphore,
+ failed_logins: FailedLoginThrottle,
+}
+
+impl SharedAuth {
+ fn new(authenticator: Option<Arc<dyn SaslAuthenticator>>, max_concurrent:
usize) -> Self {
+ Self {
+ authenticator,
+ slots: Semaphore::new(max_concurrent),
+ failed_logins: FailedLoginThrottle::default(),
+ }
+ }
}
impl KafkaGateway {
@@ -314,10 +349,10 @@ impl KafkaGateway {
self.config.sasl_enabled,
));
- let shared_auth = Arc::new(SharedAuth {
- authenticator: self.authenticator.clone(),
- slots: Semaphore::new(self.config.max_concurrent_authentications),
- });
+ let shared_auth = Arc::new(SharedAuth::new(
+ self.authenticator.clone(),
+ self.config.max_concurrent_authentications,
+ ));
let tracker = TaskTracker::new();
let conn_limiter =
Arc::new(Semaphore::new(self.config.max_connections));
// Cancelled on shutdown so connection tasks exit instead of sitting
in idle waits
@@ -442,7 +477,9 @@ struct ConnectionContext<'a> {
state: &'a GatewayState,
authenticator: Option<&'a dyn SaslAuthenticator>,
auth_slots: &'a Semaphore,
+ failed_logins: &'a FailedLoginThrottle,
peer: &'a SocketAddr,
+ cancel: &'a CancellationToken,
}
/// Decides what one decoded frame earns, advancing `sasl_state` when it
authenticates.
@@ -571,7 +608,9 @@ async fn handle_connection(
state: &state,
authenticator: shared_auth.authenticator.as_deref(),
auth_slots: &shared_auth.slots,
+ failed_logins: &shared_auth.failed_logins,
peer: &peer,
+ cancel: &cancel,
};
loop {
@@ -616,7 +655,17 @@ async fn handle_connection(
// `RequestHeader::decode` advances `body` past the header fields it
consumed via
// `Buf::advance`, so `body` is already exactly the request payload.
let outcome = route_frame(&ctx, &mut sasl_state, &req, body).await;
- if dispatch_outcome(&mut stream, &peer, &config, &req, resp_hdr_ver,
outcome).await? {
+ if dispatch_outcome(
+ &mut stream,
+ &peer,
+ &config,
+ &req,
+ resp_hdr_ver,
+ outcome,
+ &cancel,
+ )
+ .await?
+ {
return Ok(());
}
}
@@ -650,37 +699,41 @@ async fn authenticate_token(
return failed();
};
- // The wait for a slot and the verification itself share one deadline, and
it is the same
- // pre-authentication budget every other unauthenticated read gets.
Without it the queue is the
- // bound: an unauthenticated connection would sit in `acquire()` for as
long as the backlog
- // takes to drain, holding a `max_connections` permit the whole time,
which is precisely the
- // invariant `pre_auth_timeout` is documented to enforce.
- let verified = tokio::time::timeout(ctx.config.pre_auth_timeout, async {
- // Acquire fails only once the semaphore is closed, which this gateway
never does.
- let Ok(_slot) = ctx.auth_slots.acquire().await else {
- error!(%peer, "authentication slots unavailable");
- return None;
- };
- Some(authenticator.authenticate(&credentials).await)
- })
- .await;
-
- let Ok(Some(result)) = verified else {
- // Overloaded or shutting down. Close rather than answer 58: a Kafka
client treats that
- // code as fatal and surfaces it to the application, and nothing here
says the credentials
- // were wrong. A close reads as a transport failure, which is
retriable.
- warn!(%peer, "authentication did not complete within the
pre-authentication budget");
+ // Checked before a slot is taken, so a throttled peer costs neither a
slot nor a hash. Closed
+ // rather than answered 58 for the same reason an overload is: the client
would treat 58 as
+ // fatal, and a correct password retried after the delay must still get
through.
+ if ctx.failed_logins.is_blocked(peer.ip()) {
+ debug!(%peer, "SASL authentication refused: peer is throttled after a
rejected login");
+ return HandleOutcome::Close;
+ }
+
+ // The connection loop only watches the shutdown token between frames, so
a verification in
+ // flight has to watch it here or a drain waits out the whole budget.
+ let verified = tokio::select! {
+ () = ctx.cancel.cancelled() => {
+ debug!(%peer, "SASL authentication abandoned by shutdown");
+ return HandleOutcome::Close;
+ }
+ verified = verify_within_budget(ctx, authenticator, &credentials) =>
verified,
+ };
+
+ let Some(result) = verified else {
+ // Overloaded or timed out. Close rather than answer 58: a Kafka
client treats that code as
+ // fatal and surfaces it to the application, and nothing here says the
credentials were
+ // wrong. A close reads as a transport failure, which is retriable.
return HandleOutcome::Close;
};
match result {
Ok(()) => {
debug!(%peer, "SASL authentication succeeded");
+ ctx.failed_logins.record_success(peer.ip());
sasl_authenticate_outcome(api_version, ERROR_NONE, false)
}
// A rejection is the client's problem and is terminal, so it earns a
parseable 58.
Err(AuthError::Rejected) => {
debug!(%peer, "SASL authentication rejected");
+ ctx.failed_logins.record_rejection(peer.ip());
failed()
}
// An outage is not. Kafka clients treat 58 as fatal and surface it to
the application, so
@@ -695,6 +748,37 @@ async fn authenticate_token(
}
}
+/// Waits for an authentication slot, then verifies, each within the
pre-authentication budget.
+///
+/// Two budgets rather than one shared deadline. The wait needs its own bound
so a queued
+/// connection cannot hold a `max_connections` permit for as long as the
backlog takes to drain,
+/// which is the invariant `pre_auth_timeout` is documented to enforce. The
verification's budget
+/// starts once the slot is held, because cutting a verification short does
not stop the login it
+/// started inside Iggy: a verification that inherited whatever the queue left
of a shared deadline
+/// would be abandoned mid-hash and hand its slot to the next one while Iggy
is still working.
+async fn verify_within_budget(
+ ctx: &ConnectionContext<'_>,
+ authenticator: &dyn SaslAuthenticator,
+ credentials: &PlainCredentials,
+) -> Option<std::result::Result<(), AuthError>> {
+ let peer = ctx.peer;
+ let budget = ctx.config.pre_auth_timeout;
+ let Ok(acquired) = timeout(budget, ctx.auth_slots.acquire()).await else {
+ warn!(%peer, "no authentication slot came free within the
pre-authentication budget");
+ return None;
+ };
+ // Acquire fails only once the semaphore is closed, which this gateway
never does.
+ let Ok(_slot) = acquired else {
+ error!(%peer, "authentication slots unavailable");
+ return None;
+ };
+ let Ok(result) = timeout(budget,
authenticator.authenticate(credentials)).await else {
+ warn!(%peer, "authentication did not complete within the
pre-authentication budget");
+ return None;
+ };
+ Some(result)
+}
+
/// Answers a request that is well-formed but not legal in this connection's
SASL state.
///
/// `keep_open` marks the one case a real broker does not treat as fatal: a
SASL request arriving
@@ -762,6 +846,7 @@ async fn dispatch_outcome(
req: &RequestHeader,
resp_hdr_ver: i16,
outcome: HandleOutcome,
+ cancel: &CancellationToken,
) -> Result<bool> {
match outcome {
HandleOutcome::NoResponse => {
@@ -811,11 +896,35 @@ async fn dispatch_outcome(
api_version = req.request_api_version,
"closing connection after responding"
);
+ close_gracefully(stream, config.write_timeout, cancel).await;
Ok(true)
}
}
}
+/// Sends FIN, then discards input until the peer closes, `budget` runs out or
shutdown begins.
+///
+/// Dropping a socket with unread bytes in its receive queue makes Linux send
RST instead of FIN,
+/// and an RST lets the peer discard the response this close follows before
its client reads it.
+/// A real broker closes an authentication failure gracefully for the same
reason.
+async fn close_gracefully(stream: &mut TcpStream, budget: Duration, cancel:
&CancellationToken) {
+ if stream.shutdown().await.is_err() {
+ return;
+ }
+ let mut discard = [0u8; 4096];
+ let drain_input = async {
+ while let Ok(read) = stream.read(&mut discard).await {
+ if read == 0 {
+ break;
+ }
+ }
+ };
+ tokio::select! {
+ () = cancel.cancelled() => {}
+ _ = timeout(budget, drain_input) => {}
+ }
+}
+
/// Response header size for a given header version: v0 is `correlation_id`
only (4 bytes); v1
/// adds an empty tagged-fields byte (5 bytes). Kept inline rather than through
/// `kafka_protocol::messages::ResponseHeader` - that type's `Encodable` impl
needs its own
@@ -931,15 +1040,22 @@ pub async fn read_frame(
/// Returns the [`WorkerGuard`]; it must be held for the lifetime of `main`
(dropping it stops the
/// worker thread and any buffered-but-unflushed log lines are lost) - see
`main.rs`.
pub fn init_tracing() -> WorkerGuard {
- let filter = tracing_subscriber::EnvFilter::new(sdk_quieted_filter(
- std::env::var("RUST_LOG").ok().as_deref(),
- ));
+ // `EnvFilter::new` drops a directive it cannot parse, and `iggy=warn`
always parses, so one
+ // typo in `RUST_LOG` would otherwise leave every other target with no
output at all.
+ let rust_log = std::env::var("RUST_LOG").ok();
+ let (filter, rejected) = match
EnvFilter::try_new(sdk_quieted_filter(rust_log.as_deref())) {
+ Ok(filter) => (filter, None),
+ Err(error) => (EnvFilter::new(sdk_quieted_filter(None)), Some(error)),
+ };
let (non_blocking_stdout, guard) =
tracing_appender::non_blocking(io::stdout());
let _ = tracing_subscriber::fmt()
.with_env_filter(filter)
.with_writer(non_blocking_stdout)
.try_init()
.map_err(|e| error!("failed to initialize tracing: {e}"));
+ if let Some(error) = rejected {
+ warn!(%error, "RUST_LOG is invalid; logging at the default level
instead");
+ }
guard
}
@@ -949,16 +1065,53 @@ mod tests {
#[test]
fn given_no_rust_log_should_quiet_the_sdk() {
- assert_eq!(sdk_quieted_filter(None), "info,iggy=warn");
- assert_eq!(sdk_quieted_filter(Some("")), "info,iggy=warn");
+ assert_eq!(
+ sdk_quieted_filter(None),
+ "info,iggy=warn,iggy_gateway_kafka=info"
+ );
+ assert_eq!(
+ sdk_quieted_filter(Some("")),
+ "info,iggy=warn,iggy_gateway_kafka=info"
+ );
}
#[test]
fn given_a_rust_log_should_still_quiet_the_sdk() {
// The whole point: reading RUST_LOG verbatim dropped the
credential-disclosure control on
// every run command this repository's own docs give.
- assert_eq!(sdk_quieted_filter(Some("info")), "info,iggy=warn");
- assert_eq!(sdk_quieted_filter(Some("debug")), "debug,iggy=warn");
+ assert_eq!(
+ sdk_quieted_filter(Some("info")),
+ "info,iggy=warn,iggy_gateway_kafka=info"
+ );
+ assert_eq!(
+ sdk_quieted_filter(Some("debug")),
+ "debug,iggy=warn,iggy_gateway_kafka=debug"
+ );
+ }
+
+ #[test]
+ fn given_the_quieted_filter_should_keep_gateway_logs_and_drop_sdk_info() {
+ // `iggy=warn` matches `iggy_gateway_kafka` by prefix, so the built
string alone cannot
+ // show that the gateway's own lines survive. Only the filter's
verdict can.
+ let filter =
EnvFilter::try_new(sdk_quieted_filter(Some("info"))).expect("valid filter");
+ let subscriber = tracing_subscriber::fmt()
+ .with_env_filter(filter)
+ .with_writer(io::sink)
+ .finish();
+ tracing::subscriber::with_default(subscriber, || {
+ assert!(tracing::enabled!(
+ target: "iggy_gateway_kafka::server",
+ tracing::Level::INFO
+ ));
+ assert!(!tracing::enabled!(
+ target: "iggy::clients::client",
+ tracing::Level::INFO
+ ));
+ assert!(tracing::enabled!(
+ target: "iggy::clients::client",
+ tracing::Level::WARN
+ ));
+ });
}
#[test]
@@ -971,6 +1124,16 @@ mod tests {
);
}
+ #[test]
+ fn given_a_bare_gateway_target_should_not_override_it() {
+ // A bare target enables every level for it. Appending
`iggy_gateway_kafka=info` would
+ // cap what the operator explicitly asked for.
+ assert_eq!(
+ sdk_quieted_filter(Some("iggy_gateway_kafka")),
+ "iggy_gateway_kafka,iggy=warn"
+ );
+ }
+
#[test]
fn given_an_explicit_sdk_directive_should_be_left_alone() {
assert_eq!(
@@ -1304,9 +1467,9 @@ mod tests {
/// `#[serial]`, unkeyed (shares `bridge::config`'s default group - both
this module and
/// `bridge::config` compile into the same lib unit-test binary.
`main.rs`'s own `#[serial]`
- /// test does NOT share this group: `main.rs` is the separate
`iggy-gateway-kafka` bin's own
- /// test harness, a different process, and `serial_test`'s mutex is
process-local - see that
- /// test's own doc comment for the mirror-image note): `init_tracing`
reads `RUST_LOG` with
+ /// tests do NOT share this group: `main.rs` is the separate
`iggy-gateway-kafka` bin's own
+ /// test harness, a different process, and `serial_test`'s mutex is
process-local - see the
+ /// first of those tests' doc comment for the mirror-image note):
`init_tracing` reads `RUST_LOG` with
/// `std::env::var`, and edition 2024's `env::set_var`/`remove_var` are
unsound against *any*
/// concurrent env read in another thread, not just a write to the same
key - a set/remove
/// elsewhere in this binary racing this read is exactly the hazard,
regardless of which var
diff --git a/gateways/kafka/tests/common/wire.rs
b/gateways/kafka/tests/common/wire.rs
index 4d57358bc..304604ed1 100644
--- a/gateways/kafka/tests/common/wire.rs
+++ b/gateways/kafka/tests/common/wire.rs
@@ -41,7 +41,6 @@ pub const OUT_OF_SCOPE_API_KEYS: &[(i16, &str)] = &[
(14, "SyncGroup"),
(15, "DescribeGroups"),
(16, "ListGroups"),
- (17, "SaslHandshake"),
(20, "DeleteTopics"),
];
diff --git a/gateways/kafka/tests/sasl_tests.rs
b/gateways/kafka/tests/sasl_tests.rs
index 5a15a8b99..903ece36d 100644
--- a/gateways/kafka/tests/sasl_tests.rs
+++ b/gateways/kafka/tests/sasl_tests.rs
@@ -26,7 +26,7 @@ use std::time::Duration;
use async_trait::async_trait;
use bytes::{BufMut, Bytes, BytesMut};
-use tokio::io::AsyncWriteExt;
+use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpStream;
use iggy_gateway_kafka::GatewayConfig;
@@ -244,6 +244,44 @@ async fn
given_wrong_credentials_when_authenticating_should_fail_then_close() {
assert_closed(&mut stream).await;
}
+#[tokio::test]
+async fn
given_pipelined_bytes_behind_a_rejected_token_should_close_with_a_fin_not_a_reset()
{
+ // Unread input at close makes Linux send RST instead of FIN, and an RST
lets the client's
+ // stack discard the 58 before the client reads it. Loopback delivers the
58 ahead of the RST
+ // either way, so what this pins is the close itself: a clean EOF, not a
reset.
+ let addr = spawn_sasl_gateway().await;
+ let mut stream = TcpStream::connect(addr).await.expect("connect");
+ handshake_ok(&mut stream).await;
+
+ let rejected = build_request_frame(
+ API_KEY_SASL_AUTHENTICATE,
+ AUTHENTICATE_VERSION,
+ 2,
+ Some("sasl-test"),
+ &authenticate_body(&plain_token("alice", "wrong-password")),
+ );
+ let pipelined = build_request_frame(API_KEY_METADATA, 0, 3,
Some("sasl-test"), &[0, 0, 0, 0]);
+ let mut both = BytesMut::from(&rejected[..]);
+ both.extend_from_slice(&pipelined);
+ stream.write_all(&both).await.expect("write requests");
+ // Let the server answer and close before the client reads anything.
+ tokio::time::sleep(Duration::from_millis(200)).await;
+
+ let payload = tcp::read_response_frame(&mut stream, 8 * 1024 * 1024).await;
+ let (echoed, body) =
+ parse_response_payload(API_KEY_SASL_AUTHENTICATE,
AUTHENTICATE_VERSION, payload);
+ assert_eq!(echoed, 2);
+ assert_eq!(error_code(&body), ERROR_SASL_AUTHENTICATION_FAILED);
+ let mut rest = [0u8; 1];
+ let read = tokio::time::timeout(Duration::from_secs(5), stream.read(&mut
rest))
+ .await
+ .expect("the server must close the connection");
+ assert!(
+ matches!(read, Ok(0)),
+ "the close must be a FIN the client reads as EOF, not a reset:
{read:?}"
+ );
+}
+
#[tokio::test]
async fn
given_a_malformed_plain_token_when_authenticating_should_fail_like_a_bad_password()
{
let addr = spawn_sasl_gateway().await;
@@ -875,7 +913,7 @@ impl SaslAuthenticator for StallingAuthenticator {
#[tokio::test]
async fn
given_all_authentication_slots_are_busy_when_waiting_too_long_should_close_not_reject()
{
- // The permit wait shares the pre-authentication budget. Without that
bound a connection sits
+ // The permit wait is bounded by the pre-authentication budget. Without
that a connection sits
// in the queue holding a `max_connections` permit for as long as the
backlog takes to drain,
// which is the invariant `pre_auth_timeout` is documented to enforce. It
must close rather
// than answer 58, which a Kafka client treats as fatal even though
nothing was rejected.
@@ -911,6 +949,91 @@ async fn
given_all_authentication_slots_are_busy_when_waiting_too_long_should_cl
assert_closed(&mut queued).await;
}
+#[tokio::test]
+async fn
given_a_verification_in_flight_when_shutting_down_should_close_without_waiting_it_out()
{
+ // The connection loop only watches the shutdown token between frames. A
verification that
+ // did not watch it too would hold the drain for the whole
pre-authentication budget.
+ let config = GatewayConfig {
+ pre_auth_timeout: Duration::from_secs(30),
+ shutdown_drain_timeout: Duration::from_secs(30),
+ ..sasl_config()
+ };
+ let authenticator = Arc::new(StallingAuthenticator {
+ release: tokio::sync::Semaphore::new(0),
+ });
+ let (addr, shutdown) = spawn_test_server_with_authenticator(config,
authenticator).await;
+
+ let mut stream = TcpStream::connect(addr).await.expect("connect");
+ handshake_ok(&mut stream).await;
+ let token = plain_token("alice", "s3cret");
+ let frame = build_request_frame(
+ API_KEY_SASL_AUTHENTICATE,
+ AUTHENTICATE_VERSION,
+ 2,
+ Some("sasl-test"),
+ &authenticate_body(&token),
+ );
+ stream.write_all(&frame).await.expect("write request");
+ tokio::time::sleep(Duration::from_millis(100)).await;
+
+ shutdown.send(()).expect("signal shutdown");
+ assert_eq!(
+ read_byte_with_timeout(&mut stream, Duration::from_secs(2)).await,
+ ByteRead::Closed,
+ "shutdown must end a verification in flight"
+ );
+}
+
+#[tokio::test]
+async fn
given_a_rejected_login_when_the_peer_retries_at_once_should_close_until_the_delay_passes()
+{
+ let addr = spawn_sasl_gateway().await;
+
+ let mut rejected = TcpStream::connect(addr).await.expect("connect");
+ handshake_ok(&mut rejected).await;
+ let body = send(
+ &mut rejected,
+ API_KEY_SASL_AUTHENTICATE,
+ AUTHENTICATE_VERSION,
+ 2,
+ &authenticate_body(&plain_token("alice", "wrong-password")),
+ )
+ .await;
+ assert_eq!(error_code(&body), ERROR_SASL_AUTHENTICATION_FAILED);
+ assert_closed(&mut rejected).await;
+
+ // Even the right password is not checked while the peer is throttled, and
it gets a close
+ // rather than 58, which a Kafka client would treat as fatal.
+ let mut throttled = TcpStream::connect(addr).await.expect("connect");
+ handshake_ok(&mut throttled).await;
+ let frame = build_request_frame(
+ API_KEY_SASL_AUTHENTICATE,
+ AUTHENTICATE_VERSION,
+ 2,
+ Some("sasl-test"),
+ &authenticate_body(&plain_token("alice", "s3cret")),
+ );
+ throttled.write_all(&frame).await.expect("write request");
+ assert_closed(&mut throttled).await;
+
+ tokio::time::sleep(Duration::from_millis(600)).await;
+ let mut retried = TcpStream::connect(addr).await.expect("connect");
+ handshake_ok(&mut retried).await;
+ let body = send(
+ &mut retried,
+ API_KEY_SASL_AUTHENTICATE,
+ AUTHENTICATE_VERSION,
+ 2,
+ &authenticate_body(&plain_token("alice", "s3cret")),
+ )
+ .await;
+ assert_eq!(
+ error_code(&body),
+ ERROR_NONE,
+ "the delay must run out on its own"
+ );
+}
+
#[tokio::test]
async fn
given_a_pre_auth_request_above_the_firewall_should_close_rather_than_answer_it()
{
// `kafka_protocol`'s encoders reach further than this gateway's firewall,
so encoding at the
diff --git a/gateways/kafka/tests/version_firewall_tests.rs
b/gateways/kafka/tests/version_firewall_tests.rs
index c0ea9f315..fb90877d2 100644
--- a/gateways/kafka/tests/version_firewall_tests.rs
+++ b/gateways/kafka/tests/version_firewall_tests.rs
@@ -346,7 +346,7 @@ async fn
create_topics_below_min_version_closes_connection() {
#[tokio::test]
async fn unsupported_api_keys_close_connection() {
- for key in [8, 9, 10, 11, 17, 20, 42, 999] {
+ for key in [8, 9, 10, 11, 20, 42, 999] {
let outcome = handle_request(key, 0, Bytes::new(),
&default_broker()).await;
assert!(
outcome.is_close(),