use bytes::{Buf, BufMut};
use super::{VersionedDecode, VersionedEncode};
use crate::error::{ErrorCode, KrafkaError, ProtocolErrorKind, Result};
use crate::protocol::primitives::{Decode, KafkaString, TaggedFields, TryEncode};
use crate::protocol::{check_compact_array_len, decode_capacity, encode_compact_array_len};
#[derive(Debug, Clone)]
pub struct StreamsGroupDescribeRequest {
pub group_ids: Vec<String>,
pub include_authorized_operations: bool,
}
impl StreamsGroupDescribeRequest {
#[must_use]
pub fn new(group_ids: Vec<String>) -> Self {
Self {
group_ids,
include_authorized_operations: false,
}
}
#[must_use]
pub fn with_authorized_operations(mut self, include: bool) -> Self {
self.include_authorized_operations = include;
self
}
pub fn encode_v0(&self, buf: &mut impl BufMut) -> Result<()> {
encode_compact_array_len(self.group_ids.len(), buf)?;
for id in &self.group_ids {
KafkaString::new(id).try_encode_compact(buf)?;
}
buf.put_u8(u8::from(self.include_authorized_operations));
TaggedFields::default().try_encode(buf)?;
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct StreamsEndpoint {
pub host: String,
pub port: u16,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct StreamsKeyValue {
pub key: String,
pub value: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct StreamsTopicInfo {
pub name: String,
pub partitions: i32,
pub replication_factor: i16,
pub topic_configs: Vec<StreamsKeyValue>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct StreamsSubtopology {
pub subtopology_id: String,
pub source_topics: Vec<String>,
pub repartition_sink_topics: Vec<String>,
pub state_changelog_topics: Vec<StreamsTopicInfo>,
pub repartition_source_topics: Vec<StreamsTopicInfo>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct StreamsTopology {
pub epoch: i32,
pub subtopologies: Option<Vec<StreamsSubtopology>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct StreamsTaskIds {
pub subtopology_id: String,
pub partitions: Vec<i32>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct StreamsAssignment {
pub active_tasks: Vec<StreamsTaskIds>,
pub standby_tasks: Vec<StreamsTaskIds>,
pub warmup_tasks: Vec<StreamsTaskIds>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct StreamsTaskOffset {
pub subtopology_id: String,
pub partition: i32,
pub offset: i64,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct StreamsGroupMember {
pub member_id: String,
pub member_epoch: i32,
pub instance_id: Option<String>,
pub rack_id: Option<String>,
pub client_id: String,
pub client_host: String,
pub topology_epoch: i32,
pub process_id: String,
pub user_endpoint: Option<StreamsEndpoint>,
pub client_tags: Vec<StreamsKeyValue>,
pub task_offsets: Vec<StreamsTaskOffset>,
pub task_end_offsets: Vec<StreamsTaskOffset>,
pub assignment: StreamsAssignment,
pub target_assignment: StreamsAssignment,
pub is_classic: bool,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct DescribedStreamsGroup {
pub error_code: ErrorCode,
pub error_message: Option<String>,
pub group_id: String,
pub group_state: String,
pub group_epoch: i32,
pub assignment_epoch: i32,
pub topology: Option<StreamsTopology>,
pub members: Vec<StreamsGroupMember>,
pub authorized_operations: i32,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct StreamsGroupDescribeResponse {
pub throttle_time_ms: i32,
pub groups: Vec<DescribedStreamsGroup>,
}
fn read_presence(buf: &mut impl Buf, what: &'static str) -> Result<bool> {
if buf.remaining() < 1 {
return Err(KrafkaError::protocol_kind(
ProtocolErrorKind::TruncatedFrame,
format!("not enough bytes for {what} presence tag"),
));
}
let presence = buf.get_i8();
if presence < 0 {
return Ok(false);
}
if presence != 1 {
return Err(KrafkaError::protocol_kind(
ProtocolErrorKind::Malformed,
format!("invalid {what} presence tag {presence}: expected -1 (null) or 1 (present)"),
));
}
Ok(true)
}
fn read_compact_string(buf: &mut impl Buf, field: &'static str) -> Result<String> {
super::non_nullable_string(field, KafkaString::decode_compact(buf)?.0)
}
fn read_compact_string_array(buf: &mut impl Buf, field: &'static str) -> Result<Vec<String>> {
let count = check_compact_array_len(crate::util::varint::decode_unsigned_varint(buf)?)?;
let mut out = Vec::with_capacity(decode_capacity(count, buf.remaining()));
for _ in 0..count {
out.push(read_compact_string(buf, field)?);
}
Ok(out)
}
fn read_key_values(buf: &mut impl Buf) -> Result<Vec<StreamsKeyValue>> {
let count = check_compact_array_len(crate::util::varint::decode_unsigned_varint(buf)?)?;
let mut out = Vec::with_capacity(decode_capacity(count, buf.remaining()));
for _ in 0..count {
let key = read_compact_string(buf, "KeyValue.Key")?;
let value = read_compact_string(buf, "KeyValue.Value")?;
let _ = TaggedFields::decode(buf)?;
out.push(StreamsKeyValue { key, value });
}
Ok(out)
}
fn read_topic_infos(buf: &mut impl Buf) -> Result<Vec<StreamsTopicInfo>> {
let count = check_compact_array_len(crate::util::varint::decode_unsigned_varint(buf)?)?;
let mut out = Vec::with_capacity(decode_capacity(count, buf.remaining()));
for _ in 0..count {
let name = read_compact_string(buf, "TopicInfo.Name")?;
let partitions = i32::decode(buf)?;
let replication_factor = i16::decode(buf)?;
let topic_configs = read_key_values(buf)?;
let _ = TaggedFields::decode(buf)?;
out.push(StreamsTopicInfo {
name,
partitions,
replication_factor,
topic_configs,
});
}
Ok(out)
}
fn read_task_ids(buf: &mut impl Buf) -> Result<Vec<StreamsTaskIds>> {
let count = check_compact_array_len(crate::util::varint::decode_unsigned_varint(buf)?)?;
let mut out = Vec::with_capacity(decode_capacity(count, buf.remaining()));
for _ in 0..count {
let subtopology_id = read_compact_string(buf, "TaskIds.SubtopologyId")?;
let p_count = check_compact_array_len(crate::util::varint::decode_unsigned_varint(buf)?)?;
let mut partitions = Vec::with_capacity(decode_capacity(p_count, buf.remaining()));
for _ in 0..p_count {
partitions.push(i32::decode(buf)?);
}
let _ = TaggedFields::decode(buf)?;
out.push(StreamsTaskIds {
subtopology_id,
partitions,
});
}
Ok(out)
}
fn read_assignment(buf: &mut impl Buf) -> Result<StreamsAssignment> {
let active_tasks = read_task_ids(buf)?;
let standby_tasks = read_task_ids(buf)?;
let warmup_tasks = read_task_ids(buf)?;
let _ = TaggedFields::decode(buf)?;
Ok(StreamsAssignment {
active_tasks,
standby_tasks,
warmup_tasks,
})
}
fn read_task_offsets(buf: &mut impl Buf) -> Result<Vec<StreamsTaskOffset>> {
let count = check_compact_array_len(crate::util::varint::decode_unsigned_varint(buf)?)?;
let mut out = Vec::with_capacity(decode_capacity(count, buf.remaining()));
for _ in 0..count {
let subtopology_id = read_compact_string(buf, "TaskOffset.SubtopologyId")?;
let partition = i32::decode(buf)?;
let offset = i64::decode(buf)?;
let _ = TaggedFields::decode(buf)?;
out.push(StreamsTaskOffset {
subtopology_id,
partition,
offset,
});
}
Ok(out)
}
fn read_topology(buf: &mut impl Buf) -> Result<Option<StreamsTopology>> {
if !read_presence(buf, "Topology")? {
return Ok(None);
}
let epoch = i32::decode(buf)?;
let raw = crate::util::varint::decode_unsigned_varint(buf)?;
let subtopologies = if raw == 0 {
None
} else {
let count = check_compact_array_len(raw)?;
let mut out = Vec::with_capacity(decode_capacity(count, buf.remaining()));
for _ in 0..count {
let subtopology_id = read_compact_string(buf, "Subtopology.SubtopologyId")?;
let source_topics = read_compact_string_array(buf, "Subtopology.SourceTopics")?;
let repartition_sink_topics =
read_compact_string_array(buf, "Subtopology.RepartitionSinkTopics")?;
let state_changelog_topics = read_topic_infos(buf)?;
let repartition_source_topics = read_topic_infos(buf)?;
let _ = TaggedFields::decode(buf)?;
out.push(StreamsSubtopology {
subtopology_id,
source_topics,
repartition_sink_topics,
state_changelog_topics,
repartition_source_topics,
});
}
Some(out)
};
let _ = TaggedFields::decode(buf)?;
Ok(Some(StreamsTopology {
epoch,
subtopologies,
}))
}
fn read_member(buf: &mut impl Buf) -> Result<StreamsGroupMember> {
let member_id = read_compact_string(buf, "Member.MemberId")?;
let member_epoch = i32::decode(buf)?;
let instance_id = KafkaString::decode_compact(buf)?.0;
let rack_id = KafkaString::decode_compact(buf)?.0;
let client_id = read_compact_string(buf, "Member.ClientId")?;
let client_host = read_compact_string(buf, "Member.ClientHost")?;
let topology_epoch = i32::decode(buf)?;
let process_id = read_compact_string(buf, "Member.ProcessId")?;
let user_endpoint = if read_presence(buf, "UserEndpoint")? {
let host = read_compact_string(buf, "Endpoint.Host")?;
if buf.remaining() < 2 {
return Err(KrafkaError::protocol_kind(
ProtocolErrorKind::TruncatedFrame,
"not enough bytes for Endpoint.Port",
));
}
let port = buf.get_u16();
let _ = TaggedFields::decode(buf)?;
Some(StreamsEndpoint { host, port })
} else {
None
};
let client_tags = read_key_values(buf)?;
let task_offsets = read_task_offsets(buf)?;
let task_end_offsets = read_task_offsets(buf)?;
let assignment = read_assignment(buf)?;
let target_assignment = read_assignment(buf)?;
if buf.remaining() < 1 {
return Err(KrafkaError::protocol_kind(
ProtocolErrorKind::TruncatedFrame,
"not enough bytes for Member.IsClassic",
));
}
let is_classic = buf.get_u8() != 0;
let _ = TaggedFields::decode(buf)?;
Ok(StreamsGroupMember {
member_id,
member_epoch,
instance_id,
rack_id,
client_id,
client_host,
topology_epoch,
process_id,
user_endpoint,
client_tags,
task_offsets,
task_end_offsets,
assignment,
target_assignment,
is_classic,
})
}
impl StreamsGroupDescribeResponse {
pub fn decode_v0(buf: &mut impl Buf) -> Result<Self> {
let throttle_time_ms = i32::decode(buf)?;
let group_count =
check_compact_array_len(crate::util::varint::decode_unsigned_varint(buf)?)?;
let mut groups = Vec::with_capacity(decode_capacity(group_count, buf.remaining()));
for _ in 0..group_count {
let error_code = ErrorCode::from_i16(i16::decode(buf)?);
let error_message = KafkaString::decode_compact(buf)?.0;
let group_id = read_compact_string(buf, "DescribedGroup.GroupId")?;
let group_state = read_compact_string(buf, "DescribedGroup.GroupState")?;
let group_epoch = i32::decode(buf)?;
let assignment_epoch = i32::decode(buf)?;
let topology = read_topology(buf)?;
let member_count =
check_compact_array_len(crate::util::varint::decode_unsigned_varint(buf)?)?;
let mut members = Vec::with_capacity(decode_capacity(member_count, buf.remaining()));
for _ in 0..member_count {
members.push(read_member(buf)?);
}
let authorized_operations = i32::decode(buf)?;
let _ = TaggedFields::decode(buf)?;
groups.push(DescribedStreamsGroup {
error_code,
error_message,
group_id,
group_state,
group_epoch,
assignment_epoch,
topology,
members,
authorized_operations,
});
}
let _ = TaggedFields::decode(buf)?;
Ok(Self {
throttle_time_ms,
groups,
})
}
}
impl VersionedEncode for StreamsGroupDescribeRequest {
fn encode_versioned(&self, version: i16, buf: &mut impl BufMut) -> Result<()> {
match version {
0 => self.encode_v0(buf),
_ => unsupported_encode!("StreamsGroupDescribeRequest", version),
}
}
}
impl VersionedDecode for StreamsGroupDescribeResponse {
fn decode_versioned(version: i16, buf: &mut impl Buf) -> Result<Self> {
match version {
0 => Self::decode_v0(buf),
_ => unsupported_decode!("StreamsGroupDescribeResponse", version),
}
}
}