1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 61;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct DescribeProducersRequestV0 {
9 pub correlation_id: i32,
10 pub client_id: Option<String>,
11 pub topics: Vec<DescribeProducersTopicV0>,
12}
13
14impl DescribeProducersRequestV0 {
15 pub fn encode(&self) -> Result<Vec<u8>> {
16 let mut encoder = Encoder::new();
17 RequestHeader {
18 api_key: API_KEY,
19 api_version: 0,
20 correlation_id: self.correlation_id,
21 client_id: self.client_id.clone(),
22 }
23 .encode_v2(&mut encoder)?;
24 encoder.write_compact_array(Some(&self.topics), |encoder, topic| {
25 encoder.write_compact_string(&topic.name)?;
26 encoder.write_compact_array(Some(&topic.partition_indexes), |encoder, partition| {
27 encoder.write_i32(*partition);
28 Ok(())
29 })?;
30 encoder.write_empty_tagged_fields();
31 Ok(())
32 })?;
33 encoder.write_empty_tagged_fields();
34 Ok(encoder.into_bytes())
35 }
36}
37
38#[derive(Debug, Clone, PartialEq, Eq)]
39pub struct DescribeProducersTopicV0 {
40 pub name: String,
41 pub partition_indexes: Vec<i32>,
42}
43
44#[derive(Debug, Clone, PartialEq, Eq)]
45pub struct DescribeProducersResponseV0 {
46 pub throttle_time_ms: i32,
47 pub topics: Vec<DescribeProducersTopicResponseV0>,
48}
49
50impl DescribeProducersResponseV0 {
51 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
52 let throttle_time_ms = decoder.read_i32()?;
53 let topics = decoder
54 .read_compact_array("describe producers topics", |decoder| {
55 let name = decoder.read_compact_string()?;
56 let partitions = decoder
57 .read_compact_array("describe producers partitions", |decoder| {
58 let partition_index = decoder.read_i32()?;
59 let error_code = decoder.read_i16()?;
60 let error_message = decoder.read_compact_nullable_string()?;
61 let active_producers = decoder
62 .read_compact_array(
63 "describe producers active producers",
64 DescribeProducersActiveProducerV0::decode,
65 )?
66 .unwrap_or_default();
67 decoder.read_tagged_fields()?;
68 Ok(DescribeProducersPartitionResponseV0 {
69 partition_index,
70 error_code,
71 error_message,
72 active_producers,
73 })
74 })?
75 .unwrap_or_default();
76 decoder.read_tagged_fields()?;
77 Ok(DescribeProducersTopicResponseV0 { name, partitions })
78 })?
79 .unwrap_or_default();
80 decoder.read_tagged_fields()?;
81 Ok(Self {
82 throttle_time_ms,
83 topics,
84 })
85 }
86}
87
88#[derive(Debug, Clone, PartialEq, Eq)]
89pub struct DescribeProducersTopicResponseV0 {
90 pub name: String,
91 pub partitions: Vec<DescribeProducersPartitionResponseV0>,
92}
93
94#[derive(Debug, Clone, PartialEq, Eq)]
95pub struct DescribeProducersPartitionResponseV0 {
96 pub partition_index: i32,
97 pub error_code: i16,
98 pub error_message: Option<String>,
99 pub active_producers: Vec<DescribeProducersActiveProducerV0>,
100}
101
102#[derive(Debug, Clone, PartialEq, Eq)]
103pub struct DescribeProducersActiveProducerV0 {
104 pub producer_id: i64,
105 pub producer_epoch: i32,
106 pub last_sequence: i32,
107 pub last_timestamp: i64,
108 pub coordinator_epoch: i32,
109 pub current_txn_start_offset: i64,
110}
111
112impl DescribeProducersActiveProducerV0 {
113 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
114 let result = Self {
115 producer_id: decoder.read_i64()?,
116 producer_epoch: decoder.read_i32()?,
117 last_sequence: decoder.read_i32()?,
118 last_timestamp: decoder.read_i64()?,
119 coordinator_epoch: decoder.read_i32()?,
120 current_txn_start_offset: decoder.read_i64()?,
121 };
122 decoder.read_tagged_fields()?;
123 Ok(result)
124 }
125}
126
127#[cfg(test)]
128#[allow(clippy::unwrap_used)]
129mod tests {
130 use super::{
131 DescribeProducersRequestV0, DescribeProducersResponseV0, DescribeProducersTopicV0, API_KEY,
132 };
133 use crate::codec::{Decoder, Encoder};
134
135 #[test]
136 fn encodes_describe_producers_v0_request() {
137 let request = DescribeProducersRequestV0 {
138 correlation_id: 61,
139 client_id: Some("kafrust".to_owned()),
140 topics: vec![DescribeProducersTopicV0 {
141 name: "orders".to_owned(),
142 partition_indexes: vec![0, 2],
143 }],
144 };
145
146 let bytes = request.encode().unwrap();
147 assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 0]);
148 assert_eq!(&bytes[4..8], &[0, 0, 0, 61]);
149 assert_eq!(&bytes[8..10], &[0, 7]); assert_eq!(bytes[18], 2); assert_eq!(bytes[26], 3); assert_eq!(bytes.last(), Some(&0));
153 }
154
155 #[test]
156 fn decodes_describe_producers_v0_response() {
157 let mut bytes = Encoder::new();
158 bytes.write_i32(12);
159 bytes.write_unsigned_varint(2);
160 bytes.write_compact_string("orders").unwrap();
161 bytes.write_unsigned_varint(3);
162 bytes.write_i32(0);
163 bytes.write_i16(0);
164 bytes.write_compact_nullable_string(None).unwrap();
165 bytes.write_unsigned_varint(2);
166 bytes.write_i64(42);
167 bytes.write_i32(3);
168 bytes.write_i32(17);
169 bytes.write_i64(1_700_000_000_000);
170 bytes.write_i32(9);
171 bytes.write_i64(-1);
172 bytes.write_empty_tagged_fields();
173 bytes.write_empty_tagged_fields();
174 bytes.write_i32(1);
175 bytes.write_i16(29);
176 bytes.write_compact_nullable_string(Some("denied")).unwrap();
177 bytes.write_unsigned_varint(1);
178 bytes.write_empty_tagged_fields();
179 bytes.write_empty_tagged_fields();
180 bytes.write_empty_tagged_fields();
181 let bytes = bytes.into_bytes();
182 let mut decoder = Decoder::new(&bytes);
183
184 let response = DescribeProducersResponseV0::decode_body(&mut decoder).unwrap();
185
186 assert_eq!(response.throttle_time_ms, 12);
187 assert_eq!(response.topics[0].name, "orders");
188 assert_eq!(response.topics[0].partitions[0].partition_index, 0);
189 assert_eq!(
190 response.topics[0].partitions[0].active_producers[0].producer_id,
191 42
192 );
193 assert_eq!(response.topics[0].partitions[1].error_code, 29);
194 assert_eq!(
195 response.topics[0].partitions[1].error_message.as_deref(),
196 Some("denied")
197 );
198 assert!(decoder.is_empty());
199 }
200}