Skip to main content

kafrust_protocol/api/
list_transactions.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 66;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct ListTransactionsRequestV0 {
9    pub correlation_id: i32,
10    pub client_id: Option<String>,
11    pub state_filters: Vec<String>,
12    pub producer_id_filters: Vec<i64>,
13}
14
15impl ListTransactionsRequestV0 {
16    pub fn encode(&self) -> Result<Vec<u8>> {
17        let mut encoder = Encoder::new();
18        RequestHeader {
19            api_key: API_KEY,
20            api_version: 0,
21            correlation_id: self.correlation_id,
22            client_id: self.client_id.clone(),
23        }
24        .encode_v2(&mut encoder)?;
25        encode_filters(&mut encoder, &self.state_filters, &self.producer_id_filters)?;
26        encoder.write_empty_tagged_fields();
27        Ok(encoder.into_bytes())
28    }
29}
30
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct ListTransactionsRequestV1 {
33    pub correlation_id: i32,
34    pub client_id: Option<String>,
35    pub state_filters: Vec<String>,
36    pub producer_id_filters: Vec<i64>,
37    pub duration_filter_ms: i64,
38}
39
40impl ListTransactionsRequestV1 {
41    pub fn encode(&self) -> Result<Vec<u8>> {
42        let mut encoder = Encoder::new();
43        RequestHeader {
44            api_key: API_KEY,
45            api_version: 1,
46            correlation_id: self.correlation_id,
47            client_id: self.client_id.clone(),
48        }
49        .encode_v2(&mut encoder)?;
50        encode_filters(&mut encoder, &self.state_filters, &self.producer_id_filters)?;
51        encoder.write_i64(self.duration_filter_ms);
52        encoder.write_empty_tagged_fields();
53        Ok(encoder.into_bytes())
54    }
55}
56
57fn encode_filters(encoder: &mut Encoder, states: &[String], producer_ids: &[i64]) -> Result<()> {
58    encoder.write_compact_array(Some(states), |encoder, state| {
59        encoder.write_compact_string(state)
60    })?;
61    encoder.write_compact_array(Some(producer_ids), |encoder, producer_id| {
62        encoder.write_i64(*producer_id);
63        Ok(())
64    })?;
65    Ok(())
66}
67
68#[derive(Debug, Clone, PartialEq, Eq)]
69pub struct ListTransactionsResponseV0 {
70    pub throttle_time_ms: i32,
71    pub error_code: i16,
72    pub unknown_state_filters: Vec<String>,
73    pub transaction_states: Vec<ListedTransactionV0>,
74}
75
76impl ListTransactionsResponseV0 {
77    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
78        let throttle_time_ms = decoder.read_i32()?;
79        let error_code = decoder.read_i16()?;
80        let unknown_state_filters = decoder
81            .read_compact_array("list transactions unknown state filters", |decoder| {
82                decoder.read_compact_string()
83            })?
84            .unwrap_or_default();
85        let transaction_states = decoder
86            .read_compact_array("listed transactions", ListedTransactionV0::decode)?
87            .unwrap_or_default();
88        decoder.read_tagged_fields()?;
89        Ok(Self {
90            throttle_time_ms,
91            error_code,
92            unknown_state_filters,
93            transaction_states,
94        })
95    }
96}
97
98pub type ListTransactionsResponseV1 = ListTransactionsResponseV0;
99
100#[derive(Debug, Clone, PartialEq, Eq)]
101pub struct ListedTransactionV0 {
102    pub transactional_id: String,
103    pub producer_id: i64,
104    pub transaction_state: String,
105}
106
107impl ListedTransactionV0 {
108    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
109        let result = Self {
110            transactional_id: decoder.read_compact_string()?,
111            producer_id: decoder.read_i64()?,
112            transaction_state: decoder.read_compact_string()?,
113        };
114        decoder.read_tagged_fields()?;
115        Ok(result)
116    }
117}
118
119#[cfg(test)]
120#[allow(clippy::unwrap_used)]
121mod tests {
122    use super::{
123        ListTransactionsRequestV0, ListTransactionsRequestV1, ListTransactionsResponseV0, API_KEY,
124    };
125    use crate::codec::{Decoder, Encoder};
126
127    #[test]
128    fn encodes_list_transactions_v0_request() {
129        let request = ListTransactionsRequestV0 {
130            correlation_id: 66,
131            client_id: Some("kafrust".to_owned()),
132            state_filters: vec!["Ongoing".to_owned()],
133            producer_id_filters: vec![42],
134        };
135
136        let bytes = request.encode().unwrap();
137        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 0]);
138        assert_eq!(&bytes[4..8], &[0, 0, 0, 66]);
139        assert!(bytes.ends_with(&[0]));
140    }
141
142    #[test]
143    fn encodes_list_transactions_v1_duration_filter() {
144        let request = ListTransactionsRequestV1 {
145            correlation_id: 67,
146            client_id: None,
147            state_filters: Vec::new(),
148            producer_id_filters: Vec::new(),
149            duration_filter_ms: 30_000,
150        };
151
152        let bytes = request.encode().unwrap();
153        assert_eq!(&bytes[0..4], &[0, API_KEY as u8, 0, 1]);
154        assert_eq!(&bytes[4..8], &[0, 0, 0, 67]);
155        assert_eq!(&bytes[8..10], &[255, 255]);
156        assert_eq!(bytes[10], 0); // request-header tagged fields
157        assert_eq!(bytes[11], 1); // empty state filter array
158        assert_eq!(bytes[12], 1); // empty producer ID filter array
159        assert_eq!(&bytes[13..21], &30_000_i64.to_be_bytes());
160        assert!(bytes.ends_with(&[0]));
161    }
162
163    #[test]
164    fn decodes_list_transactions_response() {
165        let mut bytes = Encoder::new();
166        bytes.write_i32(8);
167        bytes.write_i16(0);
168        bytes.write_unsigned_varint(2);
169        bytes.write_compact_string("UnknownState").unwrap();
170        bytes.write_unsigned_varint(2);
171        bytes.write_compact_string("payments-tx").unwrap();
172        bytes.write_i64(42);
173        bytes.write_compact_string("Ongoing").unwrap();
174        bytes.write_empty_tagged_fields();
175        bytes.write_empty_tagged_fields();
176        let bytes = bytes.into_bytes();
177        let mut decoder = Decoder::new(&bytes);
178
179        let response = ListTransactionsResponseV0::decode_body(&mut decoder).unwrap();
180
181        assert_eq!(response.throttle_time_ms, 8);
182        assert_eq!(response.unknown_state_filters, ["UnknownState"]);
183        assert_eq!(
184            response.transaction_states[0].transactional_id,
185            "payments-tx"
186        );
187        assert_eq!(response.transaction_states[0].producer_id, 42);
188        assert_eq!(response.transaction_states[0].transaction_state, "Ongoing");
189        assert!(decoder.is_empty());
190    }
191}