1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const API_KEY: i16 = 89;
7
8#[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 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#[derive(Debug, Clone, PartialEq, Eq)]
39pub struct StreamsGroupDescribeResponseV0 {
40 pub throttle_time_ms: i32,
41 pub groups: Vec<DescribedStreamsGroup>,
42}
43
44impl StreamsGroupDescribeResponseV0 {
45 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#[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#[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#[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#[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#[derive(Debug, Clone, PartialEq, Eq)]
205pub struct StreamsGroupTopicConfig {
206 pub key: String,
207 pub value: String,
208}
209
210#[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#[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#[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#[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#[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#[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}