use crate::ProducerRecord;
use crate::native::{
KafkaClientError, KafkaClientResult,
model::{KafkaTimestamp, OffsetRange, PayloadBatchBuilder},
profile::{self, ProfileBucket},
protocol::{CompressionCodec, Decoder, Encoder, compression},
};
const RECORD_BATCH_MAGIC_V2: i8 = 2;
const ATTR_COMPRESSION_MASK: i16 = 0x0007;
const ATTR_TIMESTAMP_TYPE_MASK: i16 = 0x0008;
const ATTR_CONTROL_BATCH_MASK: i16 = 0x0020;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct BatchIdentity {
pub(crate) producer_id: i64,
pub(crate) producer_epoch: i16,
pub(crate) base_sequence: i32,
}
pub(crate) fn encode_produce_record_batch(
records: &[&ProducerRecord],
compression_codec: CompressionCodec,
zstd_context: &mut compression::ZstdEncoderContext,
identity: Option<BatchIdentity>,
) -> KafkaClientResult<Vec<u8>> {
if records.is_empty() {
return Err(KafkaClientError::protocol(
"cannot encode an empty Kafka produce record batch",
));
}
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_err(|_| KafkaClientError::protocol("system clock is before the Unix epoch"))?
.as_millis();
let now_ms = i64::try_from(now_ms)
.map_err(|_| KafkaClientError::protocol("current timestamp exceeds int64"))?;
let base_timestamp = records
.iter()
.map(|record| record.timestamp.unwrap_or(now_ms))
.min()
.unwrap_or(now_ms);
let max_timestamp = records
.iter()
.map(|record| record.timestamp.unwrap_or(now_ms))
.max()
.unwrap_or(now_ms);
let mut record_data = Encoder::with_capacity(128);
for (offset_delta, source) in records.iter().copied().enumerate() {
let timestamp = source.timestamp.unwrap_or(now_ms);
let mut record = Encoder::with_capacity(
source
.key
.as_ref()
.map_or(0, bytes::Bytes::len)
.saturating_add(source.payload.as_ref().map_or(0, bytes::Bytes::len))
.saturating_add(64),
);
record.put_i8(0);
record.put_varint_i64(timestamp.saturating_sub(base_timestamp));
record.put_varint_i32(i32_len(offset_delta, "record offset delta")?);
put_nullable_varbytes(&mut record, source.key.as_deref())?;
put_nullable_varbytes(&mut record, source.payload.as_deref())?;
record.put_varint_i32(i32_len(source.headers.len(), "record header count")?);
for header in &source.headers {
record.put_varint_i32(i32_len(header.key.len(), "record header key")?);
record.put_raw(header.key.as_bytes());
put_nullable_varbytes(&mut record, header.value.as_deref())?;
}
let record = record.into_inner();
record_data.put_varint_i32(i32_len(record.len(), "record body")?);
record_data.put_raw(&record);
}
let record_data =
compression::compress_records(compression_codec, record_data.into_inner(), zstd_context)?;
let mut crc_body = Encoder::with_capacity(record_data.len().saturating_add(40));
crc_body.put_i16(compression_codec.attribute());
crc_body.put_i32(i32_len(
records.len().saturating_sub(1),
"record offset delta",
)?);
crc_body.put_i64(base_timestamp);
crc_body.put_i64(max_timestamp);
let (producer_id, producer_epoch, base_sequence) = match identity {
Some(identity) => (
identity.producer_id,
identity.producer_epoch,
identity.base_sequence,
),
None => (-1, -1, -1),
};
crc_body.put_i64(producer_id);
crc_body.put_i16(producer_epoch);
crc_body.put_i32(base_sequence);
crc_body.put_i32(i32_len(records.len(), "record count")?);
crc_body.put_raw(&record_data);
let crc_body = crc_body.into_inner();
let crc = crc32c::crc32c(&crc_body);
let mut body = Encoder::with_capacity(crc_body.len() + 9);
body.put_i32(-1);
body.put_i8(RECORD_BATCH_MAGIC_V2);
body.put_i32(crc as i32);
body.put_raw(&crc_body);
let body = body.into_inner();
let mut batch = Encoder::with_capacity(body.len() + 12);
batch.put_i64(0);
batch.put_i32(i32_len(body.len(), "record batch")?);
batch.put_raw(&body);
Ok(batch.into_inner())
}
fn put_nullable_varbytes(encoder: &mut Encoder, value: Option<&[u8]>) -> KafkaClientResult<()> {
match value {
Some(value) => {
encoder.put_varint_i32(i32_len(value.len(), "record bytes")?);
encoder.put_raw(value);
}
None => encoder.put_varint_i32(-1),
}
Ok(())
}
fn i32_len(value: usize, what: &'static str) -> KafkaClientResult<i32> {
i32::try_from(value)
.map_err(|_| KafkaClientError::protocol(format!("Kafka {what} exceeds int32")))
}
pub(crate) fn decode_record_batches(
topic: &str,
partition: i32,
min_offset: i64,
records: &[u8],
builder: &mut PayloadBatchBuilder,
) -> KafkaClientResult<RecordDecodeStats> {
let mut decoder = Decoder::new(records);
let mut stats = RecordDecodeStats::default();
while decoder.remaining() > 0 {
if decoder.remaining() < 12 {
return Err(KafkaClientError::protocol(
"truncated Kafka record batch header",
));
}
let base_offset = decoder.get_i64()?;
let batch_length = decoder.get_i32()?;
if batch_length < 0 {
return Err(KafkaClientError::protocol(
"negative Kafka record batch length",
));
}
if decoder.remaining() < batch_length as usize {
if stats.decoded == 0 {
return Err(KafkaClientError::protocol(format!(
"truncated Kafka record batch: need {batch_length} bytes, have {}",
decoder.remaining()
)));
}
break;
}
let batch = decoder.take(batch_length as usize)?;
stats.merge(decode_one_batch(
topic,
partition,
base_offset,
min_offset,
batch,
builder,
)?);
}
Ok(stats)
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub(crate) struct RecordDecodeStats {
pub(crate) decoded: usize,
pub(crate) min_offset: Option<i64>,
pub(crate) max_offset: Option<i64>,
}
impl RecordDecodeStats {
fn observe(&mut self, offset: i64) {
self.decoded += 1;
self.min_offset = Some(
self.min_offset
.map_or(offset, |current| current.min(offset)),
);
self.max_offset = Some(
self.max_offset
.map_or(offset, |current| current.max(offset)),
);
}
fn merge(&mut self, other: Self) {
self.decoded += other.decoded;
if let Some(offset) = other.min_offset {
self.min_offset = Some(
self.min_offset
.map_or(offset, |current| current.min(offset)),
);
}
if let Some(offset) = other.max_offset {
self.max_offset = Some(
self.max_offset
.map_or(offset, |current| current.max(offset)),
);
}
}
}
fn decode_one_batch(
topic: &str,
partition: i32,
base_offset: i64,
min_offset: i64,
batch: &[u8],
builder: &mut PayloadBatchBuilder,
) -> KafkaClientResult<RecordDecodeStats> {
if batch.len() < 49 {
return Err(KafkaClientError::protocol(
"truncated Kafka record batch body",
));
}
let mut decoder = Decoder::new(batch);
let _partition_leader_epoch = decoder.get_i32()?;
let magic = decoder.get_i8()?;
if magic != RECORD_BATCH_MAGIC_V2 {
return Err(KafkaClientError::unsupported(format!(
"Kafka record batch magic {magic}; KC-1 requires magic v2"
)));
}
let expected_crc = decoder.get_i32()? as u32;
let actual_crc = profile::measure(ProfileBucket::RecordDecode, || crc32c::crc32c(&batch[9..]));
if expected_crc != actual_crc {
return Err(KafkaClientError::protocol(format!(
"Kafka record batch CRC mismatch: expected {expected_crc:#010x}, got {actual_crc:#010x}"
)));
}
let attributes = decoder.get_i16()?;
let compression_codec = CompressionCodec::from_attribute(attributes & ATTR_COMPRESSION_MASK)?;
let control_batch = (attributes & ATTR_CONTROL_BATCH_MASK) != 0;
let log_append_time = (attributes & ATTR_TIMESTAMP_TYPE_MASK) != 0;
let _last_offset_delta = decoder.get_i32()?;
let base_timestamp = decoder.get_i64()?;
let max_timestamp = decoder.get_i64()?;
let _producer_id = decoder.get_i64()?;
let _producer_epoch = decoder.get_i16()?;
let _base_sequence = decoder.get_i32()?;
let record_count = decoder.get_i32()?;
if record_count < 0 {
return Err(KafkaClientError::protocol(
"negative Kafka record count in batch",
));
}
let decompressed;
let mut record_decoder = if compression_codec == CompressionCodec::None {
decoder
} else {
let compressed_records = decoder.take(decoder.remaining())?;
decompressed = profile::measure(ProfileBucket::RecordDecode, || {
compression::decompress_records(compression_codec, compressed_records)
})?;
Decoder::new(&decompressed)
};
let stats = if profile::enabled() {
decode_records_profiled(
&mut record_decoder,
record_count,
control_batch,
log_append_time,
base_timestamp,
max_timestamp,
base_offset,
min_offset,
topic,
partition,
builder,
)?
} else {
decode_records_fast(
&mut record_decoder,
record_count,
control_batch,
log_append_time,
base_timestamp,
max_timestamp,
base_offset,
min_offset,
topic,
partition,
builder,
)?
};
if !record_decoder.is_done() {
return Err(KafkaClientError::protocol(format!(
"Kafka record batch left {} trailing record bytes after decoding {record_count} records",
record_decoder.remaining()
)));
}
if let (Some(first_offset), Some(max_offset)) = (stats.min_offset, stats.max_offset) {
builder.observe_partition(
topic,
OffsetRange {
partition,
first_offset,
next_offset: max_offset + 1,
},
);
}
Ok(stats)
}
#[derive(Debug)]
struct DecodedRecord<'a> {
offset: i64,
timestamp: KafkaTimestamp,
payload: &'a [u8],
}
#[allow(clippy::too_many_arguments)]
fn decode_records_fast(
decoder: &mut Decoder<'_>,
record_count: i32,
control_batch: bool,
log_append_time: bool,
base_timestamp: i64,
max_timestamp: i64,
base_offset: i64,
min_offset: i64,
topic: &str,
partition: i32,
builder: &mut PayloadBatchBuilder,
) -> KafkaClientResult<RecordDecodeStats> {
let mut stats = RecordDecodeStats::default();
for _ in 0..record_count {
let Some((offset, timestamp, payload)) = decode_record_for_batch(
decoder,
control_batch,
log_append_time,
base_timestamp,
max_timestamp,
base_offset,
min_offset,
builder.include_timestamps(),
)?
else {
continue;
};
builder.push_record(topic, partition, offset, timestamp, payload);
stats.observe(offset);
}
Ok(stats)
}
#[allow(clippy::too_many_arguments)]
fn decode_records_profiled<'a>(
decoder: &mut Decoder<'a>,
record_count: i32,
control_batch: bool,
log_append_time: bool,
base_timestamp: i64,
max_timestamp: i64,
base_offset: i64,
min_offset: i64,
topic: &str,
partition: i32,
builder: &mut PayloadBatchBuilder,
) -> KafkaClientResult<RecordDecodeStats> {
let include_timestamps = builder.include_timestamps();
let mut decoded = Vec::with_capacity((record_count as usize).min(decoder.remaining()));
let mut stats = RecordDecodeStats::default();
profile::measure(ProfileBucket::RecordDecode, || -> KafkaClientResult<()> {
for _ in 0..record_count {
let Some((offset, timestamp, payload)) = decode_record_for_batch(
decoder,
control_batch,
log_append_time,
base_timestamp,
max_timestamp,
base_offset,
min_offset,
include_timestamps,
)?
else {
continue;
};
stats.observe(offset);
decoded.push(DecodedRecord {
offset,
timestamp,
payload,
});
}
Ok(())
})?;
profile::measure(ProfileBucket::PayloadCopy, || {
for record in decoded {
builder.push_record(
topic,
partition,
record.offset,
record.timestamp,
record.payload,
);
}
});
Ok(stats)
}
#[allow(clippy::too_many_arguments)]
fn decode_record_for_batch<'a>(
decoder: &mut Decoder<'a>,
control_batch: bool,
log_append_time: bool,
base_timestamp: i64,
max_timestamp: i64,
base_offset: i64,
min_offset: i64,
include_timestamps: bool,
) -> KafkaClientResult<Option<(i64, KafkaTimestamp, &'a [u8])>> {
let record_len = decoder.get_varint_i32()?;
if record_len < 0 {
return Err(KafkaClientError::protocol("negative Kafka record length"));
}
let record = decoder.take(record_len as usize)?;
if !include_timestamps {
let Some((offset_delta, payload)) = decode_record_payload_only(record)? else {
return Ok(None);
};
if control_batch {
return Ok(None);
}
let offset = base_offset.saturating_add(offset_delta as i64);
if offset < min_offset {
return Ok(None);
}
return Ok(Some((offset, KafkaTimestamp::NotAvailable, payload)));
}
let Some((offset_delta, timestamp_delta, payload)) = decode_record(record)? else {
return Ok(None);
};
if control_batch {
return Ok(None);
}
let timestamp = if !include_timestamps {
KafkaTimestamp::NotAvailable
} else if log_append_time {
KafkaTimestamp::LogAppendTime(max_timestamp)
} else {
KafkaTimestamp::CreateTime(base_timestamp.saturating_add(timestamp_delta))
};
let offset = base_offset.saturating_add(offset_delta as i64);
if offset < min_offset {
return Ok(None);
}
Ok(Some((offset, timestamp, payload)))
}
fn decode_record_payload_only(record: &[u8]) -> KafkaClientResult<Option<(i32, &[u8])>> {
let mut decoder = Decoder::new(record);
let _attributes = decoder.get_i8()?;
decoder.skip_varint()?;
let offset_delta = decoder.get_varint_i32()?;
let key_len = decoder.get_varint_i32()?;
if key_len >= 0 {
decoder.take(key_len as usize)?;
}
let payload_len = decoder.get_varint_i32()?;
let payload = if payload_len >= 0 {
decoder.take(payload_len as usize)?
} else {
&[]
};
let header_count = decoder.get_varint_i32()?;
if header_count < 0 {
return Err(KafkaClientError::protocol(
"negative Kafka record header count",
));
}
for _ in 0..header_count {
let key_len = decoder.get_varint_i32()?;
if key_len < 0 {
return Err(KafkaClientError::protocol(
"Kafka record header key was null",
));
}
decoder.take(key_len as usize)?;
let value_len = decoder.get_varint_i32()?;
if value_len >= 0 {
decoder.take(value_len as usize)?;
}
}
if !decoder.is_done() {
return Err(KafkaClientError::protocol(format!(
"Kafka record left {} trailing bytes",
decoder.remaining()
)));
}
Ok(Some((offset_delta, payload)))
}
fn decode_record(record: &[u8]) -> KafkaClientResult<Option<(i32, i64, &[u8])>> {
let mut decoder = Decoder::new(record);
let _attributes = decoder.get_i8()?;
let timestamp_delta = decoder.get_varint_i64()?;
let offset_delta = decoder.get_varint_i32()?;
let key_len = decoder.get_varint_i32()?;
if key_len >= 0 {
decoder.take(key_len as usize)?;
}
let payload_len = decoder.get_varint_i32()?;
let payload = if payload_len >= 0 {
decoder.take(payload_len as usize)?
} else {
&[]
};
let header_count = decoder.get_varint_i32()?;
if header_count < 0 {
return Err(KafkaClientError::protocol(
"negative Kafka record header count",
));
}
for _ in 0..header_count {
let key_len = decoder.get_varint_i32()?;
if key_len < 0 {
return Err(KafkaClientError::protocol(
"Kafka record header key was null",
));
}
decoder.take(key_len as usize)?;
let value_len = decoder.get_varint_i32()?;
if value_len >= 0 {
decoder.take(value_len as usize)?;
}
}
if !decoder.is_done() {
return Err(KafkaClientError::protocol(format!(
"Kafka record left {} trailing bytes",
decoder.remaining()
)));
}
Ok(Some((offset_delta, timestamp_delta, payload)))
}
#[cfg(test)]
pub(crate) fn encode_test_record_batch(
base_offset: i64,
payloads: &[&[u8]],
) -> KafkaClientResult<Vec<u8>> {
encode_test_record_batch_with_attributes(base_offset, payloads, 0)
}
#[cfg(test)]
fn encode_test_record_batch_with_attributes(
base_offset: i64,
payloads: &[&[u8]],
attributes: i16,
) -> KafkaClientResult<Vec<u8>> {
let mut body = Encoder::new();
body.put_i32(0);
body.put_i8(RECORD_BATCH_MAGIC_V2);
body.put_i32(0);
body.put_i16(attributes);
body.put_i32(payloads.len().saturating_sub(1) as i32);
body.put_i64(1_000);
body.put_i64(1_000);
body.put_i64(-1);
body.put_i16(-1);
body.put_i32(-1);
body.put_i32(payloads.len() as i32);
for (index, payload) in payloads.iter().enumerate() {
let mut record = Encoder::new();
record.put_i8(0);
record.put_varint_i64(0);
record.put_varint_i32(index as i32);
record.put_varint_i32(-1);
record.put_varint_i32(
i32::try_from(payload.len())
.map_err(|_| KafkaClientError::protocol("test payload too large"))?,
);
record.put_raw(payload);
record.put_varint_i32(0);
let record = record.into_inner();
body.put_varint_i32(
i32::try_from(record.len())
.map_err(|_| KafkaClientError::protocol("test record too large"))?,
);
body.put_raw(&record);
}
let mut body = body.into_inner();
let crc = crc32c::crc32c(&body[9..]);
body[5..9].copy_from_slice(&crc.to_be_bytes());
let mut batch = Encoder::new();
batch.put_i64(base_offset);
batch.put_i32(
i32::try_from(body.len())
.map_err(|_| KafkaClientError::protocol("test batch too large"))?,
);
batch.put_raw(&body);
Ok(batch.into_inner())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::native::model::PayloadBatchBuilder;
use bytes::Bytes;
#[test]
fn produce_record_batch_v2_round_trips_payloads_and_crc() {
let records = [
ProducerRecord::new("topic", Some(Bytes::from_static(b"a")))
.with_key(Some(Bytes::from_static(b"key")))
.with_timestamp(1_000)
.with_header("trace", Some(Bytes::from_static(b"one"))),
ProducerRecord::new("topic", Some(Bytes::from_static(b"bb"))).with_timestamp(1_005),
];
let record_refs = records.iter().collect::<Vec<_>>();
let mut zstd_context = compression::ZstdEncoderContext::new().expect("zstd context");
for codec in [
CompressionCodec::None,
CompressionCodec::Lz4,
CompressionCodec::Zstd,
] {
let bytes = encode_produce_record_batch(&record_refs, codec, &mut zstd_context, None)
.expect("encode produce batch");
let mut builder = PayloadBatchBuilder::with_capacity(8, 32, true);
let decoded = decode_record_batches("topic", 3, 0, &bytes, &mut builder)
.expect("decode produce batch");
assert_eq!(decoded.decoded, 2, "codec={}", codec.name());
let batch = builder
.finish(vec![crate::native::TopicPartition::new("topic", 3)])
.expect("batch");
assert_eq!(batch.payload(&batch.records()[0]), b"a");
assert_eq!(batch.payload(&batch.records()[1]), b"bb");
assert_eq!(batch.records()[0].offset, 0);
assert_eq!(batch.records()[1].offset, 1);
}
}
#[test]
fn produce_record_batch_writes_idempotent_header_and_stays_decodable() {
let records = [
ProducerRecord::new("topic", Some(Bytes::from_static(b"a"))),
ProducerRecord::new("topic", Some(Bytes::from_static(b"bb"))),
];
let record_refs = records.iter().collect::<Vec<_>>();
let mut zstd_context = compression::ZstdEncoderContext::new().expect("zstd context");
let identity = BatchIdentity {
producer_id: 4242,
producer_epoch: 7,
base_sequence: 100,
};
let bytes = encode_produce_record_batch(
&record_refs,
CompressionCodec::None,
&mut zstd_context,
Some(identity),
)
.expect("encode idempotent batch");
let header = &bytes[12..];
let producer_id = i64::from_be_bytes(header[31..39].try_into().expect("producer id"));
let producer_epoch = i16::from_be_bytes(header[39..41].try_into().expect("producer epoch"));
let base_sequence = i32::from_be_bytes(header[41..45].try_into().expect("base sequence"));
assert_eq!(producer_id, 4242);
assert_eq!(producer_epoch, 7);
assert_eq!(base_sequence, 100);
let sentinel = encode_produce_record_batch(
&record_refs,
CompressionCodec::None,
&mut zstd_context,
None,
)
.expect("encode non-idempotent batch");
let sentinel_header = &sentinel[12..];
assert_eq!(
i64::from_be_bytes(sentinel_header[31..39].try_into().expect("producer id")),
-1
);
assert_eq!(
i32::from_be_bytes(sentinel_header[41..45].try_into().expect("base sequence")),
-1
);
let mut builder = PayloadBatchBuilder::with_capacity(8, 32, false);
let decoded = decode_record_batches("topic", 0, 0, &bytes, &mut builder)
.expect("decode idempotent batch");
assert_eq!(decoded.decoded, 2);
}
#[test]
fn record_batch_decodes_payloads_and_offsets() {
let bytes = encode_test_record_batch(42, &[b"a", b"bb"]).expect("encode test batch");
let mut builder = PayloadBatchBuilder::with_capacity(8, 32, false);
let decoded =
decode_record_batches("topic", 3, 0, &bytes, &mut builder).expect("decode batch");
assert_eq!(decoded.decoded, 2);
assert_eq!(decoded.min_offset, Some(42));
assert_eq!(decoded.max_offset, Some(43));
let batch = builder
.finish(vec![crate::native::TopicPartition::new("topic", 3)])
.expect("batch");
assert_eq!(batch.records().len(), 2);
assert_eq!(batch.records()[0].offset, 42);
assert_eq!(batch.records()[1].offset, 43);
assert_eq!(batch.payload(&batch.records()[1]), b"bb");
assert_eq!(batch.watermarks()[0].offset, 44);
}
#[test]
fn payload_only_decode_skips_control_batches() {
let bytes =
encode_test_record_batch_with_attributes(42, &[b"control"], ATTR_CONTROL_BATCH_MASK)
.expect("encode control batch");
let mut builder = PayloadBatchBuilder::with_capacity(8, 32, false);
let decoded =
decode_record_batches("topic", 3, 0, &bytes, &mut builder).expect("decode batch");
assert_eq!(decoded.decoded, 0);
assert_eq!(decoded.max_offset, None);
assert!(
builder
.finish(vec![crate::native::TopicPartition::new("topic", 3)])
.is_none()
);
}
#[test]
fn record_batch_rejects_bad_crc() {
let records = [ProducerRecord::new(
"topic",
Some(Bytes::from_static(b"crc payload")),
)];
let record_refs = records.iter().collect::<Vec<_>>();
let mut zstd_context = compression::ZstdEncoderContext::new().expect("zstd context");
for codec in [
CompressionCodec::None,
CompressionCodec::Lz4,
CompressionCodec::Zstd,
] {
let mut bytes =
encode_produce_record_batch(&record_refs, codec, &mut zstd_context, None)
.expect("encode test batch");
let last = bytes.len() - 1;
bytes[last] ^= 0xff;
let mut builder = PayloadBatchBuilder::with_capacity(8, 32, false);
let error = decode_record_batches("topic", 0, 0, &bytes, &mut builder)
.expect_err("bad CRC must fail before record decompression");
assert!(matches!(error, KafkaClientError::Protocol(_)));
}
}
#[test]
fn record_batch_rejects_truncated_response() {
let bytes = [0_u8; 5];
let mut builder = PayloadBatchBuilder::with_capacity(8, 32, false);
assert!(decode_record_batches("topic", 0, 0, &bytes, &mut builder).is_err());
}
#[test]
fn record_batch_compression_failures_are_typed() {
for attributes in [1_i16, 2] {
let bytes = encode_test_record_batch_with_attributes(0, &[b"unsupported"], attributes)
.expect("encode unsupported batch");
let mut builder = PayloadBatchBuilder::with_capacity(8, 32, false);
let error = decode_record_batches("topic", 0, 0, &bytes, &mut builder)
.expect_err("unsupported codec must fail fast");
assert!(matches!(error, KafkaClientError::Unsupported(_)));
}
for attributes in [3_i16, 4] {
let bytes = encode_test_record_batch_with_attributes(0, &[b"not a frame"], attributes)
.expect("encode malformed compressed batch");
let mut builder = PayloadBatchBuilder::with_capacity(8, 32, false);
let error = decode_record_batches("topic", 0, 0, &bytes, &mut builder)
.expect_err("malformed frame must return a typed decompression error");
assert!(matches!(error, KafkaClientError::Decompression { .. }));
}
}
}