Skip to main content

kafrust/
consumer.rs

1use std::collections::{BTreeMap, BTreeSet};
2
3use kafrust_protocol::api::fetch::{
4    AbortedTransactionV4, FetchPartitionResponseV11, FetchPartitionResponseV12,
5    FetchPartitionResponseV13, FetchPartitionResponseV4, FetchPartitionV12, FetchResponseV11,
6    FetchResponseV12, FetchResponseV13, FetchResponseV4, FetchTopicV13, MessageSetRecord,
7};
8use kafrust_protocol::api::list_offsets::{
9    ListOffsetsPartitionResponseV1, ListOffsetsPartitionV1, ListOffsetsTopicResponseV1,
10    ListOffsetsTopicV1, EARLIEST_TIMESTAMP, LATEST_TIMESTAMP,
11};
12use kafrust_protocol::api::metadata::{
13    BrokerMetadata, MetadataPartitionV12, MetadataRequestTopicV12, MetadataResponseV1,
14    MetadataResponseV12,
15};
16use kafrust_protocol::api::offset_for_leader_epoch::{
17    OffsetForLeaderEpochPartitionResponseV3, OffsetForLeaderEpochPartitionV3,
18    OffsetForLeaderEpochTopicResponseV3, OffsetForLeaderEpochTopicV3,
19};
20
21use crate::broker_client_cache::BrokerClientCache;
22use crate::client::{Client, FetchOneRequestV11, FetchOneRequestV12, FetchOneRequestV4};
23use crate::config::{ClientConfig, OAuthBearerTokenProvider, SecurityProtocol};
24use crate::error::{BrokerErrorKind, Error, Result};
25use crate::metrics::ClientMetrics;
26use tokio::sync::mpsc;
27use tokio::sync::mpsc::error::TrySendError;
28use tracing::debug;
29
30#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
31/// Controls whether a consumer can return records from aborted transactions.
32pub enum IsolationLevel {
33    /// Return all records, including records from aborted transactions.
34    #[default]
35    ReadUncommitted,
36    /// Return only committed records and hide Kafka transaction control records.
37    ReadCommitted,
38}
39
40impl IsolationLevel {
41    fn as_i8(self) -> i8 {
42        match self {
43            Self::ReadUncommitted => 0,
44            Self::ReadCommitted => 1,
45        }
46    }
47}
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
50/// Starting position used when a consumer has no usable offset.
51pub enum OffsetResetPolicy {
52    /// Start at the partition's earliest retained offset.
53    Earliest,
54    /// Start after the partition's current log end.
55    Latest,
56    /// Start at an explicit absolute offset.
57    Offset(i64),
58}
59
60impl Default for OffsetResetPolicy {
61    fn default() -> Self {
62        Self::Offset(0)
63    }
64}
65
66impl OffsetResetPolicy {
67    pub(crate) fn timestamp(self) -> Option<i64> {
68        match self {
69            Self::Earliest => Some(EARLIEST_TIMESTAMP),
70            Self::Latest => Some(LATEST_TIMESTAMP),
71            Self::Offset(_) => None,
72        }
73    }
74
75    fn is_recovery(self) -> bool {
76        matches!(self, Self::Earliest | Self::Latest)
77    }
78}
79
80#[derive(Debug, Clone, PartialEq, Eq)]
81/// Record fetched from a Kafka topic partition.
82pub struct ConsumerRecord {
83    topic: String,
84    partition: i32,
85    offset: i64,
86    leader_epoch: i32,
87    timestamp_ms: i64,
88    key: Option<Vec<u8>>,
89    value: Option<Vec<u8>>,
90    headers: Vec<ConsumerRecordHeader>,
91}
92
93#[derive(Debug, Clone, PartialEq, Eq)]
94/// Header attached to a fetched Kafka record.
95///
96/// Kafka permits a header value to be null. Use [`Self::value`] to preserve
97/// that distinction instead of treating a null value as an empty byte slice.
98pub struct ConsumerRecordHeader {
99    key: String,
100    value: Option<Vec<u8>>,
101}
102
103impl ConsumerRecordHeader {
104    fn from_protocol(header: kafrust_protocol::api::produce::RecordBatchHeader) -> Self {
105        Self {
106            key: header.key,
107            value: header.value,
108        }
109    }
110
111    /// Returns the header key.
112    pub fn key(&self) -> &str {
113        &self.key
114    }
115
116    /// Returns the nullable header value bytes.
117    pub fn value(&self) -> Option<&[u8]> {
118        self.value.as_deref()
119    }
120}
121
122/// A bounded queue containing records fetched for one topic partition.
123///
124/// A queue is created with [`Consumer::split_partition_queue`]. The owning
125/// consumer continues to perform network polling, while this handle provides
126/// an independent receive path for the selected partition. Dropping the
127/// handle closes the split and causes subsequent records to return through
128/// [`Consumer::poll`]. Records already buffered in the queue remain readable
129/// until it is drained.
130#[derive(Debug)]
131pub struct ConsumerPartitionQueue {
132    topic: String,
133    partition: i32,
134    receiver: mpsc::Receiver<ConsumerRecord>,
135}
136
137impl ConsumerPartitionQueue {
138    /// Returns the topic served by this queue.
139    pub fn topic(&self) -> &str {
140        &self.topic
141    }
142
143    /// Returns the partition served by this queue.
144    pub fn partition(&self) -> i32 {
145        self.partition
146    }
147
148    /// Waits for the next record, or returns `None` after the queue is closed
149    /// and drained.
150    pub async fn recv(&mut self) -> Option<ConsumerRecord> {
151        self.receiver.recv().await
152    }
153
154    /// Returns one immediately available record.
155    ///
156    /// `None` means that the queue is currently empty or has been closed and
157    /// drained. Use [`Self::recv`] when the distinction matters.
158    pub fn try_recv(&mut self) -> Option<ConsumerRecord> {
159        self.receiver.try_recv().ok()
160    }
161
162    /// Receives up to `max_records`, waiting for the first record when the
163    /// queue is not already closed. A zero limit returns an empty vector.
164    pub async fn recv_batch(&mut self, max_records: usize) -> Vec<ConsumerRecord> {
165        if max_records == 0 {
166            return Vec::new();
167        }
168
169        let Some(first) = self.recv().await else {
170            return Vec::new();
171        };
172        let mut records = Vec::with_capacity(max_records.min(16));
173        records.push(first);
174        while records.len() < max_records {
175            let Some(record) = self.try_recv() else {
176                break;
177            };
178            records.push(record);
179        }
180        records
181    }
182}
183
184impl ConsumerRecord {
185    pub(crate) fn from_message_set(topic: &str, partition: i32, record: MessageSetRecord) -> Self {
186        Self {
187            topic: topic.to_owned(),
188            partition,
189            offset: record.offset,
190            leader_epoch: record.leader_epoch,
191            timestamp_ms: record.timestamp_ms,
192            key: record.key,
193            value: record.value,
194            headers: record
195                .headers
196                .into_iter()
197                .map(ConsumerRecordHeader::from_protocol)
198                .collect(),
199        }
200    }
201
202    /// Returns the Kafka topic name.
203    pub fn topic(&self) -> &str {
204        &self.topic
205    }
206
207    /// Returns the Kafka partition index.
208    pub fn partition(&self) -> i32 {
209        self.partition
210    }
211
212    /// Returns the Kafka record offset.
213    pub fn offset(&self) -> i64 {
214        self.offset
215    }
216
217    /// Returns the Kafka partition leader epoch attached to this record.
218    ///
219    /// Legacy MessageSet records return `-1` because that wire format has no
220    /// leader epoch field.
221    pub fn leader_epoch(&self) -> i32 {
222        self.leader_epoch
223    }
224
225    /// Returns the Kafka record timestamp in milliseconds since the Unix epoch.
226    pub fn timestamp_ms(&self) -> i64 {
227        self.timestamp_ms
228    }
229
230    /// Returns the record key bytes.
231    pub fn key(&self) -> Option<&[u8]> {
232        self.key.as_deref()
233    }
234
235    /// Returns the record value bytes.
236    pub fn value(&self) -> Option<&[u8]> {
237        self.value.as_deref()
238    }
239
240    /// Returns the headers attached to this record in wire order.
241    pub fn headers(&self) -> &[ConsumerRecordHeader] {
242        &self.headers
243    }
244}
245
246#[derive(Debug, Clone, Copy, PartialEq, Eq)]
247/// Earliest and latest available offsets for one Kafka topic partition.
248pub struct PartitionWatermarks {
249    low: i64,
250    high: i64,
251}
252
253#[derive(Debug, Clone, Copy, PartialEq, Eq)]
254/// The end offset recorded for a requested Kafka partition leader epoch.
255pub struct LeaderEpochOffset {
256    leader_epoch: i32,
257    end_offset: i64,
258}
259
260impl LeaderEpochOffset {
261    /// Returns the leader epoch reported by Kafka.
262    pub fn leader_epoch(&self) -> i32 {
263        self.leader_epoch
264    }
265
266    /// Returns the first offset after the requested epoch's log range.
267    pub fn end_offset(&self) -> i64 {
268        self.end_offset
269    }
270}
271
272impl PartitionWatermarks {
273    /// Returns the earliest available offset.
274    pub fn low(&self) -> i64 {
275        self.low
276    }
277
278    /// Returns the latest offset, which is the next offset after the log end.
279    pub fn high(&self) -> i64 {
280        self.high
281    }
282}
283
284#[derive(Debug)]
285/// Direct Kafka consumer for manually assigned topic partitions.
286pub struct Consumer {
287    client: Client,
288    config: ConsumerConfig,
289    assignments: Vec<ConsumerAssignment>,
290    poll_cursor: usize,
291    partition_queues: BTreeMap<(String, i32), mpsc::Sender<ConsumerRecord>>,
292    metadata_cache: BTreeMap<String, MetadataResponseV1>,
293    fetch_topic_ids: BTreeMap<String, [u8; 16]>,
294    broker_clients: BrokerClientCache,
295    fetch_sessions: BTreeMap<String, FetchSessionState>,
296    preferred_read_replicas: BTreeMap<(String, i32), i32>,
297}
298
299#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
300struct FetchSessionState {
301    session_id: i32,
302    next_epoch: i32,
303}
304
305impl FetchSessionState {
306    fn next_request(self) -> (i32, i32) {
307        (self.session_id, self.next_epoch)
308    }
309
310    fn advance_with_response(self, session_id: i32) -> Self {
311        Self {
312            session_id,
313            next_epoch: self.next_epoch.saturating_add(1),
314        }
315    }
316}
317
318impl Consumer {
319    pub(crate) fn from_assignments(
320        client: Client,
321        config: ConsumerConfig,
322        assignments: Vec<ConsumerAssignment>,
323    ) -> Self {
324        Self {
325            client,
326            config,
327            assignments,
328            poll_cursor: 0,
329            partition_queues: BTreeMap::new(),
330            metadata_cache: BTreeMap::new(),
331            fetch_topic_ids: BTreeMap::new(),
332            broker_clients: BrokerClientCache::default(),
333            fetch_sessions: BTreeMap::new(),
334            preferred_read_replicas: BTreeMap::new(),
335        }
336    }
337
338    pub(crate) fn replace_assignments(&mut self, mut assignments: Vec<ConsumerAssignment>) {
339        let previous_assignments = self.assignments.clone();
340        let previous_positions = previous_assignments
341            .iter()
342            .map(|assignment| {
343                (
344                    (assignment.topic.clone(), assignment.partition),
345                    assignment.next_offset,
346                )
347            })
348            .collect::<BTreeMap<_, _>>();
349        let paused = self
350            .assignments
351            .iter()
352            .filter(|assignment| assignment.paused)
353            .map(|assignment| (assignment.topic.clone(), assignment.partition))
354            .collect::<BTreeSet<_>>();
355        for assignment in &mut assignments {
356            assignment.paused = paused.contains(&(assignment.topic.clone(), assignment.partition));
357        }
358        assignments.sort_by(|left, right| {
359            left.topic
360                .cmp(&right.topic)
361                .then_with(|| left.partition.cmp(&right.partition))
362        });
363        self.partition_queues.retain(|(topic, partition), _| {
364            assignments.iter().any(|assignment| {
365                assignment.topic == *topic
366                    && assignment.partition == *partition
367                    && previous_positions.get(&(topic.clone(), *partition))
368                        == Some(&assignment.next_offset)
369            })
370        });
371        self.assignments = assignments;
372        self.poll_cursor %= self.assignments.len().max(1);
373        self.restore_assignment_state(&previous_assignments);
374        self.metadata_cache.clear();
375        self.fetch_topic_ids.clear();
376        self.fetch_sessions.clear();
377    }
378
379    pub(crate) fn restore_assignment_state(&mut self, previous_assignments: &[ConsumerAssignment]) {
380        let previous_state = previous_assignments
381            .iter()
382            .map(|assignment| {
383                (
384                    (assignment.topic.clone(), assignment.partition),
385                    (assignment.next_offset, assignment.leader_epoch),
386                )
387            })
388            .collect::<BTreeMap<_, _>>();
389        for assignment in &mut self.assignments {
390            if let Some((next_offset, leader_epoch)) =
391                previous_state.get(&(assignment.topic.clone(), assignment.partition))
392            {
393                assignment.next_offset = *next_offset;
394                assignment.leader_epoch = assignment.leader_epoch.max(*leader_epoch);
395            }
396        }
397    }
398
399    /// Assigns a topic partition and next offset to fetch.
400    pub fn assign(&mut self, topic: impl Into<String>, partition: i32, offset: i64) {
401        let topic = topic.into();
402        self.partition_queues.remove(&(topic.clone(), partition));
403        assign_partition(&mut self.assignments, topic, partition, offset);
404        self.poll_cursor %= self.assignments.len().max(1);
405        self.fetch_sessions.clear();
406    }
407
408    /// Returns the current topic partition assignments.
409    pub fn assignments(&self) -> &[ConsumerAssignment] {
410        &self.assignments
411    }
412
413    /// Splits one assigned topic partition into a bounded receive queue.
414    ///
415    /// Records fetched for the partition are delivered to the returned queue
416    /// instead of the vector returned by [`Self::poll`]. The queue capacity is
417    /// configured with [`ConsumerConfig::partition_queue_capacity`]. When the
418    /// queue is full, `poll` returns an error and does not advance beyond the
419    /// last record accepted by the queue.
420    pub fn split_partition_queue(
421        &mut self,
422        topic: impl Into<String>,
423        partition: i32,
424    ) -> Result<ConsumerPartitionQueue> {
425        let topic = topic.into();
426        if self.assignment(&topic, partition).is_none() {
427            return Err(Error::UnassignedTopicPartition { topic, partition });
428        }
429
430        let key = (topic.clone(), partition);
431        if let Some(sender) = self.partition_queues.get(&key) {
432            if !sender.is_closed() {
433                return Err(Error::Unsupported("partition queue is already split"));
434            }
435        }
436        self.partition_queues.remove(&key);
437
438        let (sender, receiver) = mpsc::channel(self.config.partition_queue_capacity);
439        self.partition_queues.insert(key, sender);
440        Ok(ConsumerPartitionQueue {
441            topic,
442            partition,
443            receiver,
444        })
445    }
446
447    /// Returns the next offset for an assigned topic partition.
448    pub fn position(&self, topic: &str, partition: i32) -> Option<i64> {
449        self.assignment(topic, partition)
450            .map(ConsumerAssignment::next_offset)
451    }
452
453    /// Changes the next offset for an assigned topic partition.
454    pub fn seek(&mut self, topic: &str, partition: i32, offset: i64) -> Result<()> {
455        if self
456            .partition_queues
457            .get(&(topic.to_owned(), partition))
458            .is_some_and(|sender| !sender.is_closed())
459        {
460            return Err(Error::Unsupported(
461                "seek requires dropping the active partition queue",
462            ));
463        }
464        self.assignment_mut(topic, partition)?.next_offset = offset;
465        self.fetch_sessions.clear();
466        Ok(())
467    }
468
469    /// Pauses fetching from an assigned topic partition.
470    pub fn pause(&mut self, topic: &str, partition: i32) -> Result<()> {
471        self.assignment_mut(topic, partition)?.paused = true;
472        self.fetch_sessions.clear();
473        Ok(())
474    }
475
476    /// Resumes fetching from an assigned topic partition.
477    pub fn resume(&mut self, topic: &str, partition: i32) -> Result<()> {
478        self.assignment_mut(topic, partition)?.paused = false;
479        self.fetch_sessions.clear();
480        Ok(())
481    }
482
483    /// Fetches the earliest and latest available offsets for a topic partition.
484    ///
485    /// The partition does not need to be assigned to this consumer. Kafka's
486    /// latest offset is the next offset after the current log end.
487    #[tracing::instrument(
488        level = "debug",
489        name = "kafka.consumer.fetch_watermarks",
490        skip_all,
491        fields(topic = tracing::field::Empty, partition),
492        err
493    )]
494    pub async fn fetch_watermarks(
495        &mut self,
496        topic: impl Into<String>,
497        partition: i32,
498    ) -> Result<PartitionWatermarks> {
499        let topic = topic.into();
500        tracing::Span::current().record("topic", topic.as_str());
501        let mut attempt = 0;
502
503        loop {
504            match self.fetch_watermarks_once(&topic, partition).await {
505                Err(error) if attempt < self.config.max_retries && can_retry_fetch(&error) => {
506                    invalidate_metadata_cache(&mut self.metadata_cache, &topic);
507                    self.config.client.record_retry();
508                    attempt += 1;
509                }
510                result => return result,
511            }
512        }
513    }
514
515    /// Resolves the end offset for a partition leader epoch.
516    ///
517    /// `current_leader_epoch` is the epoch from the consumer's current
518    /// metadata, or `-1` when it is unknown. `leader_epoch` is the epoch whose
519    /// end offset should be returned. The partition does not need to be
520    /// assigned to this consumer.
521    #[tracing::instrument(
522        level = "debug",
523        name = "kafka.consumer.offset_for_leader_epoch",
524        skip_all,
525        fields(topic = tracing::field::Empty, partition, current_leader_epoch, leader_epoch),
526        err
527    )]
528    pub async fn offset_for_leader_epoch(
529        &mut self,
530        topic: impl Into<String>,
531        partition: i32,
532        current_leader_epoch: i32,
533        leader_epoch: i32,
534    ) -> Result<LeaderEpochOffset> {
535        let topic = topic.into();
536        tracing::Span::current().record("topic", topic.as_str());
537        let mut attempt = 0;
538
539        loop {
540            match self
541                .offset_for_leader_epoch_once(&topic, partition, current_leader_epoch, leader_epoch)
542                .await
543            {
544                Err(error) if attempt < self.config.max_retries && can_retry_fetch(&error) => {
545                    invalidate_metadata_cache(&mut self.metadata_cache, &topic);
546                    self.config.client.record_retry();
547                    attempt += 1;
548                }
549                result => return result,
550            }
551        }
552    }
553
554    async fn offset_for_leader_epoch_once(
555        &mut self,
556        topic: &str,
557        partition: i32,
558        current_leader_epoch: i32,
559        leader_epoch: i32,
560    ) -> Result<LeaderEpochOffset> {
561        let metadata = self.metadata_for_topic(topic).await?;
562        let leader = leader_for(&metadata, topic, partition)?;
563        let broker_addr = broker_addr_for(&metadata, leader)?;
564        let mut leader_client = self.connect_or_reuse_broker(&broker_addr).await?;
565        let response = leader_client
566            .offset_for_leader_epoch_v3(vec![OffsetForLeaderEpochTopicV3 {
567                name: topic.to_owned(),
568                partitions: vec![OffsetForLeaderEpochPartitionV3 {
569                    partition_index: partition,
570                    current_leader_epoch,
571                    leader_epoch,
572                }],
573            }])
574            .await?;
575        let partition_response =
576            offset_for_leader_epoch_partition_response(&response.topics, topic, partition)?;
577        if partition_response.error_code != 0 {
578            return Err(self.config.client.broker_error(
579                partition_response.error_code,
580                format!("offset for leader epoch {topic}-{partition}"),
581            ));
582        }
583        self.cache_broker_client(broker_addr, leader_client);
584        Ok(LeaderEpochOffset {
585            leader_epoch: partition_response.leader_epoch,
586            end_offset: partition_response.end_offset,
587        })
588    }
589
590    async fn fetch_watermarks_once(
591        &mut self,
592        topic: &str,
593        partition: i32,
594    ) -> Result<PartitionWatermarks> {
595        let metadata = self.metadata_for_topic(topic).await?;
596        let leader = leader_for(&metadata, topic, partition)?;
597        let broker_addr = broker_addr_for(&metadata, leader)?;
598        let mut leader_client = self.connect_or_reuse_broker(&broker_addr).await?;
599        let low = self
600            .request_partition_offset(&mut leader_client, topic, partition, EARLIEST_TIMESTAMP)
601            .await?;
602        let high = self
603            .request_partition_offset(&mut leader_client, topic, partition, LATEST_TIMESTAMP)
604            .await?;
605        self.cache_broker_client(broker_addr, leader_client);
606        Ok(PartitionWatermarks { low, high })
607    }
608
609    async fn request_partition_offset(
610        &self,
611        client: &mut Client,
612        topic: &str,
613        partition: i32,
614        timestamp: i64,
615    ) -> Result<i64> {
616        let response = client
617            .list_offsets_v1(vec![ListOffsetsTopicV1 {
618                name: topic.to_owned(),
619                partitions: vec![ListOffsetsPartitionV1 {
620                    partition_index: partition,
621                    timestamp,
622                }],
623            }])
624            .await?;
625        let partition_response =
626            list_offset_partition_response(&response.topics, topic, partition)?;
627        if partition_response.error_code != 0 {
628            return Err(self.config.client.broker_error(
629                partition_response.error_code,
630                format!("list offsets {topic}-{partition}"),
631            ));
632        }
633        Ok(partition_response.offset)
634    }
635
636    /// Polls assigned partitions and advances in-memory offsets for fetched records.
637    #[tracing::instrument(
638        level = "debug",
639        name = "kafka.consumer.poll",
640        skip_all,
641        fields(assignment_count = self.assignments.len(), max_poll_records = self.config.max_poll_records),
642        err
643    )]
644    pub async fn poll(&mut self) -> Result<Vec<ConsumerRecord>> {
645        let assignments = self.assignments.clone();
646        let mut records = Vec::new();
647        let mut delivered_record_count = 0;
648        let assignment_count = assignments.len();
649        if assignment_count == 0 {
650            return Ok(records);
651        }
652        let start_index = self.poll_cursor % assignment_count;
653        let mut next_cursor = start_index;
654        debug!(
655            assignment_count = assignments.len(),
656            max_poll_records = self.config.max_poll_records,
657            "polling kafka consumer assignments"
658        );
659
660        for assignment_index in 0..assignment_count {
661            let index = (start_index + assignment_index) % assignment_count;
662            let assignment = assignments[index].clone();
663            if delivered_record_count >= self.config.max_poll_records {
664                break;
665            }
666            next_cursor = (index + 1) % assignment_count;
667            if assignment.paused {
668                continue;
669            }
670
671            let mut fetched = self
672                .fetch_with_progress(
673                    &assignment.topic,
674                    assignment.partition,
675                    assignment.next_offset,
676                    assignment.leader_epoch,
677                    Some(self.config.offset_reset_policy),
678                )
679                .await?;
680            if let Some(leader_epoch) = fetched.leader_epoch {
681                self.update_assignment_leader_epoch(
682                    &assignment.topic,
683                    assignment.partition,
684                    leader_epoch,
685                );
686            }
687            let fetched_record_count = fetched.records.len();
688            limit_fetched_records(
689                &mut fetched.records,
690                delivered_record_count,
691                self.config.max_poll_records,
692            );
693            let was_truncated = fetched.records.len() < fetched_record_count;
694
695            let key = (assignment.topic.clone(), assignment.partition);
696            if self.partition_queues.contains_key(&key) {
697                let route = self.enqueue_partition_records(&key, fetched.records)?;
698                fetched.records = route.records;
699                delivered_record_count += route.queued_count;
700                if fetched.records.is_empty() {
701                    let next_offset = if was_truncated {
702                        route.next_offset
703                    } else {
704                        (fetched_record_count > 0).then_some(fetched.next_offset)
705                    };
706                    if let Some(next_offset) = next_offset {
707                        self.update_assignment_offset(
708                            &assignment.topic,
709                            assignment.partition,
710                            next_offset,
711                        );
712                    }
713                    continue;
714                }
715            }
716            let next_offset = if fetched.records.len() < fetched_record_count {
717                fetched
718                    .records
719                    .last()
720                    .map(|record| record.offset().saturating_add(1))
721            } else {
722                Some(fetched.next_offset)
723            };
724            if let Some(next_offset) = next_offset {
725                self.update_assignment_offset(&assignment.topic, assignment.partition, next_offset);
726            }
727            if was_truncated && fetched.records.is_empty() {
728                break;
729            }
730            delivered_record_count += fetched.records.len();
731            records.extend(fetched.records);
732        }
733
734        self.poll_cursor = next_cursor;
735
736        debug!(
737            record_count = records.len(),
738            "polled kafka consumer records"
739        );
740        self.config.client.record_consumed(delivered_record_count);
741        Ok(records)
742    }
743
744    fn enqueue_partition_records(
745        &mut self,
746        key: &(String, i32),
747        records: Vec<ConsumerRecord>,
748    ) -> Result<PartitionRoute> {
749        let Some(sender) = self.partition_queues.get(key).cloned() else {
750            return Ok(PartitionRoute {
751                records,
752                next_offset: None,
753                queued_count: 0,
754            });
755        };
756        if sender.is_closed() {
757            self.partition_queues.remove(key);
758            return Ok(PartitionRoute {
759                records,
760                next_offset: None,
761                queued_count: 0,
762            });
763        }
764
765        let mut iterator = records.into_iter();
766        let mut queued_count = 0;
767        let mut next_offset = None;
768        while let Some(record) = iterator.next() {
769            let record_next_offset = record.offset().saturating_add(1);
770            match sender.try_send(record) {
771                Ok(()) => {
772                    queued_count += 1;
773                    next_offset = Some(record_next_offset);
774                }
775                Err(TrySendError::Closed(record)) => {
776                    self.partition_queues.remove(key);
777                    let mut main_records = Vec::with_capacity(iterator.len() + 1);
778                    main_records.push(record);
779                    main_records.extend(iterator);
780                    return Ok(PartitionRoute {
781                        records: main_records,
782                        next_offset,
783                        queued_count,
784                    });
785                }
786                Err(TrySendError::Full(_record)) => {
787                    if let Some(next_offset) = next_offset {
788                        self.update_assignment_offset(key.0.as_str(), key.1, next_offset);
789                    }
790                    return Err(Error::PartitionQueueFull {
791                        topic: key.0.clone(),
792                        partition: key.1,
793                        capacity: self.config.partition_queue_capacity,
794                    });
795                }
796            }
797        }
798
799        Ok(PartitionRoute {
800            records: Vec::new(),
801            next_offset,
802            queued_count,
803        })
804    }
805
806    /// Fetches records for one topic partition without changing assignment state.
807    #[tracing::instrument(
808        level = "debug",
809        name = "kafka.consumer.fetch",
810        skip_all,
811        fields(topic = tracing::field::Empty, partition, offset),
812        err
813    )]
814    pub async fn fetch(
815        &mut self,
816        topic: impl Into<String>,
817        partition: i32,
818        offset: i64,
819    ) -> Result<Vec<ConsumerRecord>> {
820        let topic = topic.into();
821        tracing::Span::current().record("topic", topic.as_str());
822        let records = self
823            .fetch_with_progress(&topic, partition, offset, -1, None)
824            .await?
825            .records;
826        self.config.client.record_consumed(records.len());
827        Ok(records)
828    }
829
830    async fn fetch_with_progress(
831        &mut self,
832        topic: &str,
833        partition: i32,
834        offset: i64,
835        current_leader_epoch: i32,
836        offset_reset_policy: Option<OffsetResetPolicy>,
837    ) -> Result<FetchedPartition> {
838        let mut attempt = 0;
839        let mut fetch_offset = offset;
840        let mut request_leader_epoch = current_leader_epoch;
841        let mut reset_applied = false;
842        let mut truncation_checked = false;
843        let mut recovered_leader_epoch = None;
844        debug!(topic, partition, offset, "fetching kafka records");
845
846        loop {
847            let result = self
848                .fetch_once(topic, partition, fetch_offset, request_leader_epoch)
849                .await;
850            if result.is_err() {
851                self.preferred_read_replicas
852                    .remove(&(topic.to_owned(), partition));
853            }
854            match result {
855                Err(error)
856                    if !reset_applied
857                        && offset_reset_policy.is_some_and(OffsetResetPolicy::is_recovery)
858                        && is_offset_out_of_range(&error) =>
859                {
860                    let Some(policy) = offset_reset_policy else {
861                        return Err(error);
862                    };
863                    fetch_offset = self.reset_fetch_offset(topic, partition, policy).await?;
864                    request_leader_epoch = -1;
865                    reset_applied = true;
866                    self.config.client.record_retry();
867                    attempt += 1;
868                }
869                Err(error)
870                    if !truncation_checked
871                        && attempt < self.config.max_retries
872                        && request_leader_epoch >= 0
873                        && is_leader_epoch_transition_error(&error) =>
874                {
875                    let previous_leader_epoch = request_leader_epoch;
876                    let current_leader_epoch = self
877                        .current_leader_epoch(topic, partition)
878                        .await?
879                        .filter(|epoch| *epoch >= previous_leader_epoch);
880                    if let Some(current_leader_epoch) = current_leader_epoch {
881                        invalidate_metadata_cache(&mut self.metadata_cache, topic);
882                        let epoch_offset = self
883                            .offset_for_leader_epoch(
884                                topic,
885                                partition,
886                                current_leader_epoch,
887                                previous_leader_epoch,
888                            )
889                            .await?;
890                        let end_offset = epoch_offset.end_offset();
891                        if end_offset >= 0 {
892                            fetch_offset = fetch_offset.min(end_offset);
893                        }
894                        debug!(
895                            topic,
896                            partition,
897                            previous_leader_epoch,
898                            current_leader_epoch,
899                            end_offset,
900                            fetch_offset,
901                            "recovered fetch after leader epoch transition"
902                        );
903                        request_leader_epoch = current_leader_epoch;
904                        recovered_leader_epoch = Some(current_leader_epoch);
905                        truncation_checked = true;
906                    } else {
907                        request_leader_epoch = -1;
908                        invalidate_metadata_cache(&mut self.metadata_cache, topic);
909                    }
910                    self.config.client.record_retry();
911                    attempt += 1;
912                }
913                Err(error) if attempt < self.config.max_retries && can_retry_fetch(&error) => {
914                    invalidate_metadata_cache(&mut self.metadata_cache, topic);
915                    self.config.client.record_retry();
916                    attempt += 1;
917                }
918                Ok(records) => {
919                    debug!(
920                        topic,
921                        partition,
922                        offset,
923                        record_count = records.records.len(),
924                        "fetched kafka records"
925                    );
926                    let mut records = records;
927                    if records.leader_epoch.is_none() {
928                        records.leader_epoch = recovered_leader_epoch;
929                    }
930                    return Ok(records);
931                }
932                Err(error) => return Err(error),
933            }
934        }
935    }
936
937    async fn reset_fetch_offset(
938        &mut self,
939        topic: &str,
940        partition: i32,
941        policy: OffsetResetPolicy,
942    ) -> Result<i64> {
943        match policy {
944            OffsetResetPolicy::Earliest => Ok(self.fetch_watermarks(topic, partition).await?.low),
945            OffsetResetPolicy::Latest => Ok(self.fetch_watermarks(topic, partition).await?.high),
946            OffsetResetPolicy::Offset(_) => Err(Error::Unsupported(
947                "explicit offset reset policy cannot recover an out-of-range fetch",
948            )),
949        }
950    }
951
952    async fn current_leader_epoch(&mut self, topic: &str, partition: i32) -> Result<Option<i32>> {
953        if !self.client.supports_metadata_v12().await? {
954            return Ok(None);
955        }
956
957        let metadata = self
958            .client
959            .metadata_v12(Some(vec![MetadataRequestTopicV12 {
960                topic_id: [0; 16],
961                name: Some(topic.to_owned()),
962            }]))
963            .await?;
964        let partition_metadata = metadata_v12_partition(&metadata, topic, partition)?;
965        if partition_metadata.error_code != 0 {
966            return Err(self.config.client.broker_error(
967                partition_metadata.error_code,
968                format!("metadata {topic}-{partition}"),
969            ));
970        }
971        Ok(Some(partition_metadata.leader_epoch))
972    }
973
974    async fn fetch_once(
975        &mut self,
976        topic: &str,
977        partition: i32,
978        offset: i64,
979        current_leader_epoch: i32,
980    ) -> Result<FetchedPartition> {
981        let metadata = self.metadata_for_topic(topic).await?;
982        let leader = leader_for(&metadata, topic, partition)?;
983        let selected_broker = self
984            .preferred_read_replicas
985            .get(&(topic.to_owned(), partition))
986            .copied()
987            .filter(|broker| broker_addr_for(&metadata, *broker).is_ok())
988            .unwrap_or(leader);
989        let broker_addr = broker_addr_for(&metadata, selected_broker)?;
990        debug!(
991            topic = topic,
992            partition,
993            leader,
994            selected_broker,
995            broker_addr = broker_addr.as_str(),
996            "resolved fetch broker"
997        );
998        let mut leader_client = self.connect_or_reuse_broker(&broker_addr).await?;
999        let rack_id = self
1000            .config
1001            .client_config()
1002            .client_rack_ref()
1003            .map(str::to_owned)
1004            .unwrap_or_default();
1005        let session = self
1006            .fetch_sessions
1007            .get(&broker_addr)
1008            .copied()
1009            .unwrap_or_default();
1010        let (session_id, session_epoch) = session.next_request();
1011        let topic_id = if leader_client.supports_fetch_v13().await? {
1012            self.topic_id_for_fetch(topic).await?
1013        } else {
1014            None
1015        };
1016        let (
1017            error_code,
1018            preferred_read_replica,
1019            aborted_transactions,
1020            records,
1021            response_session_id,
1022        ) = if let Some(topic_id) = topic_id {
1023            let response = leader_client
1024                .fetch_v13(
1025                    None,
1026                    self.config.max_wait_ms,
1027                    self.config.min_bytes,
1028                    self.config.max_partition_bytes,
1029                    self.config.isolation_level.as_i8(),
1030                    session_id,
1031                    session_epoch,
1032                    vec![FetchTopicV13 {
1033                        topic_id,
1034                        partitions: vec![FetchPartitionV12 {
1035                            partition_index: partition,
1036                            current_leader_epoch,
1037                            fetch_offset: offset,
1038                            last_fetched_epoch: current_leader_epoch,
1039                            log_start_offset: -1,
1040                            max_bytes: self.config.max_partition_bytes,
1041                        }],
1042                    }],
1043                    Vec::new(),
1044                    rack_id,
1045                )
1046                .await?;
1047            if response.error_code != 0 {
1048                self.fetch_sessions.remove(&broker_addr);
1049                self.fetch_topic_ids.remove(topic);
1050                return Err(self.config.client.broker_error(
1051                    response.error_code,
1052                    format!("fetch {topic}-{partition}@{offset}"),
1053                ));
1054            }
1055            let partition_response = fetch_partition_response_v13(&response, topic_id, partition)?;
1056            (
1057                partition_response.error_code,
1058                Some(partition_response.preferred_read_replica),
1059                partition_response
1060                    .aborted_transactions
1061                    .iter()
1062                    .map(|transaction| AbortedTransactionV4 {
1063                        producer_id: transaction.producer_id,
1064                        first_offset: transaction.first_offset,
1065                    })
1066                    .collect(),
1067                partition_response.records.clone(),
1068                response.session_id,
1069            )
1070        } else if leader_client.supports_fetch_v12().await? {
1071            let response = leader_client
1072                .fetch_one_v12(FetchOneRequestV12 {
1073                    replica_id: -1,
1074                    max_wait_ms: self.config.max_wait_ms,
1075                    min_bytes: self.config.min_bytes,
1076                    max_bytes: self.config.max_partition_bytes,
1077                    isolation_level: self.config.isolation_level.as_i8(),
1078                    topic: topic.to_owned(),
1079                    partition_index: partition,
1080                    current_leader_epoch,
1081                    fetch_offset: offset,
1082                    last_fetched_epoch: current_leader_epoch,
1083                    max_partition_bytes: self.config.max_partition_bytes,
1084                    session_id,
1085                    session_epoch,
1086                    rack_id,
1087                })
1088                .await?;
1089            if response.error_code != 0 {
1090                self.fetch_sessions.remove(&broker_addr);
1091                return Err(self.config.client.broker_error(
1092                    response.error_code,
1093                    format!("fetch {topic}-{partition}@{offset}"),
1094                ));
1095            }
1096            let partition_response = fetch_partition_response_v12(&response, topic, partition)?;
1097            (
1098                partition_response.error_code,
1099                Some(partition_response.preferred_read_replica),
1100                partition_response
1101                    .aborted_transactions
1102                    .iter()
1103                    .map(|transaction| AbortedTransactionV4 {
1104                        producer_id: transaction.producer_id,
1105                        first_offset: transaction.first_offset,
1106                    })
1107                    .collect(),
1108                partition_response.records.clone(),
1109                response.session_id,
1110            )
1111        } else if leader_client.supports_fetch_v11().await? {
1112            let response = leader_client
1113                .fetch_one_v11(FetchOneRequestV11 {
1114                    replica_id: -1,
1115                    max_wait_ms: self.config.max_wait_ms,
1116                    min_bytes: self.config.min_bytes,
1117                    max_bytes: self.config.max_partition_bytes,
1118                    isolation_level: self.config.isolation_level.as_i8(),
1119                    topic: topic.to_owned(),
1120                    partition_index: partition,
1121                    current_leader_epoch,
1122                    fetch_offset: offset,
1123                    max_partition_bytes: self.config.max_partition_bytes,
1124                    session_id,
1125                    session_epoch,
1126                    rack_id,
1127                })
1128                .await?;
1129            if response.error_code != 0 {
1130                self.fetch_sessions.remove(&broker_addr);
1131                return Err(self.config.client.broker_error(
1132                    response.error_code,
1133                    format!("fetch {topic}-{partition}@{offset}"),
1134                ));
1135            }
1136            let partition_response = fetch_partition_response_v11(&response, topic, partition)?;
1137            (
1138                partition_response.error_code,
1139                Some(partition_response.preferred_read_replica),
1140                partition_response.aborted_transactions.clone(),
1141                partition_response.records.clone(),
1142                response.session_id,
1143            )
1144        } else {
1145            let response = leader_client
1146                .fetch_one_v4(FetchOneRequestV4 {
1147                    replica_id: -1,
1148                    max_wait_ms: self.config.max_wait_ms,
1149                    min_bytes: self.config.min_bytes,
1150                    max_bytes: self.config.max_partition_bytes,
1151                    isolation_level: self.config.isolation_level.as_i8(),
1152                    topic: topic.to_owned(),
1153                    partition_index: partition,
1154                    fetch_offset: offset,
1155                    max_partition_bytes: self.config.max_partition_bytes,
1156                })
1157                .await?;
1158            let partition_response = fetch_partition_response(&response, topic, partition)?;
1159            (
1160                partition_response.error_code,
1161                Some(-1),
1162                partition_response.aborted_transactions.clone(),
1163                partition_response.records.clone(),
1164                0,
1165            )
1166        };
1167        if error_code != 0 {
1168            self.fetch_sessions.remove(&broker_addr);
1169            self.fetch_topic_ids.remove(topic);
1170            return Err(self
1171                .config
1172                .client
1173                .broker_error(error_code, format!("fetch {topic}-{partition}@{offset}")));
1174        }
1175
1176        if response_session_id > 0 {
1177            self.fetch_sessions.insert(
1178                broker_addr.clone(),
1179                session.advance_with_response(response_session_id),
1180            );
1181        } else {
1182            self.fetch_sessions.remove(&broker_addr);
1183        }
1184
1185        if let Some(preferred_read_replica) = preferred_read_replica {
1186            if preferred_read_replica >= 0
1187                && broker_addr_for(&metadata, preferred_read_replica).is_ok()
1188            {
1189                self.preferred_read_replicas
1190                    .insert((topic.to_owned(), partition), preferred_read_replica);
1191            } else {
1192                self.preferred_read_replicas
1193                    .remove(&(topic.to_owned(), partition));
1194            }
1195        }
1196
1197        let next_offset = records
1198            .last()
1199            .map(|record| record.offset.saturating_add(1))
1200            .unwrap_or(offset);
1201        let leader_epoch = records
1202            .last()
1203            .map(|record| record.leader_epoch)
1204            .filter(|leader_epoch| *leader_epoch >= 0);
1205        let records = visible_records(&aborted_transactions, &records, self.config.isolation_level)
1206            .into_iter()
1207            .map(|record| ConsumerRecord::from_message_set(topic, partition, record))
1208            .collect();
1209        self.cache_broker_client(broker_addr, leader_client);
1210        Ok(FetchedPartition {
1211            records,
1212            next_offset,
1213            leader_epoch,
1214        })
1215    }
1216
1217    fn update_assignment_offset(&mut self, topic: &str, partition: i32, next_offset: i64) {
1218        if let Some(assignment) = self
1219            .assignments
1220            .iter_mut()
1221            .find(|assignment| assignment.topic == topic && assignment.partition == partition)
1222        {
1223            assignment.next_offset = next_offset;
1224        }
1225    }
1226
1227    fn update_assignment_leader_epoch(&mut self, topic: &str, partition: i32, leader_epoch: i32) {
1228        if let Some(assignment) = self
1229            .assignments
1230            .iter_mut()
1231            .find(|assignment| assignment.topic == topic && assignment.partition == partition)
1232        {
1233            assignment.set_leader_epoch(assignment.leader_epoch.max(leader_epoch));
1234        }
1235    }
1236
1237    fn assignment(&self, topic: &str, partition: i32) -> Option<&ConsumerAssignment> {
1238        self.assignments
1239            .iter()
1240            .find(|assignment| assignment.topic == topic && assignment.partition == partition)
1241    }
1242
1243    fn assignment_mut(&mut self, topic: &str, partition: i32) -> Result<&mut ConsumerAssignment> {
1244        self.assignments
1245            .iter_mut()
1246            .find(|assignment| assignment.topic == topic && assignment.partition == partition)
1247            .ok_or_else(|| Error::UnassignedTopicPartition {
1248                topic: topic.to_owned(),
1249                partition,
1250            })
1251    }
1252
1253    async fn metadata_for_topic(&mut self, topic: &str) -> Result<MetadataResponseV1> {
1254        if let Some(metadata) = self.metadata_cache.get(topic) {
1255            return Ok(metadata.clone());
1256        }
1257
1258        let metadata = self.request_metadata_for_topic(topic).await?;
1259        self.metadata_cache
1260            .insert(topic.to_owned(), metadata.clone());
1261        Ok(metadata)
1262    }
1263
1264    async fn topic_id_for_fetch(&mut self, topic: &str) -> Result<Option<[u8; 16]>> {
1265        if let Some(topic_id) = self.fetch_topic_ids.get(topic) {
1266            return Ok(Some(*topic_id));
1267        }
1268        if !self.client.supports_metadata_v12().await? {
1269            return Ok(None);
1270        }
1271
1272        let request = Some(vec![MetadataRequestTopicV12 {
1273            topic_id: [0; 16],
1274            name: Some(topic.to_owned()),
1275        }]);
1276        let metadata = match self.client.metadata_v12(request.clone()).await {
1277            Ok(metadata) => metadata,
1278            Err(error) if can_retry_fetch(&error) => {
1279                self.config.client.record_retry();
1280                self.client = self.config.client.clone().connect().await?;
1281                self.client.metadata_v12(request).await?
1282            }
1283            Err(error) => return Err(error),
1284        };
1285        let topic_metadata = metadata
1286            .topics
1287            .iter()
1288            .find(|metadata| metadata.name.as_deref() == Some(topic))
1289            .ok_or_else(|| Error::UnknownTopicOrPartition {
1290                topic: topic.to_owned(),
1291                partition: -1,
1292            })?;
1293        if topic_metadata.error_code != 0 {
1294            return Err(self
1295                .config
1296                .client
1297                .broker_error(topic_metadata.error_code, format!("metadata topic {topic}")));
1298        }
1299        let topic_id = (topic_metadata.topic_id != [0; 16]).then_some(topic_metadata.topic_id);
1300        if let Some(topic_id) = topic_id {
1301            self.fetch_topic_ids.insert(topic.to_owned(), topic_id);
1302        }
1303        Ok(topic_id)
1304    }
1305
1306    async fn connect_or_reuse_broker(&mut self, broker_addr: &str) -> Result<Client> {
1307        if let Some(client) = self.broker_clients.take(broker_addr) {
1308            return Ok(client);
1309        }
1310        self.fetch_sessions.remove(broker_addr);
1311        self.config
1312            .client
1313            .connect_broker(broker_addr.to_owned())
1314            .await
1315    }
1316
1317    fn cache_broker_client(&mut self, broker_addr: String, client: Client) {
1318        self.broker_clients.insert(
1319            broker_addr,
1320            client,
1321            self.config.client.max_idle_broker_connections_ref(),
1322        );
1323    }
1324
1325    async fn request_metadata_for_topic(&mut self, topic: &str) -> Result<MetadataResponseV1> {
1326        let topics = Some(vec![topic.to_owned()]);
1327        match self.client.metadata(topics.clone()).await {
1328            Ok(metadata) => Ok(metadata),
1329            Err(error) if can_retry_fetch(&error) => {
1330                self.config.client.record_retry();
1331                debug!(
1332                    topic,
1333                    error = %error,
1334                    "reconnecting metadata client after metadata request failure"
1335                );
1336                self.client = self.config.client.clone().connect().await?;
1337                self.client.metadata(topics).await
1338            }
1339            Err(error) => Err(error),
1340        }
1341    }
1342}
1343
1344#[derive(Debug)]
1345struct FetchedPartition {
1346    records: Vec<ConsumerRecord>,
1347    next_offset: i64,
1348    leader_epoch: Option<i32>,
1349}
1350
1351struct PartitionRoute {
1352    records: Vec<ConsumerRecord>,
1353    next_offset: Option<i64>,
1354    queued_count: usize,
1355}
1356
1357#[derive(Debug, Clone, PartialEq, Eq)]
1358/// Direct consumer topic partition assignment.
1359pub struct ConsumerAssignment {
1360    topic: String,
1361    partition: i32,
1362    next_offset: i64,
1363    leader_epoch: i32,
1364    paused: bool,
1365}
1366
1367impl ConsumerAssignment {
1368    pub(crate) fn new(topic: String, partition: i32, next_offset: i64) -> Self {
1369        Self {
1370            topic,
1371            partition,
1372            next_offset,
1373            leader_epoch: -1,
1374            paused: false,
1375        }
1376    }
1377
1378    pub(crate) fn set_leader_epoch(&mut self, leader_epoch: i32) {
1379        self.leader_epoch = leader_epoch;
1380    }
1381
1382    /// Returns the Kafka topic name.
1383    pub fn topic(&self) -> &str {
1384        &self.topic
1385    }
1386
1387    /// Returns the Kafka partition index.
1388    pub fn partition(&self) -> i32 {
1389        self.partition
1390    }
1391
1392    /// Returns the next offset that will be fetched or committed.
1393    pub fn next_offset(&self) -> i64 {
1394        self.next_offset
1395    }
1396
1397    /// Returns the latest partition leader epoch observed by this assignment.
1398    ///
1399    /// The initial value is `-1` until a RecordBatch response provides an
1400    /// epoch. Legacy MessageSet responses do not update this value.
1401    pub fn leader_epoch(&self) -> i32 {
1402        self.leader_epoch
1403    }
1404
1405    /// Returns whether fetching is paused for this assignment.
1406    pub fn is_paused(&self) -> bool {
1407        self.paused
1408    }
1409}
1410
1411fn assign_partition(
1412    assignments: &mut Vec<ConsumerAssignment>,
1413    topic: String,
1414    partition: i32,
1415    offset: i64,
1416) {
1417    if let Some(assignment) = assignments
1418        .iter_mut()
1419        .find(|assignment| assignment.topic == topic && assignment.partition == partition)
1420    {
1421        assignment.next_offset = offset;
1422        assignment.leader_epoch = -1;
1423        return;
1424    }
1425
1426    assignments.push(ConsumerAssignment {
1427        topic,
1428        partition,
1429        next_offset: offset,
1430        leader_epoch: -1,
1431        paused: false,
1432    });
1433}
1434
1435#[derive(Debug, Clone, PartialEq, Eq)]
1436/// Configuration builder for [`Consumer`].
1437pub struct ConsumerConfig {
1438    client: ClientConfig,
1439    max_wait_ms: i32,
1440    min_bytes: i32,
1441    max_partition_bytes: i32,
1442    max_retries: u32,
1443    max_poll_records: usize,
1444    partition_queue_capacity: usize,
1445    isolation_level: IsolationLevel,
1446    offset_reset_policy: OffsetResetPolicy,
1447}
1448
1449impl ConsumerConfig {
1450    /// Creates a consumer configuration from one or more Kafka bootstrap servers.
1451    pub fn new(bootstrap_servers: impl IntoIterator<Item = impl Into<String>>) -> Self {
1452        Self::from_client_config(ClientConfig::new(bootstrap_servers))
1453    }
1454
1455    pub(crate) fn from_client_config(client: ClientConfig) -> Self {
1456        Self {
1457            client,
1458            max_wait_ms: 500,
1459            min_bytes: 1,
1460            max_partition_bytes: 1_048_576,
1461            max_retries: 1,
1462            max_poll_records: 500,
1463            partition_queue_capacity: 1024,
1464            isolation_level: IsolationLevel::ReadUncommitted,
1465            offset_reset_policy: OffsetResetPolicy::Offset(0),
1466        }
1467    }
1468
1469    /// Replaces the shared client configuration used by consumer connections.
1470    ///
1471    /// This is useful when the same bootstrap, controller, security, limits,
1472    /// or metrics policy is shared by several kafrust clients.
1473    pub fn with_client_config(mut self, client: ClientConfig) -> Self {
1474        self.client = client;
1475        self
1476    }
1477
1478    /// Sets the Kafka client ID used by consumer requests.
1479    pub fn client_id(mut self, client_id: impl Into<String>) -> Self {
1480        self.client = self.client.client_id(client_id);
1481        self
1482    }
1483
1484    /// Sets the rack ID used by rack-aware consumer Fetch requests.
1485    pub fn client_rack(mut self, client_rack: impl Into<String>) -> Self {
1486        self.client = self.client.client_rack(client_rack);
1487        self
1488    }
1489
1490    /// Sets the request timeout in milliseconds.
1491    pub fn request_timeout_ms(mut self, request_timeout_ms: u64) -> Self {
1492        self.client = self.client.request_timeout_ms(request_timeout_ms);
1493        self
1494    }
1495
1496    /// Sets the maximum broker response payload allocated for one consumer request.
1497    pub fn max_response_bytes(mut self, max_response_bytes: usize) -> Self {
1498        self.client = self.client.max_response_bytes(max_response_bytes);
1499        self
1500    }
1501
1502    /// Sets the maximum number of idle broker connections retained by this
1503    /// consumer. A value of zero is rejected during validation.
1504    pub fn max_idle_broker_connections(mut self, max: usize) -> Self {
1505        self.client = self.client.max_idle_broker_connections(max);
1506        self
1507    }
1508
1509    /// Sets the maximum number of elements allocated for one Kafka response array.
1510    pub fn max_decode_array_elements(mut self, max: usize) -> Self {
1511        self.client = self.client.max_decode_array_elements(max);
1512        self
1513    }
1514
1515    /// Sets the maximum uncompressed size of one fetched record batch.
1516    pub fn max_decompressed_record_bytes(mut self, max: usize) -> Self {
1517        self.client = self.client.max_decompressed_record_bytes(max);
1518        self
1519    }
1520
1521    /// Sets the shared metrics handle used by consumer broker connections.
1522    pub fn metrics(mut self, metrics: ClientMetrics) -> Self {
1523        self.client = self.client.metrics(metrics);
1524        self
1525    }
1526
1527    /// Sets the Kafka security protocol used for consumer broker connections.
1528    pub fn security_protocol(mut self, security_protocol: SecurityProtocol) -> Self {
1529        self.client = self.client.security_protocol(security_protocol);
1530        self
1531    }
1532
1533    /// Sets the TLS server name used for consumer broker certificate validation.
1534    pub fn tls_server_name(mut self, server_name: impl Into<String>) -> Self {
1535        self.client = self.client.tls_server_name(server_name);
1536        self
1537    }
1538
1539    /// Adds a DER-encoded TLS root certificate for consumer broker validation.
1540    pub fn tls_root_certificate_der(mut self, certificate: impl Into<Vec<u8>>) -> Self {
1541        self.client = self.client.tls_root_certificate_der(certificate);
1542        self
1543    }
1544
1545    /// Adds a DER-encoded client certificate for TLS mutual authentication.
1546    pub fn tls_client_certificate_der(mut self, certificate: impl Into<Vec<u8>>) -> Self {
1547        self.client = self.client.tls_client_certificate_der(certificate);
1548        self
1549    }
1550
1551    /// Sets the DER-encoded private key for TLS mutual authentication.
1552    pub fn tls_client_private_key_der(mut self, key: impl Into<Vec<u8>>) -> Self {
1553        self.client = self.client.tls_client_private_key_der(key);
1554        self
1555    }
1556
1557    /// Sets SASL/PLAIN credentials for consumer broker connections.
1558    pub fn sasl_plain(mut self, username: impl Into<String>, password: impl Into<String>) -> Self {
1559        self.client = self.client.sasl_plain(username, password);
1560        self
1561    }
1562
1563    /// Sets SASL/SCRAM-SHA-256 credentials for consumer broker connections.
1564    pub fn sasl_scram_sha_256(
1565        mut self,
1566        username: impl Into<String>,
1567        password: impl Into<String>,
1568    ) -> Self {
1569        self.client = self.client.sasl_scram_sha_256(username, password);
1570        self
1571    }
1572
1573    /// Sets SASL/SCRAM-SHA-512 credentials for consumer broker connections.
1574    pub fn sasl_scram_sha_512(
1575        mut self,
1576        username: impl Into<String>,
1577        password: impl Into<String>,
1578    ) -> Self {
1579        self.client = self.client.sasl_scram_sha_512(username, password);
1580        self
1581    }
1582
1583    /// Sets SASL/OAUTHBEARER credentials for consumer broker connections.
1584    pub fn sasl_oauthbearer(mut self, token: impl Into<String>) -> Self {
1585        self.client = self.client.sasl_oauthbearer(token);
1586        self
1587    }
1588
1589    /// Sets SASL/OAUTHBEARER credentials with an authorization identity.
1590    pub fn sasl_oauthbearer_with_username(
1591        mut self,
1592        username: impl Into<String>,
1593        token: impl Into<String>,
1594    ) -> Self {
1595        self.client = self.client.sasl_oauthbearer_with_username(username, token);
1596        self
1597    }
1598
1599    /// Sets SASL/OAUTHBEARER credentials from an async token provider.
1600    pub fn sasl_oauthbearer_provider<P>(mut self, provider: P) -> Self
1601    where
1602        P: OAuthBearerTokenProvider + 'static,
1603    {
1604        self.client = self.client.sasl_oauthbearer_provider(provider);
1605        self
1606    }
1607
1608    /// Sets SASL/OAUTHBEARER credentials with an authorization identity and
1609    /// an async token provider.
1610    pub fn sasl_oauthbearer_with_username_and_provider<P>(
1611        mut self,
1612        username: impl Into<String>,
1613        provider: P,
1614    ) -> Self
1615    where
1616        P: OAuthBearerTokenProvider + 'static,
1617    {
1618        self.client = self
1619            .client
1620            .sasl_oauthbearer_with_username_and_provider(username, provider);
1621        self
1622    }
1623
1624    /// Sets the Kafka fetch max wait time in milliseconds.
1625    pub fn max_wait_ms(mut self, max_wait_ms: i32) -> Self {
1626        self.max_wait_ms = max_wait_ms;
1627        self
1628    }
1629
1630    /// Sets the Kafka fetch minimum response bytes.
1631    pub fn min_bytes(mut self, min_bytes: i32) -> Self {
1632        self.min_bytes = min_bytes;
1633        self
1634    }
1635
1636    /// Sets the maximum bytes to fetch from each partition.
1637    pub fn max_partition_bytes(mut self, max_partition_bytes: i32) -> Self {
1638        self.max_partition_bytes = max_partition_bytes;
1639        self
1640    }
1641
1642    /// Sets the maximum number of retry attempts for transient fetch failures.
1643    pub fn max_retries(mut self, max_retries: u32) -> Self {
1644        self.max_retries = max_retries;
1645        self
1646    }
1647
1648    /// Returns the configured maximum retry count.
1649    pub fn max_retries_ref(&self) -> u32 {
1650        self.max_retries
1651    }
1652
1653    /// Sets the maximum number of records returned by one poll.
1654    pub fn max_poll_records(mut self, max_poll_records: usize) -> Self {
1655        self.max_poll_records = max_poll_records;
1656        self
1657    }
1658
1659    /// Returns the configured maximum records per poll.
1660    pub fn max_poll_records_ref(&self) -> usize {
1661        self.max_poll_records
1662    }
1663
1664    /// Sets the bounded capacity of queues returned by
1665    /// [`Consumer::split_partition_queue`]. A value of zero is normalized to
1666    /// one because Tokio channels cannot be created without capacity.
1667    pub fn partition_queue_capacity(mut self, partition_queue_capacity: usize) -> Self {
1668        self.partition_queue_capacity = partition_queue_capacity.max(1);
1669        self
1670    }
1671
1672    /// Returns the configured partition queue capacity.
1673    pub fn partition_queue_capacity_ref(&self) -> usize {
1674        self.partition_queue_capacity
1675    }
1676
1677    /// Sets whether fetches expose records from aborted transactions.
1678    pub fn isolation_level(mut self, isolation_level: IsolationLevel) -> Self {
1679        self.isolation_level = isolation_level;
1680        self
1681    }
1682
1683    /// Returns the configured transaction isolation level.
1684    pub fn isolation_level_ref(&self) -> IsolationLevel {
1685        self.isolation_level
1686    }
1687
1688    /// Sets the fallback used when an assigned fetch offset is out of range.
1689    ///
1690    /// `Earliest` and `Latest` perform one bounded reset through the partition
1691    /// leader. `Offset(n)` preserves the explicit-offset behavior and returns
1692    /// the broker error instead of silently changing the requested position.
1693    pub fn offset_reset_policy(mut self, offset_reset_policy: OffsetResetPolicy) -> Self {
1694        self.offset_reset_policy = offset_reset_policy;
1695        self
1696    }
1697
1698    /// Returns the configured out-of-range offset policy.
1699    pub fn offset_reset_policy_ref(&self) -> OffsetResetPolicy {
1700        self.offset_reset_policy
1701    }
1702
1703    /// Returns the shared client configuration.
1704    pub fn client_config(&self) -> &ClientConfig {
1705        &self.client
1706    }
1707
1708    /// Validates this consumer configuration without opening a broker connection.
1709    pub fn validate(&self) -> Result<()> {
1710        self.client.validate()?;
1711        self.validate_values()
1712    }
1713
1714    /// Validates and returns this consumer configuration without opening a
1715    /// broker connection.
1716    pub fn build_config(self) -> Result<Self> {
1717        self.validate()?;
1718        Ok(self)
1719    }
1720
1721    /// Connects to Kafka and builds a direct consumer.
1722    pub async fn build(self) -> Result<Consumer> {
1723        self.validate()?;
1724        let client = self.client.clone().connect().await?;
1725        Ok(Consumer {
1726            client,
1727            config: self,
1728            assignments: Vec::new(),
1729            poll_cursor: 0,
1730            partition_queues: BTreeMap::new(),
1731            metadata_cache: BTreeMap::new(),
1732            fetch_topic_ids: BTreeMap::new(),
1733            broker_clients: BrokerClientCache::default(),
1734            fetch_sessions: BTreeMap::new(),
1735            preferred_read_replicas: BTreeMap::new(),
1736        })
1737    }
1738
1739    fn validate_values(&self) -> Result<()> {
1740        if self.max_wait_ms < 0 {
1741            return Err(Error::InvalidConfiguration {
1742                field: "max_wait_ms",
1743                reason: "must not be negative",
1744            });
1745        }
1746        if self.min_bytes < 0 {
1747            return Err(Error::InvalidConfiguration {
1748                field: "min_bytes",
1749                reason: "must not be negative",
1750            });
1751        }
1752        if self.max_partition_bytes <= 0 {
1753            return Err(Error::InvalidConfiguration {
1754                field: "max_partition_bytes",
1755                reason: "must be greater than zero",
1756            });
1757        }
1758        if self.max_poll_records == 0 {
1759            return Err(Error::InvalidConfiguration {
1760                field: "max_poll_records",
1761                reason: "must be greater than zero",
1762            });
1763        }
1764        Ok(())
1765    }
1766}
1767
1768fn leader_for(
1769    metadata: &MetadataResponseV1,
1770    topic_name: &str,
1771    partition_index: i32,
1772) -> Result<i32> {
1773    metadata
1774        .topics
1775        .iter()
1776        .find(|topic| topic.name == topic_name)
1777        .and_then(|topic| {
1778            topic
1779                .partitions
1780                .iter()
1781                .find(|partition| partition.partition_index == partition_index)
1782        })
1783        .ok_or_else(|| Error::UnknownTopicOrPartition {
1784            topic: topic_name.to_owned(),
1785            partition: partition_index,
1786        })
1787        .and_then(|partition| {
1788            (partition.leader_id >= 0)
1789                .then_some(partition.leader_id)
1790                .ok_or_else(|| Error::MissingLeader {
1791                    topic: topic_name.to_owned(),
1792                    partition: partition_index,
1793                })
1794        })
1795}
1796
1797fn broker_addr_for(metadata: &MetadataResponseV1, node_id: i32) -> Result<String> {
1798    metadata
1799        .brokers
1800        .iter()
1801        .find(|broker| broker.node_id == node_id)
1802        .map(broker_addr)
1803        .ok_or(Error::MissingBroker { node_id })
1804}
1805
1806fn broker_addr(broker: &BrokerMetadata) -> String {
1807    format!("{}:{}", broker.host, broker.port)
1808}
1809
1810fn fetch_partition_response<'a>(
1811    response: &'a FetchResponseV4,
1812    topic_name: &str,
1813    partition_index: i32,
1814) -> Result<&'a FetchPartitionResponseV4> {
1815    response
1816        .responses
1817        .iter()
1818        .find(|topic| topic.name == topic_name)
1819        .and_then(|topic| {
1820            topic
1821                .partitions
1822                .iter()
1823                .find(|partition| partition.partition_index == partition_index)
1824        })
1825        .ok_or_else(|| Error::UnknownTopicOrPartition {
1826            topic: topic_name.to_owned(),
1827            partition: partition_index,
1828        })
1829}
1830
1831fn fetch_partition_response_v11<'a>(
1832    response: &'a FetchResponseV11,
1833    topic_name: &str,
1834    partition_index: i32,
1835) -> Result<&'a FetchPartitionResponseV11> {
1836    response
1837        .responses
1838        .iter()
1839        .find(|topic| topic.name == topic_name)
1840        .and_then(|topic| {
1841            topic
1842                .partitions
1843                .iter()
1844                .find(|partition| partition.partition_index == partition_index)
1845        })
1846        .ok_or_else(|| Error::UnknownTopicOrPartition {
1847            topic: topic_name.to_owned(),
1848            partition: partition_index,
1849        })
1850}
1851
1852fn fetch_partition_response_v12<'a>(
1853    response: &'a FetchResponseV12,
1854    topic_name: &str,
1855    partition_index: i32,
1856) -> Result<&'a FetchPartitionResponseV12> {
1857    response
1858        .responses
1859        .iter()
1860        .find(|topic| topic.name == topic_name)
1861        .and_then(|topic| {
1862            topic
1863                .partitions
1864                .iter()
1865                .find(|partition| partition.partition_index == partition_index)
1866        })
1867        .ok_or_else(|| Error::UnknownTopicOrPartition {
1868            topic: topic_name.to_owned(),
1869            partition: partition_index,
1870        })
1871}
1872
1873fn fetch_partition_response_v13(
1874    response: &FetchResponseV13,
1875    topic_id: [u8; 16],
1876    partition_index: i32,
1877) -> Result<&FetchPartitionResponseV13> {
1878    response
1879        .responses
1880        .iter()
1881        .find(|topic| topic.topic_id == topic_id)
1882        .and_then(|topic| {
1883            topic
1884                .partitions
1885                .iter()
1886                .find(|partition| partition.partition_index == partition_index)
1887        })
1888        .ok_or_else(|| Error::UnknownTopicOrPartition {
1889            topic: format!("topic-id-{topic_id:02x?}"),
1890            partition: partition_index,
1891        })
1892}
1893
1894fn list_offset_partition_response<'a>(
1895    topics: &'a [ListOffsetsTopicResponseV1],
1896    topic_name: &str,
1897    partition_index: i32,
1898) -> Result<&'a ListOffsetsPartitionResponseV1> {
1899    topics
1900        .iter()
1901        .find(|topic| topic.name == topic_name)
1902        .and_then(|topic| {
1903            topic
1904                .partitions
1905                .iter()
1906                .find(|partition| partition.partition_index == partition_index)
1907        })
1908        .ok_or_else(|| Error::UnknownTopicOrPartition {
1909            topic: topic_name.to_owned(),
1910            partition: partition_index,
1911        })
1912}
1913
1914fn offset_for_leader_epoch_partition_response<'a>(
1915    topics: &'a [OffsetForLeaderEpochTopicResponseV3],
1916    topic_name: &str,
1917    partition_index: i32,
1918) -> Result<&'a OffsetForLeaderEpochPartitionResponseV3> {
1919    topics
1920        .iter()
1921        .find(|topic| topic.name == topic_name)
1922        .and_then(|topic| {
1923            topic
1924                .partitions
1925                .iter()
1926                .find(|partition| partition.partition_index == partition_index)
1927        })
1928        .ok_or_else(|| Error::UnknownTopicOrPartition {
1929            topic: topic_name.to_owned(),
1930            partition: partition_index,
1931        })
1932}
1933
1934fn metadata_v12_partition<'a>(
1935    metadata: &'a MetadataResponseV12,
1936    topic_name: &str,
1937    partition_index: i32,
1938) -> Result<&'a MetadataPartitionV12> {
1939    let topic = metadata
1940        .topics
1941        .iter()
1942        .find(|topic| topic.name.as_deref() == Some(topic_name))
1943        .ok_or_else(|| Error::UnknownTopicOrPartition {
1944            topic: topic_name.to_owned(),
1945            partition: partition_index,
1946        })?;
1947    if topic.error_code != 0 {
1948        return Err(Error::Broker {
1949            code: topic.error_code,
1950            context: format!("metadata topic {topic_name}"),
1951        });
1952    }
1953    topic
1954        .partitions
1955        .iter()
1956        .find(|partition| partition.partition_index == partition_index)
1957        .ok_or_else(|| Error::UnknownTopicOrPartition {
1958            topic: topic_name.to_owned(),
1959            partition: partition_index,
1960        })
1961}
1962
1963fn visible_records(
1964    aborted_transactions: &[kafrust_protocol::api::fetch::AbortedTransactionV4],
1965    input_records: &[MessageSetRecord],
1966    isolation_level: IsolationLevel,
1967) -> Vec<MessageSetRecord> {
1968    let mut aborted_transactions = aborted_transactions.to_vec();
1969    let mut visible = Vec::new();
1970
1971    for record in input_records {
1972        if record.control {
1973            if let Some(producer_id) = record.producer_id {
1974                if let Some(index) = aborted_transactions.iter().position(|transaction| {
1975                    transaction.producer_id == producer_id
1976                        && transaction.first_offset <= record.offset
1977                }) {
1978                    aborted_transactions.remove(index);
1979                }
1980            }
1981            continue;
1982        }
1983
1984        let aborted = isolation_level == IsolationLevel::ReadCommitted
1985            && record.transactional
1986            && record.producer_id.is_some_and(|producer_id| {
1987                aborted_transactions.iter().any(|transaction| {
1988                    transaction.producer_id == producer_id
1989                        && transaction.first_offset <= record.offset
1990                })
1991            });
1992        if !aborted {
1993            visible.push(record.clone());
1994        }
1995    }
1996
1997    visible
1998}
1999
2000fn can_retry_fetch(error: &Error) -> bool {
2001    match error {
2002        Error::Broker { code, .. } => {
2003            *code == FETCH_UNKNOWN_TOPIC_ID_ERROR_CODE
2004                || matches!(
2005                    BrokerErrorKind::from_code(*code),
2006                    BrokerErrorKind::UnknownTopicOrPartition
2007                        | BrokerErrorKind::LeaderNotAvailable
2008                        | BrokerErrorKind::NotLeaderOrFollower
2009                        | BrokerErrorKind::RequestTimedOut
2010                        | BrokerErrorKind::ReplicaNotAvailable
2011                        | BrokerErrorKind::FencedLeaderEpoch
2012                        | BrokerErrorKind::UnknownLeaderEpoch
2013                        | BrokerErrorKind::InvalidFetchSessionEpoch
2014                )
2015        }
2016        Error::Io(_)
2017        | Error::RequestTimedOut { .. }
2018        | Error::UnknownTopicOrPartition { .. }
2019        | Error::MissingLeader { .. }
2020        | Error::MissingBroker { .. } => true,
2021        Error::MissingBootstrapServer
2022        | Error::InvalidPartition { .. }
2023        | Error::UnassignedTopicPartition { .. }
2024        | Error::PartitionQueueFull { .. }
2025        | Error::MissingGroupDescription { .. }
2026        | Error::MissingDeleteGroupResult { .. }
2027        | Error::ResponseCountMismatch { .. }
2028        | Error::MissingSaslCredentials
2029        | Error::InvalidSaslResponse { .. }
2030        | Error::OAuthBearerTokenTimeout { .. }
2031        | Error::OAuthBearerTokenExpired
2032        | Error::TransactionOutcomeUnknown { .. }
2033        | Error::TransactionProducerDefunct
2034        | Error::IdempotentProducerDefunct
2035        | Error::ConsumerGroupCommitOutcomeUnknown { .. }
2036        | Error::DeliveryDeadlineExceeded { .. }
2037        | Error::ResponseTooLarge { .. }
2038        | Error::TlsConfig { .. }
2039        | Error::InvalidTlsServerName { .. }
2040        | Error::InvalidGroupInstanceId
2041        | Error::InvalidTopicPattern { .. }
2042        | Error::InvalidScramCredential { .. }
2043        | Error::InvalidConfiguration { .. }
2044        | Error::ConsumerGroupAssignmentTimeout { .. }
2045        | Error::AdminMutationOutcomeUnknown { .. }
2046        | Error::ShareAcknowledgementRequired { .. }
2047        | Error::ShareAcknowledgementOutcomeUnknown { .. }
2048        | Error::ShareAcknowledgementSessionUnavailable { .. }
2049        | Error::ShareRecordNotPending { .. }
2050        | Error::ShareRecordAlreadyAcknowledged { .. }
2051        | Error::TelemetryPayloadTooLarge { .. }
2052        | Error::ShareRecordNotAcquired { .. }
2053        | Error::Unsupported(_)
2054        | Error::StreamsGroupBackgroundTaskClosed
2055        | Error::StreamsTaskAssignmentInvalid { .. }
2056        | Error::StreamsTaskAssignmentConflict { .. }
2057        | Error::TaskJoin(_)
2058        | Error::Protocol(_) => false,
2059    }
2060}
2061
2062fn is_offset_out_of_range(error: &Error) -> bool {
2063    matches!(
2064        error,
2065        Error::Broker { code, .. }
2066            if BrokerErrorKind::from_code(*code) == BrokerErrorKind::OffsetOutOfRange
2067    )
2068}
2069
2070fn is_leader_epoch_transition_error(error: &Error) -> bool {
2071    matches!(
2072        error,
2073        Error::Broker { code, .. }
2074            if matches!(
2075                BrokerErrorKind::from_code(*code),
2076                BrokerErrorKind::NotLeaderOrFollower
2077                    | BrokerErrorKind::FencedLeaderEpoch
2078                    | BrokerErrorKind::UnknownLeaderEpoch
2079            )
2080    )
2081}
2082
2083fn limit_fetched_records(
2084    fetched: &mut Vec<ConsumerRecord>,
2085    current_record_count: usize,
2086    max_poll_records: usize,
2087) {
2088    let remaining = max_poll_records.saturating_sub(current_record_count);
2089    fetched.truncate(remaining);
2090}
2091
2092fn invalidate_metadata_cache(
2093    metadata_cache: &mut BTreeMap<String, MetadataResponseV1>,
2094    topic: &str,
2095) {
2096    metadata_cache.remove(topic);
2097}
2098
2099// Kafka error code 100 (UNKNOWN_TOPIC_ID) is retriable for Fetch. Keep this
2100// Fetch-specific protocol classification local because the shared broker error
2101// enum intentionally does not expose every newer Kafka code yet.
2102const FETCH_UNKNOWN_TOPIC_ID_ERROR_CODE: i16 = 100;
2103
2104#[cfg(test)]
2105#[allow(clippy::expect_used, clippy::unwrap_used)]
2106mod tests {
2107    use super::{
2108        assign_partition, can_retry_fetch, invalidate_metadata_cache,
2109        is_leader_epoch_transition_error, leader_for, limit_fetched_records,
2110        offset_for_leader_epoch_partition_response, visible_records, Consumer, ConsumerAssignment,
2111        ConsumerConfig, ConsumerRecord, IsolationLevel, OffsetResetPolicy, PartitionWatermarks,
2112        SecurityProtocol,
2113    };
2114    use crate::{Client, ClientMetrics, Error};
2115    use kafrust_protocol::api::fetch::{
2116        AbortedTransactionV4, FetchPartitionResponseV4, MessageSetRecord,
2117    };
2118    use kafrust_protocol::api::metadata::{
2119        BrokerMetadata, MetadataResponseV1, PartitionMetadata, TopicMetadata,
2120    };
2121    use kafrust_protocol::codec::Encoder;
2122    use std::collections::BTreeMap;
2123    use std::time::Duration;
2124    use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
2125    use tokio::net::TcpListener;
2126    use tokio::sync::oneshot;
2127
2128    #[test]
2129    fn builds_consumer_config() {
2130        let config = ConsumerConfig::new(["localhost:9092"])
2131            .client_id("orders-reader")
2132            .client_rack("rack-a")
2133            .request_timeout_ms(5_000)
2134            .max_idle_broker_connections(3)
2135            .security_protocol(SecurityProtocol::Tls)
2136            .tls_server_name("broker.example.com")
2137            .tls_root_certificate_der([1, 2, 3])
2138            .sasl_plain("alice", "secret-password")
2139            .max_wait_ms(250)
2140            .min_bytes(10)
2141            .max_partition_bytes(1024)
2142            .max_retries(3)
2143            .max_poll_records(10)
2144            .partition_queue_capacity(7)
2145            .isolation_level(IsolationLevel::ReadCommitted);
2146
2147        assert_eq!(
2148            config.client_config().client_id_ref(),
2149            Some("orders-reader")
2150        );
2151        assert_eq!(config.client_config().client_rack_ref(), Some("rack-a"));
2152        assert_eq!(
2153            config.client_config().security_protocol_ref(),
2154            SecurityProtocol::Tls
2155        );
2156        assert_eq!(
2157            config.client_config().tls_server_name_ref(),
2158            Some("broker.example.com")
2159        );
2160        assert_eq!(
2161            config.client_config().tls_root_certificates_der(),
2162            &[vec![1, 2, 3]]
2163        );
2164        assert_eq!(
2165            config
2166                .client_config()
2167                .sasl_credentials_ref()
2168                .unwrap()
2169                .username(),
2170            "alice"
2171        );
2172        assert_eq!(config.max_retries_ref(), 3);
2173        assert_eq!(config.client_config().max_idle_broker_connections_ref(), 3);
2174        assert_eq!(config.max_poll_records_ref(), 10);
2175        assert_eq!(config.partition_queue_capacity_ref(), 7);
2176        assert_eq!(config.isolation_level_ref(), IsolationLevel::ReadCommitted);
2177        assert_eq!(
2178            config.offset_reset_policy_ref(),
2179            OffsetResetPolicy::Offset(0)
2180        );
2181    }
2182
2183    #[test]
2184    fn normalizes_zero_partition_queue_capacity() {
2185        assert_eq!(
2186            ConsumerConfig::new(["localhost:9092"])
2187                .partition_queue_capacity(0)
2188                .partition_queue_capacity_ref(),
2189            1
2190        );
2191    }
2192
2193    #[tokio::test]
2194    async fn rejects_invalid_fetch_configuration_before_connecting() {
2195        let cases = [
2196            (
2197                ConsumerConfig::new(["127.0.0.1:1"])
2198                    .max_wait_ms(-1)
2199                    .build()
2200                    .await
2201                    .unwrap_err(),
2202                "max_wait_ms",
2203            ),
2204            (
2205                ConsumerConfig::new(["127.0.0.1:1"])
2206                    .min_bytes(-1)
2207                    .build()
2208                    .await
2209                    .unwrap_err(),
2210                "min_bytes",
2211            ),
2212            (
2213                ConsumerConfig::new(["127.0.0.1:1"])
2214                    .max_partition_bytes(0)
2215                    .build()
2216                    .await
2217                    .unwrap_err(),
2218                "max_partition_bytes",
2219            ),
2220            (
2221                ConsumerConfig::new(["127.0.0.1:1"])
2222                    .max_poll_records(0)
2223                    .build()
2224                    .await
2225                    .unwrap_err(),
2226                "max_poll_records",
2227            ),
2228        ];
2229
2230        for (error, field) in cases {
2231            assert!(matches!(
2232                error,
2233                Error::InvalidConfiguration {
2234                    field: actual,
2235                    ..
2236                } if actual == field
2237            ));
2238        }
2239
2240        assert!(ConsumerConfig::new(["127.0.0.1:1"]).validate().is_ok());
2241    }
2242
2243    #[tokio::test]
2244    async fn split_partition_queue_requires_assignment_and_protects_seek_state() {
2245        let (client_stream, _broker_stream) = tokio::io::duplex(64);
2246        let client = Client::from_stream(
2247            Box::new(client_stream),
2248            Some("kafrust-partition-queue-api-test".to_owned()),
2249            Some(std::time::Duration::from_millis(500)),
2250        );
2251        let config = ConsumerConfig::new(["localhost:9092"]);
2252        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
2253
2254        assert!(matches!(
2255            consumer.split_partition_queue("orders", 0).unwrap_err(),
2256            Error::UnassignedTopicPartition { topic, partition: 0 } if topic == "orders"
2257        ));
2258        consumer.assign("orders", 0, 10);
2259        let queue = consumer.split_partition_queue("orders", 0).unwrap();
2260        assert!(matches!(
2261            consumer.split_partition_queue("orders", 0).unwrap_err(),
2262            Error::Unsupported("partition queue is already split")
2263        ));
2264        assert!(matches!(
2265            consumer.seek("orders", 0, 20).unwrap_err(),
2266            Error::Unsupported("seek requires dropping the active partition queue")
2267        ));
2268        drop(queue);
2269        consumer.seek("orders", 0, 20).unwrap();
2270        assert_eq!(consumer.position("orders", 0), Some(20));
2271    }
2272
2273    #[test]
2274    fn maps_message_set_record() {
2275        let record = ConsumerRecord::from_message_set(
2276            "orders",
2277            1,
2278            MessageSetRecord {
2279                offset: 42,
2280                leader_epoch: 4,
2281                timestamp_ms: 123,
2282                key: Some(b"order-1".to_vec()),
2283                value: Some(b"created".to_vec()),
2284                headers: vec![kafrust_protocol::api::produce::RecordBatchHeader::new(
2285                    "source",
2286                    Some(b"checkout".to_vec()),
2287                )],
2288                producer_id: None,
2289                transactional: false,
2290                control: false,
2291            },
2292        );
2293
2294        assert_eq!(record.topic(), "orders");
2295        assert_eq!(record.partition(), 1);
2296        assert_eq!(record.offset(), 42);
2297        assert_eq!(record.leader_epoch(), 4);
2298        assert_eq!(record.timestamp_ms(), 123);
2299        assert_eq!(record.key().unwrap(), b"order-1");
2300        assert_eq!(record.value().unwrap(), b"created");
2301        assert_eq!(record.headers().len(), 1);
2302        assert_eq!(record.headers()[0].key(), "source");
2303        assert_eq!(record.headers()[0].value(), Some(&b"checkout"[..]));
2304    }
2305
2306    #[test]
2307    fn preserves_null_and_empty_record_fields() {
2308        let record = ConsumerRecord::from_message_set(
2309            "orders",
2310            1,
2311            MessageSetRecord {
2312                offset: 43,
2313                leader_epoch: -1,
2314                timestamp_ms: -1,
2315                key: Some(Vec::new()),
2316                value: None,
2317                headers: vec![
2318                    kafrust_protocol::api::produce::RecordBatchHeader::new(
2319                        "empty",
2320                        Some(Vec::new()),
2321                    ),
2322                    kafrust_protocol::api::produce::RecordBatchHeader::new("null", None),
2323                ],
2324                producer_id: None,
2325                transactional: false,
2326                control: false,
2327            },
2328        );
2329
2330        assert_eq!(record.key(), Some(&[][..]));
2331        assert_eq!(record.value(), None);
2332        assert_eq!(record.timestamp_ms(), -1);
2333        assert_eq!(record.headers().len(), 2);
2334        assert_eq!(record.headers()[0].value(), Some(&[][..]));
2335        assert_eq!(record.headers()[1].value(), None);
2336    }
2337
2338    #[test]
2339    fn tracks_assignments() {
2340        let mut assignments = Vec::<ConsumerAssignment>::new();
2341
2342        assign_partition(&mut assignments, "orders".to_owned(), 0, 10);
2343        assign_partition(&mut assignments, "orders".to_owned(), 0, 20);
2344
2345        assert_eq!(assignments.len(), 1);
2346        assert_eq!(assignments[0].topic(), "orders");
2347        assert_eq!(assignments[0].partition(), 0);
2348        assert_eq!(assignments[0].next_offset(), 20);
2349        assert!(!assignments[0].is_paused());
2350    }
2351
2352    #[test]
2353    fn tracks_and_resets_assignment_leader_epoch() {
2354        let (client_stream, _broker_stream) = tokio::io::duplex(64);
2355        let client = Client::from_stream(
2356            Box::new(client_stream),
2357            Some("kafrust-leader-epoch-state-test".to_owned()),
2358            Some(std::time::Duration::from_millis(500)),
2359        );
2360        let mut consumer =
2361            Consumer::from_assignments(client, ConsumerConfig::new(["localhost:9092"]), Vec::new());
2362
2363        consumer.assign("orders", 0, 10);
2364        assert_eq!(consumer.assignments()[0].leader_epoch(), -1);
2365        consumer.update_assignment_leader_epoch("orders", 0, 4);
2366        assert_eq!(consumer.assignments()[0].leader_epoch(), 4);
2367
2368        consumer.assign("orders", 0, 20);
2369        assert_eq!(consumer.assignments()[0].leader_epoch(), -1);
2370    }
2371
2372    #[test]
2373    fn preserves_assignment_state_when_assignments_are_rebuilt() {
2374        let (client_stream, _broker_stream) = tokio::io::duplex(64);
2375        let (rebuilt_stream, _rebuilt_broker_stream) = tokio::io::duplex(64);
2376        let client = Client::from_stream(
2377            Box::new(client_stream),
2378            Some("kafrust-assignment-state-test".to_owned()),
2379            Some(std::time::Duration::from_millis(500)),
2380        );
2381        let rebuilt_client = Client::from_stream(
2382            Box::new(rebuilt_stream),
2383            Some("kafrust-assignment-state-rebuilt-test".to_owned()),
2384            Some(std::time::Duration::from_millis(500)),
2385        );
2386        let config = ConsumerConfig::new(["localhost:9092"]);
2387        let mut consumer = Consumer::from_assignments(
2388            client,
2389            config.clone(),
2390            vec![ConsumerAssignment::new("orders".to_owned(), 0, 0)],
2391        );
2392        consumer.seek("orders", 0, 42).unwrap();
2393        consumer.update_assignment_leader_epoch("orders", 0, 7);
2394        let previous_assignments = consumer.assignments().to_vec();
2395
2396        consumer.replace_assignments(vec![ConsumerAssignment::new("orders".to_owned(), 0, 0)]);
2397
2398        assert_eq!(consumer.position("orders", 0), Some(42));
2399        assert_eq!(consumer.assignments()[0].leader_epoch(), 7);
2400
2401        let mut rebuilt = Consumer::from_assignments(
2402            rebuilt_client,
2403            config,
2404            vec![ConsumerAssignment::new("orders".to_owned(), 0, 0)],
2405        );
2406        rebuilt.restore_assignment_state(&previous_assignments);
2407        assert_eq!(rebuilt.position("orders", 0), Some(42));
2408        assert_eq!(rebuilt.assignments()[0].leader_epoch(), 7);
2409    }
2410
2411    #[test]
2412    fn does_not_restore_position_for_removed_and_reassigned_partition() {
2413        let (client_stream, _broker_stream) = tokio::io::duplex(64);
2414        let client = Client::from_stream(
2415            Box::new(client_stream),
2416            Some("kafrust-assignment-reset-test".to_owned()),
2417            Some(std::time::Duration::from_millis(500)),
2418        );
2419        let config = ConsumerConfig::new(["localhost:9092"]);
2420        let mut consumer = Consumer::from_assignments(
2421            client,
2422            config,
2423            vec![ConsumerAssignment::new("orders".to_owned(), 0, 0)],
2424        );
2425        consumer.seek("orders", 0, 42).unwrap();
2426
2427        consumer.replace_assignments(vec![ConsumerAssignment::new("payments".to_owned(), 0, 7)]);
2428        consumer.replace_assignments(vec![ConsumerAssignment::new("orders".to_owned(), 0, 0)]);
2429
2430        assert_eq!(consumer.position("orders", 0), Some(0));
2431    }
2432
2433    #[tokio::test]
2434    async fn controls_assignment_position_and_pause_state() {
2435        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
2436        let addr = listener.local_addr().unwrap();
2437        let server = tokio::spawn(async move {
2438            let _connection = listener.accept().await.unwrap();
2439        });
2440        let mut consumer = ConsumerConfig::new([addr.to_string()])
2441            .build()
2442            .await
2443            .unwrap();
2444        consumer.assign("orders", 0, 10);
2445
2446        assert_eq!(consumer.position("orders", 0), Some(10));
2447        consumer.seek("orders", 0, 20).unwrap();
2448        consumer.pause("orders", 0).unwrap();
2449        assert_eq!(consumer.position("orders", 0), Some(20));
2450        assert!(consumer.assignments()[0].is_paused());
2451        assert!(consumer.poll().await.unwrap().is_empty());
2452
2453        consumer.resume("orders", 0).unwrap();
2454        assert!(!consumer.assignments()[0].is_paused());
2455        assert!(matches!(
2456            consumer.seek("orders", 1, 0).unwrap_err(),
2457            Error::UnassignedTopicPartition {
2458                topic,
2459                partition: 1
2460            } if topic == "orders"
2461        ));
2462        server.await.unwrap();
2463    }
2464
2465    #[tokio::test]
2466    async fn fetches_partition_watermarks_from_partition_leader() {
2467        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
2468        let addr = listener.local_addr().unwrap();
2469        let server = tokio::spawn(async move {
2470            let (mut socket, _) = listener.accept().await.unwrap();
2471            for (correlation_id, timestamp, offset) in [(1, -2_i64, 4_i64), (2, -1_i64, 9_i64)] {
2472                let request = read_frame(&mut socket).await;
2473                assert_eq!(&request[0..4], &[0, 2, 0, 1]);
2474                assert_eq!(
2475                    i64::from_be_bytes(request[request.len() - 8..].try_into().unwrap()),
2476                    timestamp
2477                );
2478                write_frame(
2479                    &mut socket,
2480                    &list_offsets_response_frame(correlation_id, offset),
2481                )
2482                .await;
2483            }
2484        });
2485        let (client_stream, broker_stream) = tokio::io::duplex(64);
2486        let _broker_stream = broker_stream;
2487        let client = Client::from_stream(
2488            Box::new(client_stream),
2489            Some("kafrust-watermarks-test".to_owned()),
2490            Some(std::time::Duration::from_millis(500)),
2491        );
2492        let config = ConsumerConfig::new([addr.to_string()]).request_timeout_ms(500);
2493        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
2494        let mut metadata = metadata_fixture();
2495        metadata.brokers[0].host = addr.ip().to_string();
2496        metadata.brokers[0].port = i32::from(addr.port());
2497        consumer
2498            .metadata_cache
2499            .insert("orders".to_owned(), metadata);
2500
2501        let watermarks = consumer.fetch_watermarks("orders", 0).await.unwrap();
2502
2503        assert_eq!(watermarks, PartitionWatermarks { low: 4, high: 9 });
2504        assert_eq!(watermarks.low(), 4);
2505        assert_eq!(watermarks.high(), 9);
2506        server.await.unwrap();
2507    }
2508
2509    #[tokio::test]
2510    async fn resets_out_of_range_assignment_to_earliest_offset() {
2511        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
2512        let addr = listener.local_addr().unwrap();
2513        let server = tokio::spawn(async move {
2514            let (mut first_socket, _) = listener.accept().await.unwrap();
2515            let api_versions_request = read_frame(&mut first_socket).await;
2516            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2517            write_frame(&mut first_socket, &api_versions_v3_fetch_v12_response(1)).await;
2518            let first_fetch = read_frame(&mut first_socket).await;
2519            assert_eq!(&first_fetch[0..4], &[0, 1, 0, 12]);
2520            write_frame(&mut first_socket, &fetch_v12_out_of_range_response_frame(2)).await;
2521            drop(first_socket);
2522
2523            let (mut second_socket, _) = listener.accept().await.unwrap();
2524            for (correlation_id, timestamp, offset) in [(1, -2_i64, 4_i64), (2, -1_i64, 9_i64)] {
2525                let request = read_frame(&mut second_socket).await;
2526                assert_eq!(&request[0..4], &[0, 2, 0, 1]);
2527                assert_eq!(
2528                    i64::from_be_bytes(request[request.len() - 8..].try_into().unwrap()),
2529                    timestamp
2530                );
2531                write_frame(
2532                    &mut second_socket,
2533                    &list_offsets_response_frame(correlation_id, offset),
2534                )
2535                .await;
2536            }
2537            let api_versions_request = read_frame(&mut second_socket).await;
2538            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2539            write_frame(&mut second_socket, &api_versions_v3_fetch_v12_response(3)).await;
2540            let second_fetch = read_frame(&mut second_socket).await;
2541            assert_eq!(&second_fetch[0..4], &[0, 1, 0, 12]);
2542            write_frame(
2543                &mut second_socket,
2544                &fetch_v12_response_frame_with_record(4, 0),
2545            )
2546            .await;
2547        });
2548        let (client_stream, broker_stream) = tokio::io::duplex(64);
2549        let _broker_stream = broker_stream;
2550        let client = Client::from_stream(
2551            Box::new(client_stream),
2552            Some("kafrust-offset-reset-test".to_owned()),
2553            Some(std::time::Duration::from_millis(500)),
2554        );
2555        let config = ConsumerConfig::new([addr.to_string()])
2556            .request_timeout_ms(500)
2557            .offset_reset_policy(OffsetResetPolicy::Earliest);
2558        let mut consumer = Consumer::from_assignments(
2559            client,
2560            config,
2561            vec![ConsumerAssignment::new("orders".to_owned(), 0, 100)],
2562        );
2563        let mut metadata = metadata_fixture();
2564        metadata.brokers[0].host = addr.ip().to_string();
2565        metadata.brokers[0].port = i32::from(addr.port());
2566        consumer
2567            .metadata_cache
2568            .insert("orders".to_owned(), metadata);
2569
2570        let records = consumer.poll().await.unwrap();
2571
2572        assert_eq!(records.len(), 1);
2573        assert_eq!(records[0].offset(), 42);
2574        assert_eq!(consumer.position("orders", 0), Some(43));
2575        server.await.unwrap();
2576    }
2577
2578    #[tokio::test]
2579    async fn resets_out_of_range_assignment_to_latest_offset() {
2580        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
2581        let addr = listener.local_addr().unwrap();
2582        let server = tokio::spawn(async move {
2583            let (mut first_socket, _) = listener.accept().await.unwrap();
2584            let api_versions_request = read_frame(&mut first_socket).await;
2585            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2586            write_frame(&mut first_socket, &api_versions_v3_fetch_v12_response(1)).await;
2587            let first_fetch = read_frame(&mut first_socket).await;
2588            assert_eq!(&first_fetch[0..4], &[0, 1, 0, 12]);
2589            write_frame(&mut first_socket, &fetch_v12_out_of_range_response_frame(2)).await;
2590            drop(first_socket);
2591
2592            let (mut second_socket, _) = listener.accept().await.unwrap();
2593            for (correlation_id, timestamp, offset) in [(1, -2_i64, 4_i64), (2, -1_i64, 9_i64)] {
2594                let list_offsets_request = read_frame(&mut second_socket).await;
2595                assert_eq!(&list_offsets_request[0..4], &[0, 2, 0, 1]);
2596                assert_eq!(
2597                    i64::from_be_bytes(
2598                        list_offsets_request[list_offsets_request.len() - 8..]
2599                            .try_into()
2600                            .unwrap()
2601                    ),
2602                    timestamp
2603                );
2604                write_frame(
2605                    &mut second_socket,
2606                    &list_offsets_response_frame(correlation_id, offset),
2607                )
2608                .await;
2609            }
2610
2611            let api_versions_request = read_frame(&mut second_socket).await;
2612            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2613            write_frame(&mut second_socket, &api_versions_v3_fetch_v12_response(2)).await;
2614            let second_fetch = read_frame(&mut second_socket).await;
2615            assert_eq!(&second_fetch[0..4], &[0, 1, 0, 12]);
2616            assert!(second_fetch
2617                .windows(8)
2618                .any(|window| window == 9_i64.to_be_bytes()));
2619            write_frame(&mut second_socket, &fetch_v12_response_frame(3, -1)).await;
2620        });
2621        let (client_stream, broker_stream) = tokio::io::duplex(64);
2622        let _broker_stream = broker_stream;
2623        let client = Client::from_stream(
2624            Box::new(client_stream),
2625            Some("kafrust-offset-reset-latest-test".to_owned()),
2626            Some(std::time::Duration::from_millis(500)),
2627        );
2628        let config = ConsumerConfig::new([addr.to_string()])
2629            .request_timeout_ms(500)
2630            .offset_reset_policy(OffsetResetPolicy::Latest);
2631        let mut consumer = Consumer::from_assignments(
2632            client,
2633            config,
2634            vec![ConsumerAssignment::new("orders".to_owned(), 0, 100)],
2635        );
2636        let mut metadata = metadata_fixture();
2637        metadata.brokers[0].host = addr.ip().to_string();
2638        metadata.brokers[0].port = i32::from(addr.port());
2639        consumer
2640            .metadata_cache
2641            .insert("orders".to_owned(), metadata);
2642
2643        let records = consumer.poll().await.unwrap();
2644
2645        assert!(records.is_empty());
2646        assert_eq!(consumer.position("orders", 0), Some(9));
2647        server.await.unwrap();
2648    }
2649
2650    #[tokio::test]
2651    async fn sends_assignment_leader_epoch_in_fetch_v12_request() {
2652        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
2653        let addr = listener.local_addr().unwrap();
2654        let server = tokio::spawn(async move {
2655            let (mut socket, _) = listener.accept().await.unwrap();
2656            let api_versions_request = read_frame(&mut socket).await;
2657            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2658            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
2659
2660            let fetch_request = read_frame(&mut socket).await;
2661            assert_eq!(&fetch_request[0..4], &[0, 1, 0, 12]);
2662            assert_eq!(
2663                fetch_request
2664                    .windows(4)
2665                    .filter(|window| *window == [0, 0, 0, 8])
2666                    .count(),
2667                2
2668            );
2669            write_frame(&mut socket, &fetch_v12_response_frame(2, -1)).await;
2670        });
2671        let (client_stream, broker_stream) = tokio::io::duplex(64);
2672        let _broker_stream = broker_stream;
2673        let client = Client::from_stream(
2674            Box::new(client_stream),
2675            Some("kafrust-assignment-leader-epoch-test".to_owned()),
2676            Some(std::time::Duration::from_millis(500)),
2677        );
2678        let config = ConsumerConfig::new([addr.to_string()])
2679            .request_timeout_ms(500)
2680            .client_rack("rack-a");
2681        let mut consumer = Consumer::from_assignments(
2682            client,
2683            config,
2684            vec![ConsumerAssignment::new("orders".to_owned(), 0, 42)],
2685        );
2686        let mut metadata = metadata_fixture();
2687        metadata.brokers[0].host = addr.ip().to_string();
2688        metadata.brokers[0].port = i32::from(addr.port());
2689        consumer
2690            .metadata_cache
2691            .insert("orders".to_owned(), metadata);
2692        consumer.update_assignment_leader_epoch("orders", 0, 8);
2693
2694        assert!(consumer.poll().await.unwrap().is_empty());
2695        assert_eq!(consumer.assignments()[0].leader_epoch(), 8);
2696        server.await.unwrap();
2697    }
2698
2699    #[tokio::test]
2700    async fn selects_fetch_v13_when_topic_uuid_is_cached() {
2701        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
2702        let addr = listener.local_addr().unwrap();
2703        let server = tokio::spawn(async move {
2704            let (mut socket, _) = listener.accept().await.unwrap();
2705            let api_versions_request = read_frame(&mut socket).await;
2706            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2707            write_frame(&mut socket, &api_versions_v3_fetch_v13_response(1)).await;
2708
2709            let fetch_request = read_frame(&mut socket).await;
2710            assert_eq!(&fetch_request[0..4], &[0, 1, 0, 13]);
2711            assert!(fetch_request.windows(16).any(|window| window == [7; 16]));
2712            write_frame(&mut socket, &fetch_v13_response_frame(2, [7; 16])).await;
2713        });
2714        let (client_stream, broker_stream) = tokio::io::duplex(64);
2715        let _broker_stream = broker_stream;
2716        let client = Client::from_stream(
2717            Box::new(client_stream),
2718            Some("kafrust-fetch-v13-selection-test".to_owned()),
2719            Some(std::time::Duration::from_millis(500)),
2720        );
2721        let config = ConsumerConfig::new([addr.to_string()]).request_timeout_ms(500);
2722        let mut consumer = Consumer::from_assignments(
2723            client,
2724            config,
2725            vec![ConsumerAssignment::new("orders".to_owned(), 0, 42)],
2726        );
2727        let mut metadata = metadata_fixture();
2728        metadata.brokers[0].host = addr.ip().to_string();
2729        metadata.brokers[0].port = i32::from(addr.port());
2730        consumer
2731            .metadata_cache
2732            .insert("orders".to_owned(), metadata);
2733        consumer
2734            .fetch_topic_ids
2735            .insert("orders".to_owned(), [7; 16]);
2736
2737        assert!(consumer.poll().await.unwrap().is_empty());
2738        server.await.unwrap();
2739    }
2740
2741    #[tokio::test]
2742    async fn resolves_topic_uuid_before_selecting_fetch_v13() {
2743        let (client_stream, mut broker_stream) = tokio::io::duplex(2048);
2744        let server = tokio::spawn(async move {
2745            let api_versions_request = read_frame(&mut broker_stream).await;
2746            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2747            let api_versions_correlation =
2748                i32::from_be_bytes(api_versions_request[4..8].try_into().unwrap());
2749            write_frame(
2750                &mut broker_stream,
2751                &api_versions_v3_metadata_fetch_response(api_versions_correlation, 13),
2752            )
2753            .await;
2754
2755            let metadata_request = read_frame(&mut broker_stream).await;
2756            assert_eq!(&metadata_request[0..4], &[0, 3, 0, 12]);
2757            let metadata_correlation =
2758                i32::from_be_bytes(metadata_request[4..8].try_into().unwrap());
2759            let broker_addr = "127.0.0.1:9092".parse().unwrap();
2760            write_frame(
2761                &mut broker_stream,
2762                &metadata_v12_response_frame_with_topic_id(
2763                    metadata_correlation,
2764                    &broker_addr,
2765                    5,
2766                    [7; 16],
2767                ),
2768            )
2769            .await;
2770        });
2771        let client = Client::from_stream(
2772            Box::new(client_stream),
2773            Some("kafrust-fetch-v13-metadata-test".to_owned()),
2774            Some(std::time::Duration::from_millis(500)),
2775        );
2776        let config = ConsumerConfig::new(["localhost:9092"]).request_timeout_ms(500);
2777        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
2778
2779        assert_eq!(
2780            consumer.topic_id_for_fetch("orders").await.unwrap(),
2781            Some([7; 16])
2782        );
2783        server.await.unwrap();
2784    }
2785
2786    #[tokio::test]
2787    async fn refreshes_stale_fetch_v13_topic_uuid_after_unknown_topic_id() {
2788        let (result, position, topic_id) = run_v13_unknown_topic_id_retry(true).await;
2789        let records = result.unwrap();
2790        assert_eq!(records.len(), 1);
2791        assert_eq!(records[0].offset(), 42);
2792        assert_eq!(position, Some(43));
2793        assert_eq!(topic_id, Some([2; 16]));
2794    }
2795
2796    #[tokio::test]
2797    async fn preserves_fetch_position_when_unknown_topic_id_retry_is_exhausted() {
2798        let (result, position, topic_id) = run_v13_unknown_topic_id_retry(false).await;
2799        let error = result.unwrap_err();
2800        assert!(matches!(error, Error::Broker { code: 100, .. }));
2801        assert_eq!(position, Some(42));
2802        assert_eq!(topic_id, None);
2803    }
2804
2805    async fn run_v13_unknown_topic_id_retry(
2806        final_success: bool,
2807    ) -> (
2808        std::result::Result<Vec<ConsumerRecord>, Error>,
2809        Option<i64>,
2810        Option<[u8; 16]>,
2811    ) {
2812        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
2813        let addr = listener.local_addr().unwrap();
2814        let (client_stream, mut bootstrap_stream) = tokio::io::duplex(4096);
2815        let broker = tokio::spawn(async move {
2816            let (mut first_socket, _) = listener.accept().await.unwrap();
2817            let api_versions_request = read_frame(&mut first_socket).await;
2818            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2819            let api_versions_correlation = i32::from_be_bytes(
2820                api_versions_request[4..8]
2821                    .try_into()
2822                    .expect("api versions correlation must be present"),
2823            );
2824            write_frame(
2825                &mut first_socket,
2826                &api_versions_v3_fetch_v13_response(api_versions_correlation),
2827            )
2828            .await;
2829
2830            let stale_fetch_request = read_frame(&mut first_socket).await;
2831            assert_eq!(&stale_fetch_request[0..4], &[0, 1, 0, 13]);
2832            assert!(stale_fetch_request
2833                .windows(16)
2834                .any(|window| window == [1; 16]));
2835            assert!(stale_fetch_request
2836                .windows(8)
2837                .any(|window| window == 42_i64.to_be_bytes()));
2838            let stale_fetch_correlation = i32::from_be_bytes(
2839                stale_fetch_request[4..8]
2840                    .try_into()
2841                    .expect("fetch correlation must be present"),
2842            );
2843            write_frame(
2844                &mut first_socket,
2845                &fetch_v13_response_frame_with_partition_error(
2846                    stale_fetch_correlation,
2847                    [1; 16],
2848                    100,
2849                ),
2850            )
2851            .await;
2852
2853            let (mut retry_socket, _) = listener.accept().await.unwrap();
2854            let retry_api_versions_request = read_frame(&mut retry_socket).await;
2855            assert_eq!(&retry_api_versions_request[0..4], &[0, 18, 0, 3]);
2856            let retry_api_versions_correlation = i32::from_be_bytes(
2857                retry_api_versions_request[4..8]
2858                    .try_into()
2859                    .expect("retry api versions correlation must be present"),
2860            );
2861            write_frame(
2862                &mut retry_socket,
2863                &api_versions_v3_fetch_v13_response(retry_api_versions_correlation),
2864            )
2865            .await;
2866
2867            let retry_fetch_request = read_frame(&mut retry_socket).await;
2868            assert_eq!(&retry_fetch_request[0..4], &[0, 1, 0, 13]);
2869            assert!(retry_fetch_request
2870                .windows(16)
2871                .any(|window| window == [2; 16]));
2872            assert!(retry_fetch_request
2873                .windows(8)
2874                .any(|window| window == 42_i64.to_be_bytes()));
2875            let retry_fetch_correlation = i32::from_be_bytes(
2876                retry_fetch_request[4..8]
2877                    .try_into()
2878                    .expect("retry fetch correlation must be present"),
2879            );
2880            let retry_response = if final_success {
2881                fetch_v13_response_frame_with_record_at(retry_fetch_correlation, [2; 16], 42)
2882            } else {
2883                fetch_v13_response_frame_with_partition_error(retry_fetch_correlation, [2; 16], 100)
2884            };
2885            write_frame(&mut retry_socket, &retry_response).await;
2886        });
2887        let bootstrap = tokio::spawn(async move {
2888            let metadata_request = read_frame(&mut bootstrap_stream).await;
2889            assert_eq!(&metadata_request[0..4], &[0, 3, 0, 1]);
2890            let metadata_correlation = i32::from_be_bytes(
2891                metadata_request[4..8]
2892                    .try_into()
2893                    .expect("metadata correlation must be present"),
2894            );
2895            write_frame(
2896                &mut bootstrap_stream,
2897                &metadata_response_frame_for(metadata_correlation, &addr),
2898            )
2899            .await;
2900
2901            let api_versions_request = read_frame(&mut bootstrap_stream).await;
2902            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2903            let api_versions_correlation = i32::from_be_bytes(
2904                api_versions_request[4..8]
2905                    .try_into()
2906                    .expect("bootstrap api versions correlation must be present"),
2907            );
2908            write_frame(
2909                &mut bootstrap_stream,
2910                &api_versions_v3_metadata_fetch_response(api_versions_correlation, 13),
2911            )
2912            .await;
2913
2914            let metadata_v12_request = read_frame(&mut bootstrap_stream).await;
2915            assert_eq!(&metadata_v12_request[0..4], &[0, 3, 0, 12]);
2916            let metadata_v12_correlation = i32::from_be_bytes(
2917                metadata_v12_request[4..8]
2918                    .try_into()
2919                    .expect("metadata v12 correlation must be present"),
2920            );
2921            write_frame(
2922                &mut bootstrap_stream,
2923                &metadata_v12_response_frame_with_topic_id(
2924                    metadata_v12_correlation,
2925                    &addr,
2926                    5,
2927                    [2; 16],
2928                ),
2929            )
2930            .await;
2931        });
2932
2933        let client = Client::from_stream(
2934            Box::new(client_stream),
2935            Some("kafrust-fetch-v13-unknown-topic-id-exhaustion-test".to_owned()),
2936            Some(std::time::Duration::from_millis(500)),
2937        );
2938        let config = ConsumerConfig::new([addr.to_string()])
2939            .request_timeout_ms(500)
2940            .max_retries(1);
2941        let mut consumer = Consumer::from_assignments(
2942            client,
2943            config,
2944            vec![ConsumerAssignment::new("orders".to_owned(), 0, 42)],
2945        );
2946        let mut metadata = metadata_fixture();
2947        metadata.brokers[0].host = addr.ip().to_string();
2948        metadata.brokers[0].port = i32::from(addr.port());
2949        consumer
2950            .metadata_cache
2951            .insert("orders".to_owned(), metadata);
2952        consumer
2953            .fetch_topic_ids
2954            .insert("orders".to_owned(), [1; 16]);
2955
2956        assert_eq!(consumer.position("orders", 0), Some(42));
2957        let result = consumer.poll().await;
2958        let position = consumer.position("orders", 0);
2959        let topic_id = consumer.fetch_topic_ids.get("orders").copied();
2960        tokio::time::timeout(Duration::from_secs(2), broker)
2961            .await
2962            .expect("v13 broker retry script timed out")
2963            .unwrap();
2964        tokio::time::timeout(Duration::from_secs(2), bootstrap)
2965            .await
2966            .expect("v13 metadata refresh script timed out")
2967            .unwrap();
2968        (result, position, topic_id)
2969    }
2970
2971    #[tokio::test]
2972    async fn recovers_assignment_after_leader_epoch_truncation() {
2973        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
2974        let addr = listener.local_addr().unwrap();
2975        let server = tokio::spawn(async move {
2976            let (mut bootstrap, _) =
2977                tokio::time::timeout(std::time::Duration::from_secs(2), listener.accept())
2978                    .await
2979                    .expect("bootstrap accept timed out")
2980                    .unwrap();
2981            let (mut first_fetch, _) =
2982                tokio::time::timeout(std::time::Duration::from_secs(2), listener.accept())
2983                    .await
2984                    .expect("first fetch accept timed out")
2985                    .unwrap();
2986
2987            let api_versions_request = read_frame(&mut first_fetch).await;
2988            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
2989            write_frame(
2990                &mut first_fetch,
2991                &api_versions_v3_metadata_fetch_response(1, 12),
2992            )
2993            .await;
2994            let fetch_request = read_frame(&mut first_fetch).await;
2995            assert_eq!(&fetch_request[0..4], &[0, 1, 0, 12]);
2996            assert!(fetch_request
2997                .windows(8)
2998                .any(|window| window == 100_i64.to_be_bytes()));
2999            write_frame(
3000                &mut first_fetch,
3001                &fetch_v12_response_frame_with_partition_error(2, 74),
3002            )
3003            .await;
3004
3005            let api_versions_request = read_frame(&mut bootstrap).await;
3006            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3007            write_frame(
3008                &mut bootstrap,
3009                &api_versions_v3_metadata_fetch_response(1, 12),
3010            )
3011            .await;
3012            let metadata_v12_request = read_frame(&mut bootstrap).await;
3013            assert_eq!(&metadata_v12_request[0..4], &[0, 3, 0, 12]);
3014            write_frame(&mut bootstrap, &metadata_v12_response_frame(2, &addr, 5)).await;
3015
3016            let metadata_request = read_frame(&mut bootstrap).await;
3017            assert_eq!(&metadata_request[0..4], &[0, 3, 0, 1]);
3018            write_frame(&mut bootstrap, &metadata_response_frame_for(3, &addr)).await;
3019
3020            let (mut recovery, _) = listener.accept().await.unwrap();
3021            let epoch_request = read_frame(&mut recovery).await;
3022            assert_eq!(&epoch_request[0..4], &[0, 23, 0, 3]);
3023            write_frame(
3024                &mut recovery,
3025                &offset_for_leader_epoch_response_frame_with(1, 5, 50),
3026            )
3027            .await;
3028            let api_versions_request = read_frame(&mut recovery).await;
3029            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3030            write_frame(
3031                &mut recovery,
3032                &api_versions_v3_metadata_fetch_response(2, 12),
3033            )
3034            .await;
3035            let retry_fetch_request = read_frame(&mut recovery).await;
3036            assert_eq!(&retry_fetch_request[0..4], &[0, 1, 0, 12]);
3037            assert!(retry_fetch_request
3038                .windows(8)
3039                .any(|window| window == 50_i64.to_be_bytes()));
3040            assert!(retry_fetch_request
3041                .windows(4)
3042                .any(|window| window == 5_i32.to_be_bytes()));
3043            write_frame(
3044                &mut recovery,
3045                &fetch_v12_response_frame_with_record_at(3, 0, 50),
3046            )
3047            .await;
3048        });
3049
3050        let mut consumer = ConsumerConfig::new([addr.to_string()])
3051            .request_timeout_ms(500)
3052            .max_retries(2)
3053            .build()
3054            .await
3055            .unwrap();
3056        consumer.assign("orders", 0, 100);
3057        consumer.update_assignment_leader_epoch("orders", 0, 4);
3058        let mut metadata = metadata_fixture();
3059        metadata.brokers[0].host = addr.ip().to_string();
3060        metadata.brokers[0].port = i32::from(addr.port());
3061        consumer
3062            .metadata_cache
3063            .insert("orders".to_owned(), metadata);
3064
3065        let records = tokio::time::timeout(std::time::Duration::from_secs(5), consumer.poll())
3066            .await
3067            .expect("leader epoch recovery poll timed out")
3068            .unwrap();
3069        tokio::time::timeout(std::time::Duration::from_secs(5), server)
3070            .await
3071            .expect("leader epoch recovery server timed out")
3072            .unwrap();
3073
3074        assert_eq!(records.len(), 1);
3075        assert_eq!(records[0].offset(), 50);
3076        assert_eq!(consumer.position("orders", 0), Some(51));
3077        assert_eq!(consumer.assignments()[0].leader_epoch(), 5);
3078    }
3079
3080    #[tokio::test]
3081    async fn reuses_fetch_session_for_sequential_rack_aware_polls() {
3082        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3083        let addr = listener.local_addr().unwrap();
3084        let server = tokio::spawn(async move {
3085            let (mut socket, _) = listener.accept().await.unwrap();
3086            let api_versions_request = read_frame(&mut socket).await;
3087            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3088            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3089
3090            let first_fetch = read_frame(&mut socket).await;
3091            assert_eq!(&first_fetch[0..4], &[0, 1, 0, 12]);
3092            assert!(first_fetch.windows(4).any(|window| window == [0, 0, 0, 0]));
3093            write_frame(
3094                &mut socket,
3095                &fetch_v12_response_frame_with_session(2, -1, 17),
3096            )
3097            .await;
3098
3099            let second_fetch = read_frame(&mut socket).await;
3100            assert_eq!(&second_fetch[0..4], &[0, 1, 0, 12]);
3101            assert!(second_fetch
3102                .windows(4)
3103                .any(|window| window == [0, 0, 0, 17]));
3104            assert!(second_fetch.windows(4).any(|window| window == [0, 0, 0, 1]));
3105            write_frame(
3106                &mut socket,
3107                &fetch_v12_response_frame_with_session(3, -1, 17),
3108            )
3109            .await;
3110        });
3111        let (client_stream, broker_stream) = tokio::io::duplex(64);
3112        let _broker_stream = broker_stream;
3113        let client = Client::from_stream(
3114            Box::new(client_stream),
3115            Some("kafrust-fetch-session-test".to_owned()),
3116            Some(std::time::Duration::from_millis(500)),
3117        );
3118        let config = ConsumerConfig::new([addr.to_string()])
3119            .request_timeout_ms(500)
3120            .client_rack("rack-a");
3121        let mut consumer = Consumer::from_assignments(
3122            client,
3123            config,
3124            vec![ConsumerAssignment::new("orders".to_owned(), 0, 42)],
3125        );
3126        let mut metadata = metadata_fixture();
3127        metadata.brokers[0].host = addr.ip().to_string();
3128        metadata.brokers[0].port = i32::from(addr.port());
3129        consumer
3130            .metadata_cache
3131            .insert("orders".to_owned(), metadata);
3132
3133        assert!(consumer.poll().await.unwrap().is_empty());
3134        assert_eq!(consumer.fetch_sessions[&addr.to_string()].session_id, 17);
3135        assert_eq!(consumer.fetch_sessions[&addr.to_string()].next_epoch, 1);
3136        assert!(consumer.poll().await.unwrap().is_empty());
3137        assert_eq!(consumer.fetch_sessions[&addr.to_string()].next_epoch, 2);
3138        server.await.unwrap();
3139    }
3140
3141    #[tokio::test]
3142    async fn resets_invalid_fetch_session_epoch_before_retrying() {
3143        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3144        let addr = listener.local_addr().unwrap();
3145        let server = tokio::spawn(async move {
3146            let (mut first_socket, _) = listener.accept().await.unwrap();
3147            let api_versions_request = read_frame(&mut first_socket).await;
3148            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3149            write_frame(&mut first_socket, &api_versions_v3_fetch_v12_response(1)).await;
3150
3151            let first_fetch = read_frame(&mut first_socket).await;
3152            assert_eq!(&first_fetch[0..4], &[0, 1, 0, 12]);
3153            assert!(first_fetch.windows(4).any(|window| window == [0, 0, 0, 0]));
3154            write_frame(
3155                &mut first_socket,
3156                &fetch_v12_response_frame_with_session(2, -1, 17),
3157            )
3158            .await;
3159
3160            let stale_fetch = read_frame(&mut first_socket).await;
3161            assert_eq!(&stale_fetch[0..4], &[0, 1, 0, 12]);
3162            assert!(stale_fetch.windows(4).any(|window| window == [0, 0, 0, 17]));
3163            assert!(stale_fetch.windows(4).any(|window| window == [0, 0, 0, 1]));
3164            write_frame(
3165                &mut first_socket,
3166                &fetch_v12_response_frame_with_partition_error_and_session(3, 70, 17),
3167            )
3168            .await;
3169
3170            let (mut reset_socket, _) = listener.accept().await.unwrap();
3171            let metadata_request = read_frame(&mut reset_socket).await;
3172            assert_eq!(&metadata_request[0..4], &[0, 3, 0, 1]);
3173            let metadata_correlation = i32::from_be_bytes(
3174                metadata_request[4..8]
3175                    .try_into()
3176                    .expect("metadata correlation must be present"),
3177            );
3178            write_frame(
3179                &mut reset_socket,
3180                &metadata_response_frame_for(metadata_correlation, &addr),
3181            )
3182            .await;
3183            let (mut reset_leader_socket, _) = listener.accept().await.unwrap();
3184            let reset_api_versions = read_frame(&mut reset_leader_socket).await;
3185            assert_eq!(&reset_api_versions[0..4], &[0, 18, 0, 3]);
3186            write_frame(
3187                &mut reset_leader_socket,
3188                &api_versions_v3_fetch_v12_response(4),
3189            )
3190            .await;
3191            let reset_fetch = read_frame(&mut reset_leader_socket).await;
3192            assert_eq!(&reset_fetch[0..4], &[0, 1, 0, 12]);
3193            assert!(reset_fetch.windows(4).any(|window| window == [0, 0, 0, 0]));
3194            write_frame(
3195                &mut reset_leader_socket,
3196                &fetch_v12_response_frame_with_record(5, 0),
3197            )
3198            .await;
3199        });
3200
3201        let (client_stream, broker_stream) = tokio::io::duplex(64);
3202        let _broker_stream = broker_stream;
3203        let client = Client::from_stream(
3204            Box::new(client_stream),
3205            Some("kafrust-fetch-session-reset-test".to_owned()),
3206            Some(std::time::Duration::from_millis(500)),
3207        );
3208        let config = ConsumerConfig::new([addr.to_string()])
3209            .request_timeout_ms(500)
3210            .max_retries(1);
3211        let mut consumer = Consumer::from_assignments(
3212            client,
3213            config,
3214            vec![ConsumerAssignment::new("orders".to_owned(), 0, 42)],
3215        );
3216        let mut metadata = metadata_fixture();
3217        metadata.brokers[0].host = addr.ip().to_string();
3218        metadata.brokers[0].port = i32::from(addr.port());
3219        consumer
3220            .metadata_cache
3221            .insert("orders".to_owned(), metadata);
3222
3223        assert!(consumer.poll().await.unwrap().is_empty());
3224        let records = consumer.poll().await.unwrap();
3225        assert_eq!(records.len(), 1);
3226        assert_eq!(records[0].offset(), 42);
3227        assert!(consumer.fetch_sessions.is_empty());
3228        server.await.unwrap();
3229    }
3230
3231    #[tokio::test]
3232    async fn fetches_offset_for_leader_epoch_from_partition_leader() {
3233        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3234        let addr = listener.local_addr().unwrap();
3235        let server = tokio::spawn(async move {
3236            let (mut socket, _) = listener.accept().await.unwrap();
3237            let request = read_frame(&mut socket).await;
3238            assert_eq!(&request[0..4], &[0, 23, 0, 3]);
3239            write_frame(&mut socket, &offset_for_leader_epoch_response_frame()).await;
3240        });
3241        let (client_stream, broker_stream) = tokio::io::duplex(64);
3242        let _broker_stream = broker_stream;
3243        let client = Client::from_stream(
3244            Box::new(client_stream),
3245            Some("kafrust-leader-epoch-test".to_owned()),
3246            Some(std::time::Duration::from_millis(500)),
3247        );
3248        let config = ConsumerConfig::new([addr.to_string()]).request_timeout_ms(500);
3249        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3250        let mut metadata = metadata_fixture();
3251        metadata.brokers[0].host = addr.ip().to_string();
3252        metadata.brokers[0].port = i32::from(addr.port());
3253        consumer
3254            .metadata_cache
3255            .insert("orders".to_owned(), metadata);
3256
3257        let result = consumer
3258            .offset_for_leader_epoch("orders", 0, 9, 7)
3259            .await
3260            .unwrap();
3261
3262        assert_eq!(result.leader_epoch(), 8);
3263        assert_eq!(result.end_offset(), 42);
3264        server.await.unwrap();
3265    }
3266
3267    #[test]
3268    fn reports_missing_offset_for_leader_epoch_partition() {
3269        let error = offset_for_leader_epoch_partition_response(&[], "orders", 0).unwrap_err();
3270
3271        assert!(matches!(
3272            error,
3273            Error::UnknownTopicOrPartition { topic, partition: 0 } if topic == "orders"
3274        ));
3275    }
3276
3277    #[test]
3278    fn resolves_partition_leader() {
3279        assert_eq!(leader_for(&metadata_fixture(), "orders", 0).unwrap(), 1);
3280    }
3281
3282    #[test]
3283    fn classifies_retriable_fetch_errors() {
3284        assert!(can_retry_fetch(&Error::Broker {
3285            code: 6,
3286            context: "fetch orders-0@0".to_owned(),
3287        }));
3288        assert!(can_retry_fetch(&Error::RequestTimedOut { timeout_ms: 5 }));
3289        assert!(!can_retry_fetch(&Error::DeliveryDeadlineExceeded {
3290            phase: crate::DeliveryPhase::Produce,
3291            possibly_transmitted: true,
3292            timeout_ms: 5,
3293        }));
3294        assert!(can_retry_fetch(&Error::Io(std::io::Error::new(
3295            std::io::ErrorKind::ConnectionReset,
3296            "reset",
3297        ))));
3298        assert!(can_retry_fetch(&Error::UnknownTopicOrPartition {
3299            topic: "orders".to_owned(),
3300            partition: 3,
3301        }));
3302        assert!(can_retry_fetch(&Error::MissingLeader {
3303            topic: "orders".to_owned(),
3304            partition: 0,
3305        }));
3306        assert!(can_retry_fetch(&Error::MissingBroker { node_id: 2 }));
3307        assert!(can_retry_fetch(&Error::Broker {
3308            code: 74,
3309            context: "fetch orders-0@0".to_owned(),
3310        }));
3311        assert!(can_retry_fetch(&Error::Broker {
3312            code: 75,
3313            context: "fetch orders-0@0".to_owned(),
3314        }));
3315        assert!(can_retry_fetch(&Error::Broker {
3316            code: 70,
3317            context: "fetch orders-0@0".to_owned(),
3318        }));
3319        assert!(can_retry_fetch(&Error::Broker {
3320            code: 100,
3321            context: "fetch orders-0@0".to_owned(),
3322        }));
3323        assert!(!can_retry_fetch(&Error::Broker {
3324            code: 1,
3325            context: "fetch orders-0@0".to_owned(),
3326        }));
3327        assert!(!can_retry_fetch(&Error::Unsupported("fetch v99")));
3328        assert!(!can_retry_fetch(&Error::OAuthBearerTokenExpired));
3329    }
3330
3331    #[test]
3332    fn classifies_not_leader_as_a_leader_epoch_transition() {
3333        assert!(is_leader_epoch_transition_error(&Error::Broker {
3334            code: 6,
3335            context: "fetch orders-0@0".to_owned(),
3336        }));
3337        assert!(is_leader_epoch_transition_error(&Error::Broker {
3338            code: 74,
3339            context: "fetch orders-0@0".to_owned(),
3340        }));
3341        assert!(is_leader_epoch_transition_error(&Error::Broker {
3342            code: 75,
3343            context: "fetch orders-0@0".to_owned(),
3344        }));
3345        assert!(!is_leader_epoch_transition_error(&Error::Broker {
3346            code: 1,
3347            context: "fetch orders-0@0".to_owned(),
3348        }));
3349    }
3350
3351    #[test]
3352    fn limits_fetched_records_to_remaining_poll_budget() {
3353        let mut records = vec![
3354            ConsumerRecord::from_message_set("orders", 0, message(10)),
3355            ConsumerRecord::from_message_set("orders", 0, message(11)),
3356            ConsumerRecord::from_message_set("orders", 0, message(12)),
3357        ];
3358
3359        limit_fetched_records(&mut records, 1, 3);
3360
3361        assert_eq!(records.len(), 2);
3362        assert_eq!(records[1].offset(), 11);
3363    }
3364
3365    #[test]
3366    fn clears_fetched_records_when_poll_budget_is_exhausted() {
3367        let mut records = vec![ConsumerRecord::from_message_set("orders", 0, message(10))];
3368
3369        limit_fetched_records(&mut records, 3, 3);
3370
3371        assert!(records.is_empty());
3372    }
3373
3374    #[test]
3375    fn invalidates_topic_metadata_cache() {
3376        let mut cache = BTreeMap::new();
3377        cache.insert("orders".to_owned(), metadata_fixture());
3378        cache.insert("payments".to_owned(), metadata_fixture());
3379
3380        invalidate_metadata_cache(&mut cache, "orders");
3381
3382        assert!(!cache.contains_key("orders"));
3383        assert!(cache.contains_key("payments"));
3384    }
3385
3386    #[test]
3387    fn read_committed_hides_aborted_records_and_control_markers() {
3388        let partition = FetchPartitionResponseV4 {
3389            partition_index: 0,
3390            error_code: 0,
3391            high_watermark: 14,
3392            last_stable_offset: 14,
3393            aborted_transactions: vec![AbortedTransactionV4 {
3394                producer_id: 7,
3395                first_offset: 10,
3396            }],
3397            records: vec![
3398                transactional_message(10, 7, false),
3399                transactional_message(11, 8, false),
3400                transactional_message(12, 7, true),
3401                transactional_message(13, 7, false),
3402            ],
3403        };
3404
3405        let committed = visible_records(
3406            &partition.aborted_transactions,
3407            &partition.records,
3408            IsolationLevel::ReadCommitted,
3409        );
3410        let uncommitted = visible_records(
3411            &partition.aborted_transactions,
3412            &partition.records,
3413            IsolationLevel::ReadUncommitted,
3414        );
3415
3416        assert_eq!(
3417            committed
3418                .iter()
3419                .map(|record| record.offset)
3420                .collect::<Vec<_>>(),
3421            vec![11, 13]
3422        );
3423        assert_eq!(
3424            uncommitted
3425                .iter()
3426                .map(|record| record.offset)
3427                .collect::<Vec<_>>(),
3428            vec![10, 11, 13]
3429        );
3430    }
3431
3432    #[tokio::test]
3433    async fn reconnects_metadata_client_after_request_io_error() {
3434        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3435        let addr = listener.local_addr().unwrap();
3436        let server = tokio::spawn(async move {
3437            let (mut socket, _) = listener.accept().await.unwrap();
3438            let request = read_frame(&mut socket).await;
3439            assert_eq!(&request[0..2], &[0, 3]);
3440            write_frame(&mut socket, &metadata_response_frame()).await;
3441        });
3442
3443        let (client_stream, broker_stream) = tokio::io::duplex(64);
3444        drop(broker_stream);
3445        let client = Client::from_stream(
3446            Box::new(client_stream),
3447            Some("kafrust-consumer-test".to_owned()),
3448            Some(std::time::Duration::from_millis(50)),
3449        );
3450        let metrics = ClientMetrics::new();
3451        let config = ConsumerConfig::new([addr.to_string()])
3452            .request_timeout_ms(500)
3453            .metrics(metrics.clone());
3454        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3455
3456        let metadata = consumer.metadata_for_topic("orders").await.unwrap();
3457
3458        assert_eq!(metadata.brokers[0].node_id, 1);
3459        assert_eq!(metadata.topics[0].name, "orders");
3460        assert!(consumer.metadata_cache.contains_key("orders"));
3461        assert_eq!(metrics.snapshot().retries, 1);
3462        server.await.unwrap();
3463    }
3464
3465    #[tokio::test]
3466    async fn records_consumed_metrics_for_fetch_results() {
3467        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3468        let addr = listener.local_addr().unwrap();
3469        let server = tokio::spawn(async move {
3470            let (mut socket, _) = listener.accept().await.unwrap();
3471            let api_versions_request = read_frame(&mut socket).await;
3472            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3473            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3474            let request = read_frame(&mut socket).await;
3475            assert_eq!(&request[0..4], &[0, 1, 0, 12]);
3476            write_frame(&mut socket, &fetch_v12_response_frame_with_record(2, 0)).await;
3477        });
3478        let (client_stream, broker_stream) = tokio::io::duplex(64);
3479        let _broker_stream = broker_stream;
3480        let client = Client::from_stream(
3481            Box::new(client_stream),
3482            Some("kafrust-consumer-metrics-test".to_owned()),
3483            Some(std::time::Duration::from_millis(500)),
3484        );
3485        let metrics = ClientMetrics::new();
3486        let config = ConsumerConfig::new([addr.to_string()])
3487            .request_timeout_ms(500)
3488            .metrics(metrics.clone());
3489        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3490        let mut metadata = metadata_fixture();
3491        metadata.brokers[0].host = addr.ip().to_string();
3492        metadata.brokers[0].port = i32::from(addr.port());
3493        consumer
3494            .metadata_cache
3495            .insert("orders".to_owned(), metadata);
3496
3497        let records = consumer.fetch("orders", 0, 42).await.unwrap();
3498
3499        assert_eq!(records.len(), 1);
3500        assert_eq!(records[0].offset(), 42);
3501        assert_eq!(metrics.snapshot().consumed_records, 1);
3502        server.await.unwrap();
3503    }
3504
3505    #[tokio::test]
3506    async fn cancels_fetch_after_transmission_and_reconnects_without_reusing_session() {
3507        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3508        let addr = listener.local_addr().unwrap();
3509        let (fetch_seen_tx, fetch_seen_rx) = oneshot::channel();
3510        let server = tokio::spawn(async move {
3511            let (mut first_socket, _) = listener.accept().await.unwrap();
3512            let api_versions_request = read_frame(&mut first_socket).await;
3513            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3514            write_frame(&mut first_socket, &api_versions_v3_fetch_v12_response(1)).await;
3515            let fetch_request = read_frame(&mut first_socket).await;
3516            assert_eq!(&fetch_request[0..4], &[0, 1, 0, 12]);
3517            let _ = fetch_seen_tx.send(());
3518
3519            let (mut second_socket, _) = listener.accept().await.unwrap();
3520            let second_api_versions_request = read_frame(&mut second_socket).await;
3521            assert_eq!(&second_api_versions_request[0..4], &[0, 18, 0, 3]);
3522            write_frame(&mut second_socket, &api_versions_v3_fetch_v12_response(1)).await;
3523            let second_fetch_request = read_frame(&mut second_socket).await;
3524            assert_eq!(&second_fetch_request[0..4], &[0, 1, 0, 12]);
3525            write_frame(
3526                &mut second_socket,
3527                &fetch_v12_response_frame_with_record(2, 0),
3528            )
3529            .await;
3530        });
3531
3532        let (client_stream, broker_stream) = tokio::io::duplex(64);
3533        let _broker_stream = broker_stream;
3534        let client = Client::from_stream(
3535            Box::new(client_stream),
3536            Some("kafrust-consumer-cancel-test".to_owned()),
3537            Some(Duration::from_millis(500)),
3538        );
3539        let config = ConsumerConfig::new([addr.to_string()])
3540            .request_timeout_ms(500)
3541            .max_retries(0);
3542        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3543        let mut metadata = metadata_fixture();
3544        metadata.brokers[0].host = addr.ip().to_string();
3545        metadata.brokers[0].port = i32::from(addr.port());
3546        consumer
3547            .metadata_cache
3548            .insert("orders".to_owned(), metadata);
3549
3550        let mut canceled_fetch = Box::pin(consumer.fetch("orders", 0, 42));
3551        tokio::select! {
3552            result = &mut canceled_fetch => assert!(result.is_err(), "fetch completed before cancellation"),
3553            _ = fetch_seen_rx => {}
3554        }
3555        drop(canceled_fetch);
3556
3557        assert!(consumer.fetch_sessions.is_empty());
3558        let records = consumer.fetch("orders", 0, 42).await.unwrap();
3559        assert_eq!(records.len(), 1);
3560        assert_eq!(records[0].offset(), 42);
3561
3562        server.await.unwrap();
3563    }
3564
3565    #[tokio::test]
3566    async fn split_partition_queue_routes_poll_records_and_advances_position() {
3567        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3568        let addr = listener.local_addr().unwrap();
3569        let server = tokio::spawn(async move {
3570            let (mut socket, _) = listener.accept().await.unwrap();
3571            let api_versions_request = read_frame(&mut socket).await;
3572            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3573            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3574            let request = read_frame(&mut socket).await;
3575            assert_eq!(&request[0..4], &[0, 1, 0, 12]);
3576            write_frame(&mut socket, &fetch_v12_response_frame_with_record(2, 0)).await;
3577        });
3578
3579        let (client_stream, broker_stream) = tokio::io::duplex(64);
3580        let _broker_stream = broker_stream;
3581        let client = Client::from_stream(
3582            Box::new(client_stream),
3583            Some("kafrust-partition-queue-test".to_owned()),
3584            Some(std::time::Duration::from_millis(500)),
3585        );
3586        let config = ConsumerConfig::new([addr.to_string()])
3587            .request_timeout_ms(500)
3588            .partition_queue_capacity(2);
3589        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3590        let mut metadata = metadata_fixture();
3591        metadata.brokers[0].host = addr.ip().to_string();
3592        metadata.brokers[0].port = i32::from(addr.port());
3593        consumer
3594            .metadata_cache
3595            .insert("orders".to_owned(), metadata);
3596        consumer.assign("orders", 0, 42);
3597
3598        let mut queue = consumer.split_partition_queue("orders", 0).unwrap();
3599        assert_eq!(queue.topic(), "orders");
3600        assert_eq!(queue.partition(), 0);
3601        assert!(consumer.poll().await.unwrap().is_empty());
3602        assert_eq!(consumer.position("orders", 0), Some(43));
3603
3604        let record = queue.recv().await.unwrap();
3605        assert_eq!(record.offset(), 42);
3606        assert!(queue.try_recv().is_none());
3607        server.await.unwrap();
3608    }
3609
3610    #[tokio::test]
3611    async fn split_partition_queue_falls_back_to_poll_when_receiver_is_dropped() {
3612        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3613        let addr = listener.local_addr().unwrap();
3614        let server = tokio::spawn(async move {
3615            let (mut socket, _) = listener.accept().await.unwrap();
3616            let api_versions_request = read_frame(&mut socket).await;
3617            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3618            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3619            let request = read_frame(&mut socket).await;
3620            assert_eq!(&request[0..4], &[0, 1, 0, 12]);
3621            write_frame(&mut socket, &fetch_v12_response_frame_with_record(2, 0)).await;
3622        });
3623
3624        let (client_stream, broker_stream) = tokio::io::duplex(64);
3625        let _broker_stream = broker_stream;
3626        let client = Client::from_stream(
3627            Box::new(client_stream),
3628            Some("kafrust-partition-queue-drop-test".to_owned()),
3629            Some(std::time::Duration::from_millis(500)),
3630        );
3631        let config = ConsumerConfig::new([addr.to_string()])
3632            .request_timeout_ms(500)
3633            .partition_queue_capacity(2);
3634        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3635        let mut metadata = metadata_fixture();
3636        metadata.brokers[0].host = addr.ip().to_string();
3637        metadata.brokers[0].port = i32::from(addr.port());
3638        consumer
3639            .metadata_cache
3640            .insert("orders".to_owned(), metadata);
3641        consumer.assign("orders", 0, 42);
3642        let queue = consumer.split_partition_queue("orders", 0).unwrap();
3643        drop(queue);
3644
3645        let records = consumer.poll().await.unwrap();
3646
3647        assert_eq!(records.len(), 1);
3648        assert_eq!(records[0].offset(), 42);
3649        assert_eq!(consumer.position("orders", 0), Some(43));
3650        server.await.unwrap();
3651    }
3652
3653    #[tokio::test]
3654    async fn split_partition_queue_reports_backpressure_without_skipping_records() {
3655        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3656        let addr = listener.local_addr().unwrap();
3657        let server = tokio::spawn(async move {
3658            let (mut socket, _) = listener.accept().await.unwrap();
3659            let api_versions_request = read_frame(&mut socket).await;
3660            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3661            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3662            for correlation_id in [2, 3] {
3663                let request = read_frame(&mut socket).await;
3664                assert_eq!(&request[0..4], &[0, 1, 0, 12]);
3665                write_frame(
3666                    &mut socket,
3667                    &fetch_v12_response_frame_with_record(correlation_id, 0),
3668                )
3669                .await;
3670            }
3671        });
3672
3673        let (client_stream, broker_stream) = tokio::io::duplex(64);
3674        let _broker_stream = broker_stream;
3675        let client = Client::from_stream(
3676            Box::new(client_stream),
3677            Some("kafrust-partition-queue-capacity-test".to_owned()),
3678            Some(std::time::Duration::from_millis(500)),
3679        );
3680        let config = ConsumerConfig::new([addr.to_string()])
3681            .request_timeout_ms(500)
3682            .partition_queue_capacity(1);
3683        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3684        let mut metadata = metadata_fixture();
3685        metadata.brokers[0].host = addr.ip().to_string();
3686        metadata.brokers[0].port = i32::from(addr.port());
3687        consumer
3688            .metadata_cache
3689            .insert("orders".to_owned(), metadata);
3690        consumer.assign("orders", 0, 42);
3691        let mut queue = consumer.split_partition_queue("orders", 0).unwrap();
3692
3693        consumer.poll().await.unwrap();
3694        let error = consumer.poll().await.unwrap_err();
3695        assert!(matches!(
3696            error,
3697            Error::PartitionQueueFull {
3698                topic,
3699                partition: 0,
3700                capacity: 1
3701            } if topic == "orders"
3702        ));
3703        assert_eq!(consumer.position("orders", 0), Some(43));
3704        assert_eq!(queue.recv().await.unwrap().offset(), 42);
3705        server.await.unwrap();
3706    }
3707
3708    #[tokio::test]
3709    async fn poll_rotates_partitions_when_max_poll_records_is_one() {
3710        let first_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3711        let first_addr = first_listener.local_addr().unwrap();
3712        let first_server = tokio::spawn(async move {
3713            let (mut socket, _) = first_listener.accept().await.unwrap();
3714            let api_versions_request = read_frame(&mut socket).await;
3715            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3716            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3717            let request = read_frame(&mut socket).await;
3718            assert_eq!(&request[0..4], &[0, 1, 0, 12]);
3719            write_frame(
3720                &mut socket,
3721                &fetch_v12_response_frame_with_record_at_partition(2, 0, 0, 42),
3722            )
3723            .await;
3724        });
3725        let second_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3726        let second_addr = second_listener.local_addr().unwrap();
3727        let second_server = tokio::spawn(async move {
3728            let (mut socket, _) = second_listener.accept().await.unwrap();
3729            let api_versions_request = read_frame(&mut socket).await;
3730            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3731            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3732            let request = read_frame(&mut socket).await;
3733            assert_eq!(&request[0..4], &[0, 1, 0, 12]);
3734            write_frame(
3735                &mut socket,
3736                &fetch_v12_response_frame_with_record_at_partition(2, 0, 1, 42),
3737            )
3738            .await;
3739        });
3740
3741        let (client_stream, broker_stream) = tokio::io::duplex(64);
3742        let _broker_stream = broker_stream;
3743        let client = Client::from_stream(
3744            Box::new(client_stream),
3745            Some("kafrust-partition-poll-fairness-test".to_owned()),
3746            Some(std::time::Duration::from_millis(500)),
3747        );
3748        let config = ConsumerConfig::new([first_addr.to_string()])
3749            .request_timeout_ms(500)
3750            .max_poll_records(1);
3751        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3752        let mut metadata = metadata_fixture();
3753        metadata.brokers[0].host = first_addr.ip().to_string();
3754        metadata.brokers[0].port = i32::from(first_addr.port());
3755        metadata.brokers.push(BrokerMetadata {
3756            node_id: 2,
3757            host: second_addr.ip().to_string(),
3758            port: i32::from(second_addr.port()),
3759            rack: None,
3760        });
3761        metadata.topics[0].partitions.push(PartitionMetadata {
3762            error_code: 0,
3763            partition_index: 1,
3764            leader_id: 2,
3765            replica_nodes: vec![2],
3766            isr_nodes: vec![2],
3767        });
3768        consumer
3769            .metadata_cache
3770            .insert("orders".to_owned(), metadata);
3771        consumer.assign("orders", 0, 42);
3772        consumer.assign("orders", 1, 42);
3773
3774        let first = consumer.poll().await.unwrap();
3775        assert_eq!(first.len(), 1);
3776        assert_eq!(first[0].partition(), 0);
3777        let second = consumer.poll().await.unwrap();
3778        assert_eq!(second.len(), 1);
3779        assert_eq!(second[0].partition(), 1);
3780        first_server.await.unwrap();
3781        second_server.await.unwrap();
3782    }
3783
3784    #[tokio::test]
3785    async fn reuses_partition_leader_connection_for_sequential_fetches() {
3786        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3787        let addr = listener.local_addr().unwrap();
3788        let server = tokio::spawn(async move {
3789            let (mut socket, _) = listener.accept().await.unwrap();
3790            let api_versions_request = read_frame(&mut socket).await;
3791            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3792            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3793            for (correlation_id, session_id) in [(2, 17), (3, 17)] {
3794                let request = read_frame(&mut socket).await;
3795                assert_eq!(&request[0..4], &[0, 1, 0, 12]);
3796                write_frame(
3797                    &mut socket,
3798                    &fetch_v12_response_frame_with_record(correlation_id, session_id),
3799                )
3800                .await;
3801            }
3802        });
3803
3804        let (client_stream, broker_stream) = tokio::io::duplex(64);
3805        let _broker_stream = broker_stream;
3806        let client = Client::from_stream(
3807            Box::new(client_stream),
3808            Some("kafrust-consumer-reuse-test".to_owned()),
3809            Some(std::time::Duration::from_millis(500)),
3810        );
3811        let config = ConsumerConfig::new([addr.to_string()]).request_timeout_ms(500);
3812        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3813        let mut metadata = metadata_fixture();
3814        metadata.brokers[0].host = addr.ip().to_string();
3815        metadata.brokers[0].port = i32::from(addr.port());
3816        consumer
3817            .metadata_cache
3818            .insert("orders".to_owned(), metadata);
3819
3820        let first = consumer.fetch("orders", 0, 42).await.unwrap();
3821        let second = consumer.fetch("orders", 0, 42).await.unwrap();
3822
3823        assert_eq!(first.len(), 1);
3824        assert_eq!(second.len(), 1);
3825        assert_eq!(consumer.broker_clients.len(), 1);
3826        assert_eq!(consumer.fetch_sessions[&addr.to_string()].session_id, 17);
3827        assert_eq!(consumer.fetch_sessions[&addr.to_string()].next_epoch, 2);
3828        server.await.unwrap();
3829    }
3830
3831    #[tokio::test]
3832    async fn falls_back_to_fetch_v4_when_broker_lacks_session_versions() {
3833        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3834        let addr = listener.local_addr().unwrap();
3835        let server = tokio::spawn(async move {
3836            let (mut socket, _) = listener.accept().await.unwrap();
3837            let api_versions_request = read_frame(&mut socket).await;
3838            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3839            write_frame(&mut socket, &api_versions_v3_fetch_response(1, 10)).await;
3840
3841            let fetch_request = read_frame(&mut socket).await;
3842            assert_eq!(&fetch_request[0..4], &[0, 1, 0, 4]);
3843            write_frame(&mut socket, &fetch_v4_response_frame()).await;
3844        });
3845        let (client_stream, broker_stream) = tokio::io::duplex(64);
3846        let _broker_stream = broker_stream;
3847        let client = Client::from_stream(
3848            Box::new(client_stream),
3849            Some("kafrust-fetch-v4-fallback-test".to_owned()),
3850            Some(std::time::Duration::from_millis(500)),
3851        );
3852        let config = ConsumerConfig::new([addr.to_string()]).request_timeout_ms(500);
3853        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3854        let mut metadata = metadata_fixture();
3855        metadata.brokers[0].host = addr.ip().to_string();
3856        metadata.brokers[0].port = i32::from(addr.port());
3857        consumer
3858            .metadata_cache
3859            .insert("orders".to_owned(), metadata);
3860
3861        let records = consumer.fetch("orders", 0, 42).await.unwrap();
3862
3863        assert_eq!(records.len(), 1);
3864        assert!(consumer.fetch_sessions.is_empty());
3865        server.await.unwrap();
3866    }
3867
3868    #[tokio::test]
3869    async fn negotiates_rack_aware_fetch_and_routes_to_preferred_replica() {
3870        let leader_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3871        let leader_addr = leader_listener.local_addr().unwrap();
3872        let preferred_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3873        let preferred_addr = preferred_listener.local_addr().unwrap();
3874
3875        let leader_server = tokio::spawn(async move {
3876            let (mut socket, _) = leader_listener.accept().await.unwrap();
3877            let api_versions_request = read_frame(&mut socket).await;
3878            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3879            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3880
3881            let fetch_request = read_frame(&mut socket).await;
3882            assert_eq!(&fetch_request[0..4], &[0, 1, 0, 12]);
3883            assert!(fetch_request
3884                .windows(b"rack-a".len())
3885                .any(|window| window == b"rack-a"));
3886            assert_eq!(fetch_request.last(), Some(&0));
3887            write_frame(&mut socket, &fetch_v12_response_frame(2, 2)).await;
3888        });
3889        let preferred_server = tokio::spawn(async move {
3890            let (mut socket, _) = preferred_listener.accept().await.unwrap();
3891            let api_versions_request = read_frame(&mut socket).await;
3892            assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
3893            write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
3894
3895            let fetch_request = read_frame(&mut socket).await;
3896            assert_eq!(&fetch_request[0..4], &[0, 1, 0, 12]);
3897            assert!(fetch_request
3898                .windows(b"rack-a".len())
3899                .any(|window| window == b"rack-a"));
3900            assert_eq!(fetch_request.last(), Some(&0));
3901            write_frame(&mut socket, &fetch_v12_response_frame(2, -1)).await;
3902        });
3903
3904        let (client_stream, broker_stream) = tokio::io::duplex(64);
3905        let _broker_stream = broker_stream;
3906        let client = Client::from_stream(
3907            Box::new(client_stream),
3908            Some("kafrust-rack-routing-test".to_owned()),
3909            Some(std::time::Duration::from_millis(500)),
3910        );
3911        let config = ConsumerConfig::new([leader_addr.to_string()])
3912            .request_timeout_ms(500)
3913            .client_rack("rack-a");
3914        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3915        let mut metadata = metadata_fixture();
3916        metadata.brokers[0].host = leader_addr.ip().to_string();
3917        metadata.brokers[0].port = i32::from(leader_addr.port());
3918        metadata.brokers.push(BrokerMetadata {
3919            node_id: 2,
3920            host: preferred_addr.ip().to_string(),
3921            port: i32::from(preferred_addr.port()),
3922            rack: Some("rack-a".to_owned()),
3923        });
3924        consumer
3925            .metadata_cache
3926            .insert("orders".to_owned(), metadata);
3927
3928        assert!(consumer.fetch("orders", 0, 42).await.unwrap().is_empty());
3929        assert_eq!(
3930            consumer
3931                .preferred_read_replicas
3932                .get(&("orders".to_owned(), 0)),
3933            Some(&2)
3934        );
3935        assert!(consumer.fetch("orders", 0, 42).await.unwrap().is_empty());
3936        assert!(!consumer
3937            .preferred_read_replicas
3938            .contains_key(&("orders".to_owned(), 0)));
3939
3940        leader_server.await.unwrap();
3941        preferred_server.await.unwrap();
3942    }
3943
3944    #[tokio::test]
3945    async fn clears_preferred_replica_after_exhausted_fetch_failure() {
3946        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3947        let addr = listener.local_addr().unwrap();
3948        let server = tokio::spawn(async move {
3949            let (mut socket, _) = listener.accept().await.unwrap();
3950            let _api_versions_request = read_frame(&mut socket).await;
3951            write_frame(&mut socket, &api_versions_v3_fetch_v11_response(1)).await;
3952            let _fetch_request = read_frame(&mut socket).await;
3953            write_frame(&mut socket, &fetch_v11_error_response_frame(2, 6)).await;
3954        });
3955
3956        let (client_stream, broker_stream) = tokio::io::duplex(64);
3957        let _broker_stream = broker_stream;
3958        let client = Client::from_stream(
3959            Box::new(client_stream),
3960            Some("kafrust-rack-failure-test".to_owned()),
3961            Some(std::time::Duration::from_millis(500)),
3962        );
3963        let config = ConsumerConfig::new([addr.to_string()])
3964            .request_timeout_ms(500)
3965            .max_retries(0)
3966            .client_rack("rack-a");
3967        let mut consumer = Consumer::from_assignments(client, config, Vec::new());
3968        let mut metadata = metadata_fixture();
3969        metadata.brokers[0].host = addr.ip().to_string();
3970        metadata.brokers[0].port = i32::from(addr.port());
3971        consumer
3972            .metadata_cache
3973            .insert("orders".to_owned(), metadata);
3974        consumer
3975            .preferred_read_replicas
3976            .insert(("orders".to_owned(), 0), 1);
3977
3978        assert!(consumer.fetch("orders", 0, 42).await.is_err());
3979        assert!(!consumer
3980            .preferred_read_replicas
3981            .contains_key(&("orders".to_owned(), 0)));
3982        server.await.unwrap();
3983    }
3984
3985    fn message(offset: i64) -> MessageSetRecord {
3986        MessageSetRecord {
3987            offset,
3988            leader_epoch: -1,
3989            timestamp_ms: 123,
3990            key: None,
3991            value: None,
3992            headers: Vec::new(),
3993            producer_id: None,
3994            transactional: false,
3995            control: false,
3996        }
3997    }
3998
3999    fn transactional_message(offset: i64, producer_id: i64, control: bool) -> MessageSetRecord {
4000        MessageSetRecord {
4001            offset,
4002            leader_epoch: -1,
4003            timestamp_ms: 123,
4004            key: None,
4005            value: None,
4006            headers: Vec::new(),
4007            producer_id: Some(producer_id),
4008            transactional: true,
4009            control,
4010        }
4011    }
4012
4013    fn metadata_fixture() -> MetadataResponseV1 {
4014        MetadataResponseV1 {
4015            brokers: vec![BrokerMetadata {
4016                node_id: 1,
4017                host: "localhost".to_owned(),
4018                port: 9092,
4019                rack: None,
4020            }],
4021            controller_id: 1,
4022            topics: vec![TopicMetadata {
4023                error_code: 0,
4024                name: "orders".to_owned(),
4025                is_internal: false,
4026                partitions: vec![PartitionMetadata {
4027                    error_code: 0,
4028                    partition_index: 0,
4029                    leader_id: 1,
4030                    replica_nodes: vec![1],
4031                    isr_nodes: vec![1],
4032                }],
4033            }],
4034        }
4035    }
4036
4037    async fn read_frame<T>(stream: &mut T) -> Vec<u8>
4038    where
4039        T: AsyncRead + Unpin,
4040    {
4041        let mut size = [0u8; 4];
4042        stream.read_exact(&mut size).await.unwrap();
4043        let size = usize::try_from(i32::from_be_bytes(size)).unwrap();
4044        let mut request = vec![0u8; size];
4045        stream.read_exact(&mut request).await.unwrap();
4046        request
4047    }
4048
4049    async fn write_frame<T>(stream: &mut T, frame: &[u8])
4050    where
4051        T: AsyncWrite + Unpin,
4052    {
4053        stream
4054            .write_all(&(frame.len() as i32).to_be_bytes())
4055            .await
4056            .unwrap();
4057        stream.write_all(frame).await.unwrap();
4058        stream.flush().await.unwrap();
4059    }
4060
4061    fn metadata_response_frame() -> Vec<u8> {
4062        vec![
4063            0, 0, 0, 1, // correlation id
4064            0, 0, 0, 1, // brokers count
4065            0, 0, 0, 1, // node id
4066            0, 9, b'l', b'o', b'c', b'a', b'l', b'h', b'o', b's', b't', // host
4067            0, 0, 35, 132, // port 9092
4068            0xff, 0xff, // null rack
4069            0, 0, 0, 1, // controller id
4070            0, 0, 0, 1, // topics count
4071            0, 0, // topic error code
4072            0, 6, b'o', b'r', b'd', b'e', b'r', b's', // topic name
4073            0,    // is internal false
4074            0, 0, 0, 1, // partition count
4075            0, 0, // partition error code
4076            0, 0, 0, 0, // partition index
4077            0, 0, 0, 1, // leader id
4078            0, 0, 0, 1, // replica count
4079            0, 0, 0, 1, // replica node
4080            0, 0, 0, 1, // isr count
4081            0, 0, 0, 1, // isr node
4082        ]
4083    }
4084
4085    fn fetch_v4_response_frame() -> Vec<u8> {
4086        fetch_v4_response_frame_with_correlation(1)
4087    }
4088
4089    fn fetch_v4_response_frame_with_correlation(correlation_id: i32) -> Vec<u8> {
4090        let mut message = Encoder::new();
4091        message.write_i32(0);
4092        message.write_i8(1);
4093        message.write_i8(0);
4094        message.write_i64(123);
4095        message.write_nullable_bytes(Some(b"order-1")).unwrap();
4096        message.write_nullable_bytes(Some(b"created")).unwrap();
4097        let message = message.into_bytes();
4098
4099        let mut records = Encoder::new();
4100        records.write_i64(42);
4101        records.write_i32(i32::try_from(message.len()).unwrap());
4102        records.write_raw(&message);
4103        let records = records.into_bytes();
4104
4105        let mut response = Encoder::new();
4106        response.write_i32(correlation_id);
4107        response.write_i32(0);
4108        response.write_i32(1);
4109        response.write_string("orders").unwrap();
4110        response.write_i32(1);
4111        response.write_i32(0);
4112        response.write_i16(0);
4113        response.write_i64(43);
4114        response.write_i64(43);
4115        response.write_i32(0);
4116        response.write_bytes(&records).unwrap();
4117        response.into_bytes()
4118    }
4119
4120    fn api_versions_v3_fetch_v11_response(correlation_id: i32) -> Vec<u8> {
4121        api_versions_v3_fetch_response(correlation_id, 11)
4122    }
4123
4124    fn api_versions_v3_fetch_v12_response(correlation_id: i32) -> Vec<u8> {
4125        api_versions_v3_fetch_response(correlation_id, 12)
4126    }
4127
4128    fn api_versions_v3_fetch_v13_response(correlation_id: i32) -> Vec<u8> {
4129        api_versions_v3_fetch_response(correlation_id, 13)
4130    }
4131
4132    fn api_versions_v3_fetch_response(correlation_id: i32, max_version: i16) -> Vec<u8> {
4133        let mut response = Encoder::new();
4134        response.write_i32(correlation_id);
4135        response.write_i16(0);
4136        response.write_i8(2); // one compact API key entry
4137        response.write_i16(1);
4138        response.write_i16(0);
4139        response.write_i16(max_version);
4140        response.write_i8(0); // API key entry tags
4141        response.write_i32(0);
4142        response.write_i8(0); // response tags
4143        response.into_bytes()
4144    }
4145
4146    fn api_versions_v3_metadata_fetch_response(
4147        correlation_id: i32,
4148        fetch_max_version: i16,
4149    ) -> Vec<u8> {
4150        let mut response = Encoder::new();
4151        response.write_i32(correlation_id);
4152        response.write_i16(0);
4153        response.write_unsigned_varint(3); // two compact API key entries
4154        response.write_i16(3); // Metadata
4155        response.write_i16(0);
4156        response.write_i16(12);
4157        response.write_unsigned_varint(0);
4158        response.write_i16(1); // Fetch
4159        response.write_i16(0);
4160        response.write_i16(fetch_max_version);
4161        response.write_unsigned_varint(0);
4162        response.write_i32(0);
4163        response.write_unsigned_varint(0);
4164        response.into_bytes()
4165    }
4166
4167    fn fetch_v11_error_response_frame(correlation_id: i32, error_code: i16) -> Vec<u8> {
4168        fetch_v11_response_frame_with_error(correlation_id, error_code, -1)
4169    }
4170
4171    fn fetch_v11_response_frame_with_error(
4172        correlation_id: i32,
4173        error_code: i16,
4174        preferred_read_replica: i32,
4175    ) -> Vec<u8> {
4176        let mut response = Encoder::new();
4177        response.write_i32(correlation_id);
4178        response.write_i32(0);
4179        response.write_i16(0);
4180        response.write_i32(0);
4181        response.write_i32(1);
4182        response.write_string("orders").unwrap();
4183        response.write_i32(1);
4184        response.write_i32(0);
4185        response.write_i16(error_code);
4186        response.write_i64(43);
4187        response.write_i64(43);
4188        response.write_i64(42);
4189        response.write_i32(0);
4190        response.write_i32(preferred_read_replica);
4191        response.write_bytes(&[]).unwrap();
4192        response.into_bytes()
4193    }
4194
4195    fn fetch_v12_response_frame(correlation_id: i32, preferred_read_replica: i32) -> Vec<u8> {
4196        fetch_v12_response_frame_with_session(correlation_id, preferred_read_replica, 0)
4197    }
4198
4199    fn fetch_v12_response_frame_with_session(
4200        correlation_id: i32,
4201        preferred_read_replica: i32,
4202        session_id: i32,
4203    ) -> Vec<u8> {
4204        let mut response = Encoder::new();
4205        response.write_i32(correlation_id);
4206        response.write_unsigned_varint(0); // response header tags
4207        response.write_i32(0); // throttle time
4208        response.write_i16(0);
4209        response.write_i32(session_id);
4210        response.write_unsigned_varint(2); // one compact topic
4211        response.write_compact_string("orders").unwrap();
4212        response.write_unsigned_varint(2); // one compact partition
4213        response.write_i32(0);
4214        response.write_i16(0);
4215        response.write_i64(43);
4216        response.write_i64(43);
4217        response.write_i64(42);
4218        response.write_unsigned_varint(1); // no aborted transactions
4219        response.write_i32(preferred_read_replica);
4220        response.write_compact_nullable_bytes(Some(&[])).unwrap();
4221        response.write_unsigned_varint(0); // partition tags
4222        response.write_unsigned_varint(0); // topic tags
4223        response.write_unsigned_varint(0); // response tags
4224        response.into_bytes()
4225    }
4226
4227    fn fetch_v13_response_frame(correlation_id: i32, topic_id: [u8; 16]) -> Vec<u8> {
4228        let mut response = Encoder::new();
4229        response.write_i32(correlation_id);
4230        response.write_unsigned_varint(0); // response header tags
4231        response.write_i32(0); // throttle time
4232        response.write_i16(0); // top-level error
4233        response.write_i32(0); // fetch session id
4234        response.write_unsigned_varint(2); // one compact topic
4235        response.write_uuid(&topic_id);
4236        response.write_unsigned_varint(2); // one compact partition
4237        response.write_i32(0);
4238        response.write_i16(0);
4239        response.write_i64(43);
4240        response.write_i64(43);
4241        response.write_i64(42);
4242        response.write_unsigned_varint(1); // no aborted transactions
4243        response.write_i32(-1); // preferred read replica
4244        response.write_compact_nullable_bytes(Some(&[])).unwrap();
4245        response.write_unsigned_varint(0); // partition tags
4246        response.write_unsigned_varint(0); // topic tags
4247        response.write_unsigned_varint(0); // response tags
4248        response.into_bytes()
4249    }
4250
4251    fn fetch_v13_response_frame_with_partition_error(
4252        correlation_id: i32,
4253        topic_id: [u8; 16],
4254        error_code: i16,
4255    ) -> Vec<u8> {
4256        let mut response = Encoder::new();
4257        response.write_i32(correlation_id);
4258        response.write_unsigned_varint(0); // response header tags
4259        response.write_i32(0); // throttle time
4260        response.write_i16(0); // top-level error
4261        response.write_i32(0); // fetch session id
4262        response.write_unsigned_varint(2); // one compact topic
4263        response.write_uuid(&topic_id);
4264        response.write_unsigned_varint(2); // one compact partition
4265        response.write_i32(0);
4266        response.write_i16(error_code);
4267        response.write_i64(43);
4268        response.write_i64(43);
4269        response.write_i64(42);
4270        response.write_unsigned_varint(1); // no aborted transactions
4271        response.write_i32(-1); // preferred read replica
4272        response.write_compact_nullable_bytes(Some(&[])).unwrap();
4273        response.write_unsigned_varint(0); // partition tags
4274        response.write_unsigned_varint(0); // topic tags
4275        response.write_unsigned_varint(0); // response tags
4276        response.into_bytes()
4277    }
4278
4279    fn fetch_v13_response_frame_with_record_at(
4280        correlation_id: i32,
4281        topic_id: [u8; 16],
4282        offset: i64,
4283    ) -> Vec<u8> {
4284        let mut message = Encoder::new();
4285        message.write_i32(0);
4286        message.write_i8(1);
4287        message.write_i8(0);
4288        message.write_i64(123);
4289        message.write_nullable_bytes(Some(b"order-1")).unwrap();
4290        message.write_nullable_bytes(Some(b"created")).unwrap();
4291        let message = message.into_bytes();
4292
4293        let mut records = Encoder::new();
4294        records.write_i64(offset);
4295        records.write_i32(i32::try_from(message.len()).unwrap());
4296        records.write_raw(&message);
4297        let records = records.into_bytes();
4298
4299        let mut response = Encoder::new();
4300        response.write_i32(correlation_id);
4301        response.write_unsigned_varint(0); // response header tags
4302        response.write_i32(0); // throttle time
4303        response.write_i16(0); // top-level error
4304        response.write_i32(0); // fetch session id
4305        response.write_unsigned_varint(2); // one compact topic
4306        response.write_uuid(&topic_id);
4307        response.write_unsigned_varint(2); // one compact partition
4308        response.write_i32(0);
4309        response.write_i16(0);
4310        response.write_i64(43);
4311        response.write_i64(43);
4312        response.write_i64(0);
4313        response.write_unsigned_varint(1); // no aborted transactions
4314        response.write_i32(-1); // preferred read replica
4315        response
4316            .write_compact_nullable_bytes(Some(&records))
4317            .unwrap();
4318        response.write_unsigned_varint(0); // partition tags
4319        response.write_unsigned_varint(0); // topic tags
4320        response.write_unsigned_varint(0); // response tags
4321        response.into_bytes()
4322    }
4323
4324    fn fetch_v12_response_frame_with_record(correlation_id: i32, session_id: i32) -> Vec<u8> {
4325        fetch_v12_response_frame_with_record_at(correlation_id, session_id, 42)
4326    }
4327
4328    fn fetch_v12_response_frame_with_record_at(
4329        correlation_id: i32,
4330        session_id: i32,
4331        offset: i64,
4332    ) -> Vec<u8> {
4333        fetch_v12_response_frame_with_record_at_partition(correlation_id, session_id, 0, offset)
4334    }
4335
4336    fn fetch_v12_response_frame_with_record_at_partition(
4337        correlation_id: i32,
4338        session_id: i32,
4339        partition: i32,
4340        offset: i64,
4341    ) -> Vec<u8> {
4342        let mut message = Encoder::new();
4343        message.write_i32(0);
4344        message.write_i8(1);
4345        message.write_i8(0);
4346        message.write_i64(123);
4347        message.write_nullable_bytes(Some(b"order-1")).unwrap();
4348        message.write_nullable_bytes(Some(b"created")).unwrap();
4349        let message = message.into_bytes();
4350
4351        let mut records = Encoder::new();
4352        records.write_i64(offset);
4353        records.write_i32(i32::try_from(message.len()).unwrap());
4354        records.write_raw(&message);
4355        let records = records.into_bytes();
4356
4357        let mut response = Encoder::new();
4358        response.write_i32(correlation_id);
4359        response.write_unsigned_varint(0); // response header tags
4360        response.write_i32(0); // throttle time
4361        response.write_i16(0);
4362        response.write_i32(session_id);
4363        response.write_unsigned_varint(2); // one compact topic
4364        response.write_compact_string("orders").unwrap();
4365        response.write_unsigned_varint(2); // one compact partition
4366        response.write_i32(partition);
4367        response.write_i16(0);
4368        response.write_i64(43);
4369        response.write_i64(43);
4370        response.write_i64(0);
4371        response.write_unsigned_varint(1); // no aborted transactions
4372        response.write_i32(-1);
4373        response
4374            .write_compact_nullable_bytes(Some(&records))
4375            .unwrap();
4376        response.write_unsigned_varint(0); // partition tags
4377        response.write_unsigned_varint(0); // topic tags
4378        response.write_unsigned_varint(0); // response tags
4379        response.into_bytes()
4380    }
4381
4382    fn fetch_v12_response_frame_with_partition_error(
4383        correlation_id: i32,
4384        error_code: i16,
4385    ) -> Vec<u8> {
4386        fetch_v12_response_frame_with_partition_error_and_session(correlation_id, error_code, 0)
4387    }
4388
4389    fn fetch_v12_response_frame_with_partition_error_and_session(
4390        correlation_id: i32,
4391        error_code: i16,
4392        session_id: i32,
4393    ) -> Vec<u8> {
4394        let mut response = Encoder::new();
4395        response.write_i32(correlation_id);
4396        response.write_unsigned_varint(0); // response header tags
4397        response.write_i32(0); // throttle time
4398        response.write_i16(0);
4399        response.write_i32(session_id);
4400        response.write_unsigned_varint(2); // one compact topic
4401        response.write_compact_string("orders").unwrap();
4402        response.write_unsigned_varint(2); // one compact partition
4403        response.write_i32(0);
4404        response.write_i16(error_code);
4405        response.write_i64(43);
4406        response.write_i64(43);
4407        response.write_i64(0);
4408        response.write_unsigned_varint(1); // no aborted transactions
4409        response.write_i32(-1);
4410        response.write_compact_nullable_bytes(Some(&[])).unwrap();
4411        response.write_unsigned_varint(0); // partition tags
4412        response.write_unsigned_varint(0); // topic tags
4413        response.write_unsigned_varint(0); // response tags
4414        response.into_bytes()
4415    }
4416
4417    fn fetch_v12_out_of_range_response_frame(correlation_id: i32) -> Vec<u8> {
4418        let mut response = Encoder::new();
4419        response.write_i32(correlation_id);
4420        response.write_unsigned_varint(0); // response header tags
4421        response.write_i32(0); // throttle time
4422        response.write_i16(0);
4423        response.write_i32(0);
4424        response.write_unsigned_varint(2); // one compact topic
4425        response.write_compact_string("orders").unwrap();
4426        response.write_unsigned_varint(2); // one compact partition
4427        response.write_i32(0);
4428        response.write_i16(1); // OFFSET_OUT_OF_RANGE
4429        response.write_i64(-1);
4430        response.write_i64(-1);
4431        response.write_i64(0);
4432        response.write_unsigned_varint(1); // no aborted transactions
4433        response.write_i32(-1);
4434        response.write_compact_nullable_bytes(Some(&[])).unwrap();
4435        response.write_unsigned_varint(0); // partition tags
4436        response.write_unsigned_varint(0); // topic tags
4437        response.write_unsigned_varint(0); // response tags
4438        response.into_bytes()
4439    }
4440
4441    fn list_offsets_response_frame(correlation_id: i32, offset: i64) -> Vec<u8> {
4442        let mut response = Encoder::new();
4443        response.write_i32(correlation_id);
4444        response.write_i32(1);
4445        response.write_string("orders").unwrap();
4446        response.write_i32(1);
4447        response.write_i32(0);
4448        response.write_i16(0);
4449        response.write_i64(-1);
4450        response.write_i64(offset);
4451        response.into_bytes()
4452    }
4453
4454    fn offset_for_leader_epoch_response_frame() -> Vec<u8> {
4455        offset_for_leader_epoch_response_frame_with(1, 8, 42)
4456    }
4457
4458    fn offset_for_leader_epoch_response_frame_with(
4459        correlation_id: i32,
4460        leader_epoch: i32,
4461        end_offset: i64,
4462    ) -> Vec<u8> {
4463        let mut response = Encoder::new();
4464        response.write_i32(correlation_id);
4465        response.write_i32(0);
4466        response.write_i32(1);
4467        response.write_string("orders").unwrap();
4468        response.write_i32(1);
4469        response.write_i16(0);
4470        response.write_i32(0);
4471        response.write_i32(leader_epoch);
4472        response.write_i64(end_offset);
4473        response.into_bytes()
4474    }
4475
4476    fn metadata_response_frame_for(correlation_id: i32, addr: &std::net::SocketAddr) -> Vec<u8> {
4477        let mut response = Encoder::new();
4478        response.write_i32(correlation_id);
4479        response.write_i32(1);
4480        response.write_i32(1);
4481        response.write_string(&addr.ip().to_string()).unwrap();
4482        response.write_i32(i32::from(addr.port()));
4483        response.write_nullable_string(None).unwrap();
4484        response.write_i32(1);
4485        response.write_i32(1);
4486        response.write_i16(0);
4487        response.write_string("orders").unwrap();
4488        response.write_i8(0);
4489        response.write_i32(1);
4490        response.write_i16(0);
4491        response.write_i32(0);
4492        response.write_i32(1);
4493        response.write_i32(1);
4494        response.write_i32(1);
4495        response.write_i32(1);
4496        response.write_i32(1);
4497        response.into_bytes()
4498    }
4499
4500    fn metadata_v12_response_frame(
4501        correlation_id: i32,
4502        addr: &std::net::SocketAddr,
4503        leader_epoch: i32,
4504    ) -> Vec<u8> {
4505        metadata_v12_response_frame_with_topic_id(correlation_id, addr, leader_epoch, [0; 16])
4506    }
4507
4508    fn metadata_v12_response_frame_with_topic_id(
4509        correlation_id: i32,
4510        addr: &std::net::SocketAddr,
4511        leader_epoch: i32,
4512        topic_id: [u8; 16],
4513    ) -> Vec<u8> {
4514        let mut response = Encoder::new();
4515        response.write_i32(correlation_id);
4516        response.write_unsigned_varint(0); // response header tags
4517        response.write_i32(0); // throttle time
4518        response.write_unsigned_varint(2); // one compact broker
4519        response.write_i32(1);
4520        response
4521            .write_compact_string(&addr.ip().to_string())
4522            .unwrap();
4523        response.write_i32(i32::from(addr.port()));
4524        response.write_compact_nullable_string(None).unwrap();
4525        response.write_unsigned_varint(0); // broker tags
4526        response.write_compact_nullable_string(None).unwrap(); // cluster id
4527        response.write_i32(1); // controller id
4528        response.write_unsigned_varint(2); // one compact topic
4529        response.write_i16(0);
4530        response
4531            .write_compact_nullable_string(Some("orders"))
4532            .unwrap();
4533        response.write_uuid(&topic_id);
4534        response.write_i8(0); // internal
4535        response.write_unsigned_varint(2); // one compact partition
4536        response.write_i16(0);
4537        response.write_i32(0);
4538        response.write_i32(1);
4539        response.write_i32(leader_epoch);
4540        response.write_unsigned_varint(2); // one replica
4541        response.write_i32(1);
4542        response.write_unsigned_varint(2); // one ISR
4543        response.write_i32(1);
4544        response.write_unsigned_varint(1); // no offline replicas
4545        response.write_unsigned_varint(0); // partition tags
4546        response.write_i32(-2_147_483_648); // authorized operations unknown
4547        response.write_unsigned_varint(0); // topic tags
4548        response.write_unsigned_varint(0); // response tags
4549        response.into_bytes()
4550    }
4551}