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(),

Reply via email to