This is an automated email from the ASF dual-hosted git repository.

krishvishal pushed a commit to branch kafka-produce
in repository https://gitbox.apache.org/repos/asf/iggy.git


The following commit(s) were added to refs/heads/kafka-produce by this push:
     new e6b7466b6 fix: address review comments
e6b7466b6 is described below

commit e6b7466b6a5ab7f72a54da0416a4ef9864bde588
Author: Krishna Vishal <[email protected]>
AuthorDate: Sun Sep 27 08:50:59 2026 +0530

    fix: address review comments
---
 gateways/kafka/README.md                          |  17 +-
 gateways/kafka/docs/BRIDGE_MAPPING.md             |  43 +-
 gateways/kafka/docs/MANUAL_TESTING.md             |   4 +-
 gateways/kafka/docs/SCOPE.md                      |   6 +-
 gateways/kafka/src/protocol/api.rs                |   7 +-
 gateways/kafka/src/protocol/handlers/produce.rs   | 247 ++++------
 gateways/kafka/src/records.rs                     | 575 +++++++++++++++-------
 gateways/kafka/tests/produce_real_bridge_tests.rs |  39 +-
 8 files changed, 550 insertions(+), 388 deletions(-)

diff --git a/gateways/kafka/README.md b/gateways/kafka/README.md
index d08969b83..b718d910d 100644
--- a/gateways/kafka/README.md
+++ b/gateways/kafka/README.md
@@ -99,13 +99,12 @@ start a real `iggy-server`.
 
 ### Produce ([#3535](https://github.com/apache/iggy/issues/3535))
 
-One partition, one Iggy send, or more if its timestamps span over about 71 
min. Each partition
-answers for itself.
+One partition, one Iggy send. Each partition answers for itself.
 
 | Field | Gateway |
 | ----- | ------- |
 | Partition | Index as sent. Both count from 0. |
-| Base offset | From the send confirmation. `-1` if none, or if another writer 
wrote between split sends. |
+| Base offset | From the send confirmation. `-1` if none. |
 | `log_start_offset` | Always `-1`. |
 | `acks` | `0`, `1`, `-1` write the same. Other values: 21, before any 
per-partition check. |
 | `acks=0` | Writes, answers nothing. Any failed partition closes the 
connection. |
@@ -113,8 +112,9 @@ answers for itself.
 | `timeout_ms` | Honored, max 20 s. Past it: 7. |
 | Compression | gzip, snappy, lz4. zstd from v7, else 76. |
 | Producer id, epoch, sequence | Ignored, so a retry writes twice. |
-| Timestamps over about 71 min apart | Several sends. A failure after the 
first can duplicate on retry. |
+| Timestamps over about 71 min apart | Clamped into one send. `kafka.ts` keeps 
the real one. |
 | Several batches in one partition | 87, as Kafka. |
+| Bytes after the batch | 87. |
 | More records than the batch declares | 87, as Kafka. |
 | Same partition twice in one request | Written twice. Kafka keeps the last. |
 | Repeated header name in a record | 87. |
@@ -126,11 +126,14 @@ answers for itself.
 | 10 | Record, send or partition too large, even alone | Java splits 
multi-record batches. Else fails. |
 | 87 | Record the gateway cannot map. Reason in `error_message` from v8. | 
Fails. |
 | 35 | Transactional or control batch | Fails. |
-| 6 | Earlier partitions used the request budget, and this one fits alone. 
Nothing written. | Retries. |
+| 6 | Request budget ran out. Nothing written. | Retries. |
 | 7 | Deadline passed, or connection lost mid-send. May be written. | Retries. 
Can duplicate. |
 
-Request budget: `max_frame_size` decompressed bytes, `max_frame_size / 64` 
record slots, 3 headers
-per slot. 4 requests decode at once.
+Partition cap: `max_frame_size` decompressed bytes, `max_frame_size / 64` 
record slots, 3 headers
+per slot. Past it: 10.
+
+Request budget: 8 partition caps, refused partitions included. Past it: 6 for 
the rest, not
+decoded. 4 requests decode at once.
 
 [docs/BRIDGE_MAPPING.md](docs/BRIDGE_MAPPING.md) describes what a record 
becomes once it is stored.
 
diff --git a/gateways/kafka/docs/BRIDGE_MAPPING.md 
b/gateways/kafka/docs/BRIDGE_MAPPING.md
index 99079bea6..8ae3cac77 100644
--- a/gateways/kafka/docs/BRIDGE_MAPPING.md
+++ b/gateways/kafka/docs/BRIDGE_MAPPING.md
@@ -39,7 +39,7 @@ Produce, per record:
 | 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 timestamp (CreateTime), milliseconds | `origin_timestamp`, 
microseconds, plus `kafka.ts` when that cannot hold it |
 | record offset | partition offset, assigned by Iggy |
 | partition index | partition index, both 0-based |
 | topic | stream and topic per `TopicMapping` |
@@ -100,28 +100,25 @@ A record that arrived through Produce survives the round 
trip exactly, because i
 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.
-
-A record stamped at exactly `0` ms, the Unix epoch, is a different record that 
stores the same
-`0`. It carries a `kafka.ts` header holding `epoch` to say so, and Fetch reads 
that header before
-it reads the origin timestamp. Without it a producer that stamps a record `0` 
gets the server's
-clock back instead, and never learns that the value changed.
-
-One Iggy batch holds timestamps that span at most 
`MAX_TIMESTAMP_DELTA_MICROS`, which is
-`u32::MAX` microseconds, about 71.6 minutes 
(`core/binary_protocol/src/batch.rs:55`). The send
-encoder stores each message as a `u32` delta from the batch minimum and 
refuses a larger one
-(`core/binary_protocol/src/requests/messages/send_messages.rs:133`). Kafka 
puts no such bound on
-one produce batch, so two shapes fail:
-
-- a batch whose CreateTime values span more than 71.6 minutes, which replay 
and mirror producers
-  reach
-- a batch that mixes a record with no timestamp, stored as `0`, with normally 
stamped records
-
-Neither is fixable in the record mapping, which sees one record and has no 
batch minimum to work
-from. Produce ([#3535](https://github.com/apache/iggy/issues/3535)) owns 
splitting a produce
-batch into sends the server accepts.
+Kafka sends `-1` for a record with no timestamp. That is stored as `0`, and 
Fetch reads a zero
+origin timestamp as the server-assigned one. A real broker does the same under 
`LogAppendTime`.
+
+One Iggy send holds origin timestamps at most `u32::MAX` µs apart, about 71.6 
min
+(`core/binary_protocol/src/batch.rs:55`). Kafka has no such bound. So each 
produce batch gets one
+window, from its earliest real timestamp. A timestamp outside it is clamped 
in, and the `kafka.ts`
+header (`Int64`, ms) keeps the real one:
+
+| Record timestamp | `origin_timestamp` | `kafka.ts` |
+| ---------------- | ------------------ | ---------- |
+| In the window | timestamp × 1000 | none |
+| After the window | window end | timestamp |
+| `-1`, no real timestamp in the batch | `0` | none |
+| `-1`, next to real timestamps | window start | `-1` |
+| `0`, the epoch | window start, or `0` | `0` |
+
+Fetch reads `kafka.ts` first. `-1` there means the server-assigned timestamp. 
An Iggy client sees
+the clamped `origin_timestamp`. Only a batch over 71.6 min, or one mixing `-1` 
with real
+timestamps, has one.
 
 ## Records Iggy cannot hold natively
 
diff --git a/gateways/kafka/docs/MANUAL_TESTING.md 
b/gateways/kafka/docs/MANUAL_TESTING.md
index 0e94f643a..7baf3a8c2 100644
--- a/gateways/kafka/docs/MANUAL_TESTING.md
+++ b/gateways/kafka/docs/MANUAL_TESTING.md
@@ -203,7 +203,7 @@ Record kcat version and exact error strings in your test 
log. G1 passing is the
 | -1 | UNKNOWN_SERVER_ERROR | Produce with a bridge: Iggy error with no closer 
code, or bad bridge login |
 | 0 | NONE | Fetch top-level error field only (`ec=0` there does not mean 
per-partition success - see A6) |
 | 3 | UNKNOWN_TOPIC_OR_PARTITION | Metadata stub, per topic. Produce with a 
bridge: missing topic or partition |
-| 6 | NOT_LEADER_OR_FOLLOWER | Produce/Fetch/ListOffsets stub (not stored). 
Produce with a bridge: Iggy unreachable, or request budget used by earlier 
partitions and the entry fits alone |
+| 6 | NOT_LEADER_OR_FOLLOWER | Produce/Fetch/ListOffsets stub (not stored). 
Produce with a bridge: Iggy unreachable, or the request budget ran out |
 | 7 | REQUEST_TIMED_OUT | Produce with a bridge: deadline passed, or 
connection lost mid-send (may be stored) |
 | 10 | MESSAGE_TOO_LARGE | Produce with a bridge: record, send or partition 
too large, even alone |
 | 17 | INVALID_TOPIC_EXCEPTION | Produce with a bridge: bad topic name |
@@ -215,7 +215,7 @@ Record kcat version and exact error strings in your test 
log. G1 passing is the
 | 41 | NOT_CONTROLLER | CreateTopics stub (topic not created) |
 | 42 | INVALID_REQUEST | Fetch/ListOffsets/CreateTopics/ApiVersions decode 
failure. Not Produce: it closes the connection (H1) |
 | 76 | UNSUPPORTED_COMPRESSION_TYPE | Produce with a bridge: zstd before v7 |
-| 87 | INVALID_RECORD | Produce with a bridge: a record batch this gateway 
cannot map, including a missing, empty or second one |
+| 87 | INVALID_RECORD | Produce with a bridge: a record batch this gateway 
cannot map, including a missing, empty or second one, or bytes after it |
 
 A malformed request header (before any API-specific body is even reached) has 
no parsed header to
 build a version-correct response against, so it closes the connection rather 
than returning any
diff --git a/gateways/kafka/docs/SCOPE.md b/gateways/kafka/docs/SCOPE.md
index 29cb5d119..53d831da6 100644
--- a/gateways/kafka/docs/SCOPE.md
+++ b/gateways/kafka/docs/SCOPE.md
@@ -44,7 +44,7 @@ it knows the server supports flexible encoding.
 | --------- | ------ | ------------- | ------------- | ---------------- | 
---------- |
 | 18 | ApiVersions | 0 | 3 | 0, 1, 2, 3 | Advertise supported ranges; flexible 
encoding at v3+ |
 | 3 | Metadata | 0 | 9 | 0, 1, 2, 3, 4, 5, 6, 7, 8, 9 | Decode topic list 
count; stub broker host from `advertised_host` or the bound `local_addr` IP; 
flexible encoding at v9+ |
-| 0 | Produce | 3 | 9 | 3, 4, 5, 6, 7, 8, 9 | With a bridge: one 
`send_messages` per partition, more for timestamps over ~71 min apart. Without 
one: stub returns `NOT_LEADER_OR_FOLLOWER` (6) |
+| 0 | Produce | 3 | 9 | 3, 4, 5, 6, 7, 8, 9 | With a bridge: one 
`send_messages` per partition. Without one: stub returns 
`NOT_LEADER_OR_FOLLOWER` (6) |
 | 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+ |
@@ -105,8 +105,8 @@ below it are still open for the issues that build on top of 
it.
 - [x] Add `bridge/` module (`iggy_bridge`) - connection lifecycle, topic 
mapping, provisioning,
       high watermark, error mapping. See 
[README.md](../README.md#iggy-bridge-3533).
 - [x] Produce → `send_messages` 
([#3535](https://github.com/apache/iggy/issues/3535)) - one call
-      per partition (more for timestamps over ~71 min apart), base offset from 
the send
-      confirmation, one error code per partition. See 
[README.md](../README.md#produce-3535).
+      per partition, base offset from the send confirmation, one error code 
per partition. See
+      [README.md](../README.md#produce-3535).
 - [ ] Produce: `IggyClient` pool. Pin each partition to one client, so order 
holds.
 - [ ] Produce: write keyed records' header TLVs into one buffer 
(`records::to_iggy`). Benchmark
       first. Keep Iggy's TLV layout and the 100 KB header check.
diff --git a/gateways/kafka/src/protocol/api.rs 
b/gateways/kafka/src/protocol/api.rs
index 5c69e3e20..96e716e47 100644
--- a/gateways/kafka/src/protocol/api.rs
+++ b/gateways/kafka/src/protocol/api.rs
@@ -40,8 +40,8 @@ pub const DEFAULT_KAFKA_PORT: u16 = 9093;
 pub const ERROR_UNKNOWN_SERVER_ERROR: i16 = 
ResponseError::UnknownServerError.code();
 pub const ERROR_NONE: i16 = 0;
 pub const ERROR_UNKNOWN_TOPIC_OR_PARTITION: i16 = 
ResponseError::UnknownTopicOrPartition.code();
-/// Retriable, nothing written. The Produce stub, and a partition refused only 
because earlier
-/// partitions used the request budget.
+/// Retriable, nothing written. The Produce stub, and a partition refused 
because the request
+/// budget ran out.
 pub const ERROR_NOT_LEADER_OR_FOLLOWER: i16 = 
ResponseError::NotLeaderOrFollower.code();
 /// Outcome unknown: the write may have landed. Retriable, so a retry can 
duplicate it.
 /// Idempotent produce (#3545) closes that.
@@ -183,7 +183,8 @@ pub struct GatewayState {
     pub(crate) produce_slots: Semaphore,
 }
 
-/// Four request budgets, about 160 MB at the default 8 MiB frame. Sends run 
one at a time anyway.
+/// Each holds one decoded partition at a time, so about 160 MB at the default 
8 MiB frame. Sends
+/// run one at a time anyway.
 const PRODUCE_SLOTS: usize = 4;
 
 impl GatewayState {
diff --git a/gateways/kafka/src/protocol/handlers/produce.rs 
b/gateways/kafka/src/protocol/handlers/produce.rs
index 5f0d08a3e..9e6fe418e 100644
--- a/gateways/kafka/src/protocol/handlers/produce.rs
+++ b/gateways/kafka/src/protocol/handlers/produce.rs
@@ -17,7 +17,6 @@
 
 //! Produce (API key 0).
 
-use std::ops::Range;
 use std::time::Duration;
 
 use bytes::Bytes;
@@ -42,7 +41,8 @@ use crate::protocol::handlers::{
     decode_guarded, encode_message, respond_or_close, 
unsupported_version_response,
 };
 use crate::records::{
-    DecompressionBudget, RecordCodecError, Zstd, decode_batch, is_compressed, 
to_iggy,
+    Allowance, DecompressionBudget, RecordCodecError, TimestampWindow, Zstd, 
decode_batch,
+    is_compressed, to_iggy,
 };
 
 pub const RANGE: ApiVersionRange = ApiVersionRange {
@@ -71,17 +71,21 @@ const UNKNOWN_OFFSET: i64 = -1;
 /// to a wide topic, and a client set on hogging the shared bridge just sends 
more requests.
 const MAX_REQUEST_DEADLINE: Duration = Duration::from_secs(20);
 
-/// Frame bytes per record slot. Caps one 8 MiB request at 131,072 slots, 
about 40 MB.
+/// Frame bytes per record slot. Caps one partition at 131,072 slots, about 40 
MB, at the default
+/// 8 MiB frame.
 ///
 /// `GatewayState` caps how many requests hold that at once.
 const FRAME_BYTES_PER_RECORD: usize = 64;
 
+/// Partition allowances one request may decode, refused partitions included.
+///
+/// 64 MiB at the default 8 MiB frame, so a 1 MB request can compress 64x 
before the rest of it
+/// answers 6.
+const REQUEST_ALLOWANCES: usize = 8;
+
 /// Blobs this large, or compressed ones, plan off the async worker.
 const BLOCKING_PLAN_BYTES: usize = 64 * 1024;
 
-/// Widest `origin_timestamp` span one Iggy send takes 
(`MAX_TIMESTAMP_DELTA_MICROS`).
-const MAX_SEND_SPAN_MICROS: u64 = u32::MAX as u64;
-
 const RESPONSE_BASE_BYTES: usize = 512;
 
 /// Produce is the only request the wire protocol allows to go unanswered
@@ -162,10 +166,11 @@ pub async fn handle(state: &GatewayState, api_version: 
i16, body: Bytes) -> Hand
             )
         };
     };
-    let budget = DecompressionBudget::new(
-        state.max_frame_size,
-        state.max_frame_size / FRAME_BYTES_PER_RECORD,
-    );
+    let partition = Allowance {
+        bytes: state.max_frame_size,
+        records: state.max_frame_size / FRAME_BYTES_PER_RECORD,
+    };
+    let budget = DecompressionBudget::new(partition, 
partition.times(REQUEST_ALLOWANCES));
     let zstd = if api_version >= ZSTD_MIN_VERSION {
         Zstd::Allowed
     } else {
@@ -244,9 +249,9 @@ impl From<RecordCodecError> for Refusal {
     }
 }
 
-/// Converts and sends one partition at a time, in request order.
+/// Converts and sends one partition at a time, in request order, one Iggy 
send each.
 ///
-/// One budget for the whole request, so many partitions cannot each take the 
full allowance.
+/// One budget for the whole request, so a small request cannot make this 
inflate without end.
 /// Only one partition's messages live at a time. The budget is owned here and 
borrowed only in
 /// sync calls: a borrow held across an await makes the connection task 
`!Send`.
 ///
@@ -272,11 +277,9 @@ async fn write_request(
                 Ok(target) => {
                     let planned = plan(&budget, zstd, name, partition, 
max_send);
                     match planned {
-                        Ok((id, messages)) => {
-                            send_runs(bridge, target, name, id, messages, 
deadline)
-                                .await
-                                .map_err(Refusal::from)
-                        }
+                        Ok((id, messages)) => send_by(bridge, target, name, 
id, messages, deadline)
+                            .await
+                            .map_err(Refusal::from),
                         Err(refusal) => Err(refusal),
                     }
                 }
@@ -305,10 +308,12 @@ fn plan(
     max_send: u64,
 ) -> std::result::Result<(u32, Vec<IggyMessage>), Refusal> {
     let run = || plan_partition(budget, zstd, kafka_topic, partition, 
max_send);
-    let heavy = partition
-        .records
-        .as_ref()
-        .is_some_and(|blob| blob.len() >= BLOCKING_PLAN_BYTES || 
is_compressed(blob));
+    // A spent budget refuses before it reads a byte, so there is nothing to 
move off the worker.
+    let heavy = !budget.is_spent()
+        && partition
+            .records
+            .as_ref()
+            .is_some_and(|blob| blob.len() >= BLOCKING_PLAN_BYTES || 
is_compressed(blob));
     if heavy && Handle::current().runtime_flavor() == 
RuntimeFlavor::MultiThread {
         tokio::task::block_in_place(run)
     } else {
@@ -363,10 +368,12 @@ fn plan_partition(
         return Err(ERROR_INVALID_RECORD.into());
     }
 
+    // One window for the whole batch, so it fits one send.
+    let window = TimestampWindow::of(&records);
     let mut messages = Vec::with_capacity(records.len());
     let mut size = 0u64;
     for record in &records {
-        let message = to_iggy(record).map_err(refuse)?;
+        let message = to_iggy(record, window).map_err(refuse)?;
         size = size.saturating_add(message.get_size_bytes().as_bytes_u64());
         messages.push(message);
     }
@@ -396,76 +403,7 @@ fn request_timeout(timeout_ms: i32) -> Duration {
         .min(MAX_REQUEST_DEADLINE)
 }
 
-/// Splits `messages` into in-order runs whose timestamp span fits one Iggy 
send.
-///
-/// A record with no timestamp stores 0, so next to real timestamps it starts 
a new run.
-fn send_spans(messages: &[IggyMessage]) -> Vec<Range<usize>> {
-    let mut spans = Vec::with_capacity(1);
-    let mut start = 0;
-    let (mut low, mut high) = (u64::MAX, 0u64);
-    for (index, message) in messages.iter().enumerate() {
-        let timestamp = message.header.origin_timestamp;
-        let (next_low, next_high) = (low.min(timestamp), high.max(timestamp));
-        if next_high - next_low > MAX_SEND_SPAN_MICROS {
-            spans.push(start..index);
-            start = index;
-            (low, high) = (timestamp, timestamp);
-        } else {
-            (low, high) = (next_low, next_high);
-        }
-    }
-    spans.push(start..messages.len());
-    spans
-}
-
-/// `messages` cut into the runs of [`send_spans`], in order.
-fn split_runs(messages: Vec<IggyMessage>) -> Vec<Vec<IggyMessage>> {
-    let spans = send_spans(&messages);
-    if spans.len() == 1 {
-        return vec![messages];
-    }
-    let mut rest = messages.into_iter();
-    spans
-        .into_iter()
-        .map(|span| rest.by_ref().take(span.len()).collect())
-        .collect()
-}
-
-/// Sends each run of one partition. A failed run stops the rest. Runs already 
sent stay stored.
-async fn send_runs(
-    bridge: &IggyBridge,
-    target: &TopicTarget,
-    kafka_topic: &str,
-    partition_id: u32,
-    messages: Vec<IggyMessage>,
-    deadline: Deadline,
-) -> std::result::Result<Option<u64>, i16> {
-    let runs = split_runs(messages);
-    let mut sent = Vec::with_capacity(runs.len());
-    for run in runs {
-        let count = run.len();
-        let base_offset = send_by(bridge, target, kafka_topic, partition_id, 
run, deadline).await?;
-        sent.push((base_offset, count));
-    }
-    Ok(joined_base_offset(&sent))
-}
-
-/// The first run's base offset, if each later run starts where the one before 
it ended.
-///
-/// Another writer can append between runs, and a client counts every offset 
from the base.
-fn joined_base_offset(runs: &[(Option<u64>, usize)]) -> Option<u64> {
-    let base = runs.first()?.0?;
-    let mut next = base;
-    for &(base_offset, count) in runs {
-        if base_offset? != next {
-            return None;
-        }
-        next = next.checked_add(u64::try_from(count).ok()?)?;
-    }
-    Some(base)
-}
-
-/// Sends one run, or answers 7 once the deadline has passed.
+/// Sends one partition's messages, or answers 7 once the deadline has passed.
 ///
 /// 7 means nothing was written, or the send may still land and a retry 
duplicate it.
 async fn send_by(
@@ -529,8 +467,7 @@ const fn record_error_code(error: &RecordCodecError) -> i16 
{
         | RecordCodecError::BudgetExceeded { .. }
         | RecordCodecError::RecordBudgetExceeded { .. }
         | RecordCodecError::EnvelopeTooLarge { .. } => ERROR_MESSAGE_TOO_LARGE,
-        // Earlier partitions spent the request budget. This one fits alone, 
and nothing was
-        // written.
+        // The request budget ran out. Nothing was written.
         RecordCodecError::RequestBudgetSpent => ERROR_NOT_LEADER_OR_FOLLOWER,
         // Transactional or control. 35 is a fatal code for these producers, 
so they stop instead
         // of retrying forever.
@@ -544,6 +481,7 @@ const fn record_error_code(error: &RecordCodecError) -> i16 
{
         | RecordCodecError::TimestampOutOfRange(_)
         | RecordCodecError::Batch(_)
         | RecordCodecError::SeveralBatches(_)
+        | RecordCodecError::BatchTrailingBytes(_)
         | RecordCodecError::RecordCountTooLarge { .. }
         | RecordCodecError::RecordCountMismatch { .. }
         | RecordCodecError::HeaderCountTooLarge { .. }
@@ -555,7 +493,7 @@ const fn record_error_code(error: &RecordCodecError) -> i16 
{
         | RecordCodecError::MappingVersion(_)
         | RecordCodecError::ValueMarker(_)
         | RecordCodecError::HeaderNameCollision(_)
-        | RecordCodecError::TimestampMarker(_)
+        | RecordCodecError::TimestampHeader(_)
         | RecordCodecError::EnvelopeTruncated { .. }
         | RecordCodecError::EnvelopeTrailingBytes(_)
         | RecordCodecError::EnvelopeVersion(_)
@@ -687,19 +625,16 @@ mod tests {
 
     fn planned(entry: &PartitionProduceData) -> std::result::Result<(u32, 
Vec<IggyMessage>), i16> {
         plan_with(
-            &DecompressionBudget::new(TEST_BUDGET, usize::MAX),
+            &partition_budget(TEST_BUDGET, usize::MAX),
             Zstd::Allowed,
             entry,
         )
     }
 
-    fn message_at(micros: u64) -> IggyMessage {
-        let mut message = IggyMessage::builder()
-            .payload(Bytes::from_static(b"v"))
-            .build()
-            .unwrap();
-        message.header.origin_timestamp = micros;
-        message
+    /// A request allowance of one partition's.
+    fn partition_budget(bytes: usize, records: usize) -> DecompressionBudget {
+        let allowance = Allowance { bytes, records };
+        DecompressionBudget::new(allowance, allowance)
     }
 
     #[test]
@@ -770,7 +705,7 @@ mod tests {
     #[test]
     fn 
given_zstd_before_v7_when_planned_should_answer_unsupported_compression() {
         let batch = compressed_batch(&[record(0, b"v")], Compression::Zstd);
-        let budget = DecompressionBudget::new(TEST_BUDGET, usize::MAX);
+        let budget = partition_budget(TEST_BUDGET, usize::MAX);
 
         assert_eq!(
             plan_with(&budget, Zstd::Refused, &entry(0, 
Some(batch.clone()))).unwrap_err(),
@@ -807,7 +742,7 @@ mod tests {
     #[test]
     fn 
given_a_batch_that_decompresses_past_the_budget_when_planned_should_answer_too_large()
 {
         let batch = gzip_batch(&[record(0, &[b'a'; 4096])]);
-        let budget = DecompressionBudget::new(16, usize::MAX);
+        let budget = partition_budget(16, usize::MAX);
 
         assert_eq!(
             plan_with(&budget, Zstd::Allowed, &entry(0, 
Some(batch))).unwrap_err(),
@@ -816,9 +751,9 @@ mod tests {
     }
 
     #[test]
-    fn 
given_more_records_than_the_request_allows_when_planned_should_answer_too_large()
 {
+    fn 
given_more_records_than_a_partition_allows_when_planned_should_answer_too_large()
 {
         let batch = encode_batch(&mut [record(0, b"a"), record(1, 
b"b")]).unwrap();
-        let budget = DecompressionBudget::new(TEST_BUDGET, 1);
+        let budget = partition_budget(TEST_BUDGET, 1);
 
         assert_eq!(
             plan_with(&budget, Zstd::Allowed, &entry(0, 
Some(batch))).unwrap_err(),
@@ -829,7 +764,7 @@ mod tests {
 
     #[test]
     fn 
given_an_earlier_partition_spent_the_budget_when_planned_should_answer_retriable()
 {
-        let budget = DecompressionBudget::new(6144, usize::MAX);
+        let budget = partition_budget(6144, usize::MAX);
         let first = entry(0, Some(gzip_batch(&[record(0, &[b'a'; 4096])])));
         let second = entry(1, Some(gzip_batch(&[record(0, &[b'b'; 4096])])));
 
@@ -842,9 +777,13 @@ mod tests {
     }
 
     #[test]
-    fn 
given_an_entry_over_the_whole_budget_after_another_when_planned_should_answer_too_large()
 {
-        // gzip writes 32 KiB at a time, so the second entry passes the budget 
in several writes.
-        let budget = DecompressionBudget::new(64 * 1024, usize::MAX);
+    fn 
given_an_entry_over_its_cap_after_another_when_planned_should_answer_too_large()
 {
+        // gzip writes 32 KiB at a time, so the second entry passes its cap in 
several writes.
+        let partition = Allowance {
+            bytes: 64 * 1024,
+            records: usize::MAX,
+        };
+        let budget = DecompressionBudget::new(partition, partition.times(8));
         let first = entry(0, Some(gzip_batch(&[record(0, &vec![b'a'; 60 * 
1024])])));
         let second = entry(1, Some(gzip_batch(&[record(0, &vec![b'b'; 200 * 
1024])])));
 
@@ -858,7 +797,7 @@ mod tests {
 
     #[test]
     fn 
given_an_earlier_partition_spent_the_record_budget_when_planned_should_answer_retriable()
 {
-        let budget = DecompressionBudget::new(TEST_BUDGET, 3);
+        let budget = partition_budget(TEST_BUDGET, 3);
         let two = || encode_batch(&mut [record(0, b"a"), record(1, 
b"b")]).unwrap();
 
         assert!(plan_with(&budget, Zstd::Allowed, &entry(0, 
Some(two()))).is_ok());
@@ -869,21 +808,41 @@ mod tests {
     }
 
     #[test]
-    fn given_close_timestamps_when_split_should_send_once() {
-        let messages = [message_at(5), message_at(MAX_SEND_SPAN_MICROS + 5)];
-        assert_eq!(send_spans(&messages), vec![0..2]);
+    fn 
given_a_spent_request_when_a_later_partition_is_planned_should_answer_retry_undecoded()
 {
+        let budget = partition_budget(6144, usize::MAX);
+        let fill = |index| entry(index, Some(gzip_batch(&[record(0, &[b'a'; 
4096])])));
+        assert!(plan_with(&budget, Zstd::Allowed, &fill(0)).is_ok());
+        assert_eq!(
+            plan_with(&budget, Zstd::Allowed, &fill(1)).unwrap_err(),
+            ERROR_NOT_LEADER_OR_FOLLOWER
+        );
+
+        // Malformed, so a decode would answer 87.
+        let garbage = entry(2, Some(Bytes::from_static(b"not a record 
batch")));
+        assert_eq!(
+            plan_with(&budget, Zstd::Allowed, &garbage).unwrap_err(),
+            ERROR_NOT_LEADER_OR_FOLLOWER,
+            "nothing decodes once the request budget is gone"
+        );
     }
 
     #[test]
-    fn given_a_wide_timestamp_span_when_split_should_start_a_new_run() {
-        let far = MAX_SEND_SPAN_MICROS + 1;
-        let messages = [
-            message_at(far),
-            message_at(0),
-            message_at(1),
-            message_at(far + 1),
-        ];
-        assert_eq!(send_spans(&messages), vec![0..1, 1..3, 3..4]);
+    fn given_timestamps_past_one_send_when_planned_should_fit_one_send() {
+        let mut records = [record(0, b"a"), record(1, b"b"), record(2, b"c")];
+        records[1].timestamp = CREATE_TIME + 72 * 60 * 1000;
+        records[2].timestamp = -1;
+        let (_, messages) = planned(&entry(0, Some(encode_batch(&mut 
records).unwrap()))).unwrap();
+
+        let stamps: Vec<u64> = messages
+            .iter()
+            .map(|message| message.header.origin_timestamp)
+            .collect();
+        let low = stamps.iter().min().unwrap();
+        let high = stamps.iter().max().unwrap();
+        assert!(
+            u32::try_from(high - low).is_ok(),
+            "one Iggy send stores each as a u32 delta: {stamps:?}"
+        );
     }
 
     #[test]
@@ -955,6 +914,7 @@ mod tests {
                 walked: 2,
             },
             RecordCodecError::SeveralBatches(2),
+            RecordCodecError::BatchTrailingBytes(1),
             RecordCodecError::HeaderCountTooLarge {
                 count: i32::MAX,
                 limit: 1,
@@ -1007,46 +967,9 @@ mod tests {
     #[test]
     fn 
given_a_partition_past_the_send_cap_when_planned_should_answer_too_large() {
         let batch = encode_batch(&mut [record(0, b"value")]).unwrap();
-        let budget = DecompressionBudget::new(TEST_BUDGET, usize::MAX);
+        let budget = partition_budget(TEST_BUDGET, usize::MAX);
         let refusal =
             plan_partition(&budget, Zstd::Allowed, TOPIC, &entry(0, 
Some(batch)), 8).unwrap_err();
         assert_eq!(refusal.code, ERROR_MESSAGE_TOO_LARGE);
     }
-
-    #[test]
-    fn 
given_runs_that_follow_on_when_joined_should_keep_the_first_base_offset() {
-        assert_eq!(joined_base_offset(&[(Some(5), 3)]), Some(5));
-        assert_eq!(joined_base_offset(&[(Some(5), 3), (Some(8), 1)]), Some(5));
-    }
-
-    #[test]
-    fn given_a_gap_between_runs_when_joined_should_name_no_offset() {
-        assert_eq!(
-            joined_base_offset(&[(Some(5), 3), (Some(9), 1)]),
-            None,
-            "another writer appended between the runs"
-        );
-        assert_eq!(joined_base_offset(&[(Some(5), 3), (None, 1)]), None);
-        assert_eq!(joined_base_offset(&[(None, 3)]), None);
-    }
-
-    #[test]
-    fn given_a_wide_timestamp_span_when_split_should_keep_the_record_order() {
-        let far = MAX_SEND_SPAN_MICROS + 1;
-        let runs = split_runs(vec![
-            message_at(far),
-            message_at(0),
-            message_at(1),
-            message_at(far + 1),
-        ]);
-        let stamps: Vec<Vec<u64>> = runs
-            .iter()
-            .map(|run| {
-                run.iter()
-                    .map(|message| message.header.origin_timestamp)
-                    .collect()
-            })
-            .collect();
-        assert_eq!(stamps, vec![vec![far], vec![0, 1], vec![far + 1]]);
-    }
 }
diff --git a/gateways/kafka/src/records.rs b/gateways/kafka/src/records.rs
index 0044c4773..53640b755 100644
--- a/gateways/kafka/src/records.rs
+++ b/gateways/kafka/src/records.rs
@@ -48,8 +48,8 @@ pub const MAPPING_VERSION: u8 = 1;
 pub const KEY_HEADER: &str = "kafka.key";
 /// Iggy header naming which of null or empty a placeholder payload stands for.
 pub const VALUE_MARKER_HEADER: &str = "kafka.value";
-/// Iggy header marking a record stamped at the Unix epoch.
-pub const TIMESTAMP_MARKER_HEADER: &str = "kafka.ts";
+/// Iggy header holding the Kafka timestamp, `Int64` ms, when 
`origin_timestamp` cannot.
+pub const TIMESTAMP_HEADER: &str = "kafka.ts";
 /// Prefix every Kafka record header name is stored under.
 pub const HEADER_PREFIX: &str = "kafka.h.";
 /// Iggy header whose one-byte value is the envelope byte layout version.
@@ -61,6 +61,8 @@ pub const ENVELOPE_VERSION: u8 = 1;
 const NO_TIMESTAMP: i64 = -1;
 /// The one Kafka timestamp an `origin_timestamp` of zero cannot be told apart 
from.
 const EPOCH_TIMESTAMP: i64 = 0;
+/// Widest `origin_timestamp` span one Iggy send holds 
(`MAX_TIMESTAMP_DELTA_MICROS`).
+const MAX_SEND_SPAN_MICROS: u64 = u32::MAX as u64;
 /// Stored in place of a null or empty value, discarded on the way back.
 const PLACEHOLDER: &[u8] = &[0x00];
 /// Iggy caps one header name and one header value at this many bytes.
@@ -68,7 +70,6 @@ const MAX_FIELD: usize = 255;
 
 const MARKER_NULL: &[u8] = b"null";
 const MARKER_EMPTY: &[u8] = b"empty";
-const MARKER_EPOCH: &[u8] = b"epoch";
 
 const FLAG_KEY: u8 = 0b01;
 const FLAG_VALUE: u8 = 0b10;
@@ -120,8 +121,8 @@ pub enum RecordCodecError {
     ValueMarker(Bytes),
     #[error("two stored header keys both name {0}")]
     HeaderNameCollision(String),
-    #[error("timestamp marker {0:?} is not epoch")]
-    TimestampMarker(Bytes),
+    #[error("timestamp header {0:?} is not a Kafka timestamp")]
+    TimestampHeader(Bytes),
     #[error("envelope for this record is {size} bytes, over Iggy's 
{MAX_PAYLOAD_SIZE} byte limit")]
     EnvelopeTooLarge { size: usize },
     #[error("envelope is truncated: needed {needed} bytes, {remaining} 
remain")]
@@ -136,6 +137,8 @@ pub enum RecordCodecError {
     Batch(String),
     #[error("{0} record batches in one partition, Kafka allows 1")]
     SeveralBatches(usize),
+    #[error("{0} bytes follow the record batch")]
+    BatchTrailingBytes(usize),
     #[error("batch declares {count} records, and {limit} bytes can hold 
fewer")]
     RecordCountTooLarge { count: i32, limit: usize },
     #[error("batch declares {declared} records and holds {walked}")]
@@ -150,11 +153,11 @@ pub enum RecordCodecError {
     TransactionalBatch,
     #[error("control batches are not supported")]
     ControlBatch,
-    #[error("partition decompresses to {size} bytes, over the {limit} byte 
request budget")]
+    #[error("partition decompresses to at least {size} bytes, over its {limit} 
byte cap")]
     BudgetExceeded { size: usize, limit: usize },
-    #[error("partition needs {count} record slots, over the {limit} slot 
request budget")]
+    #[error("partition needs {count} record slots, over its {limit} slot cap")]
     RecordBudgetExceeded { count: usize, limit: usize },
-    #[error("earlier partitions used the request budget, retry")]
+    #[error("the request budget ran out, retry")]
     RequestBudgetSpent,
     #[error("zstd batch in a request older than Produce v7")]
     ZstdTooEarly,
@@ -164,7 +167,7 @@ pub enum RecordCodecError {
 
 type Result<T> = std::result::Result<T, RecordCodecError>;
 
-/// Encodes one Kafka record as one Iggy message.
+/// Encodes one Kafka record as one Iggy message, for a send whose timestamps 
fit `window`.
 ///
 /// Takes the native path when Iggy can hold every field, and the envelope 
otherwise. A caller
 /// cannot tell which from the return value, which is the point: `from_iggy` 
reverses both.
@@ -173,15 +176,16 @@ type Result<T> = std::result::Result<T, RecordCodecError>;
 ///
 /// Returns an error when the timestamp does not fit, when the envelope would 
exceed
 /// `MAX_PAYLOAD_SIZE`, or when Iggy rejects the message for a reason the 
envelope does not fix.
-pub fn to_iggy(record: &Record) -> Result<IggyMessage> {
-    if let Some(value) = plain_value(record) {
-        return plain_message(value.clone(), record.timestamp);
+pub fn to_iggy(record: &Record, window: TimestampWindow) -> 
Result<IggyMessage> {
+    let stamp = window.stamp(record.timestamp)?;
+    if let Some(value) = plain_value(record, stamp) {
+        return plain_message(value.clone(), stamp.origin);
     }
     if needs_envelope(record) {
-        return envelope_message(record);
+        return envelope_message(record, stamp);
     }
     let (payload, marker) = split_value(record.value.as_ref());
-    let mut headers = gateway_headers(record.timestamp);
+    let mut headers = gateway_headers(stamp.header);
     if let Some(marker) = marker {
         headers.insert(header_key(VALUE_MARKER_HEADER), header_value(marker));
     }
@@ -201,7 +205,46 @@ pub fn to_iggy(record: &Record) -> Result<IggyMessage> {
 
     // The only limit left is the 100 KB budget over all headers together, 
which no per-field
     // check can see. Let the constructor rule on it rather than duplicating 
its arithmetic.
-    build(payload, headers, record.timestamp)?.map_or_else(|| 
envelope_message(record), Ok)
+    build(payload, headers, stamp.origin)?.map_or_else(|| 
envelope_message(record, stamp), Ok)
+}
+
+/// The `origin_timestamp` range one Iggy send holds: `u32::MAX` µs, about 
71.6 min.
+///
+/// Starts at the batch's earliest real timestamp. A timestamp outside it is 
clamped into it, and
+/// `kafka.ts` keeps the real one, so a Kafka batch always fits one send.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub struct TimestampWindow {
+    start: u64,
+}
+
+impl TimestampWindow {
+    /// The window for `records`, the batch one send carries. `-1` and the 
epoch do not move it.
+    #[must_use]
+    pub fn of(records: &[Record]) -> Self {
+        let start = records
+            .iter()
+            .filter(|record| record.timestamp > EPOCH_TIMESTAMP)
+            .filter_map(|record| timestamp_in(record.timestamp).ok())
+            .min()
+            .unwrap_or(0);
+        Self { start }
+    }
+
+    fn stamp(self, millis: i64) -> Result<Stamp> {
+        let native = timestamp_in(millis)?;
+        let origin = native.clamp(self.start, 
self.start.saturating_add(MAX_SEND_SPAN_MICROS));
+        // A zero origin reads back as "no timestamp", so the epoch needs the 
header unclamped too.
+        let header = (origin != native || millis == 
EPOCH_TIMESTAMP).then_some(millis);
+        Ok(Stamp { origin, header })
+    }
+}
+
+/// Where one record's timestamp is stored.
+#[derive(Clone, Copy)]
+struct Stamp {
+    origin: u64,
+    /// The `kafka.ts` value, when `origin` does not read back as the record's 
timestamp.
+    header: Option<i64>,
 }
 
 /// Decodes one Iggy message as one Kafka record at `offset`.
@@ -238,11 +281,12 @@ pub fn from_iggy(message: &IggyMessage, offset: i64) -> 
Result<Record> {
         Some(envelope) => decode_envelope(envelope.as_bytes(), 
&message.payload)?,
         None => gateway_fields(message, &stored)?,
     };
-    let timestamp = match stored.get(&header_key(TIMESTAMP_MARKER_HEADER)) {
+    let timestamp = match stored.get(&header_key(TIMESTAMP_HEADER)) {
         None => timestamp_out(message),
-        Some(marker) => match marker.as_bytes() {
-            MARKER_EPOCH => EPOCH_TIMESTAMP,
-            _ => return Err(RecordCodecError::TimestampMarker(marker.value())),
+        Some(header) => match header.as_int64() {
+            Ok(NO_TIMESTAMP) => server_timestamp(message),
+            Ok(millis) if millis >= EPOCH_TIMESTAMP => millis,
+            _ => return Err(RecordCodecError::TimestampHeader(header.value())),
         },
     };
     Ok(record(key, value, headers, offset, timestamp))
@@ -276,15 +320,17 @@ fn timestamp_in(millis: i64) -> Result<u64> {
 
 /// Zero means the producer sent no timestamp, so the server-assigned one 
stands in.
 ///
-/// A record stamped at the epoch stores that same zero, and `from_iggy` reads 
the `kafka.ts`
-/// marker before it calls this, because Iggy has no other way to hold the 
difference.
+/// `from_iggy` reads `kafka.ts` first: an epoch or clamped record stores a 
value that does not
+/// read back as its own.
 fn timestamp_out(message: &IggyMessage) -> i64 {
-    let micros = if message.header.origin_timestamp == 0 {
-        message.header.timestamp
-    } else {
-        message.header.origin_timestamp
-    };
-    i64::try_from(micros / 1000).unwrap_or(NO_TIMESTAMP)
+    match message.header.origin_timestamp {
+        0 => server_timestamp(message),
+        micros => i64::try_from(micros / 1000).unwrap_or(NO_TIMESTAMP),
+    }
+}
+
+fn server_timestamp(message: &IggyMessage) -> i64 {
+    i64::try_from(message.header.timestamp / 1000).unwrap_or(NO_TIMESTAMP)
 }
 
 /// Whether any field of `record` is one Iggy refuses to hold natively.
@@ -307,8 +353,8 @@ fn needs_envelope(record: &Record) -> bool {
 }
 
 /// The value, when `kafka.v` is the only header the record needs.
-fn plain_value(record: &Record) -> Option<&Bytes> {
-    if record.key.is_some() || !record.headers.is_empty() || record.timestamp 
== EPOCH_TIMESTAMP {
+fn plain_value(record: &Record, stamp: Stamp) -> Option<&Bytes> {
+    if record.key.is_some() || !record.headers.is_empty() || 
stamp.header.is_some() {
         return None;
     }
     record.value.as_ref().filter(|value| !value.is_empty())
@@ -317,12 +363,12 @@ fn plain_value(record: &Record) -> Option<&Bytes> {
 /// Reuses one encoded `kafka.v` header block, so the common record skips a 
map and an encode.
 ///
 /// Static bytes: a clone touches no refcount shared across workers.
-fn plain_message(value: Bytes, timestamp: i64) -> Result<IggyMessage> {
+fn plain_message(value: Bytes, origin: u64) -> Result<IggyMessage> {
     static VERSION_ONLY: OnceLock<&'static [u8]> = OnceLock::new();
     let headers = *VERSION_ONLY.get_or_init(|| {
         let block = IggyMessage::builder()
             .payload(Bytes::from_static(PLACEHOLDER))
-            .user_headers(gateway_headers(NO_TIMESTAMP))
+            .user_headers(gateway_headers(None))
             .build()
             .ok()
             .and_then(|message| message.user_headers)
@@ -333,7 +379,7 @@ fn plain_message(value: Bytes, timestamp: i64) -> 
Result<IggyMessage> {
     message.header.user_headers_length =
         u32::try_from(headers.len()).unwrap_or_else(|_| unreachable!("the 
kafka.v header is tiny"));
     message.user_headers = Some(Bytes::from_static(headers));
-    message.header.origin_timestamp = timestamp_in(timestamp)?;
+    message.header.origin_timestamp = origin;
     Ok(message)
 }
 
@@ -346,18 +392,13 @@ fn split_value(value: Option<&Bytes>) -> (Bytes, 
Option<&'static [u8]>) {
     }
 }
 
-/// The headers every gateway-written message carries, whichever path it takes.
-///
-/// Iggy reads an `origin_timestamp` of zero as no timestamp at all, so a 
record that really was
-/// stamped at the epoch needs a marker to hold the difference.
-fn gateway_headers(timestamp: i64) -> BTreeMap<HeaderKey, HeaderValue> {
+/// The headers every gateway-written message carries, whichever path it 
takes, plus `kafka.ts`
+/// when `timestamp` holds one.
+fn gateway_headers(timestamp: Option<i64>) -> BTreeMap<HeaderKey, HeaderValue> 
{
     let mut headers = BTreeMap::new();
     headers.insert(header_key(VERSION_HEADER), 
header_value(&[MAPPING_VERSION]));
-    if timestamp == EPOCH_TIMESTAMP {
-        headers.insert(
-            header_key(TIMESTAMP_MARKER_HEADER),
-            header_value(MARKER_EPOCH),
-        );
+    if let Some(millis) = timestamp {
+        headers.insert(header_key(TIMESTAMP_HEADER), 
HeaderValue::from(millis));
     }
     headers
 }
@@ -366,7 +407,7 @@ fn gateway_headers(timestamp: i64) -> BTreeMap<HeaderKey, 
HeaderValue> {
 fn build(
     payload: Bytes,
     headers: BTreeMap<HeaderKey, HeaderValue>,
-    timestamp: i64,
+    origin: u64,
 ) -> Result<Option<IggyMessage>> {
     let mut message = match IggyMessage::builder()
         .payload(payload)
@@ -377,11 +418,11 @@ fn build(
         Err(IggyError::TooBigUserHeaders) => return Ok(None),
         Err(error) => return Err(error.into()),
     };
-    message.header.origin_timestamp = timestamp_in(timestamp)?;
+    message.header.origin_timestamp = origin;
     Ok(Some(message))
 }
 
-fn envelope_message(record: &Record) -> Result<IggyMessage> {
+fn envelope_message(record: &Record, stamp: Stamp) -> Result<IggyMessage> {
     // The envelope moves the key and the headers into the payload, so a 
record whose value alone
     // clears `MAX_PAYLOAD_SIZE` can be one the fallback cannot hold. Say so 
before spending the
     // allocation, since the native path has already been ruled out and 
nothing else is left.
@@ -390,12 +431,12 @@ fn envelope_message(record: &Record) -> 
Result<IggyMessage> {
         return Err(RecordCodecError::EnvelopeTooLarge { size });
     }
 
-    let mut headers = gateway_headers(record.timestamp);
+    let mut headers = gateway_headers(stamp.header);
     headers.insert(
         header_key(ENVELOPE_HEADER),
         header_value(&[ENVELOPE_VERSION]),
     );
-    build(encode_envelope(record, size), headers, record.timestamp)?
+    build(encode_envelope(record, size), headers, stamp.origin)?
         .ok_or(IggyError::TooBigUserHeaders)
         .map_err(Into::into)
 }
@@ -627,19 +668,43 @@ fn header_value(value: &[u8]) -> HeaderValue {
     HeaderValue::try_from(value).unwrap_or_else(|_| unreachable!("header value 
is out of range"))
 }
 
-/// What one Produce request may decompress to, in total.
+/// Decompressed bytes and record slots, the two things a decode spends.
 ///
-/// Charged across every partition in the request, so many partitions cannot 
each take it all.
-/// Writes go through `BudgetedWriter`, so an over-budget frame never inflates 
whole in memory.
+/// A record slot costs a `Record` and an `IggyMessage` in memory.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub struct Allowance {
+    pub bytes: usize,
+    pub records: usize,
+}
+
+impl Allowance {
+    #[must_use]
+    pub const fn times(self, factor: usize) -> Self {
+        Self {
+            bytes: self.bytes.saturating_mul(factor),
+            records: self.records.saturating_mul(factor),
+        }
+    }
+}
+
+/// What one Produce request may decode.
+///
+/// - One partition may take `partition`. Past it, it is too large alone (10).
+/// - All partitions together may take `request`. It counts work done, refused 
partitions too, so
+///   no refusal lets the same inflate run again for free. Once it runs out, 
the partition that
+///   ran it out and every later one answer 6, undecoded. A retry that comes 
first in its request
+///   gets the full allowance, so a partition too large alone still ends at 10.
+///
+/// Writes go through `BudgetedWriter`, so no partition inflates past either 
one in memory.
 pub struct DecompressionBudget {
-    bytes: usize,
-    remaining: Cell<usize>,
-    /// Records the request may decode. Each one costs a `Record` and an 
`IggyMessage` in memory.
-    records: usize,
-    records_remaining: Cell<usize>,
-    /// What the current partition charged. Tells "too large" (10) from 
"retry" (6).
+    partition: Allowance,
+    /// What the request has left.
+    bytes_left: Cell<usize>,
+    records_left: Cell<usize>,
+    /// What the current partition took.
     entry_bytes: Cell<usize>,
     entry_records: Cell<usize>,
+    spent: Cell<bool>,
     /// Why this module refused a batch, when it did. Everything this module 
reports from inside
     /// the decoder's decompression hook leaves as an `io::Error` or an 
`anyhow::Error`, which the
     /// decoder stringifies, so the typed reason is parked here and taken by 
`decode_batch`.
@@ -648,71 +713,71 @@ pub struct DecompressionBudget {
 
 impl DecompressionBudget {
     #[must_use]
-    pub const fn new(bytes: usize, records: usize) -> Self {
+    pub const fn new(partition: Allowance, request: Allowance) -> Self {
         Self {
-            bytes,
-            remaining: Cell::new(bytes),
-            records,
-            records_remaining: Cell::new(records),
+            partition,
+            bytes_left: Cell::new(request.bytes),
+            records_left: Cell::new(request.records),
             entry_bytes: Cell::new(0),
             entry_records: Cell::new(0),
+            spent: Cell::new(false),
             reason: Cell::new(None),
         }
     }
 
-    fn start_entry(&self) {
-        self.entry_bytes.set(0);
-        self.entry_records.set(0);
-    }
-
-    /// Takes `len` bytes if the request has them left.
-    fn take(&self, len: usize) -> bool {
-        let Some(remaining) = self.remaining.get().checked_sub(len) else {
-            return false;
-        };
-        self.remaining.set(remaining);
-        self.entry_bytes.set(self.entry_bytes.get() + len);
-        true
+    /// Whether the request allowance ran out, so nothing more decodes.
+    #[must_use]
+    pub const fn is_spent(&self) -> bool {
+        self.spent.get()
     }
 
-    /// Takes `len` bytes at once, or refuses.
-    fn charge(&self, len: usize) -> io::Result<()> {
-        if self.take(len) {
-            return Ok(());
+    fn start_entry(&self) -> Result<()> {
+        if self.is_spent() {
+            return Err(RecordCodecError::RequestBudgetSpent);
         }
-        self.check_alone(len)?;
-        Err(self.refuse(RecordCodecError::RequestBudgetSpent))
+        self.entry_bytes.set(0);
+        self.entry_records.set(0);
+        Ok(())
     }
 
-    /// Refuses once this partition alone passes the budget. No retry can help 
it.
-    fn check_alone(&self, spilled: usize) -> io::Result<()> {
-        let size = self.entry_bytes.get().saturating_add(spilled);
-        if size > self.bytes {
+    /// Takes `len` decoded bytes for the current partition, or refuses.
+    fn charge(&self, len: usize) -> io::Result<()> {
+        let entry = self.entry_bytes.get().saturating_add(len);
+        if entry > self.partition.bytes {
             return Err(self.refuse(RecordCodecError::BudgetExceeded {
-                size,
-                limit: self.bytes,
+                size: entry,
+                limit: self.partition.bytes,
             }));
         }
+        let Some(left) = self.bytes_left.get().checked_sub(len) else {
+            return Err(self.refuse(self.spend()));
+        };
+        self.bytes_left.set(left);
+        self.entry_bytes.set(entry);
         Ok(())
     }
 
     fn charge_records(&self, count: usize) -> Result<()> {
         let entry = self.entry_records.get().saturating_add(count);
-        let Some(remaining) = self.records_remaining.get().checked_sub(count) 
else {
-            return Err(if entry > self.records {
-                RecordCodecError::RecordBudgetExceeded {
-                    count: entry,
-                    limit: self.records,
-                }
-            } else {
-                RecordCodecError::RequestBudgetSpent
+        if entry > self.partition.records {
+            return Err(RecordCodecError::RecordBudgetExceeded {
+                count: entry,
+                limit: self.partition.records,
             });
+        }
+        let Some(left) = self.records_left.get().checked_sub(count) else {
+            return Err(self.spend());
         };
-        self.records_remaining.set(remaining);
+        self.records_left.set(left);
         self.entry_records.set(entry);
         Ok(())
     }
 
+    fn spend(&self) -> RecordCodecError {
+        self.spent.set(true);
+        RecordCodecError::RequestBudgetSpent
+    }
+
     /// Parks the typed reason and returns the error the decoder stringifies.
     fn refuse(&self, reason: RecordCodecError) -> io::Error {
         let error = io::Error::other(reason.to_string());
@@ -734,11 +799,9 @@ impl DecompressionBudget {
 /// An `io::Write` sink that stops at the budget rather than after it.
 ///
 /// Charging the output once it exists is too late: the decompressor has 
already allocated it.
-/// Past the budget it counts and drops the bytes, to tell "too large" (10) 
from "retry" (6).
 struct BudgetedWriter<'a> {
     out: BytesMut,
     budget: &'a DecompressionBudget,
-    spilled: usize,
 }
 
 impl<'a> BudgetedWriter<'a> {
@@ -746,34 +809,18 @@ impl<'a> BudgetedWriter<'a> {
         Self {
             out: BytesMut::new(),
             budget,
-            spilled: 0,
-        }
-    }
-
-    /// Charges `len` bytes. `false` once past the budget: counted, not kept.
-    fn admit(&mut self, len: usize) -> io::Result<bool> {
-        if self.spilled == 0 && self.budget.take(len) {
-            return Ok(true);
         }
-        self.spilled = self.spilled.saturating_add(len);
-        self.budget.check_alone(self.spilled)?;
-        Ok(false)
     }
 
-    /// The output, or a refusal if any of it spilled.
-    fn into_output(self) -> io::Result<Bytes> {
-        if self.spilled > 0 {
-            return 
Err(self.budget.refuse(RecordCodecError::RequestBudgetSpent));
-        }
-        Ok(self.out.freeze())
+    fn into_output(self) -> Bytes {
+        self.out.freeze()
     }
 }
 
 impl Write for BudgetedWriter<'_> {
     fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
-        if self.admit(buf.len())? {
-            self.out.extend_from_slice(buf);
-        }
+        self.budget.charge(buf.len())?;
+        self.out.extend_from_slice(buf);
         Ok(buf.len())
     }
 
@@ -800,15 +847,16 @@ pub fn is_compressed(blob: &[u8]) -> bool {
 ///
 /// # Errors
 ///
-/// Returns an error when the blob holds more than one batch, when the batch 
is malformed or holds
-/// a record count other than the one it declares, when it is a control or 
transactional batch,
-/// when it is zstd and `zstd` refuses it, or when it passes `budget`.
+/// Returns an error when the blob holds more than one batch or bytes after 
it, when the batch is
+/// malformed or holds a record count other than the one it declares, when it 
is a control or
+/// transactional batch, when it is zstd and `zstd` refuses it, or when it 
passes `budget` or the
+/// budget is already spent.
 pub fn decode_batch(
     buf: &mut Bytes,
     budget: &DecompressionBudget,
     zstd: Zstd,
 ) -> Result<Vec<Record>> {
-    budget.start_entry();
+    budget.start_entry()?;
     preflight(buf, budget, zstd)?;
     if !buf.has_remaining() {
         return Ok(Vec::new());
@@ -838,20 +886,26 @@ pub fn decode_batch(
 
 /// Checks the batch header before anything decodes a record.
 ///
-/// - One batch per partition, as Kafka requires from Produce v3.
+/// - One batch per partition, as Kafka requires from Produce v3, and no bytes 
after it.
 /// - The decoder reserves from `record_count` 
(`kafka-protocol-0.18.0/src/records.rs:517`), so
-///   the count must fit the blob, plus the full budget when compressed.
+///   the count must fit the blob, plus the partition cap when compressed.
 /// - Control and transactional batches are refused: a stored message cannot 
carry their flags.
 ///   An idempotent batch passes, its producer id, epoch and sequence ignored, 
because a stock
 ///   Java producer is idempotent. `IDEMPOTENCE.md` has why a retry is then 
not deduplicated.
 fn preflight(buf: &Bytes, budget: &DecompressionBudget, zstd: Zstd) -> 
Result<()> {
-    let infos = RecordBatchDecoder::decode_batch_info(&mut buf.clone())
+    let mut rest = buf.clone();
+    let infos = RecordBatchDecoder::decode_batch_info(&mut rest)
         .map_err(|error| RecordCodecError::Batch(error.to_string()))?;
     let info = match infos.as_slice() {
         [] => return Ok(()),
         [info] => info,
         several => return Err(RecordCodecError::SeveralBatches(several.len())),
     };
+    // `decode_batch_info` stops at a magic byte other than 2 and leaves the 
rest unread, and so
+    // does the decoder. Unchecked, the tail is dropped and the partition 
answers success.
+    if rest.has_remaining() {
+        return Err(RecordCodecError::BatchTrailingBytes(rest.remaining()));
+    }
     if info.transactional {
         return Err(RecordCodecError::TransactionalBatch);
     }
@@ -865,7 +919,7 @@ fn preflight(buf: &Bytes, budget: &DecompressionBudget, 
zstd: Zstd) -> Result<()
     let limit = if info.compression == Compression::None {
         buf.len()
     } else {
-        buf.len().saturating_add(budget.bytes)
+        buf.len().saturating_add(budget.partition.bytes)
     };
     let declared = usize::try_from(info.record_count).unwrap_or(0);
     if declared.saturating_mul(MIN_RECORD_BYTES) > limit {
@@ -938,7 +992,7 @@ fn batch_size(records: &[Record]) -> usize {
             .sum::<usize>()
 }
 
-/// Decompresses one batch, refusing the write that would pass the request 
budget.
+/// Decompresses one batch, refusing the write that would pass the budget.
 ///
 /// `kafka_protocol`'s own decompressors write the whole stream into a growing 
buffer before they
 /// hand it over, so the four codecs are driven from here instead. Each one 
writes through
@@ -964,21 +1018,21 @@ fn decompress(
             let mut decoder = flate2::write::GzDecoder::new(&mut writer);
             decoder.write_all(&body)?;
             decoder.finish()?;
-            writer.into_output()?
+            writer.into_output()
         }
         Compression::Zstd => {
             zstd::stream::copy_decode(body.as_ref(), &mut writer)?;
-            writer.into_output()?
+            writer.into_output()
         }
         Compression::Lz4 => {
             let mut decoder = lz4::Decoder::new(body.as_ref())?;
             io::copy(&mut decoder, &mut writer)?;
             decoder.finish().1?;
-            writer.into_output()?
+            writer.into_output()
         }
         Compression::Snappy => {
             inflate_snappy(&body, &mut writer)?;
-            writer.into_output()?
+            writer.into_output()
         }
     };
 
@@ -1101,17 +1155,15 @@ fn skip_field(buf: &mut &[u8]) -> Result<()> {
 ///
 /// Snappy is the one codec that states its output size up front, which `snap` 
reads without
 /// allocating. Charge that number before the decoder runs, because the 
decoder needs the whole
-/// block laid out to write into and cannot be fed a bounded sink. Past the 
budget, blocks are
-/// counted and not inflated.
+/// block laid out to write into and cannot be fed a bounded sink.
 fn inflate_snappy(body: &Bytes, writer: &mut BudgetedWriter<'_>) -> 
anyhow::Result<()> {
     let mut decoder = snap::raw::Decoder::new();
     let mut inflate = |block: &[u8], writer: &mut BudgetedWriter<'_>| -> 
anyhow::Result<()> {
         let declared = snap::raw::decompress_len(block)?;
-        if writer.admit(declared)? {
-            let start = writer.out.len();
-            writer.out.resize(start.saturating_add(declared), 0);
-            decoder.decompress(block, &mut writer.out[start..])?;
-        }
+        writer.budget.charge(declared)?;
+        let start = writer.out.len();
+        writer.out.resize(start.saturating_add(declared), 0);
+        decoder.decompress(block, &mut writer.out[start..])?;
         Ok(())
     };
 
@@ -1191,9 +1243,20 @@ mod tests {
         message_with(payload, &all)
     }
 
+    /// `to_iggy` for a record sent on its own, so nothing clamps it.
+    fn to_iggy_alone(record: &Record) -> Result<IggyMessage> {
+        to_iggy(record, TimestampWindow::of(std::slice::from_ref(record)))
+    }
+
+    /// A request allowance of one partition's.
+    fn partition_budget(bytes: usize, records: usize) -> DecompressionBudget {
+        let allowance = Allowance { bytes, records };
+        DecompressionBudget::new(allowance, allowance)
+    }
+
     #[test]
     fn given_a_value_only_record_when_encoded_should_match_the_general_path() {
-        let plain = to_iggy(&record_with(None, Some(b"v"), &[])).unwrap();
+        let plain = to_iggy_alone(&record_with(None, Some(b"v"), 
&[])).unwrap();
         let general = gateway_message(b"v", &[]);
 
         assert_eq!(plain.user_headers, general.user_headers);
@@ -1228,7 +1291,7 @@ mod tests {
     #[test]
     fn 
given_a_plain_record_when_round_tripped_should_keep_key_value_and_headers() {
         let original = record_with(Some(b"k"), Some(b"v"), &[("trace", 
Some(b"abc"))]);
-        let message = to_iggy(&original).unwrap();
+        let message = to_iggy_alone(&original).unwrap();
         assert!(!is_enveloped(&message));
         assert_eq!(message.payload.as_ref(), b"v");
 
@@ -1245,7 +1308,7 @@ mod tests {
 
     #[test]
     fn given_a_null_value_when_round_tripped_should_stay_null() {
-        let message = to_iggy(&record_with(Some(b"k"), None, &[])).unwrap();
+        let message = to_iggy_alone(&record_with(Some(b"k"), None, 
&[])).unwrap();
         assert!(
             !is_enveloped(&message),
             "a tombstone stays on the fast path"
@@ -1256,7 +1319,7 @@ mod tests {
 
     #[test]
     fn 
given_an_empty_value_when_round_tripped_should_stay_empty_and_not_null() {
-        let message = to_iggy(&record_with(Some(b"k"), Some(b""), 
&[])).unwrap();
+        let message = to_iggy_alone(&record_with(Some(b"k"), Some(b""), 
&[])).unwrap();
         assert_eq!(message.payload.as_ref(), PLACEHOLDER);
         assert_eq!(
             from_iggy(&message, 0).unwrap().value.as_deref(),
@@ -1267,7 +1330,7 @@ mod tests {
     #[test]
     fn given_an_empty_key_when_stored_should_take_the_envelope() {
         let original = record_with(Some(b""), Some(b"v"), &[]);
-        let message = to_iggy(&original).unwrap();
+        let message = to_iggy_alone(&original).unwrap();
         assert!(is_enveloped(&message));
         let back = from_iggy(&message, 0).unwrap();
         assert_eq!(back.key.as_deref(), Some(&[][..]), "empty, not null");
@@ -1277,7 +1340,7 @@ mod tests {
     #[test]
     fn given_an_oversized_key_when_stored_should_take_the_envelope() {
         let key = vec![b'x'; MAX_FIELD + 1];
-        let message = to_iggy(&record_with(Some(&key), Some(b"v"), 
&[])).unwrap();
+        let message = to_iggy_alone(&record_with(Some(&key), Some(b"v"), 
&[])).unwrap();
         assert!(is_enveloped(&message));
         assert_eq!(
             from_iggy(&message, 0).unwrap().key.as_deref(),
@@ -1287,7 +1350,8 @@ mod tests {
 
     #[test]
     fn given_a_null_header_value_when_stored_should_take_the_envelope() {
-        let message = to_iggy(&record_with(Some(b"k"), Some(b"v"), &[("flag", 
None)])).unwrap();
+        let message =
+            to_iggy_alone(&record_with(Some(b"k"), Some(b"v"), &[("flag", 
None)])).unwrap();
         assert!(is_enveloped(&message));
         assert_eq!(
             from_iggy(&message, 0)
@@ -1303,7 +1367,7 @@ mod tests {
     fn given_an_oversized_header_name_when_stored_should_take_the_envelope() {
         let name = "n".repeat(MAX_FIELD - HEADER_PREFIX.len() + 1);
         let message =
-            to_iggy(&record_with(Some(b"k"), Some(b"v"), &[(&name, 
Some(b"v"))])).unwrap();
+            to_iggy_alone(&record_with(Some(b"k"), Some(b"v"), &[(&name, 
Some(b"v"))])).unwrap();
         assert!(is_enveloped(&message));
         assert!(
             from_iggy(&message, 0)
@@ -1378,7 +1442,7 @@ mod tests {
 
     #[test]
     fn given_unreadable_user_headers_when_read_should_fail() {
-        let mut message = to_iggy(&record_with(Some(b"k"), Some(b"v"), 
&[])).unwrap();
+        let mut message = to_iggy_alone(&record_with(Some(b"k"), Some(b"v"), 
&[])).unwrap();
         message.user_headers = Some(Bytes::from_static(b"not a header block"));
 
         assert!(
@@ -1412,17 +1476,89 @@ mod tests {
     }
 
     #[test]
-    fn 
given_an_unknown_timestamp_marker_on_a_gateway_message_when_read_should_fail() {
-        let message = gateway_message(b"v", &[(TIMESTAMP_MARKER_HEADER, 
b"later")]);
+    fn 
given_a_timestamp_header_that_is_not_int64_on_a_gateway_message_when_read_should_fail()
 {
+        let message = gateway_message(b"v", &[(TIMESTAMP_HEADER, b"epoch")]);
         assert!(matches!(
             from_iggy(&message, 0),
-            Err(RecordCodecError::TimestampMarker(_))
+            Err(RecordCodecError::TimestampHeader(_))
         ));
     }
 
+    #[test]
+    fn given_a_timestamp_header_below_minus_one_when_read_should_fail() {
+        let mut headers = BTreeMap::new();
+        headers.insert(header_key(VERSION_HEADER), 
header_value(&[MAPPING_VERSION]));
+        headers.insert(header_key(TIMESTAMP_HEADER), HeaderValue::from(-2i64));
+        let message = IggyMessage::builder()
+            .payload(Bytes::from_static(b"v"))
+            .user_headers(headers)
+            .build()
+            .unwrap();
+
+        assert!(
+            matches!(
+                from_iggy(&message, 0),
+                Err(RecordCodecError::TimestampHeader(_))
+            ),
+            "to_iggy refuses these, so no gateway message holds one"
+        );
+    }
+
+    #[test]
+    fn 
given_timestamps_past_one_send_when_stored_should_clamp_and_keep_the_real_one() 
{
+        let later = CREATE_TIME + 72 * 60 * 1000;
+        let records = [record_at(CREATE_TIME), record_at(later)];
+        let window = TimestampWindow::of(&records);
+        let first = to_iggy(&records[0], window).unwrap();
+        let second = to_iggy(&records[1], window).unwrap();
+
+        assert_eq!(
+            second.header.origin_timestamp - first.header.origin_timestamp,
+            MAX_SEND_SPAN_MICROS,
+            "clamped to the window end, so one send holds both"
+        );
+        assert_eq!(from_iggy(&first, 0).unwrap().timestamp, CREATE_TIME);
+        assert_eq!(from_iggy(&second, 1).unwrap().timestamp, later);
+    }
+
+    #[test]
+    fn 
given_no_timestamp_among_real_ones_when_round_tripped_should_read_the_server_timestamp()
 {
+        let records = [record_at(CREATE_TIME), record_at(NO_TIMESTAMP)];
+        let mut message = to_iggy(&records[1], 
TimestampWindow::of(&records)).unwrap();
+        assert_eq!(
+            message.header.origin_timestamp,
+            CREATE_TIME.cast_unsigned() * 1000,
+            "clamped to the window start"
+        );
+        message.header.timestamp = 5_000_000;
+        assert_eq!(from_iggy(&message, 0).unwrap().timestamp, 5_000);
+    }
+
+    #[test]
+    fn 
given_an_epoch_timestamp_among_real_ones_when_round_tripped_should_stay_at_the_epoch()
 {
+        let records = [record_at(CREATE_TIME), record_at(EPOCH_TIMESTAMP)];
+        let message = to_iggy(&records[1], 
TimestampWindow::of(&records)).unwrap();
+        assert_eq!(from_iggy(&message, 0).unwrap().timestamp, EPOCH_TIMESTAMP);
+    }
+
+    #[test]
+    fn 
given_a_batch_when_windowed_should_start_at_its_earliest_real_timestamp() {
+        let records = [
+            record_at(NO_TIMESTAMP),
+            record_at(EPOCH_TIMESTAMP),
+            record_at(CREATE_TIME + 5),
+            record_at(CREATE_TIME),
+        ];
+        assert_eq!(
+            TimestampWindow::of(&records).start,
+            CREATE_TIME.cast_unsigned() * 1000
+        );
+        assert_eq!(TimestampWindow::of(&[record_at(NO_TIMESTAMP)]).start, 0);
+    }
+
     #[test]
     fn given_a_create_time_when_round_tripped_should_come_back_unchanged() {
-        let message = to_iggy(&record_at(CREATE_TIME)).unwrap();
+        let message = to_iggy_alone(&record_at(CREATE_TIME)).unwrap();
         assert_eq!(
             message.header.origin_timestamp,
             CREATE_TIME.cast_unsigned() * 1000
@@ -1449,7 +1585,7 @@ mod tests {
 
     #[test]
     fn given_an_epoch_timestamp_when_round_tripped_should_stay_at_the_epoch() {
-        let mut message = to_iggy(&record_at(EPOCH_TIMESTAMP)).unwrap();
+        let mut message = to_iggy_alone(&record_at(EPOCH_TIMESTAMP)).unwrap();
         message.header.timestamp = 5_000_000;
         assert_eq!(from_iggy(&message, 0).unwrap().timestamp, EPOCH_TIMESTAMP);
     }
@@ -1463,7 +1599,7 @@ mod tests {
             0,
             EPOCH_TIMESTAMP,
         );
-        let mut message = to_iggy(&original).unwrap();
+        let mut message = to_iggy_alone(&original).unwrap();
         assert!(is_enveloped(&message));
         message.header.timestamp = 5_000_000;
         assert_eq!(from_iggy(&message, 0).unwrap().timestamp, EPOCH_TIMESTAMP);
@@ -1471,14 +1607,14 @@ mod tests {
 
     #[test]
     fn given_no_timestamp_when_read_should_use_the_server_timestamp() {
-        let mut message = to_iggy(&record_at(NO_TIMESTAMP)).unwrap();
+        let mut message = to_iggy_alone(&record_at(NO_TIMESTAMP)).unwrap();
         message.header.timestamp = 5_000_000;
         assert_eq!(from_iggy(&message, 0).unwrap().timestamp, 5_000);
     }
 
     #[test]
     fn given_a_truncated_envelope_when_decoded_should_fail() {
-        let message = to_iggy(&record_with(Some(b""), Some(b"v"), 
&[])).unwrap();
+        let message = to_iggy_alone(&record_with(Some(b""), Some(b"v"), 
&[])).unwrap();
         assert!(matches!(
             decode_envelope(&[ENVELOPE_VERSION], &message.payload.slice(0..3)),
             Err(RecordCodecError::EnvelopeTruncated { .. })
@@ -1538,7 +1674,7 @@ mod tests {
             CREATE_TIME,
         );
         assert!(matches!(
-            to_iggy(&oversized),
+            to_iggy_alone(&oversized),
             Err(RecordCodecError::EnvelopeTooLarge { size })
                 if size == MAX_PAYLOAD_SIZE as usize + ENVELOPE_OVERHEAD
         ));
@@ -1586,7 +1722,7 @@ mod tests {
     fn 
given_an_uncompressed_batch_when_round_tripped_should_keep_every_record() {
         let mut records = vec![record_at_offset(0, b"1"), record_at_offset(1, 
b"2")];
         let mut encoded = encode_batch(&mut records).unwrap();
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         let decoded = decode_batch(&mut encoded, &budget, 
Zstd::Allowed).unwrap();
         assert_eq!(decoded.len(), 2);
         assert_eq!(decoded[1].value.as_deref(), Some(&b"2"[..]));
@@ -1621,7 +1757,7 @@ mod tests {
         blob.extend_from_slice(&encode_batch(&mut [record_at_offset(0, 
b"1")]).unwrap());
         blob.extend_from_slice(&encode_batch(&mut [record_at_offset(1, 
b"2")]).unwrap());
 
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         assert!(
             matches!(
                 decode_batch(&mut blob.freeze(), &budget, Zstd::Allowed),
@@ -1631,12 +1767,29 @@ mod tests {
         );
     }
 
+    #[test]
+    fn given_bytes_after_the_batch_when_decoded_should_reject() {
+        let mut blob = BytesMut::new();
+        blob.extend_from_slice(&encode_batch(&mut [record_at_offset(0, 
b"1")]).unwrap());
+        // Byte 16 of the tail is its magic. Not 2, so `decode_batch_info` 
stops there, no error.
+        blob.extend_from_slice(&[0; 32]);
+        let budget = partition_budget(1024, usize::MAX);
+
+        assert!(
+            matches!(
+                decode_batch(&mut blob.freeze(), &budget, Zstd::Allowed),
+                Err(RecordCodecError::BatchTrailingBytes(32))
+            ),
+            "the decoder never reads the tail, so the partition would answer 
success"
+        );
+    }
+
     #[test]
     fn given_more_records_than_the_batch_declares_when_decoded_should_reject() 
{
         let batch =
             encode_batch(&mut [record_at_offset(0, b"1"), record_at_offset(1, 
b"2")]).unwrap();
         let mut patched = patch_header(&batch, 57, &1i32.to_be_bytes());
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
 
         assert!(
             matches!(
@@ -1662,7 +1815,7 @@ mod tests {
     fn given_a_gzip_batch_when_decoded_should_read_it() {
         let records = vec![record_at_offset(0, b"compressed")];
         let mut encoded = encode_with(&records, Compression::Gzip);
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         let decoded = decode_batch(&mut encoded, &budget, 
Zstd::Allowed).unwrap();
         assert_eq!(decoded[0].value.as_deref(), Some(&b"compressed"[..]));
     }
@@ -1671,7 +1824,7 @@ mod tests {
     fn given_a_snappy_batch_when_decoded_should_read_it() {
         let records = vec![record_at_offset(0, b"compressed")];
         let mut encoded = encode_with(&records, Compression::Snappy);
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         let decoded = decode_batch(&mut encoded, &budget, 
Zstd::Allowed).unwrap();
         assert_eq!(decoded[0].value.as_deref(), Some(&b"compressed"[..]));
     }
@@ -1680,7 +1833,7 @@ mod tests {
     fn given_an_lz4_batch_when_decoded_should_read_it() {
         let records = vec![record_at_offset(0, b"compressed")];
         let mut encoded = encode_with(&records, Compression::Lz4);
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         let decoded = decode_batch(&mut encoded, &budget, 
Zstd::Allowed).unwrap();
         assert_eq!(decoded[0].value.as_deref(), Some(&b"compressed"[..]));
     }
@@ -1689,7 +1842,7 @@ mod tests {
     fn given_a_zstd_batch_when_decoded_should_read_it() {
         let records = vec![record_at_offset(0, b"compressed")];
         let mut encoded = encode_with(&records, Compression::Zstd);
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         let decoded = decode_batch(&mut encoded, &budget, 
Zstd::Allowed).unwrap();
         assert_eq!(decoded[0].value.as_deref(), Some(&b"compressed"[..]));
     }
@@ -1698,7 +1851,7 @@ mod tests {
     fn given_a_budget_smaller_than_the_batch_when_decoded_should_reject() {
         let records = vec![record_at_offset(0, &[b'x'; 512])];
         let mut encoded = encode_with(&records, Compression::Gzip);
-        let budget = DecompressionBudget::new(8, usize::MAX);
+        let budget = partition_budget(8, usize::MAX);
         assert!(matches!(
             decode_batch(&mut encoded, &budget, Zstd::Allowed),
             Err(RecordCodecError::BudgetExceeded { .. })
@@ -1708,7 +1861,7 @@ mod tests {
     #[test]
     fn 
given_a_spent_budget_when_a_partition_fits_alone_should_ask_for_a_retry() {
         // Enough for one partition, not for both: the budget is per request.
-        let budget = DecompressionBudget::new(400, usize::MAX);
+        let budget = partition_budget(400, usize::MAX);
         let big = || encode_batch(&mut [record_at_offset(0, &[b'x'; 
256])]).unwrap();
 
         assert!(decode_batch(&mut big(), &budget, Zstd::Allowed).is_ok());
@@ -1719,9 +1872,59 @@ mod tests {
     }
 
     #[test]
-    fn 
given_a_spent_budget_when_a_partition_is_too_large_alone_should_refuse_it() {
+    fn 
given_a_spent_budget_when_a_later_partition_arrives_should_refuse_it_undecoded()
 {
+        let budget = partition_budget(400, usize::MAX);
+        let big = || encode_batch(&mut [record_at_offset(0, &[b'x'; 
256])]).unwrap();
+        assert!(decode_batch(&mut big(), &budget, Zstd::Allowed).is_ok());
+        assert!(decode_batch(&mut big(), &budget, Zstd::Allowed).is_err());
+
+        // Malformed, so a decode would refuse it as a bad batch instead.
+        let mut garbage = Bytes::from_static(b"not a record batch");
+        assert!(
+            matches!(
+                decode_batch(&mut garbage, &budget, Zstd::Allowed),
+                Err(RecordCodecError::RequestBudgetSpent)
+            ),
+            "one spent request inflates nothing more, so it cannot hold a slot 
inflating"
+        );
+    }
+
+    #[test]
+    fn 
given_a_partition_too_large_alone_when_others_follow_should_leave_them_room() {
+        let partition = Allowance {
+            bytes: 64 * 1024,
+            records: usize::MAX,
+        };
+        let budget = DecompressionBudget::new(partition, partition.times(2));
+        let mut bomb = encode_with(
+            &[record_at_offset(0, &vec![0; 1024 * 1024])],
+            Compression::Gzip,
+        );
+        let normal = || {
+            encode_with(
+                &[record_at_offset(0, &vec![b'n'; 32 * 1024])],
+                Compression::Gzip,
+            )
+        };
+
+        assert!(matches!(
+            decode_batch(&mut bomb, &budget, Zstd::Allowed),
+            Err(RecordCodecError::BudgetExceeded { .. })
+        ));
+        assert!(
+            decode_batch(&mut normal(), &budget, Zstd::Allowed).is_ok(),
+            "the bomb stops at its own cap, so the request has room left"
+        );
+    }
+
+    #[test]
+    fn given_a_partition_too_large_alone_when_decoded_should_refuse_it() {
         // gzip writes 32 KiB at a time, so the refusal comes several writes 
before the end.
-        let budget = DecompressionBudget::new(64 * 1024, usize::MAX);
+        let partition = Allowance {
+            bytes: 64 * 1024,
+            records: usize::MAX,
+        };
+        let budget = DecompressionBudget::new(partition, partition.times(8));
         let mut first = encode_with(
             &[record_at_offset(0, &vec![b'a'; 60 * 1024])],
             Compression::Gzip,
@@ -1737,7 +1940,7 @@ mod tests {
                 decode_batch(&mut second, &budget, Zstd::Allowed),
                 Err(RecordCodecError::BudgetExceeded { limit, .. }) if limit 
== 64 * 1024
             ),
-            "a retry cannot help a partition over the whole budget"
+            "a retry cannot help a partition over its own cap"
         );
     }
 
@@ -1745,7 +1948,7 @@ mod tests {
     fn given_a_compression_bomb_when_decoded_should_stop_before_it_is_whole() {
         let bomb = vec![0u8; 4 * 1024 * 1024];
         let mut compressed = Bytes::from(encode_gzip(&bomb));
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
 
         let Err(error) = decompress(&mut compressed, Compression::Gzip, 
&budget) else {
             panic!("a 4 MB output against a 1 KB budget has to be refused");
@@ -1767,7 +1970,7 @@ mod tests {
         // A raw snappy stream whose leading varint claims u32::MAX bytes of 
output. `snap` reads
         // that length without allocating, so the declared size is checkable 
before the decode.
         let mut compressed = Bytes::from_static(&[0xff, 0xff, 0xff, 0xff, 
0x0f, 0x00]);
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         assert!(decompress(&mut compressed, Compression::Snappy, 
&budget).is_err());
         assert!(matches!(
             budget.reason.take(),
@@ -1779,7 +1982,7 @@ mod tests {
     fn given_a_record_count_past_the_frame_when_decoded_should_reject() {
         let batch = encode_batch(&mut [record_at_offset(0, b"v")]).unwrap();
         let mut patched = patch_header(&batch, 57, &i32::MAX.to_be_bytes());
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
 
         assert!(
             matches!(
@@ -1791,13 +1994,13 @@ mod tests {
     }
 
     #[test]
-    fn 
given_a_compressed_batch_when_preflighting_should_grant_it_the_whole_budget() {
+    fn 
given_a_compressed_batch_when_preflighting_should_grant_it_the_partition_cap() {
         // 100 tiny records declare more than the blob and the bytes left can 
hold, and fit the
-        // whole budget. Earlier partitions spending it must not make them 
look malformed.
-        // One offset for all, so the encoder writes one batch.
+        // partition cap. Earlier partitions spending the request must not 
make them look
+        // malformed. One offset for all, so the encoder writes one batch.
         let records: Vec<_> = (0..100).map(|_| record_at_offset(0, 
b"")).collect();
         let mut batch = encode_with(&records, Compression::Gzip);
-        let budget = DecompressionBudget::new(4096, usize::MAX);
+        let budget = partition_budget(4096, usize::MAX);
         let mut spend = encode_batch(&mut [record_at_offset(0, &[b'x'; 
4000])]).unwrap();
 
         assert!(decode_batch(&mut spend, &budget, Zstd::Allowed).is_ok());
@@ -1819,12 +2022,12 @@ mod tests {
         let mut batch = encode_batch(&mut [record_with(None, Some(b"v"), 
&headers)]).unwrap();
 
         // One slot for the record and two for its six headers.
-        let tight = DecompressionBudget::new(1024, 2);
+        let tight = partition_budget(1024, 2);
         assert!(matches!(
             decode_batch(&mut batch.clone(), &tight, Zstd::Allowed),
             Err(RecordCodecError::RecordBudgetExceeded { count: 3, limit: 2 })
         ));
-        let enough = DecompressionBudget::new(1024, 3);
+        let enough = partition_budget(1024, 3);
         assert!(decode_batch(&mut batch, &enough, Zstd::Allowed).is_ok());
     }
 
@@ -1887,7 +2090,7 @@ mod tests {
             .position(|window| window == [0x02, b'b', 0x02, b'y'])
             .unwrap();
         let mut batch = patch_header(&batch, at + 1, b"a");
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
 
         assert!(matches!(
             decode_batch(&mut batch, &budget, Zstd::Allowed),
@@ -1900,7 +2103,7 @@ mod tests {
         // The batch header says one record, so the record count bound passes. 
The count that
         // matters is the one inside the record, and no batch header reports 
it.
         let mut batch = batch_declaring_headers(i32::MAX);
-        let budget = DecompressionBudget::new(8 * 1024 * 1024, usize::MAX);
+        let budget = partition_budget(8 * 1024 * 1024, usize::MAX);
 
         assert!(
             matches!(
@@ -1916,7 +2119,7 @@ mod tests {
         // The same builder with a count the record can hold, so the scan 
cannot be passing the
         // test above by rejecting every hand-built batch.
         let mut batch = batch_declaring_headers(0);
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         assert_eq!(
             decode_batch(&mut batch, &budget, Zstd::Allowed)
                 .unwrap()
@@ -1934,7 +2137,7 @@ mod tests {
             &[("trace", Some(b"abc"))],
         )];
         let mut encoded = encode_batch(&mut with_headers).unwrap();
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
 
         let decoded = decode_batch(&mut encoded, &budget, 
Zstd::Allowed).unwrap();
         assert_eq!(
@@ -1949,7 +2152,7 @@ mod tests {
         // the documented 8 MiB, which an uncompressed batch has no claim on.
         let batch = encode_batch(&mut [record_at_offset(0, b"v")]).unwrap();
         let mut patched = patch_header(&batch, 57, &100_000i32.to_be_bytes());
-        let budget = DecompressionBudget::new(8 * 1024 * 1024, usize::MAX);
+        let budget = partition_budget(8 * 1024 * 1024, usize::MAX);
 
         assert!(matches!(
             decode_batch(&mut patched, &budget, Zstd::Allowed),
@@ -2014,7 +2217,7 @@ mod tests {
         let mut transactional = record_at_offset(0, b"v");
         transactional.transactional = true;
         let mut encoded = encode_with(&[transactional], Compression::None);
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
 
         assert!(matches!(
             decode_batch(&mut encoded, &budget, Zstd::Allowed),
@@ -2027,7 +2230,7 @@ mod tests {
         let mut control = record_at_offset(0, b"v");
         control.control = true;
         let mut encoded = encode_with(&[control], Compression::None);
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
 
         assert!(
             matches!(
@@ -2042,7 +2245,7 @@ mod tests {
     fn given_a_truncated_batch_when_decoded_should_reject() {
         let batch = encode_batch(&mut [record_at_offset(0, 
b"value")]).unwrap();
         let mut truncated = batch.slice(0..batch.len() - 4);
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         assert!(decode_batch(&mut truncated, &budget, Zstd::Allowed).is_err());
     }
 
@@ -2054,7 +2257,7 @@ mod tests {
         bytes[last] ^= 0xff;
         let mut corrupt = Bytes::from(bytes);
 
-        let budget = DecompressionBudget::new(1024, usize::MAX);
+        let budget = partition_budget(1024, usize::MAX);
         assert!(matches!(
             decode_batch(&mut corrupt, &budget, Zstd::Allowed),
             Err(RecordCodecError::Batch(_))
@@ -2064,7 +2267,7 @@ mod tests {
     #[test]
     fn 
given_an_earlier_overrun_when_a_later_batch_fails_should_not_reuse_the_budget_reason()
 {
         // A failed charge deducts nothing, so the budget still has room for 
the batch below.
-        let budget = DecompressionBudget::new(64, usize::MAX);
+        let budget = partition_budget(64, usize::MAX);
         let mut bomb = encode_with(&[record_at_offset(0, &[b'x'; 512])], 
Compression::Gzip);
         assert!(matches!(
             decode_batch(&mut bomb, &budget, Zstd::Allowed),
diff --git a/gateways/kafka/tests/produce_real_bridge_tests.rs 
b/gateways/kafka/tests/produce_real_bridge_tests.rs
index 97737c572..6172071a0 100644
--- a/gateways/kafka/tests/produce_real_bridge_tests.rs
+++ b/gateways/kafka/tests/produce_real_bridge_tests.rs
@@ -400,8 +400,43 @@ async fn 
given_timestamps_72_minutes_apart_when_handled_should_store_both() {
     );
     let messages = stored(&server, 0, 10).await;
     assert_eq!(messages.len(), 2);
-    let micros = u64::try_from(later * 1000).expect("positive");
-    assert_eq!(messages[1].header.origin_timestamp, micros);
+    let start = u64::try_from(CREATE_TIME * 1000).expect("positive");
+    assert_eq!(
+        messages[1].header.origin_timestamp,
+        start + u64::from(u32::MAX),
+        "clamped into the first record's window, so one send holds both"
+    );
+    let back = from_iggy(&messages[1], 1).expect("stored message must decode 
as a record");
+    assert_eq!(back.timestamp, later, "kafka.ts keeps the real one");
+}
+
+#[tokio::test]
+#[serial]
+async fn 
given_no_timestamp_among_stamped_records_when_handled_should_store_one_send() {
+    let data_dir = tempfile::tempdir().expect("tempdir");
+    let server = TestServer::spawn(data_dir.path()).await;
+    let state = gateway_with_topic(&server, 1).await;
+
+    let mut records = [record(0, None, b"a", &[]), record(1, None, b"b", &[])];
+    records[1].timestamp = -1;
+    let entries = [(0, batch(&records, Compression::None))];
+
+    assert_eq!(
+        produce(&state, 3, 1, TOPIC, &entries).await,
+        vec![(0, ERROR_NONE, 0)]
+    );
+    let messages = stored(&server, 0, 10).await;
+    assert_eq!(messages.len(), 2);
+    assert_eq!(
+        messages[1].header.origin_timestamp, 
messages[0].header.origin_timestamp,
+        "clamped to the window start, not stored as 0"
+    );
+    let back = from_iggy(&messages[1], 1).expect("stored message must decode 
as a record");
+    let server_millis = i64::try_from(messages[1].header.timestamp / 
1000).expect("fits");
+    assert_eq!(
+        back.timestamp, server_millis,
+        "no timestamp reads as the server's"
+    );
 }
 
 #[tokio::test]

Reply via email to