Skip to main content

kafrust_protocol/api/
share.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::{Error, Result};
3use crate::header::RequestHeader;
4
5/// Kafka ShareGroupHeartbeat API key.
6pub const SHARE_GROUP_HEARTBEAT_API_KEY: i16 = 76;
7/// Kafka ShareFetch API key.
8pub const SHARE_FETCH_API_KEY: i16 = 78;
9/// Kafka ShareAcknowledge API key.
10pub const SHARE_ACKNOWLEDGE_API_KEY: i16 = 79;
11
12/// A topic and its partition indexes in a share-group assignment.
13#[derive(Debug, Clone, PartialEq, Eq)]
14pub struct ShareTopicPartitionsV1 {
15    pub topic_id: [u8; 16],
16    pub partitions: Vec<i32>,
17}
18
19impl ShareTopicPartitionsV1 {
20    #[cfg(test)]
21    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
22        encoder.write_uuid(&self.topic_id);
23        encoder.write_compact_array(Some(&self.partitions), |encoder, partition| {
24            encoder.write_i32(*partition);
25            Ok(())
26        })?;
27        encoder.write_empty_tagged_fields();
28        Ok(())
29    }
30
31    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
32        let topic_id = decoder.read_uuid()?;
33        let partitions = decoder
34            .read_compact_array("share topic partitions", |decoder| decoder.read_i32())?
35            .unwrap_or_default();
36        decoder.read_tagged_fields()?;
37        Ok(Self {
38            topic_id,
39            partitions,
40        })
41    }
42}
43
44/// ShareGroupHeartbeat v1 request.
45///
46/// Version 1 is the stable KIP-932 wire shape. Kafka 4.0's early-access v0 is
47/// intentionally not exposed here because it was removed in Kafka 4.1.
48#[derive(Debug, Clone, PartialEq, Eq)]
49pub struct ShareGroupHeartbeatRequestV1 {
50    pub correlation_id: i32,
51    pub client_id: Option<String>,
52    pub group_id: String,
53    pub member_id: String,
54    pub member_epoch: i32,
55    pub rack_id: Option<String>,
56    pub subscribed_topic_names: Option<Vec<String>>,
57}
58
59impl ShareGroupHeartbeatRequestV1 {
60    /// Encodes the flexible request, including its request header.
61    pub fn encode(&self) -> Result<Vec<u8>> {
62        let mut encoder = Encoder::new();
63        RequestHeader {
64            api_key: SHARE_GROUP_HEARTBEAT_API_KEY,
65            api_version: 1,
66            correlation_id: self.correlation_id,
67            client_id: self.client_id.clone(),
68        }
69        .encode_v2(&mut encoder)?;
70        encoder.write_compact_string(&self.group_id)?;
71        encoder.write_compact_string(&self.member_id)?;
72        encoder.write_i32(self.member_epoch);
73        encoder.write_compact_nullable_string(self.rack_id.as_deref())?;
74        encoder.write_compact_array(self.subscribed_topic_names.as_deref(), |encoder, topic| {
75            encoder.write_compact_string(topic)
76        })?;
77        encoder.write_empty_tagged_fields();
78        Ok(encoder.into_bytes())
79    }
80}
81
82/// Assignment returned by ShareGroupHeartbeat.
83#[derive(Debug, Clone, PartialEq, Eq)]
84pub struct ShareGroupHeartbeatAssignmentV1 {
85    pub topic_partitions: Vec<ShareTopicPartitionsV1>,
86}
87
88impl ShareGroupHeartbeatAssignmentV1 {
89    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
90        let topic_partitions = decoder
91            .read_compact_array("share heartbeat assignment", ShareTopicPartitionsV1::decode)?
92            .unwrap_or_default();
93        decoder.read_tagged_fields()?;
94        Ok(Self { topic_partitions })
95    }
96}
97
98/// ShareGroupHeartbeat v1 response.
99#[derive(Debug, Clone, PartialEq, Eq)]
100pub struct ShareGroupHeartbeatResponseV1 {
101    pub throttle_time_ms: i32,
102    pub error_code: i16,
103    pub error_message: Option<String>,
104    pub member_id: Option<String>,
105    pub member_epoch: i32,
106    pub heartbeat_interval_ms: i32,
107    pub assignment: Option<ShareGroupHeartbeatAssignmentV1>,
108}
109
110impl ShareGroupHeartbeatResponseV1 {
111    /// Decodes the flexible response body after the response header.
112    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
113        let throttle_time_ms = decoder.read_i32()?;
114        let error_code = decoder.read_i16()?;
115        let error_message = decoder.read_compact_nullable_string()?;
116        let member_id = decoder.read_compact_nullable_string()?;
117        let member_epoch = decoder.read_i32()?;
118        let heartbeat_interval_ms = decoder.read_i32()?;
119        let assignment = match decoder.read_i8()? {
120            -1 => None,
121            1 => Some(ShareGroupHeartbeatAssignmentV1::decode(decoder)?),
122            marker => return Err(Error::InvalidNullableStruct(marker)),
123        };
124        decoder.read_tagged_fields()?;
125        Ok(Self {
126            throttle_time_ms,
127            error_code,
128            error_message,
129            member_id,
130            member_epoch,
131            heartbeat_interval_ms,
132            assignment,
133        })
134    }
135}
136
137/// One batch of records acknowledged by a share consumer.
138#[derive(Debug, Clone, PartialEq, Eq)]
139pub struct ShareAcknowledgementBatchV1 {
140    pub first_offset: i64,
141    pub last_offset: i64,
142    /// Kafka acknowledgement types: 0 gap, 1 accept, 2 release, 3 reject.
143    pub acknowledgement_types: Vec<i8>,
144}
145
146impl ShareAcknowledgementBatchV1 {
147    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
148        encoder.write_i64(self.first_offset);
149        encoder.write_i64(self.last_offset);
150        encoder.write_compact_array(
151            Some(&self.acknowledgement_types),
152            |encoder, acknowledgement_type| {
153                encoder.write_i8(*acknowledgement_type);
154                Ok(())
155            },
156        )?;
157        encoder.write_empty_tagged_fields();
158        Ok(())
159    }
160
161    #[cfg(test)]
162    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
163        let first_offset = decoder.read_i64()?;
164        let last_offset = decoder.read_i64()?;
165        let acknowledgement_types = decoder
166            .read_compact_array("share acknowledgement types", |decoder| decoder.read_i8())?
167            .unwrap_or_default();
168        decoder.read_tagged_fields()?;
169        Ok(Self {
170            first_offset,
171            last_offset,
172            acknowledgement_types,
173        })
174    }
175}
176
177/// Acknowledgement batches for one share-group partition.
178#[derive(Debug, Clone, PartialEq, Eq)]
179pub struct ShareFetchPartitionV1 {
180    pub partition_index: i32,
181    pub acknowledgement_batches: Vec<ShareAcknowledgementBatchV1>,
182}
183
184impl ShareFetchPartitionV1 {
185    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
186        encoder.write_i32(self.partition_index);
187        encoder.write_compact_array(Some(&self.acknowledgement_batches), |encoder, batch| {
188            batch.encode(encoder)
189        })?;
190        encoder.write_empty_tagged_fields();
191        Ok(())
192    }
193}
194
195/// One topic in a ShareFetch request.
196#[derive(Debug, Clone, PartialEq, Eq)]
197pub struct ShareFetchTopicV1 {
198    pub topic_id: [u8; 16],
199    pub partitions: Vec<ShareFetchPartitionV1>,
200}
201
202impl ShareFetchTopicV1 {
203    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
204        encoder.write_uuid(&self.topic_id);
205        encoder.write_compact_array(Some(&self.partitions), |encoder, partition| {
206            partition.encode(encoder)
207        })?;
208        encoder.write_empty_tagged_fields();
209        Ok(())
210    }
211}
212
213/// A topic and partitions to remove from a share fetch session.
214#[derive(Debug, Clone, PartialEq, Eq)]
215pub struct ShareForgottenTopicV1 {
216    pub topic_id: [u8; 16],
217    pub partitions: Vec<i32>,
218}
219
220impl ShareForgottenTopicV1 {
221    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
222        encoder.write_uuid(&self.topic_id);
223        encoder.write_compact_array(Some(&self.partitions), |encoder, partition| {
224            encoder.write_i32(*partition);
225            Ok(())
226        })?;
227        encoder.write_empty_tagged_fields();
228        Ok(())
229    }
230}
231
232/// ShareFetch v1 request.
233#[derive(Debug, Clone, PartialEq, Eq)]
234pub struct ShareFetchRequestV1 {
235    pub correlation_id: i32,
236    pub client_id: Option<String>,
237    pub group_id: Option<String>,
238    pub member_id: Option<String>,
239    pub share_session_epoch: i32,
240    pub max_wait_ms: i32,
241    pub min_bytes: i32,
242    pub max_bytes: i32,
243    pub max_records: i32,
244    pub batch_size: i32,
245    pub topics: Vec<ShareFetchTopicV1>,
246    pub forgotten_topics: Vec<ShareForgottenTopicV1>,
247}
248
249impl ShareFetchRequestV1 {
250    /// Encodes the flexible request, including its request header.
251    pub fn encode(&self) -> Result<Vec<u8>> {
252        encode_share_fetch_request(
253            1,
254            self.correlation_id,
255            self.client_id.clone(),
256            self.group_id.clone(),
257            self.member_id.clone(),
258            self.share_session_epoch,
259            self.max_wait_ms,
260            self.min_bytes,
261            self.max_bytes,
262            self.max_records,
263            self.batch_size,
264            None,
265            false,
266            &self.topics,
267            &self.forgotten_topics,
268        )
269    }
270}
271
272/// ShareFetch v2 request with KIP-1206 acquisition and KIP-1222 renewal fields.
273///
274/// The response schema remains the ShareFetch v1 response shape.
275#[derive(Debug, Clone, PartialEq, Eq)]
276pub struct ShareFetchRequestV2 {
277    pub correlation_id: i32,
278    pub client_id: Option<String>,
279    pub group_id: Option<String>,
280    pub member_id: Option<String>,
281    pub share_session_epoch: i32,
282    pub max_wait_ms: i32,
283    pub min_bytes: i32,
284    pub max_bytes: i32,
285    pub max_records: i32,
286    pub batch_size: i32,
287    /// KIP-1206 mode: `0` for batch-optimized or `1` for record-limit.
288    pub share_acquire_mode: i8,
289    /// KIP-1222 renew marker. KIP-1206 callers must send `false`.
290    pub is_renew_ack: bool,
291    pub topics: Vec<ShareFetchTopicV1>,
292    pub forgotten_topics: Vec<ShareForgottenTopicV1>,
293}
294
295impl ShareFetchRequestV2 {
296    /// Encodes the flexible request, including its request header.
297    pub fn encode(&self) -> Result<Vec<u8>> {
298        encode_share_fetch_request(
299            2,
300            self.correlation_id,
301            self.client_id.clone(),
302            self.group_id.clone(),
303            self.member_id.clone(),
304            self.share_session_epoch,
305            self.max_wait_ms,
306            self.min_bytes,
307            self.max_bytes,
308            self.max_records,
309            self.batch_size,
310            Some(self.share_acquire_mode),
311            self.is_renew_ack,
312            &self.topics,
313            &self.forgotten_topics,
314        )
315    }
316}
317
318#[allow(clippy::too_many_arguments)]
319fn encode_share_fetch_request(
320    api_version: i16,
321    correlation_id: i32,
322    client_id: Option<String>,
323    group_id: Option<String>,
324    member_id: Option<String>,
325    share_session_epoch: i32,
326    max_wait_ms: i32,
327    min_bytes: i32,
328    max_bytes: i32,
329    max_records: i32,
330    batch_size: i32,
331    share_acquire_mode: Option<i8>,
332    is_renew_ack: bool,
333    topics: &[ShareFetchTopicV1],
334    forgotten_topics: &[ShareForgottenTopicV1],
335) -> Result<Vec<u8>> {
336    let mut encoder = Encoder::new();
337    RequestHeader {
338        api_key: SHARE_FETCH_API_KEY,
339        api_version,
340        correlation_id,
341        client_id,
342    }
343    .encode_v2(&mut encoder)?;
344    encoder.write_compact_nullable_string(group_id.as_deref())?;
345    encoder.write_compact_nullable_string(member_id.as_deref())?;
346    encoder.write_i32(share_session_epoch);
347    encoder.write_i32(max_wait_ms);
348    encoder.write_i32(min_bytes);
349    encoder.write_i32(max_bytes);
350    encoder.write_i32(max_records);
351    encoder.write_i32(batch_size);
352    if let Some(share_acquire_mode) = share_acquire_mode {
353        encoder.write_i8(share_acquire_mode);
354        encoder.write_bool(is_renew_ack);
355    }
356    encoder.write_compact_array(Some(topics), |encoder, topic| topic.encode(encoder))?;
357    encoder.write_compact_array(Some(forgotten_topics), |encoder, topic| {
358        topic.encode(encoder)
359    })?;
360    encoder.write_empty_tagged_fields();
361    Ok(encoder.into_bytes())
362}
363
364/// Current leader information returned for a share partition.
365#[derive(Debug, Clone, PartialEq, Eq)]
366pub struct ShareLeaderIdAndEpochV1 {
367    pub leader_id: i32,
368    pub leader_epoch: i32,
369}
370
371impl ShareLeaderIdAndEpochV1 {
372    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
373        let leader_id = decoder.read_i32()?;
374        let leader_epoch = decoder.read_i32()?;
375        decoder.read_tagged_fields()?;
376        Ok(Self {
377            leader_id,
378            leader_epoch,
379        })
380    }
381}
382
383/// Acquired record range returned by ShareFetch.
384#[derive(Debug, Clone, PartialEq, Eq)]
385pub struct ShareAcquiredRecordsV1 {
386    pub first_offset: i64,
387    pub last_offset: i64,
388    pub delivery_count: i16,
389}
390
391impl ShareAcquiredRecordsV1 {
392    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
393        let first_offset = decoder.read_i64()?;
394        let last_offset = decoder.read_i64()?;
395        let delivery_count = decoder.read_i16()?;
396        decoder.read_tagged_fields()?;
397        Ok(Self {
398            first_offset,
399            last_offset,
400            delivery_count,
401        })
402    }
403}
404
405/// One partition returned by ShareFetch.
406#[derive(Debug, Clone, PartialEq, Eq)]
407pub struct ShareFetchPartitionResponseV1 {
408    pub partition_index: i32,
409    pub error_code: i16,
410    pub error_message: Option<String>,
411    pub acknowledgement_error_code: i16,
412    pub acknowledgement_error_message: Option<String>,
413    pub current_leader: ShareLeaderIdAndEpochV1,
414    /// Raw Kafka record bytes. The value is nullable when no records were fetched.
415    pub records: Option<Vec<u8>>,
416    pub acquired_records: Vec<ShareAcquiredRecordsV1>,
417}
418
419impl ShareFetchPartitionResponseV1 {
420    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
421        let partition_index = decoder.read_i32()?;
422        let error_code = decoder.read_i16()?;
423        let error_message = decoder.read_compact_nullable_string()?;
424        let acknowledgement_error_code = decoder.read_i16()?;
425        let acknowledgement_error_message = decoder.read_compact_nullable_string()?;
426        let current_leader = ShareLeaderIdAndEpochV1::decode(decoder)?;
427        let records = decoder.read_compact_nullable_bytes()?;
428        let acquired_records = decoder
429            .read_compact_array("share acquired records", ShareAcquiredRecordsV1::decode)?
430            .unwrap_or_default();
431        decoder.read_tagged_fields()?;
432        Ok(Self {
433            partition_index,
434            error_code,
435            error_message,
436            acknowledgement_error_code,
437            acknowledgement_error_message,
438            current_leader,
439            records,
440            acquired_records,
441        })
442    }
443}
444
445/// One topic returned by ShareFetch.
446#[derive(Debug, Clone, PartialEq, Eq)]
447pub struct ShareFetchTopicResponseV1 {
448    pub topic_id: [u8; 16],
449    pub partitions: Vec<ShareFetchPartitionResponseV1>,
450}
451
452impl ShareFetchTopicResponseV1 {
453    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
454        let topic_id = decoder.read_uuid()?;
455        let partitions = decoder
456            .read_compact_array(
457                "share fetch partitions",
458                ShareFetchPartitionResponseV1::decode,
459            )?
460            .unwrap_or_default();
461        decoder.read_tagged_fields()?;
462        Ok(Self {
463            topic_id,
464            partitions,
465        })
466    }
467}
468
469/// Broker endpoint returned when a share partition leader changes.
470#[derive(Debug, Clone, PartialEq, Eq)]
471pub struct ShareNodeEndpointV1 {
472    pub node_id: i32,
473    pub host: String,
474    pub port: i32,
475    pub rack: Option<String>,
476}
477
478impl ShareNodeEndpointV1 {
479    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
480        let node_id = decoder.read_i32()?;
481        let host = decoder.read_compact_string()?;
482        let port = decoder.read_i32()?;
483        let rack = decoder.read_compact_nullable_string()?;
484        decoder.read_tagged_fields()?;
485        Ok(Self {
486            node_id,
487            host,
488            port,
489            rack,
490        })
491    }
492}
493
494/// ShareFetch v1 response.
495#[derive(Debug, Clone, PartialEq, Eq)]
496pub struct ShareFetchResponseV1 {
497    pub throttle_time_ms: i32,
498    pub error_code: i16,
499    pub error_message: Option<String>,
500    pub acquisition_lock_timeout_ms: i32,
501    pub responses: Vec<ShareFetchTopicResponseV1>,
502    pub node_endpoints: Vec<ShareNodeEndpointV1>,
503}
504
505impl ShareFetchResponseV1 {
506    /// Decodes the flexible response body after the response header.
507    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
508        let throttle_time_ms = decoder.read_i32()?;
509        let error_code = decoder.read_i16()?;
510        let error_message = decoder.read_compact_nullable_string()?;
511        let acquisition_lock_timeout_ms = decoder.read_i32()?;
512        let responses = decoder
513            .read_compact_array("share fetch responses", ShareFetchTopicResponseV1::decode)?
514            .unwrap_or_default();
515        let node_endpoints = decoder
516            .read_compact_array("share node endpoints", ShareNodeEndpointV1::decode)?
517            .unwrap_or_default();
518        decoder.read_tagged_fields()?;
519        Ok(Self {
520            throttle_time_ms,
521            error_code,
522            error_message,
523            acquisition_lock_timeout_ms,
524            responses,
525            node_endpoints,
526        })
527    }
528}
529
530/// One topic and its acknowledgement batches in a ShareAcknowledge request.
531#[derive(Debug, Clone, PartialEq, Eq)]
532pub struct ShareAcknowledgeTopicV1 {
533    pub topic_id: [u8; 16],
534    pub partitions: Vec<ShareAcknowledgePartitionV1>,
535}
536
537impl ShareAcknowledgeTopicV1 {
538    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
539        encoder.write_uuid(&self.topic_id);
540        encoder.write_compact_array(Some(&self.partitions), |encoder, partition| {
541            partition.encode(encoder)
542        })?;
543        encoder.write_empty_tagged_fields();
544        Ok(())
545    }
546}
547
548/// One partition's acknowledgement batches.
549#[derive(Debug, Clone, PartialEq, Eq)]
550pub struct ShareAcknowledgePartitionV1 {
551    pub partition_index: i32,
552    pub acknowledgement_batches: Vec<ShareAcknowledgementBatchV1>,
553}
554
555impl ShareAcknowledgePartitionV1 {
556    fn encode(&self, encoder: &mut Encoder) -> Result<()> {
557        encoder.write_i32(self.partition_index);
558        encoder.write_compact_array(Some(&self.acknowledgement_batches), |encoder, batch| {
559            batch.encode(encoder)
560        })?;
561        encoder.write_empty_tagged_fields();
562        Ok(())
563    }
564}
565
566/// ShareAcknowledge v1 request.
567#[derive(Debug, Clone, PartialEq, Eq)]
568pub struct ShareAcknowledgeRequestV1 {
569    pub correlation_id: i32,
570    pub client_id: Option<String>,
571    pub group_id: Option<String>,
572    pub member_id: Option<String>,
573    pub share_session_epoch: i32,
574    pub topics: Vec<ShareAcknowledgeTopicV1>,
575}
576
577impl ShareAcknowledgeRequestV1 {
578    /// Encodes the flexible request, including its request header.
579    pub fn encode(&self) -> Result<Vec<u8>> {
580        encode_share_acknowledge_request(ShareAcknowledgeRequestParts {
581            api_version: 1,
582            correlation_id: self.correlation_id,
583            client_id: self.client_id.as_deref(),
584            group_id: self.group_id.as_deref(),
585            member_id: self.member_id.as_deref(),
586            share_session_epoch: self.share_session_epoch,
587            is_renew_ack: false,
588            topics: &self.topics,
589        })
590    }
591}
592
593/// ShareAcknowledge v2 request with KIP-1222 renewal support.
594#[derive(Debug, Clone, PartialEq, Eq)]
595pub struct ShareAcknowledgeRequestV2 {
596    pub correlation_id: i32,
597    pub client_id: Option<String>,
598    pub group_id: Option<String>,
599    pub member_id: Option<String>,
600    pub share_session_epoch: i32,
601    /// True when one or more acknowledgement batches contain `Renew` (4).
602    pub is_renew_ack: bool,
603    pub topics: Vec<ShareAcknowledgeTopicV1>,
604}
605
606impl ShareAcknowledgeRequestV2 {
607    /// Encodes the flexible request, including its request header.
608    pub fn encode(&self) -> Result<Vec<u8>> {
609        encode_share_acknowledge_request(ShareAcknowledgeRequestParts {
610            api_version: 2,
611            correlation_id: self.correlation_id,
612            client_id: self.client_id.as_deref(),
613            group_id: self.group_id.as_deref(),
614            member_id: self.member_id.as_deref(),
615            share_session_epoch: self.share_session_epoch,
616            is_renew_ack: self.is_renew_ack,
617            topics: &self.topics,
618        })
619    }
620}
621
622struct ShareAcknowledgeRequestParts<'a> {
623    api_version: i16,
624    correlation_id: i32,
625    client_id: Option<&'a str>,
626    group_id: Option<&'a str>,
627    member_id: Option<&'a str>,
628    share_session_epoch: i32,
629    is_renew_ack: bool,
630    topics: &'a [ShareAcknowledgeTopicV1],
631}
632
633fn encode_share_acknowledge_request(parts: ShareAcknowledgeRequestParts<'_>) -> Result<Vec<u8>> {
634    let mut encoder = Encoder::new();
635    RequestHeader {
636        api_key: SHARE_ACKNOWLEDGE_API_KEY,
637        api_version: parts.api_version,
638        correlation_id: parts.correlation_id,
639        client_id: parts.client_id.map(str::to_owned),
640    }
641    .encode_v2(&mut encoder)?;
642    encoder.write_compact_nullable_string(parts.group_id)?;
643    encoder.write_compact_nullable_string(parts.member_id)?;
644    encoder.write_i32(parts.share_session_epoch);
645    if parts.api_version >= 2 {
646        encoder.write_bool(parts.is_renew_ack);
647    }
648    encoder.write_compact_array(Some(parts.topics), |encoder, topic| topic.encode(encoder))?;
649    encoder.write_empty_tagged_fields();
650    Ok(encoder.into_bytes())
651}
652
653/// One partition result returned by ShareAcknowledge.
654#[derive(Debug, Clone, PartialEq, Eq)]
655pub struct ShareAcknowledgePartitionResponseV1 {
656    pub partition_index: i32,
657    pub error_code: i16,
658    pub error_message: Option<String>,
659    pub current_leader: ShareLeaderIdAndEpochV1,
660}
661
662impl ShareAcknowledgePartitionResponseV1 {
663    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
664        let partition_index = decoder.read_i32()?;
665        let error_code = decoder.read_i16()?;
666        let error_message = decoder.read_compact_nullable_string()?;
667        let current_leader = ShareLeaderIdAndEpochV1::decode(decoder)?;
668        decoder.read_tagged_fields()?;
669        Ok(Self {
670            partition_index,
671            error_code,
672            error_message,
673            current_leader,
674        })
675    }
676}
677
678/// One topic result returned by ShareAcknowledge.
679#[derive(Debug, Clone, PartialEq, Eq)]
680pub struct ShareAcknowledgeTopicResponseV1 {
681    pub topic_id: [u8; 16],
682    pub partitions: Vec<ShareAcknowledgePartitionResponseV1>,
683}
684
685impl ShareAcknowledgeTopicResponseV1 {
686    fn decode(decoder: &mut Decoder<'_>) -> Result<Self> {
687        let topic_id = decoder.read_uuid()?;
688        let partitions = decoder
689            .read_compact_array(
690                "share acknowledgement partitions",
691                ShareAcknowledgePartitionResponseV1::decode,
692            )?
693            .unwrap_or_default();
694        decoder.read_tagged_fields()?;
695        Ok(Self {
696            topic_id,
697            partitions,
698        })
699    }
700}
701
702/// ShareAcknowledge v1 response.
703#[derive(Debug, Clone, PartialEq, Eq)]
704pub struct ShareAcknowledgeResponseV1 {
705    pub throttle_time_ms: i32,
706    pub error_code: i16,
707    pub error_message: Option<String>,
708    pub responses: Vec<ShareAcknowledgeTopicResponseV1>,
709    pub node_endpoints: Vec<ShareNodeEndpointV1>,
710}
711
712impl ShareAcknowledgeResponseV1 {
713    /// Decodes the flexible response body after the response header.
714    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
715        let throttle_time_ms = decoder.read_i32()?;
716        let error_code = decoder.read_i16()?;
717        let error_message = decoder.read_compact_nullable_string()?;
718        let responses = decoder
719            .read_compact_array(
720                "share acknowledgement responses",
721                ShareAcknowledgeTopicResponseV1::decode,
722            )?
723            .unwrap_or_default();
724        let node_endpoints = decoder
725            .read_compact_array("share node endpoints", ShareNodeEndpointV1::decode)?
726            .unwrap_or_default();
727        decoder.read_tagged_fields()?;
728        Ok(Self {
729            throttle_time_ms,
730            error_code,
731            error_message,
732            responses,
733            node_endpoints,
734        })
735    }
736}
737
738/// ShareAcknowledge v2 response with the current acquisition lock timeout.
739#[derive(Debug, Clone, PartialEq, Eq)]
740pub struct ShareAcknowledgeResponseV2 {
741    pub throttle_time_ms: i32,
742    pub error_code: i16,
743    pub error_message: Option<String>,
744    pub acquisition_lock_timeout_ms: i32,
745    pub responses: Vec<ShareAcknowledgeTopicResponseV1>,
746    pub node_endpoints: Vec<ShareNodeEndpointV1>,
747}
748
749impl ShareAcknowledgeResponseV2 {
750    /// Decodes the flexible response body after the response header.
751    pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
752        let throttle_time_ms = decoder.read_i32()?;
753        let error_code = decoder.read_i16()?;
754        let error_message = decoder.read_compact_nullable_string()?;
755        let acquisition_lock_timeout_ms = decoder.read_i32()?;
756        let responses = decoder
757            .read_compact_array(
758                "share acknowledgement responses",
759                ShareAcknowledgeTopicResponseV1::decode,
760            )?
761            .unwrap_or_default();
762        let node_endpoints = decoder
763            .read_compact_array("share node endpoints", ShareNodeEndpointV1::decode)?
764            .unwrap_or_default();
765        decoder.read_tagged_fields()?;
766        Ok(Self {
767            throttle_time_ms,
768            error_code,
769            error_message,
770            acquisition_lock_timeout_ms,
771            responses,
772            node_endpoints,
773        })
774    }
775}
776
777#[cfg(test)]
778#[allow(clippy::unwrap_used)]
779mod tests {
780    use super::*;
781    use crate::codec::{Decoder, Encoder};
782
783    #[test]
784    fn encodes_share_group_heartbeat_v1_request() {
785        let request = ShareGroupHeartbeatRequestV1 {
786            correlation_id: 11,
787            client_id: Some("kafrust".to_owned()),
788            group_id: "share-orders".to_owned(),
789            member_id: "member-1".to_owned(),
790            member_epoch: 3,
791            rack_id: Some("rack-a".to_owned()),
792            subscribed_topic_names: Some(vec!["orders".to_owned()]),
793        };
794
795        let encoded = request.encode().unwrap();
796        let mut decoder = Decoder::new(&encoded);
797        assert_eq!(decoder.read_i16().unwrap(), SHARE_GROUP_HEARTBEAT_API_KEY);
798        assert_eq!(decoder.read_i16().unwrap(), 1);
799        assert_eq!(decoder.read_i32().unwrap(), 11);
800        assert_eq!(
801            decoder.read_nullable_string().unwrap(),
802            Some("kafrust".to_owned())
803        );
804        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
805        assert_eq!(decoder.read_compact_string().unwrap(), "share-orders");
806        assert_eq!(decoder.read_compact_string().unwrap(), "member-1");
807        assert_eq!(decoder.read_i32().unwrap(), 3);
808        assert_eq!(
809            decoder.read_compact_nullable_string().unwrap(),
810            Some("rack-a".to_owned())
811        );
812        assert_eq!(
813            decoder
814                .read_compact_array("topics", |decoder| decoder.read_compact_string())
815                .unwrap(),
816            Some(vec!["orders".to_owned()])
817        );
818        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
819        assert!(decoder.is_empty());
820    }
821
822    #[test]
823    fn decodes_share_group_heartbeat_v1_assignment() {
824        let mut bytes = Encoder::new();
825        bytes.write_i32(9);
826        bytes.write_i16(0);
827        bytes.write_compact_nullable_string(None).unwrap();
828        bytes
829            .write_compact_nullable_string(Some("member-1"))
830            .unwrap();
831        bytes.write_i32(4);
832        bytes.write_i32(2500);
833        bytes.write_i8(1);
834        bytes
835            .write_compact_array(
836                Some(&[ShareTopicPartitionsV1 {
837                    topic_id: [7; 16],
838                    partitions: vec![0, 2],
839                }]),
840                |encoder, assignment| assignment.encode(encoder),
841            )
842            .unwrap();
843        bytes.write_empty_tagged_fields();
844        bytes.write_empty_tagged_fields();
845
846        let encoded = bytes.into_bytes();
847        let mut decoder = Decoder::new(&encoded);
848        let response = ShareGroupHeartbeatResponseV1::decode_body(&mut decoder).unwrap();
849        assert_eq!(response.member_id.as_deref(), Some("member-1"));
850        assert_eq!(response.member_epoch, 4);
851        assert_eq!(
852            response.assignment.unwrap().topic_partitions[0].partitions,
853            vec![0, 2]
854        );
855        assert!(decoder.is_empty());
856    }
857
858    #[test]
859    fn encodes_share_fetch_v1_request_with_acknowledgements() {
860        let request = ShareFetchRequestV1 {
861            correlation_id: 22,
862            client_id: Some("kafrust".to_owned()),
863            group_id: Some("share-orders".to_owned()),
864            member_id: Some("member-1".to_owned()),
865            share_session_epoch: 2,
866            max_wait_ms: 500,
867            min_bytes: 1,
868            max_bytes: 1024,
869            max_records: 100,
870            batch_size: 10,
871            topics: vec![ShareFetchTopicV1 {
872                topic_id: [3; 16],
873                partitions: vec![ShareFetchPartitionV1 {
874                    partition_index: 0,
875                    acknowledgement_batches: vec![ShareAcknowledgementBatchV1 {
876                        first_offset: 10,
877                        last_offset: 12,
878                        acknowledgement_types: vec![1, 2, 3],
879                    }],
880                }],
881            }],
882            forgotten_topics: vec![ShareForgottenTopicV1 {
883                topic_id: [4; 16],
884                partitions: vec![1],
885            }],
886        };
887
888        let encoded = request.encode().unwrap();
889        assert_eq!(&encoded[0..4], &[0, 78, 0, 1]);
890        assert!(encoded.ends_with(&[0]));
891        let mut decoder = Decoder::new(&encoded[4..]);
892        assert_eq!(decoder.read_i32().unwrap(), 22);
893        assert_eq!(
894            decoder.read_nullable_string().unwrap(),
895            Some("kafrust".to_owned())
896        );
897        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
898        assert_eq!(
899            decoder.read_compact_nullable_string().unwrap(),
900            Some("share-orders".to_owned())
901        );
902        assert_eq!(
903            decoder.read_compact_nullable_string().unwrap(),
904            Some("member-1".to_owned())
905        );
906        assert_eq!(decoder.read_i32().unwrap(), 2);
907        assert_eq!(decoder.read_i32().unwrap(), 500);
908        assert_eq!(decoder.read_i32().unwrap(), 1);
909        assert_eq!(decoder.read_i32().unwrap(), 1024);
910        assert_eq!(decoder.read_i32().unwrap(), 100);
911        assert_eq!(decoder.read_i32().unwrap(), 10);
912        assert!(decoder
913            .read_compact_array("topics", |decoder| {
914                let topic_id = decoder.read_uuid()?;
915                let partitions = decoder.read_compact_array("partitions", |decoder| {
916                    let partition_index = decoder.read_i32()?;
917                    let batches = decoder
918                        .read_compact_array("batches", ShareAcknowledgementBatchV1::decode)?;
919                    decoder.read_tagged_fields()?;
920                    Ok((partition_index, batches))
921                })?;
922                decoder.read_tagged_fields()?;
923                Ok((topic_id, partitions))
924            })
925            .unwrap()
926            .is_some());
927    }
928
929    #[test]
930    fn encodes_share_fetch_v2_request_with_record_limit_mode() {
931        let request = ShareFetchRequestV2 {
932            correlation_id: 23,
933            client_id: Some("kafrust".to_owned()),
934            group_id: Some("share-orders".to_owned()),
935            member_id: Some("member-1".to_owned()),
936            share_session_epoch: 3,
937            max_wait_ms: 500,
938            min_bytes: 1,
939            max_bytes: 1024,
940            max_records: 6,
941            batch_size: 2,
942            share_acquire_mode: 1,
943            is_renew_ack: false,
944            topics: Vec::new(),
945            forgotten_topics: Vec::new(),
946        };
947
948        let encoded = request.encode().unwrap();
949        assert_eq!(&encoded[0..4], &[0, 78, 0, 2]);
950        let mut decoder = Decoder::new(&encoded[4..]);
951        assert_eq!(decoder.read_i32().unwrap(), 23);
952        assert_eq!(
953            decoder.read_nullable_string().unwrap(),
954            Some("kafrust".to_owned())
955        );
956        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
957        assert_eq!(
958            decoder.read_compact_nullable_string().unwrap(),
959            Some("share-orders".to_owned())
960        );
961        assert_eq!(
962            decoder.read_compact_nullable_string().unwrap(),
963            Some("member-1".to_owned())
964        );
965        assert_eq!(decoder.read_i32().unwrap(), 3);
966        assert_eq!(decoder.read_i32().unwrap(), 500);
967        assert_eq!(decoder.read_i32().unwrap(), 1);
968        assert_eq!(decoder.read_i32().unwrap(), 1024);
969        assert_eq!(decoder.read_i32().unwrap(), 6);
970        assert_eq!(decoder.read_i32().unwrap(), 2);
971        assert_eq!(decoder.read_i8().unwrap(), 1);
972        assert!(!decoder.read_bool().unwrap());
973        assert!(decoder
974            .read_compact_array("topics", |_decoder| Ok::<(), Error>(()))
975            .unwrap()
976            .is_some());
977        assert!(decoder
978            .read_compact_array("forgotten topics", |_decoder| Ok::<(), Error>(()))
979            .unwrap()
980            .is_some());
981        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
982        assert!(decoder.is_empty());
983    }
984
985    #[test]
986    fn encodes_share_fetch_v2_renew_request_with_zero_fetch_limits() {
987        let request = ShareFetchRequestV2 {
988            correlation_id: 24,
989            client_id: Some("kafrust".to_owned()),
990            group_id: Some("share-orders".to_owned()),
991            member_id: Some("member-1".to_owned()),
992            share_session_epoch: 7,
993            max_wait_ms: 0,
994            min_bytes: 0,
995            max_bytes: 0,
996            max_records: 0,
997            batch_size: 0,
998            share_acquire_mode: 0,
999            is_renew_ack: true,
1000            topics: Vec::new(),
1001            forgotten_topics: Vec::new(),
1002        };
1003
1004        let encoded = request.encode().unwrap();
1005        let mut decoder = Decoder::new(&encoded[4..]);
1006        assert_eq!(decoder.read_i32().unwrap(), 24);
1007        assert_eq!(
1008            decoder.read_nullable_string().unwrap(),
1009            Some("kafrust".to_owned())
1010        );
1011        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
1012        assert_eq!(
1013            decoder.read_compact_nullable_string().unwrap(),
1014            Some("share-orders".to_owned())
1015        );
1016        assert_eq!(
1017            decoder.read_compact_nullable_string().unwrap(),
1018            Some("member-1".to_owned())
1019        );
1020        assert_eq!(decoder.read_i32().unwrap(), 7);
1021        assert_eq!(decoder.read_i32().unwrap(), 0);
1022        assert_eq!(decoder.read_i32().unwrap(), 0);
1023        assert_eq!(decoder.read_i32().unwrap(), 0);
1024        assert_eq!(decoder.read_i32().unwrap(), 0);
1025        assert_eq!(decoder.read_i32().unwrap(), 0);
1026        assert_eq!(decoder.read_i8().unwrap(), 0);
1027        assert!(decoder.read_bool().unwrap());
1028        assert!(decoder
1029            .read_compact_array("topics", |_decoder| Ok::<(), Error>(()))
1030            .unwrap()
1031            .is_some());
1032        assert!(decoder
1033            .read_compact_array("forgotten topics", |_decoder| Ok::<(), Error>(()))
1034            .unwrap()
1035            .is_some());
1036        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
1037        assert!(decoder.is_empty());
1038    }
1039
1040    #[test]
1041    fn decodes_share_fetch_v1_response_and_preserves_record_bytes() {
1042        let mut bytes = Encoder::new();
1043        bytes.write_i32(7);
1044        bytes.write_i16(0);
1045        bytes.write_compact_nullable_string(None).unwrap();
1046        bytes.write_i32(30_000);
1047        bytes
1048            .write_compact_array(
1049                Some(&[ShareFetchTopicResponseV1 {
1050                    topic_id: [5; 16],
1051                    partitions: vec![ShareFetchPartitionResponseV1 {
1052                        partition_index: 0,
1053                        error_code: 0,
1054                        error_message: None,
1055                        acknowledgement_error_code: 0,
1056                        acknowledgement_error_message: None,
1057                        current_leader: ShareLeaderIdAndEpochV1 {
1058                            leader_id: 1,
1059                            leader_epoch: 8,
1060                        },
1061                        records: Some(vec![1, 2, 3]),
1062                        acquired_records: vec![ShareAcquiredRecordsV1 {
1063                            first_offset: 10,
1064                            last_offset: 12,
1065                            delivery_count: 1,
1066                        }],
1067                    }],
1068                }]),
1069                |encoder, topic| {
1070                    encoder.write_uuid(&topic.topic_id);
1071                    encoder.write_compact_array(
1072                        Some(&topic.partitions),
1073                        |encoder, partition| {
1074                            encoder.write_i32(partition.partition_index);
1075                            encoder.write_i16(partition.error_code);
1076                            encoder.write_compact_nullable_string(
1077                                partition.error_message.as_deref(),
1078                            )?;
1079                            encoder.write_i16(partition.acknowledgement_error_code);
1080                            encoder.write_compact_nullable_string(
1081                                partition.acknowledgement_error_message.as_deref(),
1082                            )?;
1083                            encoder.write_i32(partition.current_leader.leader_id);
1084                            encoder.write_i32(partition.current_leader.leader_epoch);
1085                            encoder.write_empty_tagged_fields();
1086                            encoder.write_compact_nullable_bytes(partition.records.as_deref())?;
1087                            encoder.write_compact_array(
1088                                Some(&partition.acquired_records),
1089                                |encoder, acquired| {
1090                                    encoder.write_i64(acquired.first_offset);
1091                                    encoder.write_i64(acquired.last_offset);
1092                                    encoder.write_i16(acquired.delivery_count);
1093                                    encoder.write_empty_tagged_fields();
1094                                    Ok(())
1095                                },
1096                            )?;
1097                            encoder.write_empty_tagged_fields();
1098                            Ok(())
1099                        },
1100                    )?;
1101                    encoder.write_empty_tagged_fields();
1102                    Ok(())
1103                },
1104            )
1105            .unwrap();
1106        bytes
1107            .write_compact_array::<ShareNodeEndpointV1>(Some(&[]), |_encoder, _| Ok(()))
1108            .unwrap();
1109        bytes.write_empty_tagged_fields();
1110
1111        let encoded = bytes.into_bytes();
1112        let mut decoder = Decoder::new(&encoded);
1113        let response = ShareFetchResponseV1::decode_body(&mut decoder).unwrap();
1114        let partition = &response.responses[0].partitions[0];
1115        assert_eq!(partition.records.as_deref(), Some(&[1, 2, 3][..]));
1116        assert_eq!(partition.acquired_records[0].delivery_count, 1);
1117        assert!(decoder.is_empty());
1118    }
1119
1120    #[test]
1121    fn encodes_share_acknowledge_v1_request() {
1122        let request = ShareAcknowledgeRequestV1 {
1123            correlation_id: 33,
1124            client_id: None,
1125            group_id: Some("share-orders".to_owned()),
1126            member_id: Some("member-1".to_owned()),
1127            share_session_epoch: 3,
1128            topics: vec![ShareAcknowledgeTopicV1 {
1129                topic_id: [9; 16],
1130                partitions: vec![ShareAcknowledgePartitionV1 {
1131                    partition_index: 2,
1132                    acknowledgement_batches: vec![ShareAcknowledgementBatchV1 {
1133                        first_offset: 20,
1134                        last_offset: 20,
1135                        acknowledgement_types: vec![1],
1136                    }],
1137                }],
1138            }],
1139        };
1140
1141        let encoded = request.encode().unwrap();
1142        assert_eq!(&encoded[0..4], &[0, 79, 0, 1]);
1143        assert!(encoded.ends_with(&[0]));
1144    }
1145
1146    #[test]
1147    fn encodes_share_acknowledge_v2_renew_request() {
1148        let request = ShareAcknowledgeRequestV2 {
1149            correlation_id: 34,
1150            client_id: None,
1151            group_id: Some("share-orders".to_owned()),
1152            member_id: Some("member-1".to_owned()),
1153            share_session_epoch: 3,
1154            is_renew_ack: true,
1155            topics: vec![ShareAcknowledgeTopicV1 {
1156                topic_id: [9; 16],
1157                partitions: vec![ShareAcknowledgePartitionV1 {
1158                    partition_index: 2,
1159                    acknowledgement_batches: vec![ShareAcknowledgementBatchV1 {
1160                        first_offset: 20,
1161                        last_offset: 20,
1162                        acknowledgement_types: vec![4],
1163                    }],
1164                }],
1165            }],
1166        };
1167
1168        let encoded = request.encode().unwrap();
1169        assert_eq!(&encoded[0..4], &[0, 79, 0, 2]);
1170        let mut decoder = Decoder::new(&encoded[4..]);
1171        assert_eq!(decoder.read_i32().unwrap(), 34);
1172        assert_eq!(decoder.read_nullable_string().unwrap(), None);
1173        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
1174        assert_eq!(
1175            decoder.read_compact_nullable_string().unwrap(),
1176            Some("share-orders".to_owned())
1177        );
1178        assert_eq!(
1179            decoder.read_compact_nullable_string().unwrap(),
1180            Some("member-1".to_owned())
1181        );
1182        assert_eq!(decoder.read_i32().unwrap(), 3);
1183        assert!(decoder.read_bool().unwrap());
1184        let topics = decoder
1185            .read_compact_array("topics", |decoder| {
1186                let topic_id = decoder.read_uuid()?;
1187                let partitions = decoder
1188                    .read_compact_array("partitions", |decoder| {
1189                        let partition_index = decoder.read_i32()?;
1190                        let batches = decoder
1191                            .read_compact_array("batches", ShareAcknowledgementBatchV1::decode)?;
1192                        decoder.read_tagged_fields()?;
1193                        Ok((partition_index, batches.unwrap_or_default()))
1194                    })?
1195                    .unwrap_or_default();
1196                decoder.read_tagged_fields()?;
1197                Ok((topic_id, partitions))
1198            })
1199            .unwrap()
1200            .unwrap();
1201        assert_eq!(topics[0].1[0].1[0].acknowledgement_types, vec![4]);
1202        assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
1203        assert!(decoder.is_empty());
1204    }
1205
1206    #[test]
1207    fn decodes_share_acknowledge_v1_response_with_leader_endpoint() {
1208        let mut bytes = Encoder::new();
1209        bytes.write_i32(4);
1210        bytes.write_i16(0);
1211        bytes.write_compact_nullable_string(None).unwrap();
1212        bytes
1213            .write_compact_array::<ShareAcknowledgeTopicResponseV1>(Some(&[]), |_encoder, _| Ok(()))
1214            .unwrap();
1215        bytes
1216            .write_compact_array(
1217                Some(&[ShareNodeEndpointV1 {
1218                    node_id: 2,
1219                    host: "broker".to_owned(),
1220                    port: 9092,
1221                    rack: None,
1222                }]),
1223                |encoder, endpoint| {
1224                    encoder.write_i32(endpoint.node_id);
1225                    encoder.write_compact_string(&endpoint.host)?;
1226                    encoder.write_i32(endpoint.port);
1227                    encoder.write_compact_nullable_string(endpoint.rack.as_deref())?;
1228                    encoder.write_empty_tagged_fields();
1229                    Ok(())
1230                },
1231            )
1232            .unwrap();
1233        bytes.write_empty_tagged_fields();
1234
1235        let encoded = bytes.into_bytes();
1236        let mut decoder = Decoder::new(&encoded);
1237        let response = ShareAcknowledgeResponseV1::decode_body(&mut decoder).unwrap();
1238        assert_eq!(response.node_endpoints[0].host, "broker");
1239        assert_eq!(response.node_endpoints[0].port, 9092);
1240        assert!(decoder.is_empty());
1241    }
1242
1243    #[test]
1244    fn decodes_share_acknowledge_v2_response_with_lock_timeout() {
1245        let mut bytes = Encoder::new();
1246        bytes.write_i32(4);
1247        bytes.write_i16(0);
1248        bytes.write_compact_nullable_string(None).unwrap();
1249        bytes.write_i32(45_000);
1250        bytes
1251            .write_compact_array::<ShareAcknowledgeTopicResponseV1>(Some(&[]), |_encoder, _| Ok(()))
1252            .unwrap();
1253        bytes
1254            .write_compact_array::<ShareNodeEndpointV1>(Some(&[]), |_encoder, _| Ok(()))
1255            .unwrap();
1256        bytes.write_empty_tagged_fields();
1257
1258        let encoded = bytes.into_bytes();
1259        let mut decoder = Decoder::new(&encoded);
1260        let response = ShareAcknowledgeResponseV2::decode_body(&mut decoder).unwrap();
1261        assert_eq!(response.acquisition_lock_timeout_ms, 45_000);
1262        assert!(decoder.is_empty());
1263    }
1264}