use std::collections::{BTreeMap, BTreeSet};
use kafrust_protocol::api::fetch::{
AbortedTransactionV4, FetchPartitionResponseV11, FetchPartitionResponseV12,
FetchPartitionResponseV4, FetchResponseV11, FetchResponseV12, FetchResponseV4,
MessageSetRecord,
};
use kafrust_protocol::api::list_offsets::{
ListOffsetsPartitionResponseV1, ListOffsetsPartitionV1, ListOffsetsTopicResponseV1,
ListOffsetsTopicV1, EARLIEST_TIMESTAMP, LATEST_TIMESTAMP,
};
use kafrust_protocol::api::metadata::{
BrokerMetadata, MetadataPartitionV12, MetadataRequestTopicV12, MetadataResponseV1,
MetadataResponseV12,
};
use kafrust_protocol::api::offset_for_leader_epoch::{
OffsetForLeaderEpochPartitionResponseV3, OffsetForLeaderEpochPartitionV3,
OffsetForLeaderEpochTopicResponseV3, OffsetForLeaderEpochTopicV3,
};
use crate::client::{Client, FetchOneRequestV11, FetchOneRequestV12, FetchOneRequestV4};
use crate::config::{ClientConfig, OAuthBearerTokenProvider, SecurityProtocol};
use crate::error::{BrokerErrorKind, Error, Result};
use crate::metrics::ClientMetrics;
use tokio::sync::mpsc;
use tokio::sync::mpsc::error::TrySendError;
use tracing::debug;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum IsolationLevel {
#[default]
ReadUncommitted,
ReadCommitted,
}
impl IsolationLevel {
fn as_i8(self) -> i8 {
match self {
Self::ReadUncommitted => 0,
Self::ReadCommitted => 1,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OffsetResetPolicy {
Earliest,
Latest,
Offset(i64),
}
impl Default for OffsetResetPolicy {
fn default() -> Self {
Self::Offset(0)
}
}
impl OffsetResetPolicy {
pub(crate) fn timestamp(self) -> Option<i64> {
match self {
Self::Earliest => Some(EARLIEST_TIMESTAMP),
Self::Latest => Some(LATEST_TIMESTAMP),
Self::Offset(_) => None,
}
}
fn is_recovery(self) -> bool {
matches!(self, Self::Earliest | Self::Latest)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerRecord {
topic: String,
partition: i32,
offset: i64,
leader_epoch: i32,
timestamp_ms: i64,
key: Option<Vec<u8>>,
value: Option<Vec<u8>>,
headers: Vec<ConsumerRecordHeader>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerRecordHeader {
key: String,
value: Option<Vec<u8>>,
}
impl ConsumerRecordHeader {
fn from_protocol(header: kafrust_protocol::api::produce::RecordBatchHeader) -> Self {
Self {
key: header.key,
value: header.value,
}
}
pub fn key(&self) -> &str {
&self.key
}
pub fn value(&self) -> Option<&[u8]> {
self.value.as_deref()
}
}
#[derive(Debug)]
pub struct ConsumerPartitionQueue {
topic: String,
partition: i32,
receiver: mpsc::Receiver<ConsumerRecord>,
}
impl ConsumerPartitionQueue {
pub fn topic(&self) -> &str {
&self.topic
}
pub fn partition(&self) -> i32 {
self.partition
}
pub async fn recv(&mut self) -> Option<ConsumerRecord> {
self.receiver.recv().await
}
pub fn try_recv(&mut self) -> Option<ConsumerRecord> {
self.receiver.try_recv().ok()
}
pub async fn recv_batch(&mut self, max_records: usize) -> Vec<ConsumerRecord> {
if max_records == 0 {
return Vec::new();
}
let Some(first) = self.recv().await else {
return Vec::new();
};
let mut records = Vec::with_capacity(max_records.min(16));
records.push(first);
while records.len() < max_records {
let Some(record) = self.try_recv() else {
break;
};
records.push(record);
}
records
}
}
impl ConsumerRecord {
fn from_message_set(topic: &str, partition: i32, record: MessageSetRecord) -> Self {
Self {
topic: topic.to_owned(),
partition,
offset: record.offset,
leader_epoch: record.leader_epoch,
timestamp_ms: record.timestamp_ms,
key: record.key,
value: record.value,
headers: record
.headers
.into_iter()
.map(ConsumerRecordHeader::from_protocol)
.collect(),
}
}
pub fn topic(&self) -> &str {
&self.topic
}
pub fn partition(&self) -> i32 {
self.partition
}
pub fn offset(&self) -> i64 {
self.offset
}
pub fn leader_epoch(&self) -> i32 {
self.leader_epoch
}
pub fn timestamp_ms(&self) -> i64 {
self.timestamp_ms
}
pub fn key(&self) -> Option<&[u8]> {
self.key.as_deref()
}
pub fn value(&self) -> Option<&[u8]> {
self.value.as_deref()
}
pub fn headers(&self) -> &[ConsumerRecordHeader] {
&self.headers
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PartitionWatermarks {
low: i64,
high: i64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LeaderEpochOffset {
leader_epoch: i32,
end_offset: i64,
}
impl LeaderEpochOffset {
pub fn leader_epoch(&self) -> i32 {
self.leader_epoch
}
pub fn end_offset(&self) -> i64 {
self.end_offset
}
}
impl PartitionWatermarks {
pub fn low(&self) -> i64 {
self.low
}
pub fn high(&self) -> i64 {
self.high
}
}
#[derive(Debug)]
pub struct Consumer {
client: Client,
config: ConsumerConfig,
assignments: Vec<ConsumerAssignment>,
partition_queues: BTreeMap<(String, i32), mpsc::Sender<ConsumerRecord>>,
metadata_cache: BTreeMap<String, MetadataResponseV1>,
broker_clients: BTreeMap<String, Client>,
fetch_sessions: BTreeMap<String, FetchSessionState>,
preferred_read_replicas: BTreeMap<(String, i32), i32>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
struct FetchSessionState {
session_id: i32,
next_epoch: i32,
}
impl FetchSessionState {
fn next_request(self) -> (i32, i32) {
(self.session_id, self.next_epoch)
}
fn advance_with_response(self, session_id: i32) -> Self {
Self {
session_id,
next_epoch: self.next_epoch.saturating_add(1),
}
}
}
impl Consumer {
pub(crate) fn from_assignments(
client: Client,
config: ConsumerConfig,
assignments: Vec<ConsumerAssignment>,
) -> Self {
Self {
client,
config,
assignments,
partition_queues: BTreeMap::new(),
metadata_cache: BTreeMap::new(),
broker_clients: BTreeMap::new(),
fetch_sessions: BTreeMap::new(),
preferred_read_replicas: BTreeMap::new(),
}
}
pub(crate) fn replace_assignments(&mut self, mut assignments: Vec<ConsumerAssignment>) {
let previous_positions = self
.assignments
.iter()
.map(|assignment| {
(
(assignment.topic.clone(), assignment.partition),
assignment.next_offset,
)
})
.collect::<BTreeMap<_, _>>();
let previous_leader_epochs = self
.assignments
.iter()
.map(|assignment| {
(
(assignment.topic.clone(), assignment.partition),
assignment.leader_epoch,
)
})
.collect::<BTreeMap<_, _>>();
let paused = self
.assignments
.iter()
.filter(|assignment| assignment.paused)
.map(|assignment| (assignment.topic.clone(), assignment.partition))
.collect::<BTreeSet<_>>();
for assignment in &mut assignments {
assignment.paused = paused.contains(&(assignment.topic.clone(), assignment.partition));
if let Some(previous_epoch) =
previous_leader_epochs.get(&(assignment.topic.clone(), assignment.partition))
{
assignment.leader_epoch = *previous_epoch;
}
}
assignments.sort_by(|left, right| {
left.topic
.cmp(&right.topic)
.then_with(|| left.partition.cmp(&right.partition))
});
self.partition_queues.retain(|(topic, partition), _| {
assignments.iter().any(|assignment| {
assignment.topic == *topic
&& assignment.partition == *partition
&& previous_positions.get(&(topic.clone(), *partition))
== Some(&assignment.next_offset)
})
});
self.assignments = assignments;
self.metadata_cache.clear();
self.fetch_sessions.clear();
}
pub fn assign(&mut self, topic: impl Into<String>, partition: i32, offset: i64) {
let topic = topic.into();
self.partition_queues.remove(&(topic.clone(), partition));
assign_partition(&mut self.assignments, topic, partition, offset);
self.fetch_sessions.clear();
}
pub fn assignments(&self) -> &[ConsumerAssignment] {
&self.assignments
}
pub fn split_partition_queue(
&mut self,
topic: impl Into<String>,
partition: i32,
) -> Result<ConsumerPartitionQueue> {
let topic = topic.into();
if self.assignment(&topic, partition).is_none() {
return Err(Error::UnassignedTopicPartition { topic, partition });
}
let key = (topic.clone(), partition);
if let Some(sender) = self.partition_queues.get(&key) {
if !sender.is_closed() {
return Err(Error::Unsupported("partition queue is already split"));
}
}
self.partition_queues.remove(&key);
let (sender, receiver) = mpsc::channel(self.config.partition_queue_capacity);
self.partition_queues.insert(key, sender);
Ok(ConsumerPartitionQueue {
topic,
partition,
receiver,
})
}
pub fn position(&self, topic: &str, partition: i32) -> Option<i64> {
self.assignment(topic, partition)
.map(ConsumerAssignment::next_offset)
}
pub fn seek(&mut self, topic: &str, partition: i32, offset: i64) -> Result<()> {
if self
.partition_queues
.get(&(topic.to_owned(), partition))
.is_some_and(|sender| !sender.is_closed())
{
return Err(Error::Unsupported(
"seek requires dropping the active partition queue",
));
}
self.assignment_mut(topic, partition)?.next_offset = offset;
self.fetch_sessions.clear();
Ok(())
}
pub fn pause(&mut self, topic: &str, partition: i32) -> Result<()> {
self.assignment_mut(topic, partition)?.paused = true;
self.fetch_sessions.clear();
Ok(())
}
pub fn resume(&mut self, topic: &str, partition: i32) -> Result<()> {
self.assignment_mut(topic, partition)?.paused = false;
self.fetch_sessions.clear();
Ok(())
}
#[tracing::instrument(
level = "debug",
name = "kafka.consumer.fetch_watermarks",
skip_all,
fields(topic = tracing::field::Empty, partition),
err
)]
pub async fn fetch_watermarks(
&mut self,
topic: impl Into<String>,
partition: i32,
) -> Result<PartitionWatermarks> {
let topic = topic.into();
tracing::Span::current().record("topic", topic.as_str());
let mut attempt = 0;
loop {
match self.fetch_watermarks_once(&topic, partition).await {
Err(error) if attempt < self.config.max_retries && can_retry_fetch(&error) => {
invalidate_metadata_cache(&mut self.metadata_cache, &topic);
self.config.client.record_retry();
attempt += 1;
}
result => return result,
}
}
}
#[tracing::instrument(
level = "debug",
name = "kafka.consumer.offset_for_leader_epoch",
skip_all,
fields(topic = tracing::field::Empty, partition, current_leader_epoch, leader_epoch),
err
)]
pub async fn offset_for_leader_epoch(
&mut self,
topic: impl Into<String>,
partition: i32,
current_leader_epoch: i32,
leader_epoch: i32,
) -> Result<LeaderEpochOffset> {
let topic = topic.into();
tracing::Span::current().record("topic", topic.as_str());
let mut attempt = 0;
loop {
match self
.offset_for_leader_epoch_once(&topic, partition, current_leader_epoch, leader_epoch)
.await
{
Err(error) if attempt < self.config.max_retries && can_retry_fetch(&error) => {
invalidate_metadata_cache(&mut self.metadata_cache, &topic);
self.config.client.record_retry();
attempt += 1;
}
result => return result,
}
}
}
async fn offset_for_leader_epoch_once(
&mut self,
topic: &str,
partition: i32,
current_leader_epoch: i32,
leader_epoch: i32,
) -> Result<LeaderEpochOffset> {
let metadata = self.metadata_for_topic(topic).await?;
let leader = leader_for(&metadata, topic, partition)?;
let broker_addr = broker_addr_for(&metadata, leader)?;
let mut leader_client = self.connect_or_reuse_broker(&broker_addr).await?;
let response = leader_client
.offset_for_leader_epoch_v3(vec![OffsetForLeaderEpochTopicV3 {
name: topic.to_owned(),
partitions: vec![OffsetForLeaderEpochPartitionV3 {
partition_index: partition,
current_leader_epoch,
leader_epoch,
}],
}])
.await?;
let partition_response =
offset_for_leader_epoch_partition_response(&response.topics, topic, partition)?;
if partition_response.error_code != 0 {
return Err(self.config.client.broker_error(
partition_response.error_code,
format!("offset for leader epoch {topic}-{partition}"),
));
}
self.broker_clients.insert(broker_addr, leader_client);
Ok(LeaderEpochOffset {
leader_epoch: partition_response.leader_epoch,
end_offset: partition_response.end_offset,
})
}
async fn fetch_watermarks_once(
&mut self,
topic: &str,
partition: i32,
) -> Result<PartitionWatermarks> {
let metadata = self.metadata_for_topic(topic).await?;
let leader = leader_for(&metadata, topic, partition)?;
let broker_addr = broker_addr_for(&metadata, leader)?;
let mut leader_client = self.connect_or_reuse_broker(&broker_addr).await?;
let low = self
.request_partition_offset(&mut leader_client, topic, partition, EARLIEST_TIMESTAMP)
.await?;
let high = self
.request_partition_offset(&mut leader_client, topic, partition, LATEST_TIMESTAMP)
.await?;
self.broker_clients.insert(broker_addr, leader_client);
Ok(PartitionWatermarks { low, high })
}
async fn request_partition_offset(
&self,
client: &mut Client,
topic: &str,
partition: i32,
timestamp: i64,
) -> Result<i64> {
let response = client
.list_offsets_v1(vec![ListOffsetsTopicV1 {
name: topic.to_owned(),
partitions: vec![ListOffsetsPartitionV1 {
partition_index: partition,
timestamp,
}],
}])
.await?;
let partition_response =
list_offset_partition_response(&response.topics, topic, partition)?;
if partition_response.error_code != 0 {
return Err(self.config.client.broker_error(
partition_response.error_code,
format!("list offsets {topic}-{partition}"),
));
}
Ok(partition_response.offset)
}
#[tracing::instrument(
level = "debug",
name = "kafka.consumer.poll",
skip_all,
fields(assignment_count = self.assignments.len(), max_poll_records = self.config.max_poll_records),
err
)]
pub async fn poll(&mut self) -> Result<Vec<ConsumerRecord>> {
let assignments = self.assignments.clone();
let mut records = Vec::new();
let mut delivered_record_count = 0;
debug!(
assignment_count = assignments.len(),
max_poll_records = self.config.max_poll_records,
"polling kafka consumer assignments"
);
for assignment in assignments {
if delivered_record_count >= self.config.max_poll_records {
break;
}
if assignment.paused {
continue;
}
let mut fetched = self
.fetch_with_progress(
&assignment.topic,
assignment.partition,
assignment.next_offset,
assignment.leader_epoch,
Some(self.config.offset_reset_policy),
)
.await?;
if let Some(leader_epoch) = fetched.leader_epoch {
self.update_assignment_leader_epoch(
&assignment.topic,
assignment.partition,
leader_epoch,
);
}
let fetched_record_count = fetched.records.len();
limit_fetched_records(
&mut fetched.records,
delivered_record_count,
self.config.max_poll_records,
);
let was_truncated = fetched.records.len() < fetched_record_count;
let key = (assignment.topic.clone(), assignment.partition);
if self.partition_queues.contains_key(&key) {
let route = self.enqueue_partition_records(&key, fetched.records)?;
fetched.records = route.records;
delivered_record_count += route.queued_count;
if fetched.records.is_empty() {
let next_offset = if was_truncated {
route.next_offset
} else {
(fetched_record_count > 0).then_some(fetched.next_offset)
};
if let Some(next_offset) = next_offset {
self.update_assignment_offset(
&assignment.topic,
assignment.partition,
next_offset,
);
}
continue;
}
}
let next_offset = if fetched.records.len() < fetched_record_count {
fetched
.records
.last()
.map(|record| record.offset().saturating_add(1))
} else {
Some(fetched.next_offset)
};
if let Some(next_offset) = next_offset {
self.update_assignment_offset(&assignment.topic, assignment.partition, next_offset);
}
if was_truncated && fetched.records.is_empty() {
break;
}
delivered_record_count += fetched.records.len();
records.extend(fetched.records);
}
debug!(
record_count = records.len(),
"polled kafka consumer records"
);
self.config.client.record_consumed(delivered_record_count);
Ok(records)
}
fn enqueue_partition_records(
&mut self,
key: &(String, i32),
records: Vec<ConsumerRecord>,
) -> Result<PartitionRoute> {
let Some(sender) = self.partition_queues.get(key).cloned() else {
return Ok(PartitionRoute {
records,
next_offset: None,
queued_count: 0,
});
};
if sender.is_closed() {
self.partition_queues.remove(key);
return Ok(PartitionRoute {
records,
next_offset: None,
queued_count: 0,
});
}
let mut iterator = records.into_iter();
let mut queued_count = 0;
let mut next_offset = None;
while let Some(record) = iterator.next() {
let record_next_offset = record.offset().saturating_add(1);
match sender.try_send(record) {
Ok(()) => {
queued_count += 1;
next_offset = Some(record_next_offset);
}
Err(TrySendError::Closed(record)) => {
self.partition_queues.remove(key);
let mut main_records = Vec::with_capacity(iterator.len() + 1);
main_records.push(record);
main_records.extend(iterator);
return Ok(PartitionRoute {
records: main_records,
next_offset,
queued_count,
});
}
Err(TrySendError::Full(_record)) => {
if let Some(next_offset) = next_offset {
self.update_assignment_offset(key.0.as_str(), key.1, next_offset);
}
return Err(Error::PartitionQueueFull {
topic: key.0.clone(),
partition: key.1,
capacity: self.config.partition_queue_capacity,
});
}
}
}
Ok(PartitionRoute {
records: Vec::new(),
next_offset,
queued_count,
})
}
#[tracing::instrument(
level = "debug",
name = "kafka.consumer.fetch",
skip_all,
fields(topic = tracing::field::Empty, partition, offset),
err
)]
pub async fn fetch(
&mut self,
topic: impl Into<String>,
partition: i32,
offset: i64,
) -> Result<Vec<ConsumerRecord>> {
let topic = topic.into();
tracing::Span::current().record("topic", topic.as_str());
let records = self
.fetch_with_progress(&topic, partition, offset, -1, None)
.await?
.records;
self.config.client.record_consumed(records.len());
Ok(records)
}
async fn fetch_with_progress(
&mut self,
topic: &str,
partition: i32,
offset: i64,
current_leader_epoch: i32,
offset_reset_policy: Option<OffsetResetPolicy>,
) -> Result<FetchedPartition> {
let mut attempt = 0;
let mut fetch_offset = offset;
let mut request_leader_epoch = current_leader_epoch;
let mut reset_applied = false;
let mut truncation_checked = false;
let mut recovered_leader_epoch = None;
debug!(topic, partition, offset, "fetching kafka records");
loop {
let result = self
.fetch_once(topic, partition, fetch_offset, request_leader_epoch)
.await;
if result.is_err() {
self.preferred_read_replicas
.remove(&(topic.to_owned(), partition));
}
match result {
Err(error)
if !reset_applied
&& offset_reset_policy.is_some_and(OffsetResetPolicy::is_recovery)
&& is_offset_out_of_range(&error) =>
{
let Some(policy) = offset_reset_policy else {
return Err(error);
};
fetch_offset = self.reset_fetch_offset(topic, partition, policy).await?;
request_leader_epoch = -1;
reset_applied = true;
self.config.client.record_retry();
attempt += 1;
}
Err(error)
if !truncation_checked
&& attempt < self.config.max_retries
&& request_leader_epoch >= 0
&& is_leader_epoch_transition_error(&error) =>
{
let previous_leader_epoch = request_leader_epoch;
let current_leader_epoch = self
.current_leader_epoch(topic, partition)
.await?
.filter(|epoch| *epoch >= previous_leader_epoch);
if let Some(current_leader_epoch) = current_leader_epoch {
invalidate_metadata_cache(&mut self.metadata_cache, topic);
let epoch_offset = self
.offset_for_leader_epoch(
topic,
partition,
current_leader_epoch,
previous_leader_epoch,
)
.await?;
let end_offset = epoch_offset.end_offset();
if end_offset >= 0 {
fetch_offset = fetch_offset.min(end_offset);
}
debug!(
topic,
partition,
previous_leader_epoch,
current_leader_epoch,
end_offset,
fetch_offset,
"recovered fetch after leader epoch transition"
);
request_leader_epoch = current_leader_epoch;
recovered_leader_epoch = Some(current_leader_epoch);
truncation_checked = true;
} else {
request_leader_epoch = -1;
invalidate_metadata_cache(&mut self.metadata_cache, topic);
}
self.config.client.record_retry();
attempt += 1;
}
Err(error) if attempt < self.config.max_retries && can_retry_fetch(&error) => {
invalidate_metadata_cache(&mut self.metadata_cache, topic);
self.config.client.record_retry();
attempt += 1;
}
Ok(records) => {
debug!(
topic,
partition,
offset,
record_count = records.records.len(),
"fetched kafka records"
);
let mut records = records;
if records.leader_epoch.is_none() {
records.leader_epoch = recovered_leader_epoch;
}
return Ok(records);
}
Err(error) => return Err(error),
}
}
}
async fn reset_fetch_offset(
&mut self,
topic: &str,
partition: i32,
policy: OffsetResetPolicy,
) -> Result<i64> {
match policy {
OffsetResetPolicy::Earliest => Ok(self.fetch_watermarks(topic, partition).await?.low),
OffsetResetPolicy::Latest => Ok(self.fetch_watermarks(topic, partition).await?.high),
OffsetResetPolicy::Offset(_) => Err(Error::Unsupported(
"explicit offset reset policy cannot recover an out-of-range fetch",
)),
}
}
async fn current_leader_epoch(&mut self, topic: &str, partition: i32) -> Result<Option<i32>> {
if !self.client.supports_metadata_v12().await? {
return Ok(None);
}
let metadata = self
.client
.metadata_v12(Some(vec![MetadataRequestTopicV12 {
topic_id: [0; 16],
name: Some(topic.to_owned()),
}]))
.await?;
let partition_metadata = metadata_v12_partition(&metadata, topic, partition)?;
if partition_metadata.error_code != 0 {
return Err(self.config.client.broker_error(
partition_metadata.error_code,
format!("metadata {topic}-{partition}"),
));
}
Ok(Some(partition_metadata.leader_epoch))
}
async fn fetch_once(
&mut self,
topic: &str,
partition: i32,
offset: i64,
current_leader_epoch: i32,
) -> Result<FetchedPartition> {
let metadata = self.metadata_for_topic(topic).await?;
let leader = leader_for(&metadata, topic, partition)?;
let selected_broker = self
.preferred_read_replicas
.get(&(topic.to_owned(), partition))
.copied()
.filter(|broker| broker_addr_for(&metadata, *broker).is_ok())
.unwrap_or(leader);
let broker_addr = broker_addr_for(&metadata, selected_broker)?;
debug!(
topic = topic,
partition,
leader,
selected_broker,
broker_addr = broker_addr.as_str(),
"resolved fetch broker"
);
let mut leader_client = self.connect_or_reuse_broker(&broker_addr).await?;
let rack_id = self
.config
.client_config()
.client_rack_ref()
.map(str::to_owned)
.unwrap_or_default();
let session = self
.fetch_sessions
.get(&broker_addr)
.copied()
.unwrap_or_default();
let (session_id, session_epoch) = session.next_request();
let (
error_code,
preferred_read_replica,
aborted_transactions,
records,
response_session_id,
) = if leader_client.supports_fetch_v12().await? {
let response = leader_client
.fetch_one_v12(FetchOneRequestV12 {
replica_id: -1,
max_wait_ms: self.config.max_wait_ms,
min_bytes: self.config.min_bytes,
max_bytes: self.config.max_partition_bytes,
isolation_level: self.config.isolation_level.as_i8(),
topic: topic.to_owned(),
partition_index: partition,
current_leader_epoch,
fetch_offset: offset,
last_fetched_epoch: current_leader_epoch,
max_partition_bytes: self.config.max_partition_bytes,
session_id,
session_epoch,
rack_id,
})
.await?;
if response.error_code != 0 {
self.fetch_sessions.remove(&broker_addr);
return Err(self.config.client.broker_error(
response.error_code,
format!("fetch {topic}-{partition}@{offset}"),
));
}
let partition_response = fetch_partition_response_v12(&response, topic, partition)?;
(
partition_response.error_code,
Some(partition_response.preferred_read_replica),
partition_response
.aborted_transactions
.iter()
.map(|transaction| AbortedTransactionV4 {
producer_id: transaction.producer_id,
first_offset: transaction.first_offset,
})
.collect(),
partition_response.records.clone(),
response.session_id,
)
} else if leader_client.supports_fetch_v11().await? {
let response = leader_client
.fetch_one_v11(FetchOneRequestV11 {
replica_id: -1,
max_wait_ms: self.config.max_wait_ms,
min_bytes: self.config.min_bytes,
max_bytes: self.config.max_partition_bytes,
isolation_level: self.config.isolation_level.as_i8(),
topic: topic.to_owned(),
partition_index: partition,
current_leader_epoch,
fetch_offset: offset,
max_partition_bytes: self.config.max_partition_bytes,
session_id,
session_epoch,
rack_id,
})
.await?;
if response.error_code != 0 {
self.fetch_sessions.remove(&broker_addr);
return Err(self.config.client.broker_error(
response.error_code,
format!("fetch {topic}-{partition}@{offset}"),
));
}
let partition_response = fetch_partition_response_v11(&response, topic, partition)?;
(
partition_response.error_code,
Some(partition_response.preferred_read_replica),
partition_response.aborted_transactions.clone(),
partition_response.records.clone(),
response.session_id,
)
} else {
let response = leader_client
.fetch_one_v4(FetchOneRequestV4 {
replica_id: -1,
max_wait_ms: self.config.max_wait_ms,
min_bytes: self.config.min_bytes,
max_bytes: self.config.max_partition_bytes,
isolation_level: self.config.isolation_level.as_i8(),
topic: topic.to_owned(),
partition_index: partition,
fetch_offset: offset,
max_partition_bytes: self.config.max_partition_bytes,
})
.await?;
let partition_response = fetch_partition_response(&response, topic, partition)?;
(
partition_response.error_code,
Some(-1),
partition_response.aborted_transactions.clone(),
partition_response.records.clone(),
0,
)
};
if error_code != 0 {
self.fetch_sessions.remove(&broker_addr);
return Err(self
.config
.client
.broker_error(error_code, format!("fetch {topic}-{partition}@{offset}")));
}
if response_session_id > 0 {
self.fetch_sessions.insert(
broker_addr.clone(),
session.advance_with_response(response_session_id),
);
} else {
self.fetch_sessions.remove(&broker_addr);
}
if let Some(preferred_read_replica) = preferred_read_replica {
if preferred_read_replica >= 0
&& broker_addr_for(&metadata, preferred_read_replica).is_ok()
{
self.preferred_read_replicas
.insert((topic.to_owned(), partition), preferred_read_replica);
} else {
self.preferred_read_replicas
.remove(&(topic.to_owned(), partition));
}
}
let next_offset = records
.last()
.map(|record| record.offset.saturating_add(1))
.unwrap_or(offset);
let leader_epoch = records
.last()
.map(|record| record.leader_epoch)
.filter(|leader_epoch| *leader_epoch >= 0);
let records = visible_records(&aborted_transactions, &records, self.config.isolation_level)
.into_iter()
.map(|record| ConsumerRecord::from_message_set(topic, partition, record))
.collect();
self.broker_clients.insert(broker_addr, leader_client);
Ok(FetchedPartition {
records,
next_offset,
leader_epoch,
})
}
fn update_assignment_offset(&mut self, topic: &str, partition: i32, next_offset: i64) {
if let Some(assignment) = self
.assignments
.iter_mut()
.find(|assignment| assignment.topic == topic && assignment.partition == partition)
{
assignment.next_offset = next_offset;
}
}
fn update_assignment_leader_epoch(&mut self, topic: &str, partition: i32, leader_epoch: i32) {
if let Some(assignment) = self
.assignments
.iter_mut()
.find(|assignment| assignment.topic == topic && assignment.partition == partition)
{
assignment.set_leader_epoch(assignment.leader_epoch.max(leader_epoch));
}
}
fn assignment(&self, topic: &str, partition: i32) -> Option<&ConsumerAssignment> {
self.assignments
.iter()
.find(|assignment| assignment.topic == topic && assignment.partition == partition)
}
fn assignment_mut(&mut self, topic: &str, partition: i32) -> Result<&mut ConsumerAssignment> {
self.assignments
.iter_mut()
.find(|assignment| assignment.topic == topic && assignment.partition == partition)
.ok_or_else(|| Error::UnassignedTopicPartition {
topic: topic.to_owned(),
partition,
})
}
async fn metadata_for_topic(&mut self, topic: &str) -> Result<MetadataResponseV1> {
if let Some(metadata) = self.metadata_cache.get(topic) {
return Ok(metadata.clone());
}
let metadata = self.request_metadata_for_topic(topic).await?;
self.metadata_cache
.insert(topic.to_owned(), metadata.clone());
Ok(metadata)
}
async fn connect_or_reuse_broker(&mut self, broker_addr: &str) -> Result<Client> {
if let Some(client) = self.broker_clients.remove(broker_addr) {
return Ok(client);
}
self.fetch_sessions.remove(broker_addr);
self.config
.client
.connect_broker(broker_addr.to_owned())
.await
}
async fn request_metadata_for_topic(&mut self, topic: &str) -> Result<MetadataResponseV1> {
let topics = Some(vec![topic.to_owned()]);
match self.client.metadata(topics.clone()).await {
Ok(metadata) => Ok(metadata),
Err(error) if can_retry_fetch(&error) => {
self.config.client.record_retry();
debug!(
topic,
error = %error,
"reconnecting metadata client after metadata request failure"
);
self.client = self.config.client.clone().connect().await?;
self.client.metadata(topics).await
}
Err(error) => Err(error),
}
}
}
#[derive(Debug)]
struct FetchedPartition {
records: Vec<ConsumerRecord>,
next_offset: i64,
leader_epoch: Option<i32>,
}
struct PartitionRoute {
records: Vec<ConsumerRecord>,
next_offset: Option<i64>,
queued_count: usize,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerAssignment {
topic: String,
partition: i32,
next_offset: i64,
leader_epoch: i32,
paused: bool,
}
impl ConsumerAssignment {
pub(crate) fn new(topic: String, partition: i32, next_offset: i64) -> Self {
Self {
topic,
partition,
next_offset,
leader_epoch: -1,
paused: false,
}
}
pub(crate) fn set_leader_epoch(&mut self, leader_epoch: i32) {
self.leader_epoch = leader_epoch;
}
pub fn topic(&self) -> &str {
&self.topic
}
pub fn partition(&self) -> i32 {
self.partition
}
pub fn next_offset(&self) -> i64 {
self.next_offset
}
pub fn leader_epoch(&self) -> i32 {
self.leader_epoch
}
pub fn is_paused(&self) -> bool {
self.paused
}
}
fn assign_partition(
assignments: &mut Vec<ConsumerAssignment>,
topic: String,
partition: i32,
offset: i64,
) {
if let Some(assignment) = assignments
.iter_mut()
.find(|assignment| assignment.topic == topic && assignment.partition == partition)
{
assignment.next_offset = offset;
assignment.leader_epoch = -1;
return;
}
assignments.push(ConsumerAssignment {
topic,
partition,
next_offset: offset,
leader_epoch: -1,
paused: false,
});
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerConfig {
client: ClientConfig,
max_wait_ms: i32,
min_bytes: i32,
max_partition_bytes: i32,
max_retries: u32,
max_poll_records: usize,
partition_queue_capacity: usize,
isolation_level: IsolationLevel,
offset_reset_policy: OffsetResetPolicy,
}
impl ConsumerConfig {
pub fn new(bootstrap_servers: impl IntoIterator<Item = impl Into<String>>) -> Self {
Self::from_client_config(ClientConfig::new(bootstrap_servers))
}
pub(crate) fn from_client_config(client: ClientConfig) -> Self {
Self {
client,
max_wait_ms: 500,
min_bytes: 1,
max_partition_bytes: 1_048_576,
max_retries: 1,
max_poll_records: 500,
partition_queue_capacity: 1024,
isolation_level: IsolationLevel::ReadUncommitted,
offset_reset_policy: OffsetResetPolicy::Offset(0),
}
}
pub fn client_id(mut self, client_id: impl Into<String>) -> Self {
self.client = self.client.client_id(client_id);
self
}
pub fn client_rack(mut self, client_rack: impl Into<String>) -> Self {
self.client = self.client.client_rack(client_rack);
self
}
pub fn request_timeout_ms(mut self, request_timeout_ms: u64) -> Self {
self.client = self.client.request_timeout_ms(request_timeout_ms);
self
}
pub fn max_response_bytes(mut self, max_response_bytes: usize) -> Self {
self.client = self.client.max_response_bytes(max_response_bytes);
self
}
pub fn max_decode_array_elements(mut self, max: usize) -> Self {
self.client = self.client.max_decode_array_elements(max);
self
}
pub fn max_decompressed_record_bytes(mut self, max: usize) -> Self {
self.client = self.client.max_decompressed_record_bytes(max);
self
}
pub fn metrics(mut self, metrics: ClientMetrics) -> Self {
self.client = self.client.metrics(metrics);
self
}
pub fn security_protocol(mut self, security_protocol: SecurityProtocol) -> Self {
self.client = self.client.security_protocol(security_protocol);
self
}
pub fn tls_server_name(mut self, server_name: impl Into<String>) -> Self {
self.client = self.client.tls_server_name(server_name);
self
}
pub fn tls_root_certificate_der(mut self, certificate: impl Into<Vec<u8>>) -> Self {
self.client = self.client.tls_root_certificate_der(certificate);
self
}
pub fn sasl_plain(mut self, username: impl Into<String>, password: impl Into<String>) -> Self {
self.client = self.client.sasl_plain(username, password);
self
}
pub fn sasl_scram_sha_256(
mut self,
username: impl Into<String>,
password: impl Into<String>,
) -> Self {
self.client = self.client.sasl_scram_sha_256(username, password);
self
}
pub fn sasl_scram_sha_512(
mut self,
username: impl Into<String>,
password: impl Into<String>,
) -> Self {
self.client = self.client.sasl_scram_sha_512(username, password);
self
}
pub fn sasl_oauthbearer(mut self, token: impl Into<String>) -> Self {
self.client = self.client.sasl_oauthbearer(token);
self
}
pub fn sasl_oauthbearer_with_username(
mut self,
username: impl Into<String>,
token: impl Into<String>,
) -> Self {
self.client = self.client.sasl_oauthbearer_with_username(username, token);
self
}
pub fn sasl_oauthbearer_provider<P>(mut self, provider: P) -> Self
where
P: OAuthBearerTokenProvider + 'static,
{
self.client = self.client.sasl_oauthbearer_provider(provider);
self
}
pub fn sasl_oauthbearer_with_username_and_provider<P>(
mut self,
username: impl Into<String>,
provider: P,
) -> Self
where
P: OAuthBearerTokenProvider + 'static,
{
self.client = self
.client
.sasl_oauthbearer_with_username_and_provider(username, provider);
self
}
pub fn max_wait_ms(mut self, max_wait_ms: i32) -> Self {
self.max_wait_ms = max_wait_ms;
self
}
pub fn min_bytes(mut self, min_bytes: i32) -> Self {
self.min_bytes = min_bytes;
self
}
pub fn max_partition_bytes(mut self, max_partition_bytes: i32) -> Self {
self.max_partition_bytes = max_partition_bytes;
self
}
pub fn max_retries(mut self, max_retries: u32) -> Self {
self.max_retries = max_retries;
self
}
pub fn max_retries_ref(&self) -> u32 {
self.max_retries
}
pub fn max_poll_records(mut self, max_poll_records: usize) -> Self {
self.max_poll_records = max_poll_records;
self
}
pub fn max_poll_records_ref(&self) -> usize {
self.max_poll_records
}
pub fn partition_queue_capacity(mut self, partition_queue_capacity: usize) -> Self {
self.partition_queue_capacity = partition_queue_capacity.max(1);
self
}
pub fn partition_queue_capacity_ref(&self) -> usize {
self.partition_queue_capacity
}
pub fn isolation_level(mut self, isolation_level: IsolationLevel) -> Self {
self.isolation_level = isolation_level;
self
}
pub fn isolation_level_ref(&self) -> IsolationLevel {
self.isolation_level
}
pub fn offset_reset_policy(mut self, offset_reset_policy: OffsetResetPolicy) -> Self {
self.offset_reset_policy = offset_reset_policy;
self
}
pub fn offset_reset_policy_ref(&self) -> OffsetResetPolicy {
self.offset_reset_policy
}
pub fn client_config(&self) -> &ClientConfig {
&self.client
}
pub fn validate(&self) -> Result<()> {
self.client.validate()?;
self.validate_values()
}
pub async fn build(self) -> Result<Consumer> {
self.validate()?;
let client = self.client.clone().connect().await?;
Ok(Consumer {
client,
config: self,
assignments: Vec::new(),
partition_queues: BTreeMap::new(),
metadata_cache: BTreeMap::new(),
broker_clients: BTreeMap::new(),
fetch_sessions: BTreeMap::new(),
preferred_read_replicas: BTreeMap::new(),
})
}
fn validate_values(&self) -> Result<()> {
if self.max_wait_ms < 0 {
return Err(Error::InvalidConfiguration {
field: "max_wait_ms",
reason: "must not be negative",
});
}
if self.min_bytes < 0 {
return Err(Error::InvalidConfiguration {
field: "min_bytes",
reason: "must not be negative",
});
}
if self.max_partition_bytes <= 0 {
return Err(Error::InvalidConfiguration {
field: "max_partition_bytes",
reason: "must be greater than zero",
});
}
if self.max_poll_records == 0 {
return Err(Error::InvalidConfiguration {
field: "max_poll_records",
reason: "must be greater than zero",
});
}
Ok(())
}
}
fn leader_for(
metadata: &MetadataResponseV1,
topic_name: &str,
partition_index: i32,
) -> Result<i32> {
metadata
.topics
.iter()
.find(|topic| topic.name == topic_name)
.and_then(|topic| {
topic
.partitions
.iter()
.find(|partition| partition.partition_index == partition_index)
})
.ok_or_else(|| Error::UnknownTopicOrPartition {
topic: topic_name.to_owned(),
partition: partition_index,
})
.and_then(|partition| {
(partition.leader_id >= 0)
.then_some(partition.leader_id)
.ok_or_else(|| Error::MissingLeader {
topic: topic_name.to_owned(),
partition: partition_index,
})
})
}
fn broker_addr_for(metadata: &MetadataResponseV1, node_id: i32) -> Result<String> {
metadata
.brokers
.iter()
.find(|broker| broker.node_id == node_id)
.map(broker_addr)
.ok_or(Error::MissingBroker { node_id })
}
fn broker_addr(broker: &BrokerMetadata) -> String {
format!("{}:{}", broker.host, broker.port)
}
fn fetch_partition_response<'a>(
response: &'a FetchResponseV4,
topic_name: &str,
partition_index: i32,
) -> Result<&'a FetchPartitionResponseV4> {
response
.responses
.iter()
.find(|topic| topic.name == topic_name)
.and_then(|topic| {
topic
.partitions
.iter()
.find(|partition| partition.partition_index == partition_index)
})
.ok_or_else(|| Error::UnknownTopicOrPartition {
topic: topic_name.to_owned(),
partition: partition_index,
})
}
fn fetch_partition_response_v11<'a>(
response: &'a FetchResponseV11,
topic_name: &str,
partition_index: i32,
) -> Result<&'a FetchPartitionResponseV11> {
response
.responses
.iter()
.find(|topic| topic.name == topic_name)
.and_then(|topic| {
topic
.partitions
.iter()
.find(|partition| partition.partition_index == partition_index)
})
.ok_or_else(|| Error::UnknownTopicOrPartition {
topic: topic_name.to_owned(),
partition: partition_index,
})
}
fn fetch_partition_response_v12<'a>(
response: &'a FetchResponseV12,
topic_name: &str,
partition_index: i32,
) -> Result<&'a FetchPartitionResponseV12> {
response
.responses
.iter()
.find(|topic| topic.name == topic_name)
.and_then(|topic| {
topic
.partitions
.iter()
.find(|partition| partition.partition_index == partition_index)
})
.ok_or_else(|| Error::UnknownTopicOrPartition {
topic: topic_name.to_owned(),
partition: partition_index,
})
}
fn list_offset_partition_response<'a>(
topics: &'a [ListOffsetsTopicResponseV1],
topic_name: &str,
partition_index: i32,
) -> Result<&'a ListOffsetsPartitionResponseV1> {
topics
.iter()
.find(|topic| topic.name == topic_name)
.and_then(|topic| {
topic
.partitions
.iter()
.find(|partition| partition.partition_index == partition_index)
})
.ok_or_else(|| Error::UnknownTopicOrPartition {
topic: topic_name.to_owned(),
partition: partition_index,
})
}
fn offset_for_leader_epoch_partition_response<'a>(
topics: &'a [OffsetForLeaderEpochTopicResponseV3],
topic_name: &str,
partition_index: i32,
) -> Result<&'a OffsetForLeaderEpochPartitionResponseV3> {
topics
.iter()
.find(|topic| topic.name == topic_name)
.and_then(|topic| {
topic
.partitions
.iter()
.find(|partition| partition.partition_index == partition_index)
})
.ok_or_else(|| Error::UnknownTopicOrPartition {
topic: topic_name.to_owned(),
partition: partition_index,
})
}
fn metadata_v12_partition<'a>(
metadata: &'a MetadataResponseV12,
topic_name: &str,
partition_index: i32,
) -> Result<&'a MetadataPartitionV12> {
let topic = metadata
.topics
.iter()
.find(|topic| topic.name.as_deref() == Some(topic_name))
.ok_or_else(|| Error::UnknownTopicOrPartition {
topic: topic_name.to_owned(),
partition: partition_index,
})?;
if topic.error_code != 0 {
return Err(Error::Broker {
code: topic.error_code,
context: format!("metadata topic {topic_name}"),
});
}
topic
.partitions
.iter()
.find(|partition| partition.partition_index == partition_index)
.ok_or_else(|| Error::UnknownTopicOrPartition {
topic: topic_name.to_owned(),
partition: partition_index,
})
}
fn visible_records(
aborted_transactions: &[kafrust_protocol::api::fetch::AbortedTransactionV4],
input_records: &[MessageSetRecord],
isolation_level: IsolationLevel,
) -> Vec<MessageSetRecord> {
let mut aborted_transactions = aborted_transactions.to_vec();
let mut visible = Vec::new();
for record in input_records {
if record.control {
if let Some(producer_id) = record.producer_id {
if let Some(index) = aborted_transactions.iter().position(|transaction| {
transaction.producer_id == producer_id
&& transaction.first_offset <= record.offset
}) {
aborted_transactions.remove(index);
}
}
continue;
}
let aborted = isolation_level == IsolationLevel::ReadCommitted
&& record.transactional
&& record.producer_id.is_some_and(|producer_id| {
aborted_transactions.iter().any(|transaction| {
transaction.producer_id == producer_id
&& transaction.first_offset <= record.offset
})
});
if !aborted {
visible.push(record.clone());
}
}
visible
}
fn can_retry_fetch(error: &Error) -> bool {
match error {
Error::Broker { code, .. } => matches!(
BrokerErrorKind::from_code(*code),
BrokerErrorKind::UnknownTopicOrPartition
| BrokerErrorKind::LeaderNotAvailable
| BrokerErrorKind::NotLeaderOrFollower
| BrokerErrorKind::RequestTimedOut
| BrokerErrorKind::ReplicaNotAvailable
| BrokerErrorKind::FencedLeaderEpoch
| BrokerErrorKind::UnknownLeaderEpoch
| BrokerErrorKind::InvalidFetchSessionEpoch
),
Error::Io(_)
| Error::RequestTimedOut { .. }
| Error::UnknownTopicOrPartition { .. }
| Error::MissingLeader { .. }
| Error::MissingBroker { .. } => true,
Error::MissingBootstrapServer
| Error::InvalidPartition { .. }
| Error::UnassignedTopicPartition { .. }
| Error::PartitionQueueFull { .. }
| Error::MissingGroupDescription { .. }
| Error::MissingDeleteGroupResult { .. }
| Error::ResponseCountMismatch { .. }
| Error::MissingSaslCredentials
| Error::InvalidSaslResponse { .. }
| Error::OAuthBearerTokenTimeout { .. }
| Error::TransactionOutcomeUnknown { .. }
| Error::TransactionProducerDefunct
| Error::ResponseTooLarge { .. }
| Error::TlsConfig { .. }
| Error::InvalidTlsServerName { .. }
| Error::InvalidGroupInstanceId
| Error::InvalidTopicPattern { .. }
| Error::InvalidScramCredential { .. }
| Error::InvalidConfiguration { .. }
| Error::Unsupported(_)
| Error::TaskJoin(_)
| Error::Protocol(_) => false,
}
}
fn is_offset_out_of_range(error: &Error) -> bool {
matches!(
error,
Error::Broker { code, .. }
if BrokerErrorKind::from_code(*code) == BrokerErrorKind::OffsetOutOfRange
)
}
fn is_leader_epoch_transition_error(error: &Error) -> bool {
matches!(
error,
Error::Broker { code, .. }
if matches!(
BrokerErrorKind::from_code(*code),
BrokerErrorKind::NotLeaderOrFollower
| BrokerErrorKind::FencedLeaderEpoch
| BrokerErrorKind::UnknownLeaderEpoch
)
)
}
fn limit_fetched_records(
fetched: &mut Vec<ConsumerRecord>,
current_record_count: usize,
max_poll_records: usize,
) {
let remaining = max_poll_records.saturating_sub(current_record_count);
fetched.truncate(remaining);
}
fn invalidate_metadata_cache(
metadata_cache: &mut BTreeMap<String, MetadataResponseV1>,
topic: &str,
) {
metadata_cache.remove(topic);
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::unwrap_used)]
mod tests {
use super::{
assign_partition, can_retry_fetch, invalidate_metadata_cache,
is_leader_epoch_transition_error, leader_for, limit_fetched_records,
offset_for_leader_epoch_partition_response, visible_records, Consumer, ConsumerAssignment,
ConsumerConfig, ConsumerRecord, IsolationLevel, OffsetResetPolicy, PartitionWatermarks,
SecurityProtocol,
};
use crate::{Client, ClientMetrics, Error};
use kafrust_protocol::api::fetch::{
AbortedTransactionV4, FetchPartitionResponseV4, MessageSetRecord,
};
use kafrust_protocol::api::metadata::{
BrokerMetadata, MetadataResponseV1, PartitionMetadata, TopicMetadata,
};
use kafrust_protocol::codec::Encoder;
use std::collections::BTreeMap;
use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
use tokio::net::TcpListener;
#[test]
fn builds_consumer_config() {
let config = ConsumerConfig::new(["localhost:9092"])
.client_id("orders-reader")
.client_rack("rack-a")
.request_timeout_ms(5_000)
.security_protocol(SecurityProtocol::Tls)
.tls_server_name("broker.example.com")
.tls_root_certificate_der([1, 2, 3])
.sasl_plain("alice", "secret-password")
.max_wait_ms(250)
.min_bytes(10)
.max_partition_bytes(1024)
.max_retries(3)
.max_poll_records(10)
.partition_queue_capacity(7)
.isolation_level(IsolationLevel::ReadCommitted);
assert_eq!(
config.client_config().client_id_ref(),
Some("orders-reader")
);
assert_eq!(config.client_config().client_rack_ref(), Some("rack-a"));
assert_eq!(
config.client_config().security_protocol_ref(),
SecurityProtocol::Tls
);
assert_eq!(
config.client_config().tls_server_name_ref(),
Some("broker.example.com")
);
assert_eq!(
config.client_config().tls_root_certificates_der(),
&[vec![1, 2, 3]]
);
assert_eq!(
config
.client_config()
.sasl_credentials_ref()
.unwrap()
.username(),
"alice"
);
assert_eq!(config.max_retries_ref(), 3);
assert_eq!(config.max_poll_records_ref(), 10);
assert_eq!(config.partition_queue_capacity_ref(), 7);
assert_eq!(config.isolation_level_ref(), IsolationLevel::ReadCommitted);
assert_eq!(
config.offset_reset_policy_ref(),
OffsetResetPolicy::Offset(0)
);
}
#[test]
fn normalizes_zero_partition_queue_capacity() {
assert_eq!(
ConsumerConfig::new(["localhost:9092"])
.partition_queue_capacity(0)
.partition_queue_capacity_ref(),
1
);
}
#[tokio::test]
async fn rejects_invalid_fetch_configuration_before_connecting() {
let cases = [
(
ConsumerConfig::new(["127.0.0.1:1"])
.max_wait_ms(-1)
.build()
.await
.unwrap_err(),
"max_wait_ms",
),
(
ConsumerConfig::new(["127.0.0.1:1"])
.min_bytes(-1)
.build()
.await
.unwrap_err(),
"min_bytes",
),
(
ConsumerConfig::new(["127.0.0.1:1"])
.max_partition_bytes(0)
.build()
.await
.unwrap_err(),
"max_partition_bytes",
),
(
ConsumerConfig::new(["127.0.0.1:1"])
.max_poll_records(0)
.build()
.await
.unwrap_err(),
"max_poll_records",
),
];
for (error, field) in cases {
assert!(matches!(
error,
Error::InvalidConfiguration {
field: actual,
..
} if actual == field
));
}
assert!(ConsumerConfig::new(["127.0.0.1:1"]).validate().is_ok());
}
#[tokio::test]
async fn split_partition_queue_requires_assignment_and_protects_seek_state() {
let (client_stream, _broker_stream) = tokio::io::duplex(64);
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-partition-queue-api-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new(["localhost:9092"]);
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
assert!(matches!(
consumer.split_partition_queue("orders", 0).unwrap_err(),
Error::UnassignedTopicPartition { topic, partition: 0 } if topic == "orders"
));
consumer.assign("orders", 0, 10);
let queue = consumer.split_partition_queue("orders", 0).unwrap();
assert!(matches!(
consumer.split_partition_queue("orders", 0).unwrap_err(),
Error::Unsupported("partition queue is already split")
));
assert!(matches!(
consumer.seek("orders", 0, 20).unwrap_err(),
Error::Unsupported("seek requires dropping the active partition queue")
));
drop(queue);
consumer.seek("orders", 0, 20).unwrap();
assert_eq!(consumer.position("orders", 0), Some(20));
}
#[test]
fn maps_message_set_record() {
let record = ConsumerRecord::from_message_set(
"orders",
1,
MessageSetRecord {
offset: 42,
leader_epoch: 4,
timestamp_ms: 123,
key: Some(b"order-1".to_vec()),
value: Some(b"created".to_vec()),
headers: vec![kafrust_protocol::api::produce::RecordBatchHeader::new(
"source",
Some(b"checkout".to_vec()),
)],
producer_id: None,
transactional: false,
control: false,
},
);
assert_eq!(record.topic(), "orders");
assert_eq!(record.partition(), 1);
assert_eq!(record.offset(), 42);
assert_eq!(record.leader_epoch(), 4);
assert_eq!(record.timestamp_ms(), 123);
assert_eq!(record.key().unwrap(), b"order-1");
assert_eq!(record.value().unwrap(), b"created");
assert_eq!(record.headers().len(), 1);
assert_eq!(record.headers()[0].key(), "source");
assert_eq!(record.headers()[0].value(), Some(&b"checkout"[..]));
}
#[test]
fn tracks_assignments() {
let mut assignments = Vec::<ConsumerAssignment>::new();
assign_partition(&mut assignments, "orders".to_owned(), 0, 10);
assign_partition(&mut assignments, "orders".to_owned(), 0, 20);
assert_eq!(assignments.len(), 1);
assert_eq!(assignments[0].topic(), "orders");
assert_eq!(assignments[0].partition(), 0);
assert_eq!(assignments[0].next_offset(), 20);
assert!(!assignments[0].is_paused());
}
#[test]
fn tracks_and_resets_assignment_leader_epoch() {
let (client_stream, _broker_stream) = tokio::io::duplex(64);
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-leader-epoch-state-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let mut consumer =
Consumer::from_assignments(client, ConsumerConfig::new(["localhost:9092"]), Vec::new());
consumer.assign("orders", 0, 10);
assert_eq!(consumer.assignments()[0].leader_epoch(), -1);
consumer.update_assignment_leader_epoch("orders", 0, 4);
assert_eq!(consumer.assignments()[0].leader_epoch(), 4);
consumer.assign("orders", 0, 20);
assert_eq!(consumer.assignments()[0].leader_epoch(), -1);
}
#[tokio::test]
async fn controls_assignment_position_and_pause_state() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let _connection = listener.accept().await.unwrap();
});
let mut consumer = ConsumerConfig::new([addr.to_string()])
.build()
.await
.unwrap();
consumer.assign("orders", 0, 10);
assert_eq!(consumer.position("orders", 0), Some(10));
consumer.seek("orders", 0, 20).unwrap();
consumer.pause("orders", 0).unwrap();
assert_eq!(consumer.position("orders", 0), Some(20));
assert!(consumer.assignments()[0].is_paused());
assert!(consumer.poll().await.unwrap().is_empty());
consumer.resume("orders", 0).unwrap();
assert!(!consumer.assignments()[0].is_paused());
assert!(matches!(
consumer.seek("orders", 1, 0).unwrap_err(),
Error::UnassignedTopicPartition {
topic,
partition: 1
} if topic == "orders"
));
server.await.unwrap();
}
#[tokio::test]
async fn fetches_partition_watermarks_from_partition_leader() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
for (correlation_id, timestamp, offset) in [(1, -2_i64, 4_i64), (2, -1_i64, 9_i64)] {
let request = read_frame(&mut socket).await;
assert_eq!(&request[0..4], &[0, 2, 0, 1]);
assert_eq!(
i64::from_be_bytes(request[request.len() - 8..].try_into().unwrap()),
timestamp
);
write_frame(
&mut socket,
&list_offsets_response_frame(correlation_id, offset),
)
.await;
}
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-watermarks-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()]).request_timeout_ms(500);
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
let watermarks = consumer.fetch_watermarks("orders", 0).await.unwrap();
assert_eq!(watermarks, PartitionWatermarks { low: 4, high: 9 });
assert_eq!(watermarks.low(), 4);
assert_eq!(watermarks.high(), 9);
server.await.unwrap();
}
#[tokio::test]
async fn resets_out_of_range_assignment_to_earliest_offset() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut first_socket, _) = listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut first_socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut first_socket, &api_versions_v3_fetch_v12_response(1)).await;
let first_fetch = read_frame(&mut first_socket).await;
assert_eq!(&first_fetch[0..4], &[0, 1, 0, 12]);
write_frame(&mut first_socket, &fetch_v12_out_of_range_response_frame(2)).await;
drop(first_socket);
let (mut second_socket, _) = listener.accept().await.unwrap();
for (correlation_id, timestamp, offset) in [(1, -2_i64, 4_i64), (2, -1_i64, 9_i64)] {
let request = read_frame(&mut second_socket).await;
assert_eq!(&request[0..4], &[0, 2, 0, 1]);
assert_eq!(
i64::from_be_bytes(request[request.len() - 8..].try_into().unwrap()),
timestamp
);
write_frame(
&mut second_socket,
&list_offsets_response_frame(correlation_id, offset),
)
.await;
}
let api_versions_request = read_frame(&mut second_socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut second_socket, &api_versions_v3_fetch_v12_response(3)).await;
let second_fetch = read_frame(&mut second_socket).await;
assert_eq!(&second_fetch[0..4], &[0, 1, 0, 12]);
write_frame(
&mut second_socket,
&fetch_v12_response_frame_with_record(4, 0),
)
.await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-offset-reset-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()])
.request_timeout_ms(500)
.offset_reset_policy(OffsetResetPolicy::Earliest);
let mut consumer = Consumer::from_assignments(
client,
config,
vec![ConsumerAssignment::new("orders".to_owned(), 0, 100)],
);
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
let records = consumer.poll().await.unwrap();
assert_eq!(records.len(), 1);
assert_eq!(records[0].offset(), 42);
assert_eq!(consumer.position("orders", 0), Some(43));
server.await.unwrap();
}
#[tokio::test]
async fn sends_assignment_leader_epoch_in_fetch_v12_request() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
let fetch_request = read_frame(&mut socket).await;
assert_eq!(&fetch_request[0..4], &[0, 1, 0, 12]);
assert_eq!(
fetch_request
.windows(4)
.filter(|window| *window == [0, 0, 0, 8])
.count(),
2
);
write_frame(&mut socket, &fetch_v12_response_frame(2, -1)).await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-assignment-leader-epoch-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()])
.request_timeout_ms(500)
.client_rack("rack-a");
let mut consumer = Consumer::from_assignments(
client,
config,
vec![ConsumerAssignment::new("orders".to_owned(), 0, 42)],
);
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
consumer.update_assignment_leader_epoch("orders", 0, 8);
assert!(consumer.poll().await.unwrap().is_empty());
assert_eq!(consumer.assignments()[0].leader_epoch(), 8);
server.await.unwrap();
}
#[tokio::test]
async fn recovers_assignment_after_leader_epoch_truncation() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut bootstrap, _) =
tokio::time::timeout(std::time::Duration::from_secs(2), listener.accept())
.await
.expect("bootstrap accept timed out")
.unwrap();
let (mut first_fetch, _) =
tokio::time::timeout(std::time::Duration::from_secs(2), listener.accept())
.await
.expect("first fetch accept timed out")
.unwrap();
let api_versions_request = read_frame(&mut first_fetch).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(
&mut first_fetch,
&api_versions_v3_metadata_fetch_response(1, 12),
)
.await;
let fetch_request = read_frame(&mut first_fetch).await;
assert_eq!(&fetch_request[0..4], &[0, 1, 0, 12]);
assert!(fetch_request
.windows(8)
.any(|window| window == 100_i64.to_be_bytes()));
write_frame(
&mut first_fetch,
&fetch_v12_response_frame_with_partition_error(2, 74),
)
.await;
let api_versions_request = read_frame(&mut bootstrap).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(
&mut bootstrap,
&api_versions_v3_metadata_fetch_response(1, 12),
)
.await;
let metadata_v12_request = read_frame(&mut bootstrap).await;
assert_eq!(&metadata_v12_request[0..4], &[0, 3, 0, 12]);
write_frame(&mut bootstrap, &metadata_v12_response_frame(2, &addr, 5)).await;
let metadata_request = read_frame(&mut bootstrap).await;
assert_eq!(&metadata_request[0..4], &[0, 3, 0, 1]);
write_frame(&mut bootstrap, &metadata_response_frame_for(3, &addr)).await;
let (mut recovery, _) = listener.accept().await.unwrap();
let epoch_request = read_frame(&mut recovery).await;
assert_eq!(&epoch_request[0..4], &[0, 23, 0, 3]);
write_frame(
&mut recovery,
&offset_for_leader_epoch_response_frame_with(1, 5, 50),
)
.await;
let api_versions_request = read_frame(&mut recovery).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(
&mut recovery,
&api_versions_v3_metadata_fetch_response(2, 12),
)
.await;
let retry_fetch_request = read_frame(&mut recovery).await;
assert_eq!(&retry_fetch_request[0..4], &[0, 1, 0, 12]);
assert!(retry_fetch_request
.windows(8)
.any(|window| window == 50_i64.to_be_bytes()));
assert!(retry_fetch_request
.windows(4)
.any(|window| window == 5_i32.to_be_bytes()));
write_frame(
&mut recovery,
&fetch_v12_response_frame_with_record_at(3, 0, 50),
)
.await;
});
let mut consumer = ConsumerConfig::new([addr.to_string()])
.request_timeout_ms(500)
.max_retries(2)
.build()
.await
.unwrap();
consumer.assign("orders", 0, 100);
consumer.update_assignment_leader_epoch("orders", 0, 4);
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
let records = tokio::time::timeout(std::time::Duration::from_secs(5), consumer.poll())
.await
.expect("leader epoch recovery poll timed out")
.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(5), server)
.await
.expect("leader epoch recovery server timed out")
.unwrap();
assert_eq!(records.len(), 1);
assert_eq!(records[0].offset(), 50);
assert_eq!(consumer.position("orders", 0), Some(51));
assert_eq!(consumer.assignments()[0].leader_epoch(), 5);
}
#[tokio::test]
async fn reuses_fetch_session_for_sequential_rack_aware_polls() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
let first_fetch = read_frame(&mut socket).await;
assert_eq!(&first_fetch[0..4], &[0, 1, 0, 12]);
assert!(first_fetch.windows(4).any(|window| window == [0, 0, 0, 0]));
write_frame(
&mut socket,
&fetch_v12_response_frame_with_session(2, -1, 17),
)
.await;
let second_fetch = read_frame(&mut socket).await;
assert_eq!(&second_fetch[0..4], &[0, 1, 0, 12]);
assert!(second_fetch
.windows(4)
.any(|window| window == [0, 0, 0, 17]));
assert!(second_fetch.windows(4).any(|window| window == [0, 0, 0, 1]));
write_frame(
&mut socket,
&fetch_v12_response_frame_with_session(3, -1, 17),
)
.await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-fetch-session-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()])
.request_timeout_ms(500)
.client_rack("rack-a");
let mut consumer = Consumer::from_assignments(
client,
config,
vec![ConsumerAssignment::new("orders".to_owned(), 0, 42)],
);
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
assert!(consumer.poll().await.unwrap().is_empty());
assert_eq!(consumer.fetch_sessions[&addr.to_string()].session_id, 17);
assert_eq!(consumer.fetch_sessions[&addr.to_string()].next_epoch, 1);
assert!(consumer.poll().await.unwrap().is_empty());
assert_eq!(consumer.fetch_sessions[&addr.to_string()].next_epoch, 2);
server.await.unwrap();
}
#[tokio::test]
async fn fetches_offset_for_leader_epoch_from_partition_leader() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let request = read_frame(&mut socket).await;
assert_eq!(&request[0..4], &[0, 23, 0, 3]);
write_frame(&mut socket, &offset_for_leader_epoch_response_frame()).await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-leader-epoch-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()]).request_timeout_ms(500);
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
let result = consumer
.offset_for_leader_epoch("orders", 0, 9, 7)
.await
.unwrap();
assert_eq!(result.leader_epoch(), 8);
assert_eq!(result.end_offset(), 42);
server.await.unwrap();
}
#[test]
fn reports_missing_offset_for_leader_epoch_partition() {
let error = offset_for_leader_epoch_partition_response(&[], "orders", 0).unwrap_err();
assert!(matches!(
error,
Error::UnknownTopicOrPartition { topic, partition: 0 } if topic == "orders"
));
}
#[test]
fn resolves_partition_leader() {
assert_eq!(leader_for(&metadata_fixture(), "orders", 0).unwrap(), 1);
}
#[test]
fn classifies_retriable_fetch_errors() {
assert!(can_retry_fetch(&Error::Broker {
code: 6,
context: "fetch orders-0@0".to_owned(),
}));
assert!(can_retry_fetch(&Error::RequestTimedOut { timeout_ms: 5 }));
assert!(can_retry_fetch(&Error::Io(std::io::Error::new(
std::io::ErrorKind::ConnectionReset,
"reset",
))));
assert!(can_retry_fetch(&Error::UnknownTopicOrPartition {
topic: "orders".to_owned(),
partition: 3,
}));
assert!(can_retry_fetch(&Error::MissingLeader {
topic: "orders".to_owned(),
partition: 0,
}));
assert!(can_retry_fetch(&Error::MissingBroker { node_id: 2 }));
assert!(can_retry_fetch(&Error::Broker {
code: 74,
context: "fetch orders-0@0".to_owned(),
}));
assert!(can_retry_fetch(&Error::Broker {
code: 75,
context: "fetch orders-0@0".to_owned(),
}));
assert!(can_retry_fetch(&Error::Broker {
code: 70,
context: "fetch orders-0@0".to_owned(),
}));
assert!(!can_retry_fetch(&Error::Broker {
code: 1,
context: "fetch orders-0@0".to_owned(),
}));
assert!(!can_retry_fetch(&Error::Unsupported("fetch v99")));
}
#[test]
fn classifies_not_leader_as_a_leader_epoch_transition() {
assert!(is_leader_epoch_transition_error(&Error::Broker {
code: 6,
context: "fetch orders-0@0".to_owned(),
}));
assert!(is_leader_epoch_transition_error(&Error::Broker {
code: 74,
context: "fetch orders-0@0".to_owned(),
}));
assert!(is_leader_epoch_transition_error(&Error::Broker {
code: 75,
context: "fetch orders-0@0".to_owned(),
}));
assert!(!is_leader_epoch_transition_error(&Error::Broker {
code: 1,
context: "fetch orders-0@0".to_owned(),
}));
}
#[test]
fn limits_fetched_records_to_remaining_poll_budget() {
let mut records = vec![
ConsumerRecord::from_message_set("orders", 0, message(10)),
ConsumerRecord::from_message_set("orders", 0, message(11)),
ConsumerRecord::from_message_set("orders", 0, message(12)),
];
limit_fetched_records(&mut records, 1, 3);
assert_eq!(records.len(), 2);
assert_eq!(records[1].offset(), 11);
}
#[test]
fn clears_fetched_records_when_poll_budget_is_exhausted() {
let mut records = vec![ConsumerRecord::from_message_set("orders", 0, message(10))];
limit_fetched_records(&mut records, 3, 3);
assert!(records.is_empty());
}
#[test]
fn invalidates_topic_metadata_cache() {
let mut cache = BTreeMap::new();
cache.insert("orders".to_owned(), metadata_fixture());
cache.insert("payments".to_owned(), metadata_fixture());
invalidate_metadata_cache(&mut cache, "orders");
assert!(!cache.contains_key("orders"));
assert!(cache.contains_key("payments"));
}
#[test]
fn read_committed_hides_aborted_records_and_control_markers() {
let partition = FetchPartitionResponseV4 {
partition_index: 0,
error_code: 0,
high_watermark: 14,
last_stable_offset: 14,
aborted_transactions: vec![AbortedTransactionV4 {
producer_id: 7,
first_offset: 10,
}],
records: vec![
transactional_message(10, 7, false),
transactional_message(11, 8, false),
transactional_message(12, 7, true),
transactional_message(13, 7, false),
],
};
let committed = visible_records(
&partition.aborted_transactions,
&partition.records,
IsolationLevel::ReadCommitted,
);
let uncommitted = visible_records(
&partition.aborted_transactions,
&partition.records,
IsolationLevel::ReadUncommitted,
);
assert_eq!(
committed
.iter()
.map(|record| record.offset)
.collect::<Vec<_>>(),
vec![11, 13]
);
assert_eq!(
uncommitted
.iter()
.map(|record| record.offset)
.collect::<Vec<_>>(),
vec![10, 11, 13]
);
}
#[tokio::test]
async fn reconnects_metadata_client_after_request_io_error() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let request = read_frame(&mut socket).await;
assert_eq!(&request[0..2], &[0, 3]);
write_frame(&mut socket, &metadata_response_frame()).await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
drop(broker_stream);
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-consumer-test".to_owned()),
Some(std::time::Duration::from_millis(50)),
);
let metrics = ClientMetrics::new();
let config = ConsumerConfig::new([addr.to_string()])
.request_timeout_ms(500)
.metrics(metrics.clone());
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let metadata = consumer.metadata_for_topic("orders").await.unwrap();
assert_eq!(metadata.brokers[0].node_id, 1);
assert_eq!(metadata.topics[0].name, "orders");
assert!(consumer.metadata_cache.contains_key("orders"));
assert_eq!(metrics.snapshot().retries, 1);
server.await.unwrap();
}
#[tokio::test]
async fn records_consumed_metrics_for_fetch_results() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
let request = read_frame(&mut socket).await;
assert_eq!(&request[0..4], &[0, 1, 0, 12]);
write_frame(&mut socket, &fetch_v12_response_frame_with_record(2, 0)).await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-consumer-metrics-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let metrics = ClientMetrics::new();
let config = ConsumerConfig::new([addr.to_string()])
.request_timeout_ms(500)
.metrics(metrics.clone());
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
let records = consumer.fetch("orders", 0, 42).await.unwrap();
assert_eq!(records.len(), 1);
assert_eq!(records[0].offset(), 42);
assert_eq!(metrics.snapshot().consumed_records, 1);
server.await.unwrap();
}
#[tokio::test]
async fn split_partition_queue_routes_poll_records_and_advances_position() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
let request = read_frame(&mut socket).await;
assert_eq!(&request[0..4], &[0, 1, 0, 12]);
write_frame(&mut socket, &fetch_v12_response_frame_with_record(2, 0)).await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-partition-queue-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()])
.request_timeout_ms(500)
.partition_queue_capacity(2);
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
consumer.assign("orders", 0, 42);
let mut queue = consumer.split_partition_queue("orders", 0).unwrap();
assert_eq!(queue.topic(), "orders");
assert_eq!(queue.partition(), 0);
assert!(consumer.poll().await.unwrap().is_empty());
assert_eq!(consumer.position("orders", 0), Some(43));
let record = queue.recv().await.unwrap();
assert_eq!(record.offset(), 42);
assert!(queue.try_recv().is_none());
server.await.unwrap();
}
#[tokio::test]
async fn split_partition_queue_reports_backpressure_without_skipping_records() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
for correlation_id in [2, 3] {
let request = read_frame(&mut socket).await;
assert_eq!(&request[0..4], &[0, 1, 0, 12]);
write_frame(
&mut socket,
&fetch_v12_response_frame_with_record(correlation_id, 0),
)
.await;
}
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-partition-queue-capacity-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()])
.request_timeout_ms(500)
.partition_queue_capacity(1);
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
consumer.assign("orders", 0, 42);
let mut queue = consumer.split_partition_queue("orders", 0).unwrap();
consumer.poll().await.unwrap();
let error = consumer.poll().await.unwrap_err();
assert!(matches!(
error,
Error::PartitionQueueFull {
topic,
partition: 0,
capacity: 1
} if topic == "orders"
));
assert_eq!(consumer.position("orders", 0), Some(43));
assert_eq!(queue.recv().await.unwrap().offset(), 42);
server.await.unwrap();
}
#[tokio::test]
async fn reuses_partition_leader_connection_for_sequential_fetches() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
for (correlation_id, session_id) in [(2, 17), (3, 17)] {
let request = read_frame(&mut socket).await;
assert_eq!(&request[0..4], &[0, 1, 0, 12]);
write_frame(
&mut socket,
&fetch_v12_response_frame_with_record(correlation_id, session_id),
)
.await;
}
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-consumer-reuse-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()]).request_timeout_ms(500);
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
let first = consumer.fetch("orders", 0, 42).await.unwrap();
let second = consumer.fetch("orders", 0, 42).await.unwrap();
assert_eq!(first.len(), 1);
assert_eq!(second.len(), 1);
assert_eq!(consumer.broker_clients.len(), 1);
assert_eq!(consumer.fetch_sessions[&addr.to_string()].session_id, 17);
assert_eq!(consumer.fetch_sessions[&addr.to_string()].next_epoch, 2);
server.await.unwrap();
}
#[tokio::test]
async fn falls_back_to_fetch_v4_when_broker_lacks_session_versions() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut socket, &api_versions_v3_fetch_response(1, 10)).await;
let fetch_request = read_frame(&mut socket).await;
assert_eq!(&fetch_request[0..4], &[0, 1, 0, 4]);
write_frame(&mut socket, &fetch_v4_response_frame()).await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-fetch-v4-fallback-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()]).request_timeout_ms(500);
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
let records = consumer.fetch("orders", 0, 42).await.unwrap();
assert_eq!(records.len(), 1);
assert!(consumer.fetch_sessions.is_empty());
server.await.unwrap();
}
#[tokio::test]
async fn negotiates_rack_aware_fetch_and_routes_to_preferred_replica() {
let leader_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let leader_addr = leader_listener.local_addr().unwrap();
let preferred_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let preferred_addr = preferred_listener.local_addr().unwrap();
let leader_server = tokio::spawn(async move {
let (mut socket, _) = leader_listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
let fetch_request = read_frame(&mut socket).await;
assert_eq!(&fetch_request[0..4], &[0, 1, 0, 12]);
assert!(fetch_request
.windows(b"rack-a".len())
.any(|window| window == b"rack-a"));
assert_eq!(fetch_request.last(), Some(&0));
write_frame(&mut socket, &fetch_v12_response_frame(2, 2)).await;
});
let preferred_server = tokio::spawn(async move {
let (mut socket, _) = preferred_listener.accept().await.unwrap();
let api_versions_request = read_frame(&mut socket).await;
assert_eq!(&api_versions_request[0..4], &[0, 18, 0, 3]);
write_frame(&mut socket, &api_versions_v3_fetch_v12_response(1)).await;
let fetch_request = read_frame(&mut socket).await;
assert_eq!(&fetch_request[0..4], &[0, 1, 0, 12]);
assert!(fetch_request
.windows(b"rack-a".len())
.any(|window| window == b"rack-a"));
assert_eq!(fetch_request.last(), Some(&0));
write_frame(&mut socket, &fetch_v12_response_frame(2, -1)).await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-rack-routing-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([leader_addr.to_string()])
.request_timeout_ms(500)
.client_rack("rack-a");
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let mut metadata = metadata_fixture();
metadata.brokers[0].host = leader_addr.ip().to_string();
metadata.brokers[0].port = i32::from(leader_addr.port());
metadata.brokers.push(BrokerMetadata {
node_id: 2,
host: preferred_addr.ip().to_string(),
port: i32::from(preferred_addr.port()),
rack: Some("rack-a".to_owned()),
});
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
assert!(consumer.fetch("orders", 0, 42).await.unwrap().is_empty());
assert_eq!(
consumer
.preferred_read_replicas
.get(&("orders".to_owned(), 0)),
Some(&2)
);
assert!(consumer.fetch("orders", 0, 42).await.unwrap().is_empty());
assert!(!consumer
.preferred_read_replicas
.contains_key(&("orders".to_owned(), 0)));
leader_server.await.unwrap();
preferred_server.await.unwrap();
}
#[tokio::test]
async fn clears_preferred_replica_after_exhausted_fetch_failure() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let _api_versions_request = read_frame(&mut socket).await;
write_frame(&mut socket, &api_versions_v3_fetch_v11_response(1)).await;
let _fetch_request = read_frame(&mut socket).await;
write_frame(&mut socket, &fetch_v11_error_response_frame(2, 6)).await;
});
let (client_stream, broker_stream) = tokio::io::duplex(64);
let _broker_stream = broker_stream;
let client = Client::from_stream(
Box::new(client_stream),
Some("kafrust-rack-failure-test".to_owned()),
Some(std::time::Duration::from_millis(500)),
);
let config = ConsumerConfig::new([addr.to_string()])
.request_timeout_ms(500)
.max_retries(0)
.client_rack("rack-a");
let mut consumer = Consumer::from_assignments(client, config, Vec::new());
let mut metadata = metadata_fixture();
metadata.brokers[0].host = addr.ip().to_string();
metadata.brokers[0].port = i32::from(addr.port());
consumer
.metadata_cache
.insert("orders".to_owned(), metadata);
consumer
.preferred_read_replicas
.insert(("orders".to_owned(), 0), 1);
assert!(consumer.fetch("orders", 0, 42).await.is_err());
assert!(!consumer
.preferred_read_replicas
.contains_key(&("orders".to_owned(), 0)));
server.await.unwrap();
}
fn message(offset: i64) -> MessageSetRecord {
MessageSetRecord {
offset,
leader_epoch: -1,
timestamp_ms: 123,
key: None,
value: None,
headers: Vec::new(),
producer_id: None,
transactional: false,
control: false,
}
}
fn transactional_message(offset: i64, producer_id: i64, control: bool) -> MessageSetRecord {
MessageSetRecord {
offset,
leader_epoch: -1,
timestamp_ms: 123,
key: None,
value: None,
headers: Vec::new(),
producer_id: Some(producer_id),
transactional: true,
control,
}
}
fn metadata_fixture() -> MetadataResponseV1 {
MetadataResponseV1 {
brokers: vec![BrokerMetadata {
node_id: 1,
host: "localhost".to_owned(),
port: 9092,
rack: None,
}],
controller_id: 1,
topics: vec![TopicMetadata {
error_code: 0,
name: "orders".to_owned(),
is_internal: false,
partitions: vec![PartitionMetadata {
error_code: 0,
partition_index: 0,
leader_id: 1,
replica_nodes: vec![1],
isr_nodes: vec![1],
}],
}],
}
}
async fn read_frame<T>(stream: &mut T) -> Vec<u8>
where
T: AsyncRead + Unpin,
{
let mut size = [0u8; 4];
stream.read_exact(&mut size).await.unwrap();
let size = usize::try_from(i32::from_be_bytes(size)).unwrap();
let mut request = vec![0u8; size];
stream.read_exact(&mut request).await.unwrap();
request
}
async fn write_frame<T>(stream: &mut T, frame: &[u8])
where
T: AsyncWrite + Unpin,
{
stream
.write_all(&(frame.len() as i32).to_be_bytes())
.await
.unwrap();
stream.write_all(frame).await.unwrap();
stream.flush().await.unwrap();
}
fn metadata_response_frame() -> Vec<u8> {
vec![
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, ]
}
fn fetch_v4_response_frame() -> Vec<u8> {
fetch_v4_response_frame_with_correlation(1)
}
fn fetch_v4_response_frame_with_correlation(correlation_id: i32) -> Vec<u8> {
let mut message = Encoder::new();
message.write_i32(0);
message.write_i8(1);
message.write_i8(0);
message.write_i64(123);
message.write_nullable_bytes(Some(b"order-1")).unwrap();
message.write_nullable_bytes(Some(b"created")).unwrap();
let message = message.into_bytes();
let mut records = Encoder::new();
records.write_i64(42);
records.write_i32(i32::try_from(message.len()).unwrap());
records.write_raw(&message);
let records = records.into_bytes();
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_i32(0);
response.write_i32(1);
response.write_string("orders").unwrap();
response.write_i32(1);
response.write_i32(0);
response.write_i16(0);
response.write_i64(43);
response.write_i64(43);
response.write_i32(0);
response.write_bytes(&records).unwrap();
response.into_bytes()
}
fn api_versions_v3_fetch_v11_response(correlation_id: i32) -> Vec<u8> {
api_versions_v3_fetch_response(correlation_id, 11)
}
fn api_versions_v3_fetch_v12_response(correlation_id: i32) -> Vec<u8> {
api_versions_v3_fetch_response(correlation_id, 12)
}
fn api_versions_v3_fetch_response(correlation_id: i32, max_version: i16) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_i16(0);
response.write_i8(2); response.write_i16(1);
response.write_i16(0);
response.write_i16(max_version);
response.write_i8(0); response.write_i32(0);
response.write_i8(0); response.into_bytes()
}
fn api_versions_v3_metadata_fetch_response(
correlation_id: i32,
fetch_max_version: i16,
) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_i16(0);
response.write_unsigned_varint(3); response.write_i16(3); response.write_i16(0);
response.write_i16(12);
response.write_unsigned_varint(0);
response.write_i16(1); response.write_i16(0);
response.write_i16(fetch_max_version);
response.write_unsigned_varint(0);
response.write_i32(0);
response.write_unsigned_varint(0);
response.into_bytes()
}
fn fetch_v11_error_response_frame(correlation_id: i32, error_code: i16) -> Vec<u8> {
fetch_v11_response_frame_with_error(correlation_id, error_code, -1)
}
fn fetch_v11_response_frame_with_error(
correlation_id: i32,
error_code: i16,
preferred_read_replica: i32,
) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_i32(0);
response.write_i16(0);
response.write_i32(0);
response.write_i32(1);
response.write_string("orders").unwrap();
response.write_i32(1);
response.write_i32(0);
response.write_i16(error_code);
response.write_i64(43);
response.write_i64(43);
response.write_i64(42);
response.write_i32(0);
response.write_i32(preferred_read_replica);
response.write_bytes(&[]).unwrap();
response.into_bytes()
}
fn fetch_v12_response_frame(correlation_id: i32, preferred_read_replica: i32) -> Vec<u8> {
fetch_v12_response_frame_with_session(correlation_id, preferred_read_replica, 0)
}
fn fetch_v12_response_frame_with_session(
correlation_id: i32,
preferred_read_replica: i32,
session_id: i32,
) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_unsigned_varint(0); response.write_i32(0); response.write_i16(0);
response.write_i32(session_id);
response.write_unsigned_varint(2); response.write_compact_string("orders").unwrap();
response.write_unsigned_varint(2); response.write_i32(0);
response.write_i16(0);
response.write_i64(43);
response.write_i64(43);
response.write_i64(42);
response.write_unsigned_varint(1); response.write_i32(preferred_read_replica);
response.write_compact_nullable_bytes(Some(&[])).unwrap();
response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.into_bytes()
}
fn fetch_v12_response_frame_with_record(correlation_id: i32, session_id: i32) -> Vec<u8> {
fetch_v12_response_frame_with_record_at(correlation_id, session_id, 42)
}
fn fetch_v12_response_frame_with_record_at(
correlation_id: i32,
session_id: i32,
offset: i64,
) -> Vec<u8> {
let mut message = Encoder::new();
message.write_i32(0);
message.write_i8(1);
message.write_i8(0);
message.write_i64(123);
message.write_nullable_bytes(Some(b"order-1")).unwrap();
message.write_nullable_bytes(Some(b"created")).unwrap();
let message = message.into_bytes();
let mut records = Encoder::new();
records.write_i64(offset);
records.write_i32(i32::try_from(message.len()).unwrap());
records.write_raw(&message);
let records = records.into_bytes();
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_unsigned_varint(0); response.write_i32(0); response.write_i16(0);
response.write_i32(session_id);
response.write_unsigned_varint(2); response.write_compact_string("orders").unwrap();
response.write_unsigned_varint(2); response.write_i32(0);
response.write_i16(0);
response.write_i64(43);
response.write_i64(43);
response.write_i64(0);
response.write_unsigned_varint(1); response.write_i32(-1);
response
.write_compact_nullable_bytes(Some(&records))
.unwrap();
response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.into_bytes()
}
fn fetch_v12_response_frame_with_partition_error(
correlation_id: i32,
error_code: i16,
) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_unsigned_varint(0); response.write_i32(0); response.write_i16(0);
response.write_i32(0);
response.write_unsigned_varint(2); response.write_compact_string("orders").unwrap();
response.write_unsigned_varint(2); response.write_i32(0);
response.write_i16(error_code);
response.write_i64(43);
response.write_i64(43);
response.write_i64(0);
response.write_unsigned_varint(1); response.write_i32(-1);
response.write_compact_nullable_bytes(Some(&[])).unwrap();
response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.into_bytes()
}
fn fetch_v12_out_of_range_response_frame(correlation_id: i32) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_unsigned_varint(0); response.write_i32(0); response.write_i16(0);
response.write_i32(0);
response.write_unsigned_varint(2); response.write_compact_string("orders").unwrap();
response.write_unsigned_varint(2); response.write_i32(0);
response.write_i16(1); response.write_i64(-1);
response.write_i64(-1);
response.write_i64(0);
response.write_unsigned_varint(1); response.write_i32(-1);
response.write_compact_nullable_bytes(Some(&[])).unwrap();
response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.write_unsigned_varint(0); response.into_bytes()
}
fn list_offsets_response_frame(correlation_id: i32, offset: i64) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_i32(1);
response.write_string("orders").unwrap();
response.write_i32(1);
response.write_i32(0);
response.write_i16(0);
response.write_i64(-1);
response.write_i64(offset);
response.into_bytes()
}
fn offset_for_leader_epoch_response_frame() -> Vec<u8> {
offset_for_leader_epoch_response_frame_with(1, 8, 42)
}
fn offset_for_leader_epoch_response_frame_with(
correlation_id: i32,
leader_epoch: i32,
end_offset: i64,
) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_i32(0);
response.write_i32(1);
response.write_string("orders").unwrap();
response.write_i32(1);
response.write_i16(0);
response.write_i32(0);
response.write_i32(leader_epoch);
response.write_i64(end_offset);
response.into_bytes()
}
fn metadata_response_frame_for(correlation_id: i32, addr: &std::net::SocketAddr) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_i32(1);
response.write_i32(1);
response.write_string(&addr.ip().to_string()).unwrap();
response.write_i32(i32::from(addr.port()));
response.write_nullable_string(None).unwrap();
response.write_i32(1);
response.write_i32(1);
response.write_i16(0);
response.write_string("orders").unwrap();
response.write_i8(0);
response.write_i32(1);
response.write_i16(0);
response.write_i32(0);
response.write_i32(1);
response.write_i32(1);
response.write_i32(1);
response.write_i32(1);
response.write_i32(1);
response.into_bytes()
}
fn metadata_v12_response_frame(
correlation_id: i32,
addr: &std::net::SocketAddr,
leader_epoch: i32,
) -> Vec<u8> {
let mut response = Encoder::new();
response.write_i32(correlation_id);
response.write_unsigned_varint(0); response.write_i32(0); response.write_unsigned_varint(2); response.write_i32(1);
response
.write_compact_string(&addr.ip().to_string())
.unwrap();
response.write_i32(i32::from(addr.port()));
response.write_compact_nullable_string(None).unwrap();
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);
response
.write_compact_nullable_string(Some("orders"))
.unwrap();
response.write_uuid(&[0; 16]);
response.write_i8(0); response.write_unsigned_varint(2); response.write_i16(0);
response.write_i32(0);
response.write_i32(1);
response.write_i32(leader_epoch);
response.write_unsigned_varint(2); response.write_i32(1);
response.write_unsigned_varint(2); response.write_i32(1);
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()
}
}