#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
use std::collections::HashSet;
use std::time::Duration;
use crate::admin::{AdminClient, NewTopic};
use crate::error::ErrorCode;
use crate::producer::Producer;
use crate::protocol::ApiKey;
use super::{Control, FakeBroker};
const SETTLE: Duration = Duration::from_secs(15);
const SHORT_REQUEST_TIMEOUT: Duration = Duration::from_secs(2);
const SHORT_CONNECT_TIMEOUT: Duration = Duration::from_secs(2);
async fn admin_for(broker: &FakeBroker) -> AdminClient {
AdminClient::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("admin client should connect to the fake broker")
}
#[tokio::test]
async fn a_real_admin_client_completes_a_handshake_and_creates_a_topic() {
let broker = FakeBroker::start().await.unwrap();
let admin = admin_for(&broker).await;
let results = admin
.create_topics(vec![NewTopic::new("orders", 3, 1).unwrap()], SETTLE, false)
.await
.expect("CreateTopics should succeed");
assert_eq!(results.len(), 1);
assert_eq!(results[0].error, None, "topic creation reported an error");
broker.with_state(|s| {
let topic = s
.topics
.get("orders")
.expect("broker should hold the topic");
assert_eq!(topic.partitions.len(), 3);
});
assert!(broker.request_count(ApiKey::ApiVersions) >= 1);
assert_eq!(broker.request_count(ApiKey::CreateTopics), 1);
}
#[tokio::test]
async fn a_real_producer_appends_records_at_broker_assigned_offsets() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.linger(Duration::from_millis(5))
.build()
.await
.expect("producer should connect");
for i in 0..3u8 {
let _ = producer
.send("events", None, &[b'v', i])
.await
.expect("send should be acknowledged");
}
assert_eq!(
broker.next_offset("events", 0),
Some(3),
"three records should have been appended"
);
}
#[tokio::test]
async fn not_controller_on_create_topics_makes_the_client_refresh_and_retry() {
let broker = FakeBroker::start().await.unwrap();
let admin = admin_for(&broker).await;
broker.on_once(ApiKey::CreateTopics, |_| {
Control::Error(ErrorCode::NotController)
});
let metadata_before = broker.request_count(ApiKey::Metadata);
let results = admin
.create_topics(vec![NewTopic::new("orders", 1, 1).unwrap()], SETTLE, false)
.await
.expect("the client should retry past NOT_CONTROLLER, not fail");
assert_eq!(results[0].error, None, "the retry should have succeeded");
assert_eq!(
broker.request_count(ApiKey::CreateTopics),
2,
"expected exactly one retry after NOT_CONTROLLER"
);
assert!(
broker.request_count(ApiKey::Metadata) > metadata_before,
"the client must refresh metadata to re-resolve the controller"
);
}
#[tokio::test]
async fn a_permanently_missing_controller_gives_up_instead_of_looping() {
let broker = FakeBroker::start().await.unwrap();
let admin = admin_for(&broker).await;
broker.on(ApiKey::CreateTopics, |_| {
Control::Error(ErrorCode::NotController)
});
let outcome = admin
.create_topics(vec![NewTopic::new("orders", 1, 1).unwrap()], SETTLE, false)
.await;
assert!(
outcome.is_err(),
"a permanent NOT_CONTROLLER must terminate, got {outcome:?}"
);
let attempts = broker.request_count(ApiKey::CreateTopics);
assert!(
(2..=10).contains(&attempts),
"retries should be bounded, saw {attempts} attempts"
);
}
#[tokio::test]
async fn a_producer_follows_a_partition_leader_to_another_broker() {
let broker = FakeBroker::start_cluster(2).await.unwrap();
broker.create_topic("events", 1);
broker.with_state(|s| {
if let Some(p) = s.partition_mut("events", 0) {
p.leader = 0;
}
});
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.metadata_max_age(Duration::from_millis(500))
.linger(Duration::from_millis(5))
.build()
.await
.expect("producer should connect");
let _ = producer
.send("events", None, b"before-the-move")
.await
.expect("the first send should land on the original leader");
broker.clear_requests();
assert!(broker.set_leader("events", 0, 1));
let _ = producer
.send("events", None, b"after-the-move")
.await
.expect("the producer should follow the leader rather than fail");
assert_eq!(
broker.next_offset("events", 0),
Some(2),
"both records should be in the log"
);
assert!(
broker.request_nodes(ApiKey::Produce).contains(&1),
"the producer never reached the new leader; hits were {:?}",
broker.request_nodes(ApiKey::Produce)
);
assert_eq!(
broker.request_count(ApiKey::Metadata),
0,
"the leader was named in the produce response, so no refresh was needed; \
requests were {:?}",
broker.requests()
);
}
#[tokio::test]
async fn a_leader_move_is_followed_without_waiting_for_the_cache_to_age() {
let broker = FakeBroker::start_cluster(2).await.unwrap();
broker.create_topic("events", 1);
broker.with_state(|s| {
if let Some(p) = s.partition_mut("events", 0) {
p.leader = 0;
}
});
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.delivery_timeout(Duration::from_secs(20))
.linger(Duration::from_millis(5))
.build()
.await
.expect("producer should connect");
let _ = producer.send("events", None, b"before").await.unwrap();
broker.clear_requests();
assert!(broker.set_leader("events", 0, 1));
let _ = producer
.send("events", None, b"after")
.await
.expect("a leader move must be followed without waiting out metadata_max_age");
assert!(
broker.request_nodes(ApiKey::Produce).contains(&1),
"the producer never reached the new leader; hits were {:?}",
broker.request_nodes(ApiKey::Produce)
);
assert_eq!(
broker.request_count(ApiKey::Metadata),
0,
"the broker-supplied leader must be enough on its own"
);
assert_eq!(broker.next_offset("events", 0), Some(2));
}
#[tokio::test]
async fn an_error_without_a_leader_hint_still_forces_a_metadata_refresh() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.metadata_max_age(Duration::from_millis(500))
.linger(Duration::from_millis(5))
.build()
.await
.expect("producer should connect");
let _ = producer.send("events", None, b"before").await.unwrap();
broker.clear_requests();
broker.on_once(ApiKey::Produce, |_| {
Control::Error(ErrorCode::NotLeaderForPartition)
});
let _ = producer
.send("events", None, b"after")
.await
.expect("the retry should succeed once the injected error is spent");
assert_eq!(broker.request_count(ApiKey::Produce), 2);
assert!(
broker.request_count(ApiKey::Metadata) >= 1,
"with no leader named, the client must refresh metadata to find one"
);
}
#[tokio::test]
async fn a_group_coordinator_move_is_rediscovered_on_the_new_broker() {
let broker = FakeBroker::start_cluster(2).await.unwrap();
broker.create_topic("events", 1);
broker.set_group_coordinator("analytics", 0);
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("analytics")
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.metadata_max_age(Duration::from_millis(500))
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should succeed");
assert!(
broker.wait_for_requests(ApiKey::JoinGroup, 1, SETTLE).await,
"the consumer should join the group against the original coordinator"
);
assert_eq!(
broker.request_nodes(ApiKey::JoinGroup),
vec![0],
"the first join must go to the original coordinator"
);
broker.clear_requests();
broker.set_group_coordinator("analytics", 1);
let poller = tokio::spawn(async move {
loop {
let _ = tokio::time::timeout(Duration::from_millis(200), consumer.recv()).await;
}
});
let rediscovered = broker
.wait_for_requests(ApiKey::FindCoordinator, 1, SETTLE)
.await;
let reached_new_node = broker
.wait_for_request_on_node(ApiKey::JoinGroup, 1, SETTLE)
.await;
poller.abort();
assert!(
rediscovered,
"the client must re-run FindCoordinator after NOT_COORDINATOR"
);
assert!(
reached_new_node,
"the client never re-joined after the coordinator moved"
);
assert!(
broker.request_nodes(ApiKey::JoinGroup).contains(&1),
"the client kept talking to the old coordinator; join hits were {:?}",
broker.request_nodes(ApiKey::JoinGroup)
);
}
#[tokio::test]
async fn a_response_arriving_after_the_client_timeout_leaves_the_connection_usable() {
let broker = FakeBroker::start().await.unwrap();
let admin = AdminClient::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("admin client should connect");
broker.on_once(ApiKey::CreateTopics, |_| {
Control::Delay(SHORT_REQUEST_TIMEOUT + Duration::from_secs(1))
});
let timed_out = admin
.create_topics(vec![NewTopic::new("slow", 1, 1).unwrap()], SETTLE, false)
.await;
assert!(
timed_out.is_err(),
"the request should have timed out client-side, got {timed_out:?}"
);
tokio::time::sleep(Duration::from_secs(2)).await;
let after = admin
.create_topics(vec![NewTopic::new("after", 1, 1).unwrap()], SETTLE, false)
.await
.expect("the client must still be usable after an orphaned response");
assert_eq!(after[0].error, None);
broker.with_state(|s| {
assert!(
s.topics.contains_key("after"),
"the follow-up request should have been served normally"
);
});
}
#[tokio::test]
async fn a_delayed_response_does_not_fail_requests_on_other_connections() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let slow = AdminClient::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("admin client should connect");
let healthy = admin_for(&broker).await;
broker.on_once(ApiKey::CreateTopics, |_| {
Control::Delay(SHORT_REQUEST_TIMEOUT + Duration::from_secs(1))
});
let (slow_result, healthy_result) = tokio::join!(
slow.create_topics(vec![NewTopic::new("slow", 1, 1).unwrap()], SETTLE, false),
async {
tokio::time::sleep(Duration::from_millis(50)).await;
healthy.list_topics().await
}
);
assert!(slow_result.is_err(), "the delayed request should time out");
assert!(
healthy_result.is_ok(),
"an unrelated connection must be unaffected, got {healthy_result:?}"
);
}
#[tokio::test]
async fn a_dropped_connection_is_re_established() {
let broker = FakeBroker::start().await.unwrap();
let admin = admin_for(&broker).await;
admin
.create_topics(vec![NewTopic::new("first", 1, 1).unwrap()], SETTLE, false)
.await
.expect("the first request should succeed");
broker.on_once(ApiKey::Metadata, |_| Control::Disconnect);
let outcome = admin
.create_topics(vec![NewTopic::new("second", 1, 1).unwrap()], SETTLE, false)
.await;
assert!(
outcome.is_ok(),
"the client should recover from a dropped connection, got {outcome:?}"
);
}
#[tokio::test]
async fn on_times_applies_to_exactly_that_many_requests() {
let broker = FakeBroker::start().await.unwrap();
let admin = admin_for(&broker).await;
broker.on_times(ApiKey::CreateTopics, 2, |_| {
Control::Error(ErrorCode::NotController)
});
admin
.create_topics(vec![NewTopic::new("orders", 1, 1).unwrap()], SETTLE, false)
.await
.expect("two rejections should still be within the retry budget");
assert_eq!(
broker.request_count(ApiKey::CreateTopics),
3,
"two injected failures then one success"
);
}
#[tokio::test]
async fn the_handshake_negotiates_a_flexible_api_versions_version() {
let broker = FakeBroker::start().await.unwrap();
let _admin = admin_for(&broker).await;
let negotiated: Vec<i16> = broker
.requests()
.into_iter()
.filter(|r| r.api_key == ApiKey::ApiVersions)
.map(|r| r.api_version)
.collect();
assert!(
!negotiated.is_empty(),
"the client must send at least one ApiVersions request"
);
let settled = negotiated.last().copied().unwrap_or(-1);
assert!(
settled >= 3,
"ApiVersions must settle on v3+ so KIP-511 client software identity is \
actually on the wire; got {negotiated:?}"
);
let client_ceiling = crate::protocol::versions::API_VERSIONS_MAX;
let broker_ceiling = super::handlers::API_VERSIONS_RANGE.1;
let expected_attempts = if client_ceiling > broker_ceiling {
2
} else {
1
};
assert_eq!(
negotiated.len(),
expected_attempts,
"client ceiling v{client_ceiling}, broker ceiling v{broker_ceiling}; \
got attempts {negotiated:?}"
);
assert_eq!(
settled,
client_ceiling.min(broker_ceiling),
"the handshake must settle on the highest mutually supported version"
);
}
#[tokio::test]
async fn an_unsupported_api_versions_ceiling_falls_back_instead_of_failing() {
let broker = FakeBroker::start().await.unwrap();
broker.on_once(ApiKey::ApiVersions, |_| {
Control::Error(ErrorCode::UnsupportedVersion)
});
let admin = admin_for(&broker).await;
admin
.create_topics(vec![NewTopic::new("orders", 1, 1).unwrap()], SETTLE, false)
.await
.expect("the client should fall back and complete the handshake");
assert!(
broker.request_count(ApiKey::ApiVersions) >= 2,
"a rejected ceiling must be retried at a lower version, not surfaced \
as a connection failure"
);
}
async fn producer_for(broker: &FakeBroker) -> Producer {
Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.linger(Duration::from_millis(5))
.build()
.await
.expect("producer should connect to the fake broker")
}
#[tokio::test]
async fn a_corrupt_record_batch_surfaces_from_poll_instead_of_stalling() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = producer_for(&broker).await;
let _ = producer
.send("events", None, b"payload")
.await
.expect("produce should succeed");
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.assign("events", vec![0])
.await
.expect("assign should succeed");
broker.on(ApiKey::Fetch, |_| Control::CorruptRecords);
let deadline = tokio::time::Instant::now() + SETTLE;
let mut surfaced = None;
while tokio::time::Instant::now() < deadline {
match consumer.poll(Duration::from_millis(200)).await {
Ok(records) => assert!(
records.is_empty(),
"no record may be delivered from a batch that failed its CRC"
),
Err(e) => {
surfaced = Some(e);
break;
}
}
}
let err = surfaced.expect(
"a partition stuck on an undecodable batch must surface an error from poll(), \
not stall silently",
);
let text = err.to_string();
assert!(
text.contains("events-0"),
"the error must name the stuck partition so it is actionable: {text}"
);
assert!(
text.contains("seek") && text.contains("pause"),
"the error must state both remedies: {text}"
);
assert_eq!(
err.protocol_error_kind(),
Some(crate::error::ProtocolErrorKind::CrcMismatch),
"the underlying decode failure kind must survive out to the caller"
);
assert!(
!err.is_retriable(),
"a CRC failure is not retriable: re-fetching returns the same bytes"
);
assert!(
consumer.metrics().batch_decode_errors.get() > 0,
"the corruption must be counted, so it is alertable without log scraping"
);
}
#[tokio::test]
async fn pausing_a_corrupt_partition_lets_the_others_keep_flowing() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 2);
let producer = producer_for(&broker).await;
for partition in 0..2 {
let _ = producer
.send_record(
crate::producer::ProducerRecord::new("events", &b"payload"[..])
.with_partition(partition),
)
.await
.expect("produce should succeed");
}
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.assign("events", vec![0, 1])
.await
.expect("assign should succeed");
broker.on(ApiKey::Fetch, |_| Control::CorruptRecords);
let deadline = tokio::time::Instant::now() + SETTLE;
let mut saw_fault = false;
while tokio::time::Instant::now() < deadline {
if consumer.poll(Duration::from_millis(200)).await.is_err() {
saw_fault = true;
break;
}
}
assert!(saw_fault, "the corrupt fetch should have surfaced an error");
broker.clear_hooks();
consumer.pause("events", &[0]).await;
let deadline = tokio::time::Instant::now() + SETTLE;
let mut delivered = Vec::new();
while tokio::time::Instant::now() < deadline && delivered.is_empty() {
match consumer.poll(Duration::from_millis(200)).await {
Ok(records) => delivered.extend(records),
Err(e) => panic!("the unpaused partition must not be affected: {e}"),
}
}
assert!(
!delivered.is_empty(),
"pausing the stuck partition must let the healthy one keep delivering"
);
assert!(
delivered.iter().all(|r| r.partition == 1),
"only the unpaused partition should deliver"
);
}
#[tokio::test]
async fn a_commit_carries_the_leader_epoch_it_was_read_at() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
assert!(broker.bump_leader_epoch("events", 0));
assert!(broker.bump_leader_epoch("events", 0));
let expected_epoch = broker.with_state(|s| {
s.topics
.get("events")
.and_then(|t| t.partitions.first())
.map(|p| p.leader_epoch)
.expect("partition should exist")
});
assert!(
expected_epoch > 0,
"the test needs a non-zero epoch to be meaningful, got {expected_epoch}"
);
let producer = producer_for(&broker).await;
let _ = producer
.send("events", None, b"payload")
.await
.expect("produce should succeed");
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("readers")
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.enable_auto_commit(false)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should succeed");
let deadline = tokio::time::Instant::now() + SETTLE;
let mut got = 0usize;
while tokio::time::Instant::now() < deadline && got == 0 {
got = consumer
.poll(Duration::from_millis(200))
.await
.expect("poll should succeed")
.len();
}
assert_eq!(got, 1, "the consumer should have read the record");
consumer.commit().await.expect("commit should succeed");
let committed = broker.with_state(|s| {
s.groups
.get("readers")
.and_then(|g| g.offsets.get(&("events".to_string(), 0)))
.cloned()
.expect("the broker should have recorded a commit")
});
assert_eq!(committed.offset, 1, "commit should be next-offset");
assert_eq!(
committed.leader_epoch, expected_epoch,
"the commit must carry the leader epoch the record was read at, not -1; \
without it KIP-320 truncation detection is lost across restarts"
);
}
#[tokio::test]
async fn a_stale_leader_epoch_on_list_offsets_recovers_after_a_refresh() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = producer_for(&broker).await;
let _ = producer
.send("events", None, b"payload")
.await
.expect("produce should succeed");
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.metadata_max_age(Duration::from_secs(300))
.build()
.await
.expect("consumer should connect");
consumer
.assign("events", vec![0])
.await
.expect("assign should succeed");
assert!(broker.bump_leader_epoch("events", 0));
let deadline = tokio::time::Instant::now() + SETTLE;
let mut delivered = Vec::new();
while tokio::time::Instant::now() < deadline && delivered.is_empty() {
match consumer.poll(Duration::from_millis(200)).await {
Ok(records) => delivered.extend(records),
Err(e) => panic!("the client should recover from a fenced epoch, got {e}"),
}
}
assert!(
!delivered.is_empty(),
"a fenced ListOffsets must trigger a metadata refresh and then succeed, \
not leave the partition unable to resolve its start offset"
);
}
#[tokio::test]
async fn a_kip848_consumer_joins_and_receives_a_server_side_assignment() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 3);
let producer = producer_for(&broker).await;
let _ = producer
.send("events", None, b"payload")
.await
.expect("produce should succeed");
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("modern")
.group_protocol(crate::consumer::GroupProtocol::Consumer)
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should succeed");
let deadline = tokio::time::Instant::now() + SETTLE;
let mut delivered = Vec::new();
while tokio::time::Instant::now() < deadline && delivered.is_empty() {
delivered.extend(
consumer
.poll(Duration::from_millis(200))
.await
.expect("poll should succeed"),
);
}
assert!(
!delivered.is_empty(),
"a KIP-848 consumer should receive its assignment and consume"
);
assert!(
broker.request_count(ApiKey::ConsumerGroupHeartbeat) >= 1,
"membership must be driven by ConsumerGroupHeartbeat"
);
assert_eq!(
broker.request_count(ApiKey::JoinGroup),
0,
"KIP-848 must not fall back to the classic JoinGroup protocol"
);
assert_eq!(
broker.request_count(ApiKey::SyncGroup),
0,
"KIP-848 must not fall back to the classic SyncGroup protocol"
);
let assignment = consumer.assignment().await;
assert_eq!(
assignment.get("events").map(|p| p.len()),
Some(3),
"the sole member should own every partition; got {assignment:?}"
);
}
#[tokio::test]
async fn a_fenced_kip848_member_gives_up_its_partitions() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 2);
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("modern")
.group_protocol(crate::consumer::GroupProtocol::Consumer)
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should succeed");
let deadline = tokio::time::Instant::now() + SETTLE;
while tokio::time::Instant::now() < deadline {
let _ = consumer.poll(Duration::from_millis(100)).await;
if !consumer.assignment().await.is_empty() {
break;
}
}
assert_eq!(
consumer.assignment().await.get("events").map(|p| p.len()),
Some(2),
"the consumer must hold an assignment before fencing is meaningful"
);
broker.on(ApiKey::ConsumerGroupHeartbeat, |_| {
Control::Error(ErrorCode::FencedMemberEpoch)
});
let deadline = tokio::time::Instant::now() + SETTLE;
let mut dropped = false;
while tokio::time::Instant::now() < deadline {
let _ = consumer.poll(Duration::from_millis(100)).await;
if consumer.assignment().await.is_empty() {
dropped = true;
break;
}
}
assert!(
dropped,
"a fenced member must drop its assignment; keeping it means consuming \
partitions the coordinator has reassigned to someone else"
);
}
#[tokio::test]
async fn a_fenced_kip848_member_rejoins_once_the_fencing_clears() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 2);
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("modern")
.group_protocol(crate::consumer::GroupProtocol::Consumer)
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should succeed");
let deadline = tokio::time::Instant::now() + SETTLE;
while tokio::time::Instant::now() < deadline {
let _ = consumer.poll(Duration::from_millis(100)).await;
if !consumer.assignment().await.is_empty() {
break;
}
}
let epoch_before = broker.with_state(|s| {
s.groups
.get("modern")
.map(|g| g.group_epoch)
.expect("group should exist")
});
broker.on_once(ApiKey::ConsumerGroupHeartbeat, |_| {
Control::Error(ErrorCode::FencedMemberEpoch)
});
let deadline = tokio::time::Instant::now() + SETTLE;
let mut rejoined = false;
while tokio::time::Instant::now() < deadline {
let _ = consumer.poll(Duration::from_millis(100)).await;
let epoch_now =
broker.with_state(|s| s.groups.get("modern").map(|g| g.group_epoch).unwrap_or(-1));
if epoch_now > epoch_before && !consumer.assignment().await.is_empty() {
rejoined = true;
break;
}
}
assert!(
rejoined,
"a fenced member must rejoin at epoch 0 and be re-assigned; \
group epoch was {epoch_before} before fencing"
);
}
#[tokio::test]
async fn a_settled_kip848_member_does_not_spin_on_heartbeats() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("modern")
.group_protocol(crate::consumer::GroupProtocol::Consumer)
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should succeed");
let deadline = tokio::time::Instant::now() + SETTLE;
while tokio::time::Instant::now() < deadline {
let _ = consumer.poll(Duration::from_millis(50)).await;
if !consumer.assignment().await.is_empty() {
break;
}
}
assert!(
!consumer.assignment().await.is_empty(),
"the consumer must settle before rate can be measured"
);
let settled = broker.request_count(ApiKey::ConsumerGroupHeartbeat);
let until = tokio::time::Instant::now() + Duration::from_secs(1);
while tokio::time::Instant::now() < until {
let _ = consumer.poll(Duration::from_millis(10)).await;
}
let sent = broker.request_count(ApiKey::ConsumerGroupHeartbeat) - settled;
assert!(
sent < 25,
"a settled member sent {sent} heartbeats in one second; it is spinning \
rather than heartbeating on the coordinator's interval"
);
}
#[tokio::test]
async fn the_fake_coordinator_fences_a_stale_member_epoch() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("modern")
.group_protocol(crate::consumer::GroupProtocol::Consumer)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should succeed");
let deadline = tokio::time::Instant::now() + SETTLE;
while tokio::time::Instant::now() < deadline {
let _ = consumer.poll(Duration::from_millis(100)).await;
if !consumer.assignment().await.is_empty() {
break;
}
}
let (epoch, members) = broker.with_state(|s| {
let g = s.groups.get("modern").expect("group should exist");
(g.group_epoch, g.consumer_members.len())
});
assert!(
epoch > 0,
"joining must advance the group epoch, got {epoch}"
);
assert_eq!(
members, 1,
"exactly one KIP-848 member should be registered"
);
}
#[tokio::test]
async fn two_kip848_members_converge_on_a_disjoint_split() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 4);
let build = |name: &'static str| {
let servers = broker.bootstrap_servers();
async move {
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(servers)
.group_id("modern")
.client_id(name)
.group_protocol(crate::consumer::GroupProtocol::Consumer)
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should succeed");
consumer
}
};
let first = build("first").await;
let deadline = tokio::time::Instant::now() + SETTLE;
while tokio::time::Instant::now() < deadline {
let _ = first.poll(Duration::from_millis(50)).await;
if first.assignment().await.get("events").map(|p| p.len()) == Some(4) {
break;
}
}
assert_eq!(
first.assignment().await.get("events").map(|p| p.len()),
Some(4),
"the sole member should own the whole topic before the second joins"
);
let second = build("second").await;
let deadline = tokio::time::Instant::now() + SETTLE;
let mut converged = false;
while tokio::time::Instant::now() < deadline {
let _ = first.poll(Duration::from_millis(50)).await;
let _ = second.poll(Duration::from_millis(50)).await;
let a = first.assignment().await;
let b = second.assignment().await;
let a_parts: HashSet<i32> = a.get("events").into_iter().flatten().copied().collect();
let b_parts: HashSet<i32> = b.get("events").into_iter().flatten().copied().collect();
let overlap: Vec<i32> = a_parts.intersection(&b_parts).copied().collect();
assert!(
overlap.is_empty(),
"partitions {overlap:?} were owned by both members at once; a partition \
must not reach its new owner before the previous owner released it"
);
if a_parts.len() == 2 && b_parts.len() == 2 {
converged = true;
break;
}
}
assert!(
converged,
"two members subscribed to a 4-partition topic should converge on 2 each; \
got {:?} and {:?}",
first.assignment().await,
second.assignment().await
);
let a = first.assignment().await;
let b = second.assignment().await;
let mut all: Vec<i32> = a
.get("events")
.into_iter()
.flatten()
.chain(b.get("events").into_iter().flatten())
.copied()
.collect();
all.sort_unstable();
assert_eq!(
all,
vec![0, 1, 2, 3],
"every partition must be owned by exactly one member"
);
}
#[tokio::test]
async fn a_departing_kip848_member_hands_its_partitions_back() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 4);
let build = |name: &'static str| {
let servers = broker.bootstrap_servers();
async move {
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(servers)
.group_id("modern")
.client_id(name)
.group_protocol(crate::consumer::GroupProtocol::Consumer)
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should succeed");
consumer
}
};
let survivor = build("survivor").await;
let leaver = build("leaver").await;
let deadline = tokio::time::Instant::now() + SETTLE;
while tokio::time::Instant::now() < deadline {
let _ = survivor.poll(Duration::from_millis(50)).await;
let _ = leaver.poll(Duration::from_millis(50)).await;
let a = survivor.assignment().await;
let b = leaver.assignment().await;
if a.get("events").map(|p| p.len()) == Some(2)
&& b.get("events").map(|p| p.len()) == Some(2)
{
break;
}
}
assert_eq!(
survivor.assignment().await.get("events").map(|p| p.len()),
Some(2),
"the group must split before a departure is meaningful"
);
let _ = leaver.close().await;
let deadline = tokio::time::Instant::now() + SETTLE;
let mut reclaimed = false;
while tokio::time::Instant::now() < deadline {
let _ = survivor.poll(Duration::from_millis(50)).await;
if survivor.assignment().await.get("events").map(|p| p.len()) == Some(4) {
reclaimed = true;
break;
}
}
assert!(
reclaimed,
"the survivor should reclaim the whole topic after the other member \
leaves; got {:?}",
survivor.assignment().await
);
}