This is an automated email from the ASF dual-hosted git repository.
krishvishal 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 63d97a66d docs(gateways): settle the Kafka record mapping, producer
ids, offsets (#4205)
63d97a66d is described below
commit 63d97a66d85a173657b003ce389beeb7ab43adef
Author: Krishna Vishal <[email protected]>
AuthorDate: Fri Sep 18 13:33:12 2026 +0530
docs(gateways): settle the Kafka record mapping, producer ids, offsets
(#4205)
---
gateways/kafka/README.md | 39 +++++
gateways/kafka/docs/BRIDGE_MAPPING.md | 275 ++++++++++++++++++++++++++++++++++
gateways/kafka/docs/IDEMPOTENCE.md | 179 ++++++++++++++++++++++
gateways/kafka/docs/OFFSET_STORAGE.md | 163 ++++++++++++++++++++
gateways/kafka/docs/SCOPE.md | 12 +-
5 files changed, 665 insertions(+), 3 deletions(-)
diff --git a/gateways/kafka/README.md b/gateways/kafka/README.md
index a1e10c010..f20292390 100644
--- a/gateways/kafka/README.md
+++ b/gateways/kafka/README.md
@@ -58,6 +58,23 @@ Before check-in, run the procedure in
[docs/MANUAL_TESTING.md](docs/MANUAL_TESTI
See [docs/SCOPE.md](docs/SCOPE.md) for
[#3421](https://github.com/apache/iggy/issues/3421) deliverables, supported API
key/version table, and post-foundation TODO backlog.
+## Design decisions
+
+- [docs/BRIDGE_MAPPING.md](docs/BRIDGE_MAPPING.md) — how a Kafka record
becomes an Iggy message, and back
+- [docs/IDEMPOTENCE.md](docs/IDEMPOTENCE.md) — InitProducerId, and why
delivery is at-least-once
+- [docs/OFFSET_STORAGE.md](docs/OFFSET_STORAGE.md) — where Kafka consumer
group offsets live
+
+### 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. An idempotent Kafka
producer is given
+a producer id so that it starts, but its retries are not deduplicated: 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.
+
## Iggy bridge ([#3533](https://github.com/apache/iggy/issues/3533))
`src/bridge/` is the SDK integration layer: connects to Iggy, maps Kafka
topics to Iggy
@@ -188,6 +205,28 @@ rely on.
through `ensure_topic`'s `partition_count` argument once it exceeds the
server's cap
- Anything else → `UNKNOWN_SERVER_ERROR` (-1)
+### Server limits the gateway inherits
+
+These are Iggy server limits, not gateway settings. A Kafka client cannot act
on any of them, so
+an operator has to.
+
+| Limit | Default | Where |
+| ------- | --------- | ------- |
+| Consumer offset keys per partition, per consumer kind | 4096, ceiling 262144
| `partition.consumer_offsets_max` |
+| One user header name, and one header value | 255 bytes | fixed,
`user_headers.rs` |
+| All user headers of one message | 100 KB | fixed, `MAX_USER_HEADERS_SIZE` |
+| Message payload | 64 MB | fixed, `MAX_PAYLOAD_SIZE` |
+
+Only the first is configurable. A Kafka consumer group commits one offset key
per partition it
+holds, so `partition.consumer_offsets_max` is what bounds the number of groups
that can commit
+against one partition. Passing it returns `TooManyConsumerOffsets` (3024),
which reaches the
+client as `UNKNOWN_SERVER_ERROR` because Kafka has no code for the condition.
The gateway logs
+the real Iggy error, so the server log is where an operator diagnoses it.
+
+The other three decide when a Kafka record goes into the envelope instead of
being stored
+natively. See [docs/OFFSET_STORAGE.md](docs/OFFSET_STORAGE.md) and
+[docs/BRIDGE_MAPPING.md](docs/BRIDGE_MAPPING.md).
+
## Wire fixture tool
See [tools/kafka-tool/README.md](tools/kafka-tool/README.md).
diff --git a/gateways/kafka/docs/BRIDGE_MAPPING.md
b/gateways/kafka/docs/BRIDGE_MAPPING.md
new file mode 100644
index 000000000..da213deeb
--- /dev/null
+++ b/gateways/kafka/docs/BRIDGE_MAPPING.md
@@ -0,0 +1,275 @@
+# Kafka to Iggy record mapping
+
+Status: proposed. Direction agreed with @spetz and @hubcio on 2026-09-16
("messages should be
+stored in iggy format"); the escape hatches below still need sign-off. Closes
the last open
+scope item of [#3533](https://github.com/apache/iggy/issues/3533) and blocks
+[#3535](https://github.com/apache/iggy/issues/3535) (Produce) and
+[#3536](https://github.com/apache/iggy/issues/3536) (Fetch).
+
+## Decision
+
+One Kafka record becomes one Iggy message, in Iggy's own format: the record
value is the
+message payload, the key and the Kafka headers become Iggy user headers. The
gateway rebuilds
+a Kafka record batch on Fetch.
+
+Two properties drive this.
+
+A Kafka consumer must be able to read a topic an Iggy producer wrote. This is
the staged
+migration the maintainers described: rewrite producers to the Iggy SDK first,
leave consumers
+on the gateway until later. The gateway can only encode an arbitrary Iggy
message as a Kafka
+record if the stored form has no Kafka framing in it.
+
+Kafka offsets must line up with Iggy offsets. A Kafka record batch carries N
records under one
+base offset, while an Iggy message consumes exactly one offset. Storing a
batch whole makes
+every offset the gateway reports wrong by the batch size, and fixing that
inside Iggy means
+teaching the server to count records inside an opaque payload.
+
+Storing the Kafka payload and headers as a dump inside the Iggy payload was
the alternative
+raised in the same thread. It does not satisfy the first property on its own:
the gateway would
+still need a native path for Iggy-written messages, so it would carry two
storage formats
+instead of one. Native storage with a narrow fallback (see below) keeps that
to one.
+
+## Field mapping
+
+Produce, per record:
+
+| Kafka | Iggy |
+| ------- | ------ |
+| record value | `payload` |
+| record key | `kafka.key` user header, `Raw` |
+| record header `name` | `kafka.h.<name>` user header, `Raw` |
+| record timestamp (CreateTime), milliseconds | `origin_timestamp`,
microseconds |
+| record offset | partition offset, assigned by Iggy |
+| partition index | partition index, both 0-based |
+| topic | stream and topic per `TopicMapping` |
+
+Fetch reverses it. A message with no `kafka.*` headers is a message an Iggy
client wrote, and
+encodes as a Kafka record with a null key, its Iggy user headers as Kafka
headers, and
+`origin_timestamp` as the record timestamp (falling back to the
server-assigned `timestamp`
+when the origin timestamp is zero). Iggy header kinds other than `Raw` and
`String` are emitted
+as their raw value bytes.
+
+### Timestamps
+
+Kafka carries a record timestamp in milliseconds. Iggy carries
`origin_timestamp` in
+microseconds (`core/common/src/types/message/iggy_message.rs:191`). Produce
multiplies by 1000.
+Fetch divides by 1000 and truncates toward zero.
+
+A record that arrived through Produce survives the round trip exactly, because
its microsecond
+value is always a whole number of milliseconds. A message an Iggy client wrote
does not. Its
+sub-millisecond digits are lost on the way out, and Kafka has no field to keep
them in.
+
+Kafka sends `-1` for a record with no timestamp. That is stored as `0`, and
Fetch already reads
+a zero origin timestamp as an instruction to use the server-assigned timestamp
instead. A real
+broker does the same thing under `LogAppendTime`, so the two agree.
+
+## Records Iggy cannot hold natively
+
+Iggy rejects an empty payload
(`core/common/src/types/message/iggy_message.rs:169`), caps a
+user header value at 255 bytes
(`core/common/src/types/message/user_headers.rs:631`), keys
+headers in a `BTreeMap` so a name cannot repeat, and caps all user headers of
one message at
+100 KB (`MAX_USER_HEADERS_SIZE`, `iggy_message.rs:58`). Kafka allows all of
the shapes those
+rules exclude, so two mechanisms cover them.
+
+None of those four numbers is a gateway setting. They are fixed constants on
the server's
+message type, with no configuration knob and no recorded rationale. The header
budget works out
+to roughly 350 headers at the 255-byte value cap, so in practice the
key-length rule below is
+what sends a record into the fallback and the budget is not.
+
+### Null and empty values
+
+A record with a null value (a tombstone) or a zero-length value is stored with
a single `0x00`
+byte payload and a `kafka.value` header holding `null` or `empty`. Fetch reads
that header and
+restores the original, discarding the placeholder byte.
+
+This keeps tombstones on the fast path rather than pushing them into the
fallback, because they
+are ordinary traffic on compacted Kafka topics. Iggy has no compaction, so a
tombstone is stored
+and served like any other record and nothing acts on it.
+
+### Everything else: the envelope fallback
+
+A record takes the fallback when any of these hold:
+
+- the key is present and zero-length, which is not the same as a null key.
Iggy rejects an
+ empty header value on the same rule that caps it at 255 bytes, so
`kafka.key` cannot carry it
+- the key is longer than 255 bytes
+- a header name, prefixed with `kafka.h.`, is longer than 255 bytes
+- a header value is null, empty, or longer than 255 bytes
+- two headers share a name
+- the headers together would exceed the 100 KB user-header budget
+
+Such a record is stored with a `kafka.envelope` header whose value is one
byte, the format
+version, currently `1`. The payload holds the key, the value and the headers
in the layout
+below. Fetch checks for that header first and takes the plain path only when
it is absent.
+
+The layout is fixed here rather than left to the implementation, because
`kafka_protocol` hands
+a handler a decoded `Record` and never a raw slice of the record, so there is
no verbatim body
+to copy. All integers are little-endian.
+
+```text
+u8 flags bit 0 key present, bit 1 value present
+u32 key_len 0 when the key is absent
+.. key
+u32 value_len 0 when the value is absent
+.. value
+u32 header_count
+ repeated header_count times:
+ u32 name_len
+ .. name
+ u8 value_present
+ u32 value_len 0 when the header value is absent
+ .. value
+```
+
+That is 13 bytes of fixed overhead plus 9 bytes per header. Re-encoding the
record as a
+one-record Kafka batch would also work and would cost less code, but it puts
batch framing back
+into storage, which is the thing this document decided against.
+
+The cost is that these messages are opaque to Iggy consumers and connectors.
That is the point
+of confining the fallback to record shapes that are rare in practice, rather
than making it the
+default storage form.
+
+## Batch-level fields
+
+Per-record storage drops what the Kafka record batch header carries: producer
id, producer
+epoch, base sequence, the transactional flag, compression and the batch CRC.
Fetch synthesizes
+a batch with producer id `-1`, epoch `-1`, base sequence `-1`, no compression,
`CreateTime`
+timestamps, and a recomputed CRC32C.
+
+Two consequences worth stating before they surprise someone:
+
+- Idempotent-producer deduplication cannot be reconstructed from stored data
later. If
+ [#3545](https://github.com/apache/iggy/issues/3545) ever grows past a stub,
producer id,
+ epoch and sequence need their own tracking.
+- The bytes a consumer receives are not the bytes the producer sent, so
anything comparing
+ batches byte for byte across the gateway will differ.
+
+Whether producer id `-1` is what Fetch actually sends depends on the
InitProducerId decision in
+[`IDEMPOTENCE.md`](IDEMPOTENCE.md). Allocating producer ids does not change
what is stored, only
+what Produce accepts, so this section holds under either answer.
+
+Produce decompresses gzip, snappy, lz4 and zstd batches, which means turning
those features back
+on for the `kafka-protocol` dependency (`Cargo.toml:217` currently builds it
with
+`default-features = false, features = ["broker"]`). Fetch emits uncompressed
batches.
+
+Decompression needs its own bound, and the bound is per request rather than
per batch.
+`max_frame_size` bounds the frame a client sent, which is the compressed size,
and zstd reaches
+1000 to 1 on repetitive input without being asked, so an 8 MiB frame can
expand to gigabytes. One
+frame also carries many batches: the request holds up to
`MAX_REQUEST_ELEMENTS` (4096) topic and
+partition entries, each with its own records blob. A cap applied to one batch
at a time would
+therefore still admit 4096 times that much output.
+
+Produce keeps a single decompression budget for the whole request, set to
`max_frame_size`, so a
+compressed request can never yield more than the same client could have sent
uncompressed. A
+batch that exhausts the budget is rejected with `MESSAGE_TOO_LARGE` (10)
before the output is
+allocated. Each decompressed record value has to clear Iggy's own
`MAX_PAYLOAD_SIZE` (64 MB,
+`iggy_message.rs:44`) separately, since one record becomes one message.
+
+Nothing decompresses today. The record batch stays an opaque `Bytes` on both
paths, so the bound
+above is a requirement on [#3535](https://github.com/apache/iggy/issues/3535)
rather than a
+description of current behavior.
+
+## Offsets
+
+Kafka offset and Iggy offset are the same number for the same record, and both
partition spaces
+are 0-based, so neither direction converts.
+
+Produce takes the base offset from the send confirmation
+(`SendMessagesConfirmationResponse::base_offset`). The server may return no
confirmation, for
+example for a request it classifies as a duplicate, in which case the response
carries `-1`
+rather than a guessed offset; Kafka clients surface that as an unknown offset.
+
+ListOffsets LATEST is the high watermark from `IggyBridge::high_watermarks`.
EARLIEST has no
+server-side field today (`Partition` carries no log start offset), so it reads
the first
+retained message instead, and the `(messages_count, current_offset) == (0, 0)`
ambiguity
+documented on `high_watermarks` applies to both.
+
+### Header order
+
+Kafka carries record headers as an ordered list. Iggy keys them in a
`BTreeMap`, so Fetch emits
+them sorted by name and the producer's order is gone.
+
+There is no fallback for this, because the gateway cannot tell whether a
record's header order
+carries meaning. The envelope does preserve order, since it stores the headers
as a list, but a
+record only reaches the envelope for one of the reasons above. A consumer that
depends on header
+order therefore sees a different order through the gateway than a real broker
would give it.
+
+## Partitioning
+
+Both systems number partitions from 0, so the partition index passes through
unchanged in each
+direction and neither side converts.
+
+Produce sends to the partition the request names,
`Partitioning::partition_id(index)`. A Kafka
+producer resolves the partition itself before it builds the request, so every
partition index in
+a `ProduceRequest` is a real one and `Partitioning::balanced()` has no trigger
on this path. The
+`-1` that `SCOPE.md` refers to belongs to CreateTopics, where it means "use
the broker default
+partition count", and it is handled there rather than here.
+
+Kafka consumer groups are not mapped onto Iggy consumer groups. The gateway
assigns partitions
+to group members the way Kafka does, in the client, and polls every partition
by explicit offset.
+Iggy's group registry is used as an offset key and for nothing else, which
+[`OFFSET_STORAGE.md`](OFFSET_STORAGE.md) covers.
+
+## Reserved header namespace
+
+`kafka.` is reserved on messages the gateway writes and reads. An Iggy
producer that sets a
+header in that namespace on a topic a Kafka consumer reads will have it
interpreted as gateway
+metadata.
+
+## Open questions
+
+Four questions need an answer before Produce
([#3535](https://github.com/apache/iggy/issues/3535))
+is written. Each one carries a default. If no answer lands by 2026-09-22, the
default is taken,
+this document is updated to record that it was decided by default, and the
work proceeds.
+
+### 1. Envelope fallback, or reject the record?
+
+A record that Iggy cannot hold natively goes into the envelope described
above. The alternative
+is to reject it with `MESSAGE_TOO_LARGE` (10), so that nothing an Iggy
consumer cannot read ever
+reaches a stream.
+
+Rejecting is the stricter guarantee and the worse compatibility story: a Kafka
producer that
+sends a 300-byte key works against a real broker and fails against the gateway.
+
+Default: keep the envelope.
+
+### 2. Is `kafka.` the right prefix?
+
+Every Kafka header name is stored as `kafka.h.<name>`, which spends 8 of the
255 bytes an Iggy
+header name has, on every header of every record. A shorter prefix buys those
bytes back and
+costs readability for anyone reading a stream by hand.
+
+Default: keep `kafka.`.
+
+### 3. Is the placeholder byte acceptable for tombstones?
+
+A null or empty value is stored as one `0x00` byte plus a `kafka.value` marker
header. The
+payload a native Iggy consumer sees is therefore a byte the producer never
sent.
+
+The alternative is the envelope, which costs a tombstone the fast path.
Tombstones are ordinary
+traffic on compacted Kafka topics, and Iggy has no compaction, so they are
stored and served
+like any other record.
+
+Default: keep the placeholder byte.
+
+### 4. Recompress on Fetch, or always emit uncompressed?
+
+Produce decompresses, and Fetch currently rebuilds an uncompressed batch.
Recompressing per
+topic costs CPU on the read path and saves bytes on the wire to the consumer.
+
+Default: always emit uncompressed, and revisit when a benchmark says it
matters.
+
+### Not asked here
+
+A Kafka `retention.ms` topic config could map onto Iggy's `message_expiry` at
creation time.
+`ensure_stream_and_topic` leaves topics on the server default, which never
expires. That belongs
+to CreateTopics ([#3538](https://github.com/apache/iggy/issues/3538)), which
owns topic
+configuration, rather than to the record mapping.
+
+## References
+
+- Scope and phases: [`SCOPE.md`](SCOPE.md)
+- Bridge API: `gateways/kafka/src/bridge/iggy_bridge.rs`
+- Iggy message limits: `core/common/src/types/message/iggy_message.rs`,
+ `core/common/src/types/message/user_headers.rs`
+- Produce confirmations:
`core/binary_protocol/src/responses/messages/send_messages.rs`
diff --git a/gateways/kafka/docs/IDEMPOTENCE.md
b/gateways/kafka/docs/IDEMPOTENCE.md
new file mode 100644
index 000000000..3ff1f7f31
--- /dev/null
+++ b/gateways/kafka/docs/IDEMPOTENCE.md
@@ -0,0 +1,179 @@
+# 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
+`kafka-console-producer.sh`.
+
+## The problem
+
+A stock Java producer sets `enable.idempotence=true` without being asked. That
default arrived
+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.
+
+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.
+
+## What Iggy already deduplicates
+
+Iggy deduplicates on the partition plane, and did so before this gateway
existed. The key is
+`(client_id, user_id, request)`. `client_id` is a `u128` the client generates,
and `request` is a
+monotonic counter on the client's session. Each partition group holds a
per-client watermark with
+a 128-bit `committed_window` under it. A request below the watermark with its
bit clear is a
+reordered arrival still to execute. A request below the window reads as
committed.
+`partition.dedup_clients_max` caps the distinct clients one partition tracks,
at 4096 by default
+and 65536 at the ceiling.
+
+The shape is closer to Kafka's than it looks. A producer id with a
per-partition sequence maps
+onto a session id with a request id. The window is also deeper than the five
batches a Kafka
+producer keeps in flight.
+
+What breaks the mapping is the hop each side protects. Iggy's covers gateway
to Iggy. Kafka's
+covers producer to gateway. A retrying producer sends a fresh Produce request,
and the gateway
+turns it into a fresh Iggy send carrying a fresh request id. Nothing
recognises the replay.
+
+## Closing that gap
+
+Both halves of Iggy's key have to come from the Kafka request rather than from
the SDK.
+
+The client id is not a per-request field. `dispatch_partition_request` takes
it from the
+transport's bound session, and the header transmute at
+`core/server/src/dispatch/partition.rs:247` overwrites whatever the SDK sent.
Deriving it from a
+Kafka producer id therefore needs one bound session per producer.
+
+The request id is different. The transmute copies the header and replaces only
`group`, `client`,
+`session` and `user_id`, so a caller-supplied request id survives to the
`ClientTable`.
+
+One session per producer is a connection pool keyed by producer id. That is
the pool the
+README's "Concurrency ceiling" section already owes before
+[#3535](https://github.com/apache/iggy/issues/3535) and
+[#3536](https://github.com/apache/iggy/issues/3536).
+
+Two SDK seams are missing for it:
+
+- a caller-chosen client id at build time. `ConsensusSession::with_client_id`
is public, but both
+ construction sites build with `ConsensusSession::new()`
+ (`core/sdk/src/tcp/tcp_client.rs:378` and `:583`)
+- a caller-supplied request id. `send_raw_with_response` already takes
+ `preencoded: Option<RequestHeader>` for transient replay, and it is private
+
+None of this lands in this phase. Produce and Fetch do not need it, delivery
is at-least-once
+before and after, and the additions belong to whoever owns that SDK surface.
+
+### Restart
+
+A session registers with a fresh random client id per gateway process. A
restarted gateway
+therefore cannot collide with a watermark that outlived it, and a retry that
spans a restart is
+not deduplicated. That is still at-least-once, which the README states plainly.
+
+The alternative is a stable client id derived from the producer id. It
deduplicates across a
+restart, and it is unsafe without also persisting the last request id per
producer. Request ids
+restarting at 1 under a live watermark read as duplicates, so fresh writes are
discarded. Silent
+loss is worse than duplicate delivery, so the fresh random id wins.
+
+### Confirmed and not
+
+The hop mismatch, one session per producer, and the restart choice are agreed
in the maintainer
+thread on [#3545](https://github.com/apache/iggy/issues/3545) and in Discord.
+
+Two further readings are ours and are not confirmed yet. One session per
producer is enough,
+rather than one per producer and partition, because the `ClientTable` is per
partition group. The
+same request number on two partitions is two entries, so `request_id =
base_sequence + 1` stays
+monotonic inside each. And the gateway has to preserve per-partition ordering
per producer,
+because sequences advance by record count. A lower sequence arriving after a
higher one lands
+below the watermark, outside the window, and reads as committed.
+
+## Options
+
+| Option | Cost | What a stock producer does |
+| -------- | ------ | ---------------------------- |
+| Stub with `UNSUPPORTED_VERSION` | none | fails at startup unless the user
sets `enable.idempotence=false` |
+| Allocate only | about a day | works untouched, at-least-once delivery |
+| Allocate and pool | weeks, blocked on the SDK | works untouched, retries
absorbed on both hops |
+
+## Decision
+
+Allocate only, and defer the pool.
+
+Rejecting the stock producer to avoid the pool trades away the one requirement
the maintainers
+named. It buys a guarantee nobody is asking for yet. Allocating costs about a
day, leaves
+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.
+
+InitProducerId with no `transactional_id`:
+
+- allocate the next producer id, return it with epoch 0 and error code 0
+- build the id from an instance number in the high 16 bits and a counter in
the low 47 bits,
+ leaving bit 63 clear. `producer_id` is an `i64` and `-1` means no producer
id, so the value has
+ to stay non-negative. That leaves room for 65536 instances holding 140
trillion ids each
+
+The instance number comes from configuration (`IGGY_KAFKA_INSTANCE_ID`,
default 0), not from a
+draw at startup. A random 16-bit number collides with even odds at around 300
instances. That is
+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.
+
+InitProducerId with a `transactional_id`:
+
+- answer `UNSUPPORTED_VERSION` (35), unchanged. 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.
+
+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.
+
+## What this does not give you
+
+A producer that holds an id believes its retries are deduplicated. They are
not. A retry after a
+network timeout writes the record twice, and both copies reach the stream with
their own offsets.
+
+Iggy's own deduplication does not help, because it guards the other hop.
Delivery through the
+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
+
+Allocate only, as above, or stub and document `enable.idempotence=false`?
+
+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.
+
+## References
+
+- Record mapping: [`BRIDGE_MAPPING.md`](BRIDGE_MAPPING.md), batch-level fields
+- Scope and phases: [`SCOPE.md`](SCOPE.md)
+- Version firewall: `src/protocol/api.rs`, `SUPPORTED_RANGES`
+- 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
diff --git a/gateways/kafka/docs/OFFSET_STORAGE.md
b/gateways/kafka/docs/OFFSET_STORAGE.md
new file mode 100644
index 000000000..cd8421f6d
--- /dev/null
+++ b/gateways/kafka/docs/OFFSET_STORAGE.md
@@ -0,0 +1,163 @@
+# Consumer offset storage
+
+Status: proposed. Answers [#3540](https://github.com/apache/iggy/issues/3540)
and blocks
+[#3542](https://github.com/apache/iggy/issues/3542), OffsetCommit and
OffsetFetch.
+
+## Decision
+
+Store Kafka group offsets as Iggy consumer offsets, one key per partition,
under a consumer
+group whose name is derived from the Kafka group id.
+
+The issue lists three options. None of them is this one.
+
+| Option | Why not |
+| -------- | --------- |
+| A, an Iggy-backed `__consumer_offsets` topic | Rebuilds what Iggy already
has. A compacted offset topic needs compaction, which Iggy does not have, so
the gateway replays the whole topic at every startup |
+| B, a SQLite file on the gateway host | A second durability story, a second
backup story, and offsets that do not survive moving the gateway |
+| C, in memory only | Fails the acceptance criterion in #3542, which is that
offsets survive a restart |
+
+Iggy already stores a durable offset per consumer and per partition,
replicated with the
+partition itself. Using it costs one call per partition on commit and one on
fetch.
+
+## The key
+
+One Iggy consumer offset per Kafka `(group, topic, partition)`.
+
+- consumer kind: `ConsumerKind::ConsumerGroup`
+- consumer id: `Identifier::named("kafka.cg.<group>")`
+- stream and topic: whatever `TopicMapping` resolves the Kafka topic to
+- partition: the Kafka partition index, unchanged, because both sides number
from 0
+
+The gateway calls `create_consumer_group(stream, topic, "kafka.cg.<group>")`
before the first
+commit for a group on a topic. If the group does not resolve in metadata, the
server rejects the
+offset write. The group has to exist first. The gateway never joins the group.
Offsets are
+readable by any client, member or not.
+
+### Why the group kind and not a named consumer
+
+`ConsumerKind::Consumer` with a name looks simpler, because it needs no
registration call. It is
+not. The server hashes a named consumer id to a `u32` with `XxHash32`
+(`core/server/src/dispatch/partition.rs:916`), and that hash is the offset
key. Two different
+group names can collide and silently share one offset.
+
+A consumer group name resolves through metadata to a monotonic id instead. No
hash, no
+collision, and `get_consumer_groups(stream, topic)` lists what exists.
+
+### Why the prefix
+
+`kafka.cg.` keeps a Kafka group called `orders` off the key that a native Iggy
consumer group
+called `orders` uses. Without it the two share an offset and each one moves
the other.
+
+The prefix does not make the offsets safe to poll with. That is the next
section.
+
+## What is stored
+
+The Kafka committed offset, verbatim, with no conversion.
+
+The two systems mean different things by the number. Kafka commits the next
offset to read.
+Iggy stores the last offset processed, and `PollingKind::Next` resumes at the
stored value plus
+one (`core/partitions/src/iggy_partition.rs:3835`). A Kafka offset stored in
an Iggy key is
+therefore one greater than Iggy's own convention for that key.
+
+This is inert because the gateway never polls that way. Fetch always polls
with an explicit
+offset, `PollingKind::Offset`, taken from the Kafka request. Nothing in the
gateway reads the
+stored value to decide where to resume. It is returned to the client on
OffsetFetch and
+otherwise untouched.
+
+The rule this creates: no code path polls a `kafka.cg.*` key with
`PollingKind::Next`. Doing so
+skips one record per partition. The prefix is what keeps a native Iggy
consumer from doing it by
+accident.
+
+Converting on write instead, and storing the Kafka offset minus one, breaks at
offset 0. It also
+makes an empty commit look the same as a commit of the first record. Storing
verbatim is the
+smaller problem.
+
+### A commit of -1
+
+Kafka does not validate the sign of a committed offset.
`OffsetMetadataManager` checks the
+metadata length and nothing else, so a real broker stores `-1` and hands it
back on the next
+OffsetFetch, where a consumer reads it as no committed offset.
+
+Iggy cannot store that value, because `store_consumer_offset` takes a `u64`. A
commit of `-1`
+therefore calls `delete_consumer_offset` on the key. What a client can observe
is the same: the
+next OffsetFetch finds nothing and the gateway answers `-1`.
+
+Deleting a key that is not there returns `ConsumerOffsetNotFound` (3021). On
this path that
+counts as success, because the client asked for the offset to be absent and it
is absent. Any
+other negative offset is rejected with `OFFSET_OUT_OF_RANGE` (1).
+
+## OffsetFetch with no topics named
+
+OffsetFetch v2 and later let a client pass a null topic list, which asks for
every offset the
+group holds. `kafka-consumer-groups.sh --describe` does this.
+
+Iggy has no lookup by consumer. Offsets are read one partition at a time
+(`core/common/src/traits/consumer_offset_client.rs:41`). The gateway answers
by enumerating the
+topics in the mapped stream and querying each partition of each one.
+
+Answering also means naming each topic the way the client named it, and
`TopicMapping::resolve`
+runs the wrong way. Its own doc says it is not injective, so it cannot be
inverted in general.
+Two cases divide it. A topic with no override resolves to `(default_stream,
kafka_topic)`, so the
+Iggy topic name is the Kafka name and reverses for free. A topic with an
override needs a reverse
+index, built once at config load, which `TopicMapping::new` already makes safe
by rejecting two
+overrides that share a target.
+
+What neither case covers is an unlisted Kafka topic whose name collides with
an override's target
+inside the default stream. `new` rejects the shapes it can check, but the
space of unlisted names
+is unbounded, so the collision is disclosed rather than prevented. An offset
under a colliding
+name is reported against whichever Kafka name the reverse index holds.
+
+That is one round trip per partition on an admin call. The cost is bounded by
the topic and
+partition count of one stream. This is an admin path and not a data path, so
the cost is
+acceptable. It is written here so nobody discovers it in a test.
+
+## What is dropped
+
+Kafka lets a client attach a metadata string to a commit. Iggy stores a number
and nothing else.
+The string is dropped on commit, and OffsetFetch returns an empty string.
+
+`committed_leader_epoch` is dropped the same way. The gateway reports `-1`.
+
+## Limits
+
+A partition admits a bounded number of distinct offset keys per consumer kind.
The default is
+4096, set by `partition.consumer_offsets_max`
+(`core/configs/src/server_config/partition.rs:151`). The configurable ceiling
is 262144. Passing
+the limit returns `TooManyConsumerOffsets` (3024). Kafka has no error code for
this condition, so
+it maps to `UNKNOWN_SERVER_ERROR`, which is what `bridge/error.rs` already
does where no honest
+code exists.
+
+Consumer groups and plain consumers count against separate limits, so Kafka
groups do not
+compete with native Iggy consumers for the same 4096.
+
+A client cannot act on that error, so the operator has to.
`partition.consumer_offsets_max` is
+named in the gateway README for that reason, and a handler that hits the limit
logs the Iggy
+error at `error!` level, which is what `bridge/error.rs` already asks handlers
to do wherever the
+Kafka code it sends is less specific than the Iggy error it received.
+
+An Iggy name is capped at 255 bytes (`core/common/src/lib.rs:168`), which
leaves 246 for a Kafka
+group id after the prefix. A longer group id is rejected with
`INVALID_GROUP_ID` (24).
+
+## More than one gateway instance
+
+Two gateway instances that share an Iggy cluster read and write the same
offset keys. The key
+comes from the Kafka group id and nothing else. Two instances that serve one
group therefore
+agree on committed offsets without talking to each other.
+
+They do not agree on group membership. That belongs to the coordinator
+([#3541](https://github.com/apache/iggy/issues/3541)) and is not settled here.
+
+## Open question
+
+Iggy consumer offsets keyed by group, as above, or one of A, B and C from the
issue?
+
+If no answer lands by 2026-09-22, the design above is taken and the work
proceeds. This document
+is then updated to record that it was decided by default.
+
+## References
+
+- Record mapping: [`BRIDGE_MAPPING.md`](BRIDGE_MAPPING.md)
+- Scope and phases: [`SCOPE.md`](SCOPE.md)
+- Offset API: `core/common/src/traits/consumer_offset_client.rs`
+- Group API: `core/common/src/traits/consumer_group_client.rs`
+- Offset key resolution: `core/server/src/dispatch/partition.rs`
diff --git a/gateways/kafka/docs/SCOPE.md b/gateways/kafka/docs/SCOPE.md
index de7a343de..30289a934 100644
--- a/gateways/kafka/docs/SCOPE.md
+++ b/gateways/kafka/docs/SCOPE.md
@@ -109,10 +109,10 @@ below it are still open for the issues that build on top
of it.
not part of `bridge/`'s own scope.
- [x] Idempotent `ensure_stream_and_topic()` (create-if-not-exists) -
`src/bridge/iggy_bridge.rs`,
exercised end-to-end in `tests/bridge_iggy_integration_tests.rs`.
-- [ ] Document partition mapping in `docs/BRIDGE_MAPPING.md`:
+- [x] Document partition mapping in [`BRIDGE_MAPPING.md`](BRIDGE_MAPPING.md):
- Iggy partitions are **0-based** (same as Kafka) — direct `partition_id`
mapping, no offset conversion
- - Iggy **consumer groups exist** — map Kafka group APIs to Iggy consumer
group APIs
- - Use `Partitioning::balanced()` only when Kafka sends `partition == -1`;
otherwise use request partition ID
+ - Kafka consumer groups do **not** map onto Iggy consumer groups. Assignment
stays client-side, and Iggy's group registry is used as an offset key only
([`OFFSET_STORAGE.md`](OFFSET_STORAGE.md))
+ - `Partitioning::partition_id(index)` on every Produce. A Kafka producer
resolves the partition before it builds the request, so
`Partitioning::balanced()` has no trigger there. The `-1`
default-partition-count case belongs to CreateTopics
- [ ] Real Metadata topology (brokers, partitions, leaders) backed by Iggy
state
### `kafka-protocol` crate adoption — superseded, done differently
@@ -130,12 +130,18 @@ above).
### Phase 3 — Consumer groups (~7 API keys)
+Offset persistence design
([#3540](https://github.com/apache/iggy/issues/3540)):
+[`OFFSET_STORAGE.md`](OFFSET_STORAGE.md).
+
- [ ] OffsetCommit (8), OffsetFetch (9), FindCoordinator (10)
- [ ] JoinGroup (11), Heartbeat (12), LeaveGroup (13), SyncGroup (14)
- [ ] DescribeGroups (15), ListGroups (16) as needed by target clients
### Phase 3+ — Auth, admin, tuning
+InitProducerId and idempotent producers
+([#3545](https://github.com/apache/iggy/issues/3545)):
[`IDEMPOTENCE.md`](IDEMPOTENCE.md).
+
- [ ] SASL (17, 36) if required by deployment
- [ ] Tune `max_frame_size` per workload (Kafka defaults: ~1 MiB produce, ~50
MiB fetch; current default 8 MiB)
- [ ] Target **~15–20 API keys** total for a functional bridge — not all 74+
admin keys