This is an automated email from the ASF dual-hosted git repository.
numinnex pushed a commit to branch kafka_idempotence_stub
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/kafka_idempotence_stub by this
push:
new baf56cbe1 address review
baf56cbe1 is described below
commit baf56cbe1d82510d6bd97586ae8419fc396f71c4
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Wed Sep 23 14:29:13 2026 +0200
address review
---
gateways/kafka/docs/IDEMPOTENCE.md | 23 ++++++++++----
gateways/kafka/docs/MANUAL_TESTING.md | 2 +-
gateways/kafka/docs/SCOPE.md | 3 +-
gateways/kafka/src/protocol/api.rs | 4 +--
.../src/protocol/handlers/init_producer_id.rs | 36 ++++++++++++++++++++--
gateways/kafka/src/protocol/handlers/produce.rs | 9 +++++-
gateways/kafka/tests/idempotence_tests.rs | 13 ++++----
7 files changed, 70 insertions(+), 20 deletions(-)
diff --git a/gateways/kafka/docs/IDEMPOTENCE.md
b/gateways/kafka/docs/IDEMPOTENCE.md
index 2d135739a..ea8f48a23 100644
--- a/gateways/kafka/docs/IDEMPOTENCE.md
+++ b/gateways/kafka/docs/IDEMPOTENCE.md
@@ -124,11 +124,13 @@ 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 while nothing keys
state on it,
-which is true of allocate-only: the gateway hands an id out and forgets it.
That stops being true
-the moment the pool lands or Produce persists, because a restarted allocator
replays ids a live
-producer still holds, and `producer_epoch` is always 0 so the pair cannot tell
the generations
-apart. Treat generation reuse as a pool blocker, not a detail.
+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. Nothing is persisted for
this, and 2^47
+milliseconds leaves the counter space thousands of years from running out.
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
@@ -162,7 +164,9 @@ 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)
- reject a non-empty `transactional_id` (v3+) with `UNSUPPORTED_VERSION` (35),
at the partition
- level, keeping the connection open and keeping the `acks=0` silence rule
+ 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
@@ -181,6 +185,13 @@ Without this guard a transactional batch would land as
ordinary records once
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
diff --git a/gateways/kafka/docs/MANUAL_TESTING.md
b/gateways/kafka/docs/MANUAL_TESTING.md
index b77d864ee..6a561362e 100644
--- a/gateways/kafka/docs/MANUAL_TESTING.md
+++ b/gateways/kafka/docs/MANUAL_TESTING.md
@@ -192,7 +192,7 @@ Record kcat version and exact error strings in your test
log. G1 passing is the
| 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 at all |
+| 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 |
---
diff --git a/gateways/kafka/docs/SCOPE.md b/gateways/kafka/docs/SCOPE.md
index d5084c1f9..a31930a0f 100644
--- a/gateways/kafka/docs/SCOPE.md
+++ b/gateways/kafka/docs/SCOPE.md
@@ -107,7 +107,8 @@ Three things enforce that, in the order a client meets them:
[`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. `acks=0` stays silent, and no case closes the connection.
+ 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).
diff --git a/gateways/kafka/src/protocol/api.rs
b/gateways/kafka/src/protocol/api.rs
index 087fc4d18..217b7527b 100644
--- a/gateways/kafka/src/protocol/api.rs
+++ b/gateways/kafka/src/protocol/api.rs
@@ -176,7 +176,7 @@ pub struct GatewayState {
impl GatewayState {
#[must_use]
- pub const fn new(
+ pub fn new(
broker: BrokerAdvertise,
bridge: Option<Arc<IggyBridge>>,
max_frame_size: usize,
@@ -192,7 +192,7 @@ impl GatewayState {
/// State with no bridge, so every handler takes its stub path.
#[must_use]
- pub const fn stub(broker: BrokerAdvertise, max_frame_size: usize) -> Self {
+ pub fn stub(broker: BrokerAdvertise, max_frame_size: usize) -> Self {
Self::new(broker, None, max_frame_size, 0)
}
}
diff --git a/gateways/kafka/src/protocol/handlers/init_producer_id.rs
b/gateways/kafka/src/protocol/handlers/init_producer_id.rs
index 89ca2763e..979bcde54 100644
--- a/gateways/kafka/src/protocol/handlers/init_producer_id.rs
+++ b/gateways/kafka/src/protocol/handlers/init_producer_id.rs
@@ -23,6 +23,7 @@
//! 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};
@@ -57,6 +58,13 @@ const PRODUCER_EPOCH: i16 = 0;
///
/// `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, or the clock stepped back across
the restart, with
+/// nothing persisted.
#[derive(Debug)]
pub struct ProducerIdAllocator {
instance_id: u16,
@@ -65,10 +73,10 @@ pub struct ProducerIdAllocator {
impl ProducerIdAllocator {
#[must_use]
- pub const fn new(instance_id: u16) -> Self {
+ pub fn new(instance_id: u16) -> Self {
Self {
instance_id,
- next_counter: AtomicU64::new(0),
+ next_counter: AtomicU64::new(clock_counter()),
}
}
@@ -132,6 +140,15 @@ pub fn encode_response(
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)
@@ -143,6 +160,8 @@ fn encode_inner(version: i16, error_code: i16, producer_id:
i64) -> Result<Bytes
#[cfg(test)]
mod tests {
use std::sync::atomic::Ordering;
+ use std::thread;
+ use std::time::Duration;
use super::{COUNTER_BITS, MAX_COUNTER, ProducerIdAllocator};
@@ -162,6 +181,19 @@ mod tests {
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);
diff --git a/gateways/kafka/src/protocol/handlers/produce.rs
b/gateways/kafka/src/protocol/handlers/produce.rs
index 8e7a1a2f4..b5264c90d 100644
--- a/gateways/kafka/src/protocol/handlers/produce.rs
+++ b/gateways/kafka/src/protocol/handlers/produce.rs
@@ -80,6 +80,12 @@ pub async fn handle(state: &GatewayState, api_version: i16,
body: Bytes) -> Hand
match decode_guarded::<ProduceRequest>(api_version, body, |v, b| {
validate_produce_shape(v, b, state.max_frame_size)
}) {
+ // The transactional refusal runs before the acks=0 silence, so no
write path added
+ // after it can store a transactional batch. acks=0 has no response to
carry 35, so the
+ // refusal is a close, which is what a Kafka broker does on an acks=0
produce error.
+ Ok(req) if req.acks == 0 &&
is_transactional(req.transactional_id.as_ref()) => {
+ HandleOutcome::Close
+ }
// acks=0 is fire-and-forget: the client isn't reading a response, so
// sending one desyncs the next correlation id it expects.
Ok(req) if req.acks == 0 => HandleOutcome::NoResponse,
@@ -152,7 +158,8 @@ pub fn encode_response(version: i16, req: &ProduceRequest)
-> Result<Bytes> {
/// transactional bit in attributes, and the records stay opaque bytes here,
so a hand-built
/// frame setting the bit without the request field still gets the retriable
stub error. Java
/// and librdkafka both set the request field whenever they set the batch bit,
so no real client
-/// reaches that gap; it has to close before records are ever persisted.
+/// reaches that gap. The stub does not parse batches to close it: the persist
path (#3535)
+/// decodes every batch anyway and has to refuse the bit there.
fn partition_error_code(req: &ProduceRequest) -> i16 {
if is_transactional(req.transactional_id.as_ref()) {
ERROR_UNSUPPORTED_VERSION
diff --git a/gateways/kafka/tests/idempotence_tests.rs
b/gateways/kafka/tests/idempotence_tests.rs
index 2a03b7e67..86964e70f 100644
--- a/gateways/kafka/tests/idempotence_tests.rs
+++ b/gateways/kafka/tests/idempotence_tests.rs
@@ -237,11 +237,10 @@ async fn
given_no_transactional_id_when_producing_should_keep_the_retriable_stub
}
#[tokio::test]
-async fn
given_acks_zero_and_a_transactional_id_when_producing_should_stay_silent() {
- // Paired on purpose. Asserting silence alone passes even with the
transactional guard gone,
- // because acks=0 returns before the guard ever runs, so the assertion
would hold over a
- // gateway that refuses nothing. The acks=1 half establishes that this
body is refused at all,
- // which is what makes the acks=0 half a statement about suppressing a
real error.
+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,
@@ -264,8 +263,8 @@ async fn
given_acks_zero_and_a_transactional_id_when_producing_should_stay_silen
)
.await;
assert!(
- outcome.is_no_response(),
- "acks=0 is fire-and-forget: refusing a transaction must not put a
frame on the wire"
+ outcome.is_close(),
+ "acks=0 cannot carry 35, so a transactional batch must be refused by
closing"
);
}