Skip to main content

kafrust_protocol/api/
streams_group_describe.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5/// Kafka StreamsGroupDescribe API key.
6pub const API_KEY: i16 = 89;
7
8/// StreamsGroupDescribe v0 request.
9#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct StreamsGroupDescribeRequestV0 {
11    pub correlation_id: i32,
12    pub client_id: Option<String>,
13    pub group_ids: Vec<String>,
14    pub include_authorized_operations: bool,
15}
16
17impl StreamsGroupDescribeRequestV0 {
18    /// Encodes the flexible request, including its request header.
19    pub fn encode(&self) -> Result<Vec<u8>> {
20        let mut encoder = Encoder::new();
21        RequestHeader {
22            api_key: API_KEY,
23            api_version: 0,
24            correlation_id: self.correlation_id,
25            client_id: self.client_id.clone(),
26        }
27        .encode_v2(&mut encoder)?;
28        encoder.write_compact_array(Some(&self.group_ids), |encoder, group_id| {
29            encoder.write_compact_string(group_id)
30        })?;
31        encoder.write_bool(self.include_authorized_operations);
32        encoder.write_empty_tagged_fields();
33        Ok(encoder.into_bytes())
34    }
35}
36
37/// StreamsGroupDescribe v0 response.
38#[derive(Debug, Clone, PartialEq, Eq)]
39pub struct StreamsGroupDescribeResponseV0 {
40    pub throttle_time_ms: i32,
41    pub groups: Vec<DescribedStreamsGroup>,
42}
43
44impl StreamsGroupDescribeResponseV0 {
45    /// Decodes the flexible response body after the response header.
46    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
47        let throttle_time_ms = decoder.read_i32()?;
48        let groups = decoder
49            .read_compact_array("streams group descriptions", DescribedStreamsGroup::decode)?
50            .unwrap_or_default();
51        decoder.read_tagged_fields()?;
52        Ok(Self {
53            throttle_time_ms,
54            groups,
55        })
56    }
57}
58
59/// One Streams group returned by StreamsGroupDescribe.
60#[derive(Debug, Clone, PartialEq, Eq)]
61pub struct DescribedStreamsGroup {
62    pub error_code: i16,
63    pub error_message: Option<String>,
64    pub group_id: String,
65    pub group_state: String,
66    pub group_epoch: i32,
67    pub assignment_epoch: i32,
68    pub topology: Option<StreamsGroupTopology>,
69    pub members: Vec<DescribedStreamsGroupMember>,
70    pub authorized_operations: i32,
71}
72
73impl DescribedStreamsGroup {
74    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
75        let error_code = decoder.read_i16()?;
76        let error_message = decoder.read_compact_nullable_string()?;
77        let group_id = decoder.read_compact_string()?;
78        let group_state = decoder.read_compact_string()?;
79        let group_epoch = decoder.read_i32()?;
80        let assignment_epoch = decoder.read_i32()?;
81        let topology = decode_nullable_struct(decoder, StreamsGroupTopology::decode)?;
82        let members = decoder
83            .read_compact_array("streams group members", DescribedStreamsGroupMember::decode)?
84            .unwrap_or_default();
85        let authorized_operations = decoder.read_i32()?;
86        decoder.read_tagged_fields()?;
87        Ok(Self {
88            error_code,
89            error_message,
90            group_id,
91            group_state,
92            group_epoch,
93            assignment_epoch,
94            topology,
95            members,
96            authorized_operations,
97        })
98    }
99}
100
101/// The topology currently initialized for a Streams group.
102#[derive(Debug, Clone, PartialEq, Eq)]
103pub struct StreamsGroupTopology {
104    pub epoch: i32,
105    pub subtopologies: Option<Vec<StreamsGroupSubtopology>>,
106}
107
108impl StreamsGroupTopology {
109    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
110        let epoch = decoder.read_i32()?;
111        let subtopologies = decoder.read_compact_array(
112            "streams group subtopologies",
113            StreamsGroupSubtopology::decode,
114        )?;
115        decoder.read_tagged_fields()?;
116        Ok(Self {
117            epoch,
118            subtopologies,
119        })
120    }
121}
122
123/// One subtopology in a Streams group topology.
124#[derive(Debug, Clone, PartialEq, Eq)]
125pub struct StreamsGroupSubtopology {
126    pub subtopology_id: String,
127    pub source_topics: Vec<String>,
128    pub repartition_sink_topics: Vec<String>,
129    pub state_changelog_topics: Vec<StreamsGroupTopic>,
130    pub repartition_source_topics: Vec<StreamsGroupTopic>,
131}
132
133impl StreamsGroupSubtopology {
134    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
135        let subtopology_id = decoder.read_compact_string()?;
136        let source_topics = decoder
137            .read_compact_array("streams group source topics", |decoder| {
138                decoder.read_compact_string()
139            })?
140            .unwrap_or_default();
141        let repartition_sink_topics = decoder
142            .read_compact_array("streams group repartition sink topics", |decoder| {
143                decoder.read_compact_string()
144            })?
145            .unwrap_or_default();
146        let state_changelog_topics = decoder
147            .read_compact_array(
148                "streams group state changelog topics",
149                StreamsGroupTopic::decode,
150            )?
151            .unwrap_or_default();
152        let repartition_source_topics = decoder
153            .read_compact_array(
154                "streams group repartition source topics",
155                StreamsGroupTopic::decode,
156            )?
157            .unwrap_or_default();
158        decoder.read_tagged_fields()?;
159        Ok(Self {
160            subtopology_id,
161            source_topics,
162            repartition_sink_topics,
163            state_changelog_topics,
164            repartition_source_topics,
165        })
166    }
167}
168
169/// A topic managed by a Streams group topology.
170#[derive(Debug, Clone, PartialEq, Eq)]
171pub struct StreamsGroupTopic {
172    pub name: String,
173    pub partitions: i32,
174    pub replication_factor: i16,
175    pub topic_configs: Vec<StreamsGroupTopicConfig>,
176}
177
178impl StreamsGroupTopic {
179    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
180        let name = decoder.read_compact_string()?;
181        let partitions = decoder.read_i32()?;
182        let replication_factor = decoder.read_i16()?;
183        let topic_configs = decoder
184            .read_compact_array("streams group topic configs", |decoder| {
185                let config = StreamsGroupTopicConfig {
186                    key: decoder.read_compact_string()?,
187                    value: decoder.read_compact_string()?,
188                };
189                decoder.read_tagged_fields()?;
190                Ok(config)
191            })?
192            .unwrap_or_default();
193        decoder.read_tagged_fields()?;
194        Ok(Self {
195            name,
196            partitions,
197            replication_factor,
198            topic_configs,
199        })
200    }
201}
202
203/// One configuration entry for a Streams group-managed topic.
204#[derive(Debug, Clone, PartialEq, Eq)]
205pub struct StreamsGroupTopicConfig {
206    pub key: String,
207    pub value: String,
208}
209
210/// One member returned by StreamsGroupDescribe.
211#[derive(Debug, Clone, PartialEq, Eq)]
212pub struct DescribedStreamsGroupMember {
213    pub member_id: String,
214    pub member_epoch: i32,
215    pub instance_id: Option<String>,
216    pub rack_id: Option<String>,
217    pub client_id: String,
218    pub client_host: String,
219    pub topology_epoch: i32,
220    pub process_id: String,
221    pub user_endpoint: Option<StreamsGroupEndpoint>,
222    pub client_tags: Vec<StreamsGroupKeyValue>,
223    pub task_offsets: Vec<StreamsGroupTaskOffset>,
224    pub task_end_offsets: Vec<StreamsGroupTaskOffset>,
225    pub assignment: StreamsGroupAssignment,
226    pub target_assignment: StreamsGroupAssignment,
227    pub is_classic: bool,
228}
229
230impl DescribedStreamsGroupMember {
231    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
232        let member_id = decoder.read_compact_string()?;
233        let member_epoch = decoder.read_i32()?;
234        let instance_id = decoder.read_compact_nullable_string()?;
235        let rack_id = decoder.read_compact_nullable_string()?;
236        let client_id = decoder.read_compact_string()?;
237        let client_host = decoder.read_compact_string()?;
238        let topology_epoch = decoder.read_i32()?;
239        let process_id = decoder.read_compact_string()?;
240        let user_endpoint = decode_nullable_struct(decoder, StreamsGroupEndpoint::decode)?;
241        let client_tags = decode_key_values(decoder, "streams group client tags")?;
242        let task_offsets = decoder
243            .read_compact_array("streams group task offsets", StreamsGroupTaskOffset::decode)?
244            .unwrap_or_default();
245        let task_end_offsets = decoder
246            .read_compact_array(
247                "streams group task end offsets",
248                StreamsGroupTaskOffset::decode,
249            )?
250            .unwrap_or_default();
251        let assignment = StreamsGroupAssignment::decode(decoder)?;
252        let target_assignment = StreamsGroupAssignment::decode(decoder)?;
253        let is_classic = decoder.read_bool()?;
254        decoder.read_tagged_fields()?;
255        Ok(Self {
256            member_id,
257            member_epoch,
258            instance_id,
259            rack_id,
260            client_id,
261            client_host,
262            topology_epoch,
263            process_id,
264            user_endpoint,
265            client_tags,
266            task_offsets,
267            task_end_offsets,
268            assignment,
269            target_assignment,
270            is_classic,
271        })
272    }
273}
274
275/// A host and port exposed by a Streams group member.
276#[derive(Debug, Clone, PartialEq, Eq)]
277pub struct StreamsGroupEndpoint {
278    pub host: String,
279    pub port: u16,
280}
281
282impl StreamsGroupEndpoint {
283    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
284        let host = decoder.read_compact_string()?;
285        let port = decoder.read_i16()? as u16;
286        decoder.read_tagged_fields()?;
287        Ok(Self { host, port })
288    }
289}
290
291/// A compact key/value entry attached to a Streams group member.
292#[derive(Debug, Clone, PartialEq, Eq)]
293pub struct StreamsGroupKeyValue {
294    pub key: String,
295    pub value: String,
296}
297
298fn decode_key_values(
299    decoder: &mut Decoder<'_>,
300    kind: &'static str,
301) -> Result<Vec<StreamsGroupKeyValue>> {
302    Ok(decoder
303        .read_compact_array(kind, |decoder| {
304            let key = decoder.read_compact_string()?;
305            let value = decoder.read_compact_string()?;
306            decoder.read_tagged_fields()?;
307            Ok(StreamsGroupKeyValue { key, value })
308        })?
309        .unwrap_or_default())
310}
311
312/// Current or target task assignment for a Streams group member.
313#[derive(Debug, Clone, PartialEq, Eq)]
314pub struct StreamsGroupAssignment {
315    pub active_tasks: Vec<StreamsGroupTask>,
316    pub standby_tasks: Vec<StreamsGroupTask>,
317    pub warmup_tasks: Vec<StreamsGroupTask>,
318}
319
320impl StreamsGroupAssignment {
321    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
322        let active_tasks = decoder
323            .read_compact_array("streams group active tasks", StreamsGroupTask::decode)?
324            .unwrap_or_default();
325        let standby_tasks = decoder
326            .read_compact_array("streams group standby tasks", StreamsGroupTask::decode)?
327            .unwrap_or_default();
328        let warmup_tasks = decoder
329            .read_compact_array("streams group warmup tasks", StreamsGroupTask::decode)?
330            .unwrap_or_default();
331        decoder.read_tagged_fields()?;
332        Ok(Self {
333            active_tasks,
334            standby_tasks,
335            warmup_tasks,
336        })
337    }
338}
339
340/// A Streams task assignment identified by subtopology and partitions.
341#[derive(Debug, Clone, PartialEq, Eq)]
342pub struct StreamsGroupTask {
343    pub subtopology_id: String,
344    pub partitions: Vec<i32>,
345}
346
347impl StreamsGroupTask {
348    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
349        let subtopology_id = decoder.read_compact_string()?;
350        let partitions = decoder
351            .read_compact_array("streams group task partitions", |decoder| {
352                decoder.read_i32()
353            })?
354            .unwrap_or_default();
355        decoder.read_tagged_fields()?;
356        Ok(Self {
357            subtopology_id,
358            partitions,
359        })
360    }
361}
362
363/// A cumulative changelog offset reported by a Streams group member.
364#[derive(Debug, Clone, PartialEq, Eq)]
365pub struct StreamsGroupTaskOffset {
366    pub subtopology_id: String,
367    pub partition: i32,
368    pub offset: i64,
369}
370
371impl StreamsGroupTaskOffset {
372    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
373        let subtopology_id = decoder.read_compact_string()?;
374        let partition = decoder.read_i32()?;
375        let offset = decoder.read_i64()?;
376        decoder.read_tagged_fields()?;
377        Ok(Self {
378            subtopology_id,
379            partition,
380            offset,
381        })
382    }
383}
384
385fn decode_nullable_struct<T>(
386    decoder: &mut Decoder<'_>,
387    decode: impl FnOnce(&mut Decoder<'_>) -> Result<T>,
388) -> Result<Option<T>> {
389    match decoder.read_i8()? {
390        -1 => Ok(None),
391        1 => decode(decoder).map(Some),
392        marker => Err(crate::error::Error::InvalidNullableStruct(marker)),
393    }
394}
395
396#[cfg(test)]
397#[allow(clippy::unwrap_used)]
398mod tests {
399    use super::{StreamsGroupDescribeRequestV0, StreamsGroupDescribeResponseV0, API_KEY};
400    use crate::codec::{Decoder, Encoder};
401
402    #[test]
403    fn encodes_streams_group_describe_v0_request() {
404        let request = StreamsGroupDescribeRequestV0 {
405            correlation_id: 23,
406            client_id: Some("kafrust".to_owned()),
407            group_ids: vec!["streams-orders".to_owned()],
408            include_authorized_operations: true,
409        };
410
411        let encoded = request.encode().unwrap();
412        assert_eq!(&encoded[..4], &[0, 89, 0, 0]);
413        assert_eq!(API_KEY, 89);
414    }
415
416    #[test]
417    fn decodes_streams_group_describe_v0_response() -> crate::error::Result<()> {
418        let mut bytes = Encoder::new();
419        bytes.write_i32(12);
420        bytes.write_compact_array(Some(&[1_i8]), |encoder, _| {
421            encoder.write_i16(0);
422            encoder.write_compact_nullable_string(Some("ok"))?;
423            encoder.write_compact_string("streams-orders")?;
424            encoder.write_compact_string("Stable")?;
425            encoder.write_i32(4);
426            encoder.write_i32(5);
427            encoder.write_i8(1);
428            encoder.write_i32(3);
429            encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
430                encoder.write_compact_string("subtopology-0")?;
431                encoder.write_compact_array(Some(&["orders".to_owned()]), |encoder, topic| {
432                    encoder.write_compact_string(topic)
433                })?;
434                encoder.write_compact_array::<String>(Some(&[]), |_, _| Ok(()))?;
435                encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
436                    encoder.write_compact_string("orders-store")?;
437                    encoder.write_i32(3);
438                    encoder.write_i16(1);
439                    encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
440                        encoder.write_compact_string("cleanup.policy")?;
441                        encoder.write_compact_string("compact")?;
442                        encoder.write_empty_tagged_fields();
443                        Ok(())
444                    })?;
445                    encoder.write_empty_tagged_fields();
446                    Ok(())
447                })?;
448                encoder.write_compact_array::<i8>(Some(&[]), |_, _| Ok(()))?;
449                encoder.write_empty_tagged_fields();
450                Ok(())
451            })?;
452            encoder.write_empty_tagged_fields();
453            bytes_for_member(encoder)?;
454            encoder.write_i32(-2147483648);
455            encoder.write_empty_tagged_fields();
456            Ok(())
457        })?;
458        bytes.write_empty_tagged_fields();
459
460        let encoded = bytes.into_bytes();
461        let mut decoder = Decoder::new(&encoded);
462        let response = StreamsGroupDescribeResponseV0::decode_body(&mut decoder)?;
463
464        assert_eq!(response.throttle_time_ms, 12);
465        assert_eq!(response.groups[0].group_id, "streams-orders");
466        assert_eq!(response.groups[0].topology.as_ref().unwrap().epoch, 3);
467        assert_eq!(
468            response.groups[0]
469                .topology
470                .as_ref()
471                .unwrap()
472                .subtopologies
473                .as_ref()
474                .unwrap()[0]
475                .state_changelog_topics[0]
476                .name,
477            "orders-store"
478        );
479        assert_eq!(response.groups[0].members[0].member_id, "member-1");
480        assert_eq!(
481            response.groups[0].members[0].assignment.active_tasks[0].partitions,
482            [0, 2]
483        );
484        assert!(decoder.is_empty());
485        Ok(())
486    }
487
488    fn bytes_for_member(encoder: &mut Encoder) -> crate::error::Result<()> {
489        encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
490            encoder.write_compact_string("member-1")?;
491            encoder.write_i32(7);
492            encoder.write_compact_nullable_string(None)?;
493            encoder.write_compact_nullable_string(Some("rack-a"))?;
494            encoder.write_compact_string("client-a")?;
495            encoder.write_compact_string("/127.0.0.1")?;
496            encoder.write_i32(3);
497            encoder.write_compact_string("process-1")?;
498            encoder.write_i8(1);
499            encoder.write_compact_string("127.0.0.1")?;
500            encoder.write_i16(7000);
501            encoder.write_empty_tagged_fields();
502            encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
503                encoder.write_compact_string("rack")?;
504                encoder.write_compact_string("a")?;
505                encoder.write_empty_tagged_fields();
506                Ok(())
507            })?;
508            encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
509                encoder.write_compact_string("subtopology-0")?;
510                encoder.write_i32(0);
511                encoder.write_i64(10);
512                encoder.write_empty_tagged_fields();
513                Ok(())
514            })?;
515            encoder.write_compact_array::<i8>(Some(&[]), |_, _| Ok(()))?;
516            for _ in 0..2 {
517                encoder.write_compact_array(Some(&[1_i8]), |encoder, _| {
518                    encoder.write_compact_string("subtopology-0")?;
519                    encoder.write_compact_array(Some(&[0_i32, 2]), |encoder, partition| {
520                        encoder.write_i32(*partition);
521                        Ok(())
522                    })?;
523                    encoder.write_empty_tagged_fields();
524                    Ok(())
525                })?;
526                encoder.write_compact_array::<i8>(Some(&[]), |_, _| Ok(()))?;
527                encoder.write_compact_array::<i8>(Some(&[]), |_, _| Ok(()))?;
528                encoder.write_empty_tagged_fields();
529            }
530            encoder.write_bool(false);
531            encoder.write_empty_tagged_fields();
532            Ok(())
533        })?;
534        Ok(())
535    }
536}