Skip to main content

kafrust_protocol/api/
join_group.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 11;
6
7#[derive(Debug, Clone, PartialEq, Eq)]
8pub struct JoinGroupRequestV2 {
9    pub correlation_id: i32,
10    pub client_id: Option<String>,
11    pub group_id: String,
12    pub session_timeout_ms: i32,
13    pub rebalance_timeout_ms: i32,
14    pub member_id: String,
15    pub protocol_type: String,
16    pub protocols: Vec<JoinGroupProtocol>,
17}
18
19impl JoinGroupRequestV2 {
20    pub fn encode(&self) -> Result<Vec<u8>> {
21        let mut encoder = Encoder::new();
22        RequestHeader {
23            api_key: API_KEY,
24            api_version: 2,
25            correlation_id: self.correlation_id,
26            client_id: self.client_id.clone(),
27        }
28        .encode_v1(&mut encoder)?;
29        encoder.write_string(&self.group_id)?;
30        encoder.write_i32(self.session_timeout_ms);
31        encoder.write_i32(self.rebalance_timeout_ms);
32        encoder.write_string(&self.member_id)?;
33        encoder.write_string(&self.protocol_type)?;
34        encoder.write_array(Some(self.protocols.as_slice()), |encoder, protocol| {
35            protocol.encode(encoder)
36        })?;
37        Ok(encoder.into_bytes())
38    }
39}
40
41#[derive(Debug, Clone, PartialEq, Eq)]
42pub struct JoinGroupProtocol {
43    pub name: String,
44    pub metadata: Vec<u8>,
45}
46
47impl JoinGroupProtocol {
48    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
49        encoder.write_string(&self.name)?;
50        encoder.write_bytes(&self.metadata)
51    }
52}
53
54#[derive(Debug, Clone, PartialEq, Eq)]
55pub struct JoinGroupResponseV2 {
56    pub throttle_time_ms: i32,
57    pub error_code: i16,
58    pub generation_id: i32,
59    pub protocol_name: String,
60    pub leader: String,
61    pub member_id: String,
62    pub members: Vec<JoinGroupMember>,
63}
64
65impl JoinGroupResponseV2 {
66    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
67        Ok(Self {
68            throttle_time_ms: decoder.read_i32()?,
69            error_code: decoder.read_i16()?,
70            generation_id: decoder.read_i32()?,
71            protocol_name: decoder.read_string()?,
72            leader: decoder.read_string()?,
73            member_id: decoder.read_string()?,
74            members: decoder
75                .read_array("join group members", JoinGroupMember::decode)?
76                .unwrap_or_default(),
77        })
78    }
79}
80
81#[derive(Debug, Clone, PartialEq, Eq)]
82pub struct JoinGroupMember {
83    pub member_id: String,
84    pub metadata: Vec<u8>,
85}
86
87impl JoinGroupMember {
88    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
89        Ok(Self {
90            member_id: decoder.read_string()?,
91            metadata: decoder.read_bytes()?,
92        })
93    }
94}
95
96#[cfg(test)]
97#[allow(clippy::unwrap_used)]
98mod tests {
99    use super::{JoinGroupMember, JoinGroupProtocol, JoinGroupRequestV2, JoinGroupResponseV2};
100    use crate::codec::{Decoder, Encoder};
101
102    #[test]
103    fn encodes_join_group_v2_request() {
104        let request = JoinGroupRequestV2 {
105            correlation_id: 13,
106            client_id: Some("kafrust".to_owned()),
107            group_id: "orders-group".to_owned(),
108            session_timeout_ms: 10_000,
109            rebalance_timeout_ms: 30_000,
110            member_id: String::new(),
111            protocol_type: "consumer".to_owned(),
112            protocols: vec![JoinGroupProtocol {
113                name: "range".to_owned(),
114                metadata: vec![1, 2, 3],
115            }],
116        };
117
118        assert_eq!(
119            request.encode().unwrap(),
120            [
121                0, 11, // api key
122                0, 2, // api version
123                0, 0, 0, 13, // correlation id
124                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't', // client id
125                0, 12, b'o', b'r', b'd', b'e', b'r', b's', b'-', b'g', b'r', b'o', b'u',
126                b'p', // group id
127                0, 0, 39, 16, // session timeout
128                0, 0, 117, 48, // rebalance timeout
129                0, 0, // member id
130                0, 8, b'c', b'o', b'n', b's', b'u', b'm', b'e', b'r', // protocol type
131                0, 0, 0, 1, // protocol count
132                0, 5, b'r', b'a', b'n', b'g', b'e', // protocol name
133                0, 0, 0, 3, 1, 2, 3, // metadata
134            ]
135        );
136    }
137
138    #[test]
139    fn decodes_join_group_v2_response() {
140        let mut bytes = Encoder::new();
141        bytes.write_i32(0);
142        bytes.write_i16(0);
143        bytes.write_i32(7);
144        bytes.write_string("range").unwrap();
145        bytes.write_string("member-a").unwrap();
146        bytes.write_string("member-a").unwrap();
147        bytes.write_i32(1);
148        bytes.write_string("member-a").unwrap();
149        bytes.write_bytes(&[1, 2, 3]).unwrap();
150        let bytes = bytes.into_bytes();
151
152        let mut decoder = Decoder::new(&bytes);
153        let response = JoinGroupResponseV2::decode_body(&mut decoder).unwrap();
154
155        assert_eq!(response.throttle_time_ms, 0);
156        assert_eq!(response.error_code, 0);
157        assert_eq!(response.generation_id, 7);
158        assert_eq!(response.protocol_name, "range");
159        assert_eq!(response.leader, "member-a");
160        assert_eq!(response.member_id, "member-a");
161        assert_eq!(
162            response.members,
163            vec![JoinGroupMember {
164                member_id: "member-a".to_owned(),
165                metadata: vec![1, 2, 3],
166            }]
167        );
168        assert!(decoder.is_empty());
169    }
170}