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#[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#[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, 0, 0, 0, 0, 0, 31, 0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', 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, 0, 3, 1, ]
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, 0, 47, ];
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}