Skip to main content

kafrust_protocol/api/
describe_producers.rs

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]); // fixed nullable client ID length
150        assert_eq!(bytes[18], 2); // one topic
151        assert_eq!(bytes[26], 3); // two partition indexes
152        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}