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 91df1da99 feat(gateways): support idempotent Kafka producers, refuse
transactions (#4250)
91df1da99 is described below
commit 91df1da997efd0d7512f97efeef693037cae2f33
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Mon Sep 28 13:20:50 2026 +0200
feat(gateways): support idempotent Kafka producers, refuse transactions
(#4250)
---
gateways/kafka/README.md | 49 +--
gateways/kafka/docs/IDEMPOTENCE.md | 139 +++++--
gateways/kafka/docs/MANUAL_TESTING.md | 38 +-
gateways/kafka/docs/SCOPE.md | 58 ++-
gateways/kafka/docs/TEST_SUITE.md | 1 +
gateways/kafka/docs/kafka_api_keys_reference.md | 30 +-
gateways/kafka/scripts/ci-wire-fixtures.sh | 2 +-
gateways/kafka/src/main.rs | 84 +++-
gateways/kafka/src/protocol/api.rs | 23 +-
gateways/kafka/src/protocol/bounds_guard.rs | 60 +++
.../src/protocol/handlers/init_producer_id.rs | 207 ++++++++++
gateways/kafka/src/protocol/handlers/mod.rs | 19 +-
gateways/kafka/src/protocol/handlers/produce.rs | 19 +-
gateways/kafka/src/server.rs | 58 ++-
gateways/kafka/tests/api_handler_tests.rs | 8 +-
gateways/kafka/tests/common/scope.rs | 1 +
gateways/kafka/tests/common/wire.rs | 23 ++
.../kafka/tests/create_topics_real_bridge_tests.rs | 1 +
gateways/kafka/tests/golden_wire_fixtures_tests.rs | 64 +--
gateways/kafka/tests/idempotence_tests.rs | 430 +++++++++++++++++++++
.../kafka/tests/list_offsets_real_bridge_tests.rs | 1 +
gateways/kafka/tests/listener_robustness_tests.rs | 2 +-
gateways/kafka/tests/metadata_real_bridge_tests.rs | 6 +-
gateways/kafka/tests/produce_real_bridge_tests.rs | 1 +
gateways/kafka/tests/server_e2e_tests.rs | 2 +-
gateways/kafka/tests/version_firewall_tests.rs | 11 +-
gateways/kafka/tools/kafka-tool/src/response.rs | 18 +-
27 files changed, 1183 insertions(+), 172 deletions(-)
diff --git a/gateways/kafka/README.md b/gateways/kafka/README.md
index 6d7b632e2..8930ac3cf 100644
--- a/gateways/kafka/README.md
+++ b/gateways/kafka/README.md
@@ -2,25 +2,20 @@
Foundation layer for
[apache/iggy#3421](https://github.com/apache/iggy/issues/3421): a TCP listener
on the Kafka wire port that decodes requests, validates scoped API keys and
versions. With a bridge, Produce writes to Iggy and ListOffsets reads offsets
from it. Everything else is a stub.
-> **Stub warning:** Produce and Fetch still don't persist or read real data -
they return
+> **Stub warning:** Fetch does not read real data yet: with or without a
bridge it returns
> retriable `NOT_LEADER_OR_FOLLOWER` (6) so clients keep data locally / retry
> elsewhere instead of
-> trusting a fake success. CreateTopics, Metadata, and ListOffsets are wired
to the Iggy bridge:
-> with `IGGY_KAFKA_BRIDGE_ENABLED=true`, CreateTopics creates a real Iggy
stream/topic, Metadata
-> reports real topics and partition counts (a topic not requested by name and
not found is
-> silently absent from a null-topics "list all" response, and
`UNKNOWN_TOPIC_OR_PARTITION` when
-> named explicitly), and ListOffsets answers `EARLIEST`/`LATEST` from real
partition state; with
-> the bridge off (the default), all three stay stubs - CreateTopics answers
`NOT_CONTROLLER` (41),
-> Metadata reports every requested topic unknown, and ListOffsets answers
`NOT_LEADER_OR_FOLLOWER`
-> (6). **CreateTopics has no authentication gate yet**: with the bridge on,
any client that can
-> reach this port can create topics (up to 1000 partitions each) as the
bridge's own Iggy user,
-> until SASL ([#3549](https://github.com/apache/iggy/issues/3549)) lands. See
-> **Stub warning:** When you set `IGGY_KAFKA_BRIDGE_ENABLED=true`, Produce
writes to Iggy and
-> ListOffsets answers `EARLIEST`/`LATEST` from real partition state. No other
API stores or reads
-> real data. Produce and ListOffsets without a bridge, and Fetch with or
without one, answer
-> retriable `NOT_LEADER_OR_FOLLOWER` (6). Clients then keep their data and do
not trust a fake
-> success. CreateTopics answers `NOT_CONTROLLER` (41) and creates nothing.
Metadata reports every
-> topic as unknown, so a real client cannot reach Produce or ListOffsets yet.
See
+> trusting a fake success. With `IGGY_KAFKA_BRIDGE_ENABLED=true`, Produce
writes to Iggy,
+> CreateTopics creates a real Iggy stream/topic, Metadata reports real topics
and partition counts
+> (a topic not requested by name and not found is silently absent from a
null-topics "list all"
+> response, and `UNKNOWN_TOPIC_OR_PARTITION` when named explicitly), and
ListOffsets answers
+> `EARLIEST`/`LATEST` from real partition state. With the bridge off (the
default) these stay stubs: Produce and
+> ListOffsets answer `NOT_LEADER_OR_FOLLOWER` (6), CreateTopics answers
`NOT_CONTROLLER` (41), and
+> Metadata reports every requested topic unknown. **CreateTopics runs as the
bridge's own Iggy
+> user**: with the bridge on and `IGGY_KAFKA_SASL_ENABLED` off (the default),
any client that can
+> reach this port can create topics (up to 1000 partitions each). See
> [docs/SCOPE.md](docs/SCOPE.md).
+>
+> 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.
## Run
@@ -41,6 +36,7 @@ Default bind: `127.0.0.1:9093`. Environment variables:
| `IGGY_KAFKA_READ_TIMEOUT_SECS` | `15` | Seconds allowed to read a frame body
once its length prefix arrives |
| `IGGY_KAFKA_WRITE_TIMEOUT_SECS` | `10` | Seconds allowed to write a response
frame |
| `IGGY_KAFKA_SHUTDOWN_DRAIN_TIMEOUT_SECS` | `25` | Seconds graceful shutdown
waits for in-flight connections before abandoning them |
+| `IGGY_KAFKA_INSTANCE_ID` | `0` | This gateway's number among the gateways
fronting one Iggy cluster. It is the high half of every producer id
`InitProducerId` hands out, and Kafka requires those to be cluster-unique, so
give every gateway its own value. A single gateway can leave it at `0`. |
| `IGGY_KAFKA_BRIDGE_ENABLED` | `false` | Connect the Iggy bridge at startup.
While false every API answers with its stub, and the `IGGY_KAFKA_IGGY_*`
variables below are read by nothing. A failed connection is fatal, not a
downgrade to stubs. |
| `IGGY_KAFKA_SASL_ENABLED` | `false` | Require SASL/PLAIN authentication
before serving any other API (`true` or `false`, nothing else) |
| `IGGY_KAFKA_PRE_AUTH_TIMEOUT_SECS` | `15` | Seconds an unauthenticated
connection may sit between frames. Waiting for an authentication slot and the
verification itself each get this budget, the verification's starting once it
holds a slot. Separate from the 10-minute idle timeout that applies once
authenticated |
@@ -65,7 +61,7 @@ cargo test -p iggy-gateway-kafka
Or generate only the keys the tests need:
```bash
-for key in 0 1 2 19; do
+for key in 0 1 2 19 22; do
cargo run -p kafka-message-gen -- generate \
--output gateways/kafka/tools/kafka-tool/kafka_messages \
--api-key "$key"
@@ -91,17 +87,22 @@ See [docs/SCOPE.md](docs/SCOPE.md) for
[#3421](https://github.com/apache/iggy/is
### Delivery guarantees
Delivery through this gateway is **at-least-once**, and stays at-least-once
across a gateway
-restart. Transactions are not supported, and will not be. A retry after a
timeout can write a
-record twice.
-
-Java producers must set `enable.idempotence=false` until the gateway serves
`InitProducerId`
-([#3545](https://github.com/apache/iggy/issues/3545)). Produce already stores
idempotent batches.
-It ignores their producer id, epoch and sequence, so it does not deduplicate a
retry.
+restart. Transactions are not supported, and will not be. An idempotent Kafka
producer is given
+a producer id so that it starts, but its retries are not deduplicated: Produce
stores idempotent
+batches but ignores their producer id, epoch and sequence, so a retry after a
network timeout
+writes the record twice, and both copies reach the stream with their own
offsets.
Iggy deduplicates writes on its own partition plane, and that does not close
this gap, because it
guards the hop from the gateway to Iggy rather than the hop from the producer
to the gateway.
[docs/IDEMPOTENCE.md](docs/IDEMPOTENCE.md) has the detail and what closing it
needs.
+Transactions are **not supported**, and will not be. AddPartitionsToTxn (24),
AddOffsetsToTxn
+(25), EndTxn (26) and TxnOffsetCommit (28) are never advertised, so a
conforming client never
+sends one; `InitProducerId` with a `transactional_id` answers
`UNSUPPORTED_VERSION` (35); and a
+Produce request carrying a `transactional_id` gets `UNSUPPORTED_VERSION` (35)
on every partition
+rather than having its records stored as if they were ordinary ones. None of
those closes the
+connection. `docs/SCOPE.md`'s Transactions section has the ordering and the
reasoning.
+
## Authentication ([#3549](https://github.com/apache/iggy/issues/3549))
Off by default. With `IGGY_KAFKA_SASL_ENABLED=true` the gateway requires
SASL/PLAIN before it serves
diff --git a/gateways/kafka/docs/IDEMPOTENCE.md
b/gateways/kafka/docs/IDEMPOTENCE.md
index 3ff1f7f31..2d7af0a36 100644
--- a/gateways/kafka/docs/IDEMPOTENCE.md
+++ b/gateways/kafka/docs/IDEMPOTENCE.md
@@ -1,8 +1,7 @@
# InitProducerId and idempotent producers
-Status: proposed. Answers the open half of
-[#3545](https://github.com/apache/iggy/issues/3545) and gates the Phase 1
end-to-end test
-([#3539](https://github.com/apache/iggy/issues/3539)), which drives
+Status: implemented, [#3545](https://github.com/apache/iggy/issues/3545).
Gates the Phase 1
+end-to-end test ([#3539](https://github.com/apache/iggy/issues/3539)), which
drives
`kafka-console-producer.sh`.
## The problem
@@ -11,12 +10,10 @@ A stock Java producer sets `enable.idempotence=true`
without being asked. That d
in Kafka 3.0 and took effect from 3.0.1, 3.1.1 and 3.2.0, where a bug that
suppressed it was
fixed. `kafka-console-producer.sh` leaves it on.
-An idempotent producer sends InitProducerId (key 22) before its first record.
The gateway does
-not list key 22, so ApiVersions does not advertise it, and the producer raises
-`UnsupportedVersionException`. That exception is fatal.
-`TransactionManager.maybeTransitionToErrorState` tests it above the
`isTransactional()` branch.
-The producer therefore enters a fatal error state instead of dropping back to
weaker semantics.
-It fails at startup, before it sends a record.
+An idempotent producer sends InitProducerId (key 22) before its first record.
A gateway that
+does not list key 22 does not advertise it either, and the producer raises
+`UnsupportedVersionException` rather than dropping back to weaker semantics.
It fails at
+startup, before it sends a record.
The gateway's stated purpose is that a Kafka user swaps the broker and changes
no application
code. A broker that the default producer cannot start against does not meet it.
@@ -108,9 +105,10 @@ delivery where it already is, and blocks nothing the pool
later needs.
## Behavior
-Add key 22 to `SUPPORTED_RANGES` in `src/protocol/api.rs` and advertise it
through ApiVersions.
-Without both, the producer never sends the request. `kafka-protocol` 0.18
carries the schemas,
-request v0 to v5 and response v0 to v6, flexible from v2.
+Key 22 is in `SUPPORTED_RANGES` (`src/protocol/api.rs`) and therefore
advertised through
+ApiVersions. Without both, the producer never sends the request.
`kafka-protocol` 0.18 carries
+the schemas, request v0 to v5 and response v0 to v6, flexible from v2; the
gateway serves v0 to
+v5. The handler is `src/protocol/handlers/init_producer_id.rs`.
InitProducerId with no `transactional_id`:
@@ -126,28 +124,104 @@ a birthday collision, not a remote one.
The id is a pool key, not a dedup identity. Under the design above, the dedup
identity is the
session's own random client id, minted at register. The producer id only
decides which connection
serves a producer. Kafka still requires it to be unique across the cluster,
which is what the
-instance number buys. It does not have to survive a restart.
+instance number buys. It also has to stay unique across a restart once the
pool lands or
+Produce persists, because `producer_epoch` is always 0, so a replayed id is a
replayed
+`(producer_id, producer_epoch)` pair a live producer may still hold. The
counter therefore starts
+at the wall clock in milliseconds rather than at 0. A restarted gateway starts
above every id its
+previous run handed out unless that run averaged more than one allocation per
millisecond of its
+uptime, or the clock stepped back across the restart. A clock that reads
before the Unix epoch
+seeds 0 and replays from the bottom. 2^47 milliseconds leaves the counter
space thousands of years
+from running out.
+
+Nothing is persisted for this, so the clock seed narrows the replay window
without closing it. That
+is harmless while no code path reads the producer id. It stops being harmless
once one does:
+[#3535](https://github.com/apache/iggy/issues/3535) must first either persist
a high-water mark
+and seed the counter above it, or bump `producer_epoch` on every gateway
start. Either one makes a
+replayed `(producer_id, producer_epoch)` pair impossible rather than unlikely.
+
+`IGGY_KAFKA_INSTANCE_ID` defaults to 0 so a single gateway needs no
configuration. Two gateways
+left on the default draw from the same 47-bit space and can hand out identical
ids, on the first
+allocation if they start in the same millisecond. Nothing in the cluster
detects that, so a gateway started without the
+variable logs a warning, and the running value is in the startup log line.
+
+An empty `transactional_id` reads as absent. A wire null decodes to `None`, but
+`kafka-protocol`'s own `Default` is `Some("")`, and a producer that is
idempotent-only names no
+transaction either way.
InitProducerId with a `transactional_id`:
-- answer `UNSUPPORTED_VERSION` (35), unchanged. Transactions stay out of
scope, and so do
+- answer `UNSUPPORTED_VERSION` (35). Transactions stay out of scope, and so do
AddPartitionsToTxn (24), AddOffsetsToTxn (25), EndTxn (26) and
TxnOffsetCommit (28)
-35 rather than `INVALID_REQUEST` (42), because of the same fatal set quoted
above.
-`maybeTransitionToErrorState` holds ClusterAuthorization,
TransactionalIdAuthorization,
-ProducerFenced, UnsupportedVersion and InvalidPidMapping. `INVALID_REQUEST` is
not in it, so a
-transactional producer moves to an abortable error instead. The application is
then told to abort
-and retry something that can never succeed. `COORDINATOR_NOT_AVAILABLE` (15)
is worse again. It
-is retriable, so the producer never stops trying.
+Not for the reason the Produce path uses. `maybeTransitionToErrorState`
governs a failed Produce
+batch and never sees an InitProducerId response; those reach
+`InitProducerIdHandler.handleResponse`, whose trailing `else` is
`fatalError(new
+KafkaException("Unexpected error in InitProducerIdResponse; ..."))`. Anything
it does not
+recognise is fatal there, so `INVALID_REQUEST` (42) would be equally fatal and
the "42 is
+abortable, therefore 35" argument does not apply to this API. 35 is chosen for
consistency with
+the Produce guard below, and because it is the one code that also states the
truth: the gateway
+does not implement this version of the transactional protocol. What must be
avoided is a
+*retriable* code: that same handler re-enqueues `COORDINATOR_LOAD_IN_PROGRESS`
(14) and
+`CONCURRENT_TRANSACTIONS` (51), so the producer would never stop trying.
+
+No single response code is terminal on both target clients. 35 is fatal for
the Java producer
+and an infinite retry for librdkafka, whose `rd_kafka_idemp_check_error`
treats only
+`__UNSUPPORTED_FEATURE`, `INVALID_TRANSACTION_TIMEOUT` (50), 53 and 31 as
fatal. 35 stays the
+choice, and librdkafka is stopped at FindCoordinator instead: key 10 is
unadvertised today, so
+`rd_kafka_init_transactions()` fails before InitProducerId is reached. Phase 3
advertises key 10
+for consumer groups and has to refuse a `TXN`-type coordinator lookup
explicitly.
Produce:
- accept `producer_id`, `producer_epoch` and `base_sequence` on the request
and ignore them
- never answer `OUT_OF_ORDER_SEQUENCE_NUMBER` (45) or
`DUPLICATE_SEQUENCE_NUMBER` (46)
-
-Those two codes stay unsent even once the pool lands. The watermark accepts
any request above it
-without noticing a gap, so a gap cannot be told apart from ordinary traffic.
Sending either code
-claims a detection the gateway does not have.
+- reject a non-empty `transactional_id` (v3+) with `UNSUPPORTED_VERSION` (35),
at the partition
+ level, keeping the connection open. Under `acks=0` there is no response to
carry 35, so the
+ gateway closes the connection instead, as a Kafka broker does on any
`acks=0` produce error.
+ The check runs before the `acks=0` branch, so no write path can see a
transactional batch
+
+Those first two codes stay unsent even once the pool lands. The watermark
accepts any request
+above it without noticing a gap, so a gap cannot be told apart from ordinary
traffic. Sending
+either code claims a detection the gateway does not have.
+
+The third is where `maybeTransitionToErrorState` is exact.
`Sender.completeBatch` ->
+`canRetry` false -> `failBatch` -> `handleFailedBatch` ->
`maybeTransitionToErrorState`, whose
+explicit fatal set holds ClusterAuthorization, TransactionalIdAuthorization,
ProducerFenced,
+UnsupportedVersion and InvalidPidMapping. `INVALID_TXN_STATE` (48) is
explicitly rewritten to
+abortable there, and `INVALID_REQUEST` (42) and
`UNSUPPORTED_FOR_MESSAGE_FORMAT` (43) fall
+through to abortable, so any of those would tell the application to abort and
retry something
+that can never succeed.
+
+Without this guard a transactional batch would land as ordinary records once
+[#3535](https://github.com/apache/iggy/issues/3535) wires the bridge: no last
stable offset, no
+abort markers, `read_committed` unimplementable, and an aborted transaction's
records delivered
+to every consumer.
+
+The guard reads the request-level `transactional_id` only. A record batch also
carries a
+transactional bit in its attributes, and the stub keeps records opaque, so a
hand-built frame
+that sets the bit without the request field gets the retriable stub error (6).
Java and
+librdkafka set both, so no stock client reaches it, and the stub does not
parse batches to close
+it. [#3535](https://github.com/apache/iggy/issues/3535) decodes each batch
before persisting it,
+and must refuse one with the transactional bit set there, the same way.
+
+## Invariants this design rests on
+
+Both are absences. Losing either is caught: `golden_wire_fixtures_tests.rs`
pins the ApiVersions
+v1 and v3 bodies byte-exactly, so adding a finalized feature or any advertised
key fails both
+goldens.
+
+**Never advertise `transaction.version >= 2` in the ApiVersions
`finalized_features`.**
+`TransactionManager.maybeUpdateTransactionV2Enabled` reads it, and under TV2
`maybeAddPartition`
+adds partitions client-side, so AddPartitionsToTxn and AddOffsetsToTxn are
never sent at all.
+That collapses the absence gate this design depends on.
`api_versions::encode_response` builds
+an `ApiVersionsResponse` with no finalized features, which is correct and
load-bearing.
+
+**The transactional handshake order is fixed**: FindCoordinator(TXN) ->
InitProducerId ->
+AddPartitionsToTxn. While FindCoordinator (key 10) stays unadvertised,
`initTransactions()`
+already dies before InitProducerId is reached, and the InitProducerId branch
is belt and braces.
+Advertising FindCoordinator for consumer groups
+([#3541](https://github.com/apache/iggy/issues/3541)) removes that shield, and
is what makes the
+branch load-bearing.
## What this does not give you
@@ -158,15 +232,14 @@ Iggy's own deduplication does not help, because it guards
the other hop. Deliver
gateway is at-least-once until the pool lands, and at-least-once across a
gateway restart after
that.
-State that limitation in the README, next to the transaction section, in those
words. Do not
-leave a user to infer it from the presence of key 22.
-
-## Open question
+The README states that limitation under "Delivery guarantees", so a user does
not have to infer
+it from the presence of key 22.
-Allocate only, as above, or stub and document `enable.idempotence=false`?
+## Resolution
-If no answer lands by 2026-09-22, allocate only is taken and the work
proceeds. This document
-is then updated to record that it was decided by default.
+The open question was allocate only, as above, against stub and document
+`enable.idempotence=false`. Allocate only was taken, which this document named
as the default
+outcome, and is what shipped.
## References
@@ -176,4 +249,6 @@ is then updated to record that it was decided by default.
- Dedup key and window: `core/consensus/src/client_table.rs`
- Session identity: `core/sdk/src/session.rs`, `core/sdk/src/tcp/tcp_client.rs`
- Header rewrite: `core/server/src/dispatch/partition.rs`
-- Fatal path: `TransactionManager.maybeTransitionToErrorState`, apache/kafka
trunk
+- Produce fatal path: `TransactionManager.maybeTransitionToErrorState`,
apache/kafka trunk
+- InitProducerId fatal path:
`TransactionManager.InitProducerIdHandler.handleResponse`, same file
+- librdkafka's fatal set: `rd_kafka_idemp_check_error`,
`src/rdkafka_idempotence.c`
diff --git a/gateways/kafka/docs/MANUAL_TESTING.md
b/gateways/kafka/docs/MANUAL_TESTING.md
index 94cb2326d..1ddd50f83 100644
--- a/gateways/kafka/docs/MANUAL_TESTING.md
+++ b/gateways/kafka/docs/MANUAL_TESTING.md
@@ -38,13 +38,14 @@ kafka listener bound on 127.0.0.1:9093
```bash
# Terminal 2
-# Keys 0/1/2/19 match ci-wire-fixtures.sh: the only keys any test actually
loads a .bin
-# fixture for. Metadata (3) and ApiVersions (18) requests are built
synthetically in-test
-# instead, so fixtures for those keys are generated but unused - `generate`
still accepts
-# them if you want them for manual `send`/`verify` below.
+# Keys 0/1/2/19/22 match ci-wire-fixtures.sh's FIXTURE_API_KEYS: the only keys
any test
+# actually loads a .bin fixture for. Omit key 22 and the six InitProducerId
cases skip
+# silently, which reads as a pass. Metadata (3) and ApiVersions (18) requests
are built
+# synthetically in-test instead, so fixtures for those keys are generated but
unused -
+# `generate` still accepts them if you want them for manual `send`/`verify`
below.
cargo run -p kafka-message-gen -- generate \
--output gateways/kafka/tools/kafka-tool/kafka_messages \
- --api-key 0 --api-key 1 --api-key 2 --api-key 19
+ --api-key 0 --api-key 1 --api-key 2 --api-key 19 --api-key 22
```
---
@@ -200,7 +201,8 @@ unset). With a bridge, A5 writes one record to
`test-topic`: `ec=0` if it exists
| A6 | Fetch v4 | `send --host 127.0.0.1:9093 --api-key 1 --version 4` |
Decode + stub response | Top-level `ec=0`; per-partition `ec=6`
(NOT_LEADER_OR_FOLLOWER) |
| A7 | ListOffsets v1 | `send --host 127.0.0.1:9093 --api-key 2 --version 1` |
Decode + stub offsets | Per-partition `ec=6` (NOT_LEADER_OR_FOLLOWER) - no
top-level error field on this response |
| A8 | CreateTopics v2 | `send --host 127.0.0.1:9093 --api-key 19 --version 2`
| Decode + stub non-creation ack | `ec=41` (NOT_CONTROLLER) per topic |
-| A9 | Verify all scoped keys | `cargo run -p kafka-message-gen -- verify
--host 127.0.0.1:9093 --api-key 0 --api-key 1 --api-key 2 --api-key 3 --api-key
18 --api-key 19` | Exit code 0 | No timeouts or I/O errors (`verify` already
knows each stub's expected non-zero code - see `is_acceptable_verify_error` in
`kafka-tool/src/response.rs`) |
+| A9 | InitProducerId v4 | `send --host 127.0.0.1:9093 --api-key 22 --version
4` | Producer id allocated | `ec=0`, `producer_id >= 0`, `producer_epoch=0`; a
second send returns a different `producer_id` |
+| A10 | Verify all scoped keys | `cargo run -p kafka-message-gen -- verify
--host 127.0.0.1:9093 --api-key 0 --api-key 1 --api-key 2 --api-key 3 --api-key
18 --api-key 19 --api-key 22` | Exit code 0 | No timeouts or I/O errors
(`verify` already knows each stub's expected non-zero code - see
`is_acceptable_verify_error` in `kafka-tool/src/response.rs`) |
### Category B — Version firewall (boundary validation)
@@ -214,16 +216,18 @@ For each API key, test **min−1**, **min**, **max**,
**max+1** using `kafka-mes
| 1 | Fetch | 4 | 12 | 3, 4, 12, 13 |
| 2 | ListOffsets | 1 | 6 | 0, 1, 6, 7 |
| 19 | CreateTopics | 2 | 5 | 1, 2, 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 6 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 7 API keys with
correct min/max | KIP-511 exception: still answers, `error_code=35`
(UNSUPPORTED_VERSION), v0 response header regardless of the request's own
encoding |
| B2 | Metadata out-of-range | N/A | **Connection closes**, no response sent -
Metadata has no top-level error field to carry a version-correct error in |
-| B3 | Produce/Fetch/ListOffsets/CreateTopics out-of-range | N/A |
**Connection closes** for both above-max and below-min - `kafka_protocol`'s
schema floor for each of these four messages equals `SUPPORTED_RANGES`' own
min, so there is no encodable error response below min either (see `SCOPE.md`'s
Governance model) |
-| B4 | ApiVersions lists only scoped keys | Decode response | Contains keys
0,1,2,3,18,19 only — no consumer-group keys |
+| 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,18,19,22 only — no consumer-group keys, and no transaction keys (24,
25, 26, 28) |
-Only ApiVersions (B1) ever returns `error_code=35` on this gateway. Every
other API key's
-out-of-range case closes the connection - see B2/B3.
+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
+35 in range, for a transactional request - see Category H.
**Validation tip:** Use `--hex` when generating to inspect request bytes:
@@ -307,6 +311,9 @@ Record kcat version and exact error strings in your test
log. G1 passing is the
| H1 | Truncated Produce body | Send valid header + incomplete body | No
response. Connection closed. **No panic** |
| H2 | Random bytes | `dd if=/dev/urandom bs=64 count=1 \| nc 127.0.0.1 9093`
| Connection closed or protocol error; gateway stays up |
| H3 | Empty body after header | ApiVersions with valid header, empty body |
`ec=0` (ApiVersions accepts empty body) |
+| H4 | Transactional InitProducerId | Send key 22 v4 with a non-null
`transactional_id` | `ec=35` (UNSUPPORTED_VERSION); connection stays open (send
A2 next to confirm) |
+| H5 | Transactional Produce | Send key 0 v3 with a non-null
`transactional_id` and `acks=1` | `ec=35` per partition, **not** `ec=6`;
connection stays open. With `acks=0`: no response, and the gateway closes the
connection |
+| H6 | Transaction API keys | `send --host 127.0.0.1:9093 --api-key 24` (also
25, 26, 28) | Connection closes, no response bytes - they are never advertised |
---
@@ -325,7 +332,7 @@ Record kcat version and exact error strings in your test
log. G1 passing is the
| 17 | INVALID_TOPIC_EXCEPTION | Produce with a bridge: bad topic name |
| 21 | INVALID_REQUIRED_ACKS | Produce with a bridge: `acks` is not 0, 1 or -1
|
| 29 | TOPIC_AUTHORIZATION_FAILED | Produce with a bridge: Iggy user lacks
permission |
-| 35 | UNSUPPORTED_VERSION | **ApiVersions only** (KIP-511 exception), plus
Produce with a bridge: transactional or control batch. Every other API key's
out-of-range version closes the connection instead - see Category B |
+| 35 | UNSUPPORTED_VERSION | Out of range: **ApiVersions only** (KIP-511
exception); every other API key's out-of-range version closes the connection
instead - see Category B. In range: InitProducerId with a `transactional_id`;
Produce with a `transactional_id`, on every partition with or without a bridge
(an `acks=0` one closes instead); and, with a bridge, a transactional or
control batch - transactions are not supported, see `SCOPE.md` |
| 37 | INVALID_PARTITIONS | CreateTopics: partition count `0` or `< -1` (or
any non-positive on v2–v3) |
| 38 | INVALID_REPLICATION_FACTOR | CreateTopics: replication factor `0` or `<
-1` (or any non-positive on v2–v3) |
| 41 | NOT_CONTROLLER | CreateTopics stub (topic not created) |
@@ -350,6 +357,7 @@ Header version selection now delegates entirely to
`kafka_protocol::messages::Ap
| 1 Fetch | v12+ | v1 |
| 2 ListOffsets | v6+ | v1 |
| 19 CreateTopics | v5+ | v1 |
+| 22 InitProducerId | v2+ | v1 |
### Frame layout (for manual hex inspection)
@@ -390,14 +398,14 @@ Tester: ___________
Gateway commit: ___________
kcat version (if used): ___________
-[ ] A1–A9 Smoke tests
-[ ] B1–B4 Version firewall (all 6 keys × 4 boundary versions)
+[ ] A1–A10 Smoke tests
+[ ] B1–B4 Version firewall (all 7 keys × 4 boundary versions)
[ ] C1–C4 Unsupported API keys
[ ] 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)
-[ ] H1–H3 Adversarial input
+[ ] H1–H6 Adversarial input
Automated regression:
[ ] cargo test -p iggy-gateway-kafka — all passed (see `TEST_SUITE.md` for why
this checklist
diff --git a/gateways/kafka/docs/SCOPE.md b/gateways/kafka/docs/SCOPE.md
index a9773f6b3..f878e6787 100644
--- a/gateways/kafka/docs/SCOPE.md
+++ b/gateways/kafka/docs/SCOPE.md
@@ -12,7 +12,7 @@ Foundation layer only: a TCP listener on the Kafka wire port
that decodes reques
| Length-prefixed frame read/write with `max_frame_size` cap | Done |
`src/server.rs` |
| Request header v1/v2 auto-detection | Done | `src/protocol/header.rs`
(delegates to `kafka_protocol::messages::ApiKey`) |
| Version negotiation firewall (`SUPPORTED_RANGES`) | Done |
`src/protocol/api.rs` |
-| Request decode + stub encode for 6 API keys | Done | `src/protocol/api.rs`,
`responses.rs` (via the `kafka_protocol` crate) |
+| Request decode + stub encode for 7 API keys | Done | `src/protocol/api.rs`,
`src/protocol/handlers/` (via the `kafka_protocol` crate) |
| Produce hot path: RecordBatch as opaque `Bytes` | Done |
`src/protocol/responses.rs` |
| Pre-decode bounds guard against unbounded allocation | Done |
`src/protocol/bounds_guard.rs` |
| Graceful errors (corrupt decode, invalid header) | Done |
`src/protocol/api.rs`, `src/server.rs` |
@@ -27,11 +27,12 @@ Source of truth for supported ranges: `SUPPORTED_RANGES` in
[`src/protocol/api.r
Expand `SUPPORTED_RANGES` only after a key/version pair is manually tested.
ApiVersions advertises exactly what the firewall allows.
**Every unsupported-version case closes the connection, for every listed key**
- not just above
-the encoder max. `kafka_protocol`'s schema floor for each of the six supported
messages happens
-to equal `SUPPORTED_RANGES`' own min today (Produce 3, Fetch 4, ListOffsets 1,
Metadata 0,
-ApiVersions 0, CreateTopics 2), so there is no version below an API's min that
the crate can
-actually encode a response for either - `unsupported_version_response` still
tries, but the
-encode attempt fails and the connection closes rather than sending a malformed
body.
+the encoder max. `kafka_protocol`'s schema floor for each of the seven
supported messages
+happens to equal `SUPPORTED_RANGES`' own min today (Produce 3, Fetch 4,
ListOffsets 1,
+Metadata 0, ApiVersions 0, CreateTopics 2, InitProducerId 0), so there is no
version below an
+API's min that the crate can actually encode a response for either -
`unsupported_version_response`
+still tries, but the encode attempt fails and the connection closes rather
than sending a
+malformed body.
**ApiVersions is the sole exception** (KIP-511): out of range still answers
with a v0 error body,
because a client probing an unknown server must be able to parse the discovery
response before
it knows the server supports flexible encoding.
@@ -48,8 +49,11 @@ it knows the server supports flexible encoding.
| 1 | Fetch | 4 | 12 | 4, 5, 6, 7, 8, 9, 10, 11, 12 | Decode request; stub
response |
| 2 | ListOffsets | 1 | 6 | 1, 2, 3, 4, 5, 6 | Decode request; stub response |
| 19 | CreateTopics | 2 | 5 | 2, 3, 4, 5 | Decode request; stub returns
`NOT_CONTROLLER` (41); `-1` partitions/RF = broker default on 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+ |
-A request is accepted when `min_version ≤ api_version ≤ max_version` for that
API key. Any other version for a listed key closes the connection (ApiVersions
excepted - see Governance model above). Any unlisted API key also closes the
connection: no api-specific response schema exists for it, so any body this
gateway could send would be misparsed by the client against the schema it
expected.
+A request is accepted when `min_version ≤ api_version ≤ max_version` for that
API key. Any other version for a listed key closes the connection (ApiVersions
excepted - see Governance model above).
+
+Any unlisted API key also closes the connection. The gateway declines to
define a response for a key it does not advertise, and a conforming client
never sends one: it reads ApiVersions first and the key's absence is what stops
the request. (`kafka-protocol`'s `broker` feature does ship response schemas
for keys this gateway leaves unlisted, so the reason is a deliberate refusal,
not an encoding limit.)
### Valid versions reference (by API key)
@@ -63,6 +67,7 @@ Use this table when configuring clients or generating wire
fixtures with `kafka-
| 3 | Metadata | 0–9 | v9 |
| 18 | ApiVersions | 0–3 | v3 |
| 19 | CreateTopics | 2–5 | v5 |
+| 22 | InitProducerId | 0–5 | v2 |
---
@@ -79,10 +84,37 @@ All API keys not listed above close the connection (see
Governance model above)
| 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)) |
| 36 | SaslAuthenticate | Implemented behind `IGGY_KAFKA_SASL_ENABLED`,
advertised only while it is on ([`AUTHENTICATION.md`](AUTHENTICATION.md)) |
-| 20+ | DeleteTopics, InitProducerId, transactions, etc. | Later issues |
+| 24, 25, 26, 28 | AddPartitionsToTxn, AddOffsetsToTxn, EndTxn,
TxnOffsetCommit | Transactions - not supported, see below |
+| 20, 21, 23, 27, 30–35, 37+ | DeleteTopics, DeleteRecords,
`OffsetForLeaderEpoch`, `WriteTxnMarkers`, `CreateAcls`/`DeleteAcls`, etc. |
Later issues |
Full reference for future phases:
[`kafka_api_keys_reference.md`](kafka_api_keys_reference.md).
+### Transactions
+
+Transactions are not supported and are not planned. There is no last stable
offset, no abort
+marker, and nothing that could make `read_committed` mean anything, so
accepting a transactional
+write would deliver an aborted transaction's records to every consumer.
+
+Three things enforce that, in the order a client meets them:
+
+1. **AddPartitionsToTxn (24), AddOffsetsToTxn (25), EndTxn (26) and
TxnOffsetCommit (28) stay out
+ of `SUPPORTED_RANGES`**, so ApiVersions never advertises them and a
conforming client never
+ sends one. This is the primary gate: the Java client's
`NodeApiVersions.latestUsableVersion`
+ throws and `NetworkClient.doSend` keeps the request off the wire;
librdkafka's four request
+ builders return `__UNSUPPORTED_FEATURE`, which is fatal there.
+2. **InitProducerId (22) with a `transactional_id`** answers
`UNSUPPORTED_VERSION` (35), so a
+ producer that got past step 1 fails before it can open a transaction.
Terminal on the Java
+ client, whose `InitProducerIdHandler` treats any unrecognised code as
fatal. Not terminal on
+ librdkafka, which retries 35 here indefinitely; librdkafka is stopped by
step 1 instead, and
+ [`IDEMPOTENCE.md`](IDEMPOTENCE.md) records why that matters when
FindCoordinator is advertised.
+3. **Produce with a non-empty `transactional_id`** answers
`UNSUPPORTED_VERSION` (35) per
+ partition, so a raw client that skipped both earlier gates still cannot
write transactional
+ records. Under `acks=0` there is no response to carry 35, so the connection
is closed
+ instead, the same way a Kafka broker answers an `acks=0` produce error.
+
+An idempotent (non-transactional) producer is unaffected: it gets a producer
id and works
+untouched, at at-least-once delivery. See [`IDEMPOTENCE.md`](IDEMPOTENCE.md).
+
---
## Architecture (three layers)
@@ -90,8 +122,8 @@ Full reference for future phases:
[`kafka_api_keys_reference.md`](kafka_api_keys
| Layer | #3421 | Description |
| ------- | ------- | ------------- |
| **1 — Wire framing** | In scope | `server.rs` — custom, zero-copy frame I/O;
`header.rs` delegates version selection to `kafka_protocol::messages::ApiKey` |
-| **2 — Request/response codecs** | Partial | Decode/encode via the
`kafka_protocol` crate (broker feature only) for 6 hot-path keys;
`bounds_guard.rs` pre-validates against unbounded allocation before handing a
frame to the crate; stub responses except Produce and ListOffsets with a bridge
|
-| **3 — Iggy bridge** | Produce and ListOffsets wired | `bridge/` module
(connection, topic mapping, provisioning, high watermark, `topic_target` +
`send_records`). 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
([#3536](https://github.com/apache/iggy/issues/3536)) |
+| **2 — Request/response codecs** | Partial | Decode/encode via the
`kafka_protocol` crate (broker feature only) for 7 hot-path keys;
`bounds_guard.rs` pre-validates against unbounded allocation before handing a
frame to the crate; stub responses except InitProducerId, 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 ([# [...]
---
@@ -253,7 +285,7 @@ below it are still open for the issues that build on top of
it.
This TODO originally proposed a selective, feature-gated adoption
(`kafka-protocol-cold`)
alongside the hand-rolled `requests.rs`/`responses.rs` codecs, keeping custom
code for the
Produce/Fetch hot paths. That hybrid approach was not taken: `kafka_protocol`
(broker feature
-only) now decodes/encodes all six supported message types wholesale, and the
hand-rolled
+only) now decodes/encodes all seven supported message types wholesale, and the
hand-rolled
`codec.rs`/`requests.rs` were deleted. RecordBatch bytes stay opaque
(`Option<Bytes>`, never
decoded) on the Produce/Fetch hot paths, preserving the one property this TODO
was protecting.
`bounds_guard.rs` covers the DoS-bound gap the crate itself leaves open (see
Governance model
@@ -275,6 +307,10 @@ Offset persistence design
([#3540](https://github.com/apache/iggy/issues/3540)):
InitProducerId and idempotent producers
([#3545](https://github.com/apache/iggy/issues/3545)):
[`IDEMPOTENCE.md`](IDEMPOTENCE.md).
+- [x] InitProducerId (22) allocates a producer id so a stock idempotent
producer starts; a
+ transactional request is refused. The producer-id-keyed connection pool
that would make
+ retries deduplicated is deferred - delivery stays at-least-once
+
Authentication design ([#3549](https://github.com/apache/iggy/issues/3549)):
[`AUTHENTICATION.md`](AUTHENTICATION.md).
diff --git a/gateways/kafka/docs/TEST_SUITE.md
b/gateways/kafka/docs/TEST_SUITE.md
index da897e3d0..cd1a849d4 100644
--- a/gateways/kafka/docs/TEST_SUITE.md
+++ b/gateways/kafka/docs/TEST_SUITE.md
@@ -59,6 +59,7 @@ file under `tests/` anymore.
| [`golden_wire_fixtures_tests.rs`](../tests/golden_wire_fixtures_tests.rs) |
Byte-exact golden responses (ApiVersions v1, Metadata v0) | No |
| [`fixtures_canary_tests.rs`](../tests/fixtures_canary_tests.rs) | Fails
loudly if `KAFKA_FIXTURES_REQUIRED=1` and no `.bin` fixtures exist, so a broken
generation step can't leave the fixture-backed suites green-but-empty | Canary
only |
| [`version_firewall_tests.rs`](../tests/version_firewall_tests.rs) | Version
boundary matrix, unsupported keys, corrupt bodies | Partial |
+| [`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 |
| [`server_e2e_tests.rs`](../tests/server_e2e_tests.rs) | Full `KafkaGateway`
TCP round-trips | Partial |
diff --git a/gateways/kafka/docs/kafka_api_keys_reference.md
b/gateways/kafka/docs/kafka_api_keys_reference.md
index 99c42c1e2..19c004215 100644
--- a/gateways/kafka/docs/kafka_api_keys_reference.md
+++ b/gateways/kafka/docs/kafka_api_keys_reference.md
@@ -15,6 +15,7 @@
| 🟠 Required Stub | Client state-machine API — must return a well-formed
response or clients will stall/crash |
| 🟡 Optional Stub | Admin/observability — can safely return
`UNSUPPORTED_VERSION` or `NOT_CONTROLLER` |
| ❌ Reject | Internal broker / KRaft only — return `INVALID_REQUEST` with a
well-formed frame; **do not close the connection** |
+| ❌ Unadvertised | Deliberately absent from `ApiVersions`, so a conforming
client never sends one; an arriving request closes the connection |
> This table no longer carries a per-key header-framing status column.
> `src/protocol/header.rs`
> has no per-key table of its own to be behind or caught up on: it delegates
> entirely to
@@ -122,13 +123,19 @@ Key new minimums:
| Key | API Name | Min (4.0) | Max (4.0) | Flexible From | Gateway Action |
| :---: | ---------- | :---------: | :---------: | :-------------: |
:--------------: |
-| 22 | **InitProducerId** | 2 | 5 | v2 | 🟡 Optional Stub |
+| 22 | **InitProducerId** | 2 | 5 | v2 | 🟠 Required Stub |
| 23 | **OffsetForLeaderEpoch** | 1 | 5 | v4 | 🟡 Optional Stub |
-| 24 | **AddPartitionsToTxn** | 1 | 5 | v3 | 🟡 Optional Stub |
-| 25 | **AddOffsetsToTxn** | 1 | 4 | v3 | 🟡 Optional Stub |
-| 26 | **EndTxn** | 1 | 4 | v3 | 🟡 Optional Stub |
+| 24 | **AddPartitionsToTxn** | 1 | 5 | v3 | ❌ Unadvertised |
+| 25 | **AddOffsetsToTxn** | 1 | 4 | v3 | ❌ Unadvertised |
+| 26 | **EndTxn** | 1 | 4 | v3 | ❌ Unadvertised |
| 27 | **WriteTxnMarkers** | 0 | 1 | v1 | 🟡 Optional Stub |
-| 28 | **TxnOffsetCommit** | 2 | 5 | v3 | 🟡 Optional Stub |
+| 28 | **TxnOffsetCommit** | 2 | 5 | v3 | ❌ Unadvertised |
+
+> InitProducerId is implemented, not stubbed: it allocates a producer id so a
stock idempotent
+> producer starts, and answers `UNSUPPORTED_VERSION` (35) only when the
request carries a
+> `transactional_id`. The four keys marked Unadvertised are never listed in
`ApiVersions`, which
+> is what stops a conforming client from opening a transaction at all. See
+> [`IDEMPOTENCE.md`](IDEMPOTENCE.md) and `SCOPE.md`'s Transactions section.
---
@@ -257,8 +264,9 @@ Key new minimums:
| Category | Count | Notes |
| ---------- | :-----: | ------- |
| 🔴 Bridge (data path) | 7 | Produce, Fetch, Metadata, SaslHandshake,
ApiVersions, SaslAuthenticate, ShareFetch |
-| 🟠 Required Stub (client state machine) | 12 | ListOffsets, consumer group
(8-14), CreateTopics, ConsumerGroupHeartbeat (68), ShareGroupHeartbeat (77),
ShareAcknowledge (80) |
-| 🟡 Optional Stub (admin/observability) | 47 | Can return
`UNSUPPORTED_VERSION` or `NOT_CONTROLLER` safely |
+| 🟠 Required Stub (client state machine) | 13 | ListOffsets, consumer group
(8-14), CreateTopics, InitProducerId (22), ConsumerGroupHeartbeat (68),
ShareGroupHeartbeat (77), ShareAcknowledge (80) |
+| 🟡 Optional Stub (admin/observability) | 42 | Can return
`UNSUPPORTED_VERSION` or `NOT_CONTROLLER` safely |
+| ❌ Unadvertised (transactions) | 4 | AddPartitionsToTxn (24), AddOffsetsToTxn
(25), EndTxn (26), TxnOffsetCommit (28). Absent from ApiVersions, so a
conforming client never sends one |
| ❌ Reject (broker/KRaft internal) | 22 | Return `INVALID_REQUEST` with valid
frame — never close the TCP connection |
| **Total API Keys in this document** | **88** | Key IDs 0-88 with a gap at 73
|
@@ -277,12 +285,12 @@ Key new minimums:
| ApiVersions | v0-v3 | v4 | 1 version behind |
| CreateTopics | v2-v5 | v7 | 2 versions behind |
-### Missing from `SUPPORTED_RANGES` (82 of the 88 API keys in this document)
+### Missing from `SUPPORTED_RANGES` (81 of the 88 API keys in this document)
Every key not in `SUPPORTED_RANGES` closes the connection - the same policy
applied to every
-other unlisted key, not a special case for these. No api-specific response
schema exists for an
-unlisted key, so any body the gateway could send would be misparsed by the
client against the
-schema it expected. This includes:
+other unlisted key, not a special case for these. The gateway declines to
define a response for a
+key it does not advertise, and a conforming client never sends one, so no
response shape has to
+be agreed. This includes:
- **Client bootstrap blockers**: OffsetCommit (8), OffsetFetch (9),
FindCoordinator (10)
- **Classic consumer group protocol**: JoinGroup (11), Heartbeat (12),
LeaveGroup (13), SyncGroup (14)
diff --git a/gateways/kafka/scripts/ci-wire-fixtures.sh
b/gateways/kafka/scripts/ci-wire-fixtures.sh
index df8c99f6b..920923e48 100755
--- a/gateways/kafka/scripts/ci-wire-fixtures.sh
+++ b/gateways/kafka/scripts/ci-wire-fixtures.sh
@@ -24,7 +24,7 @@ set -euo pipefail
FIXTURES_DIR="gateways/kafka/tools/kafka-tool/kafka_messages"
# API keys requested by api_handler_tests, version_firewall_tests, and
server_e2e_tests.
-FIXTURE_API_KEYS=(0 1 2 19)
+FIXTURE_API_KEYS=(0 1 2 19 22)
usage() {
echo "Usage: $0 {generate|cleanup}" >&2
diff --git a/gateways/kafka/src/main.rs b/gateways/kafka/src/main.rs
index 3b92b4875..11233607e 100644
--- a/gateways/kafka/src/main.rs
+++ b/gateways/kafka/src/main.rs
@@ -100,6 +100,7 @@ const KNOWN_KAFKA_ENV_VARS: &[&str] = &[
"IGGY_KAFKA_WRITE_TIMEOUT_SECS",
"IGGY_KAFKA_SHUTDOWN_DRAIN_TIMEOUT_SECS",
"IGGY_KAFKA_BRIDGE_ENABLED",
+ "IGGY_KAFKA_INSTANCE_ID",
"IGGY_KAFKA_SASL_ENABLED",
"IGGY_KAFKA_PRE_AUTH_TIMEOUT_SECS",
"IGGY_KAFKA_MAX_CONCURRENT_AUTHENTICATIONS",
@@ -244,6 +245,19 @@ fn load_config() -> Result<GatewayConfig, String> {
.map_err(|e| format!("invalid
IGGY_KAFKA_SHUTDOWN_DRAIN_TIMEOUT_SECS `{raw}`: {e}"))?;
config.shutdown_drain_timeout = Duration::from_secs(secs);
}
+ // Not `parse_positive`: 0 is the default, and the right value for a
single gateway.
+ if let Some(raw) = env_var("IGGY_KAFKA_INSTANCE_ID") {
+ config.instance_id = raw
+ .parse()
+ .map_err(|e| format!("invalid IGGY_KAFKA_INSTANCE_ID `{raw}`:
{e}"))?;
+ } else {
+ // A second gateway left on the default hands out the same producer
ids as the first,
+ // and nothing in the cluster can detect it, so an unset value is
surfaced at startup.
+ warn!(
+ "IGGY_KAFKA_INSTANCE_ID is not set, defaulting to 0: every gateway
fronting the same \
+ Iggy cluster needs its own value, or they hand out colliding
producer ids"
+ );
+ }
reject_iggy_tls_without_sasl(config.sasl_enabled)?;
Ok(config)
@@ -323,11 +337,14 @@ async fn shutdown_signal() {
mod tests {
use serial_test::serial;
- use super::{parse_positive, reject_iggy_tls_without_sasl,
reject_unknown_kafka_env_vars};
+ use super::{
+ load_config, parse_positive, reject_iggy_tls_without_sasl,
reject_unknown_kafka_env_vars,
+ };
/// Sequential (not two separate `#[test]` fns), and `#[serial]` (unkeyed
- this binary's
- /// default group). The `#[serial]` tests in this module are the only ones
compiled into *this*
- /// binary (`main.rs` -> the `iggy-gateway-kafka` bin's own test harness)
- `auth`'s,
+ /// default group), shared with the instance-id and TLS tests below since
all of them touch the
+ /// process environment. The `#[serial]` tests in this module are the only
ones compiled into
+ /// *this* binary (`main.rs` -> the `iggy-gateway-kafka` bin's own test
harness) - `auth`'s,
/// `bridge::config`'s and `server`'s env-touching tests compile into the
separate lib test
/// binary, and `serial_test`'s mutex is process-local, so it does not
(and does not need to)
/// coordinate with any of those; `server.rs`'s own `#[serial]` test makes
the mirror-image
@@ -382,6 +399,67 @@ mod tests {
);
}
+ /// `#[serial]` and `# Safety` as on
+ /// `reject_unknown_kafka_env_vars_flags_typo_but_accepts_known_keys`
above.
+ ///
+ /// Covers both halves of adding this var: it has to be in
`KNOWN_KAFKA_ENV_VARS` (or setting
+ /// it refuses to start the gateway) and it has to keep `0`, which
`parse_positive` rejects.
+ #[test]
+ #[serial]
+ fn
given_an_instance_id_env_var_when_loading_config_should_accept_and_parse_it() {
+ unsafe {
+ std::env::set_var("IGGY_KAFKA_INSTANCE_ID", "7");
+ }
+ let seven = load_config();
+ unsafe {
+ std::env::set_var("IGGY_KAFKA_INSTANCE_ID", "0");
+ }
+ let zero = load_config();
+ unsafe {
+ std::env::set_var("IGGY_KAFKA_INSTANCE_ID", "65536");
+ }
+ let overflow = load_config();
+ unsafe {
+ std::env::remove_var("IGGY_KAFKA_INSTANCE_ID");
+ }
+
+ assert_eq!(seven.expect("instance id 7 must load").instance_id, 7);
+ assert_eq!(zero.expect("instance id 0 must load").instance_id, 0);
+ assert!(
+ overflow.is_err(),
+ "an instance id above u16::MAX must be rejected, not truncated"
+ );
+ }
+
+ /// A bad instance id must fail startup loudly. Silently defaulting to 0
would let two
+ /// gateways mint colliding producer ids, which Kafka requires to be
unique cluster-wide,
+ /// and the message has to name the variable and the offending value or an
operator cannot
+ /// act on it.
+ #[test]
+ #[serial]
+ fn
given_an_unparsable_instance_id_when_loading_config_should_reject_and_name_it()
{
+ for raw in ["abc", "-1", "", " 7", "7.0"] {
+ unsafe {
+ std::env::set_var("IGGY_KAFKA_INSTANCE_ID", raw);
+ }
+ let loaded = load_config();
+ unsafe {
+ std::env::remove_var("IGGY_KAFKA_INSTANCE_ID");
+ }
+ let error = loaded.err().unwrap_or_else(|| {
+ panic!("instance id `{raw}` must be rejected, not silently
defaulted to 0")
+ });
+ assert!(
+ error.contains("IGGY_KAFKA_INSTANCE_ID"),
+ "`{raw}` rejection must name the variable, got: {error}"
+ );
+ assert!(
+ error.contains(raw),
+ "`{raw}` rejection must quote the offending value, got:
{error}"
+ );
+ }
+ }
+
/// `#[serial]` for the same reason as the test above: it mutates
process-wide env state.
#[test]
#[serial]
diff --git a/gateways/kafka/src/protocol/api.rs
b/gateways/kafka/src/protocol/api.rs
index 0d9217e81..a63e9a1a4 100644
--- a/gateways/kafka/src/protocol/api.rs
+++ b/gateways/kafka/src/protocol/api.rs
@@ -33,9 +33,10 @@ use crate::protocol::acl::{
use crate::protocol::bounds_guard::{
validate_describe_acls_shape, validate_sasl_authenticate_shape,
validate_sasl_handshake_shape,
};
+use crate::protocol::handlers::init_producer_id::ProducerIdAllocator;
use crate::protocol::handlers::{
- api_versions, create_topics, decode_guarded, dispatch, fetch,
list_offsets, metadata, produce,
- respond_or_close,
+ api_versions, create_topics, decode_guarded, dispatch, fetch,
init_producer_id, list_offsets,
+ metadata, produce, respond_or_close,
};
use crate::protocol::sasl::{
SaslMechanism, encode_sasl_authenticate_response,
encode_sasl_handshake_response,
@@ -48,6 +49,7 @@ pub const API_KEY_METADATA: i16 = 3;
pub const API_KEY_SASL_HANDSHAKE: i16 = 17;
pub const API_KEY_API_VERSIONS: i16 = 18;
pub const API_KEY_CREATE_TOPICS: i16 = 19;
+pub const API_KEY_INIT_PRODUCER_ID: i16 = 22;
pub const API_KEY_DESCRIBE_ACLS: i16 = 29;
pub const API_KEY_SASL_AUTHENTICATE: i16 = 36;
@@ -209,6 +211,11 @@ pub struct ApiVersionRange {
pub max_version: i16,
}
+/// The version firewall, and the exact set `ApiVersions` advertises.
+///
+/// Absence is load-bearing for the transaction keys (24, 25, 26, 28): a
conforming client that
+/// does not see a key here never sends it, which is the whole enforcement of
"transactions are
+/// unsupported". See `docs/SCOPE.md`.
static SUPPORTED_RANGES: &[ApiVersionRange] = &[
produce::RANGE,
fetch::RANGE,
@@ -216,6 +223,7 @@ static SUPPORTED_RANGES: &[ApiVersionRange] = &[
metadata::RANGE,
api_versions::RANGE,
create_topics::RANGE,
+ init_producer_id::RANGE,
];
#[must_use]
@@ -238,6 +246,9 @@ pub struct GatewayState {
/// Whether `SaslHandshake` and `SaslAuthenticate` are advertised and
routed. Kept on the
/// shared state so `ApiVersions` can answer without a widened handler
signature.
pub sasl_enabled: bool,
+ /// Shared across every connection this gateway serves: a producer id has
to be unique for
+ /// the process, not for the connection that asked for it.
+ pub producer_ids: ProducerIdAllocator,
/// Produce requests that decode and send at once. Caps their memory.
pub(crate) produce_slots: Semaphore,
}
@@ -248,25 +259,27 @@ const PRODUCE_SLOTS: usize = 4;
impl GatewayState {
#[must_use]
- pub const fn new(
+ pub fn new(
broker: BrokerAdvertise,
bridge: Option<Arc<IggyBridge>>,
max_frame_size: usize,
sasl_enabled: bool,
+ instance_id: u16,
) -> Self {
Self {
broker,
bridge,
max_frame_size,
sasl_enabled,
+ producer_ids: ProducerIdAllocator::new(instance_id),
produce_slots: Semaphore::const_new(PRODUCE_SLOTS),
}
}
/// State with no bridge, so every handler takes its stub path.
#[must_use]
- pub const fn stub(broker: BrokerAdvertise, max_frame_size: usize) -> Self {
- Self::new(broker, None, max_frame_size, false)
+ pub fn stub(broker: BrokerAdvertise, max_frame_size: usize) -> Self {
+ Self::new(broker, None, max_frame_size, false, 0)
}
}
diff --git a/gateways/kafka/src/protocol/bounds_guard.rs
b/gateways/kafka/src/protocol/bounds_guard.rs
index cd80d5321..b75d34a59 100644
--- a/gateways/kafka/src/protocol/bounds_guard.rs
+++ b/gateways/kafka/src/protocol/bounds_guard.rs
@@ -708,6 +708,36 @@ pub fn validate_metadata_shape(version: i16, body: &Bytes,
max_frame_size: usize
Ok(())
}
+/// Mirrors the field order `InitProducerIdRequest::decode` walks.
+///
+/// No response-size guard needed: the response is four fixed-width fields and
echoes nothing
+/// from the request, so `usize::MAX` disables that check rather than plumbing
`max_frame_size`
+/// through for no effect (same as [`validate_api_versions_shape`]).
+///
+/// # Errors
+///
+/// Returns an error when the declared `transactional_id` 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_init_producer_id_shape(version: i16, body: &Bytes) ->
Result<()> {
+ let mut c = ShapeCursor::new(body.clone(), usize::MAX);
+ let flexible = version >= 2;
+
+ if flexible {
+ c.compact_string(true)?;
+ } else {
+ c.legacy_string(true)?;
+ }
+ let _transaction_timeout_ms = c.read_i32()?;
+ if version >= 3 {
+ let _producer_id = c.read_i64()?;
+ let _producer_epoch = c.read_i16()?;
+ }
+ if flexible {
+ c.tagged_fields()?;
+ }
+ Ok(())
+}
+
/// Mirrors the field order `ApiVersionsRequest::decode` walks. v0-2 have an
empty body (no
/// length-prefixed fields to bound), so this is a no-op below v3.
///
@@ -947,6 +977,36 @@ mod tests {
assert!(validate_metadata_shape(0, &body,
TEST_MAX_FRAME_SIZE).is_err());
}
+ /// Every sibling guard carries a rejection POC; without one,
short-circuiting this guard to
+ /// `Ok(())` leaves the whole suite green, so nothing proved it rejected a
hostile frame.
+ #[test]
+ fn init_producer_id_v5_huge_compact_string_rejected() {
+ // Compact string length varint far past the frame: nothing follows it
to read.
+ let body = Bytes::from_static(&[0xFF, 0xFF, 0xFF, 0xFF, 0x0F]);
+ assert!(validate_init_producer_id_shape(5, &body).is_err());
+ }
+
+ #[test]
+ fn init_producer_id_v0_truncated_legacy_string_rejected() {
+ // Declares 32767 bytes of transactional id, supplies none.
+ let body = Bytes::from_static(&[0x7F, 0xFF]);
+ assert!(validate_init_producer_id_shape(0, &body).is_err());
+ }
+
+ #[test]
+ fn init_producer_id_v5_null_transactional_id_accepted() {
+ // Null compact string, transaction_timeout_ms, then the v3+ producer
id/epoch pair
+ // (both -1, "no producer id"), then tagged fields.
+ let body = Bytes::from_static(&[
+ 0x00, // transactional_id: null compact string
+ 0x00, 0x00, 0x75, 0x30, // transaction_timeout_ms: 30000
+ 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, // producer_id: -1
+ 0xFF, 0xFF, // producer_epoch: -1
+ 0x00, // tagged fields
+ ]);
+ assert!(validate_init_producer_id_shape(5, &body).is_ok());
+ }
+
#[test]
fn metadata_v0_null_array_all_topics_accepted() {
let body = Bytes::from_static(&[0xFF, 0xFF, 0xFF, 0xFF]); // -1: all
topics
diff --git a/gateways/kafka/src/protocol/handlers/init_producer_id.rs
b/gateways/kafka/src/protocol/handlers/init_producer_id.rs
new file mode 100644
index 000000000..c9718f197
--- /dev/null
+++ b/gateways/kafka/src/protocol/handlers/init_producer_id.rs
@@ -0,0 +1,207 @@
+// 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.
+
+//! `InitProducerId` (API key 22).
+//!
+//! A Java producer sets `enable.idempotence=true` without being asked
(KIP-679, default since
+//! Kafka 3.0) and sends this before its first record, so answering it is what
lets a stock
+//! producer start against this gateway at all. The id is handed out and then
ignored: delivery
+//! stays at-least-once, and no retry is deduplicated. See
`docs/IDEMPOTENCE.md`.
+
+use std::sync::atomic::{AtomicU64, Ordering};
+use std::time::{SystemTime, UNIX_EPOCH};
+
+use bytes::Bytes;
+use kafka_protocol::messages::{InitProducerIdRequest, InitProducerIdResponse,
ProducerId};
+
+use crate::error::Result;
+use crate::protocol::api::{
+ API_KEY_INIT_PRODUCER_ID, ApiVersionRange, ERROR_NONE,
ERROR_UNKNOWN_SERVER_ERROR,
+ ERROR_UNSUPPORTED_VERSION, GatewayState, HandleOutcome,
+};
+use crate::protocol::bounds_guard::validate_init_producer_id_shape;
+use crate::protocol::handlers::{
+ decode_guarded, encode_message, handle_versioned_request, is_transactional,
+};
+
+pub const RANGE: ApiVersionRange = ApiVersionRange {
+ api_key: API_KEY_INIT_PRODUCER_ID,
+ min_version: 0,
+ max_version: 5,
+};
+
+/// Width of the per-instance counter. The remaining 16 bits of the
non-negative range carry the
+/// instance number, and bit 63 stays clear because `producer_id` is an `i64`
whose `-1` means
+/// "no producer id".
+const COUNTER_BITS: u32 = 47;
+const MAX_COUNTER: u64 = (1 << COUNTER_BITS) - 1;
+
+/// The epoch every allocated id carries. Epochs only advance when a producer
is fenced, which
+/// needs the transactional state this gateway does not keep.
+const PRODUCER_EPOCH: i16 = 0;
+
+/// Hands out producer ids that are unique across gateway instances sharing
one Iggy cluster.
+///
+/// `instance_id` is configured (`IGGY_KAFKA_INSTANCE_ID`), not drawn at
startup: a random 16-bit
+/// value collides at even odds around 300 instances.
+///
+/// The counter starts at the wall clock in milliseconds, not at 0. Kafka keys
a producer on
+/// `(producer_id, producer_epoch)` and the epoch is always 0 here, so a
counter restarting at 0
+/// would hand a restarted gateway's producers the pairs its previous run gave
out. Seeding from
+/// the clock keeps every new id above the old ones unless the previous run
averaged more than one
+/// allocation per millisecond of its uptime, the clock stepped back across
the restart, or the
+/// clock reads before the Unix epoch and seeds 0. Nothing is persisted, and
nothing reads the id
+/// yet. Before Produce (#3535) keys anything on it, the allocator has to
persist a high-water mark
+/// or bump the epoch on every start.
+#[derive(Debug)]
+pub struct ProducerIdAllocator {
+ instance_id: u16,
+ next_counter: AtomicU64,
+}
+
+impl ProducerIdAllocator {
+ #[must_use]
+ pub fn new(instance_id: u16) -> Self {
+ Self {
+ instance_id,
+ next_counter: AtomicU64::new(clock_counter()),
+ }
+ }
+
+ /// The next id for this instance, or `None` once its counter space is
spent.
+ fn allocate(&self) -> Option<i64> {
+ let counter = self.next_counter.fetch_add(1, Ordering::Relaxed);
+ if counter > MAX_COUNTER {
+ return None;
+ }
+ i64::try_from((u64::from(self.instance_id) << COUNTER_BITS) |
counter).ok()
+ }
+}
+
+#[expect(
+ clippy::unused_async,
+ reason = "the shared handler signature, kept until a handler awaits the
bridge"
+)]
+pub async fn handle(state: &GatewayState, api_version: i16, body: Bytes) ->
HandleOutcome {
+ handle_versioned_request(
+ API_KEY_INIT_PRODUCER_ID,
+ api_version,
+ body,
+ |v, b| decode_guarded::<InitProducerIdRequest>(v, b,
validate_init_producer_id_shape),
+ |v, req| encode_response(v, req, &state.producer_ids),
+ encode_error_response,
+ "InitProducerId",
+ )
+}
+
+/// `InitProducerId` response carrying `error_code` and no usable producer id.
+///
+/// # 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_inner(version, error_code, -1)
+}
+
+/// Allocate a producer id, or refuse a transactional request with
`UNSUPPORTED_VERSION` (35).
+///
+/// 35 is the code the Java producer's `InitProducerIdHandler.handleResponse`
cannot recover
+/// from; a retriable code would leave `initTransactions()` looping forever
instead of failing.
+///
+/// # Errors
+///
+/// Returns an error when `kafka_protocol` cannot encode the response at
`version`.
+pub fn encode_response(
+ version: i16,
+ req: &InitProducerIdRequest,
+ allocator: &ProducerIdAllocator,
+) -> Result<Bytes> {
+ if is_transactional(req.transactional_id.as_ref()) {
+ return encode_error_response(version, ERROR_UNSUPPORTED_VERSION);
+ }
+ let Some(producer_id) = allocator.allocate() else {
+ tracing::error!(
+ "producer id space exhausted; restart the gateway with a free
IGGY_KAFKA_INSTANCE_ID"
+ );
+ return encode_error_response(version, ERROR_UNKNOWN_SERVER_ERROR);
+ };
+ encode_inner(version, ERROR_NONE, producer_id)
+}
+
+/// Milliseconds since the Unix epoch, capped at [`MAX_COUNTER`] so a
far-future clock spends the
+/// counter space instead of bleeding into the instance bits. A clock before
the epoch reads as 0.
+fn clock_counter() -> u64 {
+ let millis = SystemTime::now()
+ .duration_since(UNIX_EPOCH)
+ .map_or(0, |elapsed| elapsed.as_millis());
+ u64::try_from(millis).map_or(MAX_COUNTER, |millis| millis.min(MAX_COUNTER))
+}
+
+fn encode_inner(version: i16, error_code: i16, producer_id: i64) ->
Result<Bytes> {
+ let resp = InitProducerIdResponse::default()
+ .with_error_code(error_code)
+ .with_producer_id(ProducerId(producer_id))
+ .with_producer_epoch(PRODUCER_EPOCH);
+ encode_message(&resp, version, 32)
+}
+
+#[cfg(test)]
+mod tests {
+ use std::sync::atomic::Ordering;
+ use std::thread;
+ use std::time::Duration;
+
+ use super::{COUNTER_BITS, MAX_COUNTER, ProducerIdAllocator};
+
+ #[test]
+ fn
given_a_fresh_allocator_when_allocating_twice_should_return_distinct_ids() {
+ let allocator = ProducerIdAllocator::new(0);
+ let first = allocator.allocate().expect("first id");
+ let second = allocator.allocate().expect("second id");
+ assert!(first >= 0 && second >= 0);
+ assert_ne!(first, second);
+ }
+
+ #[test]
+ fn
given_an_instance_id_when_allocating_should_place_it_above_the_counter() {
+ let allocator = ProducerIdAllocator::new(0xBEEF);
+ let id = allocator.allocate().expect("id");
+ assert_eq!(id >> COUNTER_BITS, 0xBEEF);
+ }
+
+ #[test]
+ fn
given_a_restarted_allocator_when_allocating_should_start_above_the_previous_run()
{
+ let previous_run = ProducerIdAllocator::new(0);
+ let last_of_previous_run = previous_run.allocate().expect("previous
run id");
+ thread::sleep(Duration::from_millis(2));
+ let restarted = ProducerIdAllocator::new(0);
+ let first_after_restart = restarted.allocate().expect("restarted id");
+ assert!(
+ first_after_restart > last_of_previous_run,
+ "a restart must not replay a (producer_id, epoch 0) pair a live
producer still holds"
+ );
+ }
+
+ #[test]
+ fn
given_a_spent_counter_when_allocating_should_refuse_instead_of_bleeding_into_the_instance()
{
+ let allocator = ProducerIdAllocator::new(1);
+ allocator.next_counter.store(MAX_COUNTER, Ordering::Relaxed);
+ let last = allocator.allocate().expect("last id of the counter space");
+ assert_eq!(last >> COUNTER_BITS, 1);
+ assert_eq!(allocator.allocate(), None);
+ }
+}
diff --git a/gateways/kafka/src/protocol/handlers/mod.rs
b/gateways/kafka/src/protocol/handlers/mod.rs
index a03fa9eac..cba9c7b2e 100644
--- a/gateways/kafka/src/protocol/handlers/mod.rs
+++ b/gateways/kafka/src/protocol/handlers/mod.rs
@@ -27,18 +27,21 @@
pub mod api_versions;
pub mod create_topics;
pub mod fetch;
+pub mod init_producer_id;
pub mod list_offsets;
pub mod metadata;
pub mod produce;
use bytes::{Buf, Bytes, BytesMut};
+use kafka_protocol::messages::TransactionalId;
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_LIST_OFFSETS,
- API_KEY_METADATA, API_KEY_PRODUCE, ERROR_INVALID_REQUEST,
ERROR_UNSUPPORTED_VERSION,
- GatewayState, HandleOutcome, is_supported_version, supported_max_version,
+ API_KEY_API_VERSIONS, API_KEY_CREATE_TOPICS, API_KEY_FETCH,
API_KEY_INIT_PRODUCER_ID,
+ API_KEY_LIST_OFFSETS, API_KEY_METADATA, API_KEY_PRODUCE,
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.
@@ -58,10 +61,20 @@ pub async fn dispatch(
API_KEY_METADATA => metadata::handle(state, api_version, body).await,
API_KEY_API_VERSIONS => api_versions::handle(state, api_version,
body).await,
API_KEY_CREATE_TOPICS => create_topics::handle(state, api_version,
body).await,
+ API_KEY_INIT_PRODUCER_ID => init_producer_id::handle(state,
api_version, body).await,
_ => HandleOutcome::Close,
}
}
+/// Whether a request carries a transactional id, which this gateway never
serves.
+///
+/// An empty id reads as absent: `kafka_protocol` decodes a null wire string
to `None` but its
+/// own `Default` uses `Some("")`, and a producer that is idempotent-only has
no transaction to
+/// name either way.
+pub(crate) fn is_transactional(transactional_id: Option<&TransactionalId>) ->
bool {
+ transactional_id.is_some_and(|id| !id.is_empty())
+}
+
/// Encode a `kafka_protocol` message, mapping its `anyhow::Error` (the crate
has no stable
/// decode/encode error taxonomy) to a variant callers can log or fold into
/// [`HandleOutcome::Close`].
diff --git a/gateways/kafka/src/protocol/handlers/produce.rs
b/gateways/kafka/src/protocol/handlers/produce.rs
index 9e6fe418e..e64e22ac0 100644
--- a/gateways/kafka/src/protocol/handlers/produce.rs
+++ b/gateways/kafka/src/protocol/handlers/produce.rs
@@ -38,7 +38,8 @@ use crate::protocol::api::{
};
use crate::protocol::bounds_guard::validate_produce_shape;
use crate::protocol::handlers::{
- decode_guarded, encode_message, respond_or_close,
unsupported_version_response,
+ decode_guarded, encode_message, is_transactional, respond_or_close,
+ unsupported_version_response,
};
use crate::records::{
Allowance, DecompressionBudget, RecordCodecError, TimestampWindow, Zstd,
decode_batch,
@@ -136,6 +137,22 @@ pub async fn handle(state: &GatewayState, api_version:
i16, body: Bytes) -> Hand
}
};
+ // Ahead of the bridge check, so neither the stub nor the write path
answers a transactional
+ // request as ordinary records: nothing here tracks a last stable offset
or writes an abort
+ // marker, so an aborted transaction's records would reach every consumer.
35 is fatal for the
+ // producer. acks=0 has no response to carry it, so the refusal is a
close, which is what a
+ // Kafka broker does on an acks=0 produce error.
+ if is_transactional(request.transactional_id.as_ref()) {
+ return if request.acks == ACKS_NONE {
+ HandleOutcome::Close
+ } else {
+ respond_or_close(
+ encode_uniform_response(api_version, &request,
ERROR_UNSUPPORTED_VERSION),
+ "Produce",
+ )
+ };
+ }
+
let Some(bridge) = state.bridge.as_deref() else {
return stub_outcome(api_version, &request);
};
diff --git a/gateways/kafka/src/server.rs b/gateways/kafka/src/server.rs
index 2e6bb6686..1c7c1d697 100644
--- a/gateways/kafka/src/server.rs
+++ b/gateways/kafka/src/server.rs
@@ -118,6 +118,11 @@ pub struct GatewayConfig {
/// hold shutdown open past typical orchestrator grace periods (e.g.
Kubernetes' default
/// 30s `terminationGracePeriodSeconds`).
pub shutdown_drain_timeout: Duration,
+ /// This gateway's number among the gateways fronting one Iggy cluster
+ /// (`IGGY_KAFKA_INSTANCE_ID`). It is the high half of every producer id
handed out by
+ /// `InitProducerId`, which Kafka requires to be cluster-unique; two
gateways left on the
+ /// same number hand out the same ids.
+ pub instance_id: u16,
/// Require SASL authentication before serving any other API.
///
/// Off by default, and switching it on is a breaking change for every
client already talking
@@ -166,6 +171,7 @@ impl Default for GatewayConfig {
read_timeout: Duration::from_secs(15),
write_timeout: Duration::from_secs(10),
shutdown_drain_timeout: Duration::from_secs(25),
+ instance_id: 0,
sasl_enabled: false,
pre_auth_timeout: Duration::from_secs(15),
max_concurrent_authentications: 4,
@@ -319,36 +325,22 @@ impl KafkaGateway {
listener: TcpListener,
mut shutdown: broadcast::Receiver<()>,
) -> Result<()> {
- if !self.config.sasl_enabled && self.authenticator.is_some() {
- // The mirror of the guard below, and the quieter mistake: a
verifier attached while the
- // flag is off means every connection is served unauthenticated,
with nothing in the log
- // to say so. Refusing to start is the only way that failure is
visible.
- return Err(KafkaProtocolError::InvalidConfig(
- "an authenticator is configured but SASL is disabled; every
connection would be \
- served unauthenticated. Set IGGY_KAFKA_SASL_ENABLED=true, or
remove the \
- authenticator"
- .into(),
- ));
- }
- if self.config.sasl_enabled && self.authenticator.is_none() {
- return Err(KafkaProtocolError::InvalidConfig(
- "SASL is enabled but no authenticator is configured; every
client would be \
- rejected. Set IGGY_KAFKA_IGGY_ADDR to the Iggy server that
credentials are \
- verified against, or unset IGGY_KAFKA_SASL_ENABLED"
- .into(),
- ));
- }
+ self.check_sasl_wiring()?;
let local_addr = listener.local_addr()?;
let broker = BrokerAdvertise::from_server_config(&self.config,
local_addr)?;
+ // instance_id is logged because it is the only way to tell from a
running process which
+ // half of the producer-id space this gateway owns. Two gateways left
on the default
+ // collide silently, and a config file cannot be diffed against a live
deployment.
info!(
- "kafka listener bound on {} (advertised as {}:{})",
- local_addr, broker.host, broker.port
+ "kafka listener bound on {} (advertised as {}:{}, instance id {})",
+ local_addr, broker.host, broker.port, self.config.instance_id
);
let state = Arc::new(GatewayState::new(
broker,
self.bridge.clone(),
self.config.max_frame_size,
self.config.sasl_enabled,
+ self.config.instance_id,
));
let shared_auth = Arc::new(SharedAuth::new(
@@ -436,6 +428,30 @@ impl KafkaGateway {
}
Ok(())
}
+
+ /// Refuses to start when the SASL flag and the configured authenticator
disagree.
+ fn check_sasl_wiring(&self) -> Result<()> {
+ if !self.config.sasl_enabled && self.authenticator.is_some() {
+ // The mirror of the guard below, and the quieter mistake: a
verifier attached while the
+ // flag is off means every connection is served unauthenticated,
with nothing in the log
+ // to say so. Refusing to start is the only way that failure is
visible.
+ return Err(KafkaProtocolError::InvalidConfig(
+ "an authenticator is configured but SASL is disabled; every
connection would be \
+ served unauthenticated. Set IGGY_KAFKA_SASL_ENABLED=true, or
remove the \
+ authenticator"
+ .into(),
+ ));
+ }
+ if self.config.sasl_enabled && self.authenticator.is_none() {
+ return Err(KafkaProtocolError::InvalidConfig(
+ "SASL is enabled but no authenticator is configured; every
client would be \
+ rejected. Set IGGY_KAFKA_IGGY_ADDR to the Iggy server that
credentials are \
+ verified against, or unset IGGY_KAFKA_SASL_ENABLED"
+ .into(),
+ ));
+ }
+ Ok(())
+ }
}
/// Cancel in-flight connections, close the tracker to new spawns, and wait
for tasks to finish
diff --git a/gateways/kafka/tests/api_handler_tests.rs
b/gateways/kafka/tests/api_handler_tests.rs
index b12f9e027..98a407bd8 100644
--- a/gateways/kafka/tests/api_handler_tests.rs
+++ b/gateways/kafka/tests/api_handler_tests.rs
@@ -194,12 +194,12 @@ async fn
apiversions_unsupported_version_uses_v0_encoding_without_throttle() {
let body = handle_request(API_KEY_API_VERSIONS, 99, Bytes::new(),
&default_broker())
.await
.expect_response("test request has acks != 0 and expects a response");
- // v0: error_code(2) + api_keys i32 count(4) + 6 entries × 6 bytes = 42 -
no throttle_time_ms.
- assert_eq!(body.len(), 42);
+ // v0: error_code(2) + api_keys i32 count(4) + 7 entries × 6 bytes = 48 -
no throttle_time_ms.
+ assert_eq!(body.len(), 48);
let mut d = Decoder::new(body);
assert_eq!(d.read_i16().unwrap(), ERROR_UNSUPPORTED_VERSION);
- assert_eq!(d.read_i32().unwrap(), 6);
- assert_eq!(d.remaining(), 36);
+ assert_eq!(d.read_i32().unwrap(), 7);
+ assert_eq!(d.remaining(), 42);
}
#[tokio::test]
diff --git a/gateways/kafka/tests/common/scope.rs
b/gateways/kafka/tests/common/scope.rs
index 55b1361d7..43407306b 100644
--- a/gateways/kafka/tests/common/scope.rs
+++ b/gateways/kafka/tests/common/scope.rs
@@ -28,6 +28,7 @@ pub const SCOPED_API_KEYS: &[(i16, &str, i16, i16)] = &[
(3, "Metadata", 0, 9),
(18, "ApiVersions", 0, 3),
(19, "CreateTopics", 2, 5),
+ (22, "InitProducerId", 0, 5),
];
pub fn default_broker() -> BrokerAdvertise {
diff --git a/gateways/kafka/tests/common/wire.rs
b/gateways/kafka/tests/common/wire.rs
index 304604ed1..283767b1b 100644
--- a/gateways/kafka/tests/common/wire.rs
+++ b/gateways/kafka/tests/common/wire.rs
@@ -220,6 +220,29 @@ pub fn build_produce_flexible_empty_request(acks: i16) ->
Bytes {
enc.freeze()
}
+/// `InitProducerId` request for any supported version (v0-v5), flexible from
v2.
+pub fn build_init_producer_id_request(version: i16, transactional_id:
Option<&str>) -> Bytes {
+ let flexible = version >= 2;
+ let mut enc = Encoder::with_capacity(64);
+
+ if flexible {
+ enc.write_compact_nullable_string(transactional_id);
+ } else {
+ enc.write_nullable_string(transactional_id)
+ .expect("transactional id fits");
+ }
+ enc.write_i32(60_000); // transaction_timeout_ms
+ if version >= 3 {
+ enc.write_i64(-1); // producer_id
+ enc.write_i16(-1); // producer_epoch
+ }
+ if flexible {
+ enc.write_empty_tagged_fields();
+ }
+
+ enc.freeze()
+}
+
/// Fetch v4+ minimal empty-topic request.
pub fn build_fetch_empty_topics_request(version: i16) -> Bytes {
let flexible = version >= 12;
diff --git a/gateways/kafka/tests/create_topics_real_bridge_tests.rs
b/gateways/kafka/tests/create_topics_real_bridge_tests.rs
index 9341c6fb0..077a5efda 100644
--- a/gateways/kafka/tests/create_topics_real_bridge_tests.rs
+++ b/gateways/kafka/tests/create_topics_real_bridge_tests.rs
@@ -172,6 +172,7 @@ async fn connected_state(server: &TestServer) ->
GatewayState {
Some(Arc::new(bridge)),
TEST_MAX_FRAME_SIZE,
false,
+ 0,
)
}
diff --git a/gateways/kafka/tests/golden_wire_fixtures_tests.rs
b/gateways/kafka/tests/golden_wire_fixtures_tests.rs
index ebbd9568a..af1188bd9 100644
--- a/gateways/kafka/tests/golden_wire_fixtures_tests.rs
+++ b/gateways/kafka/tests/golden_wire_fixtures_tests.rs
@@ -41,23 +41,25 @@ 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=6 (compact array: N+1=7)
- // key 0 (Produce) min=0 max=9 (advertised)
- // key 1 (Fetch) min=4 max=12
- // key 2 (ListOffsets) min=1 max=6
- // key 3 (Metadata) min=0 max=9
- // key 18 (ApiVersions) min=0 max=3
- // key 19 (CreateTopics) min=2 max=5
+ // error_code=0, api_count=7 (compact array: N+1=8)
+ // key 0 (Produce) min=0 max=9 (advertised)
+ // key 1 (Fetch) min=4 max=12
+ // key 2 (ListOffsets) min=1 max=6
+ // key 3 (Metadata) min=0 max=9
+ // key 18 (ApiVersions) min=0 max=3
+ // key 19 (CreateTopics) min=2 max=5
+ // key 22 (InitProducerId) min=0 max=5
// each entry followed by an empty tagged-fields byte; throttle_ms=0;
top-level tagged fields
- let expected: [u8; 50] = [
+ let expected: [u8; 57] = [
0x00, 0x00, // error_code
- 0x07, // compact array count (6+1)
- 0x00, 0x00, 0x00, 0x00, 0x00, 0x09, 0x00, // key 0: Produce 0-9
(advertised)
- 0x00, 0x01, 0x00, 0x04, 0x00, 0x0C, 0x00, // key 1: Fetch 4-12
- 0x00, 0x02, 0x00, 0x01, 0x00, 0x06, 0x00, // key 2: ListOffsets 1-6
- 0x00, 0x03, 0x00, 0x00, 0x00, 0x09, 0x00, // key 3: Metadata 0-9
- 0x00, 0x12, 0x00, 0x00, 0x00, 0x03, 0x00, // key 18: ApiVersions 0-3
- 0x00, 0x13, 0x00, 0x02, 0x00, 0x05, 0x00, // key 19: CreateTopics 2-5
+ 0x08, // compact array count (7+1)
+ 0x00, 0x00, 0x00, 0x00, 0x00, 0x09, 0x00, // key 0: Produce
0-9 (advertised)
+ 0x00, 0x01, 0x00, 0x04, 0x00, 0x0C, 0x00, // key 1: Fetch
4-12
+ 0x00, 0x02, 0x00, 0x01, 0x00, 0x06, 0x00, // key 2: ListOffsets 1-6
+ 0x00, 0x03, 0x00, 0x00, 0x00, 0x09, 0x00, // key 3: Metadata 0-9
+ 0x00, 0x12, 0x00, 0x00, 0x00, 0x03, 0x00, // key 18: ApiVersions 0-3
+ 0x00, 0x13, 0x00, 0x02, 0x00, 0x05, 0x00, // key 19: CreateTopics 2-5
+ 0x00, 0x16, 0x00, 0x00, 0x00, 0x05, 0x00, // key 22: InitProducerId 0-5
0x00, 0x00, 0x00, 0x00, // throttle_ms
0x00, // top-level tagged fields
];
@@ -71,23 +73,25 @@ async fn golden_apiversions_v1_response_fixture() {
.await
.expect_response("test request has acks != 0 and expects a response");
- // error_code=0, api_count=6
- // key 0 (Produce) min=0 max=9 (KAFKA-18659 advertise min=0)
- // key 1 (Fetch) min=4 max=12
- // key 2 (ListOffsets) min=1 max=6
- // key 3 (Metadata) min=0 max=9
- // key 18 (ApiVersions) min=0 max=3
- // key 19 (CreateTopics) min=2 max=5
+ // error_code=0, api_count=7
+ // key 0 (Produce) min=0 max=9 (KAFKA-18659 advertise min=0)
+ // key 1 (Fetch) min=4 max=12
+ // key 2 (ListOffsets) min=1 max=6
+ // key 3 (Metadata) min=0 max=9
+ // key 18 (ApiVersions) min=0 max=3
+ // key 19 (CreateTopics) min=2 max=5
+ // key 22 (InitProducerId) min=0 max=5
// throttle_ms=0
- let expected: [u8; 46] = [
+ let expected: [u8; 52] = [
0x00, 0x00, // error_code
- 0x00, 0x00, 0x00, 0x06, // api count = 6
- 0x00, 0x00, 0x00, 0x00, 0x00, 0x09, // key 0: Produce 0–9
(advertised)
- 0x00, 0x01, 0x00, 0x04, 0x00, 0x0C, // key 1: Fetch 4–12
- 0x00, 0x02, 0x00, 0x01, 0x00, 0x06, // key 2: ListOffsets 1–6
- 0x00, 0x03, 0x00, 0x00, 0x00, 0x09, // key 3: Metadata 0–9
- 0x00, 0x12, 0x00, 0x00, 0x00, 0x03, // key 18: ApiVersions 0–3
- 0x00, 0x13, 0x00, 0x02, 0x00, 0x05, // key 19: CreateTopics 2–5
+ 0x00, 0x00, 0x00, 0x07, // api count = 7
+ 0x00, 0x00, 0x00, 0x00, 0x00, 0x09, // key 0: Produce 0–9
(advertised)
+ 0x00, 0x01, 0x00, 0x04, 0x00, 0x0C, // key 1: Fetch 4–12
+ 0x00, 0x02, 0x00, 0x01, 0x00, 0x06, // key 2: ListOffsets 1–6
+ 0x00, 0x03, 0x00, 0x00, 0x00, 0x09, // key 3: Metadata 0–9
+ 0x00, 0x12, 0x00, 0x00, 0x00, 0x03, // key 18: ApiVersions 0–3
+ 0x00, 0x13, 0x00, 0x02, 0x00, 0x05, // key 19: CreateTopics 2–5
+ 0x00, 0x16, 0x00, 0x00, 0x00, 0x05, // key 22: InitProducerId 0–5
0x00, 0x00, 0x00, 0x00, // throttle_ms
];
assert_eq!(actual.as_ref(), &expected);
diff --git a/gateways/kafka/tests/idempotence_tests.rs
b/gateways/kafka/tests/idempotence_tests.rs
new file mode 100644
index 000000000..c7d8296bc
--- /dev/null
+++ b/gateways/kafka/tests/idempotence_tests.rs
@@ -0,0 +1,430 @@
+// 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.
+
+//! `InitProducerId` allocation and the transactional refusals that keep
delivery honest.
+
+#[path = "common/codec.rs"]
+mod codec;
+#[path = "common/scope.rs"]
+mod scope;
+#[path = "common/server.rs"]
+mod server;
+#[path = "common/tcp.rs"]
+mod tcp;
+#[path = "common/wire.rs"]
+mod wire;
+
+use std::time::Duration;
+
+use bytes::{BufMut, Bytes, BytesMut};
+use tokio::io::AsyncWriteExt;
+use tokio::net::TcpStream;
+
+use iggy_gateway_kafka::GatewayConfig;
+use iggy_gateway_kafka::protocol::api::{
+ API_KEY_API_VERSIONS, API_KEY_INIT_PRODUCER_ID, API_KEY_PRODUCE,
ERROR_NONE,
+ ERROR_NOT_LEADER_OR_FOLLOWER, ERROR_UNSUPPORTED_VERSION, GatewayState,
handle_request,
+ handle_request_bounded, is_supported_version,
+};
+
+use codec::Decoder;
+use scope::default_broker;
+use server::{spawn_test_server, spawn_test_server_with_config};
+use tcp::{ByteRead, build_request_frame, read_byte_with_timeout, round_trip};
+use wire::{build_api_versions_flexible_request,
build_init_producer_id_request};
+
+/// Transaction APIs this gateway must never advertise. A client that cannot
see them in
+/// `ApiVersions` never sends them, which is the entire enforcement of "no
transactions".
+const TRANSACTION_API_KEYS: &[(i16, &str)] = &[
+ (24, "AddPartitionsToTxn"),
+ (25, "AddOffsetsToTxn"),
+ (26, "EndTxn"),
+ (28, "TxnOffsetCommit"),
+];
+
+const MAX_FRAME_SIZE: usize = 8 * 1024 * 1024;
+
+/// Counter width in a producer id; the instance number sits above it.
+const COUNTER_BITS: i64 = 47;
+
+struct InitProducerIdResponse {
+ error_code: i16,
+ producer_id: i64,
+ producer_epoch: i16,
+}
+
+/// Decode a response body for `version`, asserting nothing is left over so a
wrong flexible
+/// threshold (tagged fields written at v1, or omitted at v2) fails here
rather than passing.
+fn decode_init_producer_id_response(version: i16, body: &Bytes) ->
InitProducerIdResponse {
+ let mut d = Decoder::new(body.clone());
+ assert_eq!(d.read_i32().unwrap(), 0, "throttle_time_ms");
+ let decoded = InitProducerIdResponse {
+ error_code: d.read_i16().unwrap(),
+ producer_id: d.read_i64().unwrap(),
+ producer_epoch: d.read_i16().unwrap(),
+ };
+ if version >= 2 {
+ d.read_tagged_fields().unwrap();
+ }
+ assert_eq!(d.remaining(), 0, "v{version} response has trailing bytes");
+ decoded
+}
+
+async fn init_producer_id(
+ state: &GatewayState,
+ version: i16,
+ transactional_id: Option<&str>,
+) -> InitProducerIdResponse {
+ let request = build_init_producer_id_request(version, transactional_id);
+ let body = handle_request_bounded(state, API_KEY_INIT_PRODUCER_ID,
version, request)
+ .await
+ .expect_response("InitProducerId always answers, it is never
fire-and-forget");
+ decode_init_producer_id_response(version, &body)
+}
+
+fn stub_state(instance_id: u16) -> GatewayState {
+ GatewayState::new(default_broker(), None, MAX_FRAME_SIZE, false,
instance_id)
+}
+
+/// Produce v3 body with one topic and one partition, so a per-partition error
code has somewhere
+/// to land.
+fn produce_v3_body(acks: i16, transactional_id: Option<&str>) -> Bytes {
+ let mut body = BytesMut::new();
+ match transactional_id {
+ Some(id) => {
+ body.put_i16(i16::try_from(id.len()).expect("transactional id fits
i16"));
+ body.put_slice(id.as_bytes());
+ }
+ None => body.put_i16(-1),
+ }
+ body.put_i16(acks);
+ body.put_i32(1_000); // timeout_ms
+ body.put_i32(1); // one topic
+ body.put_i16(6);
+ body.put_slice(b"orders");
+ body.put_i32(1); // one partition
+ body.put_i32(0); // partition index
+ body.put_i32(4); // records length
+ body.put_slice(&[0x00, 0x00, 0x00, 0x00]);
+ body.freeze()
+}
+
+/// Skips to the first per-partition `error_code` of a Produce v3 response.
+fn first_produce_partition_error(body: &Bytes) -> i16 {
+ let mut d = Decoder::new(body.clone());
+ assert_eq!(d.read_i32().unwrap(), 1, "one topic in the response");
+ assert_eq!(d.read_nullable_string().unwrap().as_deref(), Some("orders"));
+ assert_eq!(d.read_i32().unwrap(), 1, "one partition in the response");
+ assert_eq!(d.read_i32().unwrap(), 0, "partition index");
+ d.read_i16().unwrap()
+}
+
+#[tokio::test]
+async fn
given_no_transactional_id_when_init_producer_id_should_allocate_an_id_with_epoch_zero()
{
+ let state = stub_state(0);
+ for version in 0i16..=5 {
+ let response = init_producer_id(&state, version, None).await;
+ assert_eq!(response.error_code, ERROR_NONE, "v{version} error_code");
+ assert_eq!(response.producer_epoch, 0, "v{version} producer_epoch");
+ assert!(
+ response.producer_id >= 0,
+ "v{version} producer id must stay non-negative: -1 means no
producer id"
+ );
+ }
+}
+
+#[tokio::test]
+async fn
given_one_gateway_when_two_producers_init_should_receive_distinct_ids() {
+ let state = stub_state(0);
+ let first = init_producer_id(&state, 4, None).await;
+ let second = init_producer_id(&state, 4, None).await;
+ assert_eq!(first.error_code, ERROR_NONE);
+ assert_eq!(second.error_code, ERROR_NONE);
+ assert_ne!(
+ first.producer_id, second.producer_id,
+ "each InitProducerId must advance the counter"
+ );
+}
+
+#[tokio::test]
+async fn
given_a_configured_instance_id_when_init_producer_id_should_return_it_in_the_high_bits()
{
+ let instance_id = 0x0123u16;
+ let response = init_producer_id(&stub_state(instance_id), 4, None).await;
+ assert_eq!(response.error_code, ERROR_NONE);
+ assert_eq!(
+ response.producer_id >> COUNTER_BITS,
+ i64::from(instance_id),
+ "the configured instance number is what makes ids unique across
gateways"
+ );
+}
+
+#[tokio::test]
+async fn
given_a_transactional_id_when_init_producer_id_should_answer_unsupported_version()
{
+ let state = stub_state(0);
+ for version in 0i16..=5 {
+ let response = init_producer_id(&state, version,
Some("orders-txn")).await;
+ assert_eq!(
+ response.error_code, ERROR_UNSUPPORTED_VERSION,
+ "v{version} must refuse a transactional producer"
+ );
+ // A real broker's error path leaves both at their schema defaults,
and -1 is the
+ // "no producer id" sentinel the whole bit-63-clear layout exists to
keep distinct
+ // from a real allocation.
+ assert_eq!(
+ response.producer_id, -1,
+ "v{version} a refusal must not hand back an allocated id"
+ );
+ assert_eq!(response.producer_epoch, 0, "v{version} producer_epoch");
+ }
+}
+
+#[tokio::test]
+async fn
given_an_empty_transactional_id_when_init_producer_id_should_still_allocate() {
+ let response = init_producer_id(&stub_state(0), 4, Some("")).await;
+ assert_eq!(
+ response.error_code, ERROR_NONE,
+ "an empty transactional id names no transaction"
+ );
+ assert!(response.producer_id >= 0);
+}
+
+#[tokio::test]
+async fn
given_a_transactional_id_when_producing_should_answer_unsupported_version_per_partition()
{
+ let body = handle_request(
+ API_KEY_PRODUCE,
+ 3,
+ produce_v3_body(1, Some("orders-txn")),
+ &default_broker(),
+ )
+ .await
+ .expect_response("acks=1 expects a response");
+ assert_eq!(
+ first_produce_partition_error(&body),
+ ERROR_UNSUPPORTED_VERSION,
+ "a transactional batch must be refused, not stored as ordinary records"
+ );
+}
+
+#[tokio::test]
+async fn
given_no_transactional_id_when_producing_should_keep_the_retriable_stub_error()
{
+ let body = handle_request(
+ API_KEY_PRODUCE,
+ 3,
+ produce_v3_body(1, None),
+ &default_broker(),
+ )
+ .await
+ .expect_response("acks=1 expects a response");
+ assert_eq!(
+ first_produce_partition_error(&body),
+ ERROR_NOT_LEADER_OR_FOLLOWER,
+ "a non-transactional produce keeps the retriable stub error, not the
transactional refusal"
+ );
+}
+
+#[tokio::test]
+async fn
given_acks_zero_and_a_transactional_id_when_producing_should_close_the_connection()
{
+ // acks=0 has no response to carry the refusal, and dropping the batch
silently would leave
+ // the refusal to whatever write path lands later. The acks=1 half pins
that the same body is
+ // refused when a response exists, so the close is the acks=0 form of that
refusal.
+ let refused = handle_request(
+ API_KEY_PRODUCE,
+ 3,
+ produce_v3_body(1, Some("orders-txn")),
+ &default_broker(),
+ )
+ .await
+ .expect_response("acks=1 expects a response");
+ assert_eq!(
+ first_produce_partition_error(&refused),
+ ERROR_UNSUPPORTED_VERSION,
+ "the same body must be refused when the client is reading a response"
+ );
+
+ let outcome = handle_request(
+ API_KEY_PRODUCE,
+ 3,
+ produce_v3_body(0, Some("orders-txn")),
+ &default_broker(),
+ )
+ .await;
+ assert!(
+ outcome.is_close(),
+ "acks=0 cannot carry 35, so a transactional batch must be refused by
closing"
+ );
+}
+
+#[tokio::test]
+async fn
given_the_advertised_api_list_when_a_client_reads_it_should_omit_the_transaction_keys()
{
+ let legacy = handle_request(API_KEY_API_VERSIONS, 1, Bytes::new(),
&default_broker())
+ .await
+ .expect_response("ApiVersions always answers");
+ let mut d = Decoder::new(legacy);
+ assert_eq!(d.read_i16().unwrap(), ERROR_NONE);
+ let count = d.read_i32().unwrap();
+ let mut advertised = Vec::new();
+ for _ in 0..count {
+ advertised.push(d.read_i16().unwrap());
+ d.read_i16().unwrap(); // min_version
+ d.read_i16().unwrap(); // max_version
+ }
+
+ let flexible = handle_request(
+ API_KEY_API_VERSIONS,
+ 3,
+ build_api_versions_flexible_request("iggy-test", "0.1.0"),
+ &default_broker(),
+ )
+ .await
+ .expect_response("ApiVersions always answers");
+ let mut d = Decoder::new(flexible);
+ assert_eq!(d.read_i16().unwrap(), ERROR_NONE);
+ let count = d.read_varint().unwrap() - 1;
+ let mut advertised_flexible = Vec::new();
+ for _ in 0..count {
+ advertised_flexible.push(d.read_i16().unwrap());
+ d.read_i16().unwrap(); // min_version
+ d.read_i16().unwrap(); // max_version
+ d.read_tagged_fields().unwrap();
+ }
+
+ assert!(
+ advertised.contains(&API_KEY_INIT_PRODUCER_ID),
+ "InitProducerId must be advertised or no producer ever sends it"
+ );
+ for &(api_key, name) in TRANSACTION_API_KEYS {
+ assert!(
+ !advertised.contains(&api_key),
+ "{name} (key {api_key}) must stay out of the v1 advertisement"
+ );
+ assert!(
+ !advertised_flexible.contains(&api_key),
+ "{name} (key {api_key}) must stay out of the v3 advertisement"
+ );
+ assert!(
+ !is_supported_version(api_key, 0),
+ "{name} (key {api_key}) must not pass the version firewall"
+ );
+ }
+}
+
+#[tokio::test]
+async fn
given_a_transaction_api_key_when_sent_anyway_should_close_the_connection() {
+ for &(api_key, name) in TRANSACTION_API_KEYS {
+ let outcome = handle_request(api_key, 0, Bytes::new(),
&default_broker()).await;
+ assert!(
+ outcome.is_close(),
+ "{name} (key {api_key}) has no response schema here and must close"
+ );
+ }
+}
+
+#[tokio::test]
+async fn
given_a_live_server_when_a_transactional_request_is_refused_should_keep_the_connection()
{
+ let (addr, _shutdown) = spawn_test_server().await;
+ let mut stream = TcpStream::connect(addr).await.expect("connect");
+
+ for (api_key, version, body) in [
+ (
+ API_KEY_INIT_PRODUCER_ID,
+ 4i16,
+ build_init_producer_id_request(4, Some("orders-txn")),
+ ),
+ (API_KEY_PRODUCE, 3, produce_v3_body(1, Some("orders-txn"))),
+ ] {
+ let frame = build_request_frame(api_key, version, 7_000,
Some("txn-test"), &body);
+ stream.write_all(&frame).await.expect("write request");
+ let payload = tcp::read_response_frame(&mut stream,
MAX_FRAME_SIZE).await;
+ let (correlation_id, _) = tcp::parse_response_payload(api_key,
version, payload);
+ assert_eq!(correlation_id, 7_000, "key {api_key} correlation id");
+ assert_ne!(
+ read_byte_with_timeout(&mut stream,
Duration::from_millis(250)).await,
+ ByteRead::Closed,
+ "key {api_key} must refuse the transaction without dropping the
connection"
+ );
+ }
+}
+
+/// Config value has to reach the allocator, not just the struct. Nothing else
pins that hop:
+/// one test pins env to config, another pins `GatewayState::new` to the high
bits, and the
+/// server's own `config.instance_id` argument sat between them uncovered.
+#[tokio::test]
+async fn
given_a_server_configured_with_an_instance_id_when_init_producer_id_should_reflect_it()
{
+ let instance_id = 0x0042u16;
+ let (addr, _shutdown) = spawn_test_server_with_config(GatewayConfig {
+ bind_addr: String::new(),
+ advertised_host: None,
+ advertised_port: None,
+ max_frame_size: MAX_FRAME_SIZE,
+ max_connections: 1024,
+ idle_timeout: Duration::from_secs(5),
+ read_timeout: Duration::from_secs(5),
+ write_timeout: Duration::from_secs(5),
+ shutdown_drain_timeout: Duration::from_secs(5),
+ instance_id,
+ ..GatewayConfig::default()
+ })
+ .await;
+
+ let request = build_init_producer_id_request(4, None);
+ let (_correlation_id, body) =
+ round_trip(addr, API_KEY_INIT_PRODUCER_ID, 4, 7_200, &request).await;
+ let response = decode_init_producer_id_response(4, &body);
+ assert_eq!(response.error_code, ERROR_NONE);
+ assert_eq!(
+ response.producer_id >> COUNTER_BITS,
+ i64::from(instance_id),
+ "the configured instance number must reach the allocator, not stop at
the config struct"
+ );
+}
+
+/// One allocator per process, shared across connections. Building a
`GatewayState` per accepted
+/// connection instead would restart every counter at 0 and put duplicate
producer ids on the
+/// wire, which Kafka requires to be unique; the in-process test above cannot
see that because it
+/// never opens a second connection.
+#[tokio::test]
+async fn
given_one_server_when_two_connections_init_should_receive_distinct_ids() {
+ let (addr, _shutdown) = spawn_test_server().await;
+ let request = build_init_producer_id_request(4, None);
+
+ let (_first_id, first_body) =
+ round_trip(addr, API_KEY_INIT_PRODUCER_ID, 4, 7_300, &request).await;
+ let (_second_id, second_body) =
+ round_trip(addr, API_KEY_INIT_PRODUCER_ID, 4, 7_301, &request).await;
+
+ let first = decode_init_producer_id_response(4, &first_body);
+ let second = decode_init_producer_id_response(4, &second_body);
+ assert_eq!(first.error_code, ERROR_NONE);
+ assert_eq!(second.error_code, ERROR_NONE);
+ assert_ne!(
+ first.producer_id, second.producer_id,
+ "separate connections must draw from one allocator"
+ );
+}
+
+#[tokio::test]
+async fn
given_a_live_server_when_init_producer_id_round_trips_should_return_an_allocated_id()
{
+ let (addr, _shutdown) = spawn_test_server().await;
+ let request = build_init_producer_id_request(4, None);
+ let (correlation_id, body) =
+ round_trip(addr, API_KEY_INIT_PRODUCER_ID, 4, 7_100, &request).await;
+ assert_eq!(correlation_id, 7_100);
+ let response = decode_init_producer_id_response(4, &body);
+ assert_eq!(response.error_code, ERROR_NONE);
+ assert_eq!(response.producer_epoch, 0);
+ assert!(response.producer_id >= 0);
+}
diff --git a/gateways/kafka/tests/list_offsets_real_bridge_tests.rs
b/gateways/kafka/tests/list_offsets_real_bridge_tests.rs
index 81ed93208..cd59b620e 100644
--- a/gateways/kafka/tests/list_offsets_real_bridge_tests.rs
+++ b/gateways/kafka/tests/list_offsets_real_bridge_tests.rs
@@ -173,6 +173,7 @@ async fn connected_state(server: &TestServer) ->
(GatewayState, IggyBridge) {
Some(Arc::new(bridge)),
TEST_MAX_FRAME_SIZE,
false,
+ 0,
);
(state, seed_bridge)
}
diff --git a/gateways/kafka/tests/listener_robustness_tests.rs
b/gateways/kafka/tests/listener_robustness_tests.rs
index 2643a5d0a..905891372 100644
--- a/gateways/kafka/tests/listener_robustness_tests.rs
+++ b/gateways/kafka/tests/listener_robustness_tests.rs
@@ -357,7 +357,7 @@ async fn e2e_flexible_apiversions_v3_request_succeeds() {
let mut d = Decoder::new(body);
assert_eq!(d.read_i16().unwrap(), 0);
let count = usize::try_from(d.read_varint().unwrap() - 1).unwrap();
- assert_eq!(count, 6, "must advertise all six scoped API keys");
+ assert_eq!(count, 7, "must advertise all seven scoped API keys");
}
#[tokio::test]
diff --git a/gateways/kafka/tests/metadata_real_bridge_tests.rs
b/gateways/kafka/tests/metadata_real_bridge_tests.rs
index ef81a6c09..60fb3b8b8 100644
--- a/gateways/kafka/tests/metadata_real_bridge_tests.rs
+++ b/gateways/kafka/tests/metadata_real_bridge_tests.rs
@@ -160,6 +160,7 @@ async fn connected_state(server: &TestServer) ->
(GatewayState, IggyBridge) {
Some(Arc::new(bridge)),
TEST_MAX_FRAME_SIZE,
false,
+ 0,
);
(state, seed)
}
@@ -280,6 +281,7 @@ async fn
a_named_lookup_reports_the_kafka_side_name_through_a_topic_mapping_over
Some(Arc::new(bridge)),
TEST_MAX_FRAME_SIZE,
false,
+ 0,
);
let topics = send(&state, Some(&["orders"])).await;
@@ -311,7 +313,7 @@ async fn
a_response_projected_over_max_frame_size_closes_instead_of_answering()
// 50 partitions * 64 bytes/partition (this crate's own conservative
per-partition estimate)
// = 3200 bytes, comfortably over a 512-byte max_frame_size.
- let tiny_state = GatewayState::new(state.broker, state.bridge, 512, false);
+ let tiny_state = GatewayState::new(state.broker, state.bridge, 512, false,
0);
let body = build_request(Some(&["orders"]));
let outcome = metadata::handle(&tiny_state, REQUEST_VERSION, body).await;
assert!(outcome.is_close(), "expected Close, got {outcome:?}");
@@ -337,7 +339,7 @@ async fn
an_all_topics_response_over_max_frame_size_truncates_instead_of_closing
// 50 partitions * 64 bytes/partition (this crate's own conservative
per-partition estimate)
// = 3200 bytes, comfortably over a 512-byte max_frame_size - so the
catalog as a whole cannot
// fit, but neither topic's own partition count is malformed or
attacker-shaped.
- let tiny_state = GatewayState::new(state.broker, state.bridge, 512, false);
+ let tiny_state = GatewayState::new(state.broker, state.bridge, 512, false,
0);
let topics = send(&tiny_state, None).await;
assert!(
topics.len() < 2,
diff --git a/gateways/kafka/tests/produce_real_bridge_tests.rs
b/gateways/kafka/tests/produce_real_bridge_tests.rs
index f22ae0b43..7424ea3ba 100644
--- a/gateways/kafka/tests/produce_real_bridge_tests.rs
+++ b/gateways/kafka/tests/produce_real_bridge_tests.rs
@@ -76,6 +76,7 @@ async fn gateway_with_topic(server: &TestServer, partitions:
u32) -> GatewayStat
Some(Arc::new(bridge)),
MAX_FRAME_SIZE,
false,
+ 0,
)
}
diff --git a/gateways/kafka/tests/server_e2e_tests.rs
b/gateways/kafka/tests/server_e2e_tests.rs
index bca20a8af..1d0f6eee0 100644
--- a/gateways/kafka/tests/server_e2e_tests.rs
+++ b/gateways/kafka/tests/server_e2e_tests.rs
@@ -70,7 +70,7 @@ async fn
e2e_apiversions_v3_flexible_preserves_correlation_id() {
let mut d = Decoder::new(body);
assert_eq!(d.read_i16().unwrap(), 0);
let count = usize::try_from(d.read_varint().unwrap() - 1).expect("api
count fits usize");
- assert_eq!(count, 6);
+ assert_eq!(count, 7);
}
#[tokio::test]
diff --git a/gateways/kafka/tests/version_firewall_tests.rs
b/gateways/kafka/tests/version_firewall_tests.rs
index 5f3813133..e953bc96a 100644
--- a/gateways/kafka/tests/version_firewall_tests.rs
+++ b/gateways/kafka/tests/version_firewall_tests.rs
@@ -37,8 +37,8 @@ use tokio::io::AsyncWriteExt;
use tokio::net::TcpStream;
use iggy_gateway_kafka::protocol::api::{
- API_KEY_API_VERSIONS, API_KEY_CREATE_TOPICS, API_KEY_FETCH,
API_KEY_LIST_OFFSETS,
- API_KEY_METADATA, API_KEY_PRODUCE, ERROR_INVALID_REQUEST, ERROR_NONE,
+ API_KEY_API_VERSIONS, API_KEY_CREATE_TOPICS, API_KEY_FETCH,
API_KEY_INIT_PRODUCER_ID,
+ API_KEY_LIST_OFFSETS, API_KEY_METADATA, API_KEY_PRODUCE,
ERROR_INVALID_REQUEST, ERROR_NONE,
ERROR_UNSUPPORTED_VERSION, advertised_min_version, handle_request,
is_supported_version,
supported_api_ranges,
};
@@ -54,14 +54,14 @@ use tcp::{
};
use wire::{
OUT_OF_SCOPE_API_KEYS, build_api_versions_flexible_request,
build_create_topics_empty_request,
- build_fetch_empty_topics_request, build_list_offsets_request,
+ build_fetch_empty_topics_request, build_init_producer_id_request,
build_list_offsets_request,
build_metadata_all_topics_flexible, build_metadata_all_topics_legacy,
build_metadata_flexible_request_v10,
};
#[test]
-fn supported_ranges_table_has_six_entries() {
- assert_eq!(supported_api_ranges().len(), 6);
+fn supported_ranges_table_has_seven_entries() {
+ assert_eq!(supported_api_ranges().len(), 7);
}
#[test]
@@ -490,6 +490,7 @@ fn request_body_for_scoped_api(api_key: i16, name: &str,
version: i16) -> Bytes
.flatten()
.unwrap_or_else(|| build_list_offsets_request(version,
"scope-topic", 0)),
API_KEY_CREATE_TOPICS => build_create_topics_empty_request(version),
+ API_KEY_INIT_PRODUCER_ID => build_init_producer_id_request(version,
None),
_ => Bytes::new(),
}
}
diff --git a/gateways/kafka/tools/kafka-tool/src/response.rs
b/gateways/kafka/tools/kafka-tool/src/response.rs
index 4214a2ce2..467f225fd 100644
--- a/gateways/kafka/tools/kafka-tool/src/response.rs
+++ b/gateways/kafka/tools/kafka-tool/src/response.rs
@@ -19,8 +19,8 @@
use bytes::Bytes;
use kafka_protocol::messages::{
- ApiKey, ApiVersionsResponse, CreateTopicsResponse, FetchResponse,
ListOffsetsResponse,
- MetadataResponse, ProduceResponse,
+ ApiKey, ApiVersionsResponse, CreateTopicsResponse, FetchResponse,
InitProducerIdResponse,
+ ListOffsetsResponse, MetadataResponse, ProduceResponse,
};
use kafka_protocol::protocol::Decodable;
@@ -355,6 +355,20 @@ fn decode_body(
));
}
}
+ 22 => {
+ let resp = InitProducerIdResponse::decode(&mut buf, api_version)?;
+ codes.push(resp.error_code);
+ details.push(format!("throttle_time_ms={}",
resp.throttle_time_ms));
+ details.push(format!(
+ "top_level.error_code={} ({})",
+ resp.error_code,
+ format_error_code(resp.error_code)
+ ));
+ details.push(format!(
+ "producer_id={} producer_epoch={}",
+ resp.producer_id.0, resp.producer_epoch
+ ));
+ }
other => {
details.push(format!("no schema decoder for api_key={other}"));
if body.len() >= 2 {