use std::collections::HashMap;
use bytes::{Buf, BufMut, Bytes, BytesMut};
use super::api::NodeEndpoint;
use super::buf;
use super::records::{self, RecordBatch};
use crate::error::{Error, Result};
pub const ACK_GAP: i8 = 0;
pub const ACK_ACCEPT: i8 = 1;
pub const ACK_RELEASE: i8 = 2;
pub const ACK_REJECT: i8 = 3;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareTopicPartitions {
pub topic_id: [u8; 16],
pub partitions: Vec<i32>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareGroupHeartbeatRequest {
pub group_id: String,
pub member_id: String,
pub member_epoch: i32,
pub rack_id: Option<String>,
pub subscribed_topic_names: Option<Vec<String>>,
}
impl ShareGroupHeartbeatRequest {
pub const LEAVE_GROUP_MEMBER_EPOCH: i32 = -1;
pub const JOIN_GROUP_MEMBER_EPOCH: i32 = 0;
#[must_use]
pub fn error_response(error_code: i16, throttle_time_ms: i32) -> ShareGroupHeartbeatResponse {
ShareGroupHeartbeatResponse {
throttle_time_ms,
error_code,
error_message: None,
member_id: None,
member_epoch: 0,
heartbeat_interval_ms: 0,
assignment: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareGroupHeartbeatResponse {
pub throttle_time_ms: i32,
pub error_code: i16,
pub error_message: Option<String>,
pub member_id: Option<String>,
pub member_epoch: i32,
pub heartbeat_interval_ms: i32,
pub assignment: Option<Vec<ShareTopicPartitions>>,
}
impl ShareGroupHeartbeatResponse {
#[must_use]
pub fn error_counts(&self) -> HashMap<i16, i32> {
HashMap::from([(self.error_code, 1)])
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AcknowledgementBatch {
pub first_offset: i64,
pub last_offset: i64,
pub types: Vec<i8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareFetchPartition {
pub partition: i32,
pub partition_max_bytes: i32,
pub acknowledgements: Vec<AcknowledgementBatch>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareFetchTopic {
pub topic_id: [u8; 16],
pub partitions: Vec<ShareFetchPartition>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareForgottenTopic {
pub topic_id: [u8; 16],
pub partitions: Vec<i32>,
}
pub struct ShareFetchRequest;
impl ShareFetchRequest {
#[must_use]
pub fn forgotten_topics(
forgotten: &[ShareForgottenTopic],
topic_names: &HashMap<[u8; 16], String>,
) -> Vec<([u8; 16], Option<String>, i32)> {
let mut to_forget = Vec::new();
for topic in forgotten {
let name = topic_names.get(&topic.topic_id).cloned();
for partition in &topic.partitions {
to_forget.push((topic.topic_id, name.clone(), *partition));
}
}
to_forget
}
#[must_use]
pub fn update_forgotten_data<I>(
forgotten: &[ShareForgottenTopic],
forget: I,
) -> Vec<ShareForgottenTopic>
where
I: IntoIterator<Item = ([u8; 16], i32)>,
{
let mut out = forgotten.to_vec();
let mut order: Vec<[u8; 16]> = Vec::new();
let mut by_id: HashMap<[u8; 16], Vec<i32>> = HashMap::new();
for (topic_id, partition) in forget {
by_id
.entry(topic_id)
.or_insert_with(|| {
order.push(topic_id);
Vec::new()
})
.push(partition);
}
for topic_id in order {
if let Some(partitions) = by_id.remove(&topic_id) {
out.push(ShareForgottenTopic {
topic_id,
partitions,
});
}
}
out
}
#[must_use]
pub fn share_fetch_data(
topics: &[ShareFetchTopic],
topic_names: &HashMap<[u8; 16], String>,
) -> HashMap<([u8; 16], Option<String>, i32), i32> {
let mut share_fetch_data = HashMap::new();
for topic in topics {
let name = topic_names.get(&topic.topic_id).cloned();
for partition in &topic.partitions {
let _prev = share_fetch_data.insert(
(topic.topic_id, name.clone(), partition.partition),
partition.partition_max_bytes,
);
}
}
share_fetch_data
}
#[must_use]
pub fn for_consumer<S, A>(
is_closing_share_session: bool,
fetch_size: i32,
send: S,
acknowledgements: A,
) -> Vec<ShareFetchTopic>
where
S: IntoIterator<Item = ([u8; 16], i32)>,
A: IntoIterator<Item = ([u8; 16], i32, Vec<AcknowledgementBatch>)>,
{
let ack_only_partition_max_bytes = if is_closing_share_session {
0
} else {
fetch_size
};
let mut topic_order = Vec::new();
let mut by_id = HashMap::new();
if is_closing_share_session {
drop(send);
} else {
for (topic_id, partition) in send {
let (part_order, part_map) = by_id.entry(topic_id).or_insert_with(|| {
topic_order.push(topic_id);
(Vec::new(), HashMap::new())
});
if part_map
.insert(
partition,
ShareFetchPartition {
partition,
partition_max_bytes: fetch_size,
acknowledgements: Vec::new(),
},
)
.is_none()
{
part_order.push(partition);
}
}
}
for (topic_id, partition, batches) in acknowledgements {
let (part_order, part_map) = by_id.entry(topic_id).or_insert_with(|| {
topic_order.push(topic_id);
(Vec::new(), HashMap::new())
});
if let Some(existing) = part_map.get_mut(&partition) {
existing.acknowledgements = batches;
} else {
part_order.push(partition);
let _prev = part_map.insert(
partition,
ShareFetchPartition {
partition,
partition_max_bytes: ack_only_partition_max_bytes,
acknowledgements: batches,
},
);
}
}
let mut topics = Vec::with_capacity(topic_order.len());
for topic_id in topic_order {
let Some((part_order, mut part_map)) = by_id.remove(&topic_id) else {
continue;
};
let mut partitions = Vec::with_capacity(part_order.len());
for partition in part_order {
if let Some(part) = part_map.remove(&partition) {
partitions.push(part);
}
}
topics.push(ShareFetchTopic {
topic_id,
partitions,
});
}
topics
}
pub fn error_response(
buf: &mut BytesMut,
version: i16,
error_code: i16,
throttle_time_ms: i32,
) -> crate::error::Result<()> {
encode_share_fetch_error_with_throttle(buf, version, error_code, throttle_time_ms)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AcquiredRange {
pub first_offset: i64,
pub last_offset: i64,
pub delivery_count: i16,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareFetchedPartition {
pub partition: i32,
pub error_code: i16,
pub error_message: Option<String>,
pub acknowledge_error_code: i16,
pub acknowledge_error_message: Option<String>,
pub current_leader_id: i32,
pub current_leader_epoch: i32,
pub records: Vec<RecordBatch>,
pub acquired: Vec<AcquiredRange>,
}
impl ShareFetchedPartition {
#[must_use]
pub fn partition_response(partition: i32, error_code: i16) -> Self {
Self {
partition,
error_code,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}
}
pub fn records_size(&self) -> Result<i32> {
if self.records.is_empty() {
return Ok(0);
}
let mut recs = BytesMut::new();
for batch in &self.records {
records::encode_record_batch(&mut recs, batch)?;
}
buf::i32_from_usize(recs.len())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareFetchedTopic {
pub topic_id: [u8; 16],
pub partitions: Vec<ShareFetchedPartition>,
}
pub struct ShareFetchResponse;
impl ShareFetchResponse {
#[must_use]
pub fn error_counts(error_code: i16, topics: &[ShareFetchedTopic]) -> HashMap<i16, i32> {
let mut counts = HashMap::new();
let count = counts.entry(error_code).or_insert(0);
*count += 1;
for topic in topics {
for partition in &topic.partitions {
let count = counts.entry(partition.error_code).or_insert(0);
*count += 1;
}
}
counts
}
#[must_use]
pub fn response_data(
topics: &[ShareFetchedTopic],
topic_names: &HashMap<[u8; 16], String>,
) -> HashMap<([u8; 16], String, i32), ShareFetchedPartition> {
let mut response_data = HashMap::new();
for topic in topics {
let Some(name) = topic_names.get(&topic.topic_id) else {
continue;
};
for partition in &topic.partitions {
let _prev = response_data.insert(
(topic.topic_id, name.clone(), partition.partition),
partition.clone(),
);
}
}
response_data
}
#[must_use]
pub fn to_message(
entries: &[([u8; 16], i32, ShareFetchedPartition)],
) -> Vec<ShareFetchedTopic> {
let mut topics: Vec<ShareFetchedTopic> = Vec::new();
for (topic_id, partition, body) in entries {
let mut body = body.clone();
body.partition = *partition;
if let Some(topic) = topics.iter_mut().find(|topic| topic.topic_id == *topic_id) {
topic.partitions.push(body);
} else {
topics.push(ShareFetchedTopic {
topic_id: *topic_id,
partitions: vec![body],
});
}
}
topics
}
pub fn size_of(
version: i16,
entries: &[([u8; 16], i32, ShareFetchedPartition)],
) -> Result<i32> {
let topics = Self::to_message(entries);
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &topics)?;
let n = buf
.len()
.checked_add(4)
.ok_or_else(|| Error::protocol("ShareFetchResponse.sizeOf overflow"))?;
buf::i32_from_usize(n)
}
pub fn of(
buf: &mut BytesMut,
version: i16,
error_code: i16,
throttle_time_ms: i32,
entries: &[([u8; 16], i32, ShareFetchedPartition)],
endpoints: &[NodeEndpoint],
) -> Result<()> {
let topics = Self::to_message(entries);
encode_share_fetch_response_full(
buf,
version,
&topics,
endpoints,
throttle_time_ms,
None,
0,
error_code,
)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareAcknowledgeResponsePartition {
pub partition: i32,
pub error_code: i16,
pub error_message: Option<String>,
pub current_leader_id: i32,
pub current_leader_epoch: i32,
}
impl ShareAcknowledgeResponsePartition {
#[must_use]
pub fn partition_response(partition: i32, error_code: i16) -> Self {
Self {
partition,
error_code,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareAcknowledgeResponseTopic {
pub topic_id: [u8; 16],
pub partitions: Vec<ShareAcknowledgeResponsePartition>,
}
pub struct ShareAcknowledgeResponse;
impl ShareAcknowledgeResponse {
#[must_use]
pub fn error_counts(
error_code: i16,
topics: &[ShareAcknowledgeResponseTopic],
) -> HashMap<i16, i32> {
let mut counts = HashMap::new();
let count = counts.entry(error_code).or_insert(0);
*count += 1;
for topic in topics {
for partition in &topic.partitions {
let count = counts.entry(partition.error_code).or_insert(0);
*count += 1;
}
}
counts
}
#[must_use]
pub fn to_message(
entries: &[([u8; 16], i32, ShareAcknowledgeResponsePartition)],
) -> Vec<ShareAcknowledgeResponseTopic> {
let mut topics: Vec<ShareAcknowledgeResponseTopic> = Vec::new();
for (topic_id, partition, body) in entries {
let mut body = body.clone();
body.partition = *partition;
if let Some(topic) = topics.iter_mut().find(|topic| topic.topic_id == *topic_id) {
topic.partitions.push(body);
} else {
topics.push(ShareAcknowledgeResponseTopic {
topic_id: *topic_id,
partitions: vec![body],
});
}
}
topics
}
pub fn of(
buf: &mut BytesMut,
version: i16,
error_code: i16,
throttle_time_ms: i32,
entries: &[([u8; 16], i32, ShareAcknowledgeResponsePartition)],
endpoints: &[NodeEndpoint],
) -> Result<()> {
let topics = Self::to_message(entries);
encode_share_acknowledge_topics_response_full(
buf,
version,
error_code,
&topics,
endpoints,
throttle_time_ms,
None,
)
}
}
fn share_group_heartbeat_flexible(version: i16) -> Result<bool> {
match version {
0..=1 => Ok(true),
other => Err(Error::protocol(format!(
"ShareGroupHeartbeat version {other} is not implemented"
))),
}
}
pub fn encode_share_group_heartbeat_request(
buf: &mut BytesMut,
version: i16,
req: &ShareGroupHeartbeatRequest,
) -> crate::error::Result<()> {
let flexible = share_group_heartbeat_flexible(version)?;
buf::put_string(buf, flexible, Some(&req.group_id))?;
buf::put_string(buf, flexible, Some(&req.member_id))?;
buf.put_i32(req.member_epoch);
buf::put_string(buf, flexible, req.rack_id.as_deref())?;
match &req.subscribed_topic_names {
None => buf::put_array_len(buf, flexible, None)?,
Some(names) => {
buf::put_array_len(buf, flexible, Some(names.len()))?;
for n in names {
buf::put_string(buf, flexible, Some(n))?;
}
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
Ok(())
}
pub fn decode_share_group_heartbeat_request<B: Buf>(
buf: &mut B,
version: i16,
) -> Result<ShareGroupHeartbeatRequest> {
let flexible = share_group_heartbeat_flexible(version)?;
let group_id = buf::get_string(buf, flexible)?.unwrap_or_default();
let member_id = buf::get_string(buf, flexible)?.unwrap_or_default();
let member_epoch = buf::get_i32(buf)?;
let rack_id = buf::get_string(buf, flexible)?;
let subscribed_topic_names = {
let n = buf::get_array_len(buf, flexible)?;
match n {
None => None,
Some(n) => {
let mut names = Vec::with_capacity(n);
for _ in 0..n {
names.push(buf::get_string(buf, flexible)?.unwrap_or_default());
}
Some(names)
}
}
};
if flexible {
buf::skip_tagged_fields(buf)?;
}
Ok(ShareGroupHeartbeatRequest {
group_id,
member_id,
member_epoch,
rack_id,
subscribed_topic_names,
})
}
pub fn encode_share_group_heartbeat_response(
buf: &mut BytesMut,
version: i16,
resp: &ShareGroupHeartbeatResponse,
) -> crate::error::Result<()> {
let flexible = share_group_heartbeat_flexible(version)?;
buf.put_i32(resp.throttle_time_ms);
buf.put_i16(resp.error_code);
buf::put_string(buf, flexible, resp.error_message.as_deref())?;
buf::put_string(buf, flexible, resp.member_id.as_deref())?;
buf.put_i32(resp.member_epoch);
buf.put_i32(resp.heartbeat_interval_ms);
match &resp.assignment {
None => buf.put_i8(-1),
Some(parts) => {
buf.put_i8(1);
buf::put_array_len(buf, flexible, Some(parts.len()))?;
for t in parts {
buf.extend_from_slice(&t.topic_id);
buf::put_array_len(buf, flexible, Some(t.partitions.len()))?;
for p in &t.partitions {
buf.put_i32(*p);
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
Ok(())
}
pub fn decode_share_group_heartbeat_response<B: Buf>(
buf: &mut B,
version: i16,
) -> Result<ShareGroupHeartbeatResponse> {
let flexible = share_group_heartbeat_flexible(version)?;
let throttle_time_ms = buf::get_i32(buf)?;
let error_code = buf::get_i16(buf)?;
let error_message = buf::get_string(buf, flexible)?;
let member_id = buf::get_string(buf, flexible)?;
let member_epoch = buf::get_i32(buf)?;
let heartbeat_interval_ms = buf::get_i32(buf)?;
let present = buf::get_i8(buf)?;
let assignment = if present < 0 {
None
} else if present != 1 {
return Err(Error::protocol(format!(
"invalid ShareGroupHeartbeat Assignment marker {present}"
)));
} else {
let n = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut parts = Vec::with_capacity(n);
for _ in 0..n {
let topic_id = buf::get_uuid(buf)?;
let pn = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut partitions = Vec::with_capacity(pn);
for _ in 0..pn {
partitions.push(buf::get_i32(buf)?);
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
parts.push(ShareTopicPartitions {
topic_id,
partitions,
});
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
Some(parts)
};
if flexible {
buf::skip_tagged_fields(buf)?;
}
Ok(ShareGroupHeartbeatResponse {
throttle_time_ms,
error_code,
error_message,
member_id,
member_epoch,
heartbeat_interval_ms,
assignment,
})
}
fn share_fetch_flexible(version: i16) -> Result<bool> {
match version {
0..=1 => Ok(true),
other => Err(Error::protocol(format!(
"ShareFetch version {other} is not implemented"
))),
}
}
fn encode_ack_batches(
buf: &mut BytesMut,
batches: &[AcknowledgementBatch],
) -> crate::error::Result<()> {
buf::put_array_len(buf, true, Some(batches.len()))?;
for b in batches {
buf.put_i64(b.first_offset);
buf.put_i64(b.last_offset);
buf::put_array_len(buf, true, Some(b.types.len()))?;
for t in &b.types {
buf.put_i8(*t);
}
buf::put_empty_tagged_fields(buf);
}
Ok(())
}
fn decode_ack_batches<B: Buf>(buf: &mut B) -> Result<Vec<AcknowledgementBatch>> {
let n = buf::get_array_len(buf, true)?.unwrap_or(0);
let mut out = Vec::with_capacity(n);
for _ in 0..n {
let first_offset = buf::get_i64(buf)?;
let last_offset = buf::get_i64(buf)?;
let tn = buf::get_array_len(buf, true)?.unwrap_or(0);
let mut types = Vec::with_capacity(tn);
for _ in 0..tn {
types.push(buf::get_i8(buf)?);
}
buf::skip_tagged_fields(buf)?;
out.push(AcknowledgementBatch {
first_offset,
last_offset,
types,
});
}
Ok(out)
}
#[expect(
clippy::too_many_arguments,
reason = "ShareFetch v0–v1 body fields are a single wire encode"
)]
pub fn encode_share_fetch_request(
buf: &mut BytesMut,
version: i16,
group_id: &str,
member_id: &str,
share_session_epoch: i32,
max_wait_ms: i32,
min_bytes: i32,
max_bytes: i32,
max_records: i32,
topics: &[ShareFetchTopic],
) -> crate::error::Result<()> {
encode_share_fetch_request_with_forgotten(
buf,
version,
group_id,
member_id,
share_session_epoch,
max_wait_ms,
min_bytes,
max_bytes,
max_records,
topics,
&[],
)
}
#[expect(
clippy::too_many_arguments,
reason = "ShareFetch ForgottenTopicsData is encoded with the rest of the v0–v1 body"
)]
pub fn encode_share_fetch_request_with_forgotten(
buf: &mut BytesMut,
version: i16,
group_id: &str,
member_id: &str,
share_session_epoch: i32,
max_wait_ms: i32,
min_bytes: i32,
max_bytes: i32,
max_records: i32,
topics: &[ShareFetchTopic],
forgotten: &[ShareForgottenTopic],
) -> crate::error::Result<()> {
encode_share_fetch_request_fields(
buf,
version,
group_id,
member_id,
share_session_epoch,
max_wait_ms,
min_bytes,
max_bytes,
max_records,
topics,
forgotten,
max_records,
)
}
#[expect(
clippy::too_many_arguments,
reason = "ShareFetch BatchSize is encoded with the rest of the v0–v1 body"
)]
pub fn encode_share_fetch_request_with_batch_size(
buf: &mut BytesMut,
version: i16,
group_id: &str,
member_id: &str,
share_session_epoch: i32,
max_wait_ms: i32,
min_bytes: i32,
max_bytes: i32,
max_records: i32,
topics: &[ShareFetchTopic],
batch_size: i32,
) -> crate::error::Result<()> {
encode_share_fetch_request_fields(
buf,
version,
group_id,
member_id,
share_session_epoch,
max_wait_ms,
min_bytes,
max_bytes,
max_records,
topics,
&[],
batch_size,
)
}
#[expect(
clippy::too_many_arguments,
reason = "ShareFetch request body fields are a single wire encode"
)]
fn encode_share_fetch_request_fields(
buf: &mut BytesMut,
version: i16,
group_id: &str,
member_id: &str,
share_session_epoch: i32,
max_wait_ms: i32,
min_bytes: i32,
max_bytes: i32,
max_records: i32,
topics: &[ShareFetchTopic],
forgotten: &[ShareForgottenTopic],
batch_size: i32,
) -> crate::error::Result<()> {
let flexible = share_fetch_flexible(version)?;
buf::put_string(buf, flexible, Some(group_id))?;
buf::put_string(buf, flexible, Some(member_id))?;
buf.put_i32(share_session_epoch);
buf.put_i32(max_wait_ms);
buf.put_i32(min_bytes);
buf.put_i32(max_bytes);
if version >= 1 {
buf.put_i32(max_records);
buf.put_i32(batch_size);
}
buf::put_array_len(buf, flexible, Some(topics.len()))?;
for t in topics {
buf.extend_from_slice(&t.topic_id);
buf::put_array_len(buf, flexible, Some(t.partitions.len()))?;
for p in &t.partitions {
buf.put_i32(p.partition);
if version == 0 {
buf.put_i32(p.partition_max_bytes);
}
encode_ack_batches(buf, &p.acknowledgements)?;
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
buf::put_array_len(buf, flexible, Some(forgotten.len()))?;
for t in forgotten {
buf.extend_from_slice(&t.topic_id);
buf::put_array_len(buf, flexible, Some(t.partitions.len()))?;
for p in &t.partitions {
buf.put_i32(*p);
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
Ok(())
}
#[expect(
clippy::type_complexity,
reason = "ShareFetch request decode returns group, member, epoch, max records, topics, forgotten, batch size, max wait, min bytes, and max bytes together"
)]
pub fn decode_share_fetch_request<B: Buf>(
buf: &mut B,
version: i16,
) -> Result<(
String,
String,
i32,
i32,
Vec<ShareFetchTopic>,
Vec<ShareForgottenTopic>,
i32,
i32,
i32,
i32,
)> {
let flexible = share_fetch_flexible(version)?;
let group_id = buf::get_string(buf, flexible)?.unwrap_or_default();
let member_id = buf::get_string(buf, flexible)?.unwrap_or_default();
let epoch = buf::get_i32(buf)?;
let max_wait_ms = buf::get_i32(buf)?;
let min_bytes = buf::get_i32(buf)?;
let max_bytes = buf::get_i32(buf)?;
let (max_records, batch_size) = if version >= 1 {
(buf::get_i32(buf)?, buf::get_i32(buf)?)
} else {
(0, 0)
};
let n = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut topics = Vec::with_capacity(n);
for _ in 0..n {
let topic_id = buf::get_uuid(buf)?;
let pn = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut partitions = Vec::with_capacity(pn);
for _ in 0..pn {
let partition = buf::get_i32(buf)?;
let partition_max_bytes = if version == 0 { buf::get_i32(buf)? } else { 0 };
let acknowledgements = decode_ack_batches(buf)?;
if flexible {
buf::skip_tagged_fields(buf)?;
}
partitions.push(ShareFetchPartition {
partition,
partition_max_bytes,
acknowledgements,
});
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
topics.push(ShareFetchTopic {
topic_id,
partitions,
});
}
let n = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut forgotten_out = Vec::with_capacity(n);
for _ in 0..n {
let topic_id = buf::get_uuid(buf)?;
let pn = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut partitions = Vec::with_capacity(pn);
for _ in 0..pn {
partitions.push(buf::get_i32(buf)?);
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
forgotten_out.push(ShareForgottenTopic {
topic_id,
partitions,
});
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
Ok((
group_id,
member_id,
epoch,
max_records,
topics,
forgotten_out,
batch_size,
max_wait_ms,
min_bytes,
max_bytes,
))
}
fn encode_leader(buf: &mut BytesMut, leader_id: i32, leader_epoch: i32) {
buf.put_i32(leader_id);
buf.put_i32(leader_epoch);
buf::put_empty_tagged_fields(buf);
}
fn decode_leader<B: Buf>(buf: &mut B) -> Result<(i32, i32)> {
let id = buf::get_i32(buf)?;
let epoch = buf::get_i32(buf)?;
buf::skip_tagged_fields(buf)?;
Ok((id, epoch))
}
pub fn encode_share_fetch_response(
buf: &mut BytesMut,
version: i16,
topics: &[ShareFetchedTopic],
) -> crate::error::Result<()> {
encode_share_fetch_response_with_endpoints(buf, version, topics, &[])
}
pub fn encode_share_fetch_response_with_endpoints(
buf: &mut BytesMut,
version: i16,
topics: &[ShareFetchedTopic],
endpoints: &[NodeEndpoint],
) -> crate::error::Result<()> {
encode_share_fetch_response_full(buf, version, topics, endpoints, 0, None, 15_000, 0)
}
pub fn encode_share_fetch_response_with_throttle(
buf: &mut BytesMut,
version: i16,
topics: &[ShareFetchedTopic],
throttle_time_ms: i32,
) -> crate::error::Result<()> {
encode_share_fetch_response_full(buf, version, topics, &[], throttle_time_ms, None, 15_000, 0)
}
pub fn encode_share_fetch_response_with_error_message(
buf: &mut BytesMut,
version: i16,
topics: &[ShareFetchedTopic],
error_message: Option<&str>,
) -> crate::error::Result<()> {
encode_share_fetch_response_full(buf, version, topics, &[], 0, error_message, 15_000, 0)
}
pub fn encode_share_fetch_response_with_acquisition_lock_timeout(
buf: &mut BytesMut,
version: i16,
topics: &[ShareFetchedTopic],
acquisition_lock_timeout_ms: i32,
) -> crate::error::Result<()> {
encode_share_fetch_response_full(
buf,
version,
topics,
&[],
0,
None,
acquisition_lock_timeout_ms,
0,
)
}
pub fn encode_share_fetch_response_with_error_code(
buf: &mut BytesMut,
version: i16,
topics: &[ShareFetchedTopic],
error_code: i16,
) -> crate::error::Result<()> {
encode_share_fetch_response_full(buf, version, topics, &[], 0, None, 15_000, error_code)
}
#[expect(
clippy::too_many_arguments,
reason = "ShareFetch response body fields are a single wire encode"
)]
fn encode_share_fetch_response_full(
buf: &mut BytesMut,
version: i16,
topics: &[ShareFetchedTopic],
endpoints: &[NodeEndpoint],
throttle_time_ms: i32,
error_message: Option<&str>,
acquisition_lock_timeout_ms: i32,
error_code: i16,
) -> crate::error::Result<()> {
let flexible = share_fetch_flexible(version)?;
buf.put_i32(throttle_time_ms);
buf.put_i16(error_code);
buf::put_string(buf, flexible, error_message)?;
if version >= 1 {
buf.put_i32(acquisition_lock_timeout_ms);
}
buf::put_array_len(buf, flexible, Some(topics.len()))?;
for t in topics {
buf.extend_from_slice(&t.topic_id);
buf::put_array_len(buf, flexible, Some(t.partitions.len()))?;
for p in &t.partitions {
buf.put_i32(p.partition);
buf.put_i16(p.error_code);
buf::put_string(buf, flexible, p.error_message.as_deref())?;
buf.put_i16(p.acknowledge_error_code);
buf::put_string(buf, flexible, p.acknowledge_error_message.as_deref())?;
encode_leader(buf, p.current_leader_id, p.current_leader_epoch);
let mut recs = BytesMut::new();
for batch in &p.records {
records::encode_record_batch(&mut recs, batch)?;
}
buf::put_bytes(buf, flexible, Some(&recs))?;
buf::put_array_len(buf, flexible, Some(p.acquired.len()))?;
for a in &p.acquired {
buf.put_i64(a.first_offset);
buf.put_i64(a.last_offset);
buf.put_i16(a.delivery_count);
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
super::api::put_compact_node_endpoints(buf, endpoints)?;
if flexible {
buf::put_empty_tagged_fields(buf);
}
Ok(())
}
#[expect(
clippy::type_complexity,
reason = "ShareFetch response decode returns topics, node endpoints, throttle, ErrorMessage, AcquisitionLockTimeoutMs, and ErrorCode together"
)]
pub fn decode_share_fetch_response<B: Buf>(
buf: &mut B,
version: i16,
) -> Result<(
Vec<ShareFetchedTopic>,
Vec<NodeEndpoint>,
i32,
Option<String>,
i32,
i16,
)> {
let flexible = share_fetch_flexible(version)?;
let throttle_time_ms = buf::get_i32(buf)?;
let error_code = buf::get_i16(buf)?;
let error_message = buf::get_string(buf, flexible)?;
let acquisition_lock_timeout_ms = if version >= 1 { buf::get_i32(buf)? } else { 0 };
let n = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut topics = Vec::with_capacity(n);
for _ in 0..n {
let topic_id = buf::get_uuid(buf)?;
let pn = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut partitions = Vec::with_capacity(pn);
for _ in 0..pn {
let partition = buf::get_i32(buf)?;
let error_code = buf::get_i16(buf)?;
let error_message = buf::get_string(buf, flexible)?;
let acknowledge_error_code = buf::get_i16(buf)?;
let acknowledge_error_message = buf::get_string(buf, flexible)?;
let (current_leader_id, current_leader_epoch) = decode_leader(buf)?;
let rec_opt = if flexible {
buf::take_compact_bytes(buf)?
} else {
buf::take_classic_bytes(buf)?
};
let rec_bytes = match rec_opt {
Some(b) => b,
None if version == 0 => Bytes::new(),
None => {
return Err(Error::protocol(
"non-nullable field records was serialized as null",
));
}
};
let records = if rec_bytes.is_empty() {
Vec::new()
} else {
let mut rec_buf = rec_bytes;
records::decode_record_batches(&mut rec_buf)?
};
let an = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut acquired = Vec::with_capacity(an);
for _ in 0..an {
let first_offset = buf::get_i64(buf)?;
let last_offset = buf::get_i64(buf)?;
let delivery_count = buf::get_i16(buf)?;
if flexible {
buf::skip_tagged_fields(buf)?;
}
acquired.push(AcquiredRange {
first_offset,
last_offset,
delivery_count,
});
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
partitions.push(ShareFetchedPartition {
partition,
error_code,
error_message,
acknowledge_error_code,
acknowledge_error_message,
current_leader_id,
current_leader_epoch,
records,
acquired,
});
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
topics.push(ShareFetchedTopic {
topic_id,
partitions,
});
}
let endpoints = super::api::get_compact_node_endpoints(buf)?;
if flexible {
buf::skip_tagged_fields(buf)?;
}
Ok((
topics,
endpoints,
throttle_time_ms,
error_message,
acquisition_lock_timeout_ms,
error_code,
))
}
fn share_acknowledge_flexible(version: i16) -> Result<bool> {
match version {
0..=1 => Ok(true),
other => Err(Error::protocol(format!(
"ShareAcknowledge version {other} is not implemented"
))),
}
}
pub fn encode_share_acknowledge_request(
buf: &mut BytesMut,
version: i16,
group_id: &str,
member_id: &str,
share_session_epoch: i32,
topic_id: [u8; 16],
partitions: &[(i32, Vec<AcknowledgementBatch>)],
) -> crate::error::Result<()> {
if partitions.is_empty() {
encode_share_acknowledge_topics(buf, version, group_id, member_id, share_session_epoch, &[])
} else {
encode_share_acknowledge_topics(
buf,
version,
group_id,
member_id,
share_session_epoch,
&[ShareAckTopic {
topic_id,
partitions: partitions.to_vec(),
}],
)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareAckTopic {
pub topic_id: [u8; 16],
pub partitions: Vec<(i32, Vec<AcknowledgementBatch>)>,
}
pub struct ShareAcknowledgeRequest;
impl ShareAcknowledgeRequest {
#[must_use]
pub fn for_consumer<I>(acknowledgements: I) -> Vec<ShareAckTopic>
where
I: IntoIterator<Item = ([u8; 16], i32, Vec<AcknowledgementBatch>)>,
{
let mut topic_order = Vec::new();
let mut by_id = HashMap::new();
for (topic_id, partition, batches) in acknowledgements {
let (part_order, part_map) = by_id.entry(topic_id).or_insert_with(|| {
topic_order.push(topic_id);
(Vec::new(), HashMap::new())
});
if part_map.insert(partition, batches).is_none() {
part_order.push(partition);
}
}
let mut topics = Vec::with_capacity(topic_order.len());
for topic_id in topic_order {
let Some((part_order, mut part_map)) = by_id.remove(&topic_id) else {
continue;
};
let mut partitions = Vec::with_capacity(part_order.len());
for partition in part_order {
if let Some(batches) = part_map.remove(&partition) {
partitions.push((partition, batches));
}
}
topics.push(ShareAckTopic {
topic_id,
partitions,
});
}
topics
}
pub fn error_response(
buf: &mut BytesMut,
version: i16,
error_code: i16,
throttle_time_ms: i32,
) -> crate::error::Result<()> {
encode_share_acknowledge_topics_response_with_throttle(
buf,
version,
error_code,
&[],
throttle_time_ms,
)
}
}
pub fn encode_share_acknowledge_topics(
buf: &mut BytesMut,
version: i16,
group_id: &str,
member_id: &str,
share_session_epoch: i32,
topics: &[ShareAckTopic],
) -> crate::error::Result<()> {
let flexible = share_acknowledge_flexible(version)?;
buf::put_string(buf, flexible, Some(group_id))?;
buf::put_string(buf, flexible, Some(member_id))?;
buf.put_i32(share_session_epoch);
buf::put_array_len(buf, flexible, Some(topics.len()))?;
for t in topics {
buf.extend_from_slice(&t.topic_id);
buf::put_array_len(buf, flexible, Some(t.partitions.len()))?;
for (partition, batches) in &t.partitions {
buf.put_i32(*partition);
encode_ack_batches(buf, batches)?;
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
Ok(())
}
pub fn encode_share_fetch_error(
buf: &mut BytesMut,
version: i16,
error_code: i16,
) -> crate::error::Result<()> {
encode_share_fetch_error_with_throttle(buf, version, error_code, 0)
}
pub fn encode_share_fetch_error_with_throttle(
buf: &mut BytesMut,
version: i16,
error_code: i16,
throttle_time_ms: i32,
) -> crate::error::Result<()> {
let flexible = share_fetch_flexible(version)?;
buf.put_i32(throttle_time_ms);
buf.put_i16(error_code);
buf::put_string(buf, flexible, None)?;
if version >= 1 {
buf.put_i32(0);
}
buf::put_array_len(buf, flexible, Some(0))?;
buf::put_array_len(buf, flexible, Some(0))?;
if flexible {
buf::put_empty_tagged_fields(buf);
}
Ok(())
}
#[expect(
clippy::type_complexity,
reason = "ack request is group, member, epoch, and topic-partition batches"
)]
pub fn decode_share_acknowledge_request<B: Buf>(
buf: &mut B,
version: i16,
) -> Result<(
String,
String,
i32,
Vec<([u8; 16], i32, Vec<AcknowledgementBatch>)>,
)> {
let flexible = share_acknowledge_flexible(version)?;
let group_id = buf::get_string(buf, flexible)?.unwrap_or_default();
let member_id = buf::get_string(buf, flexible)?.unwrap_or_default();
let epoch = buf::get_i32(buf)?;
let n = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut topics = Vec::new();
for _ in 0..n {
let topic_id = buf::get_uuid(buf)?;
let pn = buf::get_array_len(buf, flexible)?.unwrap_or(0);
for _ in 0..pn {
let partition = buf::get_i32(buf)?;
let batches = decode_ack_batches(buf)?;
if flexible {
buf::skip_tagged_fields(buf)?;
}
topics.push((topic_id, partition, batches));
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
Ok((group_id, member_id, epoch, topics))
}
pub fn encode_share_acknowledge_response(
buf: &mut BytesMut,
version: i16,
error_code: i16,
) -> crate::error::Result<()> {
encode_share_acknowledge_topics_response(buf, version, error_code, &[])
}
pub fn encode_share_acknowledge_topics_response(
buf: &mut BytesMut,
version: i16,
error_code: i16,
topics: &[ShareAcknowledgeResponseTopic],
) -> crate::error::Result<()> {
encode_share_acknowledge_topics_response_with_endpoints(buf, version, error_code, topics, &[])
}
pub fn encode_share_acknowledge_topics_response_with_endpoints(
buf: &mut BytesMut,
version: i16,
error_code: i16,
topics: &[ShareAcknowledgeResponseTopic],
endpoints: &[NodeEndpoint],
) -> crate::error::Result<()> {
encode_share_acknowledge_topics_response_full(
buf, version, error_code, topics, endpoints, 0, None,
)
}
pub fn encode_share_acknowledge_topics_response_with_throttle(
buf: &mut BytesMut,
version: i16,
error_code: i16,
topics: &[ShareAcknowledgeResponseTopic],
throttle_time_ms: i32,
) -> crate::error::Result<()> {
encode_share_acknowledge_topics_response_full(
buf,
version,
error_code,
topics,
&[],
throttle_time_ms,
None,
)
}
pub fn encode_share_acknowledge_topics_response_with_error_message(
buf: &mut BytesMut,
version: i16,
error_code: i16,
topics: &[ShareAcknowledgeResponseTopic],
error_message: Option<&str>,
) -> crate::error::Result<()> {
encode_share_acknowledge_topics_response_full(
buf,
version,
error_code,
topics,
&[],
0,
error_message,
)
}
fn encode_share_acknowledge_topics_response_full(
buf: &mut BytesMut,
version: i16,
error_code: i16,
topics: &[ShareAcknowledgeResponseTopic],
endpoints: &[NodeEndpoint],
throttle_time_ms: i32,
error_message: Option<&str>,
) -> crate::error::Result<()> {
let flexible = share_acknowledge_flexible(version)?;
buf.put_i32(throttle_time_ms);
buf.put_i16(error_code);
buf::put_string(buf, flexible, error_message)?;
buf::put_array_len(buf, flexible, Some(topics.len()))?;
for t in topics {
buf.extend_from_slice(&t.topic_id);
buf::put_array_len(buf, flexible, Some(t.partitions.len()))?;
for p in &t.partitions {
buf.put_i32(p.partition);
buf.put_i16(p.error_code);
buf::put_string(buf, flexible, p.error_message.as_deref())?;
encode_leader(buf, p.current_leader_id, p.current_leader_epoch);
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
if flexible {
buf::put_empty_tagged_fields(buf);
}
}
super::api::put_compact_node_endpoints(buf, endpoints)?;
if flexible {
buf::put_empty_tagged_fields(buf);
}
Ok(())
}
pub fn decode_share_acknowledge_response<B: Buf>(buf: &mut B, version: i16) -> Result<i16> {
let (error_code, ..) = decode_share_acknowledge_topics_response(buf, version)?;
Ok(error_code)
}
#[expect(
clippy::type_complexity,
reason = "ShareAcknowledge response decode returns error, topics, node endpoints, throttle, and top-level ErrorMessage together"
)]
pub fn decode_share_acknowledge_topics_response<B: Buf>(
buf: &mut B,
version: i16,
) -> Result<(
i16,
Vec<ShareAcknowledgeResponseTopic>,
Vec<NodeEndpoint>,
i32,
Option<String>,
)> {
let flexible = share_acknowledge_flexible(version)?;
let throttle_time_ms = buf::get_i32(buf)?;
let error_code = buf::get_i16(buf)?;
let error_message = buf::get_string(buf, flexible)?;
let n = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut topics = Vec::with_capacity(n);
for _ in 0..n {
let topic_id = buf::get_uuid(buf)?;
let pn = buf::get_array_len(buf, flexible)?.unwrap_or(0);
let mut partitions = Vec::with_capacity(pn);
for _ in 0..pn {
let partition = buf::get_i32(buf)?;
let part_error = buf::get_i16(buf)?;
let error_message = buf::get_string(buf, flexible)?;
let (current_leader_id, current_leader_epoch) = decode_leader(buf)?;
if flexible {
buf::skip_tagged_fields(buf)?;
}
partitions.push(ShareAcknowledgeResponsePartition {
partition,
error_code: part_error,
error_message,
current_leader_id,
current_leader_epoch,
});
}
if flexible {
buf::skip_tagged_fields(buf)?;
}
topics.push(ShareAcknowledgeResponseTopic {
topic_id,
partitions,
});
}
let endpoints = super::api::get_compact_node_endpoints(buf)?;
if flexible {
buf::skip_tagged_fields(buf)?;
}
Ok((
error_code,
topics,
endpoints,
throttle_time_ms,
error_message,
))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::protocol::records::Record;
use bytes::{Buf, Bytes};
use std::collections::HashMap;
#[test]
fn share_group_heartbeat_join_leave_roundtrip() {
let req = ShareGroupHeartbeatRequest {
group_id: "sg".into(),
member_id: "m1".into(),
member_epoch: ShareGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH,
rack_id: None,
subscribed_topic_names: Some(vec!["t".into()]),
};
let mut buf = BytesMut::new();
encode_share_group_heartbeat_request(&mut buf, 1, &req).unwrap();
let mut cur = &buf[..];
let decoded = decode_share_group_heartbeat_request(&mut cur, 1).unwrap();
assert!(!cur.has_remaining(), "v1 request leftover-empty");
assert_eq!(
decoded.member_epoch,
ShareGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH
);
assert_eq!(decoded.member_id, "m1");
assert_eq!(decoded.subscribed_topic_names, Some(vec!["t".into()]));
let resp = ShareGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: 0,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: Some(vec![ShareTopicPartitions {
topic_id: [1u8; 16],
partitions: vec![0],
}]),
};
buf.clear();
encode_share_group_heartbeat_response(&mut buf, 1, &resp).unwrap();
let mut cur = &buf[..];
assert_eq!(
decode_share_group_heartbeat_response(&mut cur, 1).unwrap(),
resp
);
assert!(!cur.has_remaining(), "v1 response leftover-empty");
let leave = ShareGroupHeartbeatRequest {
group_id: "sg".into(),
member_id: "m1".into(),
member_epoch: ShareGroupHeartbeatRequest::LEAVE_GROUP_MEMBER_EPOCH,
rack_id: None,
subscribed_topic_names: None,
};
buf.clear();
encode_share_group_heartbeat_request(&mut buf, 1, &leave).unwrap();
let mut cur = &buf[..];
assert_eq!(
decode_share_group_heartbeat_request(&mut cur, 1)
.unwrap()
.member_epoch,
ShareGroupHeartbeatRequest::LEAVE_GROUP_MEMBER_EPOCH
);
assert!(!cur.has_remaining(), "v1 leave leftover-empty");
}
#[test]
fn share_group_heartbeat_v0_matches_v1_and_does_not_speak_v2() {
let req = ShareGroupHeartbeatRequest {
group_id: "sg".into(),
member_id: "m1".into(),
member_epoch: ShareGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH,
rack_id: None,
subscribed_topic_names: Some(vec!["t".into()]),
};
let mut v0 = BytesMut::new();
encode_share_group_heartbeat_request(&mut v0, 0, &req).unwrap();
let mut v1 = BytesMut::new();
encode_share_group_heartbeat_request(&mut v1, 1, &req).unwrap();
assert_eq!(v0.as_ref(), v1.as_ref(), "v0 and v1 request bodies match");
let mut cur = v0.as_ref();
assert_eq!(
decode_share_group_heartbeat_request(&mut cur, 0).unwrap(),
req
);
assert!(!cur.has_remaining(), "v0 request leftover-empty");
let err = encode_share_group_heartbeat_request(&mut BytesMut::new(), 2, &req).unwrap_err();
assert!(
err.to_string().contains("not implemented"),
"v2 is not spoken, got {err}"
);
let mut empty: &[u8] = &[];
let err = decode_share_group_heartbeat_request(&mut empty, 2).unwrap_err();
assert!(
err.to_string().contains("not implemented"),
"v2 decode is not spoken, got {err}"
);
assert_eq!(crate::protocol::api_keys::pick_version(0, 0, 0, 1), Some(0));
assert_eq!(crate::protocol::api_keys::pick_version(1, 1, 0, 1), Some(1));
assert_eq!(crate::protocol::api_keys::pick_version(0, 1, 0, 1), Some(1));
assert_eq!(crate::protocol::api_keys::pick_version(2, 2, 0, 1), None);
assert_eq!(ShareGroupHeartbeatRequest::LEAVE_GROUP_MEMBER_EPOCH, -1);
assert_eq!(ShareGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH, 0);
let resp = ShareGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: 0,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: None,
};
v0.clear();
encode_share_group_heartbeat_response(&mut v0, 0, &resp).unwrap();
v1.clear();
encode_share_group_heartbeat_response(&mut v1, 1, &resp).unwrap();
assert_eq!(v0.as_ref(), v1.as_ref(), "v0 and v1 response bodies match");
let mut cur = v0.as_ref();
assert_eq!(
decode_share_group_heartbeat_response(&mut cur, 0).unwrap(),
resp
);
assert!(!cur.has_remaining(), "v0 response leftover-empty");
v0.clear();
let err = encode_share_group_heartbeat_response(&mut v0, 2, &resp).unwrap_err();
assert!(
err.to_string().contains("not implemented"),
"v2 response is not spoken, got {err}"
);
}
#[test]
fn share_group_heartbeat_request_rack_id_matches_java() {
let req = ShareGroupHeartbeatRequest {
group_id: "sg".into(),
member_id: "m1".into(),
member_epoch: ShareGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH,
rack_id: Some("az-east".into()),
subscribed_topic_names: Some(vec!["t".into()]),
};
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_group_heartbeat_request(&mut buf, version, &req).unwrap();
let mut cur = buf.as_ref();
let got = decode_share_group_heartbeat_request(&mut cur, version).unwrap();
assert_eq!(got.rack_id.as_deref(), Some("az-east"));
assert_eq!(got, req);
assert!(
cur.is_empty(),
"ShareGroupHeartbeat request v{version} RackId leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_group_heartbeat_request(&mut with, 0, &req).unwrap();
let none = ShareGroupHeartbeatRequest {
group_id: "sg".into(),
member_id: "m1".into(),
member_epoch: ShareGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH,
rack_id: None,
subscribed_topic_names: Some(vec!["t".into()]),
};
let mut omitted = BytesMut::new();
encode_share_group_heartbeat_request(&mut omitted, 0, &none).unwrap();
assert_ne!(
&with[..],
&omitted[..],
"v0 RackId is not always the JSON default null"
);
let got = decode_share_group_heartbeat_request(&mut omitted.as_ref(), 0).unwrap();
assert_eq!(got.rack_id, None);
let mut v1_with = BytesMut::new();
encode_share_group_heartbeat_request(&mut v1_with, 1, &req).unwrap();
assert_eq!(
&with[..],
&v1_with[..],
"v0 and v1 both write RackId (JSON 0+)"
);
}
#[test]
fn share_group_heartbeat_response_throttle_time_ms_matches_java() {
let zero = ShareGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: 0,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: None,
};
let with = ShareGroupHeartbeatResponse {
throttle_time_ms: 3_600_000,
error_code: 0,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: None,
};
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_group_heartbeat_response(&mut buf, version, &with).unwrap();
let mut cur = buf.as_ref();
let got = decode_share_group_heartbeat_response(&mut cur, version).unwrap();
assert_eq!(got, with);
assert_eq!(got.throttle_time_ms, 3_600_000);
assert!(
cur.is_empty(),
"ShareGroupHeartbeat v{version} ThrottleTimeMs leftover-empty"
);
}
let mut with_buf = BytesMut::new();
encode_share_group_heartbeat_response(&mut with_buf, 0, &with).unwrap();
let mut zero_buf = BytesMut::new();
encode_share_group_heartbeat_response(&mut zero_buf, 0, &zero).unwrap();
assert_ne!(
&with_buf[..],
&zero_buf[..],
"v0 ThrottleTimeMs is not always the JSON default 0"
);
let mut v1_with = BytesMut::new();
encode_share_group_heartbeat_response(&mut v1_with, 1, &with).unwrap();
assert_eq!(
&with_buf[..],
&v1_with[..],
"v0 and v1 both write ThrottleTimeMs (JSON 0+); ShareGroupHeartbeat has no AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_group_heartbeat_request_error_response_matches_java() {
let err = ShareGroupHeartbeatRequest::error_response(16, 3_600_000);
assert_eq!(err.throttle_time_ms, 3_600_000);
assert_eq!(err.error_code, 16);
assert!(err.error_message.is_none());
assert!(err.member_id.is_none());
assert_eq!(err.member_epoch, 0);
assert_eq!(err.heartbeat_interval_ms, 0);
assert!(err.assignment.is_none());
assert_eq!(
ShareGroupHeartbeatRequest::error_response(0, 0),
ShareGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: 0,
error_message: None,
member_id: None,
member_epoch: 0,
heartbeat_interval_ms: 0,
assignment: None,
}
);
leftover_share_group_heartbeat_error_response(0, &err);
leftover_share_group_heartbeat_error_response(
0,
&ShareGroupHeartbeatRequest::error_response(0, 0),
);
leftover_share_group_heartbeat_error_response(1, &err);
leftover_share_group_heartbeat_error_response(
1,
&ShareGroupHeartbeatRequest::error_response(0, 0),
);
}
fn leftover_share_group_heartbeat_error_response(
version: i16,
resp: &ShareGroupHeartbeatResponse,
) {
let mut buf = BytesMut::new();
encode_share_group_heartbeat_response(&mut buf, version, resp).unwrap();
let mut cur = buf.as_ref();
let got = decode_share_group_heartbeat_response(&mut cur, version).unwrap();
assert_eq!(got, *resp);
let empty = if resp.error_code == 0 && resp.throttle_time_ms == 0 {
"empty "
} else {
""
};
assert!(
cur.is_empty(),
"ShareGroupHeartbeat v{version} getErrorResponse {empty}leftover-empty; leftover {} bytes",
cur.len()
);
}
#[test]
fn share_group_heartbeat_response_error_counts_matches_java() {
let none = ShareGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: 0,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: Some(vec![ShareTopicPartitions {
topic_id: [7u8; 16],
partitions: vec![0, 1],
}]),
};
assert_eq!(
none.error_counts(),
HashMap::from([(0, 1)]),
"NONE is a singleton 1, not an empty map"
);
let full = ShareGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: crate::error::GROUP_MAX_SIZE_REACHED,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: Some(vec![ShareTopicPartitions {
topic_id: [7u8; 16],
partitions: vec![0, 1],
}]),
};
assert_eq!(
full.error_counts(),
HashMap::from([(crate::error::GROUP_MAX_SIZE_REACHED, 1)])
);
for version in 0..=1_i16 {
let mut resp = BytesMut::new();
encode_share_group_heartbeat_response(&mut resp, version, &full).unwrap();
let mut cur = &resp[..];
let decoded = decode_share_group_heartbeat_response(&mut cur, version).unwrap();
assert_eq!(
decoded.error_counts(),
HashMap::from([(crate::error::GROUP_MAX_SIZE_REACHED, 1)]),
"ShareGroupHeartbeat v{version} errorCounts must count the decoded code"
);
assert!(
cur.is_empty(),
"ShareGroupHeartbeat v{version} errorCounts leftover-empty; leftover {} bytes",
cur.len()
);
}
}
#[test]
fn share_fetch_and_acknowledge_roundtrip() {
let rec = Record {
offset: 0,
timestamp: 1,
key: None,
value: Some(Bytes::from_static(b"s")),
headers: vec![],
};
let topics = vec![ShareFetchedTopic {
topic_id: [0u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 0,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: vec![RecordBatch::from_records(vec![rec])],
acquired: vec![AcquiredRange {
first_offset: 0,
last_offset: 0,
delivery_count: 1,
}],
}],
}];
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, 1, &topics).unwrap();
let (decoded, ..) = decode_share_fetch_response(&mut &buf[..], 1).unwrap();
assert_eq!(decoded[0].partitions[0].acquired[0].first_offset, 0);
assert_eq!(
decoded[0].partitions[0].records[0].records[0]
.value
.as_deref(),
Some(&b"s"[..])
);
let req_topics = vec![ShareFetchTopic {
topic_id: [0u8; 16],
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 1024,
acknowledgements: vec![],
}],
}];
buf.clear();
encode_share_fetch_request(&mut buf, 1, "sg", "m1", 0, 10, 1, 1024, 16, &req_topics)
.unwrap();
let mut cur = &buf[..];
let (gid, mid, epoch, max_records, got, ..) =
decode_share_fetch_request(&mut cur, 1).unwrap();
assert_eq!(
(gid.as_str(), mid.as_str(), epoch, max_records),
("sg", "m1", 0, 16)
);
assert_eq!(got[0].partitions[0].partition, 0);
assert_eq!(
got[0].partitions[0].partition_max_bytes, 0,
"v1 omits PartitionMaxBytes; decode fills 0"
);
assert!(!cur.has_remaining(), "v1 request leftover-empty");
buf.clear();
encode_share_acknowledge_request(
&mut buf,
1,
"sg",
"m1",
1,
[0u8; 16],
&[(
0,
vec![AcknowledgementBatch {
first_offset: 0,
last_offset: 2,
types: vec![ACK_ACCEPT],
}],
)],
)
.unwrap();
let (gid, mid, _e, acks) = decode_share_acknowledge_request(&mut &buf[..], 1).unwrap();
assert_eq!(gid, "sg");
assert_eq!(mid, "m1");
assert_eq!(acks[0].2[0].types, vec![ACK_ACCEPT]);
assert_eq!(acks[0].2[0].last_offset, 2);
buf.clear();
encode_share_acknowledge_response(&mut buf, 1, 0).unwrap();
assert_eq!(
decode_share_acknowledge_response(&mut &buf[..], 1).unwrap(),
0
);
}
#[test]
fn share_acknowledge_encodes_several_partitions() {
let mut buf = BytesMut::new();
encode_share_acknowledge_request(
&mut buf,
1,
"sg",
"m1",
2,
[7u8; 16],
&[
(
0,
vec![AcknowledgementBatch {
first_offset: 1,
last_offset: 3,
types: vec![ACK_ACCEPT],
}],
),
(
1,
vec![AcknowledgementBatch {
first_offset: 8,
last_offset: 8,
types: vec![ACK_REJECT],
}],
),
],
)
.unwrap();
let (_gid, _mid, epoch, acks) = decode_share_acknowledge_request(&mut &buf[..], 1).unwrap();
assert_eq!(epoch, 2);
assert_eq!(acks.len(), 2);
assert_eq!(acks[0].1, 0);
assert_eq!(acks[1].1, 1);
assert_eq!(acks[1].2[0].types, vec![ACK_REJECT]);
}
#[test]
fn share_acknowledge_close_session_has_no_topics() {
let mut buf = BytesMut::new();
encode_share_acknowledge_request(&mut buf, 1, "sg", "m1", -1, [0u8; 16], &[]).unwrap();
let (_gid, _mid, epoch, acks) = decode_share_acknowledge_request(&mut &buf[..], 1).unwrap();
assert_eq!(epoch, -1);
assert!(acks.is_empty());
}
#[test]
fn share_acknowledge_v0_matches_v1_and_does_not_speak_v2() {
let partitions = [(
0,
vec![AcknowledgementBatch {
first_offset: 0,
last_offset: 1,
types: vec![ACK_ACCEPT],
}],
)];
let mut v0 = BytesMut::new();
encode_share_acknowledge_request(&mut v0, 0, "sg", "m1", 1, [0u8; 16], &partitions)
.unwrap();
let mut v1 = BytesMut::new();
encode_share_acknowledge_request(&mut v1, 1, "sg", "m1", 1, [0u8; 16], &partitions)
.unwrap();
assert_eq!(v0.as_ref(), v1.as_ref(), "v0 and v1 request bodies match");
let mut cur = v0.as_ref();
let (gid, mid, epoch, acks) = decode_share_acknowledge_request(&mut cur, 0).unwrap();
assert_eq!((gid.as_str(), mid.as_str(), epoch), ("sg", "m1", 1));
assert_eq!(acks[0].2[0].types, vec![ACK_ACCEPT]);
assert!(!cur.has_remaining(), "v0 request leftover-empty");
let err = encode_share_acknowledge_request(
&mut BytesMut::new(),
2,
"sg",
"m1",
1,
[0u8; 16],
&partitions,
)
.unwrap_err();
assert!(
err.to_string().contains("not implemented"),
"v2 is not spoken, got {err}"
);
let mut empty: &[u8] = &[];
let err = decode_share_acknowledge_request(&mut empty, 2).unwrap_err();
assert!(
err.to_string().contains("not implemented"),
"v2 decode is not spoken, got {err}"
);
assert_eq!(crate::protocol::api_keys::pick_version(0, 0, 0, 1), Some(0));
assert_eq!(crate::protocol::api_keys::pick_version(1, 1, 0, 1), Some(1));
assert_eq!(crate::protocol::api_keys::pick_version(0, 1, 0, 1), Some(1));
assert_eq!(crate::protocol::api_keys::pick_version(2, 2, 0, 1), None);
v0.clear();
encode_share_acknowledge_response(&mut v0, 0, 0).unwrap();
v1.clear();
encode_share_acknowledge_response(&mut v1, 1, 0).unwrap();
assert_eq!(v0.as_ref(), v1.as_ref(), "v0 and v1 response bodies match");
let mut cur = v0.as_ref();
assert_eq!(decode_share_acknowledge_response(&mut cur, 0).unwrap(), 0);
assert!(!cur.has_remaining(), "v0 response leftover-empty");
v0.clear();
let err = encode_share_acknowledge_response(&mut v0, 2, 0).unwrap_err();
assert!(
err.to_string().contains("not implemented"),
"v2 response is not spoken, got {err}"
);
}
#[test]
fn share_fetch_error_roundtrip() {
let mut buf = BytesMut::new();
encode_share_fetch_error(&mut buf, 1, crate::error::INVALID_SHARE_SESSION_EPOCH).unwrap();
let mut cur = &buf[..];
let _th = crate::protocol::buf::get_i32(&mut cur).unwrap();
let err = crate::protocol::buf::get_i16(&mut cur).unwrap();
assert_eq!(err, crate::error::INVALID_SHARE_SESSION_EPOCH);
}
#[test]
fn share_fetch_error_throttle_time_ms_matches_java() {
let err = crate::error::INVALID_SHARE_SESSION_EPOCH;
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_error_with_throttle(&mut buf, version, err, 3_600_000).unwrap();
let mut cur = buf.as_ref();
assert_eq!(crate::protocol::buf::get_i32(&mut cur).unwrap(), 3_600_000);
assert_eq!(crate::protocol::buf::get_i16(&mut cur).unwrap(), err);
assert_eq!(
crate::protocol::buf::get_string(&mut cur, true).unwrap(),
None
);
if version >= 1 {
assert_eq!(
crate::protocol::buf::get_i32(&mut cur).unwrap(),
0,
"ShareFetch error v1 AcquisitionLockTimeoutMs stays 0"
);
}
assert_eq!(
crate::protocol::buf::get_array_len(&mut cur, true)
.unwrap()
.unwrap_or(0),
0
);
assert_eq!(
crate::protocol::buf::get_array_len(&mut cur, true)
.unwrap()
.unwrap_or(0),
0
);
crate::protocol::buf::skip_tagged_fields(&mut cur).unwrap();
assert!(
cur.is_empty(),
"ShareFetch error v{version} ThrottleTimeMs leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_error_with_throttle(&mut with, 0, err, 3_600_000).unwrap();
let mut zero = BytesMut::new();
encode_share_fetch_error_with_throttle(&mut zero, 0, err, 0).unwrap();
assert_ne!(
&with[..],
&zero[..],
"v0 error ThrottleTimeMs is not always the JSON default 0"
);
let mut conv = BytesMut::new();
encode_share_fetch_error(&mut conv, 0, err).unwrap();
assert_eq!(
&conv[..],
&zero[..],
"encode_share_fetch_error still writes ThrottleTimeMs 0"
);
assert_eq!(
&with[..4],
&3_600_000i32.to_be_bytes(),
"ShareFetch error ThrottleTimeMs is the first field"
);
assert_eq!(
&with[4..6],
&err.to_be_bytes(),
"ShareFetch error ErrorCode is at bytes 4–5"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_error_with_throttle(&mut v1_with, 1, err, 3_600_000).unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v1 error adds AcquisitionLockTimeoutMs 0 after ErrorMessage"
);
}
#[test]
fn share_fetch_request_error_response_matches_java() {
let err = crate::error::INVALID_SHARE_SESSION_EPOCH;
leftover_share_fetch_error_response(0, err, 3_600_000);
leftover_share_fetch_error_response(0, 0, 0);
leftover_share_fetch_error_response(1, err, 3_600_000);
leftover_share_fetch_error_response(1, 0, 0);
let mut named = BytesMut::new();
ShareFetchRequest::error_response(&mut named, 0, err, 0).unwrap();
let mut conv = BytesMut::new();
encode_share_fetch_error(&mut conv, 0, err).unwrap();
assert_eq!(
&named[..],
&conv[..],
"encode_share_fetch_error still writes ThrottleTimeMs 0"
);
let mut with = BytesMut::new();
ShareFetchRequest::error_response(&mut with, 0, err, 3_600_000).unwrap();
assert_ne!(
&with[..],
&conv[..],
"getErrorResponse ThrottleTimeMs is not always the JSON default 0"
);
}
fn leftover_share_fetch_error_response(version: i16, error_code: i16, throttle_time_ms: i32) {
let mut buf = BytesMut::new();
ShareFetchRequest::error_response(&mut buf, version, error_code, throttle_time_ms).unwrap();
let mut cur = buf.as_ref();
let (topics, endpoints, got_throttle, error_message, lock, got_error) =
decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got_error, error_code);
assert!(topics.is_empty());
assert!(endpoints.is_empty());
assert_eq!(got_throttle, throttle_time_ms);
assert!(error_message.is_none());
assert_eq!(lock, 0, "Java ShareFetchResponse.of last argument is 0");
let empty = if error_code == 0 && throttle_time_ms == 0 {
"empty "
} else {
""
};
assert!(
cur.is_empty(),
"ShareFetch v{version} getErrorResponse {empty}leftover-empty; leftover {} bytes",
cur.len()
);
}
#[test]
fn share_fetch_v0_omits_v1_fields_and_does_not_speak_v2() {
let req_topics = vec![ShareFetchTopic {
topic_id: [0u8; 16],
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 1024,
acknowledgements: vec![],
}],
}];
let mut v0 = BytesMut::new();
encode_share_fetch_request(&mut v0, 0, "sg", "m1", 0, 10, 1, 1024, 16, &req_topics)
.unwrap();
let mut v1 = BytesMut::new();
encode_share_fetch_request(&mut v1, 1, "sg", "m1", 0, 10, 1, 1024, 16, &req_topics)
.unwrap();
assert_ne!(
v0.as_ref(),
v1.as_ref(),
"v0 PartitionMaxBytes and v1 MaxRecords/BatchSize differ on the wire"
);
let mut cur = v0.as_ref();
let (gid, mid, epoch, max_records, got, ..) =
decode_share_fetch_request(&mut cur, 0).unwrap();
assert_eq!(
(gid.as_str(), mid.as_str(), epoch, max_records),
("sg", "m1", 0, 0),
"v0 omits MaxRecords; decode fills 0"
);
assert_eq!(got[0].partitions[0].partition, 0);
assert_eq!(
got[0].partitions[0].partition_max_bytes, 1024,
"v0 stores PartitionMaxBytes"
);
assert!(!cur.has_remaining(), "v0 request leftover-empty");
let err = encode_share_fetch_request(
&mut BytesMut::new(),
2,
"sg",
"m1",
0,
10,
1,
1024,
16,
&req_topics,
)
.unwrap_err();
assert!(
err.to_string().contains("not implemented"),
"v2 is not spoken, got {err}"
);
let mut empty: &[u8] = &[];
let err = decode_share_fetch_request(&mut empty, 2).unwrap_err();
assert!(
err.to_string().contains("not implemented"),
"v2 decode is not spoken, got {err}"
);
assert_eq!(crate::protocol::api_keys::pick_version(0, 0, 0, 1), Some(0));
assert_eq!(crate::protocol::api_keys::pick_version(1, 1, 0, 1), Some(1));
assert_eq!(crate::protocol::api_keys::pick_version(0, 1, 0, 1), Some(1));
assert_eq!(crate::protocol::api_keys::pick_version(2, 2, 0, 1), None);
let resp = vec![ShareFetchedTopic {
topic_id: [0u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 0,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
v0.clear();
encode_share_fetch_response(&mut v0, 0, &resp).unwrap();
v1.clear();
encode_share_fetch_response(&mut v1, 1, &resp).unwrap();
assert_ne!(
v0.as_ref(),
v1.as_ref(),
"v1 AcquisitionLockTimeoutMs is absent on v0"
);
let mut cur = v0.as_ref();
let (decoded, endpoints, ..) = decode_share_fetch_response(&mut cur, 0).unwrap();
assert_eq!(decoded, resp);
assert!(endpoints.is_empty());
assert!(!cur.has_remaining(), "v0 response leftover-empty");
v0.clear();
let err = encode_share_fetch_response(&mut v0, 2, &resp).unwrap_err();
assert!(
err.to_string().contains("not implemented"),
"v2 response is not spoken, got {err}"
);
}
#[test]
fn share_fetch_response_node_endpoints_matches_java() {
let topics = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
let endpoints = [crate::protocol::api::NodeEndpoint {
node_id: 3,
host: "h".into(),
port: 1,
rack: Some("r".into()),
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response_with_endpoints(&mut buf, version, &topics, &endpoints)
.unwrap();
let mut cur = buf.as_ref();
let (got, eps, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got, topics);
assert_eq!(eps, endpoints);
assert!(
cur.is_empty(),
"ShareFetch v{version} NodeEndpoints leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_response_with_endpoints(&mut with, 0, &topics, &endpoints).unwrap();
let mut empty = BytesMut::new();
encode_share_fetch_response_with_endpoints(&mut empty, 0, &topics, &[]).unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareFetch NodeEndpoints is not always empty"
);
let mut conv = BytesMut::new();
encode_share_fetch_response(&mut conv, 0, &topics).unwrap();
assert_eq!(
&conv[..],
&empty[..],
"encode_share_fetch_response still writes empty NodeEndpoints"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_response_with_endpoints(&mut v1_with, 1, &topics, &endpoints).unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 NodeEndpoints share compact layout; v1 still adds AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_fetch_response_current_leader_matches_java() {
let defaults = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
let with_leader = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 2,
current_leader_epoch: 7,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &with_leader).unwrap();
let mut cur = buf.as_ref();
let (got, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got, with_leader);
assert!(
cur.is_empty(),
"ShareFetch v{version} CurrentLeader leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_response(&mut with, 0, &with_leader).unwrap();
let mut empty = BytesMut::new();
encode_share_fetch_response(&mut empty, 0, &defaults).unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareFetch CurrentLeader is not always 0/0"
);
let pr = ShareFetchedPartition::partition_response(0, 6);
assert_eq!(pr.current_leader_id, 0);
assert_eq!(pr.current_leader_epoch, 0);
let mut conv = BytesMut::new();
encode_share_fetch_response(
&mut conv,
0,
&[ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![pr],
}],
)
.unwrap();
assert_eq!(
&conv[..],
&empty[..],
"partition_response still writes CurrentLeader 0/0"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_response(&mut v1_with, 1, &with_leader).unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 CurrentLeader share layout; v1 still adds AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_fetch_response_error_message_matches_java() {
let defaults = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
let with_msg = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: Some("e".into()),
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
let with_empty = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: Some(String::new()),
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &with_msg).unwrap();
let mut cur = buf.as_ref();
let (got, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got, with_msg);
assert!(
cur.is_empty(),
"ShareFetch v{version} ErrorMessage leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_response(&mut with, 0, &with_msg).unwrap();
let mut empty = BytesMut::new();
encode_share_fetch_response(&mut empty, 0, &defaults).unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareFetch ErrorMessage is not always null"
);
let mut empty_present = BytesMut::new();
encode_share_fetch_response(&mut empty_present, 0, &with_empty).unwrap();
assert_ne!(
&empty_present[..],
&empty[..],
"empty-but-present ErrorMessage is not JSON null"
);
let mut cur = empty_present.as_ref();
let (got, ..) = decode_share_fetch_response(&mut cur, 0).unwrap();
assert_eq!(got, with_empty);
assert!(cur.is_empty(), "ShareFetch v0 ErrorMessage leftover-empty");
let pr = ShareFetchedPartition::partition_response(0, 6);
assert!(pr.error_message.is_none());
let mut conv = BytesMut::new();
encode_share_fetch_response(
&mut conv,
0,
&[ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![pr],
}],
)
.unwrap();
assert_eq!(
&conv[..],
&empty[..],
"partition_response still writes ErrorMessage null"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_response(&mut v1_with, 1, &with_msg).unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 ErrorMessage share compact layout; v1 still adds AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_fetch_response_acknowledge_error_code_matches_java() {
let defaults = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
let with_ack = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: crate::error::INVALID_RECORD_STATE,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &with_ack).unwrap();
let mut cur = buf.as_ref();
let (got, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got, with_ack);
assert!(
cur.is_empty(),
"ShareFetch v{version} AcknowledgeErrorCode leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_response(&mut with, 0, &with_ack).unwrap();
let mut empty = BytesMut::new();
encode_share_fetch_response(&mut empty, 0, &defaults).unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareFetch AcknowledgeErrorCode is not always 0"
);
let pr = ShareFetchedPartition::partition_response(0, 6);
assert_eq!(pr.acknowledge_error_code, 0);
let mut conv = BytesMut::new();
encode_share_fetch_response(
&mut conv,
0,
&[ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![pr],
}],
)
.unwrap();
assert_eq!(
&conv[..],
&empty[..],
"partition_response still writes AcknowledgeErrorCode 0"
);
assert_eq!(
ShareFetchResponse::error_counts(0, &with_ack),
HashMap::from([(0, 1), (6, 1)]),
"errorCounts uses ErrorCode, not AcknowledgeErrorCode"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_response(&mut v1_with, 1, &with_ack).unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 AcknowledgeErrorCode share layout; v1 still adds AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_fetch_response_acknowledge_error_message_matches_java() {
let defaults = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
let with_msg = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: Some("e".into()),
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
let with_empty = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: Some(String::new()),
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
let with_fetch_msg = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: Some("e".into()),
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &with_msg).unwrap();
let mut cur = buf.as_ref();
let (got, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got, with_msg);
assert!(
cur.is_empty(),
"ShareFetch v{version} AcknowledgeErrorMessage leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_response(&mut with, 0, &with_msg).unwrap();
let mut empty = BytesMut::new();
encode_share_fetch_response(&mut empty, 0, &defaults).unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareFetch AcknowledgeErrorMessage is not always null"
);
let mut fetch_msg = BytesMut::new();
encode_share_fetch_response(&mut fetch_msg, 0, &with_fetch_msg).unwrap();
assert_ne!(
&with[..],
&fetch_msg[..],
"AcknowledgeErrorMessage is not fetch ErrorMessage"
);
let mut empty_present = BytesMut::new();
encode_share_fetch_response(&mut empty_present, 0, &with_empty).unwrap();
assert_ne!(
&empty_present[..],
&empty[..],
"empty-but-present AcknowledgeErrorMessage is not JSON null"
);
let mut cur = empty_present.as_ref();
let (got, ..) = decode_share_fetch_response(&mut cur, 0).unwrap();
assert_eq!(got, with_empty);
assert!(
cur.is_empty(),
"ShareFetch v0 AcknowledgeErrorMessage leftover-empty"
);
let pr = ShareFetchedPartition::partition_response(0, 6);
assert!(pr.acknowledge_error_message.is_none());
let mut conv = BytesMut::new();
encode_share_fetch_response(
&mut conv,
0,
&[ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![pr],
}],
)
.unwrap();
assert_eq!(
&conv[..],
&empty[..],
"partition_response still writes AcknowledgeErrorMessage null"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_response(&mut v1_with, 1, &with_msg).unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 AcknowledgeErrorMessage share compact layout; v1 still adds AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_fetch_response_throttle_time_ms_matches_java() {
let topics = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response_with_throttle(&mut buf, version, &topics, 3_600_000)
.unwrap();
let mut cur = buf.as_ref();
let (got, endpoints, throttle, ..) =
decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got, topics);
assert!(endpoints.is_empty());
assert_eq!(throttle, 3_600_000);
assert!(
cur.is_empty(),
"ShareFetch v{version} ThrottleTimeMs leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_response_with_throttle(&mut with, 0, &topics, 3_600_000).unwrap();
let mut zero = BytesMut::new();
encode_share_fetch_response_with_throttle(&mut zero, 0, &topics, 0).unwrap();
assert_ne!(
&with[..],
&zero[..],
"v0 ThrottleTimeMs is not always the JSON default 0"
);
let mut conv = BytesMut::new();
encode_share_fetch_response(&mut conv, 0, &topics).unwrap();
assert_eq!(
&conv[..],
&zero[..],
"encode_share_fetch_response still writes ThrottleTimeMs 0"
);
let mut endpoints_zero = BytesMut::new();
encode_share_fetch_response_with_endpoints(&mut endpoints_zero, 0, &topics, &[]).unwrap();
assert_eq!(
&endpoints_zero[..],
&zero[..],
"encode_share_fetch_response_with_endpoints still writes ThrottleTimeMs 0"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_response_with_throttle(&mut v1_with, 1, &topics, 3_600_000).unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 both write ThrottleTimeMs (JSON 0+); v1 still adds AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_fetch_response_top_level_error_message_matches_java() {
let topics = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
let with_part = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: Some("e".into()),
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response_with_error_message(&mut buf, version, &topics, Some("e"))
.unwrap();
let mut cur = buf.as_ref();
let (got, endpoints, throttle, msg, ..) =
decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got, topics);
assert!(endpoints.is_empty());
assert_eq!(throttle, 0);
assert_eq!(msg.as_deref(), Some("e"));
assert!(
cur.is_empty(),
"ShareFetch v{version} top-level ErrorMessage leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_response_with_error_message(&mut with, 0, &topics, Some("e")).unwrap();
let mut empty = BytesMut::new();
encode_share_fetch_response_with_error_message(&mut empty, 0, &topics, None).unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareFetch top-level ErrorMessage is not always null"
);
let mut conv = BytesMut::new();
encode_share_fetch_response(&mut conv, 0, &topics).unwrap();
assert_eq!(
&conv[..],
&empty[..],
"encode_share_fetch_response still writes top-level ErrorMessage null"
);
let mut part_msg = BytesMut::new();
encode_share_fetch_response(&mut part_msg, 0, &with_part).unwrap();
assert_ne!(
&with[..],
&part_msg[..],
"top-level ErrorMessage is not partition ErrorMessage"
);
let mut empty_present = BytesMut::new();
encode_share_fetch_response_with_error_message(&mut empty_present, 0, &topics, Some(""))
.unwrap();
assert_ne!(
&empty_present[..],
&empty[..],
"empty-but-present top-level ErrorMessage is not JSON null"
);
let mut cur = empty_present.as_ref();
let (got, .., msg, _, _) = decode_share_fetch_response(&mut cur, 0).unwrap();
assert_eq!(got, topics);
assert_eq!(msg.as_deref(), Some(""));
assert!(
cur.is_empty(),
"ShareFetch v0 top-level ErrorMessage leftover-empty"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_response_with_error_message(&mut v1_with, 1, &topics, Some("e"))
.unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 top-level ErrorMessage share compact layout; v1 still adds AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_fetch_response_acquisition_lock_timeout_ms_matches_java() {
let topics = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response_with_acquisition_lock_timeout(
&mut buf, version, &topics, 3_600_000,
)
.unwrap();
let mut cur = buf.as_ref();
let (got, endpoints, throttle, msg, lock, ..) =
decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got, topics);
assert!(endpoints.is_empty());
assert_eq!(throttle, 0);
assert_eq!(msg, None);
if version >= 1 {
assert_eq!(lock, 3_600_000);
} else {
assert_eq!(
lock, 0,
"ShareFetch v{version} omits AcquisitionLockTimeoutMs even when the body has a non-zero value"
);
}
assert!(
cur.is_empty(),
"ShareFetch v{version} AcquisitionLockTimeoutMs leftover-empty"
);
}
let mut v0_with = BytesMut::new();
encode_share_fetch_response_with_acquisition_lock_timeout(
&mut v0_with,
0,
&topics,
3_600_000,
)
.unwrap();
let mut v0_zero = BytesMut::new();
encode_share_fetch_response_with_acquisition_lock_timeout(&mut v0_zero, 0, &topics, 0)
.unwrap();
assert_eq!(
&v0_with[..],
&v0_zero[..],
"v0 omits AcquisitionLockTimeoutMs even when the body has a non-zero value"
);
let mut conv_v0 = BytesMut::new();
encode_share_fetch_response(&mut conv_v0, 0, &topics).unwrap();
assert_eq!(
&conv_v0[..],
&v0_zero[..],
"encode_share_fetch_response v0 has no AcquisitionLockTimeoutMs"
);
let mut with = BytesMut::new();
encode_share_fetch_response_with_acquisition_lock_timeout(&mut with, 1, &topics, 3_600_000)
.unwrap();
let mut zero = BytesMut::new();
encode_share_fetch_response_with_acquisition_lock_timeout(&mut zero, 1, &topics, 0)
.unwrap();
assert_ne!(
&with[..],
&zero[..],
"v1 AcquisitionLockTimeoutMs is not always the generated Java default 0"
);
let mut conv = BytesMut::new();
encode_share_fetch_response(&mut conv, 1, &topics).unwrap();
let mut fifteen = BytesMut::new();
encode_share_fetch_response_with_acquisition_lock_timeout(&mut fifteen, 1, &topics, 15_000)
.unwrap();
assert_eq!(
&conv[..],
&fifteen[..],
"encode_share_fetch_response still writes AcquisitionLockTimeoutMs 15000"
);
assert_ne!(
&conv[..],
&with[..],
"encode_share_fetch_response still writes 15000, not the helper value"
);
let mut endpoints_fifteen = BytesMut::new();
encode_share_fetch_response_with_endpoints(&mut endpoints_fifteen, 1, &topics, &[])
.unwrap();
assert_eq!(
&endpoints_fifteen[..],
&fifteen[..],
"encode_share_fetch_response_with_endpoints still writes AcquisitionLockTimeoutMs 15000"
);
let mut throttle_zero = BytesMut::new();
encode_share_fetch_response_with_throttle(&mut throttle_zero, 1, &topics, 0).unwrap();
assert_eq!(
&throttle_zero[..],
&fifteen[..],
"encode_share_fetch_response_with_throttle still writes AcquisitionLockTimeoutMs 15000"
);
let mut throttle_same = BytesMut::new();
encode_share_fetch_response_with_throttle(&mut throttle_same, 1, &topics, 3_600_000)
.unwrap();
assert_ne!(
&with[..],
&throttle_same[..],
"AcquisitionLockTimeoutMs is not ThrottleTimeMs"
);
assert_eq!(
&with[..4],
&0i32.to_be_bytes(),
"ShareFetch success-path ThrottleTimeMs stays 0 on this helper"
);
assert_eq!(
&with[4..6],
&0i16.to_be_bytes(),
"ShareFetch success-path ErrorCode stays 0"
);
assert_eq!(
with[6], 0,
"ShareFetch success-path ErrorMessage stays null (compact)"
);
assert_eq!(
&with[7..11],
&3_600_000i32.to_be_bytes(),
"ShareFetch v1 AcquisitionLockTimeoutMs is after ErrorMessage"
);
assert_ne!(
&v0_with[..],
&with[..],
"v1 adds AcquisitionLockTimeoutMs after ErrorMessage"
);
}
#[test]
fn share_fetch_response_error_code_matches_java() {
let err = crate::error::INVALID_SHARE_SESSION_EPOCH;
let topics = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition {
partition: 0,
error_code: 6,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response_with_error_code(&mut buf, version, &topics, err).unwrap();
let mut cur = buf.as_ref();
let (got, endpoints, throttle, msg, lock, error_code) =
decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(got, topics);
assert!(endpoints.is_empty());
assert_eq!(throttle, 0);
assert_eq!(msg, None);
if version >= 1 {
assert_eq!(lock, 15_000);
} else {
assert_eq!(lock, 0);
}
assert_eq!(error_code, err);
assert_ne!(
error_code, 6,
"top-level ErrorCode is not partition ErrorCode"
);
assert!(
cur.is_empty(),
"ShareFetch v{version} ErrorCode leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_response_with_error_code(&mut with, 0, &topics, err).unwrap();
let mut zero = BytesMut::new();
encode_share_fetch_response_with_error_code(&mut zero, 0, &topics, 0).unwrap();
assert_ne!(
&with[..],
&zero[..],
"v0 ErrorCode is not always the JSON default 0"
);
let mut conv = BytesMut::new();
encode_share_fetch_response(&mut conv, 0, &topics).unwrap();
assert_eq!(
&conv[..],
&zero[..],
"encode_share_fetch_response still writes ErrorCode 0"
);
let mut lock_zero = BytesMut::new();
encode_share_fetch_response_with_acquisition_lock_timeout(&mut lock_zero, 0, &topics, 0)
.unwrap();
assert_eq!(
&lock_zero[..],
&zero[..],
"encode_share_fetch_response_with_acquisition_lock_timeout still writes ErrorCode 0"
);
assert_eq!(
&with[..4],
&0i32.to_be_bytes(),
"ShareFetch success-path ThrottleTimeMs stays 0 on this helper"
);
assert_eq!(
&with[4..6],
&err.to_be_bytes(),
"ShareFetch ErrorCode is at bytes 4–5"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_response_with_error_code(&mut v1_with, 1, &topics, err).unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 both write ErrorCode (JSON 0+); v1 still adds AcquisitionLockTimeoutMs"
);
let mut err_path = BytesMut::new();
encode_share_fetch_error(&mut err_path, 1, err).unwrap();
assert_ne!(
&v1_with[..],
&err_path[..],
"success-path ErrorCode helper is not encode_share_fetch_error (lock 15000 vs 0, topics vs empty)"
);
let mut err_cur = err_path.as_ref();
let (got, endpoints, throttle, msg, lock, error_code) =
decode_share_fetch_response(&mut err_cur, 1).unwrap();
assert!(got.is_empty());
assert!(endpoints.is_empty());
assert_eq!(throttle, 0);
assert_eq!(msg, None);
assert_eq!(
lock, 0,
"ShareFetch error v1 AcquisitionLockTimeoutMs stays 0"
);
assert_eq!(error_code, err);
assert!(
err_cur.is_empty(),
"ShareFetch error v1 ErrorCode leftover-empty"
);
}
#[test]
fn share_fetch_partition_response_leftover_empty() {
let err = ShareFetchedPartition::partition_response(3, crate::error::UNKNOWN_TOPIC_ID);
assert_eq!(
err,
ShareFetchedPartition {
partition: 3,
error_code: crate::error::UNKNOWN_TOPIC_ID,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}
);
let topics = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![
ShareFetchedPartition::partition_response(0, crate::error::UNKNOWN_TOPIC_ID),
ShareFetchedPartition::partition_response(3, crate::error::UNKNOWN_TOPIC_ID),
],
}];
for version in [0i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &topics).unwrap();
let mut cur = buf.as_ref();
let (decoded, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(decoded, topics);
assert!(
!cur.has_remaining(),
"ShareFetch v{version} partitionResponse leftover-empty; leftover {} bytes",
cur.remaining()
);
}
let empty: Vec<ShareFetchedTopic> = Vec::new();
for version in [0i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &empty).unwrap();
let mut cur = buf.as_ref();
let (decoded, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(decoded, empty);
assert!(
!cur.has_remaining(),
"ShareFetch v{version} empty partitionResponse leftover-empty; leftover {} bytes",
cur.remaining()
);
}
}
#[test]
fn share_fetch_response_records_size_matches_java() {
let rec = Record {
offset: 0,
timestamp: 1,
key: None,
value: Some(Bytes::from_static(b"s")),
headers: vec![],
};
let mut part = ShareFetchedPartition::partition_response(0, 0);
assert_eq!(part.records_size().unwrap(), 0);
part.records = vec![RecordBatch::from_records(vec![rec])];
let size = part.records_size().unwrap();
assert!(size > 0, "non-empty recordsSize must be the blob length");
let mut recs = BytesMut::new();
let batch = part.records.first().expect("one batch");
records::encode_record_batch(&mut recs, batch).unwrap();
assert_eq!(
size,
buf::i32_from_usize(recs.len()).unwrap(),
"recordsSize must match encoded records blob"
);
let topics = vec![ShareFetchedTopic {
topic_id: [0u8; 16],
partitions: vec![part.clone()],
}];
for version in [0i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &topics).unwrap();
let mut cur = buf.as_ref();
let (decoded, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert!(
!cur.has_remaining(),
"ShareFetch v{version} recordsSize leftover-empty; leftover {} bytes",
cur.remaining()
);
let got = decoded
.first()
.and_then(|t| t.partitions.first())
.expect("one partition");
assert_eq!(
got.records_size().unwrap(),
size,
"v{version} decoded recordsSize must match"
);
}
}
fn share_fetch_records_byte_index(buf: &[u8], version: i16) -> usize {
let total = buf.len();
let mut cur = buf;
let _throttle = buf::get_i32(&mut cur).unwrap();
let _error = buf::get_i16(&mut cur).unwrap();
let _msg = buf::get_string(&mut cur, true).unwrap();
if version >= 1 {
let _lock = buf::get_i32(&mut cur).unwrap();
}
let n = buf::get_array_len(&mut cur, true).unwrap().unwrap_or(0);
assert_eq!(n, 1, "one topic");
let _id = buf::get_uuid(&mut cur).unwrap();
let pn = buf::get_array_len(&mut cur, true).unwrap().unwrap_or(0);
assert_eq!(pn, 1, "one partition");
let _partition = buf::get_i32(&mut cur).unwrap();
let _p_err = buf::get_i16(&mut cur).unwrap();
let _p_msg = buf::get_string(&mut cur, true).unwrap();
let _ack = buf::get_i16(&mut cur).unwrap();
let _ack_msg = buf::get_string(&mut cur, true).unwrap();
let _leader = decode_leader(&mut cur).unwrap();
total - cur.len()
}
#[test]
fn share_fetch_response_records_nullable_versions_matches_java() {
let topics = vec![ShareFetchedTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchedPartition::partition_response(0, 0)],
}];
for version in [0_i16, 1] {
let mut empty = BytesMut::new();
encode_share_fetch_response(&mut empty, version, &topics).unwrap();
let at = share_fetch_records_byte_index(&empty, version);
assert_eq!(
empty[at], 0x01,
"encode_share_fetch_response still writes empty Records not null"
);
let mut null_recs = empty.clone();
null_recs[at] = 0x00;
assert_ne!(
&empty[..],
&null_recs[..],
"ShareFetch v{version} compact null Records is not empty"
);
if version == 0 {
let mut cur = null_recs.as_ref();
let (decoded, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
let got = decoded
.first()
.and_then(|t| t.partitions.first())
.expect("one partition");
assert!(got.records.is_empty(), "v0 null Records is empty");
assert_eq!(got.records_size().unwrap(), 0);
assert!(
cur.is_empty(),
"ShareFetch v{version} Records nullableVersions leftover-empty"
);
} else {
let mut cur = null_recs.as_ref();
let err = decode_share_fetch_response(&mut cur, version).unwrap_err();
assert!(
err.to_string()
.contains("non-nullable field records was serialized as null"),
"v1 null Records is protocol, got {err}"
);
}
}
let mut v0 = BytesMut::new();
encode_share_fetch_response(&mut v0, 0, &topics).unwrap();
let mut v1 = BytesMut::new();
encode_share_fetch_response(&mut v1, 1, &topics).unwrap();
assert_ne!(
&v0[..],
&v1[..],
"v0 and v1 both write empty Records not null; v1 still adds AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_acknowledge_partition_response_leftover_empty() {
let err = ShareAcknowledgeResponsePartition::partition_response(
3,
crate::error::UNKNOWN_TOPIC_ID,
);
assert_eq!(
err,
ShareAcknowledgeResponsePartition {
partition: 3,
error_code: crate::error::UNKNOWN_TOPIC_ID,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
}
);
let topics = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![
ShareAcknowledgeResponsePartition::partition_response(
0,
crate::error::UNKNOWN_TOPIC_ID,
),
ShareAcknowledgeResponsePartition::partition_response(
3,
crate::error::UNKNOWN_TOPIC_ID,
),
],
}];
for version in [0i16, 1] {
let mut buf = BytesMut::new();
encode_share_acknowledge_topics_response(&mut buf, version, 0, &topics).unwrap();
let mut cur = buf.as_ref();
let (top, decoded, ..) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
assert_eq!(top, 0);
assert_eq!(decoded, topics);
assert!(
!cur.has_remaining(),
"ShareAcknowledge v{version} partitionResponse leftover-empty; leftover {} bytes",
cur.remaining()
);
assert_eq!(
decode_share_acknowledge_response(&mut buf.as_ref(), version).unwrap(),
0
);
}
let empty: Vec<ShareAcknowledgeResponseTopic> = Vec::new();
for version in [0i16, 1] {
let mut via_topics = BytesMut::new();
encode_share_acknowledge_topics_response(&mut via_topics, version, 0, &empty).unwrap();
let mut via_empty = BytesMut::new();
encode_share_acknowledge_response(&mut via_empty, version, 0).unwrap();
assert_eq!(
via_topics.as_ref(),
via_empty.as_ref(),
"empty Responses matches getErrorResponse encode"
);
let mut cur = via_topics.as_ref();
let (top, decoded, ..) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
assert_eq!(top, 0);
assert_eq!(decoded, empty);
assert!(
!cur.has_remaining(),
"ShareAcknowledge v{version} empty partitionResponse leftover-empty; leftover {} bytes",
cur.remaining()
);
}
}
#[test]
fn share_acknowledge_response_node_endpoints_matches_java() {
let topics = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 6,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
}],
}];
let endpoints = [crate::protocol::api::NodeEndpoint {
node_id: 3,
host: "h".into(),
port: 1,
rack: Some("r".into()),
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_acknowledge_topics_response_with_endpoints(
&mut buf, version, 0, &topics, &endpoints,
)
.unwrap();
let mut cur = buf.as_ref();
let (top, got, eps, ..) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
assert_eq!(top, 0);
assert_eq!(got, topics);
assert_eq!(eps, endpoints);
assert!(
cur.is_empty(),
"ShareAcknowledge v{version} NodeEndpoints leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_acknowledge_topics_response_with_endpoints(
&mut with, 0, 0, &topics, &endpoints,
)
.unwrap();
let mut empty = BytesMut::new();
encode_share_acknowledge_topics_response_with_endpoints(&mut empty, 0, 0, &topics, &[])
.unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareAcknowledge NodeEndpoints is not always empty"
);
let mut conv = BytesMut::new();
encode_share_acknowledge_topics_response(&mut conv, 0, 0, &topics).unwrap();
assert_eq!(
&conv[..],
&empty[..],
"encode_share_acknowledge_topics_response still writes empty NodeEndpoints"
);
let mut v1_with = BytesMut::new();
encode_share_acknowledge_topics_response_with_endpoints(
&mut v1_with,
1,
0,
&topics,
&endpoints,
)
.unwrap();
assert_eq!(
&with[..],
&v1_with[..],
"v0 and v1 ShareAcknowledge NodeEndpoints layout match; do not confuse with ShareFetch AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_acknowledge_response_current_leader_matches_java() {
let defaults = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 6,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
}],
}];
let with_leader = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 6,
error_message: None,
current_leader_id: 2,
current_leader_epoch: 7,
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_acknowledge_topics_response(&mut buf, version, 0, &with_leader).unwrap();
let mut cur = buf.as_ref();
let (top, got, ..) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
assert_eq!(top, 0);
assert_eq!(got, with_leader);
assert!(
cur.is_empty(),
"ShareAcknowledge v{version} CurrentLeader leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_acknowledge_topics_response(&mut with, 0, 0, &with_leader).unwrap();
let mut empty = BytesMut::new();
encode_share_acknowledge_topics_response(&mut empty, 0, 0, &defaults).unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareAcknowledge CurrentLeader is not always 0/0"
);
let pr = ShareAcknowledgeResponsePartition::partition_response(0, 6);
assert_eq!(pr.current_leader_id, 0);
assert_eq!(pr.current_leader_epoch, 0);
let mut conv = BytesMut::new();
encode_share_acknowledge_topics_response(
&mut conv,
0,
0,
&[ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![pr],
}],
)
.unwrap();
assert_eq!(
&conv[..],
&empty[..],
"partition_response still writes CurrentLeader 0/0"
);
let mut v1_with = BytesMut::new();
encode_share_acknowledge_topics_response(&mut v1_with, 1, 0, &with_leader).unwrap();
assert_eq!(
&with[..],
&v1_with[..],
"v0 and v1 CurrentLeader layout match; ShareAcknowledge has no AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_acknowledge_response_error_message_matches_java() {
let defaults = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 6,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
}],
}];
let with_msg = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 6,
error_message: Some("e".into()),
current_leader_id: 0,
current_leader_epoch: 0,
}],
}];
let with_empty = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 6,
error_message: Some(String::new()),
current_leader_id: 0,
current_leader_epoch: 0,
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_acknowledge_topics_response(&mut buf, version, 0, &with_msg).unwrap();
let mut cur = buf.as_ref();
let (top, got, ..) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
assert_eq!(top, 0);
assert_eq!(got, with_msg);
assert!(
cur.is_empty(),
"ShareAcknowledge v{version} ErrorMessage leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_acknowledge_topics_response(&mut with, 0, 0, &with_msg).unwrap();
let mut empty = BytesMut::new();
encode_share_acknowledge_topics_response(&mut empty, 0, 0, &defaults).unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareAcknowledge ErrorMessage is not always null"
);
let mut empty_present = BytesMut::new();
encode_share_acknowledge_topics_response(&mut empty_present, 0, 0, &with_empty).unwrap();
assert_ne!(
&empty_present[..],
&empty[..],
"empty-but-present ErrorMessage is not JSON null"
);
let mut cur = empty_present.as_ref();
let (top, got, ..) = decode_share_acknowledge_topics_response(&mut cur, 0).unwrap();
assert_eq!(top, 0);
assert_eq!(got, with_empty);
assert!(
cur.is_empty(),
"ShareAcknowledge v0 ErrorMessage leftover-empty"
);
let pr = ShareAcknowledgeResponsePartition::partition_response(0, 6);
assert!(pr.error_message.is_none());
let mut conv = BytesMut::new();
encode_share_acknowledge_topics_response(
&mut conv,
0,
0,
&[ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![pr],
}],
)
.unwrap();
assert_eq!(
&conv[..],
&empty[..],
"partition_response still writes ErrorMessage null"
);
let mut v1_with = BytesMut::new();
encode_share_acknowledge_topics_response(&mut v1_with, 1, 0, &with_msg).unwrap();
assert_eq!(
&with[..],
&v1_with[..],
"v0 and v1 ErrorMessage layout match; ShareAcknowledge has no AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_acknowledge_response_throttle_time_ms_matches_java() {
let topics = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 6,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_acknowledge_topics_response_with_throttle(
&mut buf, version, 0, &topics, 3_600_000,
)
.unwrap();
let mut cur = buf.as_ref();
let (top, got, endpoints, throttle, ..) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
assert_eq!(top, 0);
assert_eq!(got, topics);
assert!(endpoints.is_empty());
assert_eq!(throttle, 3_600_000);
assert!(
cur.is_empty(),
"ShareAcknowledge v{version} ThrottleTimeMs leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_acknowledge_topics_response_with_throttle(&mut with, 0, 0, &topics, 3_600_000)
.unwrap();
let mut zero = BytesMut::new();
encode_share_acknowledge_topics_response_with_throttle(&mut zero, 0, 0, &topics, 0)
.unwrap();
assert_ne!(
&with[..],
&zero[..],
"v0 ThrottleTimeMs is not always the JSON default 0"
);
let mut conv = BytesMut::new();
encode_share_acknowledge_topics_response(&mut conv, 0, 0, &topics).unwrap();
assert_eq!(
&conv[..],
&zero[..],
"encode_share_acknowledge_topics_response still writes ThrottleTimeMs 0"
);
let mut endpoints_zero = BytesMut::new();
encode_share_acknowledge_topics_response_with_endpoints(
&mut endpoints_zero,
0,
0,
&topics,
&[],
)
.unwrap();
assert_eq!(
&endpoints_zero[..],
&zero[..],
"encode_share_acknowledge_topics_response_with_endpoints still writes ThrottleTimeMs 0"
);
let mut error_enc = BytesMut::new();
encode_share_acknowledge_response(&mut error_enc, 0, 0).unwrap();
let mut empty_topics = BytesMut::new();
encode_share_acknowledge_topics_response_with_throttle(&mut empty_topics, 0, 0, &[], 0)
.unwrap();
assert_eq!(
&error_enc[..],
&empty_topics[..],
"encode_share_acknowledge_response still writes ThrottleTimeMs 0"
);
let mut v1_with = BytesMut::new();
encode_share_acknowledge_topics_response_with_throttle(
&mut v1_with,
1,
0,
&topics,
3_600_000,
)
.unwrap();
assert_eq!(
&with[..],
&v1_with[..],
"v0 and v1 both write ThrottleTimeMs (JSON 0+); ShareAcknowledge has no AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_acknowledge_request_error_response_matches_java() {
leftover_error_response(0, 16, 3_600_000);
leftover_error_response(0, 0, 0);
leftover_error_response(1, 16, 3_600_000);
leftover_error_response(1, 0, 0);
let mut named = BytesMut::new();
ShareAcknowledgeRequest::error_response(&mut named, 0, 16, 0).unwrap();
let mut conv = BytesMut::new();
encode_share_acknowledge_response(&mut conv, 0, 16).unwrap();
assert_eq!(
&named[..],
&conv[..],
"encode_share_acknowledge_response still writes ThrottleTimeMs 0"
);
let mut with = BytesMut::new();
ShareAcknowledgeRequest::error_response(&mut with, 0, 16, 3_600_000).unwrap();
assert_ne!(
&with[..],
&conv[..],
"getErrorResponse ThrottleTimeMs is not always the JSON default 0"
);
}
fn leftover_error_response(version: i16, error_code: i16, throttle_time_ms: i32) {
let mut buf = BytesMut::new();
ShareAcknowledgeRequest::error_response(&mut buf, version, error_code, throttle_time_ms)
.unwrap();
let mut cur = buf.as_ref();
let (got_error, topics, endpoints, got_throttle, error_message) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
assert_eq!(got_error, error_code);
assert!(topics.is_empty());
assert!(endpoints.is_empty());
assert_eq!(got_throttle, throttle_time_ms);
assert!(error_message.is_none());
let empty = if error_code == 0 && throttle_time_ms == 0 {
"empty "
} else {
""
};
assert!(
cur.is_empty(),
"ShareAcknowledge v{version} getErrorResponse {empty}leftover-empty; leftover {} bytes",
cur.len()
);
}
#[test]
fn share_acknowledge_response_top_level_error_message_matches_java() {
let topics = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 6,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
}],
}];
let with_part = vec![ShareAcknowledgeResponseTopic {
topic_id: [7u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 6,
error_message: Some("e".into()),
current_leader_id: 0,
current_leader_epoch: 0,
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_acknowledge_topics_response_with_error_message(
&mut buf,
version,
0,
&topics,
Some("e"),
)
.unwrap();
let mut cur = buf.as_ref();
let (top, got, endpoints, throttle, msg) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
assert_eq!(top, 0);
assert_eq!(got, topics);
assert!(endpoints.is_empty());
assert_eq!(throttle, 0);
assert_eq!(msg.as_deref(), Some("e"));
assert!(
cur.is_empty(),
"ShareAcknowledge v{version} top-level ErrorMessage leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_acknowledge_topics_response_with_error_message(
&mut with,
0,
0,
&topics,
Some("e"),
)
.unwrap();
let mut empty = BytesMut::new();
encode_share_acknowledge_topics_response_with_error_message(
&mut empty, 0, 0, &topics, None,
)
.unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareAcknowledge top-level ErrorMessage is not always null"
);
let mut conv = BytesMut::new();
encode_share_acknowledge_topics_response(&mut conv, 0, 0, &topics).unwrap();
assert_eq!(
&conv[..],
&empty[..],
"encode_share_acknowledge_topics_response still writes top-level ErrorMessage null"
);
let mut part_msg = BytesMut::new();
encode_share_acknowledge_topics_response(&mut part_msg, 0, 0, &with_part).unwrap();
assert_ne!(
&with[..],
&part_msg[..],
"top-level ErrorMessage is not partition ErrorMessage"
);
let mut empty_present = BytesMut::new();
encode_share_acknowledge_topics_response_with_error_message(
&mut empty_present,
0,
0,
&topics,
Some(""),
)
.unwrap();
assert_ne!(
&empty_present[..],
&empty[..],
"empty-but-present top-level ErrorMessage is not JSON null"
);
let mut cur = empty_present.as_ref();
let (top, got, .., msg) = decode_share_acknowledge_topics_response(&mut cur, 0).unwrap();
assert_eq!(top, 0);
assert_eq!(got, topics);
assert_eq!(msg.as_deref(), Some(""));
assert!(
cur.is_empty(),
"ShareAcknowledge v0 top-level ErrorMessage leftover-empty"
);
let mut v1_with = BytesMut::new();
encode_share_acknowledge_topics_response_with_error_message(
&mut v1_with,
1,
0,
&topics,
Some("e"),
)
.unwrap();
assert_eq!(
&with[..],
&v1_with[..],
"v0 and v1 top-level ErrorMessage layout match; ShareAcknowledge has no AcquisitionLockTimeoutMs"
);
}
#[test]
fn share_acknowledge_response_error_counts_matches_java() {
assert_eq!(
ShareAcknowledgeResponse::error_counts(0, &[]),
HashMap::from([(0, 1)])
);
let topics = vec![
ShareAcknowledgeResponseTopic {
topic_id: [1u8; 16],
partitions: vec![
ShareAcknowledgeResponsePartition::partition_response(0, 0),
ShareAcknowledgeResponsePartition::partition_response(
1,
crate::error::UNKNOWN_TOPIC_ID,
),
],
},
ShareAcknowledgeResponseTopic {
topic_id: [2u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition::partition_response(0, 0)],
},
];
assert_eq!(
ShareAcknowledgeResponse::error_counts(0, &topics),
HashMap::from([(0, 3), (crate::error::UNKNOWN_TOPIC_ID, 1)])
);
let top =
ShareAcknowledgeResponse::error_counts(crate::error::GROUP_AUTHORIZATION_FAILED, &[]);
assert_eq!(
top,
HashMap::from([(crate::error::GROUP_AUTHORIZATION_FAILED, 1)])
);
let same = ShareAcknowledgeResponse::error_counts(
crate::error::UNKNOWN_TOPIC_ID,
&[ShareAcknowledgeResponseTopic {
topic_id: [3u8; 16],
partitions: vec![ShareAcknowledgeResponsePartition::partition_response(
0,
crate::error::UNKNOWN_TOPIC_ID,
)],
}],
);
assert_eq!(same, HashMap::from([(crate::error::UNKNOWN_TOPIC_ID, 2)]));
}
#[test]
fn share_acknowledge_response_to_message_matches_java() {
assert!(ShareAcknowledgeResponse::to_message(&[]).is_empty());
let a = [1u8; 16];
let b = [2u8; 16];
let body0 = ShareAcknowledgeResponsePartition::partition_response(99, 0);
let body1 = ShareAcknowledgeResponsePartition::partition_response(
1,
crate::error::UNKNOWN_TOPIC_ID,
);
let body2 = ShareAcknowledgeResponsePartition::partition_response(2, 0);
let grouped = ShareAcknowledgeResponse::to_message(&[
(a, 0, body0.clone()),
(b, 1, body1),
(a, 3, body2),
]);
assert_eq!(
grouped,
vec![
ShareAcknowledgeResponseTopic {
topic_id: a,
partitions: vec![
ShareAcknowledgeResponsePartition {
partition: 0,
error_code: 0,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
},
ShareAcknowledgeResponsePartition {
partition: 3,
error_code: 0,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
},
],
},
ShareAcknowledgeResponseTopic {
topic_id: b,
partitions: vec![ShareAcknowledgeResponsePartition {
partition: 1,
error_code: crate::error::UNKNOWN_TOPIC_ID,
error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
}],
},
]
);
assert_eq!(
grouped
.first()
.and_then(|topic| topic.partitions.first())
.map(|part| part.partition),
Some(0),
"setPartitionIndex copies the key partition onto the body"
);
assert_eq!(body0.partition, 99);
for version in [0i16, 1] {
let mut buf = BytesMut::new();
encode_share_acknowledge_topics_response(&mut buf, version, 0, &grouped).unwrap();
let mut cur = buf.as_ref();
let (top, decoded, ..) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
assert_eq!(top, 0);
assert_eq!(decoded, grouped);
assert!(
!cur.has_remaining(),
"ShareAcknowledge v{version} toMessage leftover-empty; leftover {} bytes",
cur.remaining()
);
}
}
#[test]
fn share_acknowledge_request_for_consumer_matches_java() {
assert!(ShareAcknowledgeRequest::for_consumer(std::iter::empty::<(
[u8; 16],
i32,
Vec<AcknowledgementBatch>
)>())
.is_empty());
let a = [1u8; 16];
let b = [2u8; 16];
let b0 = AcknowledgementBatch {
first_offset: 0,
last_offset: 1,
types: vec![ACK_ACCEPT],
};
let b1 = AcknowledgementBatch {
first_offset: 2,
last_offset: 3,
types: vec![ACK_RELEASE],
};
let b2 = AcknowledgementBatch {
first_offset: 4,
last_offset: 5,
types: vec![ACK_REJECT],
};
let grouped = ShareAcknowledgeRequest::for_consumer([
(a, 0, vec![b0.clone()]),
(b, 1, vec![b1.clone()]),
(a, 2, vec![b2.clone()]),
]);
assert_eq!(
grouped,
vec![
ShareAckTopic {
topic_id: a,
partitions: vec![(0, vec![b0.clone()]), (2, vec![b2.clone()])],
},
ShareAckTopic {
topic_id: b,
partitions: vec![(1, vec![b1.clone()])],
},
]
);
let last_wins = ShareAcknowledgeRequest::for_consumer([
(a, 0, vec![b0.clone()]),
(a, 0, vec![b2.clone()]),
]);
assert_eq!(
last_wins,
vec![ShareAckTopic {
topic_id: a,
partitions: vec![(0, vec![b2.clone()])],
}]
);
let order = ShareAcknowledgeRequest::for_consumer([
(a, 0, vec![b0]),
(a, 1, vec![b1.clone()]),
(a, 0, vec![b2.clone()]),
]);
assert_eq!(
order,
vec![ShareAckTopic {
topic_id: a,
partitions: vec![(0, vec![b2]), (1, vec![b1])],
}]
);
leftover_share_ack_for_consumer(0, &grouped);
leftover_share_ack_for_consumer(0, &[]);
leftover_share_ack_for_consumer(1, &grouped);
leftover_share_ack_for_consumer(1, &[]);
}
fn leftover_share_ack_for_consumer(version: i16, topics: &[ShareAckTopic]) {
let mut buf = BytesMut::new();
encode_share_acknowledge_topics(&mut buf, version, "g", "m", 1, topics).unwrap();
let mut cur = buf.as_ref();
let (gid, mid, epoch, flat) = decode_share_acknowledge_request(&mut cur, version).unwrap();
assert_eq!(gid, "g");
assert_eq!(mid, "m");
assert_eq!(epoch, 1);
assert_eq!(
ShareAcknowledgeRequest::for_consumer(flat).as_slice(),
topics
);
let empty = if topics.is_empty() { "empty " } else { "" };
assert!(
!cur.has_remaining(),
"ShareAcknowledge v{version} Builder.forConsumer {empty}leftover-empty; leftover {} bytes",
cur.remaining()
);
}
#[test]
fn share_fetch_response_error_counts_matches_java() {
assert_eq!(
ShareFetchResponse::error_counts(0, &[]),
HashMap::from([(0, 1)])
);
let topics = vec![ShareFetchedTopic {
topic_id: [1u8; 16],
partitions: vec![
ShareFetchedPartition::partition_response(0, 0),
ShareFetchedPartition::partition_response(1, crate::error::UNKNOWN_TOPIC_ID),
],
}];
assert_eq!(
ShareFetchResponse::error_counts(0, &topics),
HashMap::from([(0, 2), (crate::error::UNKNOWN_TOPIC_ID, 1)])
);
assert_eq!(
ShareFetchResponse::error_counts(crate::error::GROUP_AUTHORIZATION_FAILED, &[]),
HashMap::from([(crate::error::GROUP_AUTHORIZATION_FAILED, 1)])
);
let same = ShareFetchResponse::error_counts(
crate::error::UNKNOWN_TOPIC_ID,
&[ShareFetchedTopic {
topic_id: [2u8; 16],
partitions: vec![ShareFetchedPartition::partition_response(
0,
crate::error::UNKNOWN_TOPIC_ID,
)],
}],
);
assert_eq!(same, HashMap::from([(crate::error::UNKNOWN_TOPIC_ID, 2)]));
}
#[test]
fn share_fetch_response_response_data_matches_java() {
let topic_id = [1u8; 16];
let unknown_id = [2u8; 16];
let p0 = ShareFetchedPartition::partition_response(0, 0);
let p1 = ShareFetchedPartition::partition_response(1, crate::error::UNKNOWN_TOPIC_ID);
let overwrite =
ShareFetchedPartition::partition_response(0, crate::error::NOT_LEADER_OR_FOLLOWER);
let topics = vec![
ShareFetchedTopic {
topic_id,
partitions: vec![p0.clone(), p1.clone()],
},
ShareFetchedTopic {
topic_id,
partitions: vec![overwrite.clone()],
},
ShareFetchedTopic {
topic_id: unknown_id,
partitions: vec![ShareFetchedPartition::partition_response(0, 0)],
},
];
assert!(ShareFetchResponse::response_data(&[], &HashMap::new()).is_empty());
assert!(ShareFetchResponse::response_data(&topics, &HashMap::new()).is_empty());
let names = HashMap::from([(topic_id, "t".into())]);
assert_eq!(
ShareFetchResponse::response_data(&topics, &names),
HashMap::from([
((topic_id, "t".into(), 0), overwrite.clone()),
((topic_id, "t".into(), 1), p1.clone()),
])
);
for version in [0i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &topics).unwrap();
let mut cur = buf.as_ref();
let (decoded, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(decoded, topics);
assert_eq!(
ShareFetchResponse::response_data(&decoded, &names),
ShareFetchResponse::response_data(&topics, &names)
);
assert!(
!cur.has_remaining(),
"ShareFetch v{version} responseData leftover-empty; leftover {} bytes",
cur.remaining()
);
}
}
#[test]
fn share_fetch_request_forgotten_topics_matches_java() {
let topic_id = [1u8; 16];
let forgotten = vec![
ShareForgottenTopic {
topic_id,
partitions: vec![0, 1],
},
ShareForgottenTopic {
topic_id,
partitions: vec![0],
},
];
let empty_names = HashMap::new();
assert!(ShareFetchRequest::forgotten_topics(&[], &empty_names).is_empty());
let unresolved = ShareFetchRequest::forgotten_topics(&forgotten, &empty_names);
assert_eq!(
unresolved,
vec![
(topic_id, None, 0),
(topic_id, None, 1),
(topic_id, None, 0),
],
"missing name is still inserted"
);
let names = HashMap::from([(topic_id, "t".into())]);
assert_eq!(
ShareFetchRequest::forgotten_topics(&forgotten, &names),
vec![
(topic_id, Some("t".into()), 0),
(topic_id, Some("t".into()), 1),
(topic_id, Some("t".into()), 0),
]
);
let fetch = vec![ShareFetchTopic {
topic_id,
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 1024,
acknowledgements: vec![],
}],
}];
for version in [0_i16, 1] {
let got = ShareFetchRequest::forgotten_topics(&forgotten, &names);
assert_eq!(got.len(), 3);
let mut buf = BytesMut::new();
encode_share_fetch_request(&mut buf, version, "sg", "m1", 0, 10, 1, 1024, 16, &fetch)
.unwrap();
let mut cur = buf.as_ref();
let (_gid, _mid, _epoch, _max, _decoded, ..) =
decode_share_fetch_request(&mut cur, version).unwrap();
assert!(
!cur.has_remaining(),
"ShareFetch v{version} forgottenTopics leftover-empty; leftover {} bytes",
cur.remaining()
);
}
for version in [0_i16, 1] {
let got = ShareFetchRequest::forgotten_topics(&[], &empty_names);
assert!(got.is_empty());
let mut buf = BytesMut::new();
encode_share_fetch_request(&mut buf, version, "sg", "m1", 0, 10, 1, 1024, 16, &fetch)
.unwrap();
let mut cur = buf.as_ref();
let (_gid, _mid, _epoch, _max, _decoded, ..) =
decode_share_fetch_request(&mut cur, version).unwrap();
assert!(
!cur.has_remaining(),
"ShareFetch v{version} forgottenTopics empty leftover-empty; leftover {} bytes",
cur.remaining()
);
}
}
#[test]
fn share_fetch_request_forgotten_matches_java() {
let topics: Vec<ShareFetchTopic> = Vec::new();
let forgotten = [ShareForgottenTopic {
topic_id: [7u8; 16],
partitions: vec![1, 1],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_request_with_forgotten(
&mut buf, version, "sg", "m1", 0, 10, 1, 1024, 16, &topics, &forgotten,
)
.unwrap();
let mut cur = buf.as_ref();
let (.., got, _, _, _, _) = decode_share_fetch_request(&mut cur, version).unwrap();
assert_eq!(got.as_slice(), forgotten.as_slice());
assert!(
cur.is_empty(),
"ShareFetch v{version} ForgottenTopicsData leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_request_with_forgotten(
&mut with, 0, "sg", "m1", 0, 10, 1, 1024, 16, &topics, &forgotten,
)
.unwrap();
let mut empty = BytesMut::new();
encode_share_fetch_request_with_forgotten(
&mut empty,
0,
"sg",
"m1",
0,
10,
1,
1024,
16,
&topics,
&[],
)
.unwrap();
assert_ne!(
&with[..],
&empty[..],
"ShareFetch ForgottenTopicsData is not always empty"
);
let mut conv = BytesMut::new();
encode_share_fetch_request(&mut conv, 0, "sg", "m1", 0, 10, 1, 1024, 16, &topics).unwrap();
assert_eq!(
&conv[..],
&empty[..],
"encode_share_fetch_request still writes empty ForgottenTopicsData"
);
let mut v1_with = BytesMut::new();
encode_share_fetch_request_with_forgotten(
&mut v1_with,
1,
"sg",
"m1",
0,
10,
1,
1024,
16,
&topics,
&forgotten,
)
.unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 ForgottenTopicsData share TopicId layout; v1 still adds MaxRecords / BatchSize"
);
}
#[test]
fn share_fetch_request_batch_size_matches_java() {
let topics = vec![ShareFetchTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 0,
acknowledgements: vec![],
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_request_with_batch_size(
&mut buf, version, "sg", "m1", 0, 10, 1, 1024, 16, &topics, 3_600_000,
)
.unwrap();
let mut cur = buf.as_ref();
let (gid, mid, epoch, max_records, got, forgotten, batch_size, ..) =
decode_share_fetch_request(&mut cur, version).unwrap();
assert_eq!(gid.as_str(), "sg");
assert_eq!(mid.as_str(), "m1");
assert_eq!(epoch, 0);
assert!(forgotten.is_empty());
assert_eq!(got, topics);
if version >= 1 {
assert_eq!(max_records, 16);
assert_eq!(batch_size, 3_600_000);
} else {
assert_eq!(max_records, 0, "v0 omits MaxRecords");
assert_eq!(
batch_size, 0,
"ShareFetch v{version} omits BatchSize even when the body has a non-zero value"
);
}
assert!(
cur.is_empty(),
"ShareFetch request v{version} BatchSize leftover-empty"
);
}
let mut v0_with = BytesMut::new();
encode_share_fetch_request_with_batch_size(
&mut v0_with,
0,
"sg",
"m1",
0,
10,
1,
1024,
16,
&topics,
3_600_000,
)
.unwrap();
let mut v0_zero = BytesMut::new();
encode_share_fetch_request_with_batch_size(
&mut v0_zero,
0,
"sg",
"m1",
0,
10,
1,
1024,
16,
&topics,
0,
)
.unwrap();
assert_eq!(
&v0_with[..],
&v0_zero[..],
"v0 omits BatchSize even when the body has a non-zero value"
);
let mut conv_v0 = BytesMut::new();
encode_share_fetch_request(&mut conv_v0, 0, "sg", "m1", 0, 10, 1, 1024, 16, &topics)
.unwrap();
assert_eq!(
&conv_v0[..],
&v0_zero[..],
"encode_share_fetch_request v0 has no BatchSize"
);
let mut with = BytesMut::new();
encode_share_fetch_request_with_batch_size(
&mut with, 1, "sg", "m1", 0, 10, 1, 1024, 16, &topics, 3_600_000,
)
.unwrap();
let mut as_max = BytesMut::new();
encode_share_fetch_request_with_batch_size(
&mut as_max,
1,
"sg",
"m1",
0,
10,
1,
1024,
16,
&topics,
16,
)
.unwrap();
assert_ne!(
&with[..],
&as_max[..],
"v1 BatchSize is not always copied from MaxRecords"
);
let mut conv = BytesMut::new();
encode_share_fetch_request(&mut conv, 1, "sg", "m1", 0, 10, 1, 1024, 16, &topics).unwrap();
assert_eq!(
&conv[..],
&as_max[..],
"encode_share_fetch_request still writes BatchSize as MaxRecords"
);
let mut forgotten_copy = BytesMut::new();
encode_share_fetch_request_with_forgotten(
&mut forgotten_copy,
1,
"sg",
"m1",
0,
10,
1,
1024,
16,
&topics,
&[],
)
.unwrap();
assert_eq!(
&forgotten_copy[..],
&as_max[..],
"encode_share_fetch_request_with_forgotten still writes BatchSize as MaxRecords"
);
assert_ne!(
&v0_with[..],
&with[..],
"v1 adds MaxRecords and BatchSize after MaxBytes"
);
}
#[test]
fn share_fetch_request_max_wait_ms_matches_java() {
let topics = vec![ShareFetchTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 0,
acknowledgements: vec![],
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_request(
&mut buf, version, "sg", "m1", 0, 3_600_000, 1, 1024, 16, &topics,
)
.unwrap();
let mut cur = buf.as_ref();
let (gid, mid, epoch, max_records, got, forgotten, batch_size, max_wait, ..) =
decode_share_fetch_request(&mut cur, version).unwrap();
assert_eq!(gid.as_str(), "sg");
assert_eq!(mid.as_str(), "m1");
assert_eq!(epoch, 0);
assert!(forgotten.is_empty());
assert_eq!(got, topics);
assert_eq!(max_wait, 3_600_000);
if version >= 1 {
assert_eq!(max_records, 16);
assert_eq!(batch_size, 16);
} else {
assert_eq!(max_records, 0);
assert_eq!(batch_size, 0);
}
assert!(
cur.is_empty(),
"ShareFetch request v{version} MaxWaitMs leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_request(&mut with, 0, "sg", "m1", 0, 3_600_000, 1, 1024, 16, &topics)
.unwrap();
let mut ten = BytesMut::new();
encode_share_fetch_request(&mut ten, 0, "sg", "m1", 0, 10, 1, 1024, 16, &topics).unwrap();
assert_ne!(&with[..], &ten[..], "v0 MaxWaitMs is not always 10");
let mut cur = ten.as_ref();
let (.., max_wait, _, _) = decode_share_fetch_request(&mut cur, 0).unwrap();
assert_eq!(max_wait, 10);
let mut v1_with = BytesMut::new();
encode_share_fetch_request(
&mut v1_with,
1,
"sg",
"m1",
0,
3_600_000,
1,
1024,
16,
&topics,
)
.unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 both write MaxWaitMs (JSON 0+); v1 still adds MaxRecords / BatchSize"
);
}
#[test]
fn share_fetch_request_min_bytes_matches_java() {
let topics = vec![ShareFetchTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 0,
acknowledgements: vec![],
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_request(
&mut buf, version, "sg", "m1", 0, 10, 3_600_000, 1024, 16, &topics,
)
.unwrap();
let mut cur = buf.as_ref();
let (gid, mid, epoch, max_records, got, forgotten, batch_size, max_wait, min_bytes, ..) =
decode_share_fetch_request(&mut cur, version).unwrap();
assert_eq!(gid.as_str(), "sg");
assert_eq!(mid.as_str(), "m1");
assert_eq!(epoch, 0);
assert!(forgotten.is_empty());
assert_eq!(got, topics);
assert_eq!(max_wait, 10);
assert_eq!(min_bytes, 3_600_000);
if version >= 1 {
assert_eq!(max_records, 16);
assert_eq!(batch_size, 16);
} else {
assert_eq!(max_records, 0);
assert_eq!(batch_size, 0);
}
assert!(
cur.is_empty(),
"ShareFetch request v{version} MinBytes leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_request(
&mut with, 0, "sg", "m1", 0, 10, 3_600_000, 1024, 16, &topics,
)
.unwrap();
let mut one = BytesMut::new();
encode_share_fetch_request(&mut one, 0, "sg", "m1", 0, 10, 1, 1024, 16, &topics).unwrap();
assert_ne!(&with[..], &one[..], "v0 MinBytes is not always 1");
let mut cur = one.as_ref();
let (.., min_bytes, _) = decode_share_fetch_request(&mut cur, 0).unwrap();
assert_eq!(min_bytes, 1);
let mut v1_with = BytesMut::new();
encode_share_fetch_request(
&mut v1_with,
1,
"sg",
"m1",
0,
10,
3_600_000,
1024,
16,
&topics,
)
.unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 both write MinBytes (JSON 0+); v1 still adds MaxRecords / BatchSize"
);
}
#[test]
fn share_fetch_request_max_bytes_matches_java() {
let topics = vec![ShareFetchTopic {
topic_id: [7u8; 16],
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 0,
acknowledgements: vec![],
}],
}];
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_request(
&mut buf, version, "sg", "m1", 0, 10, 1, 3_600_000, 16, &topics,
)
.unwrap();
let mut cur = buf.as_ref();
let (
gid,
mid,
epoch,
max_records,
got,
forgotten,
batch_size,
max_wait,
min_bytes,
max_bytes,
) = decode_share_fetch_request(&mut cur, version).unwrap();
assert_eq!(gid.as_str(), "sg");
assert_eq!(mid.as_str(), "m1");
assert_eq!(epoch, 0);
assert!(forgotten.is_empty());
assert_eq!(got, topics);
assert_eq!(max_wait, 10);
assert_eq!(min_bytes, 1);
assert_eq!(max_bytes, 3_600_000);
if version >= 1 {
assert_eq!(max_records, 16);
assert_eq!(batch_size, 16);
} else {
assert_eq!(max_records, 0);
assert_eq!(batch_size, 0);
}
assert!(
cur.is_empty(),
"ShareFetch request v{version} MaxBytes leftover-empty"
);
}
let mut with = BytesMut::new();
encode_share_fetch_request(&mut with, 0, "sg", "m1", 0, 10, 1, 3_600_000, 16, &topics)
.unwrap();
let mut kilo = BytesMut::new();
encode_share_fetch_request(&mut kilo, 0, "sg", "m1", 0, 10, 1, 1024, 16, &topics).unwrap();
assert_ne!(&with[..], &kilo[..], "v0 MaxBytes is not always 1024");
let mut cur = kilo.as_ref();
let (.., max_bytes) = decode_share_fetch_request(&mut cur, 0).unwrap();
assert_eq!(max_bytes, 1024);
let mut json_default = BytesMut::new();
encode_share_fetch_request(
&mut json_default,
0,
"sg",
"m1",
0,
10,
1,
i32::MAX,
16,
&topics,
)
.unwrap();
assert_ne!(
&with[..],
&json_default[..],
"v0 MaxBytes is not always JSON default 0x7fffffff"
);
let mut cur = json_default.as_ref();
let (.., max_bytes) = decode_share_fetch_request(&mut cur, 0).unwrap();
assert_eq!(max_bytes, i32::MAX);
let mut v1_with = BytesMut::new();
encode_share_fetch_request(
&mut v1_with,
1,
"sg",
"m1",
0,
10,
1,
3_600_000,
16,
&topics,
)
.unwrap();
assert_ne!(
&with[..],
&v1_with[..],
"v0 and v1 both write MaxBytes (JSON 0+); v1 still adds MaxRecords / BatchSize"
);
}
#[test]
fn share_fetch_request_update_forgotten_data_matches_java() {
let none: Vec<([u8; 16], i32)> = Vec::new();
assert!(ShareFetchRequest::update_forgotten_data(&[], none.clone()).is_empty());
let id_a = [1u8; 16];
let id_b = [2u8; 16];
let forget = [(id_a, 0i32), (id_b, 1), (id_a, 2), (id_a, 0)];
let grouped = ShareFetchRequest::update_forgotten_data(&[], forget);
assert_eq!(
grouped,
vec![
ShareForgottenTopic {
topic_id: id_a,
partitions: vec![0, 2, 0],
},
ShareForgottenTopic {
topic_id: id_b,
partitions: vec![1],
},
]
);
let existing = [ShareForgottenTopic {
topic_id: id_a,
partitions: vec![9],
}];
let appended = ShareFetchRequest::update_forgotten_data(&existing, forget);
assert_eq!(
appended,
vec![
ShareForgottenTopic {
topic_id: id_a,
partitions: vec![9],
},
ShareForgottenTopic {
topic_id: id_a,
partitions: vec![0, 2, 0],
},
ShareForgottenTopic {
topic_id: id_b,
partitions: vec![1],
},
],
"a second call with the same id appends, it does not merge"
);
let fetch = vec![ShareFetchTopic {
topic_id: id_a,
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 1024,
acknowledgements: vec![],
}],
}];
for version in [0_i16, 1] {
let got = ShareFetchRequest::update_forgotten_data(&[], forget);
assert_eq!(got.len(), 2);
let mut buf = BytesMut::new();
encode_share_fetch_request(&mut buf, version, "sg", "m1", 0, 10, 1, 1024, 16, &fetch)
.unwrap();
let mut cur = buf.as_ref();
let (_gid, _mid, _epoch, _max, _decoded, ..) =
decode_share_fetch_request(&mut cur, version).unwrap();
assert!(
!cur.has_remaining(),
"ShareFetch v{version} updateForgottenData leftover-empty; leftover {} bytes",
cur.remaining()
);
}
for version in [0_i16, 1] {
let got = ShareFetchRequest::update_forgotten_data(&[], none.clone());
assert!(got.is_empty());
let mut buf = BytesMut::new();
encode_share_fetch_request(&mut buf, version, "sg", "m1", 0, 10, 1, 1024, 16, &fetch)
.unwrap();
let mut cur = buf.as_ref();
let (_gid, _mid, _epoch, _max, _decoded, ..) =
decode_share_fetch_request(&mut cur, version).unwrap();
assert!(
!cur.has_remaining(),
"ShareFetch v{version} updateForgottenData empty leftover-empty; leftover {} bytes",
cur.remaining()
);
}
}
#[test]
fn share_fetch_request_for_consumer_matches_java() {
assert!(ShareFetchRequest::for_consumer(
false,
1024,
std::iter::empty::<([u8; 16], i32)>(),
std::iter::empty::<([u8; 16], i32, Vec<AcknowledgementBatch>)>(),
)
.is_empty());
let a = [1u8; 16];
let b = [2u8; 16];
let b0 = AcknowledgementBatch {
first_offset: 0,
last_offset: 1,
types: vec![ACK_ACCEPT],
};
let b1 = AcknowledgementBatch {
first_offset: 2,
last_offset: 3,
types: vec![ACK_RELEASE],
};
let grouped = ShareFetchRequest::for_consumer(
false,
1024,
[(a, 0), (b, 1), (a, 2)],
std::iter::empty::<([u8; 16], i32, Vec<AcknowledgementBatch>)>(),
);
assert_eq!(
grouped,
vec![
ShareFetchTopic {
topic_id: a,
partitions: vec![
ShareFetchPartition {
partition: 0,
partition_max_bytes: 1024,
acknowledgements: vec![],
},
ShareFetchPartition {
partition: 2,
partition_max_bytes: 1024,
acknowledgements: vec![],
},
],
},
ShareFetchTopic {
topic_id: b,
partitions: vec![ShareFetchPartition {
partition: 1,
partition_max_bytes: 1024,
acknowledgements: vec![],
}],
},
]
);
let last_wins = ShareFetchRequest::for_consumer(
false,
2048,
[(a, 0), (a, 0)],
std::iter::empty::<([u8; 16], i32, Vec<AcknowledgementBatch>)>(),
);
assert_eq!(
last_wins,
vec![ShareFetchTopic {
topic_id: a,
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 2048,
acknowledgements: vec![],
}],
}]
);
let with_acks = ShareFetchRequest::for_consumer(
false,
1024,
[(a, 0)],
[(a, 0, vec![b0.clone()]), (a, 1, vec![b1.clone()])],
);
assert_eq!(
with_acks,
vec![ShareFetchTopic {
topic_id: a,
partitions: vec![
ShareFetchPartition {
partition: 0,
partition_max_bytes: 1024,
acknowledgements: vec![b0.clone()],
},
ShareFetchPartition {
partition: 1,
partition_max_bytes: 1024,
acknowledgements: vec![b1.clone()],
},
],
}]
);
let ack_last_wins = ShareFetchRequest::for_consumer(
false,
1024,
[(a, 0)],
[(a, 0, vec![b0.clone()]), (a, 0, vec![b1.clone()])],
);
assert_eq!(
ack_last_wins
.first()
.and_then(|topic| topic.partitions.first())
.map(|part| part.acknowledgements.as_slice()),
Some(std::slice::from_ref(&b1))
);
let closing = ShareFetchRequest::for_consumer(
true,
1024,
[(a, 0), (b, 1)],
[(a, 0, vec![b0.clone()])],
);
assert_eq!(
closing,
vec![ShareFetchTopic {
topic_id: a,
partitions: vec![ShareFetchPartition {
partition: 0,
partition_max_bytes: 0,
acknowledgements: vec![b0],
}],
}],
"closing skips send; ack-only PartitionMaxBytes is 0"
);
leftover_share_fetch_for_consumer(0, &grouped);
leftover_share_fetch_for_consumer(0, &[]);
leftover_share_fetch_for_consumer(1, &grouped);
leftover_share_fetch_for_consumer(1, &[]);
}
fn leftover_share_fetch_for_consumer(version: i16, topics: &[ShareFetchTopic]) {
let mut buf = BytesMut::new();
encode_share_fetch_request(&mut buf, version, "g", "m", 1, 10, 1, 1024, 16, topics)
.unwrap();
let mut cur = buf.as_ref();
let (gid, mid, epoch, _max, decoded, ..) =
decode_share_fetch_request(&mut cur, version).unwrap();
assert_eq!(gid, "g");
assert_eq!(mid, "m");
assert_eq!(epoch, 1);
if version == 0 {
assert_eq!(decoded.as_slice(), topics);
} else {
let zeroed: Vec<ShareFetchTopic> = topics
.iter()
.map(|topic| ShareFetchTopic {
topic_id: topic.topic_id,
partitions: topic
.partitions
.iter()
.map(|part| ShareFetchPartition {
partition: part.partition,
partition_max_bytes: 0,
acknowledgements: part.acknowledgements.clone(),
})
.collect(),
})
.collect();
assert_eq!(decoded, zeroed);
}
let empty = if topics.is_empty() { "empty " } else { "" };
assert!(
!cur.has_remaining(),
"ShareFetch v{version} Builder.forConsumer {empty}leftover-empty; leftover {} bytes",
cur.remaining()
);
}
#[test]
fn share_fetch_request_share_fetch_data_matches_java() {
assert!(ShareFetchRequest::share_fetch_data(&[], &HashMap::new()).is_empty());
let a = [1u8; 16];
let b = [2u8; 16];
let topics = vec![
ShareFetchTopic {
topic_id: a,
partitions: vec![
ShareFetchPartition {
partition: 0,
partition_max_bytes: 1024,
acknowledgements: vec![],
},
ShareFetchPartition {
partition: 0,
partition_max_bytes: 2048,
acknowledgements: vec![],
},
ShareFetchPartition {
partition: 1,
partition_max_bytes: 512,
acknowledgements: vec![],
},
],
},
ShareFetchTopic {
topic_id: b,
partitions: vec![ShareFetchPartition {
partition: 2,
partition_max_bytes: 256,
acknowledgements: vec![],
}],
},
];
let empty_names = HashMap::new();
let unresolved = ShareFetchRequest::share_fetch_data(&topics, &empty_names);
assert_eq!(
unresolved,
HashMap::from([
((a, None, 0), 2048),
((a, None, 1), 512),
((b, None, 2), 256),
]),
"missing name is still inserted; later partition overwrites"
);
let names = HashMap::from([(a, "t".into()), (b, "u".into())]);
assert_eq!(
ShareFetchRequest::share_fetch_data(&topics, &names),
HashMap::from([
((a, Some("t".into()), 0), 2048),
((a, Some("t".into()), 1), 512),
((b, Some("u".into()), 2), 256),
])
);
leftover_share_fetch_share_fetch_data(0, &topics);
leftover_share_fetch_share_fetch_data(0, &[]);
leftover_share_fetch_share_fetch_data(1, &topics);
leftover_share_fetch_share_fetch_data(1, &[]);
}
fn leftover_share_fetch_share_fetch_data(version: i16, topics: &[ShareFetchTopic]) {
let mut buf = BytesMut::new();
encode_share_fetch_request(&mut buf, version, "g", "m", 1, 10, 1, 1024, 16, topics)
.unwrap();
let mut cur = buf.as_ref();
let (_gid, _mid, _epoch, _max, decoded, ..) =
decode_share_fetch_request(&mut cur, version).unwrap();
let names = HashMap::new();
let _got = ShareFetchRequest::share_fetch_data(&decoded, &names);
let empty = if topics.is_empty() { "empty " } else { "" };
assert!(
!cur.has_remaining(),
"ShareFetch v{version} shareFetchData {empty}leftover-empty; leftover {} bytes",
cur.remaining()
);
}
#[test]
fn share_fetch_response_to_message_matches_java() {
assert!(ShareFetchResponse::to_message(&[]).is_empty());
let a = [1u8; 16];
let b = [2u8; 16];
let body0 = ShareFetchedPartition::partition_response(99, 0);
let body1 = ShareFetchedPartition::partition_response(1, crate::error::UNKNOWN_TOPIC_ID);
let body2 = ShareFetchedPartition::partition_response(2, 0);
let grouped = ShareFetchResponse::to_message(&[
(a, 0, body0.clone()),
(b, 1, body1.clone()),
(a, 3, body2.clone()),
]);
assert_eq!(
grouped,
vec![
ShareFetchedTopic {
topic_id: a,
partitions: vec![
ShareFetchedPartition {
partition: 0,
error_code: 0,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
},
ShareFetchedPartition {
partition: 3,
error_code: 0,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
},
],
},
ShareFetchedTopic {
topic_id: b,
partitions: vec![ShareFetchedPartition {
partition: 1,
error_code: crate::error::UNKNOWN_TOPIC_ID,
error_message: None,
acknowledge_error_code: 0,
acknowledge_error_message: None,
current_leader_id: 0,
current_leader_epoch: 0,
records: Vec::new(),
acquired: Vec::new(),
}],
},
]
);
assert_eq!(
grouped
.first()
.and_then(|topic| topic.partitions.first())
.map(|part| part.partition),
Some(0),
"setPartitionIndex copies the key partition onto the body"
);
assert_eq!(body0.partition, 99);
for version in [0i16, 1] {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, &grouped).unwrap();
let mut cur = buf.as_ref();
let (decoded, ..) = decode_share_fetch_response(&mut cur, version).unwrap();
assert_eq!(decoded, grouped);
assert!(
!cur.has_remaining(),
"ShareFetch v{version} toMessage leftover-empty; leftover {} bytes",
cur.remaining()
);
}
}
#[test]
fn share_fetch_response_size_of_matches_java() {
let empty = ShareFetchResponse::size_of(0, &[]).unwrap();
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, 0, &[]).unwrap();
let body_len = buf::i32_from_usize(buf.len()).unwrap();
assert_eq!(empty, body_len + 4, "sizeOf is 4 plus the encoded body");
assert_ne!(empty, body_len);
leftover_share_fetch_size_of(0, &[]);
let a = [1u8; 16];
let b = [2u8; 16];
let body = ShareFetchedPartition::partition_response(99, 0);
let entries = [
(a, 0, body.clone()),
(b, 1, body.clone()),
(a, 3, body.clone()),
];
let topics = ShareFetchResponse::to_message(&entries);
assert_eq!(topics.len(), 2, "non-adjacent same topicId still merges");
for version in [0i16, 1] {
let size = ShareFetchResponse::size_of(version, &entries).unwrap();
buf.clear();
encode_share_fetch_response(&mut buf, version, &topics).unwrap();
let encoded_len = buf::i32_from_usize(buf.len()).unwrap();
assert_eq!(size, encoded_len + 4);
leftover_share_fetch_size_of(version, &topics);
}
leftover_share_fetch_size_of(0, &[]);
leftover_share_fetch_size_of(1, &[]);
leftover_share_fetch_size_of(0, &topics);
leftover_share_fetch_size_of(1, &topics);
let consecutive = [
(a, 0, body.clone()),
(a, 3, body.clone()),
(b, 1, body.clone()),
];
assert_eq!(
ShareFetchResponse::size_of(0, &entries).unwrap(),
ShareFetchResponse::size_of(0, &consecutive).unwrap(),
"LinkedHashMap by topicId merges non-adjacent entries"
);
buf.clear();
encode_share_fetch_response_with_throttle(&mut buf, 1, &topics, 3_600_000).unwrap();
let mut convenience = BytesMut::new();
encode_share_fetch_response(&mut convenience, 1, &topics).unwrap();
assert_eq!(
buf.len(),
convenience.len(),
"ThrottleTimeMs is a fixed-width INT32"
);
buf.clear();
encode_share_fetch_response_with_acquisition_lock_timeout(&mut buf, 1, &topics, 0).unwrap();
assert_eq!(
buf.len(),
convenience.len(),
"AcquisitionLockTimeoutMs is a fixed-width INT32"
);
buf.clear();
encode_share_fetch_response_with_error_code(
&mut buf,
1,
&topics,
crate::error::GROUP_AUTHORIZATION_FAILED,
)
.unwrap();
assert_eq!(
buf.len(),
convenience.len(),
"ErrorCode is a fixed-width INT16"
);
}
fn leftover_share_fetch_size_of(version: i16, topics: &[ShareFetchedTopic]) {
let mut buf = BytesMut::new();
encode_share_fetch_response(&mut buf, version, topics).unwrap();
let mut cur = buf.as_ref();
drop(decode_share_fetch_response(&mut cur, version).unwrap());
leftover_share_fetch_size_of_cur(version, cur);
}
fn leftover_share_fetch_size_of_cur(version: i16, cur: &[u8]) {
let msg = match (version, cur.is_empty()) {
(0, true) => "ShareFetch v0 Response.sizeOf leftover-empty",
(1, true) => "ShareFetch v1 Response.sizeOf leftover-empty",
_ => "ShareFetch Response.sizeOf leftover-empty",
};
assert!(cur.is_empty(), "{msg}; leftover {} bytes", cur.len());
}
#[test]
fn share_fetch_response_of_matches_java() {
let a = [1u8; 16];
let b = [2u8; 16];
let body = ShareFetchedPartition::partition_response(99, 0);
let entries = [
(a, 0, body.clone()),
(b, 1, body.clone()),
(a, 3, body.clone()),
];
let topics = ShareFetchResponse::to_message(&entries);
assert_eq!(topics.len(), 2, "non-adjacent same topicId still merges");
let mut v0_none = BytesMut::new();
ShareFetchResponse::of(&mut v0_none, 0, 0, 0, &entries, &[]).unwrap();
let mut v0_conv = BytesMut::new();
encode_share_fetch_response(&mut v0_conv, 0, &topics).unwrap();
assert_eq!(
v0_none, v0_conv,
"v0 of(NONE, 0, empty) matches convenience encode"
);
let mut v1_none = BytesMut::new();
ShareFetchResponse::of(&mut v1_none, 1, 0, 0, &entries, &[]).unwrap();
let mut v1_conv = BytesMut::new();
encode_share_fetch_response(&mut v1_conv, 1, &topics).unwrap();
assert_ne!(
v1_none, v1_conv,
"v1 of writes AcquisitionLockTimeoutMs 0; convenience encode writes 15000"
);
let err = crate::error::GROUP_AUTHORIZATION_FAILED;
let throttle = 3_600_000;
for version in [0i16, 1] {
let mut buf = BytesMut::new();
ShareFetchResponse::of(&mut buf, version, err, throttle, &entries, &[]).unwrap();
let mut cur = buf.as_ref();
let (decoded, endpoints, throttle_time_ms, error_message, lock_ms, error_code) =
decode_share_fetch_response(&mut cur, version).unwrap();
leftover_share_fetch_of(version, cur);
assert_eq!(decoded, topics);
assert!(endpoints.is_empty());
assert_eq!(throttle_time_ms, throttle);
assert_eq!(error_message, None);
assert_eq!(error_code, err);
assert_eq!(lock_ms, 0);
}
leftover_share_fetch_of_encode(0, err, throttle, &entries, &[]);
leftover_share_fetch_of_encode(1, err, throttle, &entries, &[]);
leftover_share_fetch_of_encode(0, 0, 0, &[], &[]);
leftover_share_fetch_of_encode(1, 0, 0, &[], &[]);
let mut of_buf = BytesMut::new();
ShareFetchResponse::of(&mut of_buf, 1, err, throttle, &entries, &[]).unwrap();
let mut with_throttle = BytesMut::new();
encode_share_fetch_response_with_throttle(&mut with_throttle, 1, &topics, throttle)
.unwrap();
assert_ne!(
of_buf, with_throttle,
"of writes ErrorCode and AcquisitionLockTimeoutMs 0"
);
let mut with_err = BytesMut::new();
encode_share_fetch_response_with_error_code(&mut with_err, 1, &topics, err).unwrap();
assert_ne!(
of_buf, with_err,
"of writes ThrottleTimeMs from the argument"
);
let endpoint = NodeEndpoint::new(1, "h", 9092, None);
let mut with_ep = BytesMut::new();
ShareFetchResponse::of(
&mut with_ep,
1,
err,
throttle,
&entries,
std::slice::from_ref(&endpoint),
)
.unwrap();
leftover_share_fetch_of_encode(1, err, throttle, &entries, std::slice::from_ref(&endpoint));
let mut conv_ep = BytesMut::new();
encode_share_fetch_response_with_endpoints(
&mut conv_ep,
1,
&topics,
std::slice::from_ref(&endpoint),
)
.unwrap();
assert_ne!(
with_ep, conv_ep,
"of writes throttle / ErrorCode / AcquisitionLockTimeoutMs 0; with_endpoints writes 0 / 0 / 15000"
);
}
fn leftover_share_fetch_of_encode(
version: i16,
error_code: i16,
throttle_time_ms: i32,
entries: &[([u8; 16], i32, ShareFetchedPartition)],
endpoints: &[NodeEndpoint],
) {
let mut buf = BytesMut::new();
ShareFetchResponse::of(
&mut buf,
version,
error_code,
throttle_time_ms,
entries,
endpoints,
)
.unwrap();
let mut cur = buf.as_ref();
drop(decode_share_fetch_response(&mut cur, version).unwrap());
leftover_share_fetch_of(version, cur);
}
fn leftover_share_fetch_of(version: i16, cur: &[u8]) {
let msg = match (version, cur.is_empty()) {
(0, true) => "ShareFetch v0 Response.of leftover-empty",
(1, true) => "ShareFetch v1 Response.of leftover-empty",
_ => "ShareFetch Response.of leftover-empty",
};
assert!(cur.is_empty(), "{msg}; leftover {} bytes", cur.len());
}
#[test]
fn share_acknowledge_response_of_matches_java() {
let a = [1u8; 16];
let b = [2u8; 16];
let body = ShareAcknowledgeResponsePartition::partition_response(99, 0);
let entries = [
(a, 0, body.clone()),
(b, 1, body.clone()),
(a, 3, body.clone()),
];
let topics = ShareAcknowledgeResponse::to_message(&entries);
assert_eq!(topics.len(), 2, "non-adjacent same topicId still merges");
let mut none = BytesMut::new();
ShareAcknowledgeResponse::of(&mut none, 0, 0, 0, &entries, &[]).unwrap();
let mut conv = BytesMut::new();
encode_share_acknowledge_topics_response(&mut conv, 0, 0, &topics).unwrap();
assert_eq!(none, conv, "of(NONE, 0, empty) matches convenience encode");
let mut v1_none = BytesMut::new();
ShareAcknowledgeResponse::of(&mut v1_none, 1, 0, 0, &entries, &[]).unwrap();
assert_eq!(none, v1_none, "v0 and v1 of bodies match");
let err = crate::error::GROUP_AUTHORIZATION_FAILED;
let throttle = 3_600_000;
for version in [0i16, 1] {
let mut buf = BytesMut::new();
ShareAcknowledgeResponse::of(&mut buf, version, err, throttle, &entries, &[]).unwrap();
let mut cur = buf.as_ref();
let (error_code, decoded, endpoints, throttle_time_ms, error_message) =
decode_share_acknowledge_topics_response(&mut cur, version).unwrap();
leftover_share_acknowledge_of(version, cur);
assert_eq!(decoded, topics);
assert!(endpoints.is_empty());
assert_eq!(throttle_time_ms, throttle);
assert_eq!(error_message, None);
assert_eq!(error_code, err);
}
leftover_share_acknowledge_of_encode(0, err, throttle, &entries, &[]);
leftover_share_acknowledge_of_encode(1, err, throttle, &entries, &[]);
leftover_share_acknowledge_of_encode(0, 0, 0, &[], &[]);
leftover_share_acknowledge_of_encode(1, 0, 0, &[], &[]);
let mut of_throttle = BytesMut::new();
ShareAcknowledgeResponse::of(&mut of_throttle, 1, err, throttle, &entries, &[]).unwrap();
let mut with_throttle = BytesMut::new();
encode_share_acknowledge_topics_response_with_throttle(
&mut with_throttle,
1,
err,
&topics,
throttle,
)
.unwrap();
assert_eq!(
of_throttle, with_throttle,
"of with empty endpoints matches with_throttle"
);
let endpoint = NodeEndpoint::new(1, "h", 9092, None);
let mut of_ep = BytesMut::new();
ShareAcknowledgeResponse::of(
&mut of_ep,
1,
err,
0,
&entries,
std::slice::from_ref(&endpoint),
)
.unwrap();
leftover_share_acknowledge_of_encode(1, err, 0, &entries, std::slice::from_ref(&endpoint));
let mut with_ep = BytesMut::new();
encode_share_acknowledge_topics_response_with_endpoints(
&mut with_ep,
1,
err,
&topics,
std::slice::from_ref(&endpoint),
)
.unwrap();
assert_eq!(of_ep, with_ep, "of with throttle 0 matches with_endpoints");
let mut of_both = BytesMut::new();
ShareAcknowledgeResponse::of(
&mut of_both,
1,
err,
throttle,
&entries,
std::slice::from_ref(&endpoint),
)
.unwrap();
leftover_share_acknowledge_of_encode(
0,
err,
throttle,
&entries,
std::slice::from_ref(&endpoint),
);
leftover_share_acknowledge_of_encode(
1,
err,
throttle,
&entries,
std::slice::from_ref(&endpoint),
);
assert_ne!(
of_both, of_throttle,
"of writes NodeEndpoints from the argument"
);
assert_ne!(of_both, of_ep, "of writes ThrottleTimeMs from the argument");
let mut error_resp = BytesMut::new();
ShareAcknowledgeRequest::error_response(&mut error_resp, 1, err, throttle).unwrap();
assert_ne!(
of_both, error_resp,
"Request.getErrorResponse writes empty Responses"
);
}
fn leftover_share_acknowledge_of_encode(
version: i16,
error_code: i16,
throttle_time_ms: i32,
entries: &[([u8; 16], i32, ShareAcknowledgeResponsePartition)],
endpoints: &[NodeEndpoint],
) {
let mut buf = BytesMut::new();
ShareAcknowledgeResponse::of(
&mut buf,
version,
error_code,
throttle_time_ms,
entries,
endpoints,
)
.unwrap();
let mut cur = buf.as_ref();
drop(decode_share_acknowledge_topics_response(&mut cur, version).unwrap());
leftover_share_acknowledge_of(version, cur);
}
fn leftover_share_acknowledge_of(version: i16, cur: &[u8]) {
let msg = match (version, cur.is_empty()) {
(0, true) => "ShareAcknowledge v0 Response.of leftover-empty",
(1, true) => "ShareAcknowledge v1 Response.of leftover-empty",
_ => "ShareAcknowledge Response.of leftover-empty",
};
assert!(cur.is_empty(), "{msg}; leftover {} bytes", cur.len());
}
}