This is an automated email from the ASF dual-hosted git repository.

numinnex pushed a commit to branch kafka_gateway_group
in repository https://gitbox.apache.org/repos/asf/iggy.git

commit 9221034d72a9345a916ea4aa7ab70642c6a1e85c
Merge: 898caeef9 768ea8093
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Fri Sep 25 15:01:49 2026 +0200

    Merge remote-tracking branch 'origin/master' into kafka_gateway_group

 .github/workflows/_common.yml                      |    6 +-
 .github/workflows/_test.yml                        |   16 +-
 .github/workflows/_test_examples.yml               |    2 +-
 .github/workflows/coverage-baseline.yml            |   20 +-
 .github/workflows/post-merge.yml                   |    2 +-
 .github/workflows/publish.yml                      |    2 +-
 AGENTS.md                                          |   13 +-
 Cargo.lock                                         |    1 +
 bdd/python/uv.lock                                 |   64 +-
 examples/python/pyproject.toml                     |    4 +-
 examples/python/uv.lock                            |   68 +-
 foreign/java/gradle/libs.versions.toml             |    2 +-
 foreign/python/uv.lock                             |   64 +-
 gateways/kafka/Cargo.toml                          |    1 +
 gateways/kafka/README.md                           |   58 ++
 gateways/kafka/docs/AUTHENTICATION.md              |  260 +++++
 gateways/kafka/docs/MANUAL_TESTING.md              |   64 ++
 gateways/kafka/docs/SCOPE.md                       |   14 +-
 gateways/kafka/docs/TEST_SUITE.md                  |    1 +
 gateways/kafka/src/auth.rs                         |  606 +++++++++++
 gateways/kafka/src/bridge/config.rs                |    6 +-
 gateways/kafka/src/env.rs                          |   64 ++
 gateways/kafka/src/lib.rs                          |    2 +
 gateways/kafka/src/main.rs                         |  120 ++-
 gateways/kafka/src/protocol/api.rs                 |  207 +++-
 gateways/kafka/src/protocol/bounds_guard.rs        |  120 +++
 .../kafka/src/protocol/handlers/api_versions.rs    |   34 +-
 gateways/kafka/src/protocol/mod.rs                 |    1 +
 gateways/kafka/src/protocol/sasl.rs                |  573 ++++++++++
 gateways/kafka/src/server.rs                       |  605 ++++++++++-
 gateways/kafka/tests/common/scope.rs               |    9 +-
 gateways/kafka/tests/common/server.rs              |   30 +-
 gateways/kafka/tests/common/wire.rs                |    1 -
 gateways/kafka/tests/consumer_group_tests.rs       |    1 +
 gateways/kafka/tests/golden_wire_fixtures_tests.rs |    8 +-
 .../kafka/tests/list_offsets_real_bridge_tests.rs  |    1 +
 gateways/kafka/tests/listener_robustness_tests.rs  |   34 +-
 gateways/kafka/tests/sasl_tests.rs                 | 1097 ++++++++++++++++++++
 gateways/kafka/tests/version_firewall_tests.rs     |    2 +-
 39 files changed, 3967 insertions(+), 216 deletions(-)

diff --cc gateways/kafka/README.md
index 2c3ac8cb5,71af7d23c..c62a536d0
--- a/gateways/kafka/README.md
+++ b/gateways/kafka/README.md
@@@ -72,7 -73,7 +75,8 @@@ See [docs/SCOPE.md](docs/SCOPE.md) for 
  - [docs/BRIDGE_MAPPING.md](docs/BRIDGE_MAPPING.md) — how a Kafka record 
becomes an Iggy message, and back
  - [docs/IDEMPOTENCE.md](docs/IDEMPOTENCE.md) — InitProducerId, and why 
delivery is at-least-once
  - [docs/OFFSET_STORAGE.md](docs/OFFSET_STORAGE.md) — where Kafka consumer 
group offsets live
 +- [docs/CONSUMER_GROUPS.md](docs/CONSUMER_GROUPS.md) — group membership, 
rebalances, and why one gateway per bootstrap endpoint
+ - [docs/AUTHENTICATION.md](docs/AUTHENTICATION.md) — how a Kafka client 
authenticates, and why PLAIN only
  
  ### Delivery guarantees
  
diff --cc gateways/kafka/docs/SCOPE.md
index 0729b7eb1,4645d04dd..3c0a4627e
--- a/gateways/kafka/docs/SCOPE.md
+++ b/gateways/kafka/docs/SCOPE.md
@@@ -81,12 -72,11 +81,12 @@@ All API keys not listed above close th
  
  | API key | Name | Notes |
  | --------- | ------ | ------- |
 -| 8 | OffsetCommit | Consumer group — later issue |
 -| 9 | OffsetFetch | Consumer group — later issue |
 -| 10 | FindCoordinator | Consumer group — later issue |
 -| 11–16 | JoinGroup, Heartbeat, LeaveGroup, SyncGroup, DescribeGroups, 
ListGroups | Consumer group — later issue |
 +| 8 | OffsetCommit | Consumer group offsets — 
[#3542](https://github.com/apache/iggy/issues/3542) |
 +| 9 | OffsetFetch | Consumer group offsets — 
[#3542](https://github.com/apache/iggy/issues/3542); sent right after 
SyncGroup, so a joined consumer loops on it today 
([`CONSUMER_GROUPS.md`](CONSUMER_GROUPS.md)) |
 +| 13 | LeaveGroup | Graceful shutdown — 
[#3543](https://github.com/apache/iggy/issues/3543); without it a departing 
member is evicted by session expiry instead |
 +| 15, 16 | DescribeGroups, ListGroups | Admin views — 
[#3548](https://github.com/apache/iggy/issues/3548) |
- | 17 | SaslHandshake | Auth — later issue |
+ | 17 | SaslHandshake | Implemented behind `IGGY_KAFKA_SASL_ENABLED`, 
advertised only while it is on ([`AUTHENTICATION.md`](AUTHENTICATION.md)) |
 +| 68 | ConsumerGroupHeartbeat | KIP-848 protocol; a 4.0 client may need 
`group.protocol=classic` |
  | 20+ | DeleteTopics, InitProducerId, transactions, ACLs, etc. | Later issues 
|
  
  Full reference for future phases: 
[`kafka_api_keys_reference.md`](kafka_api_keys_reference.md).
diff --cc gateways/kafka/docs/TEST_SUITE.md
index d19ca3db6,e912a421a..64c94e942
--- a/gateways/kafka/docs/TEST_SUITE.md
+++ b/gateways/kafka/docs/TEST_SUITE.md
@@@ -61,9 -61,9 +61,10 @@@ file under `tests/` anymore
  | [`version_firewall_tests.rs`](../tests/version_firewall_tests.rs) | Version 
boundary matrix, unsupported keys, corrupt bodies | Partial |
  | [`broker_advertise_tests.rs`](../tests/broker_advertise_tests.rs) | 
`BrokerAdvertise::from_server_config` parsing | No |
  | [`server_integration_tests.rs`](../tests/server_integration_tests.rs) | 
`read_frame` unit-level I/O | No |
 +| [`consumer_group_tests.rs`](../tests/consumer_group_tests.rs) | 
`FindCoordinator`/`JoinGroup`/`Heartbeat`/`SyncGroup` against one shared 
`GatewayState` with paused time, plus a two-socket TCP rebalance | No |
  | [`server_e2e_tests.rs`](../tests/server_e2e_tests.rs) | Full `KafkaGateway` 
TCP round-trips | Partial |
  | [`listener_robustness_tests.rs`](../tests/listener_robustness_tests.rs) | 
TCP listener robustness — framing, pipelining, concurrency, connection limits | 
No |
+ | [`sasl_tests.rs`](../tests/sasl_tests.rs) | SASL/PLAIN over a socket — full 
handshake, every refusal path, and the disabled default. Drives a stub verifier 
implementing `SaslAuthenticator`, so no Iggy server is needed | No |
  | 
[`bridge_iggy_integration_tests.rs`](../tests/bridge_iggy_integration_tests.rs) 
| `IggyBridge` against a real, spawned `iggy-server` — provisioning 
idempotency, high watermark, credential/connection edge cases | No (needs the 
`iggy-server` binary - see Prerequisites) |
  
  `tests/common/` holds shared helpers (`codec.rs`, `fixtures.rs`, `scope.rs`, 
`server.rs`,
diff --cc gateways/kafka/src/lib.rs
index 0f020eccb,101e201e4..722ddb6cb
--- a/gateways/kafka/src/lib.rs
+++ b/gateways/kafka/src/lib.rs
@@@ -17,9 -17,10 +17,11 @@@
  
  //! Kafka wire protocol gateway foundation for Apache Iggy.
  
+ pub mod auth;
  pub mod bridge;
+ pub mod env;
  pub mod error;
 +pub mod group;
  pub mod protocol;
  pub mod records;
  pub mod server;
diff --cc gateways/kafka/src/protocol/api.rs
index e37a78cdc,40c2a1188..726439a82
--- a/gateways/kafka/src/protocol/api.rs
+++ b/gateways/kafka/src/protocol/api.rs
@@@ -19,25 -19,28 +19,35 @@@ use std::sync::Arc
  
  use bytes::Bytes;
  
+ use kafka_protocol::messages::{SaslAuthenticateRequest, SaslHandshakeRequest};
 +use tokio_util::sync::CancellationToken;
  
  use crate::bridge::IggyBridge;
+ use crate::error::Result;
 +use crate::group::{GroupCoordinator, GroupCoordinatorConfig};
+ use crate::protocol::bounds_guard::{
+     validate_sasl_authenticate_shape, validate_sasl_handshake_shape,
+ };
  use crate::protocol::handlers::{
-     api_versions, create_topics, dispatch, fetch, find_coordinator, 
heartbeat, join_group,
-     list_offsets, metadata, produce, sync_group,
 -    api_versions, create_topics, decode_guarded, dispatch, fetch, 
list_offsets, metadata, produce,
++    api_versions, create_topics, decode_guarded, dispatch, fetch, 
find_coordinator, heartbeat,
++    join_group, list_offsets, metadata, produce, sync_group,
+ };
+ use crate::protocol::sasl::{
+     SaslMechanism, encode_sasl_authenticate_response, 
encode_sasl_handshake_response,
  };
  
  pub const API_KEY_PRODUCE: i16 = 0;
  pub const API_KEY_FETCH: i16 = 1;
  pub const API_KEY_LIST_OFFSETS: i16 = 2;
  pub const API_KEY_METADATA: i16 = 3;
 +pub const API_KEY_FIND_COORDINATOR: i16 = 10;
 +pub const API_KEY_JOIN_GROUP: i16 = 11;
 +pub const API_KEY_HEARTBEAT: i16 = 12;
 +pub const API_KEY_SYNC_GROUP: i16 = 14;
+ pub const API_KEY_SASL_HANDSHAKE: i16 = 17;
  pub const API_KEY_API_VERSIONS: i16 = 18;
  pub const API_KEY_CREATE_TOPICS: i16 = 19;
+ pub const API_KEY_SASL_AUTHENTICATE: i16 = 36;
  
  pub const DEFAULT_KAFKA_PORT: u16 = 9093;
  
@@@ -65,27 -68,19 +75,35 @@@ pub const ERROR_REQUEST_TIMED_OUT: i16 
  /// call is made - a real Kafka client library validates topic names 
client-side and would never
  /// send one of these, but a raw/non-conformant client could.
  pub const ERROR_INVALID_TOPIC_EXCEPTION: i16 = 17;
 +/// Retriable. Sent when this coordinator is at one of its 
`GroupCoordinatorConfig` capacity
 +/// caps: the client should back off and retry rather than treat the group as 
unusable.
 +pub const ERROR_COORDINATOR_NOT_AVAILABLE: i16 = 15;
 +/// Retriable. Sent to a parked `JoinGroup`/`SyncGroup` waiter when the 
gateway starts draining,
 +/// so a shutdown does not hold a connection open for a full rebalance 
timeout.
 +pub const ERROR_NOT_COORDINATOR: i16 = 16;
 +/// The member's generation is not the group's current one; it must rejoin.
 +pub const ERROR_ILLEGAL_GENERATION: i16 = 22;
 +/// No protocol name is supported by every member, or the first member sent 
an empty protocol
 +/// type / empty protocol list.
 +pub const ERROR_INCONSISTENT_GROUP_PROTOCOL: i16 = 23;
 +pub const ERROR_INVALID_GROUP_ID: i16 = 24;
 +pub const ERROR_UNKNOWN_MEMBER_ID: i16 = 25;
 +pub const ERROR_INVALID_SESSION_TIMEOUT: i16 = 26;
 +/// How a follower learns to rejoin: its heartbeat is answered with this 
while the group prepares.
 +pub const ERROR_REBALANCE_IN_PROGRESS: i16 = 27;
  /// Closest fit for an Iggy permission/credential rejection in `bridge`'s 
error mapping.
  ///
- /// There is no bridge-side SASL exchange yet (`#3549`), so 
`SASL_AUTHENTICATION_FAILED` would
- /// misstate the failure point. Not sent by any stub response today.
+ /// Still not `SASL_AUTHENTICATION_FAILED`, and now for a firmer reason than 
when this was written:
+ /// a connection reaching the bridge has already completed its SASL exchange, 
so a credential or
+ /// permission rejection from Iggy at that point is not an authentication 
failure and saying so
+ /// would send an operator to the wrong hop. Not sent by any stub response 
today.
  pub const ERROR_TOPIC_AUTHORIZATION_FAILED: i16 = 29;
+ /// The mechanism a client asked for in `SaslHandshake` is not one this 
gateway enables. The
+ /// response still carries the enabled mechanism list, which is what the 
client prints.
+ pub const ERROR_UNSUPPORTED_SASL_MECHANISM: i16 = 33;
+ /// A request arrived that is legal on the wire but not in this connection's 
SASL state: a token
+ /// before a handshake, a normal request before authenticating, or a SASL 
request after.
+ pub const ERROR_ILLEGAL_SASL_STATE: i16 = 34;
  pub const ERROR_UNSUPPORTED_VERSION: i16 = 35;
  /// `bridge`'s mapping for `BridgeError::PartitionCountMismatch`: the topic 
exists, just not with
  /// the requested partition count.
@@@ -106,15 -101,10 +124,19 @@@ pub const ERROR_INVALID_REQUEST: i16 = 
  /// Non-retriable, so a Java client resolves immediately instead of retrying
  /// [`ERROR_UNKNOWN_SERVER_ERROR`] until its own `default.api.timeout.ms`.
  pub const ERROR_UNSUPPORTED_FOR_MESSAGE_FORMAT: i16 = 43;
 +/// `FindCoordinator` for a transaction coordinator, which the gateway has 
none of.
 +///
 +/// The one code both the Java client and librdkafka treat as fatal for that 
lookup: anything else,
 +/// `INVALID_REQUEST` included, sends librdkafka into a 500ms retry loop that 
never ends.
 +pub const ERROR_TRANSACTIONAL_ID_AUTHORIZATION_FAILED: i16 = 53;
+ /// A credential was refused. Deliberately undifferentiated: Iggy answers a 
bad password and an
+ /// unknown user the same way, and distinguishing them here would reintroduce 
a user-enumeration
+ /// oracle.
+ pub const ERROR_SASL_AUTHENTICATION_FAILED: i16 = 58;
 +/// KIP-394: a `JoinGroup` v4+ with an empty member id is answered with a 
freshly minted id and
 +/// this code, and the client rejoins carrying it.
 +pub const ERROR_MEMBER_ID_REQUIRED: i16 = 79;
 +pub const ERROR_GROUP_MAX_SIZE_REACHED: i16 = 81;
  
  /// Result of handling one Kafka request body.
  #[derive(Debug)]
@@@ -205,9 -199,9 +235,12 @@@ pub struct GatewayState 
      pub broker: BrokerAdvertise,
      pub bridge: Option<Arc<IggyBridge>>,
      pub max_frame_size: usize,
+     /// Whether `SaslHandshake` and `SaslAuthenticate` are advertised and 
routed. Kept on the
+     /// shared state so `ApiVersions` can answer without a widened handler 
signature.
+     pub sasl_enabled: bool,
 +    /// Consumer group membership. Process-wide and independent of the 
bridge: a member outlives
 +    /// the connection that created it, and group coordination needs no Iggy 
call.
 +    pub groups: GroupCoordinator,
  }
  
  impl GatewayState {
@@@ -216,27 -210,20 +249,30 @@@
          broker: BrokerAdvertise,
          bridge: Option<Arc<IggyBridge>>,
          max_frame_size: usize,
+         sasl_enabled: bool,
 +        groups: GroupCoordinator,
      ) -> Self {
          Self {
              broker,
              bridge,
              max_frame_size,
+             sasl_enabled,
 +            groups,
          }
      }
  
      /// State with no bridge, so every handler takes its stub path.
 +    ///
 +    /// The coordinator is real but fresh, so two calls never share group 
state.
      #[must_use]
 -    pub const fn stub(broker: BrokerAdvertise, max_frame_size: usize) -> Self 
{
 -        Self::new(broker, None, max_frame_size, false)
 +    pub fn stub(broker: BrokerAdvertise, max_frame_size: usize) -> Self {
 +        Self::new(
 +            broker,
 +            None,
 +            max_frame_size,
++            false,
 +            GroupCoordinator::new(GroupCoordinatorConfig::default(), 
CancellationToken::new()),
 +        )
      }
  }
  
diff --cc gateways/kafka/src/protocol/bounds_guard.rs
index ebb47a144,d0c2cd4e5..b476a19d3
--- a/gateways/kafka/src/protocol/bounds_guard.rs
+++ b/gateways/kafka/src/protocol/bounds_guard.rs
@@@ -769,200 -730,63 +769,255 @@@ pub fn validate_api_versions_shape(vers
      Ok(())
  }
  
 +/// Mirrors the field order `FindCoordinatorRequest::decode` walks.
 +///
 +/// # Errors
 +///
 +/// Returns an error when a declared array/string length cannot fit in the 
bytes remaining in the
 +/// frame, or the body is truncated or malformed in a way that cannot be 
walked.
 +pub fn validate_find_coordinator_shape(
 +    version: i16,
 +    body: &Bytes,
 +    max_frame_size: usize,
 +) -> Result<()> {
 +    let mut c = ShapeCursor::new(body.clone(), max_frame_size);
 +    let flexible = version >= 3;
 +
 +    if version <= 3 {
 +        if flexible {
 +            c.compact_string(false)?;
 +        } else {
 +            c.legacy_string(false)?;
 +        }
 +    }
 +    if version >= 1 {
 +        let _key_type = c.read_i8()?;
 +    }
 +    if version >= 4 {
 +        let keys_count = c.compact_array_count()?;
 +        for _ in 0..keys_count {
 +            c.compact_string(false)?;
 +        }
 +    }
 +    if flexible {
 +        c.tagged_fields()?;
 +    }
 +    Ok(())
 +}
 +
 +/// Mirrors the field order `JoinGroupRequest::decode` walks.
 +///
 +/// # Errors
 +///
 +/// Returns an error when a declared array/string/bytes length cannot fit in 
the bytes remaining
 +/// in the frame, or the body is truncated or malformed in a way that cannot 
be walked.
 +pub fn validate_join_group_shape(version: i16, body: &Bytes, max_frame_size: 
usize) -> Result<()> {
 +    let mut c = ShapeCursor::new(body.clone(), max_frame_size);
 +    let flexible = version >= 6;
 +
 +    if flexible {
 +        c.compact_string(false)?;
 +    } else {
 +        c.legacy_string(false)?;
 +    }
 +    let _session_timeout_ms = c.read_i32()?;
 +    if version >= 1 {
 +        let _rebalance_timeout_ms = c.read_i32()?;
 +    }
 +    if flexible {
 +        c.compact_string(false)?;
 +    } else {
 +        c.legacy_string(false)?;
 +    }
 +    if version >= 5 {
 +        if flexible {
 +            c.compact_string(true)?;
 +        } else {
 +            c.legacy_string(true)?;
 +        }
 +    }
 +    if flexible {
 +        c.compact_string(false)?;
 +    } else {
 +        c.legacy_string(false)?;
 +    }
 +
 +    let protocols_count = if flexible {
 +        c.compact_array_count()?
 +    } else {
 +        c.legacy_array_count()?
 +    };
 +    for _ in 0..protocols_count {
 +        if flexible {
 +            c.compact_string(false)?;
 +        } else {
 +            c.legacy_string(false)?;
 +        }
 +        c.echoed_bytes(false, flexible)?;
 +        if flexible {
 +            c.tagged_fields()?;
 +        }
 +    }
 +
 +    if version >= 8 {
 +        c.compact_string(true)?;
 +    }
 +    if flexible {
 +        c.tagged_fields()?;
 +    }
 +    Ok(())
 +}
 +
 +/// Mirrors the field order `HeartbeatRequest::decode` walks.
 +///
 +/// No response-size guard is needed: a Heartbeat response is a throttle time 
and an error code,
 +/// so nothing in the request can amplify it (same reasoning as 
`validate_api_versions_shape`).
 +///
 +/// # Errors
 +///
 +/// Returns an error when a declared string length cannot fit in the bytes 
remaining in the
 +/// frame, or the body is truncated or malformed in a way that cannot be 
walked.
 +pub fn validate_heartbeat_shape(version: i16, body: &Bytes) -> Result<()> {
 +    let mut c = ShapeCursor::new(body.clone(), usize::MAX);
 +    let flexible = version >= 4;
 +
 +    if flexible {
 +        c.compact_string(false)?;
 +    } else {
 +        c.legacy_string(false)?;
 +    }
 +    let _generation_id = c.read_i32()?;
 +    if flexible {
 +        c.compact_string(false)?;
 +    } else {
 +        c.legacy_string(false)?;
 +    }
 +    if version >= 3 {
 +        if flexible {
 +            c.compact_string(true)?;
 +        } else {
 +            c.legacy_string(true)?;
 +        }
 +    }
 +    if flexible {
 +        c.tagged_fields()?;
 +    }
 +    Ok(())
 +}
 +
 +/// Mirrors the field order `SyncGroupRequest::decode` walks.
 +///
 +/// # Errors
 +///
 +/// Returns an error when a declared array/string/bytes length cannot fit in 
the bytes remaining
 +/// in the frame, or the body is truncated or malformed in a way that cannot 
be walked.
 +pub fn validate_sync_group_shape(version: i16, body: &Bytes, max_frame_size: 
usize) -> Result<()> {
 +    let mut c = ShapeCursor::new(body.clone(), max_frame_size);
 +    let flexible = version >= 4;
 +
 +    if flexible {
 +        c.compact_string(false)?;
 +    } else {
 +        c.legacy_string(false)?;
 +    }
 +    let _generation_id = c.read_i32()?;
 +    if flexible {
 +        c.compact_string(false)?;
 +    } else {
 +        c.legacy_string(false)?;
 +    }
 +    if version >= 3 {
 +        if flexible {
 +            c.compact_string(true)?;
 +        } else {
 +            c.legacy_string(true)?;
 +        }
 +    }
 +    if version >= 5 {
 +        c.compact_string(true)?;
 +        c.compact_string(true)?;
 +    }
 +
 +    let assignments_count = if flexible {
 +        c.compact_array_count()?
 +    } else {
 +        c.legacy_array_count()?
 +    };
 +    for _ in 0..assignments_count {
 +        if flexible {
 +            c.compact_string(false)?;
 +        } else {
 +            c.legacy_string(false)?;
 +        }
 +        c.echoed_bytes(false, flexible)?;
 +        if flexible {
 +            c.tagged_fields()?;
 +        }
 +    }
 +
 +    if flexible {
 +        c.tagged_fields()?;
 +    }
 +    Ok(())
 +}
 +
+ /// Cap on a whole `SaslAuthenticate` body.
+ ///
+ /// The body is one length-prefixed blob, so capping it is the same as 
capping `auth_bytes`.
+ /// PLAIN credentials are a couple of hundred bytes at the outside (Iggy caps 
a username at 50 and
+ /// a password at 100), and this is the one frame an *unauthenticated* 
connection can send
+ /// repeatedly. Generous enough to leave room for a future multi-round 
mechanism.
+ ///
+ /// This is not a memory bound on the connection. `read_frame` buffers the 
whole frame, up to
+ /// `max_frame_size`, before a header is even decoded, so an unauthenticated 
peer can still make
+ /// the gateway hold that much. What this cap bounds is everything downstream 
of the guard: what
+ /// reaches `parse_plain`, and what a future mechanism would carry into a 
credential exchange.
+ const MAX_SASL_AUTH_BYTES: usize = 4096;
+ 
+ /// `SaslHandshake` carries one non-nullable string, the mechanism name.
+ ///
+ /// `_version` is unused: the message is never flexible, so v0 and v1 share 
this shape, and the
+ /// state machine refuses v0 before a body is ever validated.
+ ///
+ /// # Errors
+ ///
+ /// Returns an error when the declared string length does not fit the 
remaining frame.
+ pub fn validate_sasl_handshake_shape(_version: i16, body: &Bytes) -> 
Result<()> {
+     // No response-size guard: the response echoes a fixed mechanism list, 
never anything decoded
+     // from this body, so nothing here can amplify.
+     let mut c = ShapeCursor::new(body.clone(), usize::MAX);
+     c.legacy_string(false)?;
+     Ok(())
+ }
+ 
+ /// `SaslAuthenticate` carries one non-nullable bytes field, the mechanism 
token.
+ ///
+ /// # Errors
+ ///
+ /// Returns [`KafkaProtocolError::FrameTooLarge`] when the body exceeds this 
module's
+ /// `MAX_SASL_AUTH_BYTES` cap, or an error when the declared length does not 
fit the remaining
+ /// frame.
+ pub fn validate_sasl_authenticate_shape(version: i16, body: &Bytes) -> 
Result<()> {
+     if body.len() > MAX_SASL_AUTH_BYTES {
+         return Err(KafkaProtocolError::FrameTooLarge {
+             max_bytes: MAX_SASL_AUTH_BYTES,
+             actual_bytes: body.len(),
+         });
+     }
+     let mut c = ShapeCursor::new(body.clone(), usize::MAX);
+     if version >= 2 {
+         c.compact_bytes(false)?;
+         c.tagged_fields()?;
+     } else {
+         c.legacy_bytes(false)?;
+     }
+     Ok(())
+ }
+ 
  #[cfg(test)]
  mod tests {
+     use bytes::BytesMut;
+ 
      use super::*;
  
      const TEST_MAX_FRAME_SIZE: usize = 8 * 1024 * 1024;
diff --cc gateways/kafka/src/server.rs
index 9a8c91c3d,8859f0f80..f3ffae42f
--- a/gateways/kafka/src/server.rs
+++ b/gateways/kafka/src/server.rs
@@@ -30,14 -30,24 +30,25 @@@ 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, FailedLoginThrottle, SaslAuthenticator};
  use crate::bridge::IggyBridge;
  use crate::error::{KafkaProtocolError, Result};
 +use crate::group::{GroupCoordinator, GroupCoordinatorConfig};
  use crate::protocol::api::{
-     BrokerAdvertise, DEFAULT_KAFKA_PORT, GatewayState, HandleOutcome, 
handle_request_bounded,
+     API_KEY_SASL_AUTHENTICATE, API_KEY_SASL_HANDSHAKE, BrokerAdvertise, 
DEFAULT_KAFKA_PORT,
+     ERROR_ILLEGAL_SASL_STATE, ERROR_NONE, ERROR_SASL_AUTHENTICATION_FAILED,
+     ERROR_UNSUPPORTED_SASL_MECHANISM, ERROR_UNSUPPORTED_VERSION, 
GatewayState, HandleOutcome,
+     decode_sasl_auth_bytes, decode_sasl_mechanism, encode_error_for_key, 
handle_request_bounded,
+     sasl_authenticate_outcome, sasl_handshake_outcome,
  };
  use crate::protocol::header::{request_header_version, 
response_header_version};
+ use crate::protocol::sasl::{
+     PlainCredentials, SASL_AUTHENTICATE_MAX_VERSION, SASL_HANDSHAKE_VERSION, 
SaslAction, SaslState,
+     parse_plain,
+ };
  use std::io;
  
  const READ_CHUNK: usize = 65536;
@@@ -64,9 -116,40 +117,43 @@@ pub struct GatewayConfig 
      /// hold shutdown open past typical orchestrator grace periods (e.g. 
Kubernetes' default
      /// 30s `terminationGracePeriodSeconds`).
      pub shutdown_drain_timeout: Duration,
 +    /// Consumer-group timeouts and capacity caps. No environment variable 
maps onto these yet;
 +    /// Kafka's own defaults apply.
 +    pub group: GroupCoordinatorConfig,
+     /// Require SASL authentication before serving any other API.
+     ///
+     /// Off by default, and switching it on is a breaking change for every 
client already talking
+     /// to this gateway: the two SASL keys only appear in the `ApiVersions` 
advertisement while it
+     /// is on, and unauthenticated clients stop being served.
+     ///
+     /// While it is off the keys are still answered, with 
`ILLEGAL_SASL_STATE` and an empty
+     /// mechanism list, because a well-formed response exists for them unlike 
a genuinely unknown
+     /// key. They are kept out of the advertisement so nothing is invited 
into an exchange that
+     /// cannot finish.
+     pub sasl_enabled: bool,
+     /// How long an unauthenticated connection may sit between frames.
+     ///
+     /// Separate from `idle_timeout` because that one is sized for a 
well-behaved idle client (10
+     /// minutes, matching a real broker) and applies to connections that have 
proven who they are.
+     /// 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 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 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.
+     pub max_concurrent_authentications: usize,
  }
  
  impl Default for GatewayConfig {
@@@ -81,7 -164,9 +168,10 @@@
              read_timeout: Duration::from_secs(15),
              write_timeout: Duration::from_secs(10),
              shutdown_drain_timeout: Duration::from_secs(25),
 +            group: GroupCoordinatorConfig::default(),
+             sasl_enabled: false,
+             pre_auth_timeout: Duration::from_secs(15),
+             max_concurrent_authentications: 4,
          }
      }
  }
@@@ -202,6 -317,25 +322,7 @@@ impl KafkaGateway 
          listener: TcpListener,
          mut shutdown: broadcast::Receiver<()>,
      ) -> Result<()> {
 -        if !self.config.sasl_enabled && self.authenticator.is_some() {
 -            // The mirror of the guard below, and the quieter mistake: a 
verifier attached while the
 -            // flag is off means every connection is served unauthenticated, 
with nothing in the log
 -            // to say so. Refusing to start is the only way that failure is 
visible.
 -            return Err(KafkaProtocolError::InvalidConfig(
 -                "an authenticator is configured but SASL is disabled; every 
connection would be \
 -                 served unauthenticated. Set IGGY_KAFKA_SASL_ENABLED=true, or 
remove the \
 -                 authenticator"
 -                    .into(),
 -            ));
 -        }
 -        if self.config.sasl_enabled && self.authenticator.is_none() {
 -            return Err(KafkaProtocolError::InvalidConfig(
 -                "SASL is enabled but no authenticator is configured; every 
client would be \
 -                 rejected. Set IGGY_KAFKA_IGGY_ADDR to the Iggy server that 
credentials are \
 -                 verified against, or unset IGGY_KAFKA_SASL_ENABLED"
 -                    .into(),
 -            ));
 -        }
++        self.check_sasl_wiring()?;
          let local_addr = listener.local_addr()?;
          let broker = BrokerAdvertise::from_server_config(&self.config, 
local_addr)?;
          info!(
@@@ -217,11 -346,18 +338,16 @@@
              broker,
              self.bridge.clone(),
              self.config.max_frame_size,
+             self.config.sasl_enabled,
 +            GroupCoordinator::new(self.config.group.clone(), 
cancel.child_token()),
          ));
  
+         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
 -        // until `idle_timeout` (or forever if that is raised).
 -        let cancel = CancellationToken::new();
  
          let drain_timeout = self.config.shutdown_drain_timeout;
  
@@@ -296,6 -434,6 +424,30 @@@
          }
          Ok(())
      }
++
++    /// Refuses to start when the SASL flag and the configured authenticator 
disagree.
++    fn check_sasl_wiring(&self) -> Result<()> {
++        if !self.config.sasl_enabled && self.authenticator.is_some() {
++            // The mirror of the guard below, and the quieter mistake: a 
verifier attached while the
++            // flag is off means every connection is served unauthenticated, 
with nothing in the log
++            // to say so. Refusing to start is the only way that failure is 
visible.
++            return Err(KafkaProtocolError::InvalidConfig(
++                "an authenticator is configured but SASL is disabled; every 
connection would be \
++                 served unauthenticated. Set IGGY_KAFKA_SASL_ENABLED=true, or 
remove the \
++                 authenticator"
++                    .into(),
++            ));
++        }
++        if self.config.sasl_enabled && self.authenticator.is_none() {
++            return Err(KafkaProtocolError::InvalidConfig(
++                "SASL is enabled but no authenticator is configured; every 
client would be \
++                 rejected. Set IGGY_KAFKA_IGGY_ADDR to the Iggy server that 
credentials are \
++                 verified against, or unset IGGY_KAFKA_SASL_ENABLED"
++                    .into(),
++            ));
++        }
++        Ok(())
++    }
  }
  
  /// Cancel in-flight connections, close the tracker to new spawns, and wait 
for tasks to finish
diff --cc gateways/kafka/tests/common/scope.rs
index c2ef894d0,55b1361d7..c655030ef
--- a/gateways/kafka/tests/common/scope.rs
+++ b/gateways/kafka/tests/common/scope.rs
@@@ -20,21 -20,14 +20,22 @@@
  
  use iggy_gateway_kafka::protocol::api::BrokerAdvertise;
  
 -/// Scoped API keys exercised by the #3421 regression suite.
 +/// Scoped API keys exercised by the regression suite.
 +///
- /// Declaration order mirrors `SUPPORTED_RANGES` in `src/protocol/api.rs`: 
the `ApiVersions` tests
- /// compare the advertised list row by row against this one.
++/// Ascending by key, the order `ApiVersions` advertises in regardless of how 
`SUPPORTED_RANGES` in
++/// `src/protocol/api.rs` is declared: the `ApiVersions` tests compare the 
advertised list row by
++/// row against this one.
  pub const SCOPED_API_KEYS: &[(i16, &str, i16, i16)] = &[
      (0, "Produce", 3, 9),
      (1, "Fetch", 4, 12),
      (2, "ListOffsets", 1, 6),
      (3, "Metadata", 0, 9),
-     (18, "ApiVersions", 0, 3),
-     (19, "CreateTopics", 2, 5),
 +    (10, "FindCoordinator", 0, 4),
 +    (11, "JoinGroup", 0, 9),
 +    (12, "Heartbeat", 0, 4),
 +    (14, "SyncGroup", 0, 5),
+     (18, "ApiVersions", 0, 3),
+     (19, "CreateTopics", 2, 5),
  ];
  
  pub fn default_broker() -> BrokerAdvertise {
diff --cc gateways/kafka/tests/common/wire.rs
index 612e161d1,304604ed1..eb8caf3fe
--- a/gateways/kafka/tests/common/wire.rs
+++ b/gateways/kafka/tests/common/wire.rs
@@@ -34,10 -34,13 +34,9 @@@ use super::codec::Encoder
  pub const OUT_OF_SCOPE_API_KEYS: &[(i16, &str)] = &[
      (8, "OffsetCommit"),
      (9, "OffsetFetch"),
 -    (10, "FindCoordinator"),
 -    (11, "JoinGroup"),
 -    (12, "Heartbeat"),
      (13, "LeaveGroup"),
 -    (14, "SyncGroup"),
      (15, "DescribeGroups"),
      (16, "ListGroups"),
-     (17, "SaslHandshake"),
      (20, "DeleteTopics"),
  ];
  
diff --cc gateways/kafka/tests/consumer_group_tests.rs
index 8de172595,000000000..9a32c22a7
mode 100644,000000..100644
--- a/gateways/kafka/tests/consumer_group_tests.rs
+++ b/gateways/kafka/tests/consumer_group_tests.rs
@@@ -1,1039 -1,0 +1,1040 @@@
 +// Licensed to the Apache Software Foundation (ASF) under one
 +// or more contributor license agreements.  See the NOTICE file
 +// distributed with this work for additional information
 +// regarding copyright ownership.  The ASF licenses this file
 +// to you under the Apache License, Version 2.0 (the
 +// "License"); you may not use this file except in compliance
 +// with the License.  You may obtain a copy of the License at
 +//
 +//   http://www.apache.org/licenses/LICENSE-2.0
 +//
 +// Unless required by applicable law or agreed to in writing,
 +// software distributed under the License is distributed on an
 +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
 +// KIND, either express or implied.  See the License for the
 +// specific language governing permissions and limitations
 +// under the License.
 +
 +//! Consumer group coordination: `FindCoordinator`, `JoinGroup`, Heartbeat, 
`SyncGroup`.
 +//!
 +//! Requests go through `handle_request_bounded` against one shared 
`GatewayState`, because
 +//! `handle_request` builds a fresh coordinator per call and no two requests 
would ever see the
 +//! same group. Time is paused, so every rebalance and eviction deadline is 
reached by advancing
 +//! it rather than by sleeping.
 +
 +#[path = "common/codec.rs"]
 +mod codec;
 +#[path = "common/scope.rs"]
 +mod scope;
 +#[path = "common/server.rs"]
 +mod server;
 +#[path = "common/tcp.rs"]
 +mod tcp;
 +#[path = "common/wire.rs"]
 +mod wire;
 +
 +use std::sync::Arc;
 +use std::time::Duration;
 +
 +use bytes::Bytes;
 +use tokio::io::AsyncWriteExt;
 +use tokio::net::TcpStream;
 +use tokio::time::advance;
 +use tokio_util::sync::CancellationToken;
 +
 +use iggy_gateway_kafka::GatewayConfig;
 +use iggy_gateway_kafka::group::{GroupCoordinator, GroupCoordinatorConfig};
 +use iggy_gateway_kafka::protocol::api::{
 +    API_KEY_FIND_COORDINATOR, API_KEY_HEARTBEAT, API_KEY_JOIN_GROUP, 
API_KEY_SYNC_GROUP,
 +    BrokerAdvertise, ERROR_GROUP_MAX_SIZE_REACHED, ERROR_ILLEGAL_GENERATION,
 +    ERROR_INCONSISTENT_GROUP_PROTOCOL, ERROR_INVALID_GROUP_ID, 
ERROR_INVALID_REQUEST,
 +    ERROR_INVALID_SESSION_TIMEOUT, ERROR_MEMBER_ID_REQUIRED, ERROR_NONE,
 +    ERROR_REBALANCE_IN_PROGRESS, ERROR_TRANSACTIONAL_ID_AUTHORIZATION_FAILED,
 +    ERROR_UNKNOWN_MEMBER_ID, GatewayState, handle_request_bounded,
 +};
 +
 +use codec::Decoder;
 +use server::spawn_test_server_with_config;
 +use tcp::{build_request_frame, parse_response_payload, read_response_frame};
 +use wire::{
 +    JoinGroupParams, SyncGroupParams, build_find_coordinator_request, 
build_heartbeat_request,
 +    build_join_group_request, build_sync_group_request,
 +};
 +
 +const GROUP: &str = "orders";
 +const JOIN_VERSION: i16 = 9;
 +const SYNC_VERSION: i16 = 5;
 +const HEARTBEAT_VERSION: i16 = 4;
 +const SESSION_TIMEOUT_MS: i32 = 10_000;
 +const REBALANCE_TIMEOUT_MS: i32 = 20_000;
 +
 +// ── Fixtures 
────────────────────────────────────────────────────────────────
 +
 +fn test_state(config: GroupCoordinatorConfig) -> Arc<GatewayState> {
 +    Arc::new(GatewayState::new(
 +        BrokerAdvertise::default(),
 +        None,
 +        8 * 1024 * 1024,
++        false,
 +        GroupCoordinator::new(config, CancellationToken::new()),
 +    ))
 +}
 +
 +/// A coordinator that completes a new group's first join immediately, so 
tests that are not
 +/// about the initial rebalance delay do not have to step time past it.
 +fn immediate_config() -> GroupCoordinatorConfig {
 +    GroupCoordinatorConfig {
 +        initial_rebalance_delay: Duration::ZERO,
 +        ..GroupCoordinatorConfig::default()
 +    }
 +}
 +
 +fn join_params<'a>(member_id: &'a str, metadata: &'a [(&'a str, &'a [u8])]) 
-> JoinGroupParams<'a> {
 +    JoinGroupParams {
 +        group_id: GROUP,
 +        session_timeout_ms: SESSION_TIMEOUT_MS,
 +        rebalance_timeout_ms: REBALANCE_TIMEOUT_MS,
 +        member_id,
 +        protocols: metadata,
 +        ..JoinGroupParams::default()
 +    }
 +}
 +
 +async fn join(state: &GatewayState, version: i16, params: 
&JoinGroupParams<'_>) -> JoinResponse {
 +    let body = build_join_group_request(version, params);
 +    let response = handle_request_bounded(state, API_KEY_JOIN_GROUP, version, 
body)
 +        .await
 +        .expect_response("JoinGroup must answer");
 +    JoinResponse::decode(version, response)
 +}
 +
 +async fn sync(state: &GatewayState, version: i16, params: 
&SyncGroupParams<'_>) -> SyncResponse {
 +    let body = build_sync_group_request(version, params);
 +    let response = handle_request_bounded(state, API_KEY_SYNC_GROUP, version, 
body)
 +        .await
 +        .expect_response("SyncGroup must answer");
 +    SyncResponse::decode(version, response)
 +}
 +
 +async fn heartbeat(state: &GatewayState, generation_id: i32, member_id: &str) 
-> i16 {
 +    let body = build_heartbeat_request(HEARTBEAT_VERSION, GROUP, 
generation_id, member_id);
 +    let response = handle_request_bounded(state, API_KEY_HEARTBEAT, 
HEARTBEAT_VERSION, body)
 +        .await
 +        .expect_response("Heartbeat must answer");
 +    let mut decoder = Decoder::new(response);
 +    decoder.read_i32().unwrap(); // throttle_time_ms
 +    let error = decoder.read_i16().unwrap();
 +    decoder.read_tagged_fields().unwrap();
 +    assert_eq!(
 +        decoder.remaining(),
 +        0,
 +        "Heartbeat response has trailing bytes"
 +    );
 +    error
 +}
 +
 +/// Claim a member id, then join with it. Returns the id and the second 
join's answer, which is
 +/// `None` while the group's join barrier is still open.
 +async fn claim_member_id(state: &GatewayState, metadata: &'static [u8]) -> 
String {
 +    let protocols: &[(&str, &[u8])] = &[("range", metadata)];
 +    let claimed = join(state, JOIN_VERSION, &join_params("", 
protocols)).await;
 +    assert_eq!(claimed.error, ERROR_MEMBER_ID_REQUIRED);
 +    claimed.member_id
 +}
 +
 +/// Two members through one full rebalance: generation 2, leader first, both 
awaiting `SyncGroup`.
 +async fn two_member_group(state: &Arc<GatewayState>) -> (String, String) {
 +    let leader_protocols: &[(&str, &[u8])] = &[("range", 
b"leader-subscription")];
 +
 +    let leader = claim_member_id(state, b"leader-subscription").await;
 +    let first = join(state, JOIN_VERSION, &join_params(&leader, 
leader_protocols)).await;
 +    assert_eq!(first.error, ERROR_NONE);
 +    assert_eq!(first.generation_id, 1);
 +
 +    let follower = claim_member_id(state, b"follower-subscription").await;
 +    let parked = {
 +        let state = Arc::clone(state);
 +        let follower = follower.clone();
 +        tokio::spawn(async move {
 +            let protocols: &[(&str, &[u8])] = &[("range", 
b"follower-subscription")];
 +            join(&state, JOIN_VERSION, &join_params(&follower, 
protocols)).await
 +        })
 +    };
 +    yield_to_parked().await;
 +
 +    let rejoined = join(state, JOIN_VERSION, &join_params(&leader, 
leader_protocols)).await;
 +    assert_eq!(rejoined.error, ERROR_NONE);
 +    assert_eq!(rejoined.generation_id, 2);
 +
 +    let follower_join = parked.await.expect("parked JoinGroup task");
 +    assert_eq!(follower_join.error, ERROR_NONE);
 +    assert_eq!(follower_join.generation_id, 2);
 +
 +    (leader, follower)
 +}
 +
 +/// Let a task that is about to park reach its park. Paused time turns this 
into one scheduler
 +/// turn plus a 1 ms clock step, not a real wait.
 +async fn yield_to_parked() {
 +    tokio::time::sleep(Duration::from_millis(1)).await;
 +}
 +
 +// ── Response decoding 
───────────────────────────────────────────────────────
 +
 +#[derive(Debug)]
 +struct JoinResponse {
 +    error: i16,
 +    generation_id: i32,
 +    protocol_type: Option<String>,
 +    protocol_name: Option<String>,
 +    leader: String,
 +    member_id: String,
 +    members: Vec<(String, Bytes)>,
 +}
 +
 +impl JoinResponse {
 +    fn decode(version: i16, body: Bytes) -> Self {
 +        let flexible = version >= 6;
 +        let mut decoder = Decoder::new(body);
 +        if version >= 2 {
 +            decoder.read_i32().unwrap(); // throttle_time_ms
 +        }
 +        let error = decoder.read_i16().unwrap();
 +        let generation_id = decoder.read_i32().unwrap();
 +        let protocol_type = if version >= 7 {
 +            decoder.read_compact_nullable_string().unwrap()
 +        } else {
 +            None
 +        };
 +        let protocol_name = read_nullable(&mut decoder, flexible);
 +        let leader = read_nullable(&mut decoder, flexible).expect("leader is 
not nullable");
 +        if version >= 9 {
 +            decoder.read_bool().unwrap(); // skip_assignment
 +        }
 +        let member_id = read_nullable(&mut decoder, flexible).expect("member 
id is not nullable");
 +
 +        let count = read_array_count(&mut decoder, flexible);
 +        let mut members = Vec::with_capacity(count);
 +        for _ in 0..count {
 +            let id = read_nullable(&mut decoder, flexible).expect("member id 
is not nullable");
 +            if version >= 5 {
 +                read_nullable(&mut decoder, flexible); // group_instance_id
 +            }
 +            let metadata = read_bytes_field(&mut decoder, flexible);
 +            if flexible {
 +                decoder.read_tagged_fields().unwrap();
 +            }
 +            members.push((id, metadata));
 +        }
 +        if flexible {
 +            decoder.read_tagged_fields().unwrap();
 +        }
 +        assert_eq!(
 +            decoder.remaining(),
 +            0,
 +            "JoinGroup v{version} response has trailing bytes"
 +        );
 +
 +        Self {
 +            error,
 +            generation_id,
 +            protocol_type,
 +            protocol_name,
 +            leader,
 +            member_id,
 +            members,
 +        }
 +    }
 +
 +    fn metadata_of(&self, member_id: &str) -> Option<&Bytes> {
 +        self.members
 +            .iter()
 +            .find(|(id, _)| id == member_id)
 +            .map(|(_, metadata)| metadata)
 +    }
 +}
 +
 +#[derive(Debug)]
 +struct SyncResponse {
 +    error: i16,
 +    protocol_name: Option<String>,
 +    assignment: Bytes,
 +}
 +
 +impl SyncResponse {
 +    fn decode(version: i16, body: Bytes) -> Self {
 +        let mut decoder = Decoder::new(body);
 +        if version >= 1 {
 +            decoder.read_i32().unwrap(); // throttle_time_ms
 +        }
 +        let error = decoder.read_i16().unwrap();
 +        let protocol_name = if version >= 5 {
 +            decoder.read_compact_nullable_string().unwrap(); // protocol_type
 +            decoder.read_compact_nullable_string().unwrap()
 +        } else {
 +            None
 +        };
 +        let assignment = read_bytes_field(&mut decoder, version >= 4);
 +        if version >= 4 {
 +            decoder.read_tagged_fields().unwrap();
 +        }
 +        assert_eq!(
 +            decoder.remaining(),
 +            0,
 +            "SyncGroup v{version} response has trailing bytes"
 +        );
 +        Self {
 +            error,
 +            protocol_name,
 +            assignment,
 +        }
 +    }
 +}
 +
 +fn read_nullable(decoder: &mut Decoder, flexible: bool) -> Option<String> {
 +    if flexible {
 +        decoder.read_compact_nullable_string().unwrap()
 +    } else {
 +        decoder.read_nullable_string().unwrap()
 +    }
 +}
 +
 +fn read_array_count(decoder: &mut Decoder, flexible: bool) -> usize {
 +    if flexible {
 +        usize::try_from(decoder.read_varint().unwrap() - 1).expect("count 
fits usize")
 +    } else {
 +        usize::try_from(decoder.read_i32().unwrap()).expect("count fits 
usize")
 +    }
 +}
 +
 +fn read_bytes_field(decoder: &mut Decoder, flexible: bool) -> Bytes {
 +    if flexible {
 +        decoder.read_compact_nullable_bytes().unwrap()
 +    } else {
 +        decoder.read_nullable_bytes().unwrap()
 +    }
 +    .expect("bytes field is not nullable")
 +}
 +
 +// ── JoinGroup 
───────────────────────────────────────────────────────────────
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_an_empty_member_id_when_joining_at_v9_should_require_a_member_id() {
 +    let state = test_state(immediate_config());
 +    let protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +
 +    let response = join(&state, 9, &join_params("", protocols)).await;
 +
 +    assert_eq!(response.error, ERROR_MEMBER_ID_REQUIRED);
 +    assert!(!response.member_id.is_empty());
 +    assert_eq!(response.generation_id, -1);
 +    assert!(response.leader.is_empty());
 +    assert!(response.members.is_empty());
 +    assert_eq!(
 +        response.protocol_name, None,
 +        "v9 encodes an absent protocol name as null"
 +    );
 +}
 +
 +/// `protocol_name` only became nullable on the wire at v7. Below that a 
client reads it as a
 +/// non-nullable string, so an error response has to carry an empty one.
 +#[tokio::test(start_paused = true)]
 +async fn 
given_an_empty_member_id_when_joining_at_v6_should_send_an_empty_protocol_name()
 {
 +    let state = test_state(immediate_config());
 +    let protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +
 +    let response = join(&state, 6, &join_params("", protocols)).await;
 +
 +    assert_eq!(response.error, ERROR_MEMBER_ID_REQUIRED);
 +    assert_eq!(response.protocol_name.as_deref(), Some(""));
 +}
 +
 +/// KIP-394 arrived in v4, so an older client never sees `MEMBER_ID_REQUIRED` 
and is admitted on
 +/// its first request.
 +#[tokio::test(start_paused = true)]
 +async fn 
given_an_empty_member_id_when_joining_at_v3_should_admit_the_member_directly() {
 +    let state = test_state(immediate_config());
 +    let protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +
 +    let response = join(&state, 3, &join_params("", protocols)).await;
 +
 +    assert_eq!(response.error, ERROR_NONE);
 +    assert_eq!(response.generation_id, 1);
 +    assert!(!response.member_id.is_empty());
 +    assert_eq!(response.leader, response.member_id);
 +    assert_eq!(response.protocol_name.as_deref(), Some("range"));
 +    assert_eq!(response.members.len(), 1);
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_two_members_when_both_join_should_send_the_roster_only_to_the_leader() {
 +    let state = test_state(immediate_config());
 +    let leader_protocols: &[(&str, &[u8])] = &[("range", 
b"leader-subscription")];
 +
 +    let leader = claim_member_id(&state, b"leader-subscription").await;
 +    let first = join(
 +        &state,
 +        JOIN_VERSION,
 +        &join_params(&leader, leader_protocols),
 +    )
 +    .await;
 +    assert_eq!(first.generation_id, 1);
 +    assert_eq!(first.members.len(), 1);
 +
 +    let follower = claim_member_id(&state, b"follower-subscription").await;
 +    let parked = {
 +        let state = Arc::clone(&state);
 +        let follower = follower.clone();
 +        tokio::spawn(async move {
 +            let protocols: &[(&str, &[u8])] = &[("range", 
b"follower-subscription")];
 +            join(&state, JOIN_VERSION, &join_params(&follower, 
protocols)).await
 +        })
 +    };
 +    yield_to_parked().await;
 +    assert!(
 +        !parked.is_finished(),
 +        "the follower must wait for the leader to rejoin"
 +    );
 +
 +    let rejoined = join(
 +        &state,
 +        JOIN_VERSION,
 +        &join_params(&leader, leader_protocols),
 +    )
 +    .await;
 +    let follower_join = parked.await.expect("parked JoinGroup task");
 +
 +    assert_eq!(rejoined.generation_id, 2);
 +    assert_eq!(rejoined.leader, leader);
 +    assert_eq!(rejoined.members.len(), 2);
 +    assert_eq!(
 +        rejoined.metadata_of(&leader).map(Bytes::as_ref),
 +        Some(b"leader-subscription".as_slice())
 +    );
 +    assert_eq!(
 +        rejoined.metadata_of(&follower).map(Bytes::as_ref),
 +        Some(b"follower-subscription".as_slice()),
 +        "the leader must receive each member's own subscription bytes, 
verbatim"
 +    );
 +    assert_eq!(follower_join.generation_id, 2);
 +    assert_eq!(follower_join.leader, leader);
 +    assert!(
 +        follower_join.members.is_empty(),
 +        "a follower runs no assignor and must not receive the roster"
 +    );
 +    assert_eq!(follower_join.protocol_type.as_deref(), Some("consumer"));
 +}
 +
 +// ── SyncGroup 
───────────────────────────────────────────────────────────────
 +
 +/// Acceptance criterion: two consumers in one group receive a disjoint 
assignment. The gateway's
 +/// share of that is the relay - each member is handed back exactly the bytes 
the leader filed
 +/// under its id, and nothing else.
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_leader_assignment_when_syncing_should_deliver_each_member_its_own_blob()
 {
 +    let state = test_state(immediate_config());
 +    let (leader, follower) = two_member_group(&state).await;
 +
 +    let parked = {
 +        let state = Arc::clone(&state);
 +        let follower = follower.clone();
 +        tokio::spawn(async move {
 +            sync(
 +                &state,
 +                SYNC_VERSION,
 +                &SyncGroupParams {
 +                    group_id: GROUP,
 +                    generation_id: 2,
 +                    member_id: &follower,
 +                    ..SyncGroupParams::default()
 +                },
 +            )
 +            .await
 +        })
 +    };
 +    yield_to_parked().await;
 +    assert!(
 +        !parked.is_finished(),
 +        "a follower must wait for the leader's assignment"
 +    );
 +
 +    let leader_blob: &[u8] = b"partitions-0-1";
 +    let follower_blob: &[u8] = b"partitions-2-3";
 +    let assignments: &[(&str, &[u8])] = &[
 +        (leader.as_str(), leader_blob),
 +        (follower.as_str(), follower_blob),
 +    ];
 +    let leader_sync = sync(
 +        &state,
 +        SYNC_VERSION,
 +        &SyncGroupParams {
 +            group_id: GROUP,
 +            generation_id: 2,
 +            member_id: &leader,
 +            assignments,
 +            ..SyncGroupParams::default()
 +        },
 +    )
 +    .await;
 +    let follower_sync = parked.await.expect("parked SyncGroup task");
 +
 +    assert_eq!(leader_sync.error, ERROR_NONE);
 +    assert_eq!(follower_sync.error, ERROR_NONE);
 +    assert_eq!(leader_sync.assignment.as_ref(), leader_blob);
 +    assert_eq!(follower_sync.assignment.as_ref(), follower_blob);
 +    assert_ne!(leader_sync.assignment, follower_sync.assignment);
 +    assert_eq!(leader_sync.protocol_name.as_deref(), Some("range"));
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_member_the_leader_omitted_when_syncing_should_return_an_empty_assignment()
 {
 +    let state = test_state(immediate_config());
 +    let (leader, follower) = two_member_group(&state).await;
 +
 +    let assignments: &[(&str, &[u8])] = &[(leader.as_str(), b"everything")];
 +    let leader_sync = sync(
 +        &state,
 +        SYNC_VERSION,
 +        &SyncGroupParams {
 +            group_id: GROUP,
 +            generation_id: 2,
 +            member_id: &leader,
 +            assignments,
 +            ..SyncGroupParams::default()
 +        },
 +    )
 +    .await;
 +    let follower_sync = sync(
 +        &state,
 +        SYNC_VERSION,
 +        &SyncGroupParams {
 +            group_id: GROUP,
 +            generation_id: 2,
 +            member_id: &follower,
 +            ..SyncGroupParams::default()
 +        },
 +    )
 +    .await;
 +
 +    assert_eq!(leader_sync.assignment.as_ref(), b"everything");
 +    assert_eq!(follower_sync.error, ERROR_NONE);
 +    assert!(follower_sync.assignment.is_empty());
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_wrong_protocol_name_when_syncing_at_v5_should_return_inconsistent_group_protocol()
 +{
 +    let state = test_state(immediate_config());
 +    let (leader, _) = two_member_group(&state).await;
 +
 +    let response = sync(
 +        &state,
 +        SYNC_VERSION,
 +        &SyncGroupParams {
 +            group_id: GROUP,
 +            generation_id: 2,
 +            member_id: &leader,
 +            protocol_type: Some("consumer"),
 +            protocol_name: Some("sticky"),
 +            ..SyncGroupParams::default()
 +        },
 +    )
 +    .await;
 +
 +    assert_eq!(response.error, ERROR_INCONSISTENT_GROUP_PROTOCOL);
 +    assert!(response.assignment.is_empty());
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_stale_generation_when_syncing_should_return_illegal_generation() {
 +    let state = test_state(immediate_config());
 +    let (leader, _) = two_member_group(&state).await;
 +
 +    let response = sync(
 +        &state,
 +        SYNC_VERSION,
 +        &SyncGroupParams {
 +            group_id: GROUP,
 +            generation_id: 1,
 +            member_id: &leader,
 +            ..SyncGroupParams::default()
 +        },
 +    )
 +    .await;
 +
 +    assert_eq!(response.error, ERROR_ILLEGAL_GENERATION);
 +}
 +
 +// ── Heartbeat 
───────────────────────────────────────────────────────────────
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_stable_group_when_a_new_member_joins_should_answer_heartbeat_rebalance_in_progress()
 + {
 +    let state = test_state(immediate_config());
 +    let leader_protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +
 +    let leader = claim_member_id(&state, b"sub").await;
 +    join(
 +        &state,
 +        JOIN_VERSION,
 +        &join_params(&leader, leader_protocols),
 +    )
 +    .await;
 +    sync(
 +        &state,
 +        SYNC_VERSION,
 +        &SyncGroupParams {
 +            group_id: GROUP,
 +            generation_id: 1,
 +            member_id: &leader,
 +            assignments: &[],
 +            ..SyncGroupParams::default()
 +        },
 +    )
 +    .await;
 +    assert_eq!(heartbeat(&state, 1, &leader).await, ERROR_NONE);
 +
 +    let follower = claim_member_id(&state, b"sub").await;
 +    let parked = {
 +        let state = Arc::clone(&state);
 +        let follower = follower.clone();
 +        tokio::spawn(async move {
 +            let protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +            join(&state, JOIN_VERSION, &join_params(&follower, 
protocols)).await
 +        })
 +    };
 +    yield_to_parked().await;
 +
 +    assert_eq!(
 +        heartbeat(&state, 1, &leader).await,
 +        ERROR_REBALANCE_IN_PROGRESS,
 +        "a stable member learns about a rebalance through its heartbeat"
 +    );
 +
 +    let rejoined = join(
 +        &state,
 +        JOIN_VERSION,
 +        &join_params(&leader, leader_protocols),
 +    )
 +    .await;
 +    let follower_join = parked.await.expect("parked JoinGroup task");
 +    assert_eq!(rejoined.generation_id, 2);
 +    assert_eq!(follower_join.generation_id, 2);
 +}
 +
 +/// Acceptance criterion: a heartbeat timeout evicts the member and triggers 
a rebalance.
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_missed_heartbeat_when_the_session_expires_should_evict_and_bump_the_generation()
 {
 +    let state = test_state(immediate_config());
 +    let leader_protocols: &[(&str, &[u8])] = &[("range", 
b"leader-subscription")];
 +    let (leader, follower) = two_member_group(&state).await;
 +
 +    // Both sessions are still alive here; the leader refreshes its own, the 
follower goes quiet.
 +    advance(Duration::from_secs(6)).await;
 +    assert_eq!(heartbeat(&state, 2, &leader).await, ERROR_NONE);
 +
 +    // Past the follower's deadline but not the leader's refreshed one.
 +    advance(Duration::from_secs(5)).await;
 +    assert_eq!(
 +        heartbeat(&state, 2, &leader).await,
 +        ERROR_REBALANCE_IN_PROGRESS,
 +        "the expired follower must be evicted and a rebalance started"
 +    );
 +
 +    let rejoined = join(
 +        &state,
 +        JOIN_VERSION,
 +        &join_params(&leader, leader_protocols),
 +    )
 +    .await;
 +    assert_eq!(rejoined.error, ERROR_NONE);
 +    assert_eq!(rejoined.generation_id, 3);
 +    assert_eq!(rejoined.members.len(), 1);
 +    assert_eq!(rejoined.metadata_of(&follower), None);
 +
 +    assert_eq!(
 +        heartbeat(&state, 2, &follower).await,
 +        ERROR_UNKNOWN_MEMBER_ID,
 +        "the evicted member must be told to rejoin from scratch"
 +    );
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_stale_generation_when_heartbeating_should_return_illegal_generation() {
 +    let state = test_state(immediate_config());
 +    let (leader, _) = two_member_group(&state).await;
 +
 +    assert_eq!(
 +        heartbeat(&state, 1, &leader).await,
 +        ERROR_ILLEGAL_GENERATION
 +    );
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_an_unknown_group_when_heartbeating_should_return_unknown_member_id() {
 +    let state = test_state(immediate_config());
 +
 +    assert_eq!(heartbeat(&state, 1, "ghost").await, ERROR_UNKNOWN_MEMBER_ID);
 +}
 +
 +/// Every member's session lapsed while nobody was talking to the group. The 
next joiner evicts
 +/// them on arrival instead of waiting out a rebalance window for members 
that will never rejoin:
 +/// the group it lands in is brand new, so its generation is 1 and not 2.
 +#[tokio::test(start_paused = true)]
 +async fn 
given_every_member_expired_when_a_new_member_joins_should_start_a_fresh_group() 
{
 +    let state = test_state(immediate_config());
 +    let protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +
 +    let first = join(&state, 3, &join_params("", protocols)).await;
 +    assert_eq!(first.generation_id, 1);
 +
 +    advance(Duration::from_secs(11)).await;
 +    let second = join(&state, 3, &join_params("", protocols)).await;
 +
 +    assert_eq!(second.error, ERROR_NONE);
 +    assert_eq!(second.generation_id, 1);
 +    assert_eq!(second.leader, second.member_id);
 +    assert_eq!(second.members.len(), 1);
 +    assert_ne!(second.member_id, first.member_id);
 +}
 +
 +// ── Rejected requests 
───────────────────────────────────────────────────────
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_session_timeout_out_of_range_when_joining_should_return_invalid_session_timeout()
 {
 +    let state = test_state(immediate_config());
 +    let protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +
 +    let too_short = join(
 +        &state,
 +        JOIN_VERSION,
 +        &JoinGroupParams {
 +            session_timeout_ms: 1_000,
 +            ..join_params("", protocols)
 +        },
 +    )
 +    .await;
 +    let too_long = join(
 +        &state,
 +        JOIN_VERSION,
 +        &JoinGroupParams {
 +            session_timeout_ms: 3_600_000,
 +            ..join_params("", protocols)
 +        },
 +    )
 +    .await;
 +
 +    assert_eq!(too_short.error, ERROR_INVALID_SESSION_TIMEOUT);
 +    assert_eq!(too_long.error, ERROR_INVALID_SESSION_TIMEOUT);
 +}
 +
 +/// 246 bytes is what an Iggy name leaves for a group id once the `kafka.cg.` 
offset-key prefix
 +/// is accounted for, so a longer one could never have its offsets committed.
 +#[tokio::test(start_paused = true)]
 +async fn 
given_an_out_of_range_group_id_when_joining_should_return_invalid_group_id() {
 +    let state = test_state(immediate_config());
 +    let protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +    let longest = "g".repeat(246);
 +    let too_long = "g".repeat(247);
 +
 +    let empty = join(
 +        &state,
 +        JOIN_VERSION,
 +        &JoinGroupParams {
 +            group_id: "",
 +            ..join_params("", protocols)
 +        },
 +    )
 +    .await;
 +    let oversized = join(
 +        &state,
 +        JOIN_VERSION,
 +        &JoinGroupParams {
 +            group_id: &too_long,
 +            ..join_params("", protocols)
 +        },
 +    )
 +    .await;
 +    let accepted = join(
 +        &state,
 +        JOIN_VERSION,
 +        &JoinGroupParams {
 +            group_id: &longest,
 +            ..join_params("", protocols)
 +        },
 +    )
 +    .await;
 +
 +    assert_eq!(empty.error, ERROR_INVALID_GROUP_ID);
 +    assert_eq!(oversized.error, ERROR_INVALID_GROUP_ID);
 +    assert_eq!(accepted.error, ERROR_MEMBER_ID_REQUIRED);
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_conflicting_protocol_type_when_joining_should_return_inconsistent_group_protocol()
 +{
 +    let state = test_state(immediate_config());
 +    let protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +    join(&state, 3, &join_params("", protocols)).await;
 +
 +    let response = join(
 +        &state,
 +        JOIN_VERSION,
 +        &JoinGroupParams {
 +            protocol_type: "connect",
 +            ..join_params("", protocols)
 +        },
 +    )
 +    .await;
 +
 +    assert_eq!(response.error, ERROR_INCONSISTENT_GROUP_PROTOCOL);
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_no_shared_protocol_when_joining_should_return_inconsistent_group_protocol()
 {
 +    let state = test_state(immediate_config());
 +    let first: &[(&str, &[u8])] = &[("range", b"sub")];
 +    let second: &[(&str, &[u8])] = &[("sticky", b"sub")];
 +    join(&state, 3, &join_params("", first)).await;
 +
 +    let response = join(&state, 3, &join_params("", second)).await;
 +
 +    assert_eq!(response.error, ERROR_INCONSISTENT_GROUP_PROTOCOL);
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_group_at_its_size_cap_when_joining_should_return_group_max_size_reached()
 {
 +    let state = test_state(GroupCoordinatorConfig {
 +        max_members_per_group: 1,
 +        ..immediate_config()
 +    });
 +    let protocols: &[(&str, &[u8])] = &[("range", b"sub")];
 +
 +    let admitted = join(&state, 3, &join_params("", protocols)).await;
 +    let rejected = join(&state, 3, &join_params("", protocols)).await;
 +
 +    assert_eq!(admitted.error, ERROR_NONE);
 +    assert_eq!(rejected.error, ERROR_GROUP_MAX_SIZE_REACHED);
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_an_oversized_subscription_when_joining_should_return_invalid_request() {
 +    let state = test_state(GroupCoordinatorConfig {
 +        max_member_blob_bytes: 8,
 +        ..immediate_config()
 +    });
 +    let protocols: &[(&str, &[u8])] = &[("range", 
b"a-subscription-over-eight-bytes")];
 +
 +    let response = join(&state, 3, &join_params("", protocols)).await;
 +
 +    assert_eq!(response.error, ERROR_INVALID_REQUEST);
 +}
 +
 +// ── FindCoordinator 
─────────────────────────────────────────────────────────
 +
 +async fn find_coordinator(
 +    state: &GatewayState,
 +    version: i16,
 +    keys: &[&str],
 +    key_type: i8,
 +) -> Decoder {
 +    let body = build_find_coordinator_request(version, keys, key_type);
 +    let response = handle_request_bounded(state, API_KEY_FIND_COORDINATOR, 
version, body)
 +        .await
 +        .expect_response("FindCoordinator must answer");
 +    Decoder::new(response)
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_group_key_when_finding_the_coordinator_at_v0_should_return_this_broker()
 {
 +    let state = test_state(immediate_config());
 +    let broker = BrokerAdvertise::default();
 +
 +    let mut decoder = find_coordinator(&state, 0, &[GROUP], 0).await;
 +
 +    assert_eq!(decoder.read_i16().unwrap(), ERROR_NONE);
 +    assert_eq!(decoder.read_i32().unwrap(), 1, "node id");
 +    assert_eq!(decoder.read_nullable_string().unwrap(), Some(broker.host));
 +    assert_eq!(decoder.read_i32().unwrap(), broker.port);
 +    assert_eq!(decoder.remaining(), 0);
 +}
 +
 +#[tokio::test(start_paused = true)]
 +async fn 
given_several_keys_when_finding_the_coordinator_at_v4_should_return_one_entry_each()
 {
 +    let state = test_state(immediate_config());
 +    let keys = ["alpha", "beta", "gamma"];
 +
 +    let mut decoder = find_coordinator(&state, 4, &keys, 0).await;
 +
 +    decoder.read_i32().unwrap(); // throttle_time_ms
 +    assert_eq!(read_array_count(&mut decoder, true), keys.len());
 +    for key in keys {
 +        assert_eq!(
 +            decoder.read_compact_nullable_string().unwrap().as_deref(),
 +            Some(key)
 +        );
 +        assert_eq!(decoder.read_i32().unwrap(), 1, "node id");
 +        decoder.read_compact_nullable_string().unwrap(); // host
 +        decoder.read_i32().unwrap(); // port
 +        assert_eq!(decoder.read_i16().unwrap(), ERROR_NONE);
 +        assert_eq!(decoder.read_compact_nullable_string().unwrap(), None);
 +        decoder.read_tagged_fields().unwrap();
 +    }
 +    decoder.read_tagged_fields().unwrap();
 +    assert_eq!(decoder.remaining(), 0);
 +}
 +
 +/// Transactions are out of scope permanently, so the answer is the code both 
the Java client and
 +/// librdkafka treat as fatal, rather than one a transactional producer would 
spin on forever.
 +#[tokio::test(start_paused = true)]
 +async fn 
given_a_transaction_key_type_when_finding_the_coordinator_should_return_transactional_id_authorization_failed()
 + {
 +    let state = test_state(immediate_config());
 +
 +    let mut decoder = find_coordinator(&state, 1, &["txn"], 1).await;
 +
 +    decoder.read_i32().unwrap(); // throttle_time_ms
 +    assert_eq!(
 +        decoder.read_i16().unwrap(),
 +        ERROR_TRANSACTIONAL_ID_AUTHORIZATION_FAILED
 +    );
 +    assert!(decoder.read_nullable_string().unwrap().is_some());
 +    assert_eq!(decoder.read_i32().unwrap(), -1, "node id");
 +    assert_eq!(decoder.read_nullable_string().unwrap().as_deref(), Some(""));
 +    assert_eq!(decoder.read_i32().unwrap(), -1, "port");
 +    assert_eq!(decoder.remaining(), 0);
 +}
 +
 +// ── Over a real TCP listener 
────────────────────────────────────────────────
 +
 +/// The in-process tests drive one `GatewayState` directly. This one proves 
the connection loop
 +/// itself parks and resumes: a follower's `SyncGroup` sits unanswered on its 
own socket until
 +/// the leader's arrives on another.
 +#[tokio::test]
 +async fn 
given_two_tcp_clients_when_they_join_and_sync_should_each_receive_their_own_assignment()
 {
 +    let (addr, _shutdown) = spawn_test_server_with_config(GatewayConfig {
 +        group: immediate_config(),
 +        ..GatewayConfig::default()
 +    })
 +    .await;
 +
 +    let mut leader_stream = TcpStream::connect(addr).await.expect("connect 
leader");
 +    let mut follower_stream = TcpStream::connect(addr).await.expect("connect 
follower");
 +
 +    let leader_protocols: &[(&str, &[u8])] = &[("range", 
b"leader-subscription")];
 +    let follower_protocols: &[(&str, &[u8])] = &[("range", 
b"follower-subscription")];
 +
 +    let leader = tcp_join(&mut leader_stream, &join_params("", 
leader_protocols))
 +        .await
 +        .member_id;
 +    let first = tcp_join(&mut leader_stream, &join_params(&leader, 
leader_protocols)).await;
 +    assert_eq!(first.generation_id, 1);
 +
 +    let follower = tcp_join(&mut follower_stream, &join_params("", 
follower_protocols))
 +        .await
 +        .member_id;
 +    // The follower's rejoin blocks until the leader rejoins, so only the 
request is written here.
 +    write_request(
 +        &mut follower_stream,
 +        API_KEY_JOIN_GROUP,
 +        JOIN_VERSION,
 +        2,
 +        &build_join_group_request(JOIN_VERSION, &join_params(&follower, 
follower_protocols)),
 +    )
 +    .await;
 +
 +    let rejoined = tcp_join(&mut leader_stream, &join_params(&leader, 
leader_protocols)).await;
 +    assert_eq!(rejoined.generation_id, 2);
 +    assert_eq!(rejoined.members.len(), 2);
 +    let follower_join = read_join(&mut follower_stream).await;
 +    assert_eq!(follower_join.generation_id, 2);
 +    assert!(follower_join.members.is_empty());
 +
 +    write_request(
 +        &mut follower_stream,
 +        API_KEY_SYNC_GROUP,
 +        SYNC_VERSION,
 +        3,
 +        &build_sync_group_request(
 +            SYNC_VERSION,
 +            &SyncGroupParams {
 +                group_id: GROUP,
 +                generation_id: 2,
 +                member_id: &follower,
 +                ..SyncGroupParams::default()
 +            },
 +        ),
 +    )
 +    .await;
 +
 +    let leader_blob: &[u8] = b"partitions-0-1";
 +    let follower_blob: &[u8] = b"partitions-2-3";
 +    let assignments: &[(&str, &[u8])] = &[
 +        (leader.as_str(), leader_blob),
 +        (follower.as_str(), follower_blob),
 +    ];
 +    write_request(
 +        &mut leader_stream,
 +        API_KEY_SYNC_GROUP,
 +        SYNC_VERSION,
 +        4,
 +        &build_sync_group_request(
 +            SYNC_VERSION,
 +            &SyncGroupParams {
 +                group_id: GROUP,
 +                generation_id: 2,
 +                member_id: &leader,
 +                assignments,
 +                ..SyncGroupParams::default()
 +            },
 +        ),
 +    )
 +    .await;
 +
 +    let leader_sync = read_sync(&mut leader_stream).await;
 +    let follower_sync = read_sync(&mut follower_stream).await;
 +
 +    assert_eq!(leader_sync.assignment.as_ref(), leader_blob);
 +    assert_eq!(follower_sync.assignment.as_ref(), follower_blob);
 +}
 +
 +async fn write_request(
 +    stream: &mut TcpStream,
 +    api_key: i16,
 +    api_version: i16,
 +    correlation_id: i32,
 +    body: &Bytes,
 +) {
 +    let frame = build_request_frame(
 +        api_key,
 +        api_version,
 +        correlation_id,
 +        Some("consumer-group-test"),
 +        body,
 +    );
 +    stream.write_all(&frame).await.expect("write request");
 +}
 +
 +async fn tcp_join(stream: &mut TcpStream, params: &JoinGroupParams<'_>) -> 
JoinResponse {
 +    write_request(
 +        stream,
 +        API_KEY_JOIN_GROUP,
 +        JOIN_VERSION,
 +        1,
 +        &build_join_group_request(JOIN_VERSION, params),
 +    )
 +    .await;
 +    read_join(stream).await
 +}
 +
 +async fn read_join(stream: &mut TcpStream) -> JoinResponse {
 +    let payload = read_response_frame(stream, 8 * 1024 * 1024).await;
 +    let (_, body) = parse_response_payload(API_KEY_JOIN_GROUP, JOIN_VERSION, 
payload);
 +    JoinResponse::decode(JOIN_VERSION, body)
 +}
 +
 +async fn read_sync(stream: &mut TcpStream) -> SyncResponse {
 +    let payload = read_response_frame(stream, 8 * 1024 * 1024).await;
 +    let (_, body) = parse_response_payload(API_KEY_SYNC_GROUP, SYNC_VERSION, 
payload);
 +    SyncResponse::decode(SYNC_VERSION, body)
 +}
diff --cc gateways/kafka/tests/golden_wire_fixtures_tests.rs
index 5d73f1936,ebbd9568a..90401a88d
--- a/gateways/kafka/tests/golden_wire_fixtures_tests.rs
+++ b/gateways/kafka/tests/golden_wire_fixtures_tests.rs
@@@ -41,21 -41,23 +41,21 @@@ async fn golden_apiversions_v3_flexible
          .await
          .expect_response("test request has acks != 0 and expects a response");
  
 -    // error_code=0, api_count=6 (compact array: N+1=7)
 -    // key 0  (Produce)      min=0  max=9  (advertised)
 -    // key 1  (Fetch)        min=4  max=12
 -    // key 2  (ListOffsets)  min=1  max=6
 -    // key 3  (Metadata)     min=0  max=9
 -    // key 18 (ApiVersions)  min=0  max=3
 -    // key 19 (CreateTopics) min=2  max=5
 +    // error_code=0, api_count=10 (compact array: N+1=11)
      // each entry followed by an empty tagged-fields byte; throttle_ms=0; 
top-level tagged fields
 -    let expected: [u8; 50] = [
 +    let expected: [u8; 78] = [
          0x00, 0x00, // error_code
 -        0x07, // compact array count (6+1)
 -        0x00, 0x00, 0x00, 0x00, 0x00, 0x09, 0x00, // key 0:  Produce      0-9 
(advertised)
 -        0x00, 0x01, 0x00, 0x04, 0x00, 0x0C, 0x00, // key 1:  Fetch        4-12
 -        0x00, 0x02, 0x00, 0x01, 0x00, 0x06, 0x00, // key 2:  ListOffsets  1-6
 -        0x00, 0x03, 0x00, 0x00, 0x00, 0x09, 0x00, // key 3:  Metadata     0-9
 -        0x00, 0x12, 0x00, 0x00, 0x00, 0x03, 0x00, // key 18: ApiVersions  0-3
 -        0x00, 0x13, 0x00, 0x02, 0x00, 0x05, 0x00, // key 19: CreateTopics 2-5
 +        0x0B, // compact array count (10+1)
 +        0x00, 0x00, 0x00, 0x00, 0x00, 0x09, 0x00, // key 0:  Produce         
0-9 (advertised)
 +        0x00, 0x01, 0x00, 0x04, 0x00, 0x0C, 0x00, // key 1:  Fetch           
4-12
 +        0x00, 0x02, 0x00, 0x01, 0x00, 0x06, 0x00, // key 2:  ListOffsets     
1-6
 +        0x00, 0x03, 0x00, 0x00, 0x00, 0x09, 0x00, // key 3:  Metadata        
0-9
-         0x00, 0x12, 0x00, 0x00, 0x00, 0x03, 0x00, // key 18: ApiVersions     
0-3
-         0x00, 0x13, 0x00, 0x02, 0x00, 0x05, 0x00, // key 19: CreateTopics    
2-5
 +        0x00, 0x0A, 0x00, 0x00, 0x00, 0x04, 0x00, // key 10: FindCoordinator 
0-4
 +        0x00, 0x0B, 0x00, 0x00, 0x00, 0x09, 0x00, // key 11: JoinGroup       
0-9
 +        0x00, 0x0C, 0x00, 0x00, 0x00, 0x04, 0x00, // key 12: Heartbeat       
0-4
 +        0x00, 0x0E, 0x00, 0x00, 0x00, 0x05, 0x00, // key 14: SyncGroup       
0-5
++        0x00, 0x12, 0x00, 0x00, 0x00, 0x03, 0x00, // key 18: ApiVersions     
0-3
++        0x00, 0x13, 0x00, 0x02, 0x00, 0x05, 0x00, // key 19: CreateTopics    
2-5
          0x00, 0x00, 0x00, 0x00, // throttle_ms
          0x00, // top-level tagged fields
      ];
@@@ -69,20 -71,23 +69,20 @@@ async fn golden_apiversions_v1_response
          .await
          .expect_response("test request has acks != 0 and expects a response");
  
 -    // error_code=0, api_count=6
 -    // key 0  (Produce)      min=0  max=9 (KAFKA-18659 advertise min=0)
 -    // key 1  (Fetch)        min=4  max=12
 -    // key 2  (ListOffsets)  min=1  max=6
 -    // key 3  (Metadata)     min=0  max=9
 -    // key 18 (ApiVersions)  min=0  max=3
 -    // key 19 (CreateTopics) min=2  max=5
 -    // throttle_ms=0
 -    let expected: [u8; 46] = [
 +    // error_code=0, api_count=10; Produce advertises min=0 per KAFKA-18659; 
throttle_ms=0
 +    let expected: [u8; 70] = [
          0x00, 0x00, // error_code
 -        0x00, 0x00, 0x00, 0x06, // api count = 6
 -        0x00, 0x00, 0x00, 0x00, 0x00, 0x09, // key 0:  Produce      0–9 
(advertised)
 -        0x00, 0x01, 0x00, 0x04, 0x00, 0x0C, // key 1:  Fetch        4–12
 -        0x00, 0x02, 0x00, 0x01, 0x00, 0x06, // key 2:  ListOffsets  1–6
 -        0x00, 0x03, 0x00, 0x00, 0x00, 0x09, // key 3:  Metadata     0–9
 -        0x00, 0x12, 0x00, 0x00, 0x00, 0x03, // key 18: ApiVersions  0–3
 -        0x00, 0x13, 0x00, 0x02, 0x00, 0x05, // key 19: CreateTopics 2–5
 +        0x00, 0x00, 0x00, 0x0A, // api count = 10
 +        0x00, 0x00, 0x00, 0x00, 0x00, 0x09, // key 0:  Produce         0-9 
(advertised)
 +        0x00, 0x01, 0x00, 0x04, 0x00, 0x0C, // key 1:  Fetch           4-12
 +        0x00, 0x02, 0x00, 0x01, 0x00, 0x06, // key 2:  ListOffsets     1-6
 +        0x00, 0x03, 0x00, 0x00, 0x00, 0x09, // key 3:  Metadata        0-9
-         0x00, 0x12, 0x00, 0x00, 0x00, 0x03, // key 18: ApiVersions     0-3
-         0x00, 0x13, 0x00, 0x02, 0x00, 0x05, // key 19: CreateTopics    2-5
 +        0x00, 0x0A, 0x00, 0x00, 0x00, 0x04, // key 10: FindCoordinator 0-4
 +        0x00, 0x0B, 0x00, 0x00, 0x00, 0x09, // key 11: JoinGroup       0-9
 +        0x00, 0x0C, 0x00, 0x00, 0x00, 0x04, // key 12: Heartbeat       0-4
 +        0x00, 0x0E, 0x00, 0x00, 0x00, 0x05, // key 14: SyncGroup       0-5
++        0x00, 0x12, 0x00, 0x00, 0x00, 0x03, // key 18: ApiVersions     0-3
++        0x00, 0x13, 0x00, 0x02, 0x00, 0x05, // key 19: CreateTopics    2-5
          0x00, 0x00, 0x00, 0x00, // throttle_ms
      ];
      assert_eq!(actual.as_ref(), &expected);
diff --cc gateways/kafka/tests/list_offsets_real_bridge_tests.rs
index 6ad89d53b,81ed93208..b4ea2755a
--- a/gateways/kafka/tests/list_offsets_real_bridge_tests.rs
+++ b/gateways/kafka/tests/list_offsets_real_bridge_tests.rs
@@@ -174,7 -172,7 +174,8 @@@ async fn connected_state(server: &TestS
          BrokerAdvertise::default(),
          Some(Arc::new(bridge)),
          TEST_MAX_FRAME_SIZE,
+         false,
 +        GroupCoordinator::new(GroupCoordinatorConfig::default(), 
CancellationToken::new()),
      );
      (state, seed_bridge)
  }
diff --cc gateways/kafka/tests/version_firewall_tests.rs
index 9e976d987,fb90877d2..3efedcf2b
--- a/gateways/kafka/tests/version_firewall_tests.rs
+++ b/gateways/kafka/tests/version_firewall_tests.rs
@@@ -348,7 -346,7 +348,7 @@@ async fn create_topics_below_min_versio
  
  #[tokio::test]
  async fn unsupported_api_keys_close_connection() {
-     for key in [8, 9, 13, 15, 17, 20, 42, 999] {
 -    for key in [8, 9, 10, 11, 20, 42, 999] {
++    for key in [8, 9, 13, 15, 20, 42, 999] {
          let outcome = handle_request(key, 0, Bytes::new(), 
&default_broker()).await;
          assert!(
              outcome.is_close(),

Reply via email to