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, &params).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}"
+            );
+        }
+    }
 }

Reply via email to