#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use bytes::{Bytes, BytesMut};
use super::acks::{self, AckBook, AckFailure, AckRange, AckType};
use super::commit::{self, AcknowledgementCommit, CommitWaiter};
use super::completed_fetch;
use super::*;
use crate::consumer::TopicPartition;
use crate::error::ErrorCode;
use crate::protocol::{Record, RecordBatch, ShareAcquiredRecords};
fn tp(partition: PartitionId) -> TopicPartition {
TopicPartition::new("t", partition)
}
fn acquired(first: i64, last: i64, delivery_count: i16) -> ShareAcquiredRecords {
ShareAcquiredRecords {
first_offset: first,
last_offset: last,
delivery_count,
}
}
fn batch(base: i64, count: i32, control: bool) -> Bytes {
let mut batch = RecordBatch::new();
batch.base_offset = base;
batch.attributes.is_control_batch = control;
batch.last_offset_delta = count - 1;
for i in 0..count {
let mut record = Record::new(None, Some(Bytes::from(vec![i as u8])));
record.offset_delta = i;
batch.add_record(record);
}
batch.encode().unwrap()
}
fn concat(parts: &[Bytes]) -> Bytes {
let mut out = BytesMut::new();
for part in parts {
out.extend_from_slice(part);
}
out.freeze()
}
fn offsets(processed: &completed_fetch::Processed) -> Vec<i64> {
processed.fetch.records.iter().map(|r| r.offset).collect()
}
fn process(records: &Bytes, ranges: &[ShareAcquiredRecords]) -> completed_fetch::Processed {
completed_fetch::process(
"t",
[1; 16],
0,
7,
Some(records),
ranges,
RecordBatch::MAX_DECOMPRESSED_SIZE,
)
}
#[test]
fn only_acquired_offsets_of_a_whole_batch_are_delivered() {
let processed = process(&batch(0, 10, false), &[acquired(5, 9, 2)]);
assert_eq!(offsets(&processed), vec![5, 6, 7, 8, 9]);
assert!(
processed
.fetch
.records
.iter()
.all(|r| r.delivery_count == Some(2))
);
assert!(
processed.acks.is_empty(),
"every acquired offset carried a record"
);
assert_eq!(processed.fetch.node, 7);
}
#[test]
fn a_control_batch_is_skipped_and_gap_acknowledged() {
let records = concat(&[batch(0, 1, false), batch(1, 1, true)]);
let processed = process(&records, &[acquired(0, 1, 1)]);
assert_eq!(offsets(&processed), vec![0]);
assert_eq!(
processed.acks,
vec![AckRange {
first: 1,
last: 1,
kind: AckType::Gap
}]
);
}
#[test]
fn acquired_offsets_without_records_are_gap_acknowledged() {
let processed = process(&batch(10, 5, false), &[acquired(10, 19, 1)]);
assert_eq!(offsets(&processed), vec![10, 11, 12, 13, 14]);
assert_eq!(
processed.acks,
vec![AckRange {
first: 15,
last: 19,
kind: AckType::Gap
}]
);
}
#[test]
fn offsets_from_an_undecodable_batch_on_are_released() {
let mut bad = BytesMut::from(batch(3, 3, false).as_ref());
let last = bad.len() - 1;
bad[last] ^= 0xff; let records = concat(&[batch(0, 3, false), bad.freeze()]);
let processed = process(&records, &[acquired(0, 5, 1)]);
assert_eq!(offsets(&processed), vec![0, 1, 2]);
assert_eq!(
processed.acks,
vec![AckRange {
first: 3,
last: 5,
kind: AckType::Release
}]
);
}
#[test]
fn a_huge_acquired_range_costs_one_gap_range() {
let processed = process(&batch(0, 2, false), &[acquired(0, i64::MAX - 1, 1)]);
assert_eq!(offsets(&processed), vec![0, 1]);
assert_eq!(processed.acks.len(), 1);
assert_eq!(processed.acks[0].first, 2);
assert_eq!(processed.acks[0].kind, AckType::Gap);
}
#[test]
fn inverted_ranges_acquire_nothing() {
let processed = process(&batch(0, 3, false), &[acquired(2, 1, 1)]);
assert!(processed.fetch.records.is_empty());
assert!(processed.acks.is_empty());
}
#[test]
fn unsorted_ranges_are_walked_in_offset_order() {
let processed = process(&batch(0, 6, false), &[acquired(4, 5, 3), acquired(0, 1, 1)]);
assert_eq!(offsets(&processed), vec![0, 1, 4, 5]);
let counts: Vec<_> = processed
.fetch
.records
.iter()
.map(|r| r.delivery_count)
.collect();
assert_eq!(counts, vec![Some(1), Some(1), Some(3), Some(3)]);
}
#[test]
fn wire_batches_are_sorted_and_merged() {
let ranges = [
AckRange::one(3, AckType::Accept),
AckRange::one(1, AckType::Accept),
AckRange::one(2, AckType::Accept),
AckRange::one(4, AckType::Release),
];
let batches = acks::to_batches(&ranges);
let wire: Vec<(i64, i64, Vec<i8>)> = batches
.iter()
.map(|b| (b.first_offset, b.last_offset, b.acknowledge_types.clone()))
.collect();
assert_eq!(wire, vec![(1, 3, vec![1]), (4, 4, vec![2])]);
}
#[test]
fn a_later_acknowledgement_replaces_a_pending_renew() {
let mut book = AckBook::default();
book.add(1, tp(0), [1; 16], AckRange::one(5, AckType::Renew));
book.add(1, tp(0), [1; 16], AckRange::one(5, AckType::Accept));
let taken = book.take(1, None);
assert_eq!(taken.len(), 1);
assert_eq!(taken[0].2, vec![AckRange::one(5, AckType::Accept)]);
}
#[test]
fn acknowledgements_stay_with_the_acquiring_node() {
let mut book = AckBook::default();
book.add(1, tp(0), [1; 16], AckRange::one(1, AckType::Accept));
book.add(2, tp(0), [1; 16], AckRange::one(2, AckType::Accept));
let taken = book.take(2, None);
assert_eq!(taken.len(), 1);
assert_eq!(taken[0].2, vec![AckRange::one(2, AckType::Accept)]);
assert!(book.has_pending(1));
assert!(!book.has_pending(2));
}
#[test]
fn errors_are_classified() {
let broker = |code| KrafkaError::broker(code, "x");
assert_eq!(
acks::classify(&KrafkaError::timeout("t")),
AckFailure::Retry
);
for code in [
ErrorCode::InvalidShareSessionEpoch,
ErrorCode::ShareSessionNotFound,
ErrorCode::ShareSessionLimitReached,
ErrorCode::RequestTimedOut,
] {
assert_eq!(acks::classify(&broker(code)), AckFailure::Retry, "{code:?}");
}
for code in [
ErrorCode::NotLeaderForPartition,
ErrorCode::FencedLeaderEpoch,
ErrorCode::UnknownTopicOrPartition,
] {
assert_eq!(
acks::classify(&broker(code)),
AckFailure::NotLeader,
"{code:?}"
);
}
for code in [ErrorCode::InvalidRecordState, ErrorCode::InvalidRequest] {
assert_eq!(
acks::classify(&broker(code)),
AckFailure::Permanent,
"{code:?}"
);
}
}
#[test]
fn a_retried_failure_goes_back_in_front_and_a_permanent_one_resolves() {
let now = tokio::time::Instant::now();
let mut book = AckBook::default();
book.add(1, tp(0), [1; 16], AckRange::one(1, AckType::Accept));
let _ = book.take(1, None);
book.add(1, tp(0), [1; 16], AckRange::one(2, AckType::Accept));
let retried = book.settle(
1,
&tp(0),
Err(KrafkaError::timeout("t")),
Some(now - Duration::from_secs(30)),
);
assert!(retried.is_none(), "a session-class failure is retried");
let taken = book.take(1, None);
assert_eq!(
taken[0].2,
vec![AckRange {
first: 1,
last: 2,
kind: AckType::Accept
}]
);
let resolved = book
.settle(
1,
&tp(0),
Err(KrafkaError::broker(ErrorCode::InvalidRecordState, "x")),
Some(now - Duration::from_secs(30)),
)
.expect("a permanent failure resolves");
assert!(resolved.result.is_err());
assert!(!book.has_pending(1));
assert!(book.nodes().is_empty(), "the book is empty");
}
#[test]
fn a_retry_past_its_cutoff_resolves_as_a_timeout() {
let mut book = AckBook::default();
book.add(1, tp(0), [1; 16], AckRange::one(1, AckType::Accept));
let _ = book.take(1, None);
let resolved = book
.settle(
1,
&tp(0),
Err(KrafkaError::timeout("t")),
Some(tokio::time::Instant::now() + Duration::from_secs(1)),
)
.expect("past the cutoff it resolves");
assert!(matches!(resolved.result, Err(KrafkaError::Timeout { .. })));
}
#[tokio::test]
async fn a_commit_waiter_collects_one_result_per_partition() {
let mut book = AckBook::default();
book.add(1, tp(0), [1; 16], AckRange::one(1, AckType::Accept));
book.add(1, tp(1), [1; 16], AckRange::one(1, AckType::Accept));
let waiter = CommitWaiter::new();
waiter.expect(book.attach(&waiter));
let _ = book.take(1, None);
let _ = book.settle(1, &tp(0), Ok(()), None);
let _ = book.settle(
1,
&tp(1),
Err(KrafkaError::broker(ErrorCode::InvalidRecordState, "x")),
None,
);
assert!(waiter.wait(tokio::time::Instant::now()).await);
let results = waiter.finish(&[]);
assert_eq!(results.len(), 2);
assert!(results[&tp(0)].is_ok());
assert!(matches!(
results[&tp(1)],
Err(KrafkaError::Broker {
code: ErrorCode::InvalidRecordState,
..
})
));
}
#[tokio::test]
async fn a_commit_with_nothing_attached_is_done_and_empty() {
let waiter = CommitWaiter::new();
assert!(waiter.wait(tokio::time::Instant::now()).await);
assert!(waiter.finish(&[]).is_empty());
}
#[test]
fn a_panicking_callback_is_contained() {
let calls = Arc::new(AtomicUsize::new(0));
let seen = Arc::clone(&calls);
let callback: AcknowledgementCommitCallback = Arc::new(move |_c: &AcknowledgementCommit| {
seen.fetch_add(1, Ordering::SeqCst);
panic!("application bug");
});
let mut book = AckBook::default();
book.add(1, tp(0), [1; 16], AckRange::one(1, AckType::Accept));
book.add(1, tp(1), [1; 16], AckRange::one(1, AckType::Accept));
let _ = book.take(1, None);
let resolved: Vec<_> = [tp(0), tp(1)]
.iter()
.filter_map(|p| book.settle(1, p, Ok(()), None))
.collect();
commit::report(Some(&callback), &resolved);
assert_eq!(calls.load(Ordering::SeqCst), 2, "every outcome is reported");
}
#[test]
fn acknowledged_offsets_are_reported_as_ranges() {
let mut book = AckBook::default();
for offset in [1, 2, 3, 7] {
book.add(1, tp(0), [1; 16], AckRange::one(offset, AckType::Accept));
}
let _ = book.take(1, None);
let resolved = book.settle(1, &tp(0), Ok(()), None).unwrap();
let commit = AcknowledgementCommit::from_resolved(&resolved);
assert_eq!(commit.offsets, vec![1..=3, 7..=7]);
assert!(commit.result.is_ok());
}
#[test]
fn defaults() {
let config = super::config::ShareConsumerConfig::default();
assert_eq!(config.acknowledgement_mode, AcknowledgementMode::Implicit);
assert_eq!(config.acquire_mode, AcquireMode::BatchOptimized);
assert_eq!(config.max_poll_records, 500);
assert_eq!(AcquireMode::BatchOptimized.to_i8(), 0);
assert_eq!(AcquireMode::RecordLimit.to_i8(), 1);
assert_eq!(CloseOptions::default().timeout, None);
assert_eq!(
CloseOptions::new().timeout(Duration::from_secs(1)).timeout,
Some(Duration::from_secs(1))
);
}
#[tokio::test]
async fn the_builder_validates() {
let kafka = crate::Kafka::detached();
let missing_group = kafka.share_consumer("").build().await;
assert!(matches!(missing_group, Err(KrafkaError::Config { .. })));
let zero_poll = kafka.share_consumer("g").max_poll_records(0).build().await;
assert!(matches!(zero_poll, Err(KrafkaError::Config { .. })));
}