Skip to main content

kafrust_protocol/api/
leave_group.rs

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