use std::collections::HashMap;
use bytes::{Buf, BufMut, BytesMut};
use super::buf;
use crate::error::{Error, Result};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TopicPartitions {
pub topic_id: [u8; 16],
pub partitions: Vec<i32>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerGroupHeartbeatRequest {
pub group_id: String,
pub member_id: String,
pub member_epoch: i32,
pub instance_id: Option<String>,
pub rack_id: Option<String>,
pub rebalance_timeout_ms: i32,
pub subscribed_topic_names: Option<Vec<String>>,
pub subscribed_topic_regex: Option<String>,
pub server_assignor: Option<String>,
pub topic_partitions: Option<Vec<TopicPartitions>>,
}
impl ConsumerGroupHeartbeatRequest {
pub const LEAVE_GROUP_MEMBER_EPOCH: i32 = -1;
pub const LEAVE_GROUP_STATIC_MEMBER_EPOCH: i32 = -2;
pub const JOIN_GROUP_MEMBER_EPOCH: i32 = 0;
pub const UNCHANGED_REBALANCE_TIMEOUT_MS: i32 = -1;
pub const CONSUMER_GENERATED_MEMBER_ID_REQUIRED_VERSION: i16 = 1;
pub const REGEX_RESOLUTION_NOT_SUPPORTED_MSG: &'static str = "The cluster does not support regular expressions resolution on ConsumerGroupHeartbeat API version 0. It must be upgraded to use ConsumerGroupHeartbeat API version >= 1 to allow to subscribe to a SubscriptionPattern.";
#[must_use]
pub fn error_response(
error_code: i16,
throttle_time_ms: i32,
) -> ConsumerGroupHeartbeatResponse {
ConsumerGroupHeartbeatResponse {
throttle_time_ms,
error_code,
error_message: None,
member_id: None,
member_epoch: 0,
heartbeat_interval_ms: 0,
assignment: None,
}
}
pub fn build(version: i16, subscribed_topic_regex: Option<&str>) -> Result<()> {
if version == 0 && subscribed_topic_regex.is_some() {
return Err(Error::Unsupported(
Self::REGEX_RESOLUTION_NOT_SUPPORTED_MSG.into(),
));
}
Ok(())
}
#[must_use]
pub const fn leave_group_epoch(group_instance_id: Option<&str>) -> i32 {
match group_instance_id {
Some(_) => Self::LEAVE_GROUP_STATIC_MEMBER_EPOCH,
None => Self::LEAVE_GROUP_MEMBER_EPOCH,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerGroupHeartbeatResponse {
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<TopicPartitions>>,
}
impl ConsumerGroupHeartbeatResponse {
#[must_use]
pub fn error_counts(&self) -> HashMap<i16, i32> {
HashMap::from([(self.error_code, 1)])
}
}
fn consumer_group_heartbeat_spoken(version: i16) -> Result<i16> {
match version {
0..=1 => Ok(version),
other => Err(Error::protocol(format!(
"ConsumerGroupHeartbeat version {other} is not implemented"
))),
}
}
pub fn encode_consumer_group_heartbeat_request(
buf: &mut BytesMut,
version: i16,
req: &ConsumerGroupHeartbeatRequest,
) -> crate::error::Result<()> {
let _ = consumer_group_heartbeat_spoken(version)?;
ConsumerGroupHeartbeatRequest::build(version, req.subscribed_topic_regex.as_deref())?;
buf::put_compact_string(buf, Some(&req.group_id))?;
buf::put_compact_string(buf, Some(&req.member_id))?;
buf.put_i32(req.member_epoch);
buf::put_compact_string(buf, req.instance_id.as_deref())?;
buf::put_compact_string(buf, req.rack_id.as_deref())?;
buf.put_i32(req.rebalance_timeout_ms);
match &req.subscribed_topic_names {
None => buf::put_array_len(buf, true, None)?,
Some(names) => {
buf::put_array_len(buf, true, Some(names.len()))?;
for n in names {
buf::put_compact_string(buf, Some(n))?;
}
}
}
if version >= 1 {
buf::put_compact_string(buf, req.subscribed_topic_regex.as_deref())?;
}
buf::put_compact_string(buf, req.server_assignor.as_deref())?;
encode_topic_partitions(buf, req.topic_partitions.as_deref())?;
buf::put_empty_tagged_fields(buf);
Ok(())
}
pub fn decode_consumer_group_heartbeat_request<B: Buf>(
buf: &mut B,
version: i16,
) -> Result<ConsumerGroupHeartbeatRequest> {
let _ = consumer_group_heartbeat_spoken(version)?;
let group_id = buf::get_compact_string(buf)?.unwrap_or_default();
let member_id = buf::get_compact_string(buf)?.unwrap_or_default();
let member_epoch = buf::get_i32(buf)?;
let instance_id = buf::get_compact_string(buf)?;
let rack_id = buf::get_compact_string(buf)?;
let rebalance_timeout_ms = buf::get_i32(buf)?;
let subscribed_topic_names = {
let n = buf::get_array_len(buf, true)?;
match n {
None => None,
Some(n) => {
let mut names = Vec::with_capacity(n);
for _ in 0..n {
names.push(buf::get_compact_string(buf)?.unwrap_or_default());
}
Some(names)
}
}
};
let subscribed_topic_regex = if version >= 1 {
buf::get_compact_string(buf)?
} else {
None
};
let server_assignor = buf::get_compact_string(buf)?;
let topic_partitions = decode_topic_partitions(buf)?;
buf::skip_tagged_fields(buf)?;
Ok(ConsumerGroupHeartbeatRequest {
group_id,
member_id,
member_epoch,
instance_id,
rack_id,
rebalance_timeout_ms,
subscribed_topic_names,
subscribed_topic_regex,
server_assignor,
topic_partitions,
})
}
pub fn encode_consumer_group_heartbeat_response(
buf: &mut BytesMut,
version: i16,
resp: &ConsumerGroupHeartbeatResponse,
) -> crate::error::Result<()> {
let _ = consumer_group_heartbeat_spoken(version)?;
buf.put_i32(resp.throttle_time_ms);
buf.put_i16(resp.error_code);
buf::put_compact_string(buf, resp.error_message.as_deref())?;
buf::put_compact_string(buf, resp.member_id.as_deref())?;
buf.put_i32(resp.member_epoch);
buf.put_i32(resp.heartbeat_interval_ms);
match &resp.assignment {
None => buf::put_unsigned_varint(buf, 0),
Some(parts) => {
buf::put_unsigned_varint(buf, 1);
encode_topic_partitions(buf, Some(parts))?;
buf::put_empty_tagged_fields(buf);
}
}
buf::put_empty_tagged_fields(buf);
Ok(())
}
pub fn decode_consumer_group_heartbeat_response<B: Buf>(
buf: &mut B,
version: i16,
) -> Result<ConsumerGroupHeartbeatResponse> {
let _ = consumer_group_heartbeat_spoken(version)?;
let throttle_time_ms = buf::get_i32(buf)?;
let error_code = buf::get_i16(buf)?;
let error_message = buf::get_compact_string(buf)?;
let member_id = buf::get_compact_string(buf)?;
let member_epoch = buf::get_i32(buf)?;
let heartbeat_interval_ms = buf::get_i32(buf)?;
let present = buf::get_unsigned_varint(buf)?;
let assignment = if present == 0 {
None
} else {
let parts = decode_topic_partitions(buf)?;
buf::skip_tagged_fields(buf)?;
parts
};
if buf.has_remaining() {
buf::skip_tagged_fields(buf)?;
}
Ok(ConsumerGroupHeartbeatResponse {
throttle_time_ms,
error_code,
error_message,
member_id,
member_epoch,
heartbeat_interval_ms,
assignment,
})
}
fn encode_topic_partitions(
buf: &mut BytesMut,
parts: Option<&[TopicPartitions]>,
) -> crate::error::Result<()> {
match parts {
None => buf::put_array_len(buf, true, None)?,
Some(parts) => {
buf::put_array_len(buf, true, Some(parts.len()))?;
for t in parts {
buf.extend_from_slice(&t.topic_id);
buf::put_array_len(buf, true, Some(t.partitions.len()))?;
for p in &t.partitions {
buf.put_i32(*p);
}
buf::put_empty_tagged_fields(buf);
}
}
}
Ok(())
}
fn decode_topic_partitions<B: Buf>(buf: &mut B) -> Result<Option<Vec<TopicPartitions>>> {
let n = buf::get_array_len(buf, true)?;
let Some(n) = n else {
return Ok(None);
};
let mut out = Vec::with_capacity(n);
for _ in 0..n {
let topic_id = buf::get_uuid(buf)?;
let pn = buf::get_array_len(buf, true)?.unwrap_or(0);
let mut partitions = Vec::with_capacity(pn);
for _ in 0..pn {
partitions.push(buf::get_i32(buf)?);
}
buf::skip_tagged_fields(buf)?;
out.push(TopicPartitions {
topic_id,
partitions,
});
}
Ok(Some(out))
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
use bytes::{Buf, BufMut};
fn join_req() -> ConsumerGroupHeartbeatRequest {
ConsumerGroupHeartbeatRequest {
group_id: "g".into(),
member_id: String::new(),
member_epoch: ConsumerGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH,
instance_id: Some("worker-1".into()),
rack_id: Some("az1".into()),
rebalance_timeout_ms: 45_000,
subscribed_topic_names: Some(vec!["t".into()]),
subscribed_topic_regex: None,
server_assignor: None,
topic_partitions: None,
}
}
#[test]
fn join_encodes_empty_topic_partitions_array_not_null() {
let req = ConsumerGroupHeartbeatRequest {
group_id: "g".into(),
member_id: String::new(),
member_epoch: ConsumerGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH,
instance_id: None,
rack_id: None,
rebalance_timeout_ms: 300_000,
subscribed_topic_names: Some(vec!["t".into()]),
subscribed_topic_regex: None,
server_assignor: None,
topic_partitions: Some(Vec::new()),
};
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut buf, 1, &req).unwrap();
let mut cur = &buf[..];
let decoded = decode_consumer_group_heartbeat_request(&mut cur, 1).unwrap();
assert_eq!(decoded.topic_partitions.as_ref().map(|v| v.len()), Some(0));
assert!(cur.is_empty());
}
#[test]
fn decode_accepts_error_response_without_trailing_tagged_fields() {
let mut buf = BytesMut::new();
buf.put_i32(0); buf.put_i16(42); buf::put_compact_string(&mut buf, Some("must be empty when (re-)joining")).unwrap();
buf::put_compact_string(&mut buf, None).unwrap(); buf.put_i32(0); buf.put_i32(5000); buf::put_unsigned_varint(&mut buf, 0); let decoded = decode_consumer_group_heartbeat_response(&mut &buf[..], 0).unwrap();
assert_eq!(decoded.error_code, 42);
assert!(decoded.assignment.is_none());
}
#[test]
fn consumer_group_heartbeat_v0_roundtrip_join() {
let req = join_req();
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut buf, 0, &req).unwrap();
assert_eq!(
decode_consumer_group_heartbeat_request(&mut &buf[..], 0).unwrap(),
req
);
let topic_id = [7u8; 16];
let resp = ConsumerGroupHeartbeatResponse {
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![TopicPartitions {
topic_id,
partitions: vec![0, 1],
}]),
};
buf.clear();
encode_consumer_group_heartbeat_response(&mut buf, 0, &resp).unwrap();
assert_eq!(
decode_consumer_group_heartbeat_response(&mut &buf[..], 0).unwrap(),
resp
);
}
#[test]
fn consumer_group_heartbeat_leave_has_epoch_minus_one() {
let req = ConsumerGroupHeartbeatRequest {
group_id: "g".into(),
member_id: "m1".into(),
member_epoch: ConsumerGroupHeartbeatRequest::LEAVE_GROUP_MEMBER_EPOCH,
instance_id: None,
rack_id: None,
rebalance_timeout_ms: 45_000,
subscribed_topic_names: None,
subscribed_topic_regex: None,
server_assignor: None,
topic_partitions: None,
};
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut buf, 0, &req).unwrap();
let decoded = decode_consumer_group_heartbeat_request(&mut &buf[..], 0).unwrap();
assert_eq!(
decoded.member_epoch,
ConsumerGroupHeartbeatRequest::LEAVE_GROUP_MEMBER_EPOCH
);
assert_eq!(decoded.member_id, "m1");
}
#[test]
fn consumer_group_heartbeat_static_leave_has_epoch_minus_two() {
let req = ConsumerGroupHeartbeatRequest {
group_id: "g".into(),
member_id: "m1".into(),
member_epoch: ConsumerGroupHeartbeatRequest::leave_group_epoch(Some("worker-1")),
instance_id: Some("worker-1".into()),
rack_id: None,
rebalance_timeout_ms: 45_000,
subscribed_topic_names: None,
subscribed_topic_regex: None,
server_assignor: None,
topic_partitions: None,
};
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut buf, 0, &req).unwrap();
let decoded = decode_consumer_group_heartbeat_request(&mut &buf[..], 0).unwrap();
assert_eq!(
decoded.member_epoch,
ConsumerGroupHeartbeatRequest::LEAVE_GROUP_STATIC_MEMBER_EPOCH
);
assert_eq!(decoded.instance_id.as_deref(), Some("worker-1"));
}
#[test]
fn consumer_group_heartbeat_request_matches_java() {
assert_eq!(ConsumerGroupHeartbeatRequest::LEAVE_GROUP_MEMBER_EPOCH, -1);
assert_eq!(
ConsumerGroupHeartbeatRequest::LEAVE_GROUP_STATIC_MEMBER_EPOCH,
-2
);
assert_eq!(ConsumerGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH, 0);
assert_eq!(
ConsumerGroupHeartbeatRequest::CONSUMER_GENERATED_MEMBER_ID_REQUIRED_VERSION,
1
);
assert_eq!(
ConsumerGroupHeartbeatRequest::leave_group_epoch(None),
ConsumerGroupHeartbeatRequest::LEAVE_GROUP_MEMBER_EPOCH
);
assert_eq!(
ConsumerGroupHeartbeatRequest::leave_group_epoch(Some("worker-1")),
ConsumerGroupHeartbeatRequest::LEAVE_GROUP_STATIC_MEMBER_EPOCH
);
assert_eq!(
ConsumerGroupHeartbeatRequest::leave_group_epoch(Some("")),
ConsumerGroupHeartbeatRequest::LEAVE_GROUP_STATIC_MEMBER_EPOCH,
"Java Optional.isPresent is true for empty group.instance.id"
);
}
#[test]
fn consumer_group_heartbeat_v1_compact_layout_matches_independent_encode() {
const REQ_V0: &[u8] = &[
0x02, 0x67, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0xaf, 0xc8, 0x02,
0x02, 0x74, 0x00, 0x00, 0x00,
];
const REQ_V1: &[u8] = &[
0x02, 0x67, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0xaf, 0xc8, 0x02,
0x02, 0x74, 0x00, 0x00, 0x00, 0x00,
];
const REQ_V1_REGEX: &[u8] = &[
0x02, 0x67, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0xaf, 0xc8, 0x02,
0x02, 0x74, 0x04, 0x74, 0x2e, 0x2a, 0x00, 0x00, 0x00,
];
let req = ConsumerGroupHeartbeatRequest {
group_id: "g".into(),
member_id: String::new(),
member_epoch: ConsumerGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH,
instance_id: None,
rack_id: None,
rebalance_timeout_ms: 45_000,
subscribed_topic_names: Some(vec!["t".into()]),
subscribed_topic_regex: None,
server_assignor: None,
topic_partitions: None,
};
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut buf, 0, &req).unwrap();
assert_eq!(&buf[..], REQ_V0);
buf.clear();
encode_consumer_group_heartbeat_request(&mut buf, 1, &req).unwrap();
assert_eq!(&buf[..], REQ_V1);
let mut with_regex = req.clone();
with_regex.subscribed_topic_regex = Some("t.*".into());
buf.clear();
encode_consumer_group_heartbeat_request(&mut buf, 1, &with_regex).unwrap();
assert_eq!(&buf[..], REQ_V1_REGEX);
assert_eq!(
decode_consumer_group_heartbeat_request(&mut &buf[..], 1).unwrap(),
with_regex
);
let err = encode_consumer_group_heartbeat_request(&mut BytesMut::new(), 0, &with_regex)
.unwrap_err();
assert!(
matches!(err, Error::Unsupported(_)),
"regex on v0 is Java UnsupportedVersionException, got {err}"
);
assert!(
err.to_string()
.contains(ConsumerGroupHeartbeatRequest::REGEX_RESOLUTION_NOT_SUPPORTED_MSG),
"got {err}"
);
assert!(
encode_consumer_group_heartbeat_request(&mut BytesMut::new(), 2, &req).is_err(),
"ConsumerGroupHeartbeat v2+ is not spoken"
);
buf.clear();
encode_consumer_group_heartbeat_response(
&mut buf,
1,
&ConsumerGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: 0,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: None,
},
)
.unwrap();
let mut v0 = BytesMut::new();
encode_consumer_group_heartbeat_response(
&mut v0,
0,
&ConsumerGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: 0,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: None,
},
)
.unwrap();
assert_eq!(&buf[..], &v0[..], "v1 response layout matches v0");
}
#[test]
fn consumer_group_heartbeat_v1_roundtrip_is_leftover_empty() {
let req = ConsumerGroupHeartbeatRequest {
group_id: "g".into(),
member_id: "m1".into(),
member_epoch: 1,
instance_id: None,
rack_id: None,
rebalance_timeout_ms: 45_000,
subscribed_topic_names: None,
subscribed_topic_regex: Some("t.*".into()),
server_assignor: None,
topic_partitions: None,
};
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut buf, 1, &req).unwrap();
let mut cur = &buf[..];
assert_eq!(
decode_consumer_group_heartbeat_request(&mut cur, 1).unwrap(),
req
);
assert!(
!cur.has_remaining(),
"ConsumerGroupHeartbeat v1 request must be leftover-empty"
);
let resp = ConsumerGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: 0,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: None,
};
buf.clear();
encode_consumer_group_heartbeat_response(&mut buf, 1, &resp).unwrap();
let mut cur = &buf[..];
assert_eq!(
decode_consumer_group_heartbeat_response(&mut cur, 1).unwrap(),
resp
);
assert!(
!cur.has_remaining(),
"ConsumerGroupHeartbeat v1 response must be leftover-empty"
);
}
#[test]
fn consumer_group_heartbeat_request_rebalance_timeout_ms_matches_java() {
assert_eq!(
ConsumerGroupHeartbeatRequest::UNCHANGED_REBALANCE_TIMEOUT_MS,
-1
);
let mut req = join_req();
req.rebalance_timeout_ms = 300_000;
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut buf, version, &req).unwrap();
let mut cur = buf.as_ref();
let got = decode_consumer_group_heartbeat_request(&mut cur, version).unwrap();
assert_eq!(got.rebalance_timeout_ms, 300_000);
assert_eq!(got, req);
assert!(
cur.is_empty(),
"ConsumerGroupHeartbeat request v{version} RebalanceTimeoutMs leftover-empty"
);
}
let mut with = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut with, 0, &req).unwrap();
let mut unchanged = join_req();
unchanged.rebalance_timeout_ms =
ConsumerGroupHeartbeatRequest::UNCHANGED_REBALANCE_TIMEOUT_MS;
let mut minus_one = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut minus_one, 0, &unchanged).unwrap();
assert_ne!(
&with[..],
&minus_one[..],
"v0 RebalanceTimeoutMs is not always UNCHANGED"
);
let got = decode_consumer_group_heartbeat_request(&mut minus_one.as_ref(), 0).unwrap();
assert_eq!(
got.rebalance_timeout_ms,
ConsumerGroupHeartbeatRequest::UNCHANGED_REBALANCE_TIMEOUT_MS
);
}
#[test]
fn consumer_group_heartbeat_request_server_assignor_matches_java() {
let mut req = join_req();
req.server_assignor = Some("uniform".into());
for version in [0_i16, 1] {
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut buf, version, &req).unwrap();
let mut cur = buf.as_ref();
let got = decode_consumer_group_heartbeat_request(&mut cur, version).unwrap();
assert_eq!(got.server_assignor.as_deref(), Some("uniform"));
assert_eq!(got, req);
assert!(
cur.is_empty(),
"ConsumerGroupHeartbeat request v{version} ServerAssignor leftover-empty"
);
}
let mut with = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut with, 0, &req).unwrap();
let none = join_req();
let mut omitted = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut omitted, 0, &none).unwrap();
assert_ne!(
&with[..],
&omitted[..],
"v0 ServerAssignor is not always null"
);
let got = decode_consumer_group_heartbeat_request(&mut omitted.as_ref(), 0).unwrap();
assert_eq!(got.server_assignor, None);
}
#[test]
fn consumer_group_heartbeat_response_throttle_time_ms_matches_java() {
let zero = ConsumerGroupHeartbeatResponse {
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 = ConsumerGroupHeartbeatResponse {
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_consumer_group_heartbeat_response(&mut buf, version, &with).unwrap();
let mut cur = buf.as_ref();
let got = decode_consumer_group_heartbeat_response(&mut cur, version).unwrap();
assert_eq!(got, with);
assert_eq!(got.throttle_time_ms, 3_600_000);
assert!(
cur.is_empty(),
"ConsumerGroupHeartbeat v{version} ThrottleTimeMs leftover-empty"
);
}
let mut with_buf = BytesMut::new();
encode_consumer_group_heartbeat_response(&mut with_buf, 0, &with).unwrap();
let mut zero_buf = BytesMut::new();
encode_consumer_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_consumer_group_heartbeat_response(&mut v1_with, 1, &with).unwrap();
assert_eq!(
&with_buf[..],
&v1_with[..],
"v0 and v1 both write ThrottleTimeMs (JSON 0+); ConsumerGroupHeartbeat response matches v0"
);
}
#[test]
fn consumer_group_heartbeat_request_error_response_matches_java() {
let err = ConsumerGroupHeartbeatRequest::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!(
ConsumerGroupHeartbeatRequest::error_response(0, 0),
ConsumerGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: 0,
error_message: None,
member_id: None,
member_epoch: 0,
heartbeat_interval_ms: 0,
assignment: None,
}
);
leftover_consumer_group_heartbeat_error_response(0, &err);
leftover_consumer_group_heartbeat_error_response(
0,
&ConsumerGroupHeartbeatRequest::error_response(0, 0),
);
leftover_consumer_group_heartbeat_error_response(1, &err);
leftover_consumer_group_heartbeat_error_response(
1,
&ConsumerGroupHeartbeatRequest::error_response(0, 0),
);
}
fn leftover_consumer_group_heartbeat_error_response(
version: i16,
resp: &ConsumerGroupHeartbeatResponse,
) {
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_response(&mut buf, version, resp).unwrap();
let mut cur = buf.as_ref();
let got = decode_consumer_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(),
"ConsumerGroupHeartbeat v{version} getErrorResponse {empty}leftover-empty; leftover {} bytes",
cur.len()
);
}
#[test]
fn consumer_group_heartbeat_request_build_matches_java() {
assert!(ConsumerGroupHeartbeatRequest::build(0, None).is_ok());
assert!(ConsumerGroupHeartbeatRequest::build(1, None).is_ok());
assert!(ConsumerGroupHeartbeatRequest::build(1, Some("t.*")).is_ok());
let v0 = ConsumerGroupHeartbeatRequest::build(0, Some("t.*")).unwrap_err();
assert!(
matches!(v0, Error::Unsupported(_)),
"regex on v0 is Java UnsupportedVersionException, got {v0}"
);
assert!(
v0.to_string()
.contains(ConsumerGroupHeartbeatRequest::REGEX_RESOLUTION_NOT_SUPPORTED_MSG),
"got {v0}"
);
let empty = ConsumerGroupHeartbeatRequest::build(0, Some("")).unwrap_err();
assert!(
matches!(empty, Error::Unsupported(_)),
"empty regex on v0 is still present, got {empty}"
);
let req = ConsumerGroupHeartbeatRequest {
group_id: "g".into(),
member_id: String::new(),
member_epoch: ConsumerGroupHeartbeatRequest::JOIN_GROUP_MEMBER_EPOCH,
instance_id: None,
rack_id: None,
rebalance_timeout_ms: 45_000,
subscribed_topic_names: Some(vec!["t".into()]),
subscribed_topic_regex: None,
server_assignor: None,
topic_partitions: None,
};
leftover_consumer_group_heartbeat_build(0, &req);
leftover_consumer_group_heartbeat_build(1, &req);
let mut with_regex = req.clone();
with_regex.subscribed_topic_regex = Some("t.*".into());
leftover_consumer_group_heartbeat_build(1, &with_regex);
assert!(
encode_consumer_group_heartbeat_request(&mut BytesMut::new(), 0, &with_regex).is_err(),
"encode rejects regex on v0; Builder.build rejects it too"
);
}
fn leftover_consumer_group_heartbeat_build(version: i16, req: &ConsumerGroupHeartbeatRequest) {
ConsumerGroupHeartbeatRequest::build(version, req.subscribed_topic_regex.as_deref())
.unwrap();
let mut buf = BytesMut::new();
encode_consumer_group_heartbeat_request(&mut buf, version, req).unwrap();
let mut cur = buf.as_ref();
let got = decode_consumer_group_heartbeat_request(&mut cur, version).unwrap();
assert_eq!(got, *req);
let empty = if req.subscribed_topic_regex.is_none() {
"empty "
} else {
""
};
assert!(
cur.is_empty(),
"ConsumerGroupHeartbeat v{version} Builder.build {empty}leftover-empty; leftover {} bytes",
cur.len()
);
}
#[test]
fn consumer_group_heartbeat_response_error_counts_matches_java() {
let none = ConsumerGroupHeartbeatResponse {
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![TopicPartitions {
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 fenced = ConsumerGroupHeartbeatResponse {
throttle_time_ms: 0,
error_code: crate::error::FENCED_MEMBER_EPOCH,
error_message: None,
member_id: Some("m1".into()),
member_epoch: 1,
heartbeat_interval_ms: 5000,
assignment: Some(vec![TopicPartitions {
topic_id: [7u8; 16],
partitions: vec![0, 1],
}]),
};
assert_eq!(
fenced.error_counts(),
HashMap::from([(crate::error::FENCED_MEMBER_EPOCH, 1)])
);
for version in 0..=1_i16 {
let mut resp = BytesMut::new();
encode_consumer_group_heartbeat_response(&mut resp, version, &fenced).unwrap();
let mut cur = &resp[..];
let decoded = decode_consumer_group_heartbeat_response(&mut cur, version).unwrap();
assert_eq!(
decoded.error_counts(),
HashMap::from([(crate::error::FENCED_MEMBER_EPOCH, 1)]),
"ConsumerGroupHeartbeat v{version} errorCounts must count the decoded code"
);
assert!(
cur.is_empty(),
"ConsumerGroupHeartbeat v{version} errorCounts leftover-empty; leftover {} bytes",
cur.len()
);
}
}
}