use std::collections::{BTreeMap, VecDeque};
use std::time::Duration;
use tokio::time::Instant;
use ahash::{AHashMap as HashMap, AHashSet as HashSet};
use super::PartitionLag;
use super::config::{AutoOffsetReset, IsolationLevel};
use super::fetch_session::FetchSessionCache;
use super::record::{ConsumerRecord, TopicPartition};
use crate::error::{KrafkaError, Result};
use crate::{BrokerId, Offset, PartitionId};
pub(super) type PartitionKey = (String, PartitionId);
const BACKOFF_INITIAL: Duration = Duration::from_millis(100);
const BACKOFF_MAX: Duration = Duration::from_secs(30);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum OffsetReset {
Earliest,
Latest,
ByDuration(Duration),
}
impl OffsetReset {
pub(super) fn timestamp(self) -> i64 {
match self {
Self::Earliest => -2,
Self::Latest => -1,
Self::ByDuration(duration) => {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default();
i64::try_from(now.saturating_sub(duration).as_millis()).unwrap_or(i64::MAX)
}
}
}
pub(super) fn from_auto(reset: AutoOffsetReset) -> Option<Self> {
match reset {
AutoOffsetReset::Earliest => Some(Self::Earliest),
AutoOffsetReset::Latest => Some(Self::Latest),
AutoOffsetReset::ByDuration(duration) => Some(Self::ByDuration(duration)),
AutoOffsetReset::None => None,
}
}
}
#[derive(Debug)]
pub(super) struct CompletedFetch {
pub(super) records: VecDeque<ConsumerRecord>,
pub(super) next_offset: Offset,
pub(super) next_epoch: Option<i32>,
}
#[derive(Debug, Default)]
pub(super) struct PartitionRecord {
pub(super) version: u64,
pub(super) position: Option<Offset>,
pub(super) epoch: Option<i32>,
pub(super) reset: Option<OffsetReset>,
pub(super) validated: bool,
pub(super) paused: bool,
pub(super) delivering: bool,
pub(super) buffered: Option<CompletedFetch>,
pub(super) high_watermark: Option<Offset>,
pub(super) watermark_updated_at: Option<Instant>,
pub(super) last_stable_offset: Option<Offset>,
pub(super) log_start_offset: Option<Offset>,
pub(super) preferred_replica: Option<(BrokerId, Instant)>,
pub(super) backoff: Option<(Instant, Duration)>,
}
impl PartitionRecord {
pub(super) fn readable_end_offset(&self, isolation_level: IsolationLevel) -> Option<Offset> {
match isolation_level {
IsolationLevel::ReadCommitted => self.last_stable_offset.or(self.high_watermark),
IsolationLevel::ReadUncommitted => self.high_watermark,
}
}
fn backoff_elapsed(&self, now: Instant) -> bool {
self.backoff.is_none_or(|(next, _)| now >= next)
}
pub(super) fn fetch_position(&self) -> Option<Offset> {
match &self.buffered {
Some(fetch) => Some(fetch.next_offset),
None => self.position,
}
}
fn lag(&self, isolation_level: IsolationLevel) -> Option<u64> {
let end = self.readable_end_offset(isolation_level)?;
let position = self.position?;
Some((end - position).max(0) as u64)
}
}
#[derive(Debug, Clone)]
pub(super) struct FetchTarget {
pub(super) key: PartitionKey,
pub(super) version: u64,
pub(super) offset: Offset,
pub(super) last_fetched_epoch: i32,
pub(super) preferred_replica: Option<BrokerId>,
}
#[derive(Debug, Clone)]
pub(super) struct PendingPosition {
pub(super) key: PartitionKey,
pub(super) version: u64,
}
#[derive(Debug)]
struct DeliveredPartition {
key: PartitionKey,
version: u64,
count: usize,
}
#[derive(Debug, Default)]
pub(super) struct Delivery {
pub(super) records: Vec<ConsumerRecord>,
partitions: Vec<DeliveredPartition>,
}
#[derive(Debug)]
pub(super) struct SubscriptionState {
pub(super) subscription: HashSet<String>,
pub(super) standalone_topics: HashSet<String>,
pub(super) standalone_resolved: Option<Instant>,
partitions: BTreeMap<PartitionKey, PartitionRecord>,
next_version: u64,
pub(super) fetch_sessions: FetchSessionCache,
rotation: usize,
pub(super) last_auto_commit: Instant,
}
impl Default for SubscriptionState {
fn default() -> Self {
Self {
subscription: HashSet::new(),
standalone_topics: HashSet::new(),
standalone_resolved: None,
partitions: BTreeMap::new(),
next_version: 0,
fetch_sessions: FetchSessionCache::new(),
rotation: 0,
last_auto_commit: Instant::now(),
}
}
}
impl SubscriptionState {
fn bump(&mut self) -> u64 {
self.next_version += 1;
self.next_version
}
fn not_assigned(topic: &str, partition: PartitionId) -> KrafkaError {
KrafkaError::illegal_state(format!(
"no current assignment for partition {topic}-{partition}"
))
}
pub(super) fn assignment(&self) -> HashMap<String, Vec<PartitionId>> {
let mut out: HashMap<String, Vec<PartitionId>> = HashMap::new();
for (topic, partition) in self.partitions.keys() {
out.entry(topic.clone()).or_default().push(*partition);
}
out
}
pub(super) fn assigned_keys(&self) -> Vec<PartitionKey> {
self.partitions.keys().cloned().collect()
}
pub(super) fn assigned_count(&self) -> usize {
self.partitions.len()
}
pub(super) fn is_assigned(&self, key: &PartitionKey) -> bool {
self.partitions.contains_key(key)
}
pub(super) fn partition(&self, key: &PartitionKey) -> Option<&PartitionRecord> {
self.partitions.get(key)
}
pub(super) fn add_partitions(&mut self, keys: impl IntoIterator<Item = PartitionKey>) {
for key in keys {
if self.partitions.contains_key(&key) {
continue;
}
let version = self.bump();
self.partitions.insert(
key,
PartitionRecord {
version,
..PartitionRecord::default()
},
);
}
}
pub(super) fn remove_partitions<'a>(
&mut self,
keys: impl IntoIterator<Item = &'a PartitionKey>,
) {
for key in keys {
self.partitions.remove(key);
}
}
pub(super) fn clear_assignment(&mut self) {
self.partitions.clear();
}
pub(super) fn seek(&mut self, key: &PartitionKey, offset: Offset) -> Result<()> {
if !self.partitions.contains_key(key) {
return Err(Self::not_assigned(&key.0, key.1));
}
let version = self.bump();
if let Some(record) = self.partitions.get_mut(key) {
record.version = version;
record.position = Some(offset);
record.epoch = None;
record.reset = None;
record.validated = false;
record.buffered = None;
record.delivering = false;
record.backoff = None;
}
Ok(())
}
pub(super) fn request_reset(&mut self, key: &PartitionKey, reset: OffsetReset) -> Result<()> {
if !self.partitions.contains_key(key) {
return Err(Self::not_assigned(&key.0, key.1));
}
let version = self.bump();
if let Some(record) = self.partitions.get_mut(key) {
record.version = version;
record.position = None;
record.epoch = None;
record.reset = Some(reset);
record.validated = false;
record.buffered = None;
record.delivering = false;
record.backoff = None;
}
Ok(())
}
pub(super) fn reset_if_current(
&mut self,
key: &PartitionKey,
version: u64,
reset: OffsetReset,
) {
if self
.partitions
.get(key)
.is_some_and(|r| r.version == version)
{
let _ = self.request_reset(key, reset);
}
}
pub(super) fn partitions_needing_position(&self, now: Instant) -> Vec<PendingPosition> {
self.partitions
.iter()
.filter(|(_, r)| r.position.is_none() && r.reset.is_none() && r.backoff_elapsed(now))
.map(|(key, r)| PendingPosition {
key: key.clone(),
version: r.version,
})
.collect()
}
pub(super) fn partitions_awaiting_reset(
&self,
now: Instant,
) -> Vec<(PendingPosition, OffsetReset)> {
self.partitions
.iter()
.filter(|(_, r)| r.backoff_elapsed(now))
.filter_map(|(key, r)| {
r.reset.map(|reset| {
(
PendingPosition {
key: key.clone(),
version: r.version,
},
reset,
)
})
})
.collect()
}
pub(super) fn set_initial_position(
&mut self,
key: &PartitionKey,
version: u64,
offset: Offset,
epoch: Option<i32>,
) -> bool {
match self.partitions.get_mut(key) {
Some(r) if r.version == version && r.position.is_none() => {
r.position = Some(offset);
r.epoch = epoch;
r.reset = None;
r.validated = false;
r.backoff = None;
true
}
_ => false,
}
}
pub(super) fn set_pending_reset(
&mut self,
key: &PartitionKey,
version: u64,
reset: OffsetReset,
) {
if let Some(r) = self.partitions.get_mut(key)
&& r.version == version
&& r.position.is_none()
{
r.reset = Some(reset);
}
}
pub(super) fn back_off(&mut self, key: &PartitionKey, version: Option<u64>, now: Instant) {
if let Some(r) = self.partitions.get_mut(key)
&& version.is_none_or(|v| v == r.version)
{
let previous = r.backoff.map(|(_, d)| d).unwrap_or(Duration::ZERO);
let next = (previous * 2).clamp(BACKOFF_INITIAL, BACKOFF_MAX);
r.backoff = Some((now + next, next));
}
}
pub(super) fn truncate(
&mut self,
key: &PartitionKey,
version: Option<u64>,
end_offset: Offset,
) -> Option<Offset> {
let current = self.partitions.get(key)?;
if version.is_some_and(|v| v != current.version) {
return None;
}
let old = current.position;
let new_version = self.bump();
let r = self.partitions.get_mut(key)?;
r.version = new_version;
r.position = Some(end_offset);
r.epoch = None;
r.reset = None;
r.validated = true;
r.buffered = None;
r.delivering = false;
old
}
pub(super) fn mark_validated(&mut self, key: &PartitionKey, version: u64) {
if let Some(r) = self.partitions.get_mut(key)
&& r.version == version
{
r.validated = true;
}
}
pub(super) fn mark_unvalidated(&mut self, key: &PartitionKey) {
if let Some(r) = self.partitions.get_mut(key) {
r.validated = false;
}
}
pub(super) fn unvalidated(
&self,
now: Instant,
) -> Vec<(PartitionKey, u64, Offset, Option<i32>)> {
self.partitions
.iter()
.filter(|(_, r)| !r.validated && !r.paused && r.backoff_elapsed(now))
.filter_map(|(key, r)| r.position.map(|p| (key.clone(), r.version, p, r.epoch)))
.collect()
}
pub(super) fn set_paused(
&mut self,
topic: &str,
partitions: &[PartitionId],
paused: bool,
) -> usize {
let mut applied = 0;
for &partition in partitions {
if let Some(r) = self.partitions.get_mut(&(topic.to_string(), partition)) {
r.paused = paused;
applied += 1;
}
}
applied
}
pub(super) fn paused(&self) -> HashSet<PartitionKey> {
self.partitions
.iter()
.filter(|(_, r)| r.paused)
.map(|(key, _)| key.clone())
.collect()
}
pub(super) fn buffered_count(&self) -> usize {
self.partitions
.values()
.filter_map(|r| r.buffered.as_ref())
.map(|f| f.records.len())
.sum()
}
pub(super) fn has_deliverable(&self) -> bool {
self.partitions.values().any(|r| {
!r.paused && !r.delivering && r.buffered.as_ref().is_some_and(|f| !f.records.is_empty())
})
}
pub(super) fn fetch_targets(&mut self, now: Instant) -> Vec<FetchTarget> {
let mut targets: Vec<FetchTarget> = Vec::new();
for (key, r) in self.partitions.iter_mut() {
if let Some((_, expiry)) = r.preferred_replica
&& now >= expiry
{
r.preferred_replica = None;
}
let Some(offset) = r.position else { continue };
if r.paused
|| r.delivering
|| r.buffered.is_some()
|| r.reset.is_some()
|| !r.backoff_elapsed(now)
{
continue;
}
targets.push(FetchTarget {
key: key.clone(),
version: r.version,
offset,
last_fetched_epoch: r.epoch.unwrap_or(-1),
preferred_replica: r.preferred_replica.map(|(id, _)| id),
});
}
if !targets.is_empty() {
let turn = self.rotation % targets.len();
self.rotation = self.rotation.wrapping_add(1);
targets.rotate_left(turn);
}
targets
}
pub(super) fn install_fetch(&mut self, target: &FetchTarget, fetch: CompletedFetch) -> bool {
let Some(r) = self.partitions.get_mut(&target.key) else {
return false;
};
if r.version != target.version
|| r.position != Some(target.offset)
|| r.buffered.is_some()
|| r.delivering
{
return false;
}
r.backoff = None;
if fetch.records.is_empty() {
if fetch.next_offset > target.offset {
r.position = Some(fetch.next_offset);
if fetch.next_epoch.is_some() {
r.epoch = fetch.next_epoch;
}
}
} else {
r.buffered = Some(fetch);
}
true
}
pub(super) fn clear_backoff(&mut self, key: &PartitionKey, version: u64) {
if let Some(r) = self.partitions.get_mut(key)
&& r.version == version
{
r.backoff = None;
}
}
pub(super) fn update_watermarks(
&mut self,
key: &PartitionKey,
high_watermark: Option<Offset>,
last_stable_offset: Option<Offset>,
log_start_offset: Option<Offset>,
now: Instant,
) {
if let Some(r) = self.partitions.get_mut(key) {
if let Some(hw) = high_watermark {
r.high_watermark = Some(hw);
r.watermark_updated_at = Some(now);
}
if last_stable_offset.is_some() {
r.last_stable_offset = last_stable_offset;
}
if log_start_offset.is_some() {
r.log_start_offset = log_start_offset;
}
}
}
pub(super) fn set_preferred_replica(
&mut self,
key: &PartitionKey,
replica: Option<(BrokerId, Instant)>,
) {
if let Some(r) = self.partitions.get_mut(key) {
r.preferred_replica = replica;
}
}
pub(super) fn take_delivery(&mut self, max: usize) -> Option<Delivery> {
if max == 0 {
return None;
}
let mut delivery = Delivery::default();
let ready: Vec<PartitionKey> = self
.partitions
.iter()
.filter(|(_, r)| {
!r.paused
&& !r.delivering
&& r.buffered.as_ref().is_some_and(|f| !f.records.is_empty())
})
.map(|(key, _)| key.clone())
.collect();
if ready.is_empty() {
return None;
}
let start = self.rotation % ready.len();
self.rotation = self.rotation.wrapping_add(1);
for key in ready.iter().cycle().skip(start).take(ready.len()) {
let remaining = max - delivery.records.len();
if remaining == 0 {
break;
}
let Some(r) = self.partitions.get_mut(key) else {
continue;
};
let Some(fetch) = r.buffered.as_mut() else {
continue;
};
let take = remaining.min(fetch.records.len());
delivery.records.extend(fetch.records.drain(..take));
r.delivering = true;
delivery.partitions.push(DeliveredPartition {
key: key.clone(),
version: r.version,
count: take,
});
}
Some(delivery)
}
pub(super) fn complete_delivery(
&mut self,
delivery: Delivery,
handed_out: usize,
) -> Vec<ConsumerRecord> {
let Delivery {
records,
partitions,
} = delivery;
let mut records = records.into_iter();
let mut out = Vec::with_capacity(handed_out.min(records.len()));
let mut remaining_out = handed_out;
for part in partitions {
let chunk: Vec<ConsumerRecord> = records.by_ref().take(part.count).collect();
let deliver = remaining_out.min(chunk.len());
remaining_out -= deliver;
let Some(r) = self.partitions.get_mut(&part.key) else {
continue;
};
if r.version != part.version {
continue;
}
r.delivering = false;
let mut chunk = chunk.into_iter();
let delivered: Vec<ConsumerRecord> = chunk.by_ref().take(deliver).collect();
let put_back: Vec<ConsumerRecord> = chunk.collect();
if let Some(last) = delivered.last() {
r.position = Some(last.offset.saturating_add(1));
if last.leader_epoch.is_some() {
r.epoch = last.leader_epoch;
}
}
if let Some(fetch) = r.buffered.as_mut() {
for record in put_back.into_iter().rev() {
fetch.records.push_front(record);
}
if fetch.records.is_empty() {
let next = fetch.next_offset;
let next_epoch = fetch.next_epoch;
r.buffered = None;
if r.position.is_none_or(|p| next > p) {
r.position = Some(next);
if next_epoch.is_some() {
r.epoch = next_epoch;
}
}
}
}
out.extend(delivered);
}
out
}
pub(super) fn abort_delivery(&mut self, delivery: Delivery) {
let _ = self.complete_delivery(delivery, 0);
}
pub(super) fn committable(&self) -> Vec<(PartitionKey, Offset, Option<i32>)> {
self.partitions
.iter()
.filter_map(|(key, r)| r.position.map(|p| (key.clone(), p, r.epoch)))
.collect()
}
pub(super) fn aggregate_lag(&self, isolation_level: IsolationLevel) -> (u64, u64) {
let mut total: u64 = 0;
let mut max: u64 = 0;
for r in self.partitions.values() {
if let Some(lag) = r.lag(isolation_level) {
total = total.saturating_add(lag);
max = max.max(lag);
}
}
(total, max)
}
pub(super) fn lag_report(
&self,
isolation_level: IsolationLevel,
now: Instant,
threshold: Duration,
) -> HashMap<TopicPartition, PartitionLag> {
self.partitions
.iter()
.map(|((topic, partition), r)| {
let stale = r
.watermark_updated_at
.is_none_or(|t| now.saturating_duration_since(t) > threshold);
(
TopicPartition::new(topic.clone(), *partition),
PartitionLag {
position: r.position,
log_start_offset: r.log_start_offset,
high_watermark: r.high_watermark,
last_stable_offset: r.last_stable_offset,
lag: r.lag(isolation_level),
stale,
},
)
})
.collect()
}
}
pub(super) struct DeliveryGuard<'a> {
state: &'a parking_lot::Mutex<SubscriptionState>,
delivery: Option<Delivery>,
}
impl<'a> DeliveryGuard<'a> {
pub(super) fn new(
state: &'a parking_lot::Mutex<SubscriptionState>,
delivery: Delivery,
) -> Self {
Self {
state,
delivery: Some(delivery),
}
}
pub(super) fn records_mut(&mut self) -> &mut [ConsumerRecord] {
match self.delivery.as_mut() {
Some(delivery) => &mut delivery.records,
None => &mut [],
}
}
pub(super) fn finish(mut self, handed_out: usize) -> Vec<ConsumerRecord> {
match self.delivery.take() {
Some(delivery) => self.state.lock().complete_delivery(delivery, handed_out),
None => Vec::new(),
}
}
}
impl Drop for DeliveryGuard<'_> {
fn drop(&mut self) {
if let Some(delivery) = self.delivery.take() {
self.state.lock().abort_delivery(delivery);
}
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
fn key(topic: &str, partition: PartitionId) -> PartitionKey {
(topic.to_string(), partition)
}
fn record(topic: &str, partition: PartitionId, offset: Offset) -> ConsumerRecord {
let mut record = ConsumerRecord::new(topic, partition, offset, None, None);
record.leader_epoch = Some(3);
record
}
fn fetched(
topic: &str,
partition: PartitionId,
offsets: std::ops::Range<Offset>,
) -> CompletedFetch {
CompletedFetch {
next_offset: offsets.end,
records: offsets.map(|o| record(topic, partition, o)).collect(),
next_epoch: Some(3),
}
}
fn buffered_state() -> SubscriptionState {
let mut state = SubscriptionState::default();
state.add_partitions([key("t", 0)]);
state.seek(&key("t", 0), 0).unwrap();
let targets = state.fetch_targets(Instant::now());
assert_eq!(targets.len(), 1);
assert!(state.install_fetch(&targets[0], fetched("t", 0, 0..10)));
state
}
#[test]
fn fetching_does_not_move_the_position() {
let state = buffered_state();
let r = state.partition(&key("t", 0)).unwrap();
assert_eq!(r.position, Some(0));
assert_eq!(r.fetch_position(), Some(10));
}
#[test]
fn delivery_advances_the_position_only_when_finished() {
let mut state = buffered_state();
let delivery = state.take_delivery(4).unwrap();
assert_eq!(delivery.records.len(), 4);
assert_eq!(state.partition(&key("t", 0)).unwrap().position, Some(0));
let out = state.complete_delivery(delivery, 4);
assert_eq!(
out.iter().map(|r| r.offset).collect::<Vec<_>>(),
vec![0, 1, 2, 3]
);
let r = state.partition(&key("t", 0)).unwrap();
assert_eq!(r.position, Some(4));
assert_eq!(r.epoch, Some(3));
}
#[test]
fn an_aborted_delivery_puts_every_record_back() {
let mut state = buffered_state();
let delivery = state.take_delivery(4).unwrap();
state.abort_delivery(delivery);
let r = state.partition(&key("t", 0)).unwrap();
assert_eq!(r.position, Some(0));
let again = state.take_delivery(100).unwrap();
assert_eq!(
again.records.iter().map(|r| r.offset).collect::<Vec<_>>(),
(0..10).collect::<Vec<_>>()
);
}
#[test]
fn a_dropped_guard_puts_the_records_back() {
let state = parking_lot::Mutex::new(buffered_state());
let delivery = state.lock().take_delivery(5).unwrap();
drop(DeliveryGuard::new(&state, delivery));
let mut state = state.into_inner();
assert_eq!(state.partition(&key("t", 0)).unwrap().position, Some(0));
assert_eq!(state.take_delivery(100).unwrap().records.len(), 10);
}
#[test]
fn a_partial_hand_out_keeps_the_rest_in_order() {
let mut state = buffered_state();
let delivery = state.take_delivery(6).unwrap();
let out = state.complete_delivery(delivery, 2);
assert_eq!(out.len(), 2);
assert_eq!(state.partition(&key("t", 0)).unwrap().position, Some(2));
let next = state.take_delivery(100).unwrap();
assert_eq!(
next.records.iter().map(|r| r.offset).collect::<Vec<_>>(),
(2..10).collect::<Vec<_>>()
);
}
#[test]
fn draining_the_buffer_moves_to_the_next_fetch_offset() {
let mut state = SubscriptionState::default();
state.add_partitions([key("t", 0)]);
state.seek(&key("t", 0), 0).unwrap();
let target = state.fetch_targets(Instant::now()).remove(0);
let mut fetch = fetched("t", 0, 0..3);
fetch.next_offset = 7;
assert!(state.install_fetch(&target, fetch));
let delivery = state.take_delivery(10).unwrap();
let _ = state.complete_delivery(delivery, 3);
let r = state.partition(&key("t", 0)).unwrap();
assert_eq!(r.position, Some(7));
assert!(r.buffered.is_none());
}
#[test]
fn a_partition_with_buffered_records_is_not_fetched() {
let mut state = buffered_state();
assert!(state.fetch_targets(Instant::now()).is_empty());
let delivery = state.take_delivery(100).unwrap();
assert!(
state.fetch_targets(Instant::now()).is_empty(),
"nor while delivering"
);
let _ = state.complete_delivery(delivery, 10);
let targets = state.fetch_targets(Instant::now());
assert_eq!(targets.len(), 1);
assert_eq!(targets[0].offset, 10);
}
#[test]
fn a_seek_drops_the_buffer_and_a_stale_delivery() {
let mut state = buffered_state();
let delivery = state.take_delivery(3).unwrap();
state.seek(&key("t", 0), 100).unwrap();
let out = state.complete_delivery(delivery, 3);
assert!(
out.is_empty(),
"records from before the seek are not handed out"
);
let r = state.partition(&key("t", 0)).unwrap();
assert_eq!(r.position, Some(100));
assert!(r.buffered.is_none());
}
#[test]
fn a_stale_fetch_is_not_installed() {
let mut state = SubscriptionState::default();
state.add_partitions([key("t", 0)]);
state.seek(&key("t", 0), 0).unwrap();
let target = state.fetch_targets(Instant::now()).remove(0);
state.seek(&key("t", 0), 50).unwrap();
assert!(!state.install_fetch(&target, fetched("t", 0, 0..10)));
assert_eq!(state.partition(&key("t", 0)).unwrap().position, Some(50));
}
#[test]
fn a_fetch_for_a_revoked_and_reassigned_partition_is_not_installed() {
let mut state = SubscriptionState::default();
state.add_partitions([key("t", 0)]);
state.seek(&key("t", 0), 0).unwrap();
let target = state.fetch_targets(Instant::now()).remove(0);
state.remove_partitions([&key("t", 0)]);
state.add_partitions([key("t", 0)]);
assert!(!state.install_fetch(&target, fetched("t", 0, 0..10)));
}
#[test]
fn seek_rejects_an_unassigned_partition_and_stores_nothing() {
let mut state = SubscriptionState::default();
assert!(state.seek(&key("t", 0), 7).is_err());
assert!(
state
.request_reset(&key("t", 0), OffsetReset::Earliest)
.is_err()
);
state.add_partitions([key("t", 0)]);
assert_eq!(state.partition(&key("t", 0)).unwrap().position, None);
}
#[test]
fn a_new_partition_waits_for_its_position() {
let mut state = SubscriptionState::default();
state.add_partitions([key("t", 0)]);
let pending = state.partitions_needing_position(Instant::now());
assert_eq!(pending.len(), 1);
assert!(state.fetch_targets(Instant::now()).is_empty());
assert!(state.set_initial_position(&key("t", 0), pending[0].version, 42, Some(5)));
let r = state.partition(&key("t", 0)).unwrap();
assert_eq!((r.position, r.epoch), (Some(42), Some(5)));
}
#[test]
fn a_reset_lookup_for_a_moved_partition_is_discarded() {
let mut state = SubscriptionState::default();
state.add_partitions([key("t", 0)]);
state
.request_reset(&key("t", 0), OffsetReset::Earliest)
.unwrap();
let (pending, reset) = state.partitions_awaiting_reset(Instant::now()).remove(0);
assert_eq!(reset, OffsetReset::Earliest);
state.seek(&key("t", 0), 9).unwrap();
assert!(!state.set_initial_position(&pending.key, pending.version, 0, None));
assert_eq!(state.partition(&key("t", 0)).unwrap().position, Some(9));
}
#[test]
fn a_paused_partition_keeps_its_buffer_and_position() {
let mut state = buffered_state();
state.set_paused("t", &[0], true);
assert!(state.take_delivery(10).is_none());
assert_eq!(state.buffered_count(), 10);
state.set_paused("t", &[0], false);
assert_eq!(state.take_delivery(10).unwrap().records.len(), 10);
}
#[test]
fn backoff_doubles_and_is_capped() {
let mut state = SubscriptionState::default();
state.add_partitions([key("t", 0)]);
let now = Instant::now();
let mut waits = Vec::new();
for _ in 0..12 {
state.back_off(&key("t", 0), None, now);
waits.push(state.partition(&key("t", 0)).unwrap().backoff.unwrap().1);
}
assert_eq!(waits[0], Duration::from_millis(100));
assert_eq!(waits[1], Duration::from_millis(200));
assert_eq!(*waits.last().unwrap(), Duration::from_secs(30));
assert!(state.partitions_needing_position(now).is_empty());
}
#[test]
fn lag_is_measured_from_the_delivered_position() {
let mut state = buffered_state();
let now = Instant::now();
state.update_watermarks(&key("t", 0), Some(10), None, None, now);
let lag = |state: &SubscriptionState| {
state.lag_report(
IsolationLevel::ReadUncommitted,
now,
Duration::from_secs(60),
)[&TopicPartition::new("t", 0)]
.clone()
};
assert_eq!(lag(&state).lag, Some(10));
assert_eq!(lag(&state).high_watermark, Some(10));
assert!(!lag(&state).stale);
let delivery = state.take_delivery(4).unwrap();
let _ = state.complete_delivery(delivery, 4);
assert_eq!(lag(&state).lag, Some(6));
assert_eq!(lag(&state).position, Some(4));
}
#[test]
fn readable_end_offset_follows_the_isolation_level() {
let record = PartitionRecord {
high_watermark: Some(100),
last_stable_offset: Some(80),
..PartitionRecord::default()
};
assert_eq!(
record.readable_end_offset(IsolationLevel::ReadCommitted),
Some(80)
);
assert_eq!(
record.readable_end_offset(IsolationLevel::ReadUncommitted),
Some(100)
);
let no_lso = PartitionRecord {
high_watermark: Some(100),
..PartitionRecord::default()
};
assert_eq!(
no_lso.readable_end_offset(IsolationLevel::ReadCommitted),
Some(100)
);
}
#[test]
fn delivery_rotates_across_partitions() {
let mut state = SubscriptionState::default();
state.add_partitions([key("t", 0), key("t", 1)]);
for p in 0..2 {
state.seek(&key("t", p), 0).unwrap();
}
for target in state.fetch_targets(Instant::now()) {
let p = target.key.1;
assert!(state.install_fetch(&target, fetched("t", p, 0..5)));
}
let delivery = state.take_delivery(5).unwrap();
let first: HashSet<PartitionId> = delivery.records.iter().map(|r| r.partition).collect();
assert_eq!(first.len(), 1, "the cap is filled from one partition first");
let _ = state.complete_delivery(delivery, 5);
let next = state.take_delivery(5).unwrap();
assert!(next.records.iter().all(|r| !first.contains(&r.partition)));
}
#[test]
fn truncation_rewinds_and_drops_the_buffer() {
let mut state = buffered_state();
let delivery = state.take_delivery(10).unwrap();
let _ = state.complete_delivery(delivery, 10);
assert_eq!(state.truncate(&key("t", 0), None, 6), Some(10));
let r = state.partition(&key("t", 0)).unwrap();
assert_eq!(r.position, Some(6));
assert!(r.validated);
assert!(r.epoch.is_none());
}
}