kafrust_protocol/api/
leave_group.rs1use 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}