This is an automated email from the ASF dual-hosted git repository.
hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/master by this push:
new 2f6923905 feat(gateways): handle Kafka LeaveGroup for consumer groups
(#4270)
2f6923905 is described below
commit 2f6923905f7eced7c1de3a579bbe5ae095a98cdf
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Mon Sep 28 14:52:16 2026 +0200
feat(gateways): handle Kafka LeaveGroup for consumer groups (#4270)
---
gateways/kafka/README.md | 9 +-
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 | 149 +++-
gateways/kafka/src/group/state.rs | 881 ++++++++++++++++++++-
gateways/kafka/src/protocol/api.rs | 7 +-
gateways/kafka/src/protocol/bounds_guard.rs | 124 ++-
.../kafka/src/protocol/handlers/leave_group.rs | 211 +++++
gateways/kafka/src/protocol/handlers/mod.rs | 10 +-
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 | 25 +-
gateways/kafka/tools/kafka-tool/src/main.rs | 35 +-
19 files changed, 2034 insertions(+), 70 deletions(-)
diff --git a/gateways/kafka/README.md b/gateways/kafka/README.md
index 8ef3cba9d..4b4bcb708 100644
--- a/gateways/kafka/README.md
+++ b/gateways/kafka/README.md
@@ -17,7 +17,12 @@ Foundation layer for
[apache/iggy#3421](https://github.com/apache/iggy/issues/34
>
> InitProducerId does real work too, with or without the bridge: it allocates
> a producer id, so a stock idempotent producer starts instead of failing at
> startup.
>
-> Consumer group coordination is not a stub either: `FindCoordinator`,
`JoinGroup`, `Heartbeat` and `SyncGroup` are real, with real membership,
rebalances and session expiry
([docs/CONSUMER_GROUPS.md](docs/CONSUMER_GROUPS.md)). With the bridge off,
Metadata reports every topic unknown, so a consumer joins a group and is
assigned 0 partitions. With it on, partitions are assigned, but offset
commit/fetch is not implemented and Fetch is still a stub, so nothing can be
consumed yet.
+> Consumer group coordination is not a stub either: `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)). With
the bridge off,
+> Metadata reports every topic unknown, so a consumer joins a group and is
assigned 0 partitions.
+> With it on, partitions are assigned, but offset commit/fetch is not
implemented and Fetch is
+> still a stub, so nothing can be consumed yet.
## Run
@@ -63,7 +68,7 @@ cargo test -p iggy-gateway-kafka
Or generate only the keys the tests need:
```bash
-for key in 0 1 2 10 11 12 14 19 22; do
+for key in 0 1 2 10 11 12 13 14 19 22; do
cargo run -p kafka-message-gen -- generate \
--output gateways/kafka/tools/kafka-tool/kafka_messages \
--api-key "$key"
diff --git a/gateways/kafka/docs/CONSUMER_GROUPS.md
b/gateways/kafka/docs/CONSUMER_GROUPS.md
index 315169d38..e77194539 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
@@ -122,11 +180,11 @@ that consumer would loop: coordinator connection closes,
client marks the coordi
re-runs FindCoordinator, retries OffsetFetch, closes again. 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. It is opt-in: a Kafka 4.0 client still defaults to
`group.protocol=classic`, which reaches
diff --git a/gateways/kafka/docs/MANUAL_TESTING.md
b/gateways/kafka/docs/MANUAL_TESTING.md
index f88db0ef9..cce7af5f1 100644
--- a/gateways/kafka/docs/MANUAL_TESTING.md
+++ b/gateways/kafka/docs/MANUAL_TESTING.md
@@ -219,15 +219,16 @@ 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 |
| 22 | InitProducerId | 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 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 |
+| B1 | ApiVersions negotiation | `error_code=0`; body lists 12 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/InitProducerId out-of-range |
N/A | **Connection closes** for both above-max and below-min -
`kafka_protocol`'s schema floor for each of these five 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,22 only — no OffsetCommit/OffsetFetch/LeaveGroup, and
no transaction keys (24, 25, 26, 28) |
+| B4 | ApiVersions lists only scoped keys | Decode response | Contains keys
0,1,2,3,10,11,12,13,14,18,19,22 only — no OffsetCommit/OffsetFetch, and no
transaction keys (24, 25, 26, 28) |
An out-of-range version only ever produces `error_code=35` on ApiVersions
(B1); every other API
key's out-of-range case closes the connection - see B2/B3. InitProducerId and
Produce also send
@@ -307,6 +308,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 | Both complete JoinGroup and SyncGroup and are assigned 0
partitions, because the Metadata stub reports `test` as unknown and the
assignor has nothing to hand out. 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, rejoins and is again
assigned 0 partitions |
| 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. `group.protocol` defaults to `classic` in
4.x; `--consumer-property group.protocol=consumer` sends ConsumerGroupHeartbeat
(68) instead 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`. Join, leave and sync are real,
but the assignment is empty: the Metadata stub reports the topic unknown, so
the assignor has no partitions to hand out |
+| 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.
@@ -410,7 +414,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–H6 Adversarial input
Automated regression:
@@ -446,4 +450,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 be8cae455..a331ca012 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; `TRANSACTIONAL_ID_AUTHORIZATION_FAILED` (53) for the transaction
key type, `INVALID_REQUEST` (42) for share; 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+ |
| 22 | InitProducerId | 0 | 5 | 0, 1, 2, 3, 4, 5 | Allocate a producer id
(epoch 0); a `transactional_id` gets `UNSUPPORTED_VERSION` (35); flexible
encoding at v2+ |
@@ -72,6 +73,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 |
@@ -87,7 +89,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 | Implemented behind `IGGY_KAFKA_SASL_ENABLED`,
advertised only while it is on ([`AUTHENTICATION.md`](AUTHENTICATION.md)) |
| 29 | DescribeAcls | Implemented behind `IGGY_KAFKA_SASL_ENABLED`, advertised
only while it is on ([`ACL_MAPPING.md`](ACL_MAPPING.md)) |
@@ -131,7 +132,7 @@ untouched, at at-least-once delivery. See
[`IDEMPOTENCE.md`](IDEMPOTENCE.md).
| 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 11 keys; `bounds_guard.rs`
pre-validates against unbounded allocation before handing a frame to the crate;
stub responses except InitProducerId and the four consumer-group keys, and
CreateTopics, Metadata, Produce and ListOffsets with a bridge |
+| **2 — Request/response codecs** | Partial | Decode/encode via the
`kafka_protocol` crate (broker feature only) for 12 keys; `bounds_guard.rs`
pre-validates against unbounded allocation before handing a frame to the crate;
stub responses except InitProducerId and the five consumer-group keys, and
CreateTopics, Metadata, Produce and ListOffsets with a bridge |
| **3 — Iggy bridge** | CreateTopics, Metadata, Produce and ListOffsets wired
| `bridge/` module (connection, topic mapping, provisioning, high watermark,
`topic_target` + `send_records`). CreateTopics
([#3538](https://github.com/apache/iggy/issues/3538)), Metadata
([#3534](https://github.com/apache/iggy/issues/3534)), Produce
([#3535](https://github.com/apache/iggy/issues/3535)) and ListOffsets
([#3537](https://github.com/apache/iggy/issues/3537)) call it. Fetch does not
call it yet ([# [...]
---
@@ -311,9 +312,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 fd529d90c..f8c2b6816 100644
--- a/gateways/kafka/docs/TEST_SUITE.md
+++ b/gateways/kafka/docs/TEST_SUITE.md
@@ -62,7 +62,7 @@ file under `tests/` anymore.
| [`idempotence_tests.rs`](../tests/idempotence_tests.rs) | `InitProducerId`
allocation across every supported version, and the transactional refusals on
`InitProducerId`/Produce | No |
| [`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 |
| [`sasl_tests.rs`](../tests/sasl_tests.rs) | SASL/PLAIN over a socket — full
handshake, every refusal path, and the disabled default. Drives a stub verifier
implementing `SaslAuthenticator`, so no Iggy server is needed | No |
diff --git a/gateways/kafka/docs/kafka_api_keys_reference.md
b/gateways/kafka/docs/kafka_api_keys_reference.md
index 3b1e1b444..1975e71f8 100644
--- a/gateways/kafka/docs/kafka_api_keys_reference.md
+++ b/gateways/kafka/docs/kafka_api_keys_reference.md
@@ -288,9 +288,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` (77 of the 88 API keys in this document)
+### Missing from `SUPPORTED_RANGES` (76 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. The gateway declines to
define a response for a
@@ -298,12 +299,14 @@ key it does not advertise, and a conforming client never
sends one, so no respon
be agreed. 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), opt-in via
`group.protocol=consumer` (the 4.0 default is still `classic`)
- **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 a946ab0fb..ec0793a49 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 10 11 12 14 19 22)
+FIXTURE_API_KEYS=(0 1 2 10 11 12 13 14 19 22)
usage() {
echo "Usage: $0 {generate|cleanup}" >&2
diff --git a/gateways/kafka/src/group/mod.rs b/gateways/kafka/src/group/mod.rs
index 41b3caa92..07cd1a891 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);
@@ -247,6 +247,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
@@ -354,6 +445,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! {
@@ -398,6 +495,7 @@ fn millis_to_duration(millis: i32) -> Duration {
mod tests {
use kafka_protocol::messages::GroupId;
use kafka_protocol::messages::join_group_request::JoinGroupRequestProtocol;
+ use kafka_protocol::messages::leave_group_request::MemberIdentity;
use
kafka_protocol::messages::sync_group_request::SyncGroupRequestAssignment;
use super::*;
@@ -461,6 +559,47 @@ mod tests {
);
}
+ #[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"))]);
+ }
+
#[test]
fn given_a_decoded_sync_when_normalizing_should_not_point_into_the_frame()
{
let frame = Bytes::from_static(b"g m consumer range f blob");
diff --git a/gateways/kafka/src/group/state.rs
b/gateways/kafka/src/group/state.rs
index 645455f1a..1597fb922 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,13 +35,13 @@ use tokio::time::Instant;
use uuid::Uuid;
use crate::group::{
- GroupCoordinatorConfig, JoinRequest, JoinResult, JoinedMember,
SyncRequest, SyncResult,
- owned_str,
+ GroupCoordinatorConfig, JoinRequest, JoinResult, JoinedMember,
LeaveRequest, LeaveResult,
+ LeavingMember, LeftMember, SyncRequest, SyncResult, owned_str,
};
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,
};
@@ -80,6 +80,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,
@@ -321,6 +332,93 @@ 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);
+ }
+ // A join window left open with nobody in it would close on a later
tick and clear the
+ // ids still on their way back. The next admitted member reopens the
barrier.
+ if self.members.is_empty() && !self.pending.is_empty() {
+ self.join_deadline = 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 {
@@ -364,14 +462,16 @@ impl GroupState {
fn complete_join(&mut self, now: Instant) {
self.members.retain(|_, member| member.rejoined);
- self.pending.clear();
self.join_deadline = None;
self.initial = false;
+ // A window that closes with nobody in it forms no generation, so the
ids still on their
+ // way back keep their own expiry and the next admitted member reopens
the barrier.
if self.members.is_empty() {
self.leader = None;
self.bump();
return;
}
+ self.pending.clear();
let leader = match self.leader.clone() {
Some(leader) if self.members.contains_key(&leader) => leader,
@@ -965,6 +1065,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,
@@ -2086,4 +2233,724 @@ 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));
+ }
+
+ /// Kafka keeps pending join members across the round that empties the
group, so a pending
+ /// client that rejoins after the old join window would have closed must
still be admitted.
+ #[test]
+ fn
given_a_pending_member_when_the_last_member_leaves_mid_rebalance_should_admit_the_rejoin_after_the_join_window()
+ {
+ let config = config();
+ let mut groups = Groups::new();
+ let now = Instant::now();
+ let (leader, follower) = stable_two_members(&mut groups, &config, now);
+ let _ = leave(&mut groups, &[(follower.as_str(), None)], now);
+ assert_eq!(groups[&group_id()].phase, Phase::PreparingRebalance);
+ let pending = member_id_of(&join_step(&mut groups, &config,
&pending_request(), now));
+
+ let result = leave(&mut groups, &[(leader.as_str(), None)], now);
+
+ assert_eq!(codes(&result), vec![ERROR_NONE]);
+ let later = now + Duration::from_secs(6);
+ let Step::Respond(rejoined) = join_step(
+ &mut groups,
+ &config,
+ &request(pending.as_str(), &["x"]),
+ later,
+ ) else {
+ panic!("the pending rejoin must be answered, not parked");
+ };
+ assert_eq!(rejoined.error, ERROR_NONE);
+ assert!(groups[&group_id()].members.contains_key(&pending));
+ }
+
+ /// The leave ticks the group before removing anyone, so a window already
past its deadline
+ /// closes first and drops the member that never rejoined. Its own leave
then finds it gone.
+ #[test]
+ fn
given_a_pending_member_when_the_last_member_leaves_after_the_join_window_should_admit_the_rejoin()
+ {
+ let config = config();
+ let mut groups = Groups::new();
+ let now = Instant::now();
+ let (leader, follower) = stable_two_members(&mut groups, &config, now);
+ let _ = leave(&mut groups, &[(follower.as_str(), None)], now);
+ let join_deadline = groups[&group_id()]
+ .join_deadline
+ .expect("the leave must open a join window");
+ let pending = member_id_of(&join_step(&mut groups, &config,
&pending_request(), now));
+
+ let result = leave(&mut groups, &[(leader.as_str(), None)],
join_deadline);
+
+ assert_eq!(codes(&result), vec![ERROR_UNKNOWN_MEMBER_ID]);
+ assert!(groups[&group_id()].pending.contains_key(&pending));
+ let Step::Respond(rejoined) = join_step(
+ &mut groups,
+ &config,
+ &request(pending.as_str(), &["x"]),
+ join_deadline,
+ ) else {
+ panic!("the pending rejoin must be answered, not parked");
+ };
+ assert_eq!(rejoined.error, ERROR_NONE);
+ assert!(groups[&group_id()].members.contains_key(&pending));
+ }
+
+ /// Same wipe through session expiry: the tick evicts the last member and
closes the overdue
+ /// window in one pass.
+ #[test]
+ fn
given_a_pending_member_when_the_last_member_expires_after_the_join_window_should_keep_the_pending_id()
+ {
+ let config = config();
+ let mut groups = Groups::new();
+ let now = Instant::now();
+ let (_, follower) = stable_two_members(&mut groups, &config, now);
+ let _ = leave(&mut groups, &[(follower.as_str(), None)], now);
+ let pending = member_id_of(&join_step(
+ &mut groups,
+ &config,
+ &JoinRequest {
+ session_timeout: Duration::from_secs(60),
+ ..pending_request()
+ },
+ now,
+ ));
+ let later = now + Duration::from_secs(11);
+
+ assert!(tick_group(&mut groups, &group_id(), later));
+ assert!(groups[&group_id()].members.is_empty());
+ assert!(groups[&group_id()].pending.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 076d0b708..05eb9ab7b 100644
--- a/gateways/kafka/src/protocol/api.rs
+++ b/gateways/kafka/src/protocol/api.rs
@@ -38,7 +38,8 @@ use crate::protocol::bounds_guard::{
use crate::protocol::handlers::init_producer_id::ProducerIdAllocator;
use crate::protocol::handlers::{
api_versions, create_topics, decode_guarded, dispatch, fetch,
find_coordinator, heartbeat,
- init_producer_id, join_group, list_offsets, metadata, produce,
respond_or_close, sync_group,
+ init_producer_id, join_group, leave_group, list_offsets, metadata,
produce, respond_or_close,
+ sync_group,
};
use crate::protocol::sasl::{
SaslMechanism, encode_sasl_authenticate_response,
encode_sasl_handshake_response,
@@ -51,6 +52,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_SASL_HANDSHAKE: i16 = 17;
pub const API_KEY_API_VERSIONS: i16 = 18;
@@ -168,6 +170,8 @@ pub const ERROR_UNSUPPORTED_COMPRESSION_TYPE: i16 =
/// 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;
/// Produce: a record or batch this gateway cannot map.
///
/// Not `CORRUPT_MESSAGE` (2), whose text fits but which `kafka-protocol`'s
table marks
@@ -258,6 +262,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 92e39a9f7..ae265b8ed 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`.
@@ -935,6 +935,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
@@ -1080,8 +1128,10 @@ pub fn validate_sasl_authenticate_shape(version: i16,
body: &Bytes) -> Result<()
#[cfg(test)]
mod tests {
use bytes::BytesMut;
+ use kafka_protocol::messages::LeaveGroupRequest;
use super::*;
+ use crate::protocol::handlers::decode_exhaustive;
const TEST_MAX_FRAME_SIZE: usize = 8 * 1024 * 1024;
@@ -1501,4 +1551,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 35fc2eb08..3b3d47a6b 100644
--- a/gateways/kafka/src/protocol/handlers/mod.rs
+++ b/gateways/kafka/src/protocol/handlers/mod.rs
@@ -31,6 +31,7 @@ pub mod find_coordinator;
pub mod heartbeat;
pub mod init_producer_id;
pub mod join_group;
+pub mod leave_group;
pub mod list_offsets;
pub mod metadata;
pub mod produce;
@@ -43,10 +44,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_INIT_PRODUCER_ID, 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_INIT_PRODUCER_ID, 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.
@@ -69,6 +70,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,
API_KEY_INIT_PRODUCER_ID => init_producer_id::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 1493a9f59..2926b0751 100644
--- a/gateways/kafka/tests/common/scope.rs
+++ b/gateways/kafka/tests/common/scope.rs
@@ -33,6 +33,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),
(18, "ApiVersions", 0, 3),
(19, "CreateTopics", 2, 5),
diff --git a/gateways/kafka/tests/common/wire.rs
b/gateways/kafka/tests/common/wire.rs
index c31272f3f..89db5f219 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"),
(20, "DeleteTopics"),
@@ -600,6 +599,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 0a5c0cd55..cb10c6608 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
@@ -41,14 +42,14 @@ use bytes::Bytes;
use kafka_protocol::protocol::StrBytes;
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,
SyncRequest};
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_TRANSACTIONAL_ID_AUTHORIZATION_FAILED,
@@ -60,13 +61,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;
@@ -136,6 +138,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 {
@@ -294,6 +312,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()
@@ -815,6 +878,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)]
@@ -1115,6 +1509,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 00fbe6ebc..e99b77584 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=11 (compact array: N+1=12)
+ // error_code=0, api_count=12 (compact array: N+1=13)
// each entry followed by an empty tagged-fields byte; throttle_ms=0;
top-level tagged fields
- let expected: [u8; 85] = [
+ let expected: [u8; 92] = [
0x00, 0x00, // error_code
- 0x0C, // compact array count (11+1)
+ 0x0D, // compact array count (12+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
@@ -53,6 +53,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, 0x12, 0x00, 0x00, 0x00, 0x03, 0x00, // key 18: ApiVersions
0-3
0x00, 0x13, 0x00, 0x02, 0x00, 0x05, 0x00, // key 19: CreateTopics
2-5
@@ -70,10 +71,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=11; Produce advertises min=0 per KAFKA-18659;
throttle_ms=0
- let expected: [u8; 76] = [
+ // error_code=0, api_count=12; Produce advertises min=0 per KAFKA-18659;
throttle_ms=0
+ let expected: [u8; 82] = [
0x00, 0x00, // error_code
- 0x00, 0x00, 0x00, 0x0B, // api count = 11
+ 0x00, 0x00, 0x00, 0x0C, // api count = 12
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
@@ -81,6 +82,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, 0x12, 0x00, 0x00, 0x00, 0x03, // key 18: ApiVersions 0-3
0x00, 0x13, 0x00, 0x02, 0x00, 0x05, // key 19: CreateTopics 2-5
diff --git a/gateways/kafka/tests/version_firewall_tests.rs
b/gateways/kafka/tests/version_firewall_tests.rs
index fa42052d5..63199e200 100644
--- a/gateways/kafka/tests/version_firewall_tests.rs
+++ b/gateways/kafka/tests/version_firewall_tests.rs
@@ -38,10 +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_INIT_PRODUCER_ID, 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_INIT_PRODUCER_ID, 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;
@@ -57,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_init_producer_id_request,
- build_join_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,
+ build_join_group_request, 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_eleven_entries() {
- assert_eq!(supported_api_ranges().len(), 11);
+fn supported_ranges_table_has_twelve_entries() {
+ assert_eq!(supported_api_ranges().len(), 12);
}
#[test]
@@ -101,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_twelve_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() {
@@ -349,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, 20, 42, 999] {
+ for key in [8, 9, 15, 16, 20, 42, 999] {
let outcome = handle_request(key, 0, Bytes::new(),
&default_broker()).await;
assert!(
outcome.is_close(),
@@ -506,6 +506,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 b79ccf644..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::*;
@@ -401,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()
@@ -873,6 +880,8 @@ async fn main() -> Result<()> {
#[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
@@ -888,4 +897,20 @@ mod tests {
}
}
}
+
+ #[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}"
+ );
+ }
+ }
}