Skip to main content

kafrust_protocol/api/
end_txn.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 26;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct EndTxnRequestV0 {
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 committed: bool,
15}
16
17impl EndTxnRequestV0 {
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_bool(self.committed);
31        Ok(encoder.into_bytes())
32    }
33}
34
35#[derive(Debug, Clone, PartialEq, Eq)]
36pub struct EndTxnResponseV0 {
37    pub throttle_time_ms: i32,
38    pub error_code: i16,
39}
40
41impl EndTxnResponseV0 {
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/// EndTxn v3, the flexible form used by current Kafka brokers.
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct EndTxnRequestV3 {
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 committed: bool,
59}
60
61impl EndTxnRequestV3 {
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_bool(self.committed);
75        encoder.write_empty_tagged_fields();
76        Ok(encoder.into_bytes())
77    }
78}
79
80/// EndTxn v3 response with flexible tagged fields.
81#[derive(Debug, Clone, PartialEq, Eq)]
82pub struct EndTxnResponseV3 {
83    pub throttle_time_ms: i32,
84    pub error_code: i16,
85}
86
87impl EndTxnResponseV3 {
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::{EndTxnRequestV0, EndTxnRequestV3, EndTxnResponseV0, EndTxnResponseV3, API_KEY};
102    use crate::codec::Decoder;
103
104    #[test]
105    fn encodes_end_txn_v0_commit_request() {
106        let request = EndTxnRequestV0 {
107            correlation_id: 31,
108            client_id: Some("kafrust".to_owned()),
109            transactional_id: "orders-tx".to_owned(),
110            producer_id: 42,
111            producer_epoch: 3,
112            committed: true,
113        };
114
115        assert_eq!(
116            request.encode().unwrap(),
117            [
118                0, 26, // api key
119                0, 0, // api version
120                0, 0, 0, 31, // correlation id
121                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', // client id
122                0, 9, b'o', b'r', b'd', b'e', b'r', b's', b'-', b't', b'x', 0, 0, 0, 0, 0, 0, 0,
123                42, // producer id
124                0, 3, // producer epoch
125                1, // committed
126            ]
127        );
128        assert_eq!(API_KEY, 26);
129    }
130
131    #[test]
132    fn encodes_end_txn_v0_abort_request() {
133        let request = EndTxnRequestV0 {
134            correlation_id: 32,
135            client_id: None,
136            transactional_id: "orders-tx".to_owned(),
137            producer_id: 42,
138            producer_epoch: 3,
139            committed: false,
140        };
141
142        assert_eq!(request.encode().unwrap().last(), Some(&0));
143    }
144
145    #[test]
146    fn decodes_end_txn_v0_response() {
147        let bytes = [
148            0, 0, 0, 12, // throttle time
149            0, 47, // invalid producer epoch
150        ];
151        let mut decoder = Decoder::new(&bytes);
152        let response = EndTxnResponseV0::decode_body(&mut decoder).unwrap();
153
154        assert_eq!(response.throttle_time_ms, 12);
155        assert_eq!(response.error_code, 47);
156        assert!(decoder.is_empty());
157    }
158
159    #[test]
160    fn encodes_end_txn_v3_commit_request_with_flexible_fields() {
161        let request = EndTxnRequestV3 {
162            correlation_id: 33,
163            client_id: Some("kafrust".to_owned()),
164            transactional_id: "orders-tx".to_owned(),
165            producer_id: 42,
166            producer_epoch: 3,
167            committed: true,
168        };
169        let encoded = request.encode().unwrap();
170
171        assert_eq!(&encoded[0..8], &[0, 26, 0, 3, 0, 0, 0, 33]);
172        assert!(encoded
173            .windows(b"orders-tx".len())
174            .any(|window| window == b"orders-tx"));
175        assert_eq!(encoded.last(), Some(&0));
176    }
177
178    #[test]
179    fn decodes_end_txn_v3_response_with_tagged_fields() {
180        let bytes = [0, 0, 0, 12, 0, 47, 0];
181        let mut decoder = Decoder::new(&bytes);
182        let response = EndTxnResponseV3::decode_body(&mut decoder).unwrap();
183
184        assert_eq!(response.throttle_time_ms, 12);
185        assert_eq!(response.error_code, 47);
186        assert!(decoder.is_empty());
187    }
188}