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)]
31pub enum IsolationLevel {
33 #[default]
35 ReadUncommitted,
36 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)]
50pub enum OffsetResetPolicy {
52 Earliest,
54 Latest,
56 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)]
81pub 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)]
94pub 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 pub fn key(&self) -> &str {
113 &self.key
114 }
115
116 pub fn value(&self) -> Option<&[u8]> {
118 self.value.as_deref()
119 }
120}
121
122#[derive(Debug)]
131pub struct ConsumerPartitionQueue {
132 topic: String,
133 partition: i32,
134 receiver: mpsc::Receiver<ConsumerRecord>,
135}
136
137impl ConsumerPartitionQueue {
138 pub fn topic(&self) -> &str {
140 &self.topic
141 }
142
143 pub fn partition(&self) -> i32 {
145 self.partition
146 }
147
148 pub async fn recv(&mut self) -> Option<ConsumerRecord> {
151 self.receiver.recv().await
152 }
153
154 pub fn try_recv(&mut self) -> Option<ConsumerRecord> {
159 self.receiver.try_recv().ok()
160 }
161
162 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 pub fn topic(&self) -> &str {
204 &self.topic
205 }
206
207 pub fn partition(&self) -> i32 {
209 self.partition
210 }
211
212 pub fn offset(&self) -> i64 {
214 self.offset
215 }
216
217 pub fn leader_epoch(&self) -> i32 {
222 self.leader_epoch
223 }
224
225 pub fn timestamp_ms(&self) -> i64 {
227 self.timestamp_ms
228 }
229
230 pub fn key(&self) -> Option<&[u8]> {
232 self.key.as_deref()
233 }
234
235 pub fn value(&self) -> Option<&[u8]> {
237 self.value.as_deref()
238 }
239
240 pub fn headers(&self) -> &[ConsumerRecordHeader] {
242 &self.headers
243 }
244}
245
246#[derive(Debug, Clone, Copy, PartialEq, Eq)]
247pub struct PartitionWatermarks {
249 low: i64,
250 high: i64,
251}
252
253#[derive(Debug, Clone, Copy, PartialEq, Eq)]
254pub struct LeaderEpochOffset {
256 leader_epoch: i32,
257 end_offset: i64,
258}
259
260impl LeaderEpochOffset {
261 pub fn leader_epoch(&self) -> i32 {
263 self.leader_epoch
264 }
265
266 pub fn end_offset(&self) -> i64 {
268 self.end_offset
269 }
270}
271
272impl PartitionWatermarks {
273 pub fn low(&self) -> i64 {
275 self.low
276 }
277
278 pub fn high(&self) -> i64 {
280 self.high
281 }
282}
283
284#[derive(Debug)]
285pub 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 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 pub fn assignments(&self) -> &[ConsumerAssignment] {
410 &self.assignments
411 }
412
413 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 pub fn position(&self, topic: &str, partition: i32) -> Option<i64> {
449 self.assignment(topic, partition)
450 .map(ConsumerAssignment::next_offset)
451 }
452
453 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 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 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 #[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 #[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 #[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 #[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)]
1358pub 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 pub fn topic(&self) -> &str {
1384 &self.topic
1385 }
1386
1387 pub fn partition(&self) -> i32 {
1389 self.partition
1390 }
1391
1392 pub fn next_offset(&self) -> i64 {
1394 self.next_offset
1395 }
1396
1397 pub fn leader_epoch(&self) -> i32 {
1402 self.leader_epoch
1403 }
1404
1405 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)]
1436pub 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 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 pub fn with_client_config(mut self, client: ClientConfig) -> Self {
1474 self.client = client;
1475 self
1476 }
1477
1478 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 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 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 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 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 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 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 pub fn metrics(mut self, metrics: ClientMetrics) -> Self {
1523 self.client = self.client.metrics(metrics);
1524 self
1525 }
1526
1527 pub fn security_protocol(mut self, security_protocol: SecurityProtocol) -> Self {
1529 self.client = self.client.security_protocol(security_protocol);
1530 self
1531 }
1532
1533 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 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 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 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 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 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 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 pub fn sasl_oauthbearer(mut self, token: impl Into<String>) -> Self {
1585 self.client = self.client.sasl_oauthbearer(token);
1586 self
1587 }
1588
1589 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 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 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 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 pub fn min_bytes(mut self, min_bytes: i32) -> Self {
1632 self.min_bytes = min_bytes;
1633 self
1634 }
1635
1636 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 pub fn max_retries(mut self, max_retries: u32) -> Self {
1644 self.max_retries = max_retries;
1645 self
1646 }
1647
1648 pub fn max_retries_ref(&self) -> u32 {
1650 self.max_retries
1651 }
1652
1653 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 pub fn max_poll_records_ref(&self) -> usize {
1661 self.max_poll_records
1662 }
1663
1664 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 pub fn partition_queue_capacity_ref(&self) -> usize {
1674 self.partition_queue_capacity
1675 }
1676
1677 pub fn isolation_level(mut self, isolation_level: IsolationLevel) -> Self {
1679 self.isolation_level = isolation_level;
1680 self
1681 }
1682
1683 pub fn isolation_level_ref(&self) -> IsolationLevel {
1685 self.isolation_level
1686 }
1687
1688 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 pub fn offset_reset_policy_ref(&self) -> OffsetResetPolicy {
1700 self.offset_reset_policy
1701 }
1702
1703 pub fn client_config(&self) -> &ClientConfig {
1705 &self.client
1706 }
1707
1708 pub fn validate(&self) -> Result<()> {
1710 self.client.validate()?;
1711 self.validate_values()
1712 }
1713
1714 pub fn build_config(self) -> Result<Self> {
1717 self.validate()?;
1718 Ok(self)
1719 }
1720
1721 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
2099const 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, 0, 0, 0, 1, 0, 0, 0, 1, 0, 9, b'l', b'o', b'c', b'a', b'l', b'h', b'o', b's', b't', 0, 0, 35, 132, 0xff, 0xff, 0, 0, 0, 1, 0, 0, 0, 1, 0, 0, 0, 6, b'o', b'r', b'd', b'e', b'r', b's', 0, 0, 0, 0, 1, 0, 0, 0, 0, 0, 0, 0, 0, 0, 1, 0, 0, 0, 1, 0, 0, 0, 1, 0, 0, 0, 1, 0, 0, 0, 1, ]
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); response.write_i16(1);
4138 response.write_i16(0);
4139 response.write_i16(max_version);
4140 response.write_i8(0); response.write_i32(0);
4142 response.write_i8(0); 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); response.write_i16(3); response.write_i16(0);
4156 response.write_i16(12);
4157 response.write_unsigned_varint(0);
4158 response.write_i16(1); 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.write_i32(0); response.write_i16(0);
4209 response.write_i32(session_id);
4210 response.write_unsigned_varint(2); response.write_compact_string("orders").unwrap();
4212 response.write_unsigned_varint(2); 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); response.write_i32(preferred_read_replica);
4220 response.write_compact_nullable_bytes(Some(&[])).unwrap();
4221 response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); 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.write_i32(0); response.write_i16(0); response.write_i32(0); response.write_unsigned_varint(2); response.write_uuid(&topic_id);
4236 response.write_unsigned_varint(2); 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); response.write_i32(-1); response.write_compact_nullable_bytes(Some(&[])).unwrap();
4245 response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); 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.write_i32(0); response.write_i16(0); response.write_i32(0); response.write_unsigned_varint(2); response.write_uuid(&topic_id);
4264 response.write_unsigned_varint(2); 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); response.write_i32(-1); response.write_compact_nullable_bytes(Some(&[])).unwrap();
4273 response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); 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.write_i32(0); response.write_i16(0); response.write_i32(0); response.write_unsigned_varint(2); response.write_uuid(&topic_id);
4307 response.write_unsigned_varint(2); 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); response.write_i32(-1); response
4316 .write_compact_nullable_bytes(Some(&records))
4317 .unwrap();
4318 response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); 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.write_i32(0); response.write_i16(0);
4362 response.write_i32(session_id);
4363 response.write_unsigned_varint(2); response.write_compact_string("orders").unwrap();
4365 response.write_unsigned_varint(2); 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); response.write_i32(-1);
4373 response
4374 .write_compact_nullable_bytes(Some(&records))
4375 .unwrap();
4376 response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); 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.write_i32(0); response.write_i16(0);
4399 response.write_i32(session_id);
4400 response.write_unsigned_varint(2); response.write_compact_string("orders").unwrap();
4402 response.write_unsigned_varint(2); 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); response.write_i32(-1);
4410 response.write_compact_nullable_bytes(Some(&[])).unwrap();
4411 response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); 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.write_i32(0); response.write_i16(0);
4423 response.write_i32(0);
4424 response.write_unsigned_varint(2); response.write_compact_string("orders").unwrap();
4426 response.write_unsigned_varint(2); response.write_i32(0);
4428 response.write_i16(1); response.write_i64(-1);
4430 response.write_i64(-1);
4431 response.write_i64(0);
4432 response.write_unsigned_varint(1); response.write_i32(-1);
4434 response.write_compact_nullable_bytes(Some(&[])).unwrap();
4435 response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); 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.write_i32(0); response.write_unsigned_varint(2); 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); response.write_compact_nullable_string(None).unwrap(); response.write_i32(1); response.write_unsigned_varint(2); 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); response.write_unsigned_varint(2); 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); response.write_i32(1);
4542 response.write_unsigned_varint(2); response.write_i32(1);
4544 response.write_unsigned_varint(1); response.write_unsigned_varint(0); response.write_i32(-2_147_483_648); response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.into_bytes()
4550 }
4551}