Skip to main content

kafrust_protocol/api/
heartbeat.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 12;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct HeartbeatRequestV2 {
9    pub correlation_id: i32,
10    pub client_id: Option<String>,
11    pub group_id: String,
12    pub generation_id: i32,
13    pub member_id: String,
14}
15
16impl HeartbeatRequestV2 {
17    pub fn encode(&self) -> Result<Vec<u8>> {
18        let mut encoder = Encoder::new();
19        RequestHeader {
20            api_key: API_KEY,
21            api_version: 2,
22            correlation_id: self.correlation_id,
23            client_id: self.client_id.clone(),
24        }
25        .encode_v1(&mut encoder)?;
26        encoder.write_string(&self.group_id)?;
27        encoder.write_i32(self.generation_id);
28        encoder.write_string(&self.member_id)?;
29        Ok(encoder.into_bytes())
30    }
31}
32
33#[derive(Debug, Clone, PartialEq, Eq)]
34pub struct HeartbeatRequestV3 {
35    pub correlation_id: i32,
36    pub client_id: Option<String>,
37    pub group_id: String,
38    pub generation_id: i32,
39    pub member_id: String,
40    pub group_instance_id: Option<String>,
41}
42
43impl HeartbeatRequestV3 {
44    pub fn encode(&self) -> Result<Vec<u8>> {
45        let mut encoder = Encoder::new();
46        RequestHeader {
47            api_key: API_KEY,
48            api_version: 3,
49            correlation_id: self.correlation_id,
50            client_id: self.client_id.clone(),
51        }
52        .encode_v1(&mut encoder)?;
53        encoder.write_string(&self.group_id)?;
54        encoder.write_i32(self.generation_id);
55        encoder.write_string(&self.member_id)?;
56        encoder.write_nullable_string(self.group_instance_id.as_deref())?;
57        Ok(encoder.into_bytes())
58    }
59}
60
61#[derive(Debug, Clone, PartialEq, Eq)]
62pub struct HeartbeatResponseV2 {
63    pub throttle_time_ms: i32,
64    pub error_code: i16,
65}
66
67impl HeartbeatResponseV2 {
68    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
69        Ok(Self {
70            throttle_time_ms: decoder.read_i32()?,
71            error_code: decoder.read_i16()?,
72        })
73    }
74}
75
76#[cfg(test)]
77#[allow(clippy::unwrap_used)]
78mod tests {
79    use super::{HeartbeatRequestV2, HeartbeatRequestV3, HeartbeatResponseV2};
80    use crate::codec::Decoder;
81
82    #[test]
83    fn encodes_heartbeat_v2_request() {
84        let request = HeartbeatRequestV2 {
85            correlation_id: 17,
86            client_id: Some("kafrust".to_owned()),
87            group_id: "orders-group".to_owned(),
88            generation_id: 7,
89            member_id: "member-a".to_owned(),
90        };
91
92        assert_eq!(
93            request.encode().unwrap(),
94            [
95                0, 12, // api key
96                0, 2, // api version
97                0, 0, 0, 17, // correlation id
98                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', // client id
99                0, 12, b'o', b'r', b'd', b'e', b'r', b's', b'-', b'g', b'r', b'o', b'u',
100                b'p', // group id
101                0, 0, 0, 7, // generation id
102                0, 8, b'm', b'e', b'm', b'b', b'e', b'r', b'-', b'a', // member id
103            ]
104        );
105    }
106
107    #[test]
108    fn decodes_heartbeat_v2_response() {
109        let bytes = [
110            0, 0, 0, 0, // throttle time
111            0, 0, // error code
112        ];
113        let mut decoder = Decoder::new(&bytes);
114        let response = HeartbeatResponseV2::decode_body(&mut decoder).unwrap();
115
116        assert_eq!(response.throttle_time_ms, 0);
117        assert_eq!(response.error_code, 0);
118        assert!(decoder.is_empty());
119    }
120
121    #[test]
122    fn encodes_heartbeat_v3_request_with_static_member() {
123        let request = HeartbeatRequestV3 {
124            correlation_id: 17,
125            client_id: Some("kafrust".to_owned()),
126            group_id: "orders-group".to_owned(),
127            generation_id: 7,
128            member_id: "member-a".to_owned(),
129            group_instance_id: Some("orders-reader-1".to_owned()),
130        };
131
132        let encoded = request.encode().unwrap();
133        assert_eq!(&encoded[0..4], &[0, 12, 0, 3]);
134        assert!(encoded
135            .windows(17)
136            .any(|bytes| bytes == b"\0\x0forders-reader-1"));
137    }
138}