use std::collections::VecDeque;
use bytes::Bytes;
use tracing::debug;
use super::acks::{AckRange, AckType};
use crate::consumer::ConsumerRecord;
use crate::protocol::{RecordBatch, ShareAcquiredRecords};
use crate::{BrokerId, Offset, PartitionId};
#[derive(Debug)]
pub(crate) struct CompletedFetch {
pub topic: String,
pub topic_id: [u8; 16],
pub partition: PartitionId,
pub node: BrokerId,
pub records: VecDeque<ConsumerRecord>,
}
#[derive(Debug)]
pub(crate) struct Processed {
pub fetch: CompletedFetch,
pub acks: Vec<AckRange>,
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn process(
topic: &str,
topic_id: [u8; 16],
partition: PartitionId,
node: BrokerId,
records: Option<&Bytes>,
acquired: &[ShareAcquiredRecords],
max_decompressed_size: usize,
) -> Processed {
let mut ranges: Vec<&ShareAcquiredRecords> = acquired
.iter()
.filter(|r| r.last_offset >= r.first_offset)
.collect();
ranges.sort_unstable_by_key(|r| r.first_offset);
let mut delivered: VecDeque<ConsumerRecord> = VecDeque::new();
let topic_name: std::sync::Arc<str> = std::sync::Arc::from(topic);
let mut covered: Vec<(Offset, Offset, AckType)> = Vec::new();
let mut release_from: Option<Offset> = None;
let mut next_range = 0;
if let Some(raw) = records {
let mut cursor = raw.clone();
while !cursor.is_empty() {
let base = (cursor.len() >= 8).then(|| {
let mut prefix = [0u8; 8];
prefix.copy_from_slice(&cursor[..8]);
i64::from_be_bytes(prefix)
});
let batch = match RecordBatch::decode_with_limit(&mut cursor, max_decompressed_size) {
Ok(batch) => batch,
Err(error) => {
debug!("undecodable record batch in {topic}-{partition}: {error}");
release_from = Some(base.unwrap_or(Offset::MIN));
break;
}
};
if batch.attributes.is_control_batch {
let last = batch
.base_offset
.saturating_add(i64::from(batch.last_offset_delta));
covered.push((batch.base_offset, last, AckType::Gap));
continue;
}
for record in batch.records {
let offset = batch
.base_offset
.saturating_add(i64::from(record.offset_delta));
if covered.last().is_some_and(|&(_, last, _)| offset <= last) {
continue;
}
while ranges
.get(next_range)
.is_some_and(|r| r.last_offset < offset)
{
next_range += 1;
}
let Some(range) = ranges.get(next_range).filter(|r| r.first_offset <= offset)
else {
continue;
};
covered.push((offset, offset, AckType::Accept));
delivered.push_back(ConsumerRecord {
topic: std::sync::Arc::clone(&topic_name),
partition,
offset,
timestamp: batch.base_timestamp.saturating_add(record.timestamp_delta),
timestamp_type: batch.attributes.timestamp_type,
key: record.key,
value: record.value,
headers: crate::consumer::headers_from_wire(record.headers),
leader_epoch: Some(batch.partition_leader_epoch),
delivery_count: Some(range.delivery_count),
});
}
}
}
covered.sort_unstable_by_key(|c| c.0);
let mut acks: Vec<AckRange> = Vec::new();
let mut push = |first: Offset, last: Offset, kind: AckType| {
if first > last {
return;
}
match acks.last_mut() {
Some(prev) if prev.kind == kind && prev.last.checked_add(1) == Some(first) => {
prev.last = last;
}
_ => acks.push(AckRange { first, last, kind }),
}
};
let mut start = 0;
for range in ranges {
while start < covered.len() && covered[start].1 < range.first_offset {
start += 1;
}
let mut next = range.first_offset;
for &(first, last, kind) in covered[start..]
.iter()
.take_while(|c| c.0 <= range.last_offset)
{
if last < next {
continue;
}
if first > next {
uncovered(next, first - 1, release_from, &mut push);
}
if kind == AckType::Gap {
push(first.max(next), last.min(range.last_offset), AckType::Gap);
}
next = last.saturating_add(1);
}
if next <= range.last_offset {
uncovered(next, range.last_offset, release_from, &mut push);
}
}
Processed {
fetch: CompletedFetch {
topic: topic.to_string(),
topic_id,
partition,
node,
records: delivered,
},
acks,
}
}
fn uncovered(
first: Offset,
last: Offset,
release_from: Option<Offset>,
push: &mut impl FnMut(Offset, Offset, AckType),
) {
match release_from {
Some(from) if from <= first => push(first, last, AckType::Release),
Some(from) if from <= last => {
push(first, from - 1, AckType::Gap);
push(from, last, AckType::Release);
}
_ => push(first, last, AckType::Gap),
}
}