Skip to main content

kafrust_protocol/api/
sync_group.rs

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