kafrust_protocol/api/
describe_transactions.rs1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 65;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct DescribeTransactionsRequestV0 {
9 pub correlation_id: i32,
10 pub client_id: Option<String>,
11 pub transactional_ids: Vec<String>,
12}
13
14impl DescribeTransactionsRequestV0 {
15 pub fn encode(&self) -> Result<Vec<u8>> {
16 let mut encoder = Encoder::new();
17 RequestHeader {
18 api_key: API_KEY,
19 api_version: 0,
20 correlation_id: self.correlation_id,
21 client_id: self.client_id.clone(),
22 }
23 .encode_v2(&mut encoder)?;
24 encoder.write_compact_array(Some(&self.transactional_ids), |encoder, id| {
25 encoder.write_compact_string(id)
26 })?;
27 encoder.write_empty_tagged_fields();
28 Ok(encoder.into_bytes())
29 }
30}
31
32#[derive(Debug, Clone, PartialEq, Eq)]
33pub struct DescribeTransactionsResponseV0 {
34 pub throttle_time_ms: i32,
35 pub transaction_states: Vec<DescribeTransactionsStateV0>,
36}
37
38impl DescribeTransactionsResponseV0 {
39 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
40 let throttle_time_ms = decoder.read_i32()?;
41 let transaction_states = decoder
42 .read_compact_array("describe transactions states", |decoder| {
43 let error_code = decoder.read_i16()?;
44 let transactional_id = decoder.read_compact_string()?;
45 let transaction_state = decoder.read_compact_string()?;
46 let transaction_timeout_ms = decoder.read_i32()?;
47 let transaction_start_time_ms = decoder.read_i64()?;
48 let producer_id = decoder.read_i64()?;
49 let producer_epoch = decoder.read_i16()?;
50 let topics = decoder
51 .read_compact_array("describe transaction topics", |decoder| {
52 let topic = decoder.read_compact_string()?;
53 let partitions = decoder
54 .read_array("describe transaction partitions", |decoder| {
55 decoder.read_i32()
56 })?
57 .unwrap_or_default();
58 decoder.read_tagged_fields()?;
59 Ok(DescribeTransactionsTopicV0 { topic, partitions })
60 })?
61 .unwrap_or_default();
62 decoder.read_tagged_fields()?;
63 Ok(DescribeTransactionsStateV0 {
64 error_code,
65 transactional_id,
66 transaction_state,
67 transaction_timeout_ms,
68 transaction_start_time_ms,
69 producer_id,
70 producer_epoch,
71 topics,
72 })
73 })?
74 .unwrap_or_default();
75 decoder.read_tagged_fields()?;
76 Ok(Self {
77 throttle_time_ms,
78 transaction_states,
79 })
80 }
81}
82
83#[derive(Debug, Clone, PartialEq, Eq)]
84pub struct DescribeTransactionsStateV0 {
85 pub error_code: i16,
86 pub transactional_id: String,
87 pub transaction_state: String,
88 pub transaction_timeout_ms: i32,
89 pub transaction_start_time_ms: i64,
90 pub producer_id: i64,
91 pub producer_epoch: i16,
92 pub topics: Vec<DescribeTransactionsTopicV0>,
93}
94
95#[derive(Debug, Clone, PartialEq, Eq)]
96pub struct DescribeTransactionsTopicV0 {
97 pub topic: String,
98 pub partitions: Vec<i32>,
99}
100
101#[cfg(test)]
102#[allow(clippy::unwrap_used)]
103mod tests {
104 use super::{DescribeTransactionsRequestV0, DescribeTransactionsResponseV0, API_KEY};
105 use crate::codec::{Decoder, Encoder};
106
107 #[test]
108 fn encodes_describe_transactions_v0_request() {
109 let request = DescribeTransactionsRequestV0 {
110 correlation_id: 65,
111 client_id: None,
112 transactional_ids: vec!["payments-tx".to_owned()],
113 };
114
115 let bytes = request.encode().unwrap();
116 assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 0]);
117 assert_eq!(&bytes[4..8], &[0, 0, 0, 65]);
118 assert_eq!(bytes.last(), Some(&0));
119 }
120
121 #[test]
122 fn decodes_describe_transactions_v0_response() {
123 let mut bytes = Encoder::new();
124 bytes.write_i32(7);
125 bytes.write_unsigned_varint(3);
126 bytes.write_i16(0);
127 bytes.write_compact_string("payments-tx").unwrap();
128 bytes.write_compact_string("Ongoing").unwrap();
129 bytes.write_i32(60_000);
130 bytes.write_i64(1_700_000_000_000);
131 bytes.write_i64(99);
132 bytes.write_i16(4);
133 bytes.write_unsigned_varint(2);
134 bytes.write_compact_string("orders").unwrap();
135 bytes
136 .write_array(Some(&[0, 2]), |encoder, partition| {
137 encoder.write_i32(*partition);
138 Ok(())
139 })
140 .unwrap();
141 bytes.write_empty_tagged_fields();
142 bytes.write_empty_tagged_fields();
143 bytes.write_i16(15);
144 bytes.write_compact_string("missing-tx").unwrap();
145 bytes.write_compact_string("Empty").unwrap();
146 bytes.write_i32(60_000);
147 bytes.write_i64(-1);
148 bytes.write_i64(-1);
149 bytes.write_i16(-1);
150 bytes.write_unsigned_varint(1);
151 bytes.write_empty_tagged_fields();
152 bytes.write_empty_tagged_fields();
153 let bytes = bytes.into_bytes();
154 let mut decoder = Decoder::new(&bytes);
155
156 let response = DescribeTransactionsResponseV0::decode_body(&mut decoder).unwrap();
157
158 assert_eq!(response.throttle_time_ms, 7);
159 assert_eq!(
160 response.transaction_states[0].transactional_id,
161 "payments-tx"
162 );
163 assert_eq!(response.transaction_states[0].producer_id, 99);
164 assert_eq!(response.transaction_states[0].topics[0].partitions, [0, 2]);
165 assert_eq!(response.transaction_states[1].error_code, 15);
166 assert!(decoder.is_empty());
167 }
168}