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"
     );
 }
 

Reply via email to