#![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, Some(&[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 produce_order_follows_enqueue_order_not_await_order() {
const RECORDS: usize = 64;
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)
.build()
.await
.expect("producer should connect");
let mut handles = Vec::with_capacity(RECORDS);
for i in 0..RECORDS {
let handle = producer
.enqueue(
crate::producer::ProducerRecord::new("events", vec![i as u8]).with_partition(0),
)
.await
.expect("enqueue should succeed");
assert_eq!(
handle.partition(),
0,
"the partition is known at enqueue time"
);
handles.push((i, handle));
}
let mut offsets = vec![0i64; RECORDS];
for (i, handle) in handles.into_iter().rev() {
offsets[i] = handle
.await
.expect("every record must be acknowledged")
.offset;
}
producer.close().await;
assert_eq!(
broker.next_offset("events", 0),
Some(RECORDS as i64),
"every record must be appended exactly once"
);
for window in offsets.windows(2) {
assert!(
window[0] < window[1],
"records must be stored in enqueue order, got offsets {offsets:?}"
);
}
}
#[tokio::test]
async fn the_default_producer_batches_concurrent_sends_to_one_partition() {
const RECORDS: usize = 200;
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = std::sync::Arc::new(
Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("producer should connect"),
);
let mut tasks = tokio::task::JoinSet::new();
for i in 0..RECORDS {
let producer = producer.clone();
tasks.spawn(async move {
producer
.send_record(
crate::producer::ProducerRecord::new("events", format!("v{i}").into_bytes())
.with_partition(0),
)
.await
});
}
let mut acknowledged = 0usize;
while let Some(joined) = tasks.join_next().await {
let _ = joined
.expect("send task should not panic")
.expect("every send must be acknowledged");
acknowledged += 1;
}
producer.close().await;
assert_eq!(acknowledged, RECORDS);
assert_eq!(
broker.next_offset("events", 0),
Some(RECORDS as i64),
"every record must be appended exactly once"
);
let produce_requests = broker.request_count(ApiKey::Produce);
assert!(
produce_requests <= RECORDS / 10,
"the default producer must coalesce concurrent sends: {produce_requests} Produce \
requests for {RECORDS} records is barely batching (the unbatched path sent one \
request per record; the accumulator sends a handful)"
);
}
#[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 the_controller_retry_budget_is_the_configured_one() {
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)
.retries(3)
.retry_backoff(Duration::from_millis(1))
.build()
.await
.expect("admin client should connect");
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"
);
assert_eq!(
broker.request_count(ApiKey::CreateTopics),
4,
"retries(3) must mean three retries on top of the first attempt"
);
let message = outcome.expect_err("checked above").to_string();
assert!(
message.contains("retries"),
"the error must name the setting to raise, got: {message}"
);
}
#[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, Some(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, Some(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, Some(b"before"))
.await
.unwrap();
broker.clear_requests();
assert!(broker.set_leader("events", 0, 1));
let _ = producer
.send("events", None, Some(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, Some(b"before"))
.await
.unwrap();
broker.clear_requests();
broker.on_once(ApiKey::Produce, |_| {
Control::Error(ErrorCode::NotLeaderForPartition)
});
let _ = producer
.send("events", None, Some(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, Some(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, Some(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, Some(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, Some(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
);
}
#[tokio::test]
async fn transport_max_response_size_reaches_the_frame_decoder() {
let broker = FakeBroker::start().await.unwrap();
for topic in ["alpha", "bravo", "charlie", "delta"] {
broker.create_topic(topic, 32);
}
let transport = crate::network::TransportConfig::builder()
.max_response_size(1024)
.build()
.expect("valid transport config");
let result = crate::admin::AdminClient::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.transport(transport)
.build()
.await;
let err = result
.err()
.expect("a 1 KiB response ceiling must reject a 128-partition metadata response");
let message = err.to_string();
assert!(
message.contains("exceeds maximum")
|| message.contains("connection closed")
|| message.contains("Connection reset"),
"expected a frame-size rejection, got: {message}"
);
}
#[tokio::test]
async fn transport_default_response_size_still_connects() {
let broker = FakeBroker::start().await.unwrap();
for topic in ["alpha", "bravo", "charlie", "delta"] {
broker.create_topic(topic, 32);
}
let admin = crate::admin::AdminClient::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.transport(crate::network::TransportConfig::default())
.build()
.await
.expect("the default ceiling must not reject a normal metadata response");
admin.close().await;
}
#[tokio::test]
async fn broker_throttle_is_honoured_and_counted() {
let broker = FakeBroker::start().await.unwrap();
let client = crate::client::KrafkaClient::builder(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("client should connect");
let conn = client
.pool()
.get_connection(&broker.bootstrap_servers())
.await
.expect("connection to the fake broker");
assert!(
conn.throttle_remaining().is_none(),
"a connection starts un-throttled"
);
conn.notify_throttle(60);
let remaining = conn
.throttle_remaining()
.expect("the reported throttle must be pending");
assert!(
remaining <= Duration::from_millis(60) && remaining > Duration::from_millis(20),
"the pending delay must reflect what the broker asked for, got {remaining:?}"
);
conn.notify_throttle(5);
assert!(
conn.throttle_remaining().expect("still pending") > Duration::from_millis(20),
"a later, smaller throttle must not cut a longer window short"
);
let metrics = client.pool().metrics();
assert_eq!(metrics.snapshot().throttle_delays, 0, "nothing waited yet");
let waited = conn
.await_throttle()
.await
.expect("there was a delay to wait");
assert!(waited > Duration::ZERO);
assert!(
conn.throttle_remaining().is_none(),
"the window is spent once it has been waited out"
);
let snapshot = metrics.snapshot();
assert_eq!(snapshot.throttle_delays, 1);
assert!(
snapshot.throttle_delay_ms > 0,
"a counted delay with zero duration is not a measurement"
);
assert!(conn.await_throttle().await.is_none());
assert_eq!(metrics.snapshot().throttle_delays, 1);
}
#[tokio::test]
async fn a_throttle_the_producer_waits_out_is_counted() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let client = crate::client::KrafkaClient::builder(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("client should connect");
let pool = client.pool().clone();
let producer = crate::producer::Producer::builder()
.with_client(&client)
.build()
.await
.expect("producer should connect");
let _ = producer
.send("events", None, Some(b"warm-up"))
.await
.expect("send should be acknowledged");
let metrics = producer.connection_metrics();
assert_eq!(
metrics.snapshot().throttle_delays,
0,
"nothing has been throttled yet"
);
let conn = pool
.get_connection(&broker.bootstrap_servers())
.await
.expect("the connection is already open");
conn.notify_throttle(40);
let _ = producer
.send("events", None, Some(b"throttled"))
.await
.expect("a throttled send still succeeds, just later");
producer.close().await;
let snapshot = metrics.snapshot();
assert_eq!(
snapshot.throttle_delays, 1,
"the producer waits out the throttle before dispatching, and that wait \
has to be counted — it used to sleep on `throttle_remaining()` directly, \
which consumed the window before the request path could record it, so \
this metric read zero on the path most likely to be throttled"
);
assert!(
snapshot.throttle_delay_ms > 0,
"a counted delay with zero duration is not a measurement"
);
}
#[tokio::test]
async fn transport_max_connections_bounds_the_pool() {
let broker = FakeBroker::start_cluster(2).await.unwrap();
broker.create_topic("events", 2);
broker.set_leader("events", 0, 0);
broker.set_leader("events", 1, 1);
let transport = crate::network::TransportConfig::builder()
.max_connections(Some(1))
.build()
.expect("valid transport config");
let mut refusal: Option<String> = None;
match crate::producer::Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.linger(Duration::from_millis(5))
.transport(transport)
.build()
.await
{
Err(e) => refusal = Some(e.to_string()),
Ok(producer) => {
for partition in 0..2 {
let record = crate::producer::ProducerRecord::new("events", b"payload".to_vec())
.with_partition(partition);
if let Err(e) = producer.send_record(record).await {
refusal.get_or_insert_with(|| e.to_string());
}
}
producer.close().await;
}
}
let message = refusal.expect(
"max_connections(1) must refuse a second broker connection somewhere; \
if nothing was refused, the cap never reached ConnectionPool",
);
assert!(
message.contains("connection pool limit reached"),
"the refusal must come from the pool cap, not an unrelated failure: {message}"
);
}
#[tokio::test]
async fn transport_sufficient_max_connections_connects() {
let broker = FakeBroker::start_cluster(2).await.unwrap();
broker.create_topic("events", 2);
broker.set_leader("events", 0, 0);
broker.set_leader("events", 1, 1);
let transport = crate::network::TransportConfig::builder()
.max_connections(Some(8))
.build()
.expect("valid transport config");
let producer = crate::producer::Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.linger(Duration::from_millis(5))
.transport(transport)
.build()
.await
.expect("a cap of 8 must accommodate a two-broker cluster");
for partition in 0..2 {
let record = crate::producer::ProducerRecord::new("events", b"payload".to_vec())
.with_partition(partition);
let _metadata = producer
.send_record(record)
.await
.expect("both partitions should be reachable under a sufficient cap");
}
producer.close().await;
}
#[tokio::test]
async fn describe_streams_groups_decodes_topology_and_members() {
use crate::testing::state::{StreamsGroupState, StreamsMemberState};
let broker = FakeBroker::start().await.unwrap();
broker.with_state(|s| {
s.streams_groups.insert(
"wordcount".to_string(),
StreamsGroupState {
group_state: "Stable".to_string(),
group_epoch: 7,
assignment_epoch: 7,
topology_epoch: Some(3),
subtopologies: Some(vec!["0".to_string(), "1".to_string()]),
members: vec![
StreamsMemberState {
member_id: "m-1".to_string(),
member_epoch: 7,
topology_epoch: 3,
process_id: "proc-a".to_string(),
user_endpoint: Some(("iq.internal".to_string(), 61234)),
active_tasks: vec![("0".to_string(), vec![0, 1])],
target_active_tasks: vec![("0".to_string(), vec![0, 1])],
},
StreamsMemberState {
member_id: "m-2".to_string(),
member_epoch: 7,
topology_epoch: 2,
process_id: "proc-b".to_string(),
user_endpoint: None,
active_tasks: vec![("1".to_string(), vec![0])],
target_active_tasks: vec![("1".to_string(), vec![0, 1])],
},
],
},
);
});
let admin = admin_for(&broker).await;
let groups = admin
.describe_streams_groups(&["wordcount"])
.await
.expect("StreamsGroupDescribe should succeed");
assert_eq!(groups.len(), 1);
let group = &groups[0];
assert_eq!(group.group_id, "wordcount");
assert_eq!(group.group_state, "Stable");
assert_eq!(group.group_epoch, 7);
let topology = group.topology.as_ref().expect("topology must be present");
assert_eq!(topology.epoch, 3);
let subs = topology
.subtopologies
.as_ref()
.expect("subtopologies must be present, not null");
assert_eq!(subs.len(), 2);
assert_eq!(subs[0].subtopology_id, "0");
assert_eq!(subs[0].source_topics, vec!["source-topic".to_string()]);
assert_eq!(group.members.len(), 2);
let m1 = &group.members[0];
let endpoint = m1
.user_endpoint
.as_ref()
.expect("m-1 has an Interactive Queries endpoint");
assert_eq!(endpoint.host, "iq.internal");
assert_eq!(
endpoint.port, 61234,
"Endpoint.Port is uint16; decoding it signed wraps this negative"
);
assert_eq!(m1.assignment.active_tasks.len(), 1);
assert_eq!(m1.assignment.active_tasks[0].partitions, vec![0, 1]);
assert_eq!(
m1.assignment, m1.target_assignment,
"m-1 is settled on its target"
);
let m2 = &group.members[1];
assert!(m2.user_endpoint.is_none(), "m-2 configured no endpoint");
assert!(
m2.topology_epoch < topology.epoch,
"m-2 is still running an older topology"
);
assert_ne!(
m2.assignment, m2.target_assignment,
"m-2 has not finished rebalancing"
);
assert_eq!(
group.authorized_operations,
i32::MIN,
"authorized operations were not requested, so the sentinel is returned"
);
}
#[tokio::test]
async fn describe_streams_groups_distinguishes_null_from_empty() {
use crate::testing::state::StreamsGroupState;
let broker = FakeBroker::start().await.unwrap();
broker.with_state(|s| {
s.streams_groups.insert(
"no-topology".to_string(),
StreamsGroupState {
group_state: "Empty".to_string(),
topology_epoch: None,
..Default::default()
},
);
s.streams_groups.insert(
"uninitialized".to_string(),
StreamsGroupState {
group_state: "NotReady".to_string(),
topology_epoch: Some(1),
subtopologies: None,
..Default::default()
},
);
s.streams_groups.insert(
"empty-topology".to_string(),
StreamsGroupState {
group_state: "Stable".to_string(),
topology_epoch: Some(1),
subtopologies: Some(Vec::new()),
..Default::default()
},
);
});
let admin = admin_for(&broker).await;
let groups = admin
.describe_streams_groups(&["no-topology", "uninitialized", "empty-topology"])
.await
.expect("all three should decode");
let by_id: std::collections::HashMap<_, _> =
groups.iter().map(|g| (g.group_id.as_str(), g)).collect();
assert!(
by_id["no-topology"].topology.is_none(),
"a null Topology struct must decode as None"
);
assert!(
by_id["uninitialized"]
.topology
.as_ref()
.expect("topology present")
.subtopologies
.is_none(),
"a null Subtopologies array must stay None, not become an empty Vec"
);
assert_eq!(
by_id["empty-topology"]
.topology
.as_ref()
.expect("topology present")
.subtopologies
.as_ref()
.expect("present but empty")
.len(),
0,
"an empty Subtopologies array is a different state from null"
);
}
#[tokio::test]
async fn describe_streams_groups_reports_unknown_groups_individually() {
use crate::testing::state::StreamsGroupState;
let broker = FakeBroker::start().await.unwrap();
broker.with_state(|s| {
s.streams_groups.insert(
"known".to_string(),
StreamsGroupState {
group_state: "Stable".to_string(),
..Default::default()
},
);
});
let admin = admin_for(&broker).await;
let groups = admin
.describe_streams_groups(&["known", "missing"])
.await
.expect("one unknown group must not fail the call");
let by_id: std::collections::HashMap<_, _> =
groups.iter().map(|g| (g.group_id.as_str(), g)).collect();
assert!(by_id["known"].error_code.is_ok());
assert_eq!(by_id["missing"].error_code, ErrorCode::GroupIdNotFound);
}
#[tokio::test]
async fn wakeup_interrupts_a_poll_parked_on_a_fetch() {
use std::sync::Arc;
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let consumer = Arc::new(
crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("wakeup-group")
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect"),
);
consumer.subscribe(&["events"]).await.unwrap();
let _ = consumer.poll(Duration::from_secs(2)).await;
broker.on(ApiKey::Fetch, |_| Control::Delay(Duration::from_secs(20)));
let waker = Arc::clone(&consumer);
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(300)).await;
waker.wakeup();
});
let started = tokio::time::Instant::now();
let outcome = consumer.poll(Duration::from_secs(15)).await;
let elapsed = started.elapsed();
assert!(
elapsed < Duration::from_secs(10),
"wakeup() must cut the poll short, but it took {elapsed:?}"
);
assert!(
outcome.is_err(),
"an interrupted poll with no records must report the wakeup, got {outcome:?}"
);
broker.clear_hooks();
let after = consumer.poll(Duration::from_secs(2)).await;
assert!(
after.is_ok(),
"the consumer must stay usable after wakeup(), got {after:?}"
);
}
#[tokio::test]
async fn wakeup_before_poll_is_not_lost() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("wakeup-race-group")
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer.wakeup();
let outcome = consumer.poll(Duration::from_millis(500)).await;
assert!(
outcome.is_err(),
"a wakeup() before poll() must not be swallowed, got {outcome:?}"
);
let after = consumer.poll(Duration::from_millis(500)).await;
assert!(
after.is_ok(),
"the wakeup flag must be consumed by one poll, got {after:?}"
);
}
#[tokio::test]
async fn committed_reports_the_groups_offsets_from_the_coordinator() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 2);
let producer = producer_for(&broker).await;
for partition in 0..2i32 {
for i in 0..3u8 {
let record =
crate::producer::ProducerRecord::new("events", vec![i]).with_partition(partition);
let _ = producer.send_record(record).await.unwrap();
}
}
producer.close().await;
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("committed-group")
.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.unwrap();
let before = consumer
.committed(&[("events", 0), ("events", 1)])
.await
.expect("committed() should reach the coordinator");
assert!(
before
.get(&("events".to_string(), 0))
.is_none_or(|p| p.offset < 0),
"a group that has never committed must not report offset 0, got {before:?}"
);
let deadline = tokio::time::Instant::now() + SETTLE;
let mut seen = 0;
while seen < 6 && tokio::time::Instant::now() < deadline {
seen += consumer
.poll(Duration::from_millis(200))
.await
.unwrap()
.len();
}
assert_eq!(seen, 6, "all produced records should arrive");
consumer.commit_sync().await.expect("commit should succeed");
let after = consumer
.committed(&[("events", 0), ("events", 1)])
.await
.expect("committed() should reach the coordinator");
for partition in 0..2i32 {
let pos = after
.get(&("events".to_string(), partition))
.unwrap_or_else(|| panic!("partition {partition} must have a committed offset"));
assert_eq!(
pos.offset, 3,
"three records were consumed from partition {partition}"
);
}
let _ = consumer.close().await;
}
#[tokio::test]
async fn committed_without_a_group_id_is_an_error() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("an assign-only consumer needs no group");
let err = consumer
.committed(&[("events", 0)])
.await
.expect_err("no group_id means no coordinator");
assert!(
err.to_string().contains("group_id"),
"the error must name what is missing, got: {err}"
);
}
#[tokio::test]
async fn validate_only_is_refused_by_a_broker_that_predates_the_field() {
let broker = FakeBroker::start().await.unwrap();
broker.set_api_versions(ApiKey::UpdateFeatures, 0, 0);
let admin = admin_for(&broker).await;
let outcome = admin
.update_features(
vec![crate::protocol::FeatureUpdateKey::upgrade(
"metadata.version",
17,
)],
true, )
.await;
assert!(
outcome.is_err(),
"a v0 broker cannot honour validate_only, so this must not succeed"
);
assert_eq!(
broker.request_count(ApiKey::UpdateFeatures),
0,
"the request must be refused before it is sent — a dry run that \
reaches the controller has already stopped being one"
);
assert_eq!(
broker.finalized_feature("metadata.version"),
None,
"nothing may have been applied"
);
}
#[tokio::test]
async fn validate_only_reaches_a_current_broker_and_applies_nothing() {
let broker = FakeBroker::start().await.unwrap();
let admin = admin_for(&broker).await;
admin
.update_features(
vec![crate::protocol::FeatureUpdateKey::upgrade(
"metadata.version",
17,
)],
true, )
.await
.expect("a v2 broker supports validate_only");
assert_eq!(
broker.request_count(ApiKey::UpdateFeatures),
1,
"the dry run must actually be validated by the controller"
);
assert_eq!(
broker.finalized_feature("metadata.version"),
None,
"a dry run must not apply the update"
);
}
#[tokio::test]
async fn a_feature_update_is_applied_by_the_controller() {
let broker = FakeBroker::start_cluster(3).await.unwrap();
broker.set_controller(2);
let admin = admin_for(&broker).await;
admin
.update_features(
vec![crate::protocol::FeatureUpdateKey::upgrade(
"metadata.version",
17,
)],
false,
)
.await
.expect("the update should be applied");
assert_eq!(
broker.finalized_feature("metadata.version"),
Some(17),
"the controller must have applied the requested level"
);
assert_eq!(
broker.request_nodes(ApiKey::UpdateFeatures),
vec![2],
"UpdateFeatures is controller-only; reaching any other broker is the \
bug that made a controller failover look blanket-retriable"
);
}
#[tokio::test]
async fn describe_features_reports_what_update_features_applied() {
let broker = FakeBroker::start().await.unwrap();
let admin = admin_for(&broker).await;
let before = admin
.describe_features()
.await
.expect("describe_features should work on a cluster with no features");
assert!(
before.finalized_features.is_empty(),
"a cluster that has finalized nothing must report nothing"
);
assert!(
before.finalized_features_epoch < 0,
"an absent epoch means the finalized list is not to be trusted, saw {}",
before.finalized_features_epoch
);
admin
.update_features(
vec![crate::protocol::FeatureUpdateKey::upgrade(
"metadata.version",
17,
)],
false,
)
.await
.expect("the update should be applied");
let after = admin
.describe_features()
.await
.expect("describe_features should work after an update");
let finalized = after
.finalized_features
.iter()
.find(|f| f.name == "metadata.version")
.expect("the finalized feature must be reported back");
assert_eq!(
finalized.max_version_level, 17,
"the level read back must be the level applied — transposing the \
(max, min) pair here decodes without error and reports 1"
);
assert!(
after.finalized_features_epoch >= 0,
"finalized features are only valid alongside a non-negative epoch"
);
let supported = after
.supported_features
.iter()
.find(|f| f.name == "metadata.version")
.expect("the broker must also advertise what it supports");
assert!(
supported.max_version >= finalized.max_version_level,
"a broker cannot finalize a level above what it supports: {} < {}",
supported.max_version,
finalized.max_version_level
);
}
#[cfg(feature = "unstable-protocol")]
async fn share_consumer_for(
broker: &FakeBroker,
group_id: &str,
) -> crate::share_consumer::ShareConsumer {
share_consumer_with(broker, group_id, |b| b).await
}
#[cfg(feature = "unstable-protocol")]
async fn share_consumer_with(
broker: &FakeBroker,
group_id: &str,
tune: impl FnOnce(
crate::share_consumer::ShareConsumerBuilder,
) -> crate::share_consumer::ShareConsumerBuilder,
) -> crate::share_consumer::ShareConsumer {
tune(
crate::share_consumer::ShareConsumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id(group_id)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT),
)
.build()
.await
.expect("share consumer should connect")
}
#[cfg(feature = "unstable-protocol")]
async fn drain_share(
consumer: &crate::share_consumer::ShareConsumer,
want: usize,
) -> Vec<crate::consumer::ConsumerRecord> {
let deadline = tokio::time::Instant::now() + SETTLE;
let mut got = Vec::new();
while got.len() < want && tokio::time::Instant::now() < deadline {
match consumer.poll(Duration::from_millis(200)).await {
Ok(records) => got.extend(records),
Err(e) => panic!("share poll failed: {e}"),
}
}
got
}
#[cfg(feature = "unstable-protocol")]
#[tokio::test]
async fn a_share_consumer_receives_records_and_counts_them() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = producer_for(&broker).await;
for i in 0..5u8 {
let _ = producer
.send("events", None, Some(&[b'v', i]))
.await
.expect("send should be acknowledged");
}
producer.close().await;
let consumer = share_consumer_for(&broker, "delivery-group").await;
let metrics = consumer.metrics();
assert_eq!(metrics.records_received.get(), 0, "nothing polled yet");
consumer
.subscribe(&["events"])
.await
.expect("subscribe should reach the coordinator");
let records = drain_share(&consumer, 5).await;
assert_eq!(records.len(), 5, "every produced record must be delivered");
let mut payloads: Vec<Vec<u8>> = records
.iter()
.filter_map(|r| r.value.as_ref().map(|v| v.to_vec()))
.collect();
payloads.sort();
assert_eq!(
payloads,
(0..5u8).map(|i| vec![b'v', i]).collect::<Vec<_>>(),
"delivered payloads must be the produced ones"
);
assert_eq!(
metrics.records_received.get(),
5,
"records_received must match what poll() actually returned"
);
assert!(
metrics.bytes_received.get() >= 10,
"five two-byte values is at least ten bytes, saw {}",
metrics.bytes_received.get()
);
assert!(
metrics.polls.get() >= 1,
"every poll() must be counted, empty or not"
);
let _ = consumer.close().await;
}
#[cfg(feature = "unstable-protocol")]
#[tokio::test]
async fn accepting_retires_a_record_and_releasing_redelivers_it() {
use crate::share_consumer::AcknowledgeType;
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = producer_for(&broker).await;
for i in 0..2u8 {
let _ = producer.send("events", None, Some(&[i])).await.unwrap();
}
producer.close().await;
let consumer = share_consumer_with(&broker, "ack-group", |b| {
b.acknowledgement_mode(crate::share_consumer::AcknowledgementMode::Explicit)
})
.await;
consumer.subscribe(&["events"]).await.unwrap();
let first = drain_share(&consumer, 2).await;
assert_eq!(first.len(), 2);
for record in &first {
let ack = if record.offset == 0 {
AcknowledgeType::Accept
} else {
AcknowledgeType::Release
};
consumer
.acknowledge(record, ack)
.await
.expect("acknowledgement should be accepted");
}
let redelivered = drain_share(&consumer, 1).await;
assert!(
!redelivered.is_empty(),
"a released record must be handed out again"
);
assert!(
redelivered.iter().all(|r| r.offset == 1),
"only the released offset may come back, saw {:?}",
redelivered.iter().map(|r| r.offset).collect::<Vec<_>>()
);
assert!(
redelivered
.iter()
.all(|r| r.delivery_count.is_some_and(|c| c >= 2)),
"a redelivery must report a delivery count above one, saw {:?}",
redelivered
.iter()
.map(|r| r.delivery_count)
.collect::<Vec<_>>()
);
let _ = consumer.close().await;
}
#[cfg(feature = "unstable-protocol")]
#[tokio::test]
async fn an_accepted_record_is_not_redelivered_to_the_next_member() {
use crate::share_consumer::{AcknowledgeType, AcknowledgementMode};
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = producer_for(&broker).await;
for i in 0..2u8 {
let _ = producer.send("events", None, Some(&[i])).await.unwrap();
}
producer.close().await;
let first = share_consumer_with(&broker, "restart-group", |b| {
b.acknowledgement_mode(AcknowledgementMode::Explicit)
})
.await;
first.subscribe(&["events"]).await.unwrap();
let records = drain_share(&first, 2).await;
assert_eq!(
records.len(),
2,
"both records should reach the first member"
);
let accepted = records
.iter()
.find(|r| r.offset == 0)
.expect("offset 0 should have been delivered");
first
.acknowledge(accepted, AcknowledgeType::Accept)
.await
.expect("acknowledgement should be accepted");
first.close().await.expect("close should flush the ack");
let second = share_consumer_for(&broker, "restart-group").await;
second.subscribe(&["events"]).await.unwrap();
let redelivered = drain_share(&second, 1).await;
assert!(
!redelivered.is_empty(),
"the unacknowledged record must be handed to the next member"
);
assert!(
redelivered.iter().all(|r| r.offset == 1),
"an accepted record must never come back, saw offsets {:?}",
redelivered.iter().map(|r| r.offset).collect::<Vec<_>>()
);
let _ = second.close().await;
}
#[cfg(feature = "unstable-protocol")]
#[tokio::test]
async fn two_share_group_members_split_the_partitions() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 2);
let producer = producer_for(&broker).await;
for partition in 0..2i32 {
for i in 0..3u8 {
let record = crate::producer::ProducerRecord::new("events", vec![partition as u8, i])
.with_partition(partition);
let _ = producer.send_record(record).await.unwrap();
}
}
producer.close().await;
let a = share_consumer_for(&broker, "split-group").await;
let b = share_consumer_for(&broker, "split-group").await;
a.subscribe(&["events"]).await.unwrap();
b.subscribe(&["events"]).await.unwrap();
let mut seen: Vec<Vec<u8>> = Vec::new();
let drain = async |c: &crate::share_consumer::ShareConsumer, into: &mut Vec<Vec<u8>>| {
if let Ok(records) = c.poll(Duration::from_millis(100)).await {
into.extend(
records
.into_iter()
.filter_map(|r| r.value.map(|v| v.to_vec())),
);
}
};
let deadline = tokio::time::Instant::now() + SETTLE;
while tokio::time::Instant::now() < deadline {
let (assign_a, assign_b) = (a.assignment().await, b.assignment().await);
let count = |m: &ahash::AHashMap<String, Vec<crate::PartitionId>>| {
m.values().map(Vec::len).sum::<usize>()
};
if count(&assign_a) == 1 && count(&assign_b) == 1 {
break;
}
drain(&a, &mut seen).await;
drain(&b, &mut seen).await;
}
let assign_a = a.assignment().await;
let assign_b = b.assignment().await;
let partitions = |m: &ahash::AHashMap<String, Vec<crate::PartitionId>>| {
m.values().flatten().copied().collect::<HashSet<_>>()
};
let (pa, pb) = (partitions(&assign_a), partitions(&assign_b));
assert_eq!(pa.len(), 1, "each member should hold one of two partitions");
assert_eq!(pb.len(), 1);
assert!(
pa.is_disjoint(&pb),
"the coordinator must not hand one partition to both members: {pa:?} vs {pb:?}"
);
let deadline = tokio::time::Instant::now() + SETTLE;
while seen.len() < 6 && tokio::time::Instant::now() < deadline {
drain(&a, &mut seen).await;
drain(&b, &mut seen).await;
}
seen.sort();
let expected: Vec<Vec<u8>> = (0..2u8)
.flat_map(|p| (0..3u8).map(move |i| vec![p, i]))
.collect();
assert_eq!(
seen, expected,
"between them the two members must see every record exactly once"
);
let _ = a.close().await;
let _ = b.close().await;
}
#[cfg(feature = "unstable-protocol")]
#[tokio::test]
async fn share_consumer_poll_metrics_are_wired() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let consumer = share_consumer_for(&broker, "metrics-share-group").await;
let metrics = consumer.metrics();
assert_eq!(metrics.polls.get(), 0, "no poll has happened yet");
assert_eq!(metrics.empty_polls.get(), 0);
for _ in 0..3 {
let _ = consumer.poll(Duration::from_millis(20)).await;
}
assert_eq!(
metrics.polls.get(),
3,
"every poll() must be counted, empty or not"
);
assert_eq!(
metrics.empty_polls.get(),
3,
"a poll with no assignment is an empty poll"
);
assert_eq!(
metrics.records_received.get(),
0,
"nothing was delivered, so nothing may be counted as delivered"
);
let _ = consumer.close().await;
}
use crate::consumer::{AutoOffsetReset, Consumer, IsolationLevel};
use crate::producer::{TopicPartitionOffset, TransactionVersion, TransactionalProducer};
async fn txn_producer_for(broker: &FakeBroker, transactional_id: &str) -> TransactionalProducer {
TransactionalProducer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.transactional_id(transactional_id)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("transactional producer should connect")
}
async fn reader_for(broker: &FakeBroker, topic: &str, isolation: IsolationLevel) -> Consumer {
let consumer = Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.auto_offset_reset(AutoOffsetReset::Earliest)
.isolation_level(isolation)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.assign(topic, vec![0])
.await
.expect("manual assignment should succeed");
consumer
}
async fn drain(consumer: &Consumer, polls: usize) -> Vec<String> {
let mut values = Vec::new();
for _ in 0..polls {
if let Ok(records) = consumer.poll(Duration::from_millis(200)).await {
for record in records {
let value = record.value.as_deref().unwrap_or_default();
values.push(String::from_utf8_lossy(value).into_owned());
}
}
}
values
}
#[tokio::test]
async fn committed_transaction_becomes_visible_to_read_committed() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let producer = txn_producer_for(&broker, "txn-visible").await;
producer.init_transactions().await.expect("init");
producer.begin_transaction().expect("begin");
let _ = producer
.send("orders", None, Some(b"committed-1"))
.await
.expect("send");
producer.flush().await.expect("flush");
assert_eq!(
broker.last_stable_offset("orders", 0),
Some(0),
"an open transaction must pin the last stable offset at its first record"
);
assert!(broker.transaction_is_open("txn-visible"));
producer.commit_transaction().await.expect("commit");
assert!(!broker.transaction_is_open("txn-visible"));
let consumer = reader_for(&broker, "orders", IsolationLevel::ReadCommitted).await;
let values = drain(&consumer, 6).await;
assert_eq!(
values,
vec!["committed-1".to_string()],
"a committed transaction must be delivered exactly once"
);
let _ = consumer.close().await;
producer.close().await;
}
#[tokio::test]
async fn aborted_transaction_is_filtered_only_for_read_committed() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let producer = txn_producer_for(&broker, "txn-abort").await;
producer.init_transactions().await.expect("init");
producer.begin_transaction().expect("begin");
let _ = producer
.send("orders", None, Some(b"doomed"))
.await
.expect("send");
producer.abort_transaction().await.expect("abort");
let (producer_id, _) = broker
.transactional_producer("txn-abort")
.expect("the coordinator knows this transactional id");
assert_eq!(
broker.aborted_transactions("orders", 0),
vec![(producer_id, 0)],
"the abort must be recorded so a read_committed fetch can report it"
);
let committed = reader_for(&broker, "orders", IsolationLevel::ReadCommitted).await;
assert!(
drain(&committed, 6).await.is_empty(),
"read_committed must not surface records from an aborted transaction"
);
let _ = committed.close().await;
let uncommitted = reader_for(&broker, "orders", IsolationLevel::ReadUncommitted).await;
assert_eq!(
drain(&uncommitted, 6).await,
vec!["doomed".to_string()],
"the records were written — read_uncommitted proves the filtering above \
was filtering, not an empty log"
);
let _ = uncommitted.close().await;
producer.close().await;
}
#[tokio::test]
async fn an_abort_does_not_poison_the_next_transaction() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let producer = txn_producer_for(&broker, "txn-sequence").await;
producer.init_transactions().await.expect("init");
producer.begin_transaction().expect("begin 1");
let _ = producer
.send("orders", None, Some(b"aborted"))
.await
.expect("send");
producer.abort_transaction().await.expect("abort");
producer.begin_transaction().expect("begin 2");
let _ = producer
.send("orders", None, Some(b"committed"))
.await
.expect("send");
producer.commit_transaction().await.expect("commit");
let consumer = reader_for(&broker, "orders", IsolationLevel::ReadCommitted).await;
assert_eq!(
drain(&consumer, 8).await,
vec!["committed".to_string()],
"the second transaction must survive the first one's abort"
);
let _ = consumer.close().await;
producer.close().await;
}
#[tokio::test]
async fn the_fake_brokers_two_record_views_differ_by_the_aborted_records() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let producer = txn_producer_for(&broker, "txn-views").await;
producer.init_transactions().await.expect("init");
producer.begin_transaction().expect("begin 1");
for i in 0..3 {
let _ = producer
.send("orders", None, Some(format!("aborted-{i}").as_bytes()))
.await
.expect("send");
}
producer.abort_transaction().await.expect("abort");
producer.begin_transaction().expect("begin 2");
for i in 0..2 {
let _ = producer
.send("orders", None, Some(format!("committed-{i}").as_bytes()))
.await
.expect("send");
}
producer.commit_transaction().await.expect("commit");
producer.close().await;
let committed = broker.committed_records("orders").expect("log decodes");
let all = broker.all_records("orders").expect("log decodes");
let values = |records: &[crate::consumer::ConsumerRecord]| -> Vec<String> {
records
.iter()
.filter_map(|r| r.value.as_ref())
.map(|v| String::from_utf8_lossy(v).into_owned())
.collect()
};
assert_eq!(
values(&committed),
vec!["committed-0".to_string(), "committed-1".to_string()],
"an aborted transaction must be invisible to the committed view"
);
assert_eq!(
values(&all).len(),
5,
"the uncommitted view must show every record, aborted included"
);
assert!(
!values(&all).iter().any(|v| v.starts_with("__")),
"control batches are never records"
);
let consumer = reader_for(&broker, "orders", IsolationLevel::ReadCommitted).await;
assert_eq!(
drain(&consumer, 8).await,
values(&committed),
"committed_records() must match what a read_committed consumer receives"
);
let _ = consumer.close().await;
}
#[tokio::test]
async fn transactional_offsets_move_only_on_commit() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let producer = txn_producer_for(&broker, "txn-offsets").await;
producer.init_transactions().await.expect("init");
let group = crate::consumer::ConsumerGroupMetadata::new("etl-group", 7, "member-1", None);
let offsets = vec![TopicPartitionOffset::new("orders", 0, 42)];
producer.begin_transaction().expect("begin 1");
let _ = producer
.send("orders", None, Some(b"x"))
.await
.expect("send");
producer
.send_offsets_to_transaction(&offsets, &group)
.await
.expect("stage offsets");
producer.abort_transaction().await.expect("abort");
assert_eq!(
broker.committed_offset("etl-group", "orders", 0),
None,
"an aborted transaction must not commit the offsets it staged"
);
producer.begin_transaction().expect("begin 2");
let _ = producer
.send("orders", None, Some(b"y"))
.await
.expect("send");
producer
.send_offsets_to_transaction(&offsets, &group)
.await
.expect("stage offsets");
producer.commit_transaction().await.expect("commit");
assert_eq!(
broker.committed_offset("etl-group", "orders", 0),
Some(42),
"a committed transaction must apply the offsets it staged"
);
producer.close().await;
}
#[tokio::test]
async fn re_initialising_fences_the_previous_producer() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let zombie = txn_producer_for(&broker, "txn-fenced").await;
zombie.init_transactions().await.expect("init");
let (first_pid, first_epoch) = broker.transactional_producer("txn-fenced").unwrap();
let successor = txn_producer_for(&broker, "txn-fenced").await;
successor.init_transactions().await.expect("init");
let (second_pid, second_epoch) = broker.transactional_producer("txn-fenced").unwrap();
assert_eq!(
second_pid, first_pid,
"the producer ID is the fencing identity and must be stable"
);
assert!(
second_epoch > first_epoch,
"a new incarnation must get a higher epoch ({second_epoch} vs {first_epoch})"
);
zombie.begin_transaction().expect("begin");
let outcome = zombie.send("orders", None, Some(b"zombie")).await;
let outcome = match outcome {
Err(e) => Err(e),
Ok(_) => zombie.commit_transaction().await.map(|()| unreachable!()),
};
assert!(
outcome.is_err(),
"a fenced producer must not be able to complete a transaction"
);
assert_eq!(
zombie.state(),
crate::producer::TransactionState::FatalError,
"a fencing error is fatal: the producer must be recreated, not retried"
);
successor.close().await;
}
#[tokio::test]
async fn tv1_registers_partitions_with_add_partitions_to_txn() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let producer = txn_producer_for(&broker, "txn-tv1").await;
producer.init_transactions().await.expect("init");
assert_eq!(
producer.transaction_version(),
TransactionVersion::V1,
"a cluster that has not finalized transaction.version is TV1"
);
producer.begin_transaction().expect("begin");
let _ = producer
.send("orders", None, Some(b"v1"))
.await
.expect("send");
producer.commit_transaction().await.expect("commit");
assert!(
broker.request_count(ApiKey::AddPartitionsToTxn) > 0,
"TV1 must register each partition with the coordinator before writing"
);
let consumer = reader_for(&broker, "orders", IsolationLevel::ReadCommitted).await;
assert_eq!(drain(&consumer, 6).await, vec!["v1".to_string()]);
let _ = consumer.close().await;
producer.close().await;
}
#[tokio::test]
async fn tv2_commits_without_add_partitions_to_txn() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
broker.set_transaction_version(2);
let producer = txn_producer_for(&broker, "txn-tv2").await;
producer.init_transactions().await.expect("init");
assert_eq!(
producer.transaction_version(),
TransactionVersion::V2,
"a cluster finalizing transaction.version=2 must negotiate TV2"
);
let (_, epoch_before) = broker.transactional_producer("txn-tv2").unwrap();
producer.begin_transaction().expect("begin");
let _ = producer
.send("orders", None, Some(b"v2"))
.await
.expect("send");
producer.commit_transaction().await.expect("commit");
assert_eq!(
broker.request_count(ApiKey::AddPartitionsToTxn),
0,
"TV2 carries the transactional ID on Produce; the extra round trip must be gone"
);
let (_, epoch_after) = broker.transactional_producer("txn-tv2").unwrap();
assert!(
epoch_after > epoch_before,
"KIP-890 bumps the producer epoch at every transaction completion"
);
let consumer = reader_for(&broker, "orders", IsolationLevel::ReadCommitted).await;
assert_eq!(drain(&consumer, 6).await, vec!["v2".to_string()]);
let _ = consumer.close().await;
producer.close().await;
}
#[tokio::test]
async fn a_multi_partition_transaction_commits_atomically() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 2);
let producer = txn_producer_for(&broker, "txn-multi").await;
producer.init_transactions().await.expect("init");
producer.begin_transaction().expect("begin");
for partition in 0..2 {
let record = crate::producer::ProducerRecord::new(
"orders",
bytes::Bytes::from(format!("p{partition}")),
)
.with_partition(partition);
let _ = producer.send_record(record).await.expect("send");
}
producer.flush().await.expect("flush");
for partition in 0..2 {
assert_eq!(
broker.last_stable_offset("orders", partition),
Some(0),
"every partition in the transaction must be pinned until the commit"
);
}
producer.commit_transaction().await.expect("commit");
for partition in 0..2 {
assert!(
broker.last_stable_offset("orders", partition) > Some(0),
"the commit must release every partition, not just the first"
);
}
producer.close().await;
}
#[tokio::test]
async fn an_old_abort_is_not_reported_to_a_consumer_that_has_read_past_it() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let producer = txn_producer_for(&broker, "txn-stale").await;
producer.init_transactions().await.expect("init");
producer.begin_transaction().expect("begin 1");
let _ = producer
.send("orders", None, Some(b"aborted"))
.await
.expect("send");
producer.abort_transaction().await.expect("abort");
producer.begin_transaction().expect("begin 2");
let _ = producer
.send("orders", None, Some(b"committed"))
.await
.expect("send");
producer.commit_transaction().await.expect("commit");
let marker_end = 2;
assert!(
broker
.aborted_transactions("orders", 0)
.iter()
.all(|(_, first)| *first < marker_end),
"the abort under test must lie below the fetch offset"
);
let consumer = Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.isolation_level(IsolationLevel::ReadCommitted)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.assign("orders", vec![0])
.await
.expect("manual assignment");
consumer
.seek("orders", 0, marker_end)
.await
.expect("seek past the abort marker");
assert_eq!(
drain(&consumer, 8).await,
vec!["committed".to_string()],
"a consumer that has read past an abort must still receive the \
committed transaction that follows it"
);
let _ = consumer.close().await;
producer.close().await;
}
#[tokio::test]
async fn a_second_poll_is_served_from_the_prefetch_buffer_without_a_fetch() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = crate::producer::Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("producer should connect");
for i in 0..20u32 {
let _ = producer
.send("events", None, Some(format!("v{i}").as_bytes()))
.await
.expect("send");
}
producer.close().await;
let consumer = Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.auto_offset_reset(AutoOffsetReset::Earliest)
.max_poll_records(5)
.max_buffered_records(50)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer
.assign("events", vec![0])
.await
.expect("manual assignment");
let first = consumer
.poll(Duration::from_millis(500))
.await
.expect("first poll");
assert_eq!(first.len(), 5, "the delivery cap must still be honoured");
let fetches_after_first = broker.request_count(ApiKey::Fetch);
assert!(
fetches_after_first >= 1,
"the first poll has to reach the broker"
);
let second = consumer
.poll(Duration::from_millis(500))
.await
.expect("second poll");
assert_eq!(second.len(), 5, "the buffer must serve a full batch");
assert_eq!(
broker.request_count(ApiKey::Fetch),
fetches_after_first,
"a poll served from the prefetch buffer must not issue a Fetch"
);
let seen: Vec<i64> = first
.iter()
.chain(second.iter())
.map(|r| r.offset)
.collect();
assert_eq!(seen, (0..10).collect::<Vec<i64>>());
let _ = consumer.close().await;
}
#[tokio::test]
async fn a_commit_never_acknowledges_records_still_parked_in_the_buffer() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let producer = crate::producer::Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("producer should connect");
for i in 0..20u32 {
let _ = producer
.send("events", None, Some(format!("v{i}").as_bytes()))
.await
.expect("send");
}
producer.close().await;
let consumer = Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.group_id("prefetch-commit-group")
.auto_offset_reset(AutoOffsetReset::Earliest)
.enable_auto_commit(false)
.max_poll_records(5)
.max_buffered_records(50)
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("consumer should connect");
consumer.subscribe(&["events"]).await.unwrap();
let mut delivered = 0usize;
for _ in 0..10 {
let records = consumer
.poll(Duration::from_millis(300))
.await
.expect("poll");
delivered += records.len();
if delivered > 0 {
break;
}
}
assert!(delivered > 0, "the consumer must receive something");
let position = consumer
.position("events", 0)
.await
.expect("position must be tracked");
let fetch_position = consumer
.fetch_position("events", 0)
.await
.expect("fetch position must be tracked");
assert_eq!(
position, delivered as i64,
"position() must report the delivered offset, not the read-ahead"
);
assert!(
fetch_position > position,
"the consumer must have read ahead of delivery, got fetch={fetch_position} \
position={position}"
);
consumer.commit().await.expect("commit");
let committed = broker
.committed_offset("prefetch-commit-group", "events", 0)
.expect("the group must have a committed offset");
assert_eq!(
committed, delivered as i64,
"the commit must acknowledge exactly what was delivered — a commit at \
the fetch position would skip the parked surplus on restart"
);
assert_eq!(
committed, position,
"commit() and position() must never disagree"
);
let _ = consumer.close().await;
}
#[tokio::test]
async fn a_commit_stops_admitting_records_before_it_drains() {
use std::sync::Arc;
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let producer = Arc::new(
TransactionalProducer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.transactional_id("txn-commit-ordering")
.linger(Duration::from_millis(200))
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("transactional producer should connect"),
);
producer.init_transactions().await.expect("init");
producer.begin_transaction().expect("begin");
let buffered = Arc::clone(&producer);
let send = tokio::spawn(async move { buffered.send("orders", None, Some(b"first")).await });
broker.on(ApiKey::Produce, |_| {
Control::Delay(Duration::from_millis(600))
});
let committing = Arc::clone(&producer);
let commit = tokio::spawn(async move { committing.commit_transaction().await });
tokio::time::sleep(Duration::from_millis(150)).await;
assert_eq!(
producer.state(),
crate::producer::TransactionState::Committing,
"the commit must own the state while it drains, so no further record \
can be admitted into a transaction that is already closing"
);
let refused = producer.send("orders", None, Some(b"too-late")).await;
let error = refused.expect_err("a send during the commit's drain must be refused");
assert!(
error.to_string().contains("Committing"),
"the refusal must name the state that caused it, got: {error}"
);
broker.clear_hooks();
let _ = send.await.expect("send task should not panic");
let _ = commit.await.expect("commit task should not panic");
producer.close().await;
}
#[tokio::test]
async fn a_commit_waits_for_an_in_flight_offset_commit() {
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use crate::consumer::ConsumerGroupMetadata;
use crate::producer::TopicPartitionOffset;
let broker = FakeBroker::start_cluster(2).await.unwrap();
broker.create_topic("orders", 1);
broker.set_group_coordinator("g", 0);
broker.set_txn_coordinator("txn-offsets-ordering", 1);
let producer = Arc::new(
TransactionalProducer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.transactional_id("txn-offsets-ordering")
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("transactional producer should connect"),
);
producer.init_transactions().await.expect("init");
producer.begin_transaction().expect("begin");
let _ = producer
.send("orders", None, Some(b"payload"))
.await
.expect("send");
broker.on(ApiKey::TxnOffsetCommit, |_| {
Control::Delay(Duration::from_millis(500))
});
let done = Arc::new(AtomicBool::new(false));
let offsets_done = Arc::clone(&done);
let offsets_producer = Arc::clone(&producer);
let offsets = tokio::spawn(async move {
let metadata = ConsumerGroupMetadata::new("g", 1, "member-1", None);
let result = offsets_producer
.send_offsets_to_transaction(&[TopicPartitionOffset::new("orders", 0, 42)], &metadata)
.await;
offsets_done.store(true, Ordering::SeqCst);
result
});
tokio::time::sleep(Duration::from_millis(100)).await;
assert!(
!done.load(Ordering::SeqCst),
"the offset commit must still be in flight for this test to mean anything"
);
producer.commit_transaction().await.expect("commit");
assert!(
done.load(Ordering::SeqCst),
"commit_transaction() returned while TxnOffsetCommit was still in flight — the EndTxn marker would have been written with the offsets outside the transaction"
);
let offsets_result = offsets.await.expect("offset task should not panic");
assert!(
offsets_result.is_ok(),
"the offset commit should complete inside the transaction: {offsets_result:?}"
);
broker.clear_hooks();
producer.close().await;
}
#[cfg(feature = "unstable-protocol")]
#[tokio::test]
async fn a_flush_waits_for_a_poll_holding_the_acknowledgements() {
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let consumer = Arc::new(share_consumer_for(&broker, "share-flush-race").await);
consumer.subscribe(&["events"]).await.unwrap();
let _ = consumer.poll(Duration::from_millis(300)).await;
broker.on(ApiKey::ShareFetch, |_| {
Control::Delay(Duration::from_millis(600))
});
let polling = Arc::clone(&consumer);
let poll_done = Arc::new(AtomicBool::new(false));
let poll_flag = Arc::clone(&poll_done);
let poll = tokio::spawn(async move {
let out = polling.poll(Duration::from_secs(2)).await;
poll_flag.store(true, Ordering::SeqCst);
out
});
tokio::time::sleep(Duration::from_millis(150)).await;
assert!(
!poll_done.load(Ordering::SeqCst),
"the poll must still be in flight for this test to mean anything"
);
consumer.commit_sync().await.expect("commit_sync");
assert!(
poll_done.load(Ordering::SeqCst),
"commit_sync() returned while a poll was still holding the pending \
acknowledgements — it would have flushed an empty map and reported \
success, stranding them"
);
broker.clear_hooks();
let _ = poll.await.expect("poll task should not panic");
let _ = consumer.close().await;
}
#[tokio::test]
async fn a_tombstone_reaches_the_log_as_a_null_value() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("users", 1);
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("producer should connect");
let _ = producer
.send("users", Some(b"user-42"), Some(b"alice"))
.await
.expect("valued send should be acknowledged");
let _ = producer
.send("users", Some(b"user-42"), Some(b""))
.await
.expect("empty send should be acknowledged");
let _ = producer
.send_record(
crate::producer::ProducerRecord::tombstone("users", "user-42")
.with_header("X-Reason", &b"gdpr-erasure"[..])
.with_null_header("X-Flag"),
)
.await
.expect("tombstone send should be acknowledged");
let stored = broker.all_records("users").expect("log should be readable");
assert_eq!(stored.len(), 3, "all three records should be in the log");
assert_eq!(stored[0].value.as_deref(), Some(&b"alice"[..]));
assert!(!stored[0].is_tombstone());
assert_eq!(
stored[1].value.as_deref(),
Some(&b""[..]),
"an empty value must not decode as null"
);
assert!(
!stored[1].is_tombstone(),
"a zero-length value is an ordinary record"
);
assert_eq!(stored[2].value, None, "the tombstone must decode as null");
assert!(stored[2].is_tombstone());
assert_eq!(stored[2].key.as_deref(), Some(&b"user-42"[..]));
let headers = &stored[2].headers;
assert_eq!(headers[0].0.as_ref(), b"X-Reason");
assert_eq!(headers[0].1.as_deref(), Some(&b"gdpr-erasure"[..]));
assert_eq!(headers[1].0.as_ref(), b"X-Flag");
assert_eq!(headers[1].1, None, "a null header value must stay null");
}
#[tokio::test]
async fn a_tombstone_routes_to_the_same_partition_as_its_key() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("users", 8);
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("producer should connect");
for key in ["user-1", "user-2", "user-3", "user-4"] {
let valued = producer
.send("users", Some(key.as_bytes()), Some(b"payload"))
.await
.expect("valued send should be acknowledged");
let tombstone = producer
.send("users", Some(key.as_bytes()), None)
.await
.expect("tombstone send should be acknowledged");
assert_eq!(
valued.partition, tombstone.partition,
"the tombstone for {key} must share its record's partition"
);
}
}
#[tokio::test]
async fn a_produced_tombstone_deletes_a_key_from_a_compacted_table() {
use crate::consumer::CompactedTable;
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("users", 1);
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.build()
.await
.expect("producer should connect");
let _ = producer
.send("users", Some(b"keep-me"), Some(b"v"))
.await
.expect("send should be acknowledged");
let _ = producer
.send("users", Some(b"delete-me"), Some(b"v"))
.await
.expect("send should be acknowledged");
let mut table = CompactedTable::new();
table.ingest(&broker.all_records("users").expect("log should be readable"));
assert_eq!(table.len(), 2, "both keys should be present");
let _ = producer
.send_record(crate::producer::ProducerRecord::tombstone(
"users",
"delete-me",
))
.await
.expect("tombstone send should be acknowledged");
let mut table = CompactedTable::new();
let changes = table.apply(&broker.all_records("users").expect("log should be readable"));
assert_eq!(table.len(), 1, "the tombstoned key should be gone");
assert!(table.get_value(b"keep-me").is_some());
assert!(table.get_value(b"delete-me").is_none());
assert!(
changes.iter().any(|c| c.is_delete()),
"the tombstone should be reported as a deletion"
);
}
#[derive(Debug, Default)]
struct RecordingInterceptor {
sends: std::sync::atomic::AtomicUsize,
acks: std::sync::Mutex<Vec<AckObservation>>,
}
#[derive(Debug)]
struct AckObservation {
token: Option<String>,
header_keys: Vec<String>,
partition: crate::PartitionId,
offset: i64,
delivery: crate::producer::DeliveryConfirmation,
failed: bool,
}
struct SendToken(String);
impl RecordingInterceptor {
fn sends(&self) -> usize {
self.sends.load(std::sync::atomic::Ordering::SeqCst)
}
fn acks(&self) -> std::sync::MutexGuard<'_, Vec<AckObservation>> {
self.acks.lock().unwrap()
}
}
impl crate::interceptor::ProducerInterceptor for RecordingInterceptor {
fn on_send(
&self,
record: &mut crate::producer::ProducerRecord,
ctx: &mut crate::interceptor::RecordContext,
) -> crate::interceptor::InterceptorResult {
let n = self.sends.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
ctx.insert(SendToken(format!("{}#{n}", record.topic)));
Ok(())
}
fn on_acknowledgement(
&self,
metadata: &crate::producer::RecordMetadata,
error: Option<&crate::error::KrafkaError>,
headers: &crate::producer::RecordHeaders,
ctx: &mut crate::interceptor::RecordContext,
) -> crate::interceptor::InterceptorResult {
self.acks().push(AckObservation {
token: ctx.take::<SendToken>().map(|t| t.0),
header_keys: headers.iter().map(|(k, _)| k.clone()).collect(),
partition: metadata.partition,
offset: metadata.offset,
delivery: metadata.delivery,
failed: error.is_some(),
});
Ok(())
}
}
async fn producer_with(
broker: &FakeBroker,
interceptor: &std::sync::Arc<RecordingInterceptor>,
) -> Producer {
Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.max_block(SHORT_MAX_BLOCK)
.add_interceptor(std::sync::Arc::clone(interceptor) as std::sync::Arc<_>)
.build()
.await
.expect("producer should connect")
}
const SHORT_MAX_BLOCK: Duration = Duration::from_secs(2);
#[tokio::test]
async fn a_partial_refresh_for_one_topic_does_not_strand_another() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("topic-a", 1);
broker.create_topic("topic-b", 1);
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.max_block(SHORT_MAX_BLOCK)
.metadata_max_age(Duration::from_millis(50))
.metadata_topic_cache_ttl(Duration::from_millis(50))
.build()
.await
.expect("producer should connect");
let _ = producer
.send("topic-a", None, Some(b"warm-a"))
.await
.expect("warm a");
let _ = producer
.send("topic-b", None, Some(b"warm-b"))
.await
.expect("warm b");
tokio::time::sleep(Duration::from_millis(120)).await;
let _ = producer
.send("topic-a", None, Some(b"refresh-a"))
.await
.expect("a stale topic refreshes itself");
let _ = producer
.send("topic-b", None, Some(b"after-refresh"))
.await
.expect("a topic that exists must stay sendable after an unrelated refresh");
producer.close().await;
}
#[tokio::test]
async fn a_topic_in_active_use_is_not_evicted_by_an_unrelated_refresh() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("hot", 1);
broker.create_topic("cold", 1);
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.max_block(SHORT_MAX_BLOCK)
.metadata_max_age(Duration::from_secs(30))
.metadata_topic_cache_ttl(Duration::from_millis(400))
.build()
.await
.expect("producer should connect");
for _ in 0..6 {
let _ = producer
.send("hot", None, Some(b"v"))
.await
.expect("hot send");
producer
.metadata()
.refresh_for_topics_forced(Some(&["cold"]))
.await
.expect("partial refresh");
tokio::time::sleep(Duration::from_millis(60)).await;
}
assert!(
producer.metadata().partition_count("hot").is_some(),
"a topic still being produced to must never be evicted as idle"
);
producer.close().await;
}
#[tokio::test]
async fn a_standalone_subscription_picks_up_a_topic_created_later() {
let broker = FakeBroker::start().await.unwrap();
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.build()
.await
.expect("consumer should connect");
consumer
.subscribe(&["appears-later"])
.await
.expect("subscribe should not fail on an absent topic");
assert!(
consumer.assignment().await.is_empty(),
"nothing to assign while the topic does not exist"
);
broker.create_topic("appears-later", 2);
let mut assignment = consumer.assignment().await;
for _ in 0..40 {
let _ = consumer.poll(Duration::from_millis(50)).await;
assignment = consumer.assignment().await;
if !assignment.is_empty() {
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert_eq!(
assignment.get("appears-later").map(Vec::len),
Some(2),
"both partitions of the newly created topic must be assigned, got {assignment:?}"
);
consumer.close().await.expect("close");
}
#[tokio::test]
async fn a_standalone_subscription_picks_up_new_partitions() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("grows", 1);
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.metadata_max_age(Duration::from_millis(50))
.build()
.await
.expect("consumer should connect");
consumer.subscribe(&["grows"]).await.expect("subscribe");
assert_eq!(consumer.assignment().await["grows"].len(), 1);
assert_eq!(broker.add_partitions("grows", 4), 3);
let mut assigned = 1;
for _ in 0..40 {
let _ = consumer.poll(Duration::from_millis(50)).await;
assigned = consumer.assignment().await.get("grows").map_or(0, Vec::len);
if assigned == 4 {
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert_eq!(
assigned, 4,
"partitions added to a subscribed topic must be assigned"
);
consumer.close().await.expect("close");
}
#[tokio::test]
async fn assign_overrides_a_standalone_subscription_for_that_topic() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("narrow", 4);
let consumer = crate::consumer::Consumer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.auto_offset_reset(crate::consumer::AutoOffsetReset::Earliest)
.metadata_max_age(Duration::from_millis(50))
.build()
.await
.expect("consumer should connect");
consumer.subscribe(&["narrow"]).await.expect("subscribe");
assert_eq!(consumer.assignment().await["narrow"].len(), 4);
consumer.assign("narrow", vec![0]).await.expect("assign");
for _ in 0..5 {
let _ = consumer.poll(Duration::from_millis(50)).await;
tokio::time::sleep(Duration::from_millis(20)).await;
assert_eq!(
consumer.assignment().await["narrow"],
vec![0],
"a manual assign() must not be widened by the standalone resolver"
);
}
consumer.close().await.expect("close");
}
#[tokio::test]
async fn auto_create_topics_lets_a_send_materialise_its_topic() {
let broker = FakeBroker::start().await.unwrap();
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.max_block(SHORT_MAX_BLOCK)
.allow_auto_create_topics(true)
.build()
.await
.expect("producer should connect");
let metadata = producer
.send("created-on-demand", None, Some(b"v"))
.await
.expect("the broker creates the topic because the client asked it to");
assert_eq!(metadata.topic, "created-on-demand");
producer.close().await;
}
#[tokio::test]
async fn auto_create_topics_is_off_by_default() {
let broker = FakeBroker::start().await.unwrap();
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.max_block(SHORT_MAX_BLOCK)
.build()
.await
.expect("producer should connect");
let error = producer
.send("not-created-on-demand", None, Some(b"v"))
.await
.expect_err("a typo must not materialise a topic");
assert!(
matches!(
error,
crate::error::KrafkaError::Broker {
code: ErrorCode::UnknownTopicOrPartition,
..
}
),
"expected the broker's own topic error, got: {error}"
);
producer.close().await;
}
#[tokio::test]
async fn an_out_of_range_partition_is_rejected_at_send() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 2);
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.max_block(SHORT_MAX_BLOCK)
.build()
.await
.expect("producer should connect");
let error = producer
.send_record(
crate::producer::ProducerRecord::new("events", b"v".to_vec()).with_partition(7),
)
.await
.expect_err("partition 7 does not exist on a 2-partition topic");
assert!(
error.to_string().contains("not in the range [0, 2)"),
"the error must name the valid range, got: {error}"
);
producer.close().await;
}
#[tokio::test]
async fn interceptor_state_survives_from_on_send_to_on_acknowledgement() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = producer_with(&broker, &interceptor).await;
for i in 0..3u8 {
let _ = producer
.send("events", None, Some(&[b'v', i]))
.await
.expect("send should be acknowledged");
}
let acks = interceptor.acks();
assert_eq!(acks.len(), 3, "every record must reach on_acknowledgement");
let tokens: Vec<_> = acks.iter().filter_map(|a| a.token.clone()).collect();
assert_eq!(
tokens,
vec!["events#0", "events#1", "events#2"],
"each record must get *its own* context back, in order"
);
for ack in acks.iter() {
assert!(!ack.failed, "a successful send reports no error");
assert_eq!(ack.delivery, crate::producer::DeliveryConfirmation::Offset);
assert!(ack.offset >= 0, "a successful send carries a real offset");
}
}
#[tokio::test]
async fn a_record_rejected_by_validation_still_reaches_on_acknowledgement() {
let broker = FakeBroker::start().await.unwrap();
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = producer_with(&broker, &interceptor).await;
let mut record = crate::producer::ProducerRecord::new("events", b"v".to_vec());
for i in 0..(crate::protocol::MAX_RECORD_HEADERS + 1) {
record = record.with_header(format!("h{i}"), bytes::Bytes::from_static(b"x"));
}
let error = producer
.send_record(record)
.await
.expect_err("a record over the header limit must be rejected");
assert!(error.to_string().contains("headers"));
let acks = interceptor.acks();
assert_eq!(interceptor.sends(), 1);
assert_eq!(acks.len(), 1, "the rejected record still owes a callback");
assert_eq!(
acks[0].token.as_deref(),
Some("events#0"),
"the context opened in on_send must come back, not be dropped"
);
assert!(acks[0].failed);
assert_eq!(acks[0].partition, crate::producer::UNKNOWN_PARTITION);
assert_eq!(acks[0].offset, -1);
assert_eq!(
acks[0].delivery,
crate::producer::DeliveryConfirmation::Failed
);
}
#[tokio::test]
async fn a_record_rejected_by_a_serializer_still_reaches_on_acknowledgement() {
#[derive(Debug)]
struct FailingSerializer;
impl crate::serdes::Serializer for FailingSerializer {
fn serialize(
&self,
_payload: bytes::Bytes,
_topic: &str,
_record_name: Option<&str>,
_is_key: bool,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = crate::error::Result<bytes::Bytes>> + Send + '_>,
> {
Box::pin(async {
Err(crate::error::KrafkaError::invalid_state(
"schema registry rejected the payload",
))
})
}
}
let broker = FakeBroker::start().await.unwrap();
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.add_interceptor(std::sync::Arc::clone(&interceptor) as std::sync::Arc<_>)
.value_serializer(std::sync::Arc::new(FailingSerializer))
.build()
.await
.expect("producer should connect");
let _ = producer
.send("events", None, Some(b"v"))
.await
.expect_err("the serializer must reject the record");
let acks = interceptor.acks();
assert_eq!(acks.len(), 1, "the rejected record still owes a callback");
assert_eq!(acks[0].token.as_deref(), Some("events#0"));
assert!(acks[0].failed);
assert_eq!(acks[0].partition, crate::producer::UNKNOWN_PARTITION);
}
#[tokio::test]
async fn a_dropped_delivery_handle_does_not_suppress_on_acknowledgement() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = producer_with(&broker, &interceptor).await;
for i in 0..3u8 {
drop(
producer
.enqueue(crate::producer::ProducerRecord::new(
"events",
vec![b'v', i],
))
.await
.expect("enqueue should succeed"),
);
}
producer
.flush()
.await
.expect("flush should drain the buffer");
let acks = interceptor.acks();
assert_eq!(
acks.len(),
3,
"dropping the handle must not cost the interceptor its callback"
);
for ack in acks.iter() {
assert!(ack.token.is_some(), "the context must survive the drop");
assert!(!ack.failed);
}
}
#[tokio::test]
async fn a_record_for_an_unknown_topic_still_reaches_on_acknowledgement() {
let broker = FakeBroker::start().await.unwrap();
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = producer_with(&broker, &interceptor).await;
let error = producer
.send("no-such-topic", None, Some(b"v"))
.await
.expect_err("an unrouteable record must be rejected");
assert!(
matches!(
error,
crate::error::KrafkaError::Broker {
code: ErrorCode::UnknownTopicOrPartition,
..
}
),
"expected the broker's own topic error, got: {error}"
);
assert!(
error.to_string().contains("no-such-topic"),
"the error must name the topic, got: {error}"
);
let acks = interceptor.acks();
assert_eq!(acks.len(), 1, "the rejected record still owes a callback");
assert_eq!(acks[0].token.as_deref(), Some("no-such-topic#0"));
assert!(acks[0].failed);
assert_eq!(acks[0].partition, crate::producer::UNKNOWN_PARTITION);
assert_eq!(
acks[0].delivery,
crate::producer::DeliveryConfirmation::Failed
);
}
#[tokio::test]
async fn a_cancelled_send_still_reaches_on_acknowledgement() {
let broker = FakeBroker::start().await.unwrap();
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.max_block(Duration::from_secs(30))
.add_interceptor(std::sync::Arc::clone(&interceptor) as std::sync::Arc<_>)
.build()
.await
.expect("producer should connect");
let cancelled = tokio::time::timeout(
Duration::from_millis(200),
producer.send("no-such-topic", None, Some(b"v")),
)
.await;
assert!(cancelled.is_err(), "the send must still be in flight");
{
let acks = interceptor.acks();
assert_eq!(
acks.len(),
1,
"a cancelled record still owes an acknowledgement"
);
assert_eq!(acks[0].token.as_deref(), Some("no-such-topic#0"));
assert!(acks[0].failed);
assert_eq!(acks[0].partition, crate::producer::UNKNOWN_PARTITION);
assert_eq!(
acks[0].delivery,
crate::producer::DeliveryConfirmation::Failed
);
}
producer.close().await;
}
#[tokio::test]
async fn a_send_cancelled_waiting_for_buffer_memory_still_acknowledges() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.buffer_memory(256)
.batch_size(256)
.linger(Duration::from_secs(30))
.max_block(Duration::from_secs(30))
.add_interceptor(std::sync::Arc::clone(&interceptor) as std::sync::Arc<_>)
.build()
.await
.expect("producer should connect");
let first = producer
.enqueue(crate::producer::ProducerRecord::new("events", vec![0u8; 100]).with_partition(0))
.await
.expect("the first record fits");
let cancelled = tokio::time::timeout(
Duration::from_millis(200),
producer.enqueue(
crate::producer::ProducerRecord::new("events", vec![1u8; 100]).with_partition(0),
),
)
.await;
assert!(
cancelled.is_err(),
"the second record must still be waiting for memory"
);
{
let acks = interceptor.acks();
let cancelled_ack = acks
.iter()
.find(|a| a.token.as_deref() == Some("events#1"))
.expect("the cancelled record still owes an acknowledgement");
assert!(cancelled_ack.failed);
}
drop(first);
producer.close().await;
}
#[tokio::test]
async fn a_split_batch_acknowledges_every_record_exactly_once_with_its_context() {
const RECORDS: usize = 4;
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
broker.on_once(ApiKey::Produce, |_| {
Control::Error(ErrorCode::MessageTooLarge)
});
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.linger(Duration::from_millis(200))
.add_interceptor(std::sync::Arc::clone(&interceptor) as std::sync::Arc<_>)
.build()
.await
.expect("producer should connect");
for i in 0..RECORDS {
drop(
producer
.enqueue(
crate::producer::ProducerRecord::new("events", vec![b'v', i as u8])
.with_partition(0),
)
.await
.expect("enqueue should succeed"),
);
}
producer
.flush()
.await
.expect("flush should drain the buffer");
assert!(
broker.request_count(ApiKey::Produce) >= 3,
"expected the rejected batch plus two halves, saw {}",
broker.request_count(ApiKey::Produce)
);
let acks = interceptor.acks();
assert_eq!(
acks.len(),
RECORDS,
"exactly one acknowledgement per record — no losses, no duplicates"
);
let mut tokens: Vec<_> = acks.iter().filter_map(|a| a.token.clone()).collect();
tokens.sort();
assert_eq!(
tokens,
vec!["events#0", "events#1", "events#2", "events#3"],
"every record must carry its own context across the split"
);
for ack in acks.iter() {
assert!(!ack.failed, "both halves should have been accepted");
}
}
#[tokio::test]
async fn the_transactional_send_path_pairs_on_send_with_on_acknowledgement() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = crate::producer::TransactionalProducer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.transactional_id("txn-interceptor")
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.max_block(SHORT_MAX_BLOCK)
.add_interceptor(std::sync::Arc::clone(&interceptor) as std::sync::Arc<_>)
.build()
.await
.expect("transactional producer should connect");
producer.init_transactions().await.expect("init");
producer.begin_transaction().expect("begin");
let _ = producer
.send("orders", None, Some(b"committed"))
.await
.expect("send");
let _ = producer
.send("no-such-topic", None, Some(b"unrouteable"))
.await
.expect_err("an unrouteable record must be rejected");
producer.commit_transaction().await.expect("commit");
producer.close().await;
let acks = interceptor.acks();
assert_eq!(acks.len(), 2, "both records owe an acknowledgement");
let committed = acks
.iter()
.find(|a| a.token.as_deref() == Some("orders#0"))
.expect("the committed record must report with its own context");
assert!(!committed.failed);
assert_eq!(
committed.delivery,
crate::producer::DeliveryConfirmation::Offset
);
let rejected = acks
.iter()
.find(|a| a.token.as_deref() == Some("no-such-topic#1"))
.expect("the rejected record must report with its own context");
assert!(rejected.failed);
assert_eq!(rejected.partition, crate::producer::UNKNOWN_PARTITION);
assert_eq!(
rejected.delivery,
crate::producer::DeliveryConfirmation::Failed
);
}
#[tokio::test]
async fn on_acknowledgement_sees_headers_written_later_in_the_chain() {
#[derive(Debug)]
struct LateHeaderInterceptor;
impl crate::interceptor::ProducerInterceptor for LateHeaderInterceptor {
fn on_send(
&self,
record: &mut crate::producer::ProducerRecord,
_ctx: &mut crate::interceptor::RecordContext,
) -> crate::interceptor::InterceptorResult {
record.headers.push((
"added-last".to_string(),
Some(bytes::Bytes::from_static(b"1")),
));
Ok(())
}
}
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("events", 1);
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = Producer::builder()
.bootstrap_servers(broker.bootstrap_servers())
.request_timeout(SHORT_REQUEST_TIMEOUT)
.connect_timeout(SHORT_CONNECT_TIMEOUT)
.add_interceptor(std::sync::Arc::clone(&interceptor) as std::sync::Arc<_>)
.add_interceptor(std::sync::Arc::new(LateHeaderInterceptor))
.build()
.await
.expect("producer should connect");
let _ = producer
.send_record(
crate::producer::ProducerRecord::new("events", b"v".to_vec())
.with_header("added-first", bytes::Bytes::from_static(b"0")),
)
.await
.expect("send should be acknowledged");
let acks = interceptor.acks();
assert_eq!(acks.len(), 1);
assert_eq!(
acks[0].header_keys,
vec!["added-first".to_string(), "added-last".to_string()],
"the acknowledgement must carry the final header set, not the one this \
interceptor saw in on_send"
);
}
#[tokio::test]
async fn a_rejected_record_still_reports_its_headers() {
let broker = FakeBroker::start().await.unwrap();
let interceptor = std::sync::Arc::new(RecordingInterceptor::default());
let producer = producer_with(&broker, &interceptor).await;
let _ = producer
.send_record(
crate::producer::ProducerRecord::new("no-such-topic", b"v".to_vec())
.with_header("trace-id", bytes::Bytes::from_static(b"abc")),
)
.await
.expect_err("an unrouteable record must be rejected");
let acks = interceptor.acks();
assert_eq!(acks.len(), 1);
assert_eq!(
acks[0].header_keys,
vec!["trace-id".to_string()],
"a pre-accumulator rejection must still report the record's headers"
);
assert!(acks[0].failed);
}