Skip to main content

kafrust_protocol/api/
share_group_offsets.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5/// Kafka AlterShareGroupOffsets API key.
6pub const ALTER_API_KEY: i16 = 91;
7/// Kafka DeleteShareGroupOffsets API key.
8pub const DELETE_API_KEY: i16 = 92;
9
10/// One share-group partition offset to set.
11#[derive(Debug, Clone, Copy, PartialEq, Eq)]
12pub struct AlterShareGroupOffsetsPartitionV0 {
13    pub partition_index: i32,
14    pub start_offset: i64,
15}
16
17/// Offsets to set for one share-group topic.
18#[derive(Debug, Clone, PartialEq, Eq)]
19pub struct AlterShareGroupOffsetsTopicV0 {
20    pub topic_name: String,
21    pub partitions: Vec<AlterShareGroupOffsetsPartitionV0>,
22}
23
24/// AlterShareGroupOffsets v0 request.
25#[derive(Debug, Clone, PartialEq, Eq)]
26pub struct AlterShareGroupOffsetsRequestV0 {
27    pub correlation_id: i32,
28    pub client_id: Option<String>,
29    pub group_id: String,
30    pub topics: Vec<AlterShareGroupOffsetsTopicV0>,
31}
32
33impl AlterShareGroupOffsetsRequestV0 {
34    /// Encodes the flexible request, including its request header.
35    pub fn encode(&self) -> Result<Vec<u8>> {
36        let mut encoder = Encoder::new();
37        RequestHeader {
38            api_key: ALTER_API_KEY,
39            api_version: 0,
40            correlation_id: self.correlation_id,
41            client_id: self.client_id.clone(),
42        }
43        .encode_v2(&mut encoder)?;
44        encoder.write_compact_string(&self.group_id)?;
45        encoder.write_compact_array(Some(&self.topics), |encoder, topic| {
46            encoder.write_compact_string(&topic.topic_name)?;
47            encoder.write_compact_array(Some(&topic.partitions), |encoder, partition| {
48                encoder.write_i32(partition.partition_index);
49                encoder.write_i64(partition.start_offset);
50                encoder.write_empty_tagged_fields();
51                Ok(())
52            })?;
53            encoder.write_empty_tagged_fields();
54            Ok(())
55        })?;
56        encoder.write_empty_tagged_fields();
57        Ok(encoder.into_bytes())
58    }
59}
60
61/// One partition result returned by AlterShareGroupOffsets.
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub struct AlterShareGroupOffsetsPartitionResultV0 {
64    pub partition_index: i32,
65    pub error_code: i16,
66    pub error_message: Option<String>,
67}
68
69impl AlterShareGroupOffsetsPartitionResultV0 {
70    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
71        let partition_index = decoder.read_i32()?;
72        let error_code = decoder.read_i16()?;
73        let error_message = decoder.read_compact_nullable_string()?;
74        decoder.read_tagged_fields()?;
75        Ok(Self {
76            partition_index,
77            error_code,
78            error_message,
79        })
80    }
81}
82
83/// One topic result returned by AlterShareGroupOffsets.
84#[derive(Debug, Clone, PartialEq, Eq)]
85pub struct AlterShareGroupOffsetsTopicResultV0 {
86    pub topic_name: String,
87    pub topic_id: [u8; 16],
88    pub partitions: Vec<AlterShareGroupOffsetsPartitionResultV0>,
89}
90
91impl AlterShareGroupOffsetsTopicResultV0 {
92    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
93        let topic_name = decoder.read_compact_string()?;
94        let topic_id = decoder.read_uuid()?;
95        let partitions = decoder
96            .read_compact_array("alter share group offset partitions", |decoder| {
97                AlterShareGroupOffsetsPartitionResultV0::decode(decoder)
98            })?
99            .unwrap_or_default();
100        decoder.read_tagged_fields()?;
101        Ok(Self {
102            topic_name,
103            topic_id,
104            partitions,
105        })
106    }
107}
108
109/// AlterShareGroupOffsets v0 response.
110#[derive(Debug, Clone, PartialEq, Eq)]
111pub struct AlterShareGroupOffsetsResponseV0 {
112    pub throttle_time_ms: i32,
113    pub error_code: i16,
114    pub error_message: Option<String>,
115    pub responses: Vec<AlterShareGroupOffsetsTopicResultV0>,
116}
117
118impl AlterShareGroupOffsetsResponseV0 {
119    /// Decodes the flexible response body after the response header.
120    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
121        let throttle_time_ms = decoder.read_i32()?;
122        let error_code = decoder.read_i16()?;
123        let error_message = decoder.read_compact_nullable_string()?;
124        let responses = decoder
125            .read_compact_array("alter share group offset responses", |decoder| {
126                AlterShareGroupOffsetsTopicResultV0::decode(decoder)
127            })?
128            .unwrap_or_default();
129        decoder.read_tagged_fields()?;
130        Ok(Self {
131            throttle_time_ms,
132            error_code,
133            error_message,
134            responses,
135        })
136    }
137}
138
139/// One topic whose share-group offsets should be deleted.
140#[derive(Debug, Clone, PartialEq, Eq)]
141pub struct DeleteShareGroupOffsetsTopicV0 {
142    pub topic_name: String,
143}
144
145/// DeleteShareGroupOffsets v0 request.
146#[derive(Debug, Clone, PartialEq, Eq)]
147pub struct DeleteShareGroupOffsetsRequestV0 {
148    pub correlation_id: i32,
149    pub client_id: Option<String>,
150    pub group_id: String,
151    pub topics: Vec<DeleteShareGroupOffsetsTopicV0>,
152}
153
154impl DeleteShareGroupOffsetsRequestV0 {
155    /// Encodes the flexible request, including its request header.
156    pub fn encode(&self) -> Result<Vec<u8>> {
157        let mut encoder = Encoder::new();
158        RequestHeader {
159            api_key: DELETE_API_KEY,
160            api_version: 0,
161            correlation_id: self.correlation_id,
162            client_id: self.client_id.clone(),
163        }
164        .encode_v2(&mut encoder)?;
165        encoder.write_compact_string(&self.group_id)?;
166        encoder.write_compact_array(Some(&self.topics), |encoder, topic| {
167            encoder.write_compact_string(&topic.topic_name)?;
168            encoder.write_empty_tagged_fields();
169            Ok(())
170        })?;
171        encoder.write_empty_tagged_fields();
172        Ok(encoder.into_bytes())
173    }
174}
175
176/// One topic result returned by DeleteShareGroupOffsets.
177#[derive(Debug, Clone, PartialEq, Eq)]
178pub struct DeleteShareGroupOffsetsTopicResultV0 {
179    pub topic_name: String,
180    pub topic_id: [u8; 16],
181    pub error_code: i16,
182    pub error_message: Option<String>,
183}
184
185impl DeleteShareGroupOffsetsTopicResultV0 {
186    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
187        let topic_name = decoder.read_compact_string()?;
188        let topic_id = decoder.read_uuid()?;
189        let error_code = decoder.read_i16()?;
190        let error_message = decoder.read_compact_nullable_string()?;
191        decoder.read_tagged_fields()?;
192        Ok(Self {
193            topic_name,
194            topic_id,
195            error_code,
196            error_message,
197        })
198    }
199}
200
201/// DeleteShareGroupOffsets v0 response.
202#[derive(Debug, Clone, PartialEq, Eq)]
203pub struct DeleteShareGroupOffsetsResponseV0 {
204    pub throttle_time_ms: i32,
205    pub error_code: i16,
206    pub error_message: Option<String>,
207    pub responses: Vec<DeleteShareGroupOffsetsTopicResultV0>,
208}
209
210impl DeleteShareGroupOffsetsResponseV0 {
211    /// Decodes the flexible response body after the response header.
212    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
213        let throttle_time_ms = decoder.read_i32()?;
214        let error_code = decoder.read_i16()?;
215        let error_message = decoder.read_compact_nullable_string()?;
216        let responses = decoder
217            .read_compact_array("delete share group offset responses", |decoder| {
218                DeleteShareGroupOffsetsTopicResultV0::decode(decoder)
219            })?
220            .unwrap_or_default();
221        decoder.read_tagged_fields()?;
222        Ok(Self {
223            throttle_time_ms,
224            error_code,
225            error_message,
226            responses,
227        })
228    }
229}
230
231#[cfg(test)]
232#[allow(clippy::unwrap_used)]
233mod tests {
234    use super::{
235        AlterShareGroupOffsetsPartitionV0, AlterShareGroupOffsetsRequestV0,
236        AlterShareGroupOffsetsResponseV0, AlterShareGroupOffsetsTopicV0,
237        DeleteShareGroupOffsetsRequestV0, DeleteShareGroupOffsetsTopicV0, ALTER_API_KEY,
238        DELETE_API_KEY,
239    };
240    use crate::codec::{Decoder, Encoder};
241
242    #[test]
243    fn encodes_alter_share_group_offsets_v0_request() {
244        let request = AlterShareGroupOffsetsRequestV0 {
245            correlation_id: 23,
246            client_id: Some("kafrust".to_owned()),
247            group_id: "share-orders".to_owned(),
248            topics: vec![AlterShareGroupOffsetsTopicV0 {
249                topic_name: "orders".to_owned(),
250                partitions: vec![AlterShareGroupOffsetsPartitionV0 {
251                    partition_index: 2,
252                    start_offset: 42,
253                }],
254            }],
255        };
256
257        let encoded = request.encode().unwrap();
258        assert_eq!(&encoded[..4], &[0, 91, 0, 0]);
259        assert_eq!(ALTER_API_KEY, 91);
260        assert_eq!(encoded.last(), Some(&0));
261    }
262
263    #[test]
264    fn encodes_delete_share_group_offsets_v0_request() {
265        let request = DeleteShareGroupOffsetsRequestV0 {
266            correlation_id: 24,
267            client_id: Some("kafrust".to_owned()),
268            group_id: "share-orders".to_owned(),
269            topics: vec![DeleteShareGroupOffsetsTopicV0 {
270                topic_name: "orders".to_owned(),
271            }],
272        };
273
274        let encoded = request.encode().unwrap();
275        assert_eq!(&encoded[..4], &[0, 92, 0, 0]);
276        assert_eq!(DELETE_API_KEY, 92);
277        assert_eq!(encoded.last(), Some(&0));
278    }
279
280    #[test]
281    fn decodes_alter_share_group_offsets_v0_response() -> crate::error::Result<()> {
282        let mut bytes = Encoder::new();
283        bytes.write_i32(8);
284        bytes.write_i16(0);
285        bytes.write_compact_nullable_string(Some("ok"))?;
286        bytes.write_compact_array(Some(&[1_i8]), |encoder, _| {
287            encoder.write_compact_string("orders")?;
288            encoder.write_uuid(&[7; 16]);
289            encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
290                encoder.write_i32(2);
291                encoder.write_i16(0);
292                encoder.write_compact_nullable_string(None)?;
293                encoder.write_empty_tagged_fields();
294                Ok(())
295            })?;
296            encoder.write_empty_tagged_fields();
297            Ok(())
298        })?;
299        bytes.write_empty_tagged_fields();
300
301        let encoded = bytes.into_bytes();
302        let mut decoder = Decoder::new(&encoded);
303        let response = AlterShareGroupOffsetsResponseV0::decode_body(&mut decoder)?;
304        assert_eq!(response.throttle_time_ms, 8);
305        assert_eq!(response.responses[0].topic_name, "orders");
306        assert_eq!(response.responses[0].partitions[0].partition_index, 2);
307        assert!(decoder.is_empty());
308        Ok(())
309    }
310}