Skip to main content

kafrust_protocol/api/
add_offsets_to_txn.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 25;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct AddOffsetsToTxnRequestV0 {
9    pub correlation_id: i32,
10    pub client_id: Option<String>,
11    pub transactional_id: String,
12    pub producer_id: i64,
13    pub producer_epoch: i16,
14    pub group_id: String,
15}
16
17impl AddOffsetsToTxnRequestV0 {
18    pub fn encode(&self) -> Result<Vec<u8>> {
19        let mut encoder = Encoder::new();
20        RequestHeader {
21            api_key: API_KEY,
22            api_version: 0,
23            correlation_id: self.correlation_id,
24            client_id: self.client_id.clone(),
25        }
26        .encode_v1(&mut encoder)?;
27        encoder.write_string(&self.transactional_id)?;
28        encoder.write_i64(self.producer_id);
29        encoder.write_i16(self.producer_epoch);
30        encoder.write_string(&self.group_id)?;
31        Ok(encoder.into_bytes())
32    }
33}
34
35#[derive(Debug, Clone, PartialEq, Eq)]
36pub struct AddOffsetsToTxnResponseV0 {
37    pub throttle_time_ms: i32,
38    pub error_code: i16,
39}
40
41impl AddOffsetsToTxnResponseV0 {
42    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
43        Ok(Self {
44            throttle_time_ms: decoder.read_i32()?,
45            error_code: decoder.read_i16()?,
46        })
47    }
48}
49
50/// AddOffsetsToTxn v3, the flexible form used by current Kafka brokers.
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct AddOffsetsToTxnRequestV3 {
53    pub correlation_id: i32,
54    pub client_id: Option<String>,
55    pub transactional_id: String,
56    pub producer_id: i64,
57    pub producer_epoch: i16,
58    pub group_id: String,
59}
60
61impl AddOffsetsToTxnRequestV3 {
62    pub fn encode(&self) -> Result<Vec<u8>> {
63        let mut encoder = Encoder::new();
64        RequestHeader {
65            api_key: API_KEY,
66            api_version: 3,
67            correlation_id: self.correlation_id,
68            client_id: self.client_id.clone(),
69        }
70        .encode_v2(&mut encoder)?;
71        encoder.write_compact_string(&self.transactional_id)?;
72        encoder.write_i64(self.producer_id);
73        encoder.write_i16(self.producer_epoch);
74        encoder.write_compact_string(&self.group_id)?;
75        encoder.write_empty_tagged_fields();
76        Ok(encoder.into_bytes())
77    }
78}
79
80/// AddOffsetsToTxn v3 response with flexible tagged fields.
81#[derive(Debug, Clone, PartialEq, Eq)]
82pub struct AddOffsetsToTxnResponseV3 {
83    pub throttle_time_ms: i32,
84    pub error_code: i16,
85}
86
87impl AddOffsetsToTxnResponseV3 {
88    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
89        let response = Self {
90            throttle_time_ms: decoder.read_i32()?,
91            error_code: decoder.read_i16()?,
92        };
93        decoder.read_tagged_fields()?;
94        Ok(response)
95    }
96}
97
98#[cfg(test)]
99#[allow(clippy::unwrap_used)]
100mod tests {
101    use super::{
102        AddOffsetsToTxnRequestV0, AddOffsetsToTxnRequestV3, AddOffsetsToTxnResponseV0,
103        AddOffsetsToTxnResponseV3, API_KEY,
104    };
105    use crate::codec::{Decoder, Encoder};
106
107    #[test]
108    fn encodes_add_offsets_to_txn_v0_request() {
109        let request = AddOffsetsToTxnRequestV0 {
110            correlation_id: 51,
111            client_id: Some("kafrust".to_owned()),
112            transactional_id: "orders-tx".to_owned(),
113            producer_id: 42,
114            producer_epoch: 3,
115            group_id: "orders-group".to_owned(),
116        };
117        let encoded = request.encode().unwrap();
118
119        assert_eq!(&encoded[0..8], &[0, 25, 0, 0, 0, 0, 0, 51]);
120        assert!(encoded.ends_with(b"orders-group"));
121        assert_eq!(API_KEY, 25);
122    }
123
124    #[test]
125    fn decodes_add_offsets_to_txn_v0_response() {
126        let bytes = [
127            0, 0, 0, 7, // throttle time
128            0, 16, // not coordinator
129        ];
130        let mut decoder = Decoder::new(&bytes);
131        let response = AddOffsetsToTxnResponseV0::decode_body(&mut decoder).unwrap();
132
133        assert_eq!(response.throttle_time_ms, 7);
134        assert_eq!(response.error_code, 16);
135        assert!(decoder.is_empty());
136    }
137
138    #[test]
139    fn encodes_add_offsets_to_txn_v3_request_with_flexible_fields() {
140        let request = AddOffsetsToTxnRequestV3 {
141            correlation_id: 52,
142            client_id: Some("kafrust".to_owned()),
143            transactional_id: "orders-tx".to_owned(),
144            producer_id: 42,
145            producer_epoch: 3,
146            group_id: "orders-group".to_owned(),
147        };
148        let encoded = request.encode().unwrap();
149
150        assert_eq!(&encoded[0..8], &[0, 25, 0, 3, 0, 0, 0, 52]);
151        assert!(encoded
152            .windows(b"orders-tx".len())
153            .any(|window| window == b"orders-tx"));
154        assert_eq!(encoded.last(), Some(&0));
155    }
156
157    #[test]
158    fn decodes_add_offsets_to_txn_v3_response_with_tagged_fields() {
159        let mut bytes = Encoder::new();
160        bytes.write_i32(9);
161        bytes.write_i16(47);
162        bytes.write_empty_tagged_fields();
163        let encoded = bytes.into_bytes();
164        let mut decoder = Decoder::new(&encoded);
165
166        let response = AddOffsetsToTxnResponseV3::decode_body(&mut decoder).unwrap();
167
168        assert_eq!(response.throttle_time_ms, 9);
169        assert_eq!(response.error_code, 47);
170        assert!(decoder.is_empty());
171    }
172}