1#![allow(dead_code)]
2
3use crate::codec::{Decoder, Encoder};
8use crate::error::Result;
9use crate::header::RequestHeader;
10
11pub const API_KEY: i16 = 88;
13
14#[derive(Debug, Clone, PartialEq, Eq)]
16pub struct StreamsGroupHeartbeatTopology {
17 pub epoch: i32,
18 pub subtopologies: Vec<StreamsGroupHeartbeatSubtopology>,
19}
20
21impl StreamsGroupHeartbeatTopology {
22 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
23 encoder.write_i32(self.epoch);
24 encoder.write_compact_array(Some(&self.subtopologies), |encoder, subtopology| {
25 subtopology.encode(encoder)
26 })?;
27 encoder.write_empty_tagged_fields();
28 Ok(())
29 }
30
31 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
32 let epoch = decoder.read_i32()?;
33 let subtopologies = decoder
34 .read_compact_array(
35 "streams heartbeat subtopologies",
36 StreamsGroupHeartbeatSubtopology::decode,
37 )?
38 .unwrap_or_default();
39 decoder.read_tagged_fields()?;
40 Ok(Self {
41 epoch,
42 subtopologies,
43 })
44 }
45}
46
47#[derive(Debug, Clone, PartialEq, Eq)]
49pub struct StreamsGroupHeartbeatSubtopology {
50 pub subtopology_id: String,
51 pub source_topics: Vec<String>,
52 pub source_topic_regex: Vec<String>,
53 pub state_changelog_topics: Vec<StreamsGroupHeartbeatTopic>,
54 pub repartition_sink_topics: Vec<String>,
55 pub repartition_source_topics: Vec<StreamsGroupHeartbeatTopic>,
56 pub copartition_groups: Vec<StreamsGroupHeartbeatCopartitionGroup>,
57}
58
59impl StreamsGroupHeartbeatSubtopology {
60 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
61 encoder.write_compact_string(&self.subtopology_id)?;
62 write_string_array(encoder, &self.source_topics)?;
63 write_string_array(encoder, &self.source_topic_regex)?;
64 encoder.write_compact_array(Some(&self.state_changelog_topics), |encoder, topic| {
65 topic.encode(encoder)
66 })?;
67 write_string_array(encoder, &self.repartition_sink_topics)?;
68 encoder.write_compact_array(Some(&self.repartition_source_topics), |encoder, topic| {
69 topic.encode(encoder)
70 })?;
71 encoder.write_compact_array(Some(&self.copartition_groups), |encoder, group| {
72 group.encode(encoder)
73 })?;
74 encoder.write_empty_tagged_fields();
75 Ok(())
76 }
77
78 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
79 let subtopology_id = decoder.read_compact_string()?;
80 let source_topics = read_string_array(decoder, "streams heartbeat source topics")?;
81 let source_topic_regex = read_string_array(decoder, "streams heartbeat source regex")?;
82 let state_changelog_topics = decoder
83 .read_compact_array(
84 "streams heartbeat state changelog topics",
85 StreamsGroupHeartbeatTopic::decode,
86 )?
87 .unwrap_or_default();
88 let repartition_sink_topics =
89 read_string_array(decoder, "streams heartbeat repartition sinks")?;
90 let repartition_source_topics = decoder
91 .read_compact_array(
92 "streams heartbeat repartition sources",
93 StreamsGroupHeartbeatTopic::decode,
94 )?
95 .unwrap_or_default();
96 let copartition_groups = decoder
97 .read_compact_array(
98 "streams heartbeat copartition groups",
99 StreamsGroupHeartbeatCopartitionGroup::decode,
100 )?
101 .unwrap_or_default();
102 decoder.read_tagged_fields()?;
103 Ok(Self {
104 subtopology_id,
105 source_topics,
106 source_topic_regex,
107 state_changelog_topics,
108 repartition_sink_topics,
109 repartition_source_topics,
110 copartition_groups,
111 })
112 }
113}
114
115#[derive(Debug, Clone, PartialEq, Eq)]
117pub struct StreamsGroupHeartbeatTopic {
118 pub name: String,
119 pub partitions: i32,
120 pub replication_factor: i16,
121 pub topic_configs: Vec<StreamsGroupHeartbeatTopicConfig>,
122}
123
124impl StreamsGroupHeartbeatTopic {
125 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
126 encoder.write_compact_string(&self.name)?;
127 encoder.write_i32(self.partitions);
128 encoder.write_i16(self.replication_factor);
129 encoder.write_compact_array(Some(&self.topic_configs), |encoder, config| {
130 config.encode(encoder)
131 })?;
132 encoder.write_empty_tagged_fields();
133 Ok(())
134 }
135
136 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
137 let name = decoder.read_compact_string()?;
138 let partitions = decoder.read_i32()?;
139 let replication_factor = decoder.read_i16()?;
140 let topic_configs = decoder
141 .read_compact_array("streams heartbeat topic configs", |decoder| {
142 StreamsGroupHeartbeatTopicConfig::decode(decoder)
143 })?
144 .unwrap_or_default();
145 decoder.read_tagged_fields()?;
146 Ok(Self {
147 name,
148 partitions,
149 replication_factor,
150 topic_configs,
151 })
152 }
153}
154
155#[derive(Debug, Clone, PartialEq, Eq)]
157pub struct StreamsGroupHeartbeatTopicConfig {
158 pub key: String,
159 pub value: String,
160}
161
162impl StreamsGroupHeartbeatTopicConfig {
163 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
164 encoder.write_compact_string(&self.key)?;
165 encoder.write_compact_string(&self.value)?;
166 encoder.write_empty_tagged_fields();
167 Ok(())
168 }
169
170 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
171 let key = decoder.read_compact_string()?;
172 let value = decoder.read_compact_string()?;
173 decoder.read_tagged_fields()?;
174 Ok(Self { key, value })
175 }
176}
177
178#[derive(Debug, Clone, PartialEq, Eq)]
180pub struct StreamsGroupHeartbeatCopartitionGroup {
181 pub source_topics: Vec<i16>,
182 pub source_topic_regex: Vec<i16>,
183 pub repartition_source_topics: Vec<i16>,
184}
185
186impl StreamsGroupHeartbeatCopartitionGroup {
187 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
188 write_i16_array(encoder, &self.source_topics)?;
189 write_i16_array(encoder, &self.source_topic_regex)?;
190 write_i16_array(encoder, &self.repartition_source_topics)?;
191 encoder.write_empty_tagged_fields();
192 Ok(())
193 }
194
195 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
196 let source_topics = read_i16_array(decoder, "streams heartbeat copartition sources")?;
197 let source_topic_regex = read_i16_array(decoder, "streams heartbeat copartition regex")?;
198 let repartition_source_topics =
199 read_i16_array(decoder, "streams heartbeat copartition repartition sources")?;
200 decoder.read_tagged_fields()?;
201 Ok(Self {
202 source_topics,
203 source_topic_regex,
204 repartition_source_topics,
205 })
206 }
207}
208
209#[derive(Debug, Clone, PartialEq, Eq)]
211pub struct StreamsGroupHeartbeatTask {
212 pub subtopology_id: String,
213 pub partitions: Vec<i32>,
214}
215
216impl StreamsGroupHeartbeatTask {
217 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
218 encoder.write_compact_string(&self.subtopology_id)?;
219 write_i32_array(encoder, &self.partitions)?;
220 encoder.write_empty_tagged_fields();
221 Ok(())
222 }
223
224 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
225 let subtopology_id = decoder.read_compact_string()?;
226 let partitions = read_i32_array(decoder, "streams heartbeat task partitions")?;
227 decoder.read_tagged_fields()?;
228 Ok(Self {
229 subtopology_id,
230 partitions,
231 })
232 }
233}
234
235#[derive(Debug, Clone, PartialEq, Eq)]
237pub struct StreamsGroupHeartbeatTaskOffset {
238 pub subtopology_id: String,
239 pub partition: i32,
240 pub offset: i64,
241}
242
243impl StreamsGroupHeartbeatTaskOffset {
244 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
245 encoder.write_compact_string(&self.subtopology_id)?;
246 encoder.write_i32(self.partition);
247 encoder.write_i64(self.offset);
248 encoder.write_empty_tagged_fields();
249 Ok(())
250 }
251
252 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
253 let subtopology_id = decoder.read_compact_string()?;
254 let partition = decoder.read_i32()?;
255 let offset = decoder.read_i64()?;
256 decoder.read_tagged_fields()?;
257 Ok(Self {
258 subtopology_id,
259 partition,
260 offset,
261 })
262 }
263}
264
265#[derive(Debug, Clone, PartialEq, Eq)]
267pub struct StreamsGroupHeartbeatEndpoint {
268 pub host: String,
269 pub port: u16,
270}
271
272impl StreamsGroupHeartbeatEndpoint {
273 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
274 encoder.write_compact_string(&self.host)?;
275 encoder.write_i16(i16::from_be_bytes(self.port.to_be_bytes()));
276 encoder.write_empty_tagged_fields();
277 Ok(())
278 }
279
280 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
281 let host = decoder.read_compact_string()?;
282 let port = u16::from_be_bytes(decoder.read_i16()?.to_be_bytes());
283 decoder.read_tagged_fields()?;
284 Ok(Self { host, port })
285 }
286}
287
288#[derive(Debug, Clone, PartialEq, Eq)]
290pub struct StreamsGroupHeartbeatKeyValue {
291 pub key: String,
292 pub value: String,
293}
294
295impl StreamsGroupHeartbeatKeyValue {
296 fn encode(&self, encoder: &mut Encoder) -> Result<()> {
297 encoder.write_compact_string(&self.key)?;
298 encoder.write_compact_string(&self.value)?;
299 encoder.write_empty_tagged_fields();
300 Ok(())
301 }
302
303 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
304 let key = decoder.read_compact_string()?;
305 let value = decoder.read_compact_string()?;
306 decoder.read_tagged_fields()?;
307 Ok(Self { key, value })
308 }
309}
310
311#[derive(Debug, Clone, PartialEq, Eq)]
313pub struct StreamsGroupHeartbeatRequestV0 {
314 pub correlation_id: i32,
315 pub client_id: Option<String>,
316 pub group_id: String,
317 pub member_id: String,
318 pub member_epoch: i32,
319 pub endpoint_information_epoch: i32,
320 pub instance_id: Option<String>,
321 pub rack_id: Option<String>,
322 pub rebalance_timeout_ms: i32,
323 pub topology: Option<StreamsGroupHeartbeatTopology>,
324 pub active_tasks: Option<Vec<StreamsGroupHeartbeatTask>>,
325 pub standby_tasks: Option<Vec<StreamsGroupHeartbeatTask>>,
326 pub warmup_tasks: Option<Vec<StreamsGroupHeartbeatTask>>,
327 pub process_id: Option<String>,
328 pub user_endpoint: Option<StreamsGroupHeartbeatEndpoint>,
329 pub client_tags: Option<Vec<StreamsGroupHeartbeatKeyValue>>,
330 pub task_offsets: Option<Vec<StreamsGroupHeartbeatTaskOffset>>,
331 pub task_end_offsets: Option<Vec<StreamsGroupHeartbeatTaskOffset>>,
332 pub shutdown_application: bool,
333}
334
335impl StreamsGroupHeartbeatRequestV0 {
336 pub fn encode(&self) -> Result<Vec<u8>> {
338 let mut encoder = Encoder::new();
339 RequestHeader {
340 api_key: API_KEY,
341 api_version: 0,
342 correlation_id: self.correlation_id,
343 client_id: self.client_id.clone(),
344 }
345 .encode_v2(&mut encoder)?;
346 encoder.write_compact_string(&self.group_id)?;
347 encoder.write_compact_string(&self.member_id)?;
348 encoder.write_i32(self.member_epoch);
349 encoder.write_i32(self.endpoint_information_epoch);
350 encoder.write_compact_nullable_string(self.instance_id.as_deref())?;
351 encoder.write_compact_nullable_string(self.rack_id.as_deref())?;
352 encoder.write_i32(self.rebalance_timeout_ms);
353 write_nullable_struct(&mut encoder, self.topology.as_ref(), |encoder, topology| {
354 topology.encode(encoder)
355 })?;
356 encoder.write_compact_array(self.active_tasks.as_deref(), |encoder, task| {
357 task.encode(encoder)
358 })?;
359 encoder.write_compact_array(self.standby_tasks.as_deref(), |encoder, task| {
360 task.encode(encoder)
361 })?;
362 encoder.write_compact_array(self.warmup_tasks.as_deref(), |encoder, task| {
363 task.encode(encoder)
364 })?;
365 encoder.write_compact_nullable_string(self.process_id.as_deref())?;
366 write_nullable_struct(
367 &mut encoder,
368 self.user_endpoint.as_ref(),
369 |encoder, endpoint| endpoint.encode(encoder),
370 )?;
371 encoder.write_compact_array(self.client_tags.as_deref(), |encoder, tag| {
372 tag.encode(encoder)
373 })?;
374 encoder.write_compact_array(self.task_offsets.as_deref(), |encoder, offset| {
375 offset.encode(encoder)
376 })?;
377 encoder.write_compact_array(self.task_end_offsets.as_deref(), |encoder, offset| {
378 offset.encode(encoder)
379 })?;
380 encoder.write_bool(self.shutdown_application);
381 encoder.write_empty_tagged_fields();
382 Ok(encoder.into_bytes())
383 }
384}
385
386#[derive(Debug, Clone, PartialEq, Eq)]
388pub struct StreamsGroupHeartbeatStatus {
389 pub status_code: i8,
390 pub status_detail: String,
391}
392
393impl StreamsGroupHeartbeatStatus {
394 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
395 let status_code = decoder.read_i8()?;
396 let status_detail = decoder.read_compact_string()?;
397 decoder.read_tagged_fields()?;
398 Ok(Self {
399 status_code,
400 status_detail,
401 })
402 }
403}
404
405#[derive(Debug, Clone, PartialEq, Eq)]
407pub struct StreamsGroupHeartbeatTopicPartitions {
408 pub topic: String,
409 pub partitions: Vec<i32>,
410}
411
412impl StreamsGroupHeartbeatTopicPartitions {
413 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
414 let topic = decoder.read_compact_string()?;
415 let partitions = read_i32_array(decoder, "streams heartbeat endpoint partitions")?;
416 decoder.read_tagged_fields()?;
417 Ok(Self { topic, partitions })
418 }
419}
420
421#[derive(Debug, Clone, PartialEq, Eq)]
423pub struct StreamsGroupHeartbeatEndpointPartitions {
424 pub user_endpoint: StreamsGroupHeartbeatEndpoint,
425 pub active_partitions: Vec<StreamsGroupHeartbeatTopicPartitions>,
426 pub standby_partitions: Vec<StreamsGroupHeartbeatTopicPartitions>,
427}
428
429impl StreamsGroupHeartbeatEndpointPartitions {
430 fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
431 let user_endpoint = StreamsGroupHeartbeatEndpoint::decode(decoder)?;
432 let active_partitions = decoder
433 .read_compact_array(
434 "streams heartbeat active endpoint partitions",
435 StreamsGroupHeartbeatTopicPartitions::decode,
436 )?
437 .unwrap_or_default();
438 let standby_partitions = decoder
439 .read_compact_array(
440 "streams heartbeat standby endpoint partitions",
441 StreamsGroupHeartbeatTopicPartitions::decode,
442 )?
443 .unwrap_or_default();
444 decoder.read_tagged_fields()?;
445 Ok(Self {
446 user_endpoint,
447 active_partitions,
448 standby_partitions,
449 })
450 }
451}
452
453#[derive(Debug, Clone, PartialEq, Eq)]
455pub struct StreamsGroupHeartbeatResponseV0 {
456 pub throttle_time_ms: i32,
457 pub error_code: i16,
458 pub error_message: Option<String>,
459 pub member_id: String,
460 pub member_epoch: i32,
461 pub heartbeat_interval_ms: i32,
462 pub acceptable_recovery_lag: i32,
463 pub task_offset_interval_ms: i32,
464 pub status: Option<Vec<StreamsGroupHeartbeatStatus>>,
465 pub active_tasks: Option<Vec<StreamsGroupHeartbeatTask>>,
466 pub standby_tasks: Option<Vec<StreamsGroupHeartbeatTask>>,
467 pub warmup_tasks: Option<Vec<StreamsGroupHeartbeatTask>>,
468 pub endpoint_information_epoch: i32,
469 pub partitions_by_user_endpoint: Option<Vec<StreamsGroupHeartbeatEndpointPartitions>>,
470}
471
472impl StreamsGroupHeartbeatResponseV0 {
473 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
475 let throttle_time_ms = decoder.read_i32()?;
476 let error_code = decoder.read_i16()?;
477 let error_message = decoder.read_compact_nullable_string()?;
478 let member_id = decoder.read_compact_string()?;
479 let member_epoch = decoder.read_i32()?;
480 let heartbeat_interval_ms = decoder.read_i32()?;
481 let acceptable_recovery_lag = decoder.read_i32()?;
482 let task_offset_interval_ms = decoder.read_i32()?;
483 let status = decoder.read_compact_array(
484 "streams heartbeat statuses",
485 StreamsGroupHeartbeatStatus::decode,
486 )?;
487 let active_tasks = decoder.read_compact_array(
488 "streams heartbeat active tasks",
489 StreamsGroupHeartbeatTask::decode,
490 )?;
491 let standby_tasks = decoder.read_compact_array(
492 "streams heartbeat standby tasks",
493 StreamsGroupHeartbeatTask::decode,
494 )?;
495 let warmup_tasks = decoder.read_compact_array(
496 "streams heartbeat warmup tasks",
497 StreamsGroupHeartbeatTask::decode,
498 )?;
499 let endpoint_information_epoch = decoder.read_i32()?;
500 let partitions_by_user_endpoint = decoder.read_compact_array(
501 "streams heartbeat endpoint assignments",
502 StreamsGroupHeartbeatEndpointPartitions::decode,
503 )?;
504 decoder.read_tagged_fields()?;
505 Ok(Self {
506 throttle_time_ms,
507 error_code,
508 error_message,
509 member_id,
510 member_epoch,
511 heartbeat_interval_ms,
512 acceptable_recovery_lag,
513 task_offset_interval_ms,
514 status,
515 active_tasks,
516 standby_tasks,
517 warmup_tasks,
518 endpoint_information_epoch,
519 partitions_by_user_endpoint,
520 })
521 }
522}
523
524fn write_nullable_struct<T>(
525 encoder: &mut Encoder,
526 value: Option<&T>,
527 mut write: impl FnMut(&mut Encoder, &T) -> Result<()>,
528) -> Result<()> {
529 match value {
530 Some(value) => {
531 encoder.write_i8(1);
532 write(encoder, value)?;
533 }
534 None => encoder.write_i8(-1),
535 }
536 Ok(())
537}
538
539fn write_string_array(encoder: &mut Encoder, values: &[String]) -> Result<()> {
540 encoder.write_compact_array(Some(values), |encoder, value| {
541 encoder.write_compact_string(value)
542 })
543}
544
545fn read_string_array(decoder: &mut Decoder<'_>, kind: &'static str) -> Result<Vec<String>> {
546 Ok(decoder
547 .read_compact_array(kind, |decoder| decoder.read_compact_string())?
548 .unwrap_or_default())
549}
550
551fn write_i16_array(encoder: &mut Encoder, values: &[i16]) -> Result<()> {
552 encoder.write_compact_array(Some(values), |encoder, value| {
553 encoder.write_i16(*value);
554 Ok(())
555 })
556}
557
558fn read_i16_array(decoder: &mut Decoder<'_>, kind: &'static str) -> Result<Vec<i16>> {
559 Ok(decoder
560 .read_compact_array(kind, |decoder| decoder.read_i16())?
561 .unwrap_or_default())
562}
563
564fn write_i32_array(encoder: &mut Encoder, values: &[i32]) -> Result<()> {
565 encoder.write_compact_array(Some(values), |encoder, value| {
566 encoder.write_i32(*value);
567 Ok(())
568 })
569}
570
571fn read_i32_array(decoder: &mut Decoder<'_>, kind: &'static str) -> Result<Vec<i32>> {
572 Ok(decoder
573 .read_compact_array(kind, |decoder| decoder.read_i32())?
574 .unwrap_or_default())
575}
576
577#[cfg(test)]
578#[allow(clippy::unwrap_used)]
579mod tests {
580 use super::{
581 StreamsGroupHeartbeatRequestV0, StreamsGroupHeartbeatResponseV0,
582 StreamsGroupHeartbeatStatus, StreamsGroupHeartbeatTask, StreamsGroupHeartbeatTopology,
583 API_KEY,
584 };
585 use crate::codec::{Decoder, Encoder};
586
587 #[test]
588 fn encodes_streams_group_heartbeat_v0_request() {
589 let request = StreamsGroupHeartbeatRequestV0 {
590 correlation_id: 23,
591 client_id: Some("kafrust".to_owned()),
592 group_id: "streams-orders".to_owned(),
593 member_id: "member-a".to_owned(),
594 member_epoch: 0,
595 endpoint_information_epoch: 0,
596 instance_id: None,
597 rack_id: None,
598 rebalance_timeout_ms: 30_000,
599 topology: Some(StreamsGroupHeartbeatTopology {
600 epoch: 1,
601 subtopologies: Vec::new(),
602 }),
603 active_tasks: Some(vec![StreamsGroupHeartbeatTask {
604 subtopology_id: "subtopology-0".to_owned(),
605 partitions: vec![0, 2],
606 }]),
607 standby_tasks: None,
608 warmup_tasks: None,
609 process_id: Some("process-a".to_owned()),
610 user_endpoint: None,
611 client_tags: None,
612 task_offsets: None,
613 task_end_offsets: None,
614 shutdown_application: false,
615 };
616
617 let encoded = request.encode().unwrap();
618 assert_eq!(&encoded[..4], &[0, 88, 0, 0]);
619 assert_eq!(API_KEY, 88);
620 assert!(encoded.windows(14).any(|bytes| bytes == b"streams-orders"));
621 assert!(encoded.windows(8).any(|bytes| bytes == b"member-a"));
622 assert!(encoded.ends_with(&[0]));
623 }
624
625 #[test]
626 fn decodes_streams_group_heartbeat_v0_response() -> crate::error::Result<()> {
627 let mut bytes = Encoder::new();
628 bytes.write_i32(12);
629 bytes.write_i16(0);
630 bytes.write_compact_nullable_string(Some("ok"))?;
631 bytes.write_compact_string("member-a")?;
632 bytes.write_i32(3);
633 bytes.write_i32(2500);
634 bytes.write_i32(10);
635 bytes.write_i32(1000);
636 bytes.write_compact_array(
637 Some(&[StreamsGroupHeartbeatStatus {
638 status_code: 2,
639 status_detail: "running".to_owned(),
640 }]),
641 |encoder, status| {
642 encoder.write_i8(status.status_code);
643 encoder.write_compact_string(&status.status_detail)?;
644 encoder.write_empty_tagged_fields();
645 Ok(())
646 },
647 )?;
648 bytes.write_compact_array(
649 Some(&[StreamsGroupHeartbeatTask {
650 subtopology_id: "subtopology-0".to_owned(),
651 partitions: vec![0, 1],
652 }]),
653 |encoder, task| task.encode(encoder),
654 )?;
655 bytes.write_compact_array::<StreamsGroupHeartbeatTask>(Some(&[]), |_, _| Ok(()))?;
656 bytes.write_compact_array::<StreamsGroupHeartbeatTask>(None, |_, _| Ok(()))?;
657 bytes.write_i32(4);
658 bytes.write_compact_array::<i8>(None, |_, _| Ok(()))?;
659 bytes.write_empty_tagged_fields();
660
661 let encoded = bytes.into_bytes();
662 let mut decoder = Decoder::new(&encoded);
663 let response = StreamsGroupHeartbeatResponseV0::decode_body(&mut decoder)?;
664
665 assert_eq!(response.throttle_time_ms, 12);
666 assert_eq!(response.member_id, "member-a");
667 assert_eq!(response.status.as_ref().unwrap()[0].status_code, 2);
668 assert_eq!(
669 response.active_tasks.as_ref().unwrap()[0].partitions,
670 [0, 1]
671 );
672 assert!(response.standby_tasks.as_ref().unwrap().is_empty());
673 assert!(response.warmup_tasks.is_none());
674 assert!(response.partitions_by_user_endpoint.is_none());
675 assert!(decoder.is_empty());
676 Ok(())
677 }
678}