This is an automated email from the ASF dual-hosted git repository. numinnex pushed a commit to branch kafka_gateway_leave_group in repository https://gitbox.apache.org/repos/asf/iggy.git
commit 40416637604f1873f263dfb58856bbb75c8f9969 Author: Grzegorz Koszyk <[email protected]> AuthorDate: Wed Sep 23 11:20:21 2026 +0200 feat(gateways): handle Kafka LeaveGroup for consumer groups --- gateways/kafka/README.md | 2 +- gateways/kafka/docs/CONSUMER_GROUPS.md | 78 +- gateways/kafka/docs/MANUAL_TESTING.md | 12 +- gateways/kafka/docs/SCOPE.md | 9 +- gateways/kafka/docs/TEST_SUITE.md | 2 +- gateways/kafka/docs/kafka_api_keys_reference.md | 7 +- gateways/kafka/scripts/ci-wire-fixtures.sh | 2 +- gateways/kafka/src/group/mod.rs | 150 +++- gateways/kafka/src/group/state.rs | 785 ++++++++++++++++++++- gateways/kafka/src/protocol/api.rs | 6 +- gateways/kafka/src/protocol/bounds_guard.rs | 125 +++- .../kafka/src/protocol/handlers/leave_group.rs | 211 ++++++ gateways/kafka/src/protocol/handlers/mod.rs | 9 +- gateways/kafka/tests/common/scope.rs | 1 + gateways/kafka/tests/common/wire.rs | 34 +- gateways/kafka/tests/consumer_group_tests.rs | 494 ++++++++++++- gateways/kafka/tests/golden_wire_fixtures_tests.rs | 14 +- gateways/kafka/tests/version_firewall_tests.rs | 20 +- gateways/kafka/tools/kafka-tool/src/main.rs | 67 +- 19 files changed, 1962 insertions(+), 66 deletions(-) diff --git a/gateways/kafka/README.md b/gateways/kafka/README.md index a1b59f09f..87b0a3495 100644 --- a/gateways/kafka/README.md +++ b/gateways/kafka/README.md @@ -4,7 +4,7 @@ Foundation layer for [apache/iggy#3421](https://github.com/apache/iggy/issues/34 > **Stub warning:** no API persists or reads real data yet. Produce, Fetch, > and ListOffsets return retriable `NOT_LEADER_OR_FOLLOWER` (6) so clients > keep data locally / retry elsewhere instead of trusting a fake success. > CreateTopics does **not** create topics; valid requests return > `NOT_CONTROLLER` (41). Metadata still reports requested topics as unknown. > Persistence lands with the Iggy bridge (see [docs/SCOPE.md](docs/SCOPE.md)). > -> Consumer group coordination is the exception: `FindCoordinator`, `JoinGroup`, `Heartbeat` and `SyncGroup` are real, with real membership, rebalances and session expiry ([docs/CONSUMER_GROUPS.md](docs/CONSUMER_GROUPS.md)). Offset commit/fetch is not, so a consumer can join a group and be assigned partitions but cannot yet consume. +> Consumer group coordination is the exception: `FindCoordinator`, `JoinGroup`, `Heartbeat`, `LeaveGroup` and `SyncGroup` are real, with real membership, rebalances, graceful leave and session expiry ([docs/CONSUMER_GROUPS.md](docs/CONSUMER_GROUPS.md)). Offset commit/fetch is not, so a consumer can join a group and be assigned partitions but cannot yet consume. ## Run diff --git a/gateways/kafka/docs/CONSUMER_GROUPS.md b/gateways/kafka/docs/CONSUMER_GROUPS.md index b609c5470..db3f723ab 100644 --- a/gateways/kafka/docs/CONSUMER_GROUPS.md +++ b/gateways/kafka/docs/CONSUMER_GROUPS.md @@ -2,7 +2,7 @@ What [#3541](https://github.com/apache/iggy/issues/3541) added: `FindCoordinator` (10), `JoinGroup` (11), `Heartbeat` (12) and `SyncGroup` (14), backed by an in-memory coordinator in -`src/group/`. This is Kafka's *classic* group protocol. Offsets are a separate concern and live +`src/group/`. [#3543](https://github.com/apache/iggy/issues/3543) added `LeaveGroup` (13). This is Kafka's *classic* group protocol. Offsets are a separate concern and live in Iggy ([`OFFSET_STORAGE.md`](OFFSET_STORAGE.md)). | API key | Name | Versions | Notes | @@ -10,6 +10,7 @@ in Iggy ([`OFFSET_STORAGE.md`](OFFSET_STORAGE.md)). | 10 | FindCoordinator | 0-4 | Always answers "this gateway", the same node the Metadata broker list advertises | | 11 | JoinGroup | 0-9 | Parks until the group's join barrier completes | | 12 | Heartbeat | 0-4 | Refreshes a session; `REBALANCE_IN_PROGRESS` is how a follower learns to rejoin | +| 13 | LeaveGroup | 0-5 | Removes members; survivors' parked joins and syncs are released at once | | 14 | SyncGroup | 0-5 | Relays the leader's assignment blobs; a follower parks until the leader syncs | `kafka-protocol` can encode FindCoordinator v5 and v6 as well, and they are byte-identical to v4. @@ -73,14 +74,71 @@ group, at which point the joiner clears the stale ids and completes immediately. observe the difference, because there is no consumer. The only cost is stale memory, bounded by the caps below. +## Generations + +The generation id increments only when a join barrier completes. Leaving never bumps it: the +survivors keep the current generation and learn about the rebalance through +`REBALANCE_IN_PROGRESS` (27) on their next heartbeat, which the Java client handles as a graceful +rejoin rather than as lost partitions. + +When a group empties, by leave or by expiry, it is dropped, so the next group under that name +starts again at generation 1. Kafka instead keeps the `Empty` group and its generation. This is +safe because member ids carry a UUID and every generation check is preceded by a membership +check, so a stale client is told `UNKNOWN_MEMBER_ID` before its generation is ever compared. + +## Leaving a group + +A dynamic consumer that closes cleanly sends `LeaveGroup`, and the survivors rebalance on their +next heartbeat instead of after `session.timeout.ms`. A survivor already parked in `JoinGroup` is +answered as soon as the leave completes the barrier. One parked in `SyncGroup` is answered +`REBALANCE_IN_PROGRESS` and rejoins. A request left parked by the leaver itself on another +connection is answered `UNKNOWN_MEMBER_ID`. + +| Client | Dynamic member on close | Static member on close | +| --- | --- | --- | +| Java consumer, classic protocol | Leaves (best effort) | Never leaves; Java 4.1+ `CloseOptions.LEAVE_GROUP` forces it | +| Kafka Streams, classic protocol | Does not leave by default | Does not leave | +| librdkafka / kcat | Leaves on `rd_kafka_consumer_close`, v0 or v1 only | Leaves on unsubscribe, not on termination | +| Admin `removeMembersFromConsumerGroup` | | Leaves by `group_instance_id` with an empty `member_id`, v3+ | + +How each identity in a request is resolved: + +| `member_id` | `group_instance_id` | Result | +| --- | --- | --- | +| id of a member | null | removed, 0 | +| id handed out with `MEMBER_ID_REQUIRED` and not yet claimed | any | id dropped, 0 | +| id of a member holding `X` | `X` | removed, 0 | +| id of any member | `X`, held by other members | `FENCED_INSTANCE_ID` (82) | +| any | `X`, held by nobody | `UNKNOWN_MEMBER_ID` (25) | +| `""` | `X`, held by members | every holder removed, 0 | +| `""` | null or unheld | `UNKNOWN_MEMBER_ID` (25) | +| unknown | null | `UNKNOWN_MEMBER_ID` (25) | + +A group that does not exist answers `UNKNOWN_MEMBER_ID` for every identity and is not created. +From v3 the response carries one entry per request identity, in request order, echoing exactly +the `member_id` and `group_instance_id` that were sent. Below v3 it carries one top-level code: +the request's own error, else the first member's. + ## Static membership is accepted, not honoured `group.instance.id` (JoinGroup v5+) is stored and echoed back to the leader, and it seeds the generated member id so a static member is recognisable in logs. Nothing else about KIP-345 is -implemented: there is no `FENCED_INSTANCE_ID`, and a returning static member is **not** matched to -its previous identity - it is a new dynamic member and its rejoin triggers a rebalance like any -other. Full static membership belongs to -[#3543](https://github.com/apache/iggy/issues/3543). +implemented, and it would need its own issue: a returning static member is **not** matched to its +previous identity. It gets a new member id and its join triggers a rebalance, while its previous +incarnation stays a member, holding its partitions, until its session expires and triggers a +second one. Kafka replaces the old identity with no rebalance at all. + +The consequences: + +- One instance id can be held by several members at once. A `LeaveGroup` by instance id removes + every holder, since leaving by instance id means that instance is gone, and still answers with + one entry. With at most one holder this is Kafka's behaviour. +- `FENCED_INSTANCE_ID` (82) is returned only by `LeaveGroup`, never by `JoinGroup`, `SyncGroup` or + `Heartbeat`. +- A static consumer does not send `LeaveGroup` on close, with either the Java client or + librdkafka, so its partitions stay assigned until its session expires, exactly as before + `LeaveGroup` existed. For prompt release use a short `session.timeout.ms`, or Java 4.1+ + `CloseOptions.LEAVE_GROUP`. ## Capacity caps @@ -113,11 +171,11 @@ survives that loop, because group state is not per-connection and heartbeats res connection - but the consumer never fetches. Fetch is a stub in any case, so nothing can be consumed until [#3535](https://github.com/apache/iggy/issues/3535)/#3542 land. -`LeaveGroup` (13) is also out of scope ([#3543](https://github.com/apache/iggy/issues/3543)). The -cost falls on the survivors, not the leaver: a consumer that shuts down gracefully stays a member -until its session expires, so the next rebalance waits up to `session.timeout.ms` (45s for a -default Java consumer) for an eviction. That is exactly what Kafka does for a *crashed* consumer, -so it is a degraded shutdown rather than a wedge. +A dynamic consumer releases its partitions on close through `LeaveGroup`. A static one does not +send it (see [Static membership](#static-membership-is-accepted-not-honoured)), so a static +consumer's shutdown still costs the survivors up to `session.timeout.ms` before they rebalance. +That is exactly what Kafka does for a *crashed* consumer, so it is a degraded shutdown rather than +a wedge. `ConsumerGroupHeartbeat` (68), the KIP-848 protocol, is not implemented and a client cannot fall back from it. A Kafka 4.0 client may need `group.protocol=classic` to reach these keys at all. diff --git a/gateways/kafka/docs/MANUAL_TESTING.md b/gateways/kafka/docs/MANUAL_TESTING.md index 544f7d604..c6d60c825 100644 --- a/gateways/kafka/docs/MANUAL_TESTING.md +++ b/gateways/kafka/docs/MANUAL_TESTING.md @@ -97,14 +97,15 @@ For each API key, test **min−1**, **min**, **max**, **max+1** using `kafka-mes | 10 | FindCoordinator | 0 | 4 | −1, 0, 4, 5 | | 11 | JoinGroup | 0 | 9 | −1, 0, 9, 10 | | 12 | Heartbeat | 0 | 4 | −1, 0, 4, 5 | +| 13 | LeaveGroup | 0 | 5 | −1, 0, 5, 6 | | 14 | SyncGroup | 0 | 5 | −1, 0, 5, 6 | | ID | Test | Expected for in-range | Expected for out-of-range | | ---- | ------ | ---------------------- | --------------------------- | -| B1 | ApiVersions negotiation | `error_code=0`; body lists 10 API keys with correct min/max | KIP-511 exception: still answers, `error_code=35` (UNSUPPORTED_VERSION), v0 response header regardless of the request's own encoding | +| B1 | ApiVersions negotiation | `error_code=0`; body lists 11 API keys with correct min/max | KIP-511 exception: still answers, `error_code=35` (UNSUPPORTED_VERSION), v0 response header regardless of the request's own encoding | | B2 | Metadata out-of-range | N/A | **Connection closes**, no response sent - Metadata has no top-level error field to carry a version-correct error in | | B3 | Produce/Fetch/ListOffsets/CreateTopics out-of-range | N/A | **Connection closes** for both above-max and below-min - `kafka_protocol`'s schema floor for each of these four messages equals `SUPPORTED_RANGES`' own min, so there is no encodable error response below min either (see `SCOPE.md`'s Governance model) | -| B4 | ApiVersions lists only scoped keys | Decode response | Contains keys 0,1,2,3,10,11,12,14,18,19 only — no OffsetCommit/OffsetFetch/LeaveGroup | +| B4 | ApiVersions lists only scoped keys | Decode response | Contains keys 0,1,2,3,10,11,12,13,14,18,19 only — no OffsetCommit/OffsetFetch | Only ApiVersions (B1) ever returns `error_code=35` on this gateway. Every other API key's out-of-range case closes the connection - see B2/B3. @@ -183,6 +184,9 @@ Requires `kcat` installed. Gateway does **not** implement SASL or full broker se | G3 | Consumer group rebalance | `kcat -b 127.0.0.1:9093 -G g1 test` in two terminals | Each prints its assigned partitions and the two sets are disjoint; then both stall, because OffsetFetch (9) is unlisted and closes the connection — the client re-runs FindCoordinator and loops. Record the exact librdkafka log lines | | G4 | Ungraceful consumer exit | `kill -9` one of G3's kcats | Within `session.timeout.ms` the survivor logs a rebalance and is assigned every partition | | G5 | Java console consumer | `kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9093 --group g2 --topic test` | Exercises JoinGroup v9, SyncGroup v5, Heartbeat v4, FindCoordinator v4. A 4.0 client may need `--consumer-property group.protocol=classic`, or it sends ConsumerGroupHeartbeat (68) and the connection closes | +| G6 | Graceful kcat exit | G3's two kcats, then Ctrl-C one (kcat closes its consumer, which sends LeaveGroup v0/v1) | The survivor rebalances within one heartbeat interval, not `session.timeout.ms`, and is assigned every partition | +| G7 | Graceful Java exit | G5 in two terminals, then Ctrl-C one | Same as G6 | +| G8 | Static member exit | G7 with `--consumer-property group.instance.id=x` on the one stopped | No LeaveGroup is sent; the survivor waits out the session timeout before rebalancing | Record kcat version and exact error strings in your test log. G1 passing is the minimum bar for client compatibility smoke. @@ -274,7 +278,7 @@ kcat version (if used): ___________ [ ] D1–D10 Flexible vs legacy encoding [ ] E1–E4 Metadata stub semantics [ ] F1–F6 TCP / connection behavior -[ ] G1–G3 kcat client (record errors for G2/G3) +[ ] G1–G8 Real clients (record errors for G2/G3) [ ] H1–H3 Adversarial input Automated regression: @@ -310,4 +314,4 @@ These are documented as TODO in [SCOPE.md](SCOPE.md) — do not fail #3421 valid - Accurate partition leadership / ISR - Transactional produce - Real offset commit semantics -- Consumer group offset commit/fetch, LeaveGroup, and admin group views (join/sync/heartbeat themselves are covered by G3-G5) +- Consumer group offset commit/fetch and admin group views (join/sync/heartbeat/leave themselves are covered by G3-G8) diff --git a/gateways/kafka/docs/SCOPE.md b/gateways/kafka/docs/SCOPE.md index b4711c3e9..cf33c2f87 100644 --- a/gateways/kafka/docs/SCOPE.md +++ b/gateways/kafka/docs/SCOPE.md @@ -52,6 +52,7 @@ it knows the server supports flexible encoding. | 10 | FindCoordinator | 0 | 4 | 0, 1, 2, 3, 4 | Answers "this gateway" for group keys; `INVALID_REQUEST` (42) for transaction/share keys; flexible encoding at v3+ | | 11 | JoinGroup | 0 | 9 | 0 … 9 | Real membership; parks on the group's join barrier; flexible encoding at v6+ | | 12 | Heartbeat | 0 | 4 | 0, 1, 2, 3, 4 | Refreshes a session; `REBALANCE_IN_PROGRESS` (27) drives a rejoin; flexible encoding at v4+ | +| 13 | LeaveGroup | 0 | 5 | 0, 1, 2, 3, 4, 5 | Removes members, per-member errors from v3; flexible encoding at v4+ | | 14 | SyncGroup | 0 | 5 | 0, 1, 2, 3, 4, 5 | Relays the leader's assignment blobs; flexible encoding at v4+ | A request is accepted when `min_version ≤ api_version ≤ max_version` for that API key. Any other version for a listed key closes the connection (ApiVersions excepted - see Governance model above). Any unlisted API key also closes the connection: no api-specific response schema exists for it, so any body this gateway could send would be misparsed by the client against the schema it expected. @@ -69,6 +70,7 @@ Use this table when configuring clients or generating wire fixtures with `kafka- | 10 | FindCoordinator | 0–4 | v3 | | 11 | JoinGroup | 0–9 | v6 | | 12 | Heartbeat | 0–4 | v4 | +| 13 | LeaveGroup | 0–5 | v4 | | 14 | SyncGroup | 0–5 | v4 | | 18 | ApiVersions | 0–3 | v3 | | 19 | CreateTopics | 2–5 | v5 | @@ -83,7 +85,6 @@ All API keys not listed above close the connection (see Governance model above) | --------- | ------ | ------- | | 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 | | 68 | ConsumerGroupHeartbeat | KIP-848 protocol; a 4.0 client may need `group.protocol=classic` | @@ -98,7 +99,7 @@ Full reference for future phases: [`kafka_api_keys_reference.md`](kafka_api_keys | Layer | #3421 | Description | | ------- | ------- | ------------- | | **1 — Wire framing** | In scope | `server.rs` — custom, zero-copy frame I/O; `header.rs` delegates version selection to `kafka_protocol::messages::ApiKey` | -| **2 — Request/response codecs** | Partial | Decode/encode via the `kafka_protocol` crate (broker feature only) for 10 keys; `bounds_guard.rs` pre-validates against unbounded allocation before handing a frame to the crate; stub responses only | +| **2 — Request/response codecs** | Partial | Decode/encode via the `kafka_protocol` crate (broker feature only) for 11 keys; `bounds_guard.rs` pre-validates against unbounded allocation before handing a frame to the crate; stub responses only | | **3 — Iggy bridge** | Landed, not wired in | `bridge/` module (connection, topic mapping, provisioning, high watermark) landed; Produce/Fetch handler wiring itself is a follow-on ([#3535](https://github.com/apache/iggy/issues/3535)/[#3536](https://github.com/apache/iggy/issues/3536)) | --- @@ -147,9 +148,9 @@ Offset persistence design ([#3540](https://github.com/apache/iggy/issues/3540)): [`OFFSET_STORAGE.md`](OFFSET_STORAGE.md). - [x] FindCoordinator (10), JoinGroup (11), Heartbeat (12), SyncGroup (14) - - [#3541](https://github.com/apache/iggy/issues/3541), see [`CONSUMER_GROUPS.md`](CONSUMER_GROUPS.md) + [#3541](https://github.com/apache/iggy/issues/3541); LeaveGroup (13) - + [#3543](https://github.com/apache/iggy/issues/3543); see [`CONSUMER_GROUPS.md`](CONSUMER_GROUPS.md) - [ ] OffsetCommit (8), OffsetFetch (9) -- [ ] LeaveGroup (13) - [ ] DescribeGroups (15), ListGroups (16) as needed by target clients ### Phase 3+ — Auth, admin, tuning diff --git a/gateways/kafka/docs/TEST_SUITE.md b/gateways/kafka/docs/TEST_SUITE.md index d19ca3db6..7a5c31948 100644 --- a/gateways/kafka/docs/TEST_SUITE.md +++ b/gateways/kafka/docs/TEST_SUITE.md @@ -61,7 +61,7 @@ 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 | +| [`consumer_group_tests.rs`](../tests/consumer_group_tests.rs) | `FindCoordinator`/`JoinGroup`/`Heartbeat`/`LeaveGroup`/`SyncGroup` against one shared `GatewayState` with paused time, plus two-socket TCP rebalance and leave scenarios | 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 | | [`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) | diff --git a/gateways/kafka/docs/kafka_api_keys_reference.md b/gateways/kafka/docs/kafka_api_keys_reference.md index 7fb6d52c1..43ace2eb6 100644 --- a/gateways/kafka/docs/kafka_api_keys_reference.md +++ b/gateways/kafka/docs/kafka_api_keys_reference.md @@ -279,9 +279,10 @@ Key new minimums: | FindCoordinator | v0-v4 | v6 | 2 versions behind | | JoinGroup | v0-v9 | v9 | current | | Heartbeat | v0-v4 | v4 | current | +| LeaveGroup | v0-v5 | v5 | current | | SyncGroup | v0-v5 | v5 | current | -### Missing from `SUPPORTED_RANGES` (78 of the 88 API keys in this document) +### Missing from `SUPPORTED_RANGES` (77 of the 88 API keys in this document) Every key not in `SUPPORTED_RANGES` closes the connection - the same policy applied to every other unlisted key, not a special case for these. No api-specific response schema exists for an @@ -289,12 +290,14 @@ unlisted key, so any body the gateway could send would be misparsed by the clien schema it expected. This includes: - **Client bootstrap blockers**: OffsetCommit (8), OffsetFetch (9) -- **Classic consumer group protocol**: LeaveGroup (13). FindCoordinator (10), JoinGroup (11), Heartbeat (12) and SyncGroup (14) are supported - see [`CONSUMER_GROUPS.md`](CONSUMER_GROUPS.md) - **New consumer group protocol**: ConsumerGroupHeartbeat (68) — default in Kafka 4.0 - **Share groups (KIP-932)**: ShareFetch, ShareGroupHeartbeat, ShareAcknowledge - **Auth flow**: SaslHandshake (17), SaslAuthenticate (36) - **All broker/KRaft-internal keys** (Group 14) +The classic consumer group keys FindCoordinator (10), JoinGroup (11), Heartbeat (12), +LeaveGroup (13) and SyncGroup (14) are supported - see [`CONSUMER_GROUPS.md`](CONSUMER_GROUPS.md). + Remaining scope (consumer groups, auth, admin/tuning) is tracked in `SCOPE.md`'s TODO section, not duplicated here. diff --git a/gateways/kafka/scripts/ci-wire-fixtures.sh b/gateways/kafka/scripts/ci-wire-fixtures.sh index df8c99f6b..de6f3c1bb 100755 --- a/gateways/kafka/scripts/ci-wire-fixtures.sh +++ b/gateways/kafka/scripts/ci-wire-fixtures.sh @@ -24,7 +24,7 @@ set -euo pipefail FIXTURES_DIR="gateways/kafka/tools/kafka-tool/kafka_messages" # API keys requested by api_handler_tests, version_firewall_tests, and server_e2e_tests. -FIXTURE_API_KEYS=(0 1 2 19) +FIXTURE_API_KEYS=(0 1 2 10 11 12 13 14 19) usage() { echo "Usage: $0 {generate|cleanup}" >&2 diff --git a/gateways/kafka/src/group/mod.rs b/gateways/kafka/src/group/mod.rs index 996cc786d..3a734a37a 100644 --- a/gateways/kafka/src/group/mod.rs +++ b/gateways/kafka/src/group/mod.rs @@ -18,9 +18,9 @@ //! In-memory coordinator for Kafka's classic consumer group protocol. //! //! [`GroupCoordinator`] owns every group this gateway instance coordinates and is the only -//! module that awaits: `FindCoordinator`/`JoinGroup`/`Heartbeat`/`SyncGroup` handlers translate -//! wire messages into the request types here, and `state` holds the synchronous state machine -//! those requests drive. +//! module that awaits: `FindCoordinator`/`JoinGroup`/`Heartbeat`/`LeaveGroup`/`SyncGroup` handlers +//! translate wire messages into the request types here, and `state` holds the synchronous state +//! machine those requests drive. //! //! Membership is process memory, not Iggy state. Two gateway instances fronting one Iggy cluster //! therefore coordinate two independent groups under one name; see `docs/CONSUMER_GROUPS.md`. @@ -31,14 +31,14 @@ use std::collections::HashMap; use std::time::Duration; use bytes::Bytes; -use kafka_protocol::messages::{JoinGroupRequest, SyncGroupRequest}; +use kafka_protocol::messages::{JoinGroupRequest, LeaveGroupRequest, SyncGroupRequest}; use kafka_protocol::protocol::StrBytes; use tokio::sync::{Mutex, watch}; use tokio::time::Instant; use tokio_util::sync::CancellationToken; use crate::group::state::{GroupState, Step}; -use crate::protocol::api::{ERROR_NOT_COORDINATOR, ERROR_UNKNOWN_MEMBER_ID}; +use crate::protocol::api::{ERROR_NONE, ERROR_NOT_COORDINATOR, ERROR_UNKNOWN_MEMBER_ID}; /// Kafka's own `group.min.session.timeout.ms` default. const DEFAULT_MIN_SESSION_TIMEOUT: Duration = Duration::from_secs(6); @@ -228,6 +228,97 @@ impl SyncResult { } } +/// One `LeaveGroup` request, normalized across wire versions. +#[derive(Debug, Clone)] +pub struct LeaveRequest { + pub group_id: StrBytes, + /// Request order, which is also the order the v3+ response answers in. + pub members: Vec<LeavingMember>, +} + +/// One identity a `LeaveGroup` asks to remove. +#[derive(Debug, Clone)] +pub struct LeavingMember { + /// Empty when a v3+ request leaves by `group_instance_id` alone. + pub member_id: StrBytes, + pub group_instance_id: Option<StrBytes>, +} + +impl From<(i16, &LeaveGroupRequest)> for LeaveRequest { + fn from((api_version, request): (i16, &LeaveGroupRequest)) -> Self { + let members = if api_version >= 3 { + request + .members + .iter() + .map(|member| LeavingMember { + member_id: member.member_id.clone(), + group_instance_id: member.group_instance_id.clone(), + }) + .collect() + } else { + vec![LeavingMember { + member_id: request.member_id.clone(), + group_instance_id: None, + }] + }; + Self { + group_id: request.group_id.0.clone(), + members, + } + } +} + +/// Everything a `LeaveGroup` response carries, before the handler shapes it for a wire version. +#[derive(Debug, Clone)] +pub struct LeaveResult { + /// A whole-request failure. When set, `members` is empty. + pub error: i16, + /// One entry per request identity, in request order, echoing exactly what was sent. + pub members: Vec<LeftMember>, +} + +impl LeaveResult { + #[must_use] + pub const fn error(error: i16) -> Self { + Self { + error, + members: Vec::new(), + } + } + + /// The single code a v0-v2 response carries: the request's own error, else the first + /// member's, as Kafka's `LeaveGroupResponse.getError` folds them. + #[must_use] + pub fn top_level_error(&self) -> i16 { + if self.error != ERROR_NONE { + return self.error; + } + self.members + .iter() + .map(|member| member.error) + .find(|error| *error != ERROR_NONE) + .unwrap_or(ERROR_NONE) + } +} + +/// The answer for one `LeavingMember`. +#[derive(Debug, Clone)] +pub struct LeftMember { + pub member_id: StrBytes, + pub group_instance_id: Option<StrBytes>, + pub error: i16, +} + +impl From<(&LeavingMember, i16)> for LeftMember { + fn from((identity, error): (&LeavingMember, i16)) -> Self { + Self { + member_id: identity.member_id.clone(), + group_instance_id: identity.group_instance_id.clone(), + error, + } + } +} + /// Every consumer group this gateway instance coordinates. /// /// There is no timer task. A request that touches a group first expires whatever is overdue in @@ -333,6 +424,12 @@ impl GroupCoordinator { ) } + /// Removes `request`'s members and releases whoever is parked on the group. Never parks. + pub async fn leave(&self, request: &LeaveRequest) -> LeaveResult { + let mut groups = self.groups.lock().await; + state::leave_step(&mut groups, request, Instant::now()) + } + /// Sleeps until the group changes or `wake_at` passes. `false` means the gateway is draining. async fn wait_until(&self, mut receiver: watch::Receiver<u64>, wake_at: Instant) -> bool { tokio::select! { @@ -370,6 +467,8 @@ fn millis_to_duration(millis: i32) -> Duration { #[cfg(test)] mod tests { + use kafka_protocol::messages::leave_group_request::MemberIdentity; + use super::*; #[test] @@ -396,4 +495,45 @@ mod tests { Duration::from_secs(30) ); } + + #[test] + fn given_a_v2_leave_when_normalizing_should_use_the_top_level_member_id() { + let request = LeaveGroupRequest::default() + .with_group_id(StrBytes::from_static_str("g").into()) + .with_member_id(StrBytes::from_static_str("m-1")); + + let normalized = LeaveRequest::from((2, &request)); + + assert_eq!(normalized.group_id.as_str(), "g"); + assert_eq!(normalized.members.len(), 1); + assert_eq!(normalized.members[0].member_id.as_str(), "m-1"); + assert_eq!(normalized.members[0].group_instance_id, None); + } + + #[test] + fn given_a_v3_leave_when_normalizing_should_use_the_members_array() { + let request = LeaveGroupRequest::default() + .with_group_id(StrBytes::from_static_str("g").into()) + .with_member_id(StrBytes::from_static_str("ignored")) + .with_members(vec![ + MemberIdentity::default().with_member_id(StrBytes::from_static_str("m-1")), + MemberIdentity::default() + .with_member_id(StrBytes::new()) + .with_group_instance_id(Some(StrBytes::from_static_str("i-2"))), + ]); + + let normalized = LeaveRequest::from((3, &request)); + + let identities: Vec<(&str, Option<&str>)> = normalized + .members + .iter() + .map(|member| { + ( + member.member_id.as_str(), + member.group_instance_id.as_ref().map(StrBytes::as_str), + ) + }) + .collect(); + assert_eq!(identities, vec![("m-1", None), ("", Some("i-2"))]); + } } diff --git a/gateways/kafka/src/group/state.rs b/gateways/kafka/src/group/state.rs index 6ca89f68d..5967f6849 100644 --- a/gateways/kafka/src/group/state.rs +++ b/gateways/kafka/src/group/state.rs @@ -25,7 +25,7 @@ //! Kafka's `Empty` group state is "absent from the map": offsets live in Iggy, so an empty group //! holds nothing worth keeping and retaining it would be an unbounded-memory vector. -use std::collections::{BTreeMap, HashMap}; +use std::collections::{BTreeMap, HashMap, HashSet}; use std::time::Duration; use bytes::Bytes; @@ -35,12 +35,13 @@ use tokio::time::Instant; use uuid::Uuid; use crate::group::{ - GroupCoordinatorConfig, JoinRequest, JoinResult, JoinedMember, SyncRequest, SyncResult, + GroupCoordinatorConfig, JoinRequest, JoinResult, JoinedMember, LeaveRequest, LeaveResult, + LeavingMember, LeftMember, SyncRequest, SyncResult, }; use crate::protocol::api::{ - ERROR_COORDINATOR_NOT_AVAILABLE, 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_COORDINATOR_NOT_AVAILABLE, ERROR_FENCED_INSTANCE_ID, 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_UNKNOWN_MEMBER_ID, }; @@ -75,6 +76,17 @@ pub enum Step<T> { }, } +/// What one `LeaveGroup` removed, which decides how the group reacts afterwards. +#[derive(Default)] +struct Departures { + members: bool, + pending: bool, +} + +/// Member ids per `group_instance_id`. Without KIP-345 replacement one instance id can be held by +/// several members, so each maps to a set. +type InstanceHolders = HashMap<StrBytes, HashSet<StrBytes>>; + pub struct Member { group_instance_id: Option<StrBytes>, session_timeout: Duration, @@ -301,6 +313,88 @@ impl GroupState { } } + fn instance_holders(&self) -> InstanceHolders { + let mut holders = InstanceHolders::new(); + for (member_id, member) in &self.members { + if let Some(instance_id) = &member.group_instance_id { + holders + .entry(instance_id.clone()) + .or_default() + .insert(member_id.clone()); + } + } + holders + } + + /// Removes one `LeaveGroup` identity, following Kafka's `handleLeaveGroup` order. + /// + /// `holders` must stay exact across calls: a later identity in the same request is resolved + /// against it. + fn leave_one( + &mut self, + identity: &LeavingMember, + holders: &mut InstanceHolders, + departed: &mut Departures, + ) -> i16 { + let holders_of = identity + .group_instance_id + .as_ref() + .and_then(|instance_id| holders.get_mut(instance_id)) + .filter(|ids| !ids.is_empty()); + + if identity.member_id.is_empty() { + let Some(ids) = holders_of else { + return ERROR_UNKNOWN_MEMBER_ID; + }; + for id in std::mem::take(ids) { + self.remove_member(&id); + } + departed.members = true; + return ERROR_NONE; + } + if self.pending.remove(&identity.member_id).is_some() { + departed.pending = true; + return ERROR_NONE; + } + match (identity.group_instance_id.is_some(), holders_of) { + (_, Some(ids)) if ids.contains(&identity.member_id) => {} + (_, Some(_)) => return ERROR_FENCED_INSTANCE_ID, + (true, None) => return ERROR_UNKNOWN_MEMBER_ID, + (false, None) if !self.members.contains_key(&identity.member_id) => { + return ERROR_UNKNOWN_MEMBER_ID; + } + (false, None) => {} + } + if let Some(ids) = self + .members + .get(&identity.member_id) + .and_then(|member| member.group_instance_id.as_ref()) + .and_then(|instance_id| holders.get_mut(instance_id)) + { + ids.remove(&identity.member_id); + } + self.remove_member(&identity.member_id); + departed.members = true; + ERROR_NONE + } + + /// React to a `LeaveGroup` the way the session sweep reacts to an expiry. + /// + /// The rebalance opens only after the removals, so its join window is sized from the + /// members that remain. The final `bump` is unconditional on any removal: it is what answers + /// a waiter the departed member still has parked on another connection. + fn after_departure(&mut self, departed: &Departures, now: Instant) { + if departed.members && !self.members.is_empty() && self.phase != Phase::PreparingRebalance { + self.prepare_rebalance(now, None); + } + if self.phase == Phase::PreparingRebalance { + self.maybe_complete_join(now); + } + if departed.members || departed.pending { + self.bump(); + } + } + /// The earliest moment any rule in this group could fire. fn next_deadline(&self) -> Option<Instant> { let phase_deadline = match self.phase { @@ -909,6 +1003,53 @@ pub fn sync_resume_step( } } +/// Never creates a group, and removes one the leave empties rather than leaving it for +/// `reclaim_expired`. +pub fn leave_step(groups: &mut Groups, request: &LeaveRequest, now: Instant) -> LeaveResult { + if request.group_id.is_empty() { + return LeaveResult::error(ERROR_INVALID_GROUP_ID); + } + let unknown = || LeaveResult { + error: ERROR_NONE, + members: request + .members + .iter() + .map(|identity| LeftMember::from((identity, ERROR_UNKNOWN_MEMBER_ID))) + .collect(), + }; + if !tick_group(groups, &request.group_id, now) { + return unknown(); + } + let Some(group) = groups.get_mut(&request.group_id) else { + return unknown(); + }; + + let mut holders = if request + .members + .iter() + .any(|identity| identity.group_instance_id.is_some()) + { + group.instance_holders() + } else { + InstanceHolders::new() + }; + let mut departed = Departures::default(); + let mut members = Vec::with_capacity(request.members.len()); + for identity in &request.members { + let error = group.leave_one(identity, &mut holders, &mut departed); + members.push(LeftMember::from((identity, error))); + } + group.after_departure(&departed, now); + + if group.is_empty() { + groups.remove(&request.group_id); + } + LeaveResult { + error: ERROR_NONE, + members, + } +} + fn admit( group: &mut GroupState, config: &GroupCoordinatorConfig, @@ -1920,4 +2061,638 @@ mod tests { assert_eq!(error_of(&step), ERROR_UNKNOWN_MEMBER_ID); assert!(groups.is_empty(), "a rejected join must not create a group"); } + + fn leave( + groups: &mut Groups, + identities: &[(&str, Option<&str>)], + now: Instant, + ) -> LeaveResult { + let request = LeaveRequest { + group_id: group_id(), + members: identities + .iter() + .map(|(member_id, instance_id)| LeavingMember { + member_id: StrBytes::from_string((*member_id).to_owned()), + group_instance_id: instance_id.map(|id| StrBytes::from_string(id.to_owned())), + }) + .collect(), + }; + leave_step(groups, &request, now) + } + + fn codes(result: &LeaveResult) -> Vec<i16> { + result.members.iter().map(|member| member.error).collect() + } + + fn static_request(instance_id: &str) -> JoinRequest { + JoinRequest { + group_instance_id: Some(StrBytes::from_string(instance_id.to_owned())), + ..request("", &["x"]) + } + } + + fn pending_request() -> JoinRequest { + JoinRequest { + require_known_member_id: true, + ..request("", &["x"]) + } + } + + /// `two_members` carried through the leader's `SyncGroup`, so the group is Stable. + fn stable_two_members( + groups: &mut Groups, + config: &GroupCoordinatorConfig, + now: Instant, + ) -> (StrBytes, StrBytes) { + let (leader, follower) = two_members(groups, config, &["x"], &["x"], now); + let generation = groups[&group_id()].generation_id; + let _ = sync_step(groups, config, &sync_request(&leader, generation), now); + assert_eq!(groups[&group_id()].phase, Phase::Stable); + (leader, follower) + } + + /// The generation must hold: a bumped one turns the survivor's heartbeat into + /// `ILLEGAL_GENERATION`, which a Java client handles as lost partitions instead of a rejoin. + #[test] + fn given_a_member_leaving_a_stable_group_should_open_a_rebalance_without_bumping_the_generation() + { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let (leader, follower) = stable_two_members(&mut groups, &config, now); + let generation = groups[&group_id()].generation_id; + + let result = leave(&mut groups, &[(follower.as_str(), None)], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + let group = &groups[&group_id()]; + assert_eq!(group.phase, Phase::PreparingRebalance); + assert_eq!(group.generation_id, generation); + assert!(!group.members.contains_key(&follower)); + assert_eq!( + heartbeat_step(&mut groups, &group_id(), generation, &leader, now), + ERROR_REBALANCE_IN_PROGRESS + ); + } + + /// Asserted with no time advance: `reclaim_expired` would free the slot later anyway, so + /// only an immediate check proves the leave removed the group itself. + #[test] + fn given_the_last_member_leaving_should_remove_the_group_immediately() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let member = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + + let result = leave(&mut groups, &[(member.as_str(), None)], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + assert!(!groups.contains_key(&group_id())); + } + + #[test] + fn given_a_leave_for_an_unknown_group_should_not_create_it() { + let mut groups = Groups::new(); + + let result = leave( + &mut groups, + &[("ghost", None), ("", Some("instance"))], + Instant::now(), + ); + + assert_eq!(result.error, ERROR_NONE); + assert_eq!( + codes(&result), + vec![ERROR_UNKNOWN_MEMBER_ID, ERROR_UNKNOWN_MEMBER_ID] + ); + assert!(groups.is_empty()); + } + + /// Checked straight after `leave_step`: a woken waiter's own tick would also complete the + /// join, so only this proves the leave did. + #[test] + fn given_the_last_straggler_leaving_a_preparing_group_should_complete_the_join() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let (leader, follower) = stable_two_members(&mut groups, &config, now); + let generation = groups[&group_id()].generation_id; + let rejoin = join_step(&mut groups, &config, &request(leader.as_str(), &["x"]), now); + assert!(matches!(rejoin, Step::Wait { .. })); + + let _ = leave(&mut groups, &[(follower.as_str(), None)], now); + + let group = &groups[&group_id()]; + assert_eq!(group.phase, Phase::CompletingRebalance); + assert_eq!(group.generation_id, generation + 1); + let answer = group.members[&leader] + .join_response + .as_ref() + .expect("the completed join must leave the leader its answer"); + let roster: Vec<&StrBytes> = answer.members.iter().map(|m| &m.member_id).collect(); + assert_eq!(roster, vec![&leader]); + } + + /// Without the rebalance the survivor becomes leader of a generation whose roster it never + /// received, and its `SyncGroup` would fan out an empty assignment. + #[test] + fn given_the_leader_leaving_while_completing_should_open_a_rebalance() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let (leader, follower) = two_members(&mut groups, &config, &["x"], &["x"], now); + let generation = groups[&group_id()].generation_id; + assert_eq!(groups[&group_id()].phase, Phase::CompletingRebalance); + + let _ = leave(&mut groups, &[(leader.as_str(), None)], now); + + let group = &groups[&group_id()]; + assert_eq!(group.phase, Phase::PreparingRebalance); + assert_eq!(group.leader.as_ref(), Some(&follower)); + let Step::Respond(sync) = sync_step( + &mut groups, + &config, + &sync_request(&follower, generation), + now, + ) else { + panic!("a sync during a rebalance is answered, not parked"); + }; + assert_eq!(sync.error, ERROR_REBALANCE_IN_PROGRESS); + } + + #[test] + fn given_a_leave_while_stable_should_size_the_join_window_from_the_remaining_members() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let survivor = member_id_of(&join_step( + &mut groups, + &config, + &request_with("", &["x"], 600, 5), + now, + )); + let leaver = member_id_of(&join_step( + &mut groups, + &config, + &request_with("", &["x"], 600, 300), + now, + )); + let _ = join_step( + &mut groups, + &config, + &request_with(survivor.as_str(), &["x"], 600, 5), + now, + ); + let generation = groups[&group_id()].generation_id; + let _ = sync_step( + &mut groups, + &config, + &sync_request(&survivor, generation), + now, + ); + assert_eq!(groups[&group_id()].phase, Phase::Stable); + + let _ = leave(&mut groups, &[(leaver.as_str(), None)], now); + + assert_eq!( + groups[&group_id()].join_deadline, + Some(now + Duration::from_secs(5)), + "the leaver's own rebalance timeout must not stretch the survivors' join window" + ); + } + + /// A pending id never joined, so dropping it must not disturb a Stable group. + #[test] + fn given_a_pending_member_leaving_a_stable_group_should_not_open_a_rebalance() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let _ = stable_two_members(&mut groups, &config, now); + let generation = groups[&group_id()].generation_id; + let pending = member_id_of(&join_step(&mut groups, &config, &pending_request(), now)); + assert!(groups[&group_id()].pending.contains_key(&pending)); + + let result = leave(&mut groups, &[(pending.as_str(), None)], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + let group = &groups[&group_id()]; + assert_eq!(group.phase, Phase::Stable); + assert_eq!(group.generation_id, generation); + assert!(group.pending.is_empty()); + } + + #[test] + fn given_a_pending_member_leaving_a_preparing_group_should_unblock_the_join() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let (leader, follower) = stable_two_members(&mut groups, &config, now); + let generation = groups[&group_id()].generation_id; + let pending = member_id_of(&join_step(&mut groups, &config, &pending_request(), now)); + let _ = join_step(&mut groups, &config, &request(leader.as_str(), &["x"]), now); + let _ = join_step( + &mut groups, + &config, + &request(follower.as_str(), &["x"]), + now, + ); + assert_eq!( + groups[&group_id()].phase, + Phase::PreparingRebalance, + "the outstanding pending id must be what holds the barrier open" + ); + + let result = leave(&mut groups, &[(pending.as_str(), None)], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + let group = &groups[&group_id()]; + assert_eq!(group.phase, Phase::CompletingRebalance); + assert_eq!(group.generation_id, generation + 1); + } + + #[test] + fn given_a_leave_by_instance_id_without_a_member_id_should_remove_that_static_member() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let static_member = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-a"), + now, + )); + let dynamic = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + + let result = leave(&mut groups, &[("", Some("i-a"))], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + let group = &groups[&group_id()]; + assert!(!group.members.contains_key(&static_member)); + assert!(group.members.contains_key(&dynamic)); + } + + /// The admin `removeMembersFromConsumerGroup` path: removing a holder must open a rebalance + /// and wake the survivors, or the group stays Stable with the removed member's partitions + /// assigned to nobody. + #[test] + fn given_a_stable_group_when_a_holder_leaves_by_instance_id_alone_should_open_a_rebalance_and_wake_waiters() + { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let (leader, follower) = stable_two_members(&mut groups, &config, now); + let generation = groups[&group_id()].generation_id; + let group = groups.get_mut(&group_id()).expect("the stable group"); + group + .members + .get_mut(&follower) + .expect("the follower") + .group_instance_id = Some(StrBytes::from_static_str("i-f")); + let mut changed = group.subscribe(); + changed.mark_unchanged(); + + let result = leave(&mut groups, &[("", Some("i-f"))], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + assert!( + changed.has_changed().expect("the group is still live"), + "survivor waiters must be woken" + ); + let group = &groups[&group_id()]; + assert!(!group.members.contains_key(&follower)); + assert_eq!(group.phase, Phase::PreparingRebalance); + assert_eq!( + heartbeat_step(&mut groups, &group_id(), generation, &leader, now), + ERROR_REBALANCE_IN_PROGRESS + ); + } + + /// What Java `CloseOptions.LEAVE_GROUP` and Kafka Streams send: the member's own pair. It is + /// the holder, so it must leave cleanly rather than be fenced. + #[test] + fn given_a_static_member_when_it_leaves_with_its_own_member_and_instance_id_should_leave() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let static_member = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-a"), + now, + )); + let dynamic = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + + let result = leave(&mut groups, &[(static_member.as_str(), Some("i-a"))], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + let group = &groups[&group_id()]; + assert!(!group.members.contains_key(&static_member)); + assert!(group.members.contains_key(&dynamic)); + } + + /// The first instance-only leave takes every holder, so a repeat in the same request finds + /// nobody left holding the id. + #[test] + fn given_two_holders_when_one_request_leaves_their_instance_id_twice_should_answer_the_repeat_unknown_member_id() + { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let first = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-a"), + now, + )); + let second = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-a"), + now, + )); + assert_ne!(first, second); + + let result = leave(&mut groups, &[("", Some("i-a")), ("", Some("i-a"))], now); + + assert_eq!(codes(&result), vec![ERROR_NONE, ERROR_UNKNOWN_MEMBER_ID]); + assert!(!groups.contains_key(&group_id())); + } + + /// A rebalance with no members would complete at once, clear `pending` and drop the group, + /// so the pending client's rejoin would be refused. Session expiry in `tick` holds the same + /// guard. + #[test] + fn given_a_pending_member_when_the_last_real_member_leaves_should_still_admit_the_pending_rejoin() + { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let member = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + let generation = groups[&group_id()].generation_id; + let _ = sync_step( + &mut groups, + &config, + &sync_request(&member, generation), + now, + ); + let pending = member_id_of(&join_step(&mut groups, &config, &pending_request(), now)); + assert!(groups[&group_id()].pending.contains_key(&pending)); + + let result = leave(&mut groups, &[(member.as_str(), None)], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + assert!(groups[&group_id()].pending.contains_key(&pending)); + let Step::Respond(rejoined) = join_step( + &mut groups, + &config, + &request(pending.as_str(), &["x"]), + now, + ) else { + panic!("the pending rejoin must be answered, not parked"); + }; + assert_eq!(rejoined.error, ERROR_NONE); + assert!(groups[&group_id()].members.contains_key(&pending)); + } + + #[test] + fn given_a_stable_group_when_only_a_pending_member_leaves_should_wake_waiters() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let _ = stable_two_members(&mut groups, &config, now); + let pending = member_id_of(&join_step(&mut groups, &config, &pending_request(), now)); + let mut changed = groups[&group_id()].subscribe(); + changed.mark_unchanged(); + + let result = leave(&mut groups, &[(pending.as_str(), None)], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + assert!(changed.has_changed().expect("the group is still live")); + assert_eq!(groups[&group_id()].phase, Phase::Stable); + } + + #[test] + fn given_a_leave_whose_member_id_does_not_hold_the_instance_id_should_be_fenced() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let holder = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-a"), + now, + )); + let other = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-b"), + now, + )); + + let result = leave(&mut groups, &[(other.as_str(), Some("i-a"))], now); + + assert_eq!(codes(&result), vec![ERROR_FENCED_INSTANCE_ID]); + let group = &groups[&group_id()]; + assert!(group.members.contains_key(&holder)); + assert!(group.members.contains_key(&other)); + } + + #[test] + fn given_a_leave_by_an_instance_id_nobody_holds_should_return_unknown_member_id() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let dynamic = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + + let result = leave(&mut groups, &[(dynamic.as_str(), Some("i-z"))], now); + + assert_eq!(codes(&result), vec![ERROR_UNKNOWN_MEMBER_ID]); + assert!(groups[&group_id()].members.contains_key(&dynamic)); + } + + /// The Java client's `CloseOptions.LEAVE_GROUP` sends a static member's id with no instance id. + #[test] + fn given_a_static_member_leaving_by_member_id_alone_should_be_removed() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let static_member = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-a"), + now, + )); + let dynamic = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + + let result = leave(&mut groups, &[(static_member.as_str(), None)], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + let group = &groups[&group_id()]; + assert!(!group.members.contains_key(&static_member)); + assert!(group.members.contains_key(&dynamic)); + } + + /// Without identity replacement a restarted static consumer leaves its old incarnation + /// behind under the same instance id. Both go, and the answer is still one entry. + #[test] + fn given_two_holders_of_one_instance_id_when_left_by_instance_should_remove_both() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let first = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-a"), + now, + )); + let second = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-a"), + now, + )); + let bystander = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + assert_ne!(first, second); + + let result = leave(&mut groups, &[("", Some("i-a"))], now); + + assert_eq!(codes(&result), vec![ERROR_NONE]); + assert_eq!(result.members[0].member_id.as_str(), ""); + assert_eq!( + result.members[0].group_instance_id.as_deref(), + Some("i-a"), + "the answer echoes the request identity, not either holder it resolved to" + ); + let remaining: Vec<&StrBytes> = groups[&group_id()].members.keys().collect(); + assert_eq!(remaining, vec![&bystander]); + } + + #[test] + fn given_a_static_member_left_by_member_id_when_the_same_request_leaves_its_instance_should_answer_unknown_member_id() + { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let static_member = member_id_of(&join_step( + &mut groups, + &config, + &static_request("i-a"), + now, + )); + let _dynamic = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + + let result = leave( + &mut groups, + &[(static_member.as_str(), None), ("", Some("i-a"))], + now, + ); + + assert_eq!(codes(&result), vec![ERROR_NONE, ERROR_UNKNOWN_MEMBER_ID]); + } + + #[test] + fn given_a_member_left_when_its_old_join_waiter_resumes_should_answer_unknown_member_id() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let (_leader, follower) = stable_two_members(&mut groups, &config, now); + let parked = join_step( + &mut groups, + &config, + &request(follower.as_str(), &["x", "y"]), + now, + ); + assert!(matches!(parked, Step::Wait { .. })); + + let _ = leave(&mut groups, &[(follower.as_str(), None)], now); + let resumed = join_resume_step(&mut groups, &group_id(), &follower, now); + + assert_eq!(error_of(&resumed), ERROR_UNKNOWN_MEMBER_ID); + assert!(!groups[&group_id()].members.contains_key(&follower)); + } + + #[test] + fn given_a_member_that_left_when_a_new_member_joins_a_full_group_should_admit_it() { + let config = GroupCoordinatorConfig { + max_members_per_group: 1, + ..config() + }; + let mut groups = Groups::new(); + let now = Instant::now(); + let member = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + let _ = leave(&mut groups, &[(member.as_str(), None)], now); + assert_eq!( + error_of(&join_step(&mut groups, &config, &request("", &["x"]), now)), + ERROR_NONE, + "a member that left must give its slot back" + ); + + let mut groups = Groups::new(); + let pending = member_id_of(&join_step(&mut groups, &config, &pending_request(), now)); + assert_eq!( + error_of(&join_step(&mut groups, &config, &request("", &["x"]), now)), + ERROR_GROUP_MAX_SIZE_REACHED, + "the pending id must be what fills the group" + ); + let _ = leave(&mut groups, &[(pending.as_str(), None)], now); + assert_eq!( + error_of(&join_step(&mut groups, &config, &request("", &["x"]), now)), + ERROR_NONE, + "a pending id that left must give its slot back" + ); + } + + #[test] + fn given_a_leave_while_completing_should_keep_uncollected_answers_of_the_rest() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let leader = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + let second = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + let third = member_id_of(&join_step(&mut groups, &config, &request("", &["x"]), now)); + let _ = join_step(&mut groups, &config, &request(leader.as_str(), &["x"]), now); + assert_eq!(groups[&group_id()].phase, Phase::CompletingRebalance); + assert!(groups[&group_id()].members[&second].join_response.is_some()); + + let _ = leave(&mut groups, &[(third.as_str(), None)], now); + + assert_eq!(groups[&group_id()].phase, Phase::PreparingRebalance); + assert!( + groups[&group_id()].members[&second].join_response.is_some(), + "a leave must not revoke an answer another member's waiter has not collected" + ); + } + + #[test] + fn given_a_duplicate_identity_in_one_leave_should_answer_the_second_unknown_member_id() { + let config = config(); + let mut groups = Groups::new(); + let now = Instant::now(); + let (_leader, follower) = stable_two_members(&mut groups, &config, now); + + let result = leave( + &mut groups, + &[(follower.as_str(), None), (follower.as_str(), None)], + now, + ); + + assert_eq!(codes(&result), vec![ERROR_NONE, ERROR_UNKNOWN_MEMBER_ID]); + } + + #[test] + fn given_an_empty_group_id_when_leaving_should_return_invalid_group_id() { + let mut groups = Groups::new(); + let request = LeaveRequest { + group_id: StrBytes::new(), + members: vec![LeavingMember { + member_id: StrBytes::from_static_str("m"), + group_instance_id: None, + }], + }; + + let result = leave_step(&mut groups, &request, Instant::now()); + + assert_eq!(result.error, ERROR_INVALID_GROUP_ID); + assert!(result.members.is_empty()); + } } diff --git a/gateways/kafka/src/protocol/api.rs b/gateways/kafka/src/protocol/api.rs index aad116f73..82359d3c1 100644 --- a/gateways/kafka/src/protocol/api.rs +++ b/gateways/kafka/src/protocol/api.rs @@ -25,7 +25,7 @@ use crate::bridge::IggyBridge; use crate::group::{GroupCoordinator, GroupCoordinatorConfig}; use crate::protocol::handlers::{ api_versions, create_topics, dispatch, fetch, find_coordinator, heartbeat, join_group, - list_offsets, metadata, produce, sync_group, + leave_group, list_offsets, metadata, produce, sync_group, }; pub const API_KEY_PRODUCE: i16 = 0; @@ -35,6 +35,7 @@ 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_LEAVE_GROUP: i16 = 13; pub const API_KEY_SYNC_GROUP: i16 = 14; pub const API_KEY_API_VERSIONS: i16 = 18; pub const API_KEY_CREATE_TOPICS: i16 = 19; @@ -103,6 +104,8 @@ pub const ERROR_INVALID_REQUEST: i16 = 42; /// 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; +/// Sent only by `LeaveGroup`: the member id does not hold the `group_instance_id` it named. +pub const ERROR_FENCED_INSTANCE_ID: i16 = 82; /// Result of handling one Kafka request body. #[derive(Debug)] @@ -173,6 +176,7 @@ static SUPPORTED_RANGES: &[ApiVersionRange] = &[ find_coordinator::RANGE, join_group::RANGE, heartbeat::RANGE, + leave_group::RANGE, sync_group::RANGE, ]; diff --git a/gateways/kafka/src/protocol/bounds_guard.rs b/gateways/kafka/src/protocol/bounds_guard.rs index ebb47a144..513307b66 100644 --- a/gateways/kafka/src/protocol/bounds_guard.rs +++ b/gateways/kafka/src/protocol/bounds_guard.rs @@ -27,7 +27,7 @@ //! any of `SUPPORTED_RANGES`, so this is reachable by anyone who can `connect()`. //! //! This module walks the same field shape `kafka_protocol`'s real decode walks for each of the -//! ten accepted message types, but only to validate every length-prefixed field (array count, +//! eleven accepted message types, but only to validate every length-prefixed field (array count, //! string length, bytes length, tagged-field size) against what could still fit in the bytes //! remaining in the frame - it never materializes a value or allocates a collection. Call the //! matching `validate_*_shape` function before handing the body to `kafka_protocol`. @@ -905,6 +905,54 @@ pub fn validate_heartbeat_shape(version: i16, body: &Bytes) -> Result<()> { Ok(()) } +/// Mirrors the field order `LeaveGroupRequest::decode` walks. +/// +/// Unlike Heartbeat this is bounded by `max_frame_size`: a v3+ response echoes every identity's +/// `member_id` and `group_instance_id`. +/// +/// # 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_leave_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)?; + } + if version <= 2 { + c.legacy_string(false)?; + } else { + let members_count = if flexible { + c.compact_array_count()? + } else { + c.legacy_array_count()? + }; + for _ in 0..members_count { + if flexible { + c.compact_string(false)?; + c.compact_string(true)?; + } else { + c.legacy_string(false)?; + c.legacy_string(true)?; + } + if version >= 5 { + c.compact_string(true)?; + } + if flexible { + c.tagged_fields()?; + } + } + } + if flexible { + c.tagged_fields()?; + } + Ok(()) +} + /// Mirrors the field order `SyncGroupRequest::decode` walks. /// /// # Errors @@ -963,7 +1011,10 @@ pub fn validate_sync_group_shape(version: i16, body: &Bytes, max_frame_size: usi #[cfg(test)] mod tests { + use kafka_protocol::messages::LeaveGroupRequest; + use super::*; + use crate::protocol::handlers::decode_exhaustive; const TEST_MAX_FRAME_SIZE: usize = 8 * 1024 * 1024; @@ -1244,4 +1295,76 @@ mod tests { ]); assert!(validate_sync_group_shape(5, &body, TEST_MAX_FRAME_SIZE).is_ok()); } + + fn leave_group_v4_body(member_id_len: usize) -> Bytes { + let mut body = vec![0x02, b'g', 0x02]; // group_id, members: 1 + let mut length = member_id_len + 1; + while length >= 0x80 { + body.push(u8::try_from(length & 0x7F).expect("masked to 7 bits") | 0x80); + length >>= 7; + } + body.push(u8::try_from(length).expect("below 0x80 after the loop")); + body.extend(std::iter::repeat_n(b'm', member_id_len)); + body.push(0x00); // group_instance_id null + body.push(0x00); // member tagged fields + body.push(0x00); // request tagged fields + Bytes::from(body) + } + + fn assert_decodes(version: i16, body: &Bytes) { + decode_exhaustive::<LeaveGroupRequest>(version, body.clone()) + .expect("kafka_protocol must agree the body is well formed"); + } + + #[test] + fn leave_group_v5_huge_members_count_rejected() { + let body = Bytes::from_static(&[0x02, b'g', 0xFF, 0xFF, 0xFF, 0xFF, 0x0F]); + assert!(validate_leave_group_shape(5, &body, TEST_MAX_FRAME_SIZE).is_err()); + } + + #[test] + fn leave_group_v0_minimal_body_accepted() { + let body = Bytes::from_static(&[ + 0x00, 0x01, b'g', // group_id + 0x00, 0x01, b'm', // member_id + ]); + assert!(validate_leave_group_shape(0, &body, TEST_MAX_FRAME_SIZE).is_ok()); + assert_decodes(0, &body); + } + + #[test] + fn leave_group_v3_null_group_instance_id_accepted() { + let body = Bytes::from_static(&[ + 0x00, 0x01, b'g', // group_id + 0x00, 0x00, 0x00, 0x01, // members: 1 + 0x00, 0x01, b'm', // member_id + 0xFF, 0xFF, // group_instance_id null + ]); + assert!(validate_leave_group_shape(3, &body, TEST_MAX_FRAME_SIZE).is_ok()); + assert_decodes(3, &body); + } + + /// `reason` arrives at v5. Reading it at v4 would swallow each member's tagged-fields byte + /// and walk the rest of the frame out of step with the decoder. + #[test] + fn leave_group_v4_two_members_accepted() { + let body = Bytes::from_static(&[ + 0x02, b'g', // group_id + 0x03, // members: 2 + 0x02, b'a', 0x00, 0x00, // member_id, group_instance_id null, tagged fields + 0x02, b'b', 0x00, 0x00, // member_id, group_instance_id null, tagged fields + 0x00, // request tagged fields + ]); + assert!(validate_leave_group_shape(4, &body, TEST_MAX_FRAME_SIZE).is_ok()); + assert_decodes(4, &body); + } + + /// From v3 every member id is echoed back, so it has to count against the response size. + #[test] + fn leave_group_member_strings_are_charged_against_the_response_budget() { + let body = leave_group_v4_body(4_096); + assert_decodes(4, &body); + assert!(validate_leave_group_shape(4, &body, 8 * 1024 * 1024).is_ok()); + assert!(validate_leave_group_shape(4, &body, 1_024).is_err()); + } } diff --git a/gateways/kafka/src/protocol/handlers/leave_group.rs b/gateways/kafka/src/protocol/handlers/leave_group.rs new file mode 100644 index 000000000..4787dd334 --- /dev/null +++ b/gateways/kafka/src/protocol/handlers/leave_group.rs @@ -0,0 +1,211 @@ +// 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. + +//! `LeaveGroup` (API key 13). +//! +//! How a consumer that closes cleanly hands its partitions back without making the rest of its +//! group wait out its session. Never parks: the removal wakes whoever is parked on the group. + +use bytes::Bytes; +use kafka_protocol::messages::leave_group_response::MemberResponse; +use kafka_protocol::messages::{LeaveGroupRequest, LeaveGroupResponse}; + +use crate::error::Result; +use crate::group::{LeaveRequest, LeaveResult}; +use crate::protocol::api::{ + API_KEY_LEAVE_GROUP, ApiVersionRange, ERROR_INVALID_REQUEST, ERROR_UNSUPPORTED_VERSION, + GatewayState, HandleOutcome, is_supported_version, +}; +use crate::protocol::bounds_guard::validate_leave_group_shape; +use crate::protocol::handlers::{ + decode_guarded, encode_message, respond_or_close, unsupported_version_response, +}; + +pub const RANGE: ApiVersionRange = ApiVersionRange { + api_key: API_KEY_LEAVE_GROUP, + min_version: 0, + max_version: 5, +}; + +/// Below this version the response has no `members` array and the encoder refuses one. +const FIRST_BATCHED_VERSION: i16 = 3; +/// Two length prefixes and the error code around each echoed identity, as encoded at v3. +const PER_MEMBER_OVERHEAD: usize = 6; + +pub async fn handle(state: &GatewayState, api_version: i16, body: Bytes) -> HandleOutcome { + if !is_supported_version(API_KEY_LEAVE_GROUP, api_version) { + return unsupported_version_response(API_KEY_LEAVE_GROUP, api_version, |version| { + encode_error_response(version, ERROR_UNSUPPORTED_VERSION) + }); + } + let request = match decode_guarded::<LeaveGroupRequest>(api_version, body, |version, body| { + validate_leave_group_shape(version, body, state.max_frame_size) + }) { + Ok(request) => request, + Err(error) => { + // debug!, not warn!: attacker-controlled, not operator-actionable. + tracing::debug!(%error, api_version, "Failed to decode LeaveGroup request"); + return respond_or_close( + encode_error_response(api_version, ERROR_INVALID_REQUEST), + "LeaveGroup", + ); + } + }; + if tracing::enabled!(tracing::Level::DEBUG) { + for reason in request + .members + .iter() + .filter_map(|member| member.reason.as_ref()) + { + tracing::debug!(%reason, "LeaveGroup reason"); + } + } + + let result = state + .groups + .leave(&LeaveRequest::from((api_version, &request))) + .await; + respond_or_close(encode_response(api_version, &result), "LeaveGroup") +} + +/// # Errors +/// +/// Returns an error when `kafka_protocol` cannot encode the response at `version`. +pub fn encode_response(version: i16, result: &LeaveResult) -> Result<Bytes> { + let response = if version < FIRST_BATCHED_VERSION { + LeaveGroupResponse::default().with_error_code(result.top_level_error()) + } else { + LeaveGroupResponse::default() + .with_error_code(result.error) + .with_members( + result + .members + .iter() + .map(|member| { + MemberResponse::default() + .with_member_id(member.member_id.clone()) + .with_group_instance_id(member.group_instance_id.clone()) + .with_error_code(member.error) + }) + .collect(), + ) + }; + let echoed: usize = if version < FIRST_BATCHED_VERSION { + 0 + } else { + result + .members + .iter() + .map(|member| { + member.member_id.len() + + member.group_instance_id.as_ref().map_or(0, |id| id.len()) + + PER_MEMBER_OVERHEAD + }) + .sum() + }; + encode_message(&response, version, 64 + echoed) +} + +/// # Errors +/// +/// Returns an error when `kafka_protocol` cannot encode the response at `version`. +pub fn encode_error_response(version: i16, error_code: i16) -> Result<Bytes> { + encode_response(version, &LeaveResult::error(error_code)) +} + +#[cfg(test)] +mod tests { + use kafka_protocol::protocol::StrBytes; + + use super::*; + use crate::group::LeftMember; + use crate::protocol::api::{ERROR_NONE, ERROR_UNKNOWN_MEMBER_ID}; + + fn left(member_id: &'static str, instance_id: Option<&'static str>, error: i16) -> LeftMember { + LeftMember { + member_id: StrBytes::from_static_str(member_id), + group_instance_id: instance_id.map(StrBytes::from_static_str), + error, + } + } + + fn result(members: Vec<LeftMember>) -> LeaveResult { + LeaveResult { + error: ERROR_NONE, + members, + } + } + + #[test] + fn given_a_member_error_at_v0_should_carry_it_at_top_level() { + let body = + encode_response(0, &result(vec![left("m-1", None, ERROR_UNKNOWN_MEMBER_ID)])).unwrap(); + + assert_eq!(body.as_ref(), &ERROR_UNKNOWN_MEMBER_ID.to_be_bytes()); + } + + /// The encoder refuses `members` below v3, and that refusal closes the connection. + #[test] + fn given_member_results_below_v3_should_encode_without_members() { + let body = encode_response(2, &result(vec![left("m-1", None, ERROR_NONE)])).unwrap(); + + assert_eq!(body.as_ref(), &[0, 0, 0, 0, 0, 0]); + } + + #[test] + fn given_a_v1_response_should_carry_throttle_time() { + let body = encode_error_response(1, ERROR_INVALID_REQUEST).unwrap(); + + assert_eq!(body.len(), 6); + assert_eq!(&body[4..], &ERROR_INVALID_REQUEST.to_be_bytes()); + } + + #[test] + fn given_a_v5_response_should_echo_each_identity_verbatim() { + let body = encode_response( + 5, + &result(vec![ + left("m-1", Some("i-1"), ERROR_NONE), + left("", None, ERROR_UNKNOWN_MEMBER_ID), + ]), + ) + .unwrap(); + + let expected: &[u8] = &[ + 0x00, 0x00, 0x00, 0x00, // throttle_time_ms + 0x00, 0x00, // error_code + 0x03, // members: 2 + 0x04, b'm', b'-', b'1', // member_id + 0x04, b'i', b'-', b'1', // group_instance_id + 0x00, 0x00, // error_code + 0x00, // member tagged fields + 0x01, // member_id: empty + 0x00, // group_instance_id: null + 0x00, 0x19, // error_code + 0x00, // member tagged fields + 0x00, // top-level tagged fields + ]; + assert_eq!(body.as_ref(), expected); + } + + #[test] + fn given_a_request_failure_at_v5_should_answer_with_an_empty_members_array() { + let body = encode_error_response(5, ERROR_INVALID_REQUEST).unwrap(); + + assert_eq!(body.as_ref(), &[0, 0, 0, 0, 0, 42, 0x01, 0x00]); + } +} diff --git a/gateways/kafka/src/protocol/handlers/mod.rs b/gateways/kafka/src/protocol/handlers/mod.rs index 21e1e9c16..45fab1a45 100644 --- a/gateways/kafka/src/protocol/handlers/mod.rs +++ b/gateways/kafka/src/protocol/handlers/mod.rs @@ -30,6 +30,7 @@ pub mod fetch; pub mod find_coordinator; pub mod heartbeat; pub mod join_group; +pub mod leave_group; pub mod list_offsets; pub mod metadata; pub mod produce; @@ -41,9 +42,10 @@ use kafka_protocol::protocol::{Decodable, Encodable}; use crate::error::{KafkaProtocolError, Result}; use crate::protocol::api::{ API_KEY_API_VERSIONS, API_KEY_CREATE_TOPICS, API_KEY_FETCH, API_KEY_FIND_COORDINATOR, - API_KEY_HEARTBEAT, API_KEY_JOIN_GROUP, API_KEY_LIST_OFFSETS, API_KEY_METADATA, API_KEY_PRODUCE, - API_KEY_SYNC_GROUP, ERROR_INVALID_REQUEST, ERROR_UNSUPPORTED_VERSION, GatewayState, - HandleOutcome, is_supported_version, supported_max_version, + API_KEY_HEARTBEAT, API_KEY_JOIN_GROUP, API_KEY_LEAVE_GROUP, API_KEY_LIST_OFFSETS, + API_KEY_METADATA, API_KEY_PRODUCE, API_KEY_SYNC_GROUP, ERROR_INVALID_REQUEST, + ERROR_UNSUPPORTED_VERSION, GatewayState, HandleOutcome, is_supported_version, + supported_max_version, }; /// Routes one decoded request body to the module that owns its API key. @@ -66,6 +68,7 @@ pub async fn dispatch( API_KEY_FIND_COORDINATOR => find_coordinator::handle(state, api_version, body).await, API_KEY_JOIN_GROUP => join_group::handle(state, api_version, body).await, API_KEY_HEARTBEAT => heartbeat::handle(state, api_version, body).await, + API_KEY_LEAVE_GROUP => leave_group::handle(state, api_version, body).await, API_KEY_SYNC_GROUP => sync_group::handle(state, api_version, body).await, _ => HandleOutcome::Close, } diff --git a/gateways/kafka/tests/common/scope.rs b/gateways/kafka/tests/common/scope.rs index c2ef894d0..2703cf4b4 100644 --- a/gateways/kafka/tests/common/scope.rs +++ b/gateways/kafka/tests/common/scope.rs @@ -34,6 +34,7 @@ pub const SCOPED_API_KEYS: &[(i16, &str, i16, i16)] = &[ (10, "FindCoordinator", 0, 4), (11, "JoinGroup", 0, 9), (12, "Heartbeat", 0, 4), + (13, "LeaveGroup", 0, 5), (14, "SyncGroup", 0, 5), ]; diff --git a/gateways/kafka/tests/common/wire.rs b/gateways/kafka/tests/common/wire.rs index 612e161d1..0d228ecdb 100644 --- a/gateways/kafka/tests/common/wire.rs +++ b/gateways/kafka/tests/common/wire.rs @@ -34,7 +34,6 @@ use super::codec::Encoder; pub const OUT_OF_SCOPE_API_KEYS: &[(i16, &str)] = &[ (8, "OffsetCommit"), (9, "OffsetFetch"), - (13, "LeaveGroup"), (15, "DescribeGroups"), (16, "ListGroups"), (17, "SaslHandshake"), @@ -578,6 +577,39 @@ pub fn build_heartbeat_request( enc.freeze() } +/// `LeaveGroup` request. Each member is `(member_id, group_instance_id, reason)`. Below v3 the +/// body carries one member id and no identities array, so only `members[0].0` is written. +pub fn build_leave_group_request( + version: i16, + group_id: &str, + members: &[(&str, Option<&str>, Option<&str>)], +) -> Bytes { + let flexible = version >= 4; + let mut enc = Encoder::with_capacity(64); + + write_string(&mut enc, flexible, Some(group_id)); + if version <= 2 { + let member_id = members.first().map_or("", |(member_id, _, _)| *member_id); + write_string(&mut enc, false, Some(member_id)); + } else { + write_array_count(&mut enc, flexible, members.len()); + for (member_id, group_instance_id, reason) in members { + write_string(&mut enc, flexible, Some(member_id)); + write_string(&mut enc, flexible, *group_instance_id); + if version >= 5 { + enc.write_compact_nullable_string(*reason); + } + if flexible { + enc.write_empty_tagged_fields(); + } + } + } + if flexible { + enc.write_empty_tagged_fields(); + } + enc.freeze() +} + /// Everything a `SyncGroup` body carries. `protocol_type`/`protocol_name` are written from v5. pub struct SyncGroupParams<'a> { pub group_id: &'a str, diff --git a/gateways/kafka/tests/consumer_group_tests.rs b/gateways/kafka/tests/consumer_group_tests.rs index 61048e295..aae185d5c 100644 --- a/gateways/kafka/tests/consumer_group_tests.rs +++ b/gateways/kafka/tests/consumer_group_tests.rs @@ -15,7 +15,8 @@ // specific language governing permissions and limitations // under the License. -//! Consumer group coordination: `FindCoordinator`, `JoinGroup`, Heartbeat, `SyncGroup`. +//! Consumer group coordination: `FindCoordinator`, `JoinGroup`, Heartbeat, `LeaveGroup`, +//! `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 @@ -39,14 +40,14 @@ use std::time::Duration; use bytes::Bytes; use tokio::io::AsyncWriteExt; use tokio::net::TcpStream; -use tokio::time::advance; +use tokio::time::{Instant, advance, timeout}; 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, + API_KEY_FIND_COORDINATOR, API_KEY_HEARTBEAT, API_KEY_JOIN_GROUP, API_KEY_LEAVE_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_UNKNOWN_MEMBER_ID, GatewayState, handle_request_bounded, @@ -57,13 +58,14 @@ 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, + build_join_group_request, build_leave_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 LEAVE_VERSION: i16 = 5; const SESSION_TIMEOUT_MS: i32 = 10_000; const REBALANCE_TIMEOUT_MS: i32 = 20_000; @@ -131,6 +133,22 @@ async fn heartbeat(state: &GatewayState, generation_id: i32, member_id: &str) -> error } +async fn leave( + state: &GatewayState, + version: i16, + members: &[(&str, Option<&str>)], +) -> LeaveResponse { + let identities: Vec<(&str, Option<&str>, Option<&str>)> = members + .iter() + .map(|(member_id, instance_id)| (*member_id, *instance_id, None)) + .collect(); + let body = build_leave_group_request(version, GROUP, &identities); + let response = handle_request_bounded(state, API_KEY_LEAVE_GROUP, version, body) + .await + .expect_response("LeaveGroup must answer"); + LeaveResponse::decode(version, response) +} + /// 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 { @@ -289,6 +307,51 @@ impl SyncResponse { } } +#[derive(Debug)] +struct LeaveResponse { + error: i16, + /// `(member_id, group_instance_id, error)`, present from v3. + members: Vec<(String, Option<String>, i16)>, +} + +impl LeaveResponse { + fn decode(version: i16, body: Bytes) -> Self { + let flexible = version >= 4; + let mut decoder = Decoder::new(body); + if version >= 1 { + decoder.read_i32().unwrap(); // throttle_time_ms + } + let error = decoder.read_i16().unwrap(); + let mut members = Vec::new(); + if version >= 3 { + let count = read_array_count(&mut decoder, flexible); + for _ in 0..count { + let member_id = + read_nullable(&mut decoder, flexible).expect("member id is not nullable"); + let instance_id = read_nullable(&mut decoder, flexible); + let member_error = decoder.read_i16().unwrap(); + if flexible { + decoder.read_tagged_fields().unwrap(); + } + members.push((member_id, instance_id, member_error)); + } + } + if flexible { + decoder.read_tagged_fields().unwrap(); + } + assert_eq!( + decoder.remaining(), + 0, + "LeaveGroup v{version} response has trailing bytes" + ); + Self { error, members } + } + + fn codes(&self) -> Vec<i16> { + self.members.iter().map(|(_, _, error)| *error).collect() + } +} + fn read_nullable(decoder: &mut Decoder, flexible: bool) -> Option<String> { if flexible { decoder.read_compact_nullable_string().unwrap() @@ -696,6 +759,337 @@ async fn given_every_member_expired_when_a_new_member_joins_should_start_a_fresh assert_ne!(second.member_id, first.member_id); } +// ── LeaveGroup ────────────────────────────────────────────────────────────── + +/// How long a woken waiter may take to answer. Paused time auto-advances to the next timer, so +/// a waiter still asleep on its own deadline loses this race instead of answering late. +const PROMPT: Duration = Duration::from_millis(5); + +async fn sync_leader(state: &GatewayState, generation_id: i32, leader: &str) { + let response = sync( + state, + SYNC_VERSION, + &SyncGroupParams { + group_id: GROUP, + generation_id, + member_id: leader, + ..SyncGroupParams::default() + }, + ) + .await; + assert_eq!(response.error, ERROR_NONE); +} + +/// Three members through one full rebalance: generation 2, the first one leading, all three +/// awaiting `SyncGroup`. +async fn three_member_group(state: &Arc<GatewayState>) -> [String; 3] { + let protocols: &[(&str, &[u8])] = &[("range", b"sub")]; + let leader = claim_member_id(state, b"sub").await; + let first = join(state, JOIN_VERSION, &join_params(&leader, protocols)).await; + assert_eq!(first.generation_id, 1); + + let mut parked = Vec::new(); + let mut followers = Vec::new(); + for _ in 0..2 { + let follower = claim_member_id(state, b"sub").await; + followers.push(follower.clone()); + let state = Arc::clone(state); + parked.push(tokio::spawn(async move { + let protocols: &[(&str, &[u8])] = &[("range", b"sub")]; + join(&state, JOIN_VERSION, &join_params(&follower, protocols)).await + })); + yield_to_parked().await; + } + + let rejoined = join(state, JOIN_VERSION, &join_params(&leader, protocols)).await; + assert_eq!(rejoined.generation_id, 2); + assert_eq!(rejoined.members.len(), 3); + for task in parked { + assert_eq!(task.await.expect("parked JoinGroup task").generation_id, 2); + } + let [second, third] = <[String; 2]>::try_from(followers).expect("two followers"); + [leader, second, third] +} + +/// Acceptance criterion: graceful shutdown releases partitions promptly. Paused time only moves +/// when something sleeps, so the elapsed bound proves no one waited out the leaver's session. +#[tokio::test(start_paused = true)] +async fn given_a_member_leaving_a_stable_group_should_let_the_survivor_rebalance_without_a_session_wait() + { + let state = test_state(immediate_config()); + let leader_protocols: &[(&str, &[u8])] = &[("range", b"leader-subscription")]; + let (leader, follower) = two_member_group(&state).await; + sync_leader(&state, 2, &leader).await; + let start = Instant::now(); + + let left = leave(&state, LEAVE_VERSION, &[(follower.as_str(), None)]).await; + assert_eq!(left.error, ERROR_NONE); + assert_eq!(left.codes(), vec![ERROR_NONE]); + assert_eq!( + heartbeat(&state, 2, &leader).await, + ERROR_REBALANCE_IN_PROGRESS + ); + let rejoined = join( + &state, + JOIN_VERSION, + &join_params(&leader, leader_protocols), + ) + .await; + + assert_eq!(rejoined.error, ERROR_NONE); + assert_eq!(rejoined.generation_id, 3); + let roster: Vec<&str> = rejoined.members.iter().map(|(id, _)| id.as_str()).collect(); + assert_eq!(roster, vec![leader.as_str()]); + assert!( + Instant::now() - start < Duration::from_secs(1), + "the survivor must not wait for the leaver's session to expire" + ); +} + +#[tokio::test(start_paused = true)] +async fn given_a_parked_rejoin_when_the_straggler_leaves_should_answer_it_promptly() { + let state = test_state(immediate_config()); + let (leader, follower) = two_member_group(&state).await; + sync_leader(&state, 2, &leader).await; + let parked = { + let state = Arc::clone(&state); + let leader = leader.clone(); + tokio::spawn(async move { + let protocols: &[(&str, &[u8])] = &[("range", b"leader-subscription")]; + join(&state, JOIN_VERSION, &join_params(&leader, protocols)).await + }) + }; + yield_to_parked().await; + assert!(!parked.is_finished(), "the leader waits for the follower"); + + leave(&state, LEAVE_VERSION, &[(follower.as_str(), None)]).await; + + let rejoined = timeout(PROMPT, parked) + .await + .expect("the leave must release the parked rejoin") + .expect("parked JoinGroup task"); + assert_eq!(rejoined.error, ERROR_NONE); + assert_eq!(rejoined.generation_id, 3); +} + +#[tokio::test(start_paused = true)] +async fn given_a_follower_parked_in_sync_when_another_member_leaves_should_answer_rebalance_in_progress() + { + let state = test_state(immediate_config()); + let [_leader, second, third] = three_member_group(&state).await; + let parked = { + let state = Arc::clone(&state); + tokio::spawn(async move { + sync( + &state, + SYNC_VERSION, + &SyncGroupParams { + group_id: GROUP, + generation_id: 2, + member_id: &second, + ..SyncGroupParams::default() + }, + ) + .await + }) + }; + yield_to_parked().await; + assert!(!parked.is_finished(), "a follower waits for the leader"); + + leave(&state, LEAVE_VERSION, &[(third.as_str(), None)]).await; + + let synced = timeout(PROMPT, parked) + .await + .expect("the leave must release the parked sync") + .expect("parked SyncGroup task"); + assert_eq!(synced.error, ERROR_REBALANCE_IN_PROGRESS); +} + +/// The member's `LeaveGroup` arrives on another connection while its own `JoinGroup` is parked, +/// which is what a client that timed out and reconnected sends. +#[tokio::test(start_paused = true)] +async fn given_a_leaving_member_with_a_parked_join_should_answer_that_join_unknown_member_id() { + let state = test_state(immediate_config()); + let (leader, follower) = two_member_group(&state).await; + sync_leader(&state, 2, &leader).await; + let parked = { + let state = Arc::clone(&state); + let follower = follower.clone(); + tokio::spawn(async move { + let changed: &[(&str, &[u8])] = &[("range", b"new-subscription")]; + join(&state, JOIN_VERSION, &join_params(&follower, changed)).await + }) + }; + yield_to_parked().await; + assert!(!parked.is_finished(), "the follower waits for the leader"); + + leave(&state, LEAVE_VERSION, &[(follower.as_str(), None)]).await; + + let answered = timeout(PROMPT, parked) + .await + .expect("the leave must answer the member's own parked join") + .expect("parked JoinGroup task"); + assert_eq!(answered.error, ERROR_UNKNOWN_MEMBER_ID); +} + +/// librdkafka only ever sends v0 or v1 and reads nothing but the top-level code. +#[tokio::test(start_paused = true)] +async fn given_a_leave_at_v0_should_answer_with_a_top_level_code() { + let state = test_state(immediate_config()); + let (_leader, follower) = two_member_group(&state).await; + + let first = leave(&state, 0, &[(follower.as_str(), None)]).await; + let second = leave(&state, 0, &[(follower.as_str(), None)]).await; + + assert_eq!(first.error, ERROR_NONE); + assert_eq!(second.error, ERROR_UNKNOWN_MEMBER_ID); +} + +/// Every advertised version must succeed through the handler, not only the ones a fixture or a +/// batching test happens to use: librdkafka sends v0-v1 and Java v3-v5. +#[tokio::test(start_paused = true)] +async fn given_each_supported_version_when_a_member_leaves_should_succeed() { + for version in 0..=5 { + let state = test_state(immediate_config()); + let (_leader, follower) = two_member_group(&state).await; + + let left = leave(&state, version, &[(follower.as_str(), None)]).await; + + assert_eq!(left.error, ERROR_NONE, "LeaveGroup v{version}"); + if version >= 3 { + assert_eq!(left.codes(), vec![ERROR_NONE], "LeaveGroup v{version}"); + } + } +} + +#[tokio::test(start_paused = true)] +async fn given_a_v3_leave_for_two_members_should_answer_each_in_request_order() { + let state = test_state(immediate_config()); + let (_leader, follower) = two_member_group(&state).await; + + let left = leave(&state, 3, &[("ghost", None), (follower.as_str(), None)]).await; + + assert_eq!(left.error, ERROR_NONE); + assert_eq!( + left.members, + vec![ + ("ghost".to_owned(), None, ERROR_UNKNOWN_MEMBER_ID), + (follower, None, ERROR_NONE), + ] + ); +} + +#[tokio::test(start_paused = true)] +async fn given_an_unknown_group_when_leaving_at_v5_should_answer_unknown_member_id_per_member() { + let state = test_state(immediate_config()); + + let left = leave(&state, 5, &[("ghost", None), ("", Some("instance"))]).await; + + assert_eq!(left.error, ERROR_NONE); + assert_eq!( + left.codes(), + vec![ERROR_UNKNOWN_MEMBER_ID, ERROR_UNKNOWN_MEMBER_ID] + ); + let protocols: &[(&str, &[u8])] = &[("range", b"sub")]; + let rejoin = join(&state, JOIN_VERSION, &join_params("ghost", protocols)).await; + assert_eq!(rejoin.error, ERROR_UNKNOWN_MEMBER_ID); +} + +/// A Java client raises `IllegalStateException` on more than one member response, and admin +/// `removeMembersFromConsumerGroup` keys results by the echoed pair, so a leave that removes +/// two holders of one instance id still answers once, with what was sent. +#[tokio::test(start_paused = true)] +async fn given_one_identity_that_removed_two_holders_should_answer_one_member() { + let state = test_state(immediate_config()); + let protocols: &[(&str, &[u8])] = &[("range", b"sub")]; + let static_params = |member_id| JoinGroupParams { + group_instance_id: Some("instance-a"), + ..join_params(member_id, protocols) + }; + let first = join(&state, JOIN_VERSION, &static_params("")) + .await + .member_id; + let joined = join(&state, JOIN_VERSION, &static_params(&first)).await; + assert_eq!(joined.generation_id, 1); + let second = join(&state, JOIN_VERSION, &static_params("")) + .await + .member_id; + let parked = { + let state = Arc::clone(&state); + tokio::spawn(async move { + let protocols: &[(&str, &[u8])] = &[("range", b"sub")]; + let params = JoinGroupParams { + group_instance_id: Some("instance-a"), + ..join_params(&second, protocols) + }; + join(&state, JOIN_VERSION, ¶ms).await + }) + }; + yield_to_parked().await; + + let left = leave(&state, LEAVE_VERSION, &[("", Some("instance-a"))]).await; + + assert_eq!( + left.members, + vec![(String::new(), Some("instance-a".to_owned()), ERROR_NONE)] + ); + let parked = timeout(PROMPT, parked) + .await + .expect("the removed holder's parked join must be answered") + .expect("parked JoinGroup task"); + assert_eq!(parked.error, ERROR_UNKNOWN_MEMBER_ID); + assert_eq!(heartbeat(&state, 1, &first).await, ERROR_UNKNOWN_MEMBER_ID); +} + +/// Partitions are opaque assignor blobs here, so what proves reassignment is the leader's +/// roster: each generation after a leave must list exactly the members still present. +#[tokio::test(start_paused = true)] +async fn given_n_members_when_they_leave_one_by_one_should_shrink_the_roster_each_generation() { + let state = test_state(immediate_config()); + let protocols: &[(&str, &[u8])] = &[("range", b"sub")]; + let [leader, second, third] = three_member_group(&state).await; + sync_leader(&state, 2, &leader).await; + + leave(&state, LEAVE_VERSION, &[(third.as_str(), None)]).await; + let parked = { + let state = Arc::clone(&state); + let leader = leader.clone(); + tokio::spawn(async move { + let protocols: &[(&str, &[u8])] = &[("range", b"sub")]; + join(&state, JOIN_VERSION, &join_params(&leader, protocols)).await + }) + }; + yield_to_parked().await; + let second_join = join(&state, JOIN_VERSION, &join_params(&second, protocols)).await; + let leader_join = parked.await.expect("parked JoinGroup task"); + assert_eq!(second_join.generation_id, 3); + assert_eq!(leader_join.generation_id, 3); + let mut roster: Vec<&str> = leader_join + .members + .iter() + .map(|(id, _)| id.as_str()) + .collect(); + roster.sort_unstable(); + let mut expected = vec![leader.as_str(), second.as_str()]; + expected.sort_unstable(); + assert_eq!(roster, expected); + sync_leader(&state, 3, &leader).await; + + leave(&state, LEAVE_VERSION, &[(second.as_str(), None)]).await; + let alone = join(&state, JOIN_VERSION, &join_params(&leader, protocols)).await; + assert_eq!(alone.generation_id, 4); + let roster: Vec<&str> = alone.members.iter().map(|(id, _)| id.as_str()).collect(); + assert_eq!(roster, vec![leader.as_str()]); + + leave(&state, LEAVE_VERSION, &[(leader.as_str(), None)]).await; + let fresh = join(&state, 3, &join_params("", protocols)).await; + assert_eq!(fresh.error, ERROR_NONE); + assert_eq!( + fresh.generation_id, 1, + "the emptied group is dropped, so the next one starts over" + ); +} + // ── Rejected requests ─────────────────────────────────────────────────────── #[tokio::test(start_paused = true)] @@ -992,6 +1386,96 @@ async fn given_two_tcp_clients_when_they_join_and_sync_should_each_receive_their assert_eq!(follower_sync.assignment.as_ref(), follower_blob); } +/// librdkafka's version (v1, header v1) and the flexible one (v5, header v2) over real sockets. +#[tokio::test] +async fn given_two_tcp_clients_when_one_leaves_should_let_the_other_rejoin_alone() { + for leave_version in [1, 5] { + tcp_leave_scenario(leave_version).await; + } +} + +async fn tcp_leave_scenario(leave_version: i16) { + 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 protocols: &[(&str, &[u8])] = &[("range", b"sub")]; + + let leader = tcp_join(&mut leader_stream, &join_params("", protocols)) + .await + .member_id; + tcp_join(&mut leader_stream, &join_params(&leader, protocols)).await; + let follower = tcp_join(&mut follower_stream, &join_params("", protocols)) + .await + .member_id; + write_request( + &mut follower_stream, + API_KEY_JOIN_GROUP, + JOIN_VERSION, + 2, + &build_join_group_request(JOIN_VERSION, &join_params(&follower, protocols)), + ) + .await; + let rejoined = tcp_join(&mut leader_stream, &join_params(&leader, protocols)).await; + assert_eq!(rejoined.generation_id, 2); + assert_eq!(read_join(&mut follower_stream).await.generation_id, 2); + write_request( + &mut leader_stream, + API_KEY_SYNC_GROUP, + SYNC_VERSION, + 3, + &build_sync_group_request( + SYNC_VERSION, + &SyncGroupParams { + group_id: GROUP, + generation_id: 2, + member_id: &leader, + ..SyncGroupParams::default() + }, + ), + ) + .await; + assert_eq!(read_sync(&mut leader_stream).await.error, ERROR_NONE); + + write_request( + &mut follower_stream, + API_KEY_LEAVE_GROUP, + leave_version, + 4, + &build_leave_group_request(leave_version, GROUP, &[(&follower, None, None)]), + ) + .await; + let payload = read_response_frame(&mut follower_stream, 8 * 1024 * 1024).await; + let (correlation_id, body) = + parse_response_payload(API_KEY_LEAVE_GROUP, leave_version, payload); + assert_eq!(correlation_id, 4); + let left = LeaveResponse::decode(leave_version, body); + assert_eq!(left.error, ERROR_NONE, "v{leave_version}"); + + write_request( + &mut leader_stream, + API_KEY_HEARTBEAT, + HEARTBEAT_VERSION, + 5, + &build_heartbeat_request(HEARTBEAT_VERSION, GROUP, 2, &leader), + ) + .await; + let payload = read_response_frame(&mut leader_stream, 8 * 1024 * 1024).await; + let (_, body) = parse_response_payload(API_KEY_HEARTBEAT, HEARTBEAT_VERSION, payload); + assert_eq!( + &body[4..6], + &ERROR_REBALANCE_IN_PROGRESS.to_be_bytes(), + "v{leave_version}" + ); + + let alone = tcp_join(&mut leader_stream, &join_params(&leader, protocols)).await; + assert_eq!(alone.generation_id, 3, "v{leave_version}"); + assert_eq!(alone.members.len(), 1, "v{leave_version}"); +} + async fn write_request( stream: &mut TcpStream, api_key: i16, diff --git a/gateways/kafka/tests/golden_wire_fixtures_tests.rs b/gateways/kafka/tests/golden_wire_fixtures_tests.rs index 5d73f1936..725e7c6d7 100644 --- a/gateways/kafka/tests/golden_wire_fixtures_tests.rs +++ b/gateways/kafka/tests/golden_wire_fixtures_tests.rs @@ -41,11 +41,11 @@ async fn golden_apiversions_v3_flexible_response_fixture() { .await .expect_response("test request has acks != 0 and expects a response"); - // error_code=0, api_count=10 (compact array: N+1=11) + // error_code=0, api_count=11 (compact array: N+1=12) // each entry followed by an empty tagged-fields byte; throttle_ms=0; top-level tagged fields - let expected: [u8; 78] = [ + let expected: [u8; 85] = [ 0x00, 0x00, // error_code - 0x0B, // compact array count (10+1) + 0x0C, // compact array count (11+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 @@ -55,6 +55,7 @@ async fn golden_apiversions_v3_flexible_response_fixture() { 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, 0x0D, 0x00, 0x00, 0x00, 0x05, 0x00, // key 13: LeaveGroup 0-5 0x00, 0x0E, 0x00, 0x00, 0x00, 0x05, 0x00, // key 14: SyncGroup 0-5 0x00, 0x00, 0x00, 0x00, // throttle_ms 0x00, // top-level tagged fields @@ -69,10 +70,10 @@ async fn golden_apiversions_v1_response_fixture() { .await .expect_response("test request has acks != 0 and expects a response"); - // error_code=0, api_count=10; Produce advertises min=0 per KAFKA-18659; throttle_ms=0 - let expected: [u8; 70] = [ + // error_code=0, api_count=11; Produce advertises min=0 per KAFKA-18659; throttle_ms=0 + let expected: [u8; 76] = [ 0x00, 0x00, // error_code - 0x00, 0x00, 0x00, 0x0A, // api count = 10 + 0x00, 0x00, 0x00, 0x0B, // api count = 11 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 @@ -82,6 +83,7 @@ async fn golden_apiversions_v1_response_fixture() { 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, 0x0D, 0x00, 0x00, 0x00, 0x05, // key 13: LeaveGroup 0-5 0x00, 0x0E, 0x00, 0x00, 0x00, 0x05, // key 14: SyncGroup 0-5 0x00, 0x00, 0x00, 0x00, // throttle_ms ]; diff --git a/gateways/kafka/tests/version_firewall_tests.rs b/gateways/kafka/tests/version_firewall_tests.rs index 9e976d987..f0c42881f 100644 --- a/gateways/kafka/tests/version_firewall_tests.rs +++ b/gateways/kafka/tests/version_firewall_tests.rs @@ -38,9 +38,10 @@ use tokio::net::TcpStream; use iggy_gateway_kafka::protocol::api::{ API_KEY_API_VERSIONS, API_KEY_CREATE_TOPICS, API_KEY_FETCH, API_KEY_FIND_COORDINATOR, - API_KEY_HEARTBEAT, API_KEY_JOIN_GROUP, API_KEY_LIST_OFFSETS, API_KEY_METADATA, API_KEY_PRODUCE, - API_KEY_SYNC_GROUP, ERROR_INVALID_REQUEST, ERROR_NONE, ERROR_UNSUPPORTED_VERSION, - advertised_min_version, handle_request, is_supported_version, supported_api_ranges, + API_KEY_HEARTBEAT, API_KEY_JOIN_GROUP, API_KEY_LEAVE_GROUP, API_KEY_LIST_OFFSETS, + API_KEY_METADATA, API_KEY_PRODUCE, API_KEY_SYNC_GROUP, ERROR_INVALID_REQUEST, ERROR_NONE, + ERROR_UNSUPPORTED_VERSION, advertised_min_version, handle_request, is_supported_version, + supported_api_ranges, }; use codec::Decoder; @@ -56,14 +57,14 @@ use wire::{ JoinGroupParams, OUT_OF_SCOPE_API_KEYS, SyncGroupParams, build_api_versions_flexible_request, build_create_topics_empty_request, build_fetch_empty_topics_request, build_find_coordinator_request, build_heartbeat_request, build_join_group_request, - build_list_offsets_request, build_metadata_all_topics_flexible, + build_leave_group_request, build_list_offsets_request, build_metadata_all_topics_flexible, build_metadata_all_topics_legacy, build_metadata_flexible_request_v10, build_sync_group_request, }; #[test] -fn supported_ranges_table_has_ten_entries() { - assert_eq!(supported_api_ranges().len(), 10); +fn supported_ranges_table_has_eleven_entries() { + assert_eq!(supported_api_ranges().len(), 11); } #[test] @@ -100,7 +101,7 @@ fn is_supported_version_matches_scope_table() { /// /// Relies on `SUPPORTED_RANGES` (src) and `SCOPED_API_KEYS` (test) sharing declaration order /// (Produce, Fetch, `ListOffsets`, Metadata, `ApiVersions`, `CreateTopics`, `FindCoordinator`, -/// `JoinGroup`, Heartbeat, `SyncGroup`) - `supported_ranges_table_has_ten_entries` plus +/// `JoinGroup`, Heartbeat, `LeaveGroup`, `SyncGroup`) - `supported_ranges_table_has_eleven_entries` plus /// `is_supported_version_matches_scope_table` already pin that both tables cover the same keys. #[tokio::test] async fn apiversions_advertises_exact_supported_ranges_v1() { @@ -348,7 +349,7 @@ async fn create_topics_below_min_version_closes_connection() { #[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, 15, 16, 17, 20, 42, 999] { let outcome = handle_request(key, 0, Bytes::new(), &default_broker()).await; assert!( outcome.is_close(), @@ -510,6 +511,9 @@ fn request_body_for_scoped_api(api_key: i16, name: &str, version: i16) -> Bytes }, ), API_KEY_HEARTBEAT => build_heartbeat_request(version, "scope-group", 1, "scope-member"), + API_KEY_LEAVE_GROUP => { + build_leave_group_request(version, "scope-group", &[("scope-member", None, None)]) + } API_KEY_SYNC_GROUP => build_sync_group_request( version, &SyncGroupParams { diff --git a/gateways/kafka/tools/kafka-tool/src/main.rs b/gateways/kafka/tools/kafka-tool/src/main.rs index 51421cbd6..0d70a23fe 100644 --- a/gateways/kafka/tools/kafka-tool/src/main.rs +++ b/gateways/kafka/tools/kafka-tool/src/main.rs @@ -26,6 +26,7 @@ use kafka_protocol::messages::delete_topics_request::*; use kafka_protocol::messages::describe_configs_request::*; use kafka_protocol::messages::fetch_request::*; use kafka_protocol::messages::join_group_request::*; +use kafka_protocol::messages::leave_group_request::MemberIdentity; use kafka_protocol::messages::list_offsets_request::*; use kafka_protocol::messages::offset_commit_request::*; use kafka_protocol::messages::produce_request::*; @@ -365,9 +366,16 @@ fn build_payload(api_key: i16, version: i16) -> Result<Bytes> { .context("OffsetFetch")?; } 10 => { - FindCoordinatorRequest::default() - .with_key(StrBytes::from_static_str("test-group")) - .with_key_type(0) + let key = StrBytes::from_static_str("test-group"); + let request = FindCoordinatorRequest::default().with_key_type(0); + // v4 replaced the single key with `coordinator_keys`; the encoder refuses whichever + // field the version does not carry. + let request = if version >= 4 { + request.with_coordinator_keys(vec![key]) + } else { + request.with_key(key) + }; + request .encode(&mut buf, version) .context("FindCoordinator")?; } @@ -394,11 +402,17 @@ fn build_payload(api_key: i16, version: i16) -> Result<Bytes> { .context("Heartbeat")?; } 13 => { - LeaveGroupRequest::default() - .with_group_id(GroupId::from(StrBytes::from_static_str("test-group"))) - .with_member_id(StrBytes::from_static_str("test-member-1")) - .encode(&mut buf, version) - .context("LeaveGroup")?; + let member_id = StrBytes::from_static_str("test-member-1"); + let request = LeaveGroupRequest::default() + .with_group_id(GroupId::from(StrBytes::from_static_str("test-group"))); + // v3 replaced the top-level member id with an identities array; the encoder refuses + // whichever field the version does not carry. + let request = if version >= 3 { + request.with_members(vec![MemberIdentity::default().with_member_id(member_id)]) + } else { + request.with_member_id(member_id) + }; + request.encode(&mut buf, version).context("LeaveGroup")?; } 14 => { SyncGroupRequest::default() @@ -863,3 +877,40 @@ async fn main() -> Result<()> { } Ok(()) } + +#[cfg(test)] +mod tests { + use kafka_protocol::protocol::Decodable; + + use super::*; + + /// A version the gateway advertises but this tool cannot build leaves the fixture-backed + /// suites without a fixture, which `KAFKA_FIXTURES_REQUIRED=1` turns into a CI failure. + #[test] + fn given_every_gateway_scoped_version_should_build_a_request() { + for (api_key, name, min_version, max_version) in gateway_verify_registry() { + for version in min_version..=max_version { + assert!( + build_framed(api_key, version, 1).is_ok(), + "{name} v{version} must build" + ); + } + } + } + + #[test] + fn given_a_leave_group_from_v3_should_carry_the_member_in_the_identities_array() { + for version in 3..=5 { + let mut payload = build_payload(13, version).expect("LeaveGroup builds"); + let request = LeaveGroupRequest::decode(&mut payload, version).expect("decodes"); + + assert!(request.member_id.is_empty(), "v{version}"); + assert_eq!(request.members.len(), 1, "v{version}"); + assert_eq!( + request.members[0].member_id.as_str(), + "test-member-1", + "v{version}" + ); + } + } +}
