Skip to main content

kafrust_protocol/api/
streams_group_heartbeat.rs

1#![allow(dead_code)]
2
3// The request topology types carry encode-only nested structures while the
4// response has a separate, smaller assignment shape. Keeping both wire shapes
5// typed here avoids lossy reuse and leaves the protocol boundary auditable.
6
7use crate::codec::{Decoder, Encoder};
8use crate::error::Result;
9use crate::header::RequestHeader;
10
11/// Kafka StreamsGroupHeartbeat API key.
12pub const API_KEY: i16 = 88;
13
14/// A Kafka Streams topology sent while initializing a Streams group.
15#[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/// One subtopology sent while initializing a Streams group.
48#[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/// A topic created or managed by a Streams topology.
116#[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/// One topic-level configuration in a Streams topology.
156#[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/// A copartition constraint represented by subtopology-level indexes.
179#[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/// A Streams task and its input partitions.
210#[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/// A cumulative changelog offset for a Streams task.
236#[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/// A host and port exposed for Interactive Queries.
266#[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/// A rack-aware client tag.
289#[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/// StreamsGroupHeartbeat v0 request.
312#[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    /// Encodes the flexible request, including its request header.
337    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/// A status entry returned by StreamsGroupHeartbeat.
387#[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/// Topic partitions materialized by one Interactive Queries endpoint.
406#[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/// Assignment information grouped by an Interactive Queries endpoint.
422#[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/// StreamsGroupHeartbeat v0 response.
454#[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    /// Decodes the flexible response body after the response header.
474    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}