Skip to main content

kafrust_protocol/api/
add_offsets_to_txn.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 25;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct AddOffsetsToTxnRequestV0 {
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 group_id: String,
15}
16
17impl AddOffsetsToTxnRequestV0 {
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_string(&self.group_id)?;
31        Ok(encoder.into_bytes())
32    }
33}
34
35#[derive(Debug, Clone, PartialEq, Eq)]
36pub struct AddOffsetsToTxnResponseV0 {
37    pub throttle_time_ms: i32,
38    pub error_code: i16,
39}
40
41impl AddOffsetsToTxnResponseV0 {
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::{AddOffsetsToTxnRequestV0, AddOffsetsToTxnResponseV0, API_KEY};
54    use crate::codec::Decoder;
55
56    #[test]
57    fn encodes_add_offsets_to_txn_v0_request() {
58        let request = AddOffsetsToTxnRequestV0 {
59            correlation_id: 51,
60            client_id: Some("kafrust".to_owned()),
61            transactional_id: "orders-tx".to_owned(),
62            producer_id: 42,
63            producer_epoch: 3,
64            group_id: "orders-group".to_owned(),
65        };
66        let encoded = request.encode().unwrap();
67
68        assert_eq!(&encoded[0..8], &[0, 25, 0, 0, 0, 0, 0, 51]);
69        assert!(encoded.ends_with(b"orders-group"));
70        assert_eq!(API_KEY, 25);
71    }
72
73    #[test]
74    fn decodes_add_offsets_to_txn_v0_response() {
75        let bytes = [
76            0, 0, 0, 7, // throttle time
77            0, 16, // not coordinator
78        ];
79        let mut decoder = Decoder::new(&bytes);
80        let response = AddOffsetsToTxnResponseV0::decode_body(&mut decoder).unwrap();
81
82        assert_eq!(response.throttle_time_ms, 7);
83        assert_eq!(response.error_code, 16);
84        assert!(decoder.is_empty());
85    }
86}