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