mod sim;
use barnabas_client::{Producer, ProducerRecord};
use kafka_protocol::messages::ApiKey;
use sim::{drive, journal, start, Fault, Sim};
const BOOTSTRAP: &str = "broker-1:9092";
fn bootstrap() -> Vec<String> {
vec![BOOTSTRAP.to_owned()]
}
fn records(n: usize) -> Vec<ProducerRecord> {
(0..n)
.map(|i| ProducerRecord::new(None, Some(bytes::Bytes::from(format!("v{i}")))))
.collect()
}
async fn transactional() -> Producer<Sim> {
Producer::<Sim>::transactional(Sim, &bootstrap(), "sim", "txn-1")
.await
.expect("producer")
}
#[test]
fn a_moved_leader_is_retried_without_duplicating() {
start(vec![Fault {
api: ApiKey::Produce,
code: 6, times: 1,
after: 0,
}]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
producer.send("t", 0, &records(3)).await.expect("send");
producer.commit_transaction().await.expect("commit");
});
let j = journal();
assert_eq!(
j.produced,
vec![(0, 0, 3)],
"the retry wrote a second batch, or the first was lost"
);
assert_eq!(j.count(ApiKey::Produce), 2, "expected one retry");
assert!(
j.count(ApiKey::Metadata) >= 2,
"a moved leader must trigger a metadata refresh, not a blind retry"
);
}
#[test]
fn repeated_retries_keep_the_same_sequence() {
start(vec![Fault {
api: ApiKey::Produce,
code: 6,
times: 5,
after: 0,
}]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
producer.send("t", 0, &records(2)).await.expect("send");
});
let j = journal();
assert_eq!(j.produced, vec![(0, 0, 2)]);
}
#[test]
fn a_moved_coordinator_is_rediscovered() {
start(vec![Fault {
api: ApiKey::InitProducerId,
code: 16, times: 1,
after: 0,
}]);
drive(async {
let _ = transactional().await;
});
let j = journal();
assert!(
j.count(ApiKey::FindCoordinator) >= 2,
"NOT_COORDINATOR was retried in place instead of re-discovering \
(FindCoordinator seen {} times)",
j.count(ApiKey::FindCoordinator)
);
}
#[test]
fn a_warming_coordinator_is_waited_out() {
start(vec![Fault {
api: ApiKey::InitProducerId,
code: 14, times: 3,
after: 0,
}]);
drive(async {
let _ = transactional().await;
});
assert_eq!(journal().count(ApiKey::InitProducerId), 4);
}
#[test]
fn concurrent_transactions_is_waited_out() {
start(vec![Fault {
api: ApiKey::AddPartitionsToTxn,
code: 51,
times: 2,
after: 0,
}]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
producer.send("t", 0, &records(1)).await.expect("send");
});
let j = journal();
assert_eq!(j.count(ApiKey::AddPartitionsToTxn), 3);
assert_eq!(j.produced, vec![(0, 0, 1)]);
}
#[test]
fn a_fenced_producer_refuses_everything_afterwards() {
start(vec![Fault {
api: ApiKey::Produce,
code: 90, times: 1,
after: 0,
}]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
let err = producer
.send("t", 0, &records(1))
.await
.expect_err("a fenced produce must fail");
assert!(
matches!(
err,
barnabas_client::Error::Broker {
disposition: barnabas_core::Disposition::Fatal,
..
}
),
"fencing must be fatal, got: {err}"
);
let again = producer.send("t", 0, &records(1)).await;
assert!(
matches!(again, Err(barnabas_client::Error::Producer(_))),
"a fenced producer accepted another send"
);
});
assert!(
journal().produced.is_empty(),
"a fenced producer wrote records"
);
}
#[test]
fn an_out_of_sequence_error_is_not_retried() {
start(vec![Fault {
api: ApiKey::Produce,
code: 45, times: 1,
after: 0,
}]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
let err = producer.send("t", 0, &records(1)).await.expect_err("fatal");
assert!(matches!(
err,
barnabas_client::Error::Broker {
disposition: barnabas_core::Disposition::Fatal,
..
}
));
});
assert_eq!(
journal().count(ApiKey::Produce),
1,
"a sequence error was retried; that is how duplicates are written"
);
}
#[test]
fn sequences_continue_across_transactions() {
start(vec![]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin 1");
producer.send("t", 0, &records(3)).await.expect("send 1");
producer.commit_transaction().await.expect("commit 1");
producer.begin_transaction().expect("begin 2");
producer.send("t", 0, &records(2)).await.expect("send 2");
producer.commit_transaction().await.expect("commit 2");
});
assert_eq!(
journal().produced,
vec![(0, 0, 3), (0, 3, 2)],
"the second transaction restarted the sequence"
);
}
#[test]
fn a_partition_is_enrolled_once_per_transaction() {
start(vec![]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
for _ in 0..3 {
producer.send("t", 0, &records(1)).await.expect("send");
}
producer.commit_transaction().await.expect("commit");
});
let j = journal();
assert_eq!(j.count(ApiKey::AddPartitionsToTxn), 1);
assert_eq!(j.count(ApiKey::Produce), 3);
}
#[test]
fn a_closed_connection_is_reconnected() {
sim::start_with(
Vec::new(),
vec![sim::Close {
api: ApiKey::Produce,
times: 1,
}],
Vec::new(),
);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
producer
.send("t", 0, &records(2))
.await
.expect("a closed connection must be reconnected, not fatal");
});
let j = journal();
assert_eq!(
j.produced,
vec![(0, 0, 2)],
"the reconnect wrote twice, or lost the batch"
);
assert_eq!(
j.count(ApiKey::Produce),
2,
"expected exactly one reconnect"
);
}
#[test]
fn a_reconnect_does_not_renumber() {
sim::start_with(
Vec::new(),
vec![sim::Close {
api: ApiKey::Produce,
times: 1,
}],
Vec::new(),
);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
producer.send("t", 0, &records(3)).await.expect("first");
producer.send("t", 0, &records(2)).await.expect("second");
});
assert_eq!(
journal().produced,
vec![(0, 0, 3), (0, 3, 2)],
"sequences restarted across a reconnect"
);
}
#[test]
fn a_hung_broker_times_out_rather_than_hanging() {
sim::start_with(
Vec::new(),
Vec::new(),
vec![sim::Hang {
api: ApiKey::Produce,
times: 1,
}],
);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
let err = producer
.send("t", 0, &records(1))
.await
.expect_err("a hung broker must time out");
assert!(
matches!(err, barnabas_client::Error::Timeout { .. }),
"expected a timeout, got: {err}"
);
});
assert!(
journal().produced.is_empty(),
"nothing was accepted, so nothing may be reported as written"
);
}
#[test]
fn one_produce_request_covers_every_partition_on_a_broker() {
start(vec![]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
let records: Vec<ProducerRecord> = (0..40)
.map(|i| {
ProducerRecord::new(
Some(bytes::Bytes::from(format!("key-{i}"))),
Some(bytes::Bytes::from(format!("v{i}"))),
)
})
.collect();
let written = producer.send_keyed("t", &records).await.expect("send");
assert!(
written.len() > 1,
"the test needs keys spread over several partitions, got {written:?}"
);
});
let j = journal();
assert_eq!(
j.produce_widths.len(),
1,
"expected one batched request, got {} ({:?})",
j.produce_widths.len(),
j.produce_widths
);
assert!(
j.produce_widths[0] > 1,
"the single request carried only one partition: {:?}",
j.produce_widths
);
assert_eq!(
j.produced.len(),
j.produce_widths[0],
"every partition in the request must be accounted for"
);
}
#[test]
fn partitions_are_enrolled_in_one_request() {
start(vec![]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
let records: Vec<ProducerRecord> = (0..40)
.map(|i| {
ProducerRecord::new(
Some(bytes::Bytes::from(format!("key-{i}"))),
Some(bytes::Bytes::from(format!("v{i}"))),
)
})
.collect();
producer.send_keyed("t", &records).await.expect("send");
});
let j = journal();
assert_eq!(
j.count(ApiKey::AddPartitionsToTxn),
1,
"enrollment was sent per partition instead of once"
);
assert!(
j.enroll_widths[0] > 1,
"the enrollment carried one partition: {:?}",
j.enroll_widths
);
}
#[test]
fn only_the_failed_partition_is_retried() {
start(vec![Fault {
api: ApiKey::Produce,
code: 6, times: 1,
after: 0,
}]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
let records: Vec<ProducerRecord> = (0..40)
.map(|i| {
ProducerRecord::new(
Some(bytes::Bytes::from(format!("key-{i}"))),
Some(bytes::Bytes::from(format!("v{i}"))),
)
})
.collect();
producer.send_keyed("t", &records).await.expect("send");
});
let j = journal();
let mut partitions: Vec<i32> = j.produced.iter().map(|(p, _, _)| *p).collect();
partitions.sort_unstable();
let unique: std::collections::BTreeSet<i32> = partitions.iter().copied().collect();
assert_eq!(
partitions.len(),
unique.len(),
"a partition was written twice: {partitions:?}"
);
}
#[test]
fn a_window_of_batches_arrives_in_sequence() {
sim::start(Vec::new());
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
for _ in 0..5 {
producer
.enqueue("t", 0, &records(2))
.await
.expect("enqueue");
}
assert_eq!(producer.queued(), 5);
let written = producer.flush().await.expect("flush");
assert_eq!(written.len(), 5, "every batch must be acknowledged");
assert_eq!(producer.queued(), 0);
let journal = journal();
assert_eq!(
journal.produced,
vec![(0, 0, 2), (0, 2, 2), (0, 4, 2), (0, 6, 2), (0, 8, 2)],
"sequences must be contiguous and in order"
);
assert_eq!(
journal.produce_widths.len(),
5,
"expected five requests, got {:?}",
journal.produce_widths
);
});
}
#[test]
fn the_window_is_bounded_by_max_in_flight() {
sim::start(Vec::new());
drive(async {
let mut producer = transactional().await;
producer.set_max_in_flight(3);
producer.begin_transaction().expect("begin");
for _ in 0..10 {
producer
.enqueue("t", 0, &records(1))
.await
.expect("enqueue");
}
producer.flush().await.expect("flush");
let journal = journal();
assert_eq!(journal.produced.len(), 10);
for (i, (_, base, _)) in journal.produced.iter().enumerate() {
assert_eq!(*base, i as i32, "sequences must be 0..10 in order");
}
});
}
#[test]
fn a_failure_mid_window_resends_everything_behind_it() {
sim::start(vec![Fault {
api: ApiKey::Produce,
code: 6, times: 1,
after: 0,
}]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
for _ in 0..4 {
producer
.enqueue("t", 0, &records(2))
.await
.expect("enqueue");
}
let written = producer.flush().await.expect("flush");
assert_eq!(written.len(), 4, "every batch must land exactly once");
let journal = journal();
assert_eq!(
journal.produced,
vec![(0, 0, 2), (0, 2, 2), (0, 4, 2), (0, 6, 2)],
"the log must read in order with no gap and no duplicate"
);
assert!(
!matches!(producer.state(), barnabas_core::TxnState::Fatal),
"an out-of-sequence answer caused by an earlier failure must not fence"
);
});
}
#[test]
fn an_unexplained_sequence_error_still_fences() {
sim::start(vec![Fault {
api: ApiKey::Produce,
code: 45, times: 1,
after: 0,
}]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
producer
.enqueue("t", 0, &records(1))
.await
.expect("enqueue");
let err = producer.flush().await.expect_err("must not be tolerated");
assert!(
matches!(
err,
barnabas_client::Error::Broker {
disposition: barnabas_core::Disposition::Fatal,
..
}
),
"expected a fatal sequence error, got: {err}"
);
assert!(matches!(producer.state(), barnabas_core::TxnState::Fatal));
});
}
#[test]
fn a_success_after_a_failure_is_not_retired() {
sim::start(vec![Fault {
api: ApiKey::Produce,
code: 6, times: 1,
after: 1, }]);
drive(async {
let mut producer = transactional().await;
producer.begin_transaction().expect("begin");
for _ in 0..4 {
producer
.enqueue("t", 0, &records(2))
.await
.expect("enqueue");
}
let written = producer.flush().await.expect("flush");
assert_eq!(written.len(), 4, "every batch must land exactly once");
let journal = journal();
assert_eq!(
journal.produced,
vec![(0, 0, 2), (0, 2, 2), (0, 4, 2), (0, 6, 2)],
"the log must read in order with no gap and no duplicate"
);
});
}
#[test]
fn a_call_on_a_busy_connection_is_refused() {
sim::start(Vec::new());
drive(async {
let mut cluster = barnabas_client::Cluster::connect(Sim, &bootstrap(), "sim")
.await
.expect("cluster");
cluster
.send_at_for_test(ApiKey::Metadata, 12, BOOTSTRAP, &metadata_request())
.await
.expect("send");
let err = cluster
.call_at_for_test::<_, kafka_protocol::messages::MetadataResponse>(
BOOTSTRAP,
ApiKey::Metadata,
12,
&metadata_request(),
)
.await
.expect_err("a busy connection must refuse");
assert!(
matches!(
err,
barnabas_client::Error::ConnectionBusy { in_flight: 1, .. }
),
"expected ConnectionBusy, got: {err}"
);
});
}
fn metadata_request() -> kafka_protocol::messages::MetadataRequest {
let mut req = kafka_protocol::messages::MetadataRequest::default();
req.topics = Some(vec![]);
req
}