kafrust_protocol/api/
end_txn.rs1use 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#[cfg(test)]
51#[allow(clippy::unwrap_used)]
52mod tests {
53 use super::{EndTxnRequestV0, EndTxnResponseV0, API_KEY};
54 use crate::codec::Decoder;
55
56 #[test]
57 fn encodes_end_txn_v0_commit_request() {
58 let request = EndTxnRequestV0 {
59 correlation_id: 31,
60 client_id: Some("kafrust".to_owned()),
61 transactional_id: "orders-tx".to_owned(),
62 producer_id: 42,
63 producer_epoch: 3,
64 committed: true,
65 };
66
67 assert_eq!(
68 request.encode().unwrap(),
69 [
70 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,
75 42, 0, 3, 1, ]
79 );
80 assert_eq!(API_KEY, 26);
81 }
82
83 #[test]
84 fn encodes_end_txn_v0_abort_request() {
85 let request = EndTxnRequestV0 {
86 correlation_id: 32,
87 client_id: None,
88 transactional_id: "orders-tx".to_owned(),
89 producer_id: 42,
90 producer_epoch: 3,
91 committed: false,
92 };
93
94 assert_eq!(request.encode().unwrap().last(), Some(&0));
95 }
96
97 #[test]
98 fn decodes_end_txn_v0_response() {
99 let bytes = [
100 0, 0, 0, 12, 0, 47, ];
103 let mut decoder = Decoder::new(&bytes);
104 let response = EndTxnResponseV0::decode_body(&mut decoder).unwrap();
105
106 assert_eq!(response.throttle_time_ms, 12);
107 assert_eq!(response.error_code, 47);
108 assert!(decoder.is_empty());
109 }
110}