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]