#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
use std::time::Duration;
use crate::Kafka;
use crate::consumer::{AutoOffsetReset, Consumer};
use crate::testing::FakeBroker;
async fn connect(broker: &FakeBroker, client_id: &str) -> Kafka {
Kafka::builder(broker.bootstrap_servers())
.client_id(client_id)
.request_timeout(Duration::from_secs(2))
.connect_timeout(Duration::from_secs(1))
.connect()
.await
.unwrap()
}
async fn poll_n(consumer: &Consumer, n: usize) -> usize {
let deadline = tokio::time::Instant::now() + Duration::from_secs(10);
let mut got = 0;
while got < n && tokio::time::Instant::now() < deadline {
got += consumer
.poll(Duration::from_millis(200))
.await
.unwrap()
.len();
}
got
}
fn sample(text: &str, series: &str) -> Option<u64> {
text.lines()
.find_map(|l| l.strip_prefix(series)?.strip_prefix(' ')?.parse().ok())
}
#[tokio::test]
async fn what_the_clients_did_reaches_their_metrics_and_the_handle_sum() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let kafka = connect(&broker, "orders-service").await;
let producer = kafka.producer().build().await.unwrap();
for i in 0..100u32 {
let _ = producer
.send(crate::Record::new("orders", i.to_be_bytes().to_vec()))
.await
.unwrap();
}
let metrics = producer.metrics();
assert_eq!(metrics.producer.records_sent, 100);
assert!(metrics.connections.connections_created >= 1);
assert_eq!(metrics.client_id.as_deref(), Some("orders-service"));
let text = metrics.prometheus_text();
assert_eq!(
sample(
&text,
"krafka_producer_records_sent_total{client_id=\"orders-service\"}"
),
Some(100)
);
assert_eq!(
sample(
&text,
"krafka_connections_created_total{client_id=\"orders-service\"}"
),
Some(metrics.connections.connections_created)
);
for i in 0..100u32 {
let _ = producer
.send(crate::Record::new("orders", i.to_be_bytes().to_vec()))
.await
.unwrap();
}
assert_eq!(metrics.producer.records_sent, 100);
assert_eq!(producer.metrics().producer.records_sent, 200);
let consumer = kafka
.consumer("g")
.auto_offset_reset(AutoOffsetReset::Earliest)
.build()
.await
.unwrap();
consumer.subscribe(["orders"]).await.unwrap();
assert_eq!(poll_n(&consumer, 200).await, 200);
assert_eq!(consumer.metrics().consumer.records_received, 200);
let total = kafka.metrics();
assert_eq!(total.client_id, None);
assert_eq!(total.producer.records_sent, 200);
assert_eq!(total.consumer.records_received, 200);
assert_eq!(total.connections, kafka.admin().metrics().connections);
assert_eq!(
total.connections.connections_created,
consumer.metrics().connections.connections_created
);
let text = total.prometheus_text();
assert_eq!(
sample(&text, "krafka_producer_records_sent_total"),
Some(200)
);
assert!(!text.contains("client_id"));
consumer.close().await.unwrap();
producer.close().await.unwrap();
}
#[tokio::test]
async fn every_client_reports_its_own_activity() {
let broker = FakeBroker::start().await.unwrap();
broker.create_topic("orders", 1);
let kafka = connect(&broker, "svc").await;
let txn = kafka.producer().build_transactional("t").await.unwrap();
txn.begin().unwrap();
let _ = txn.send(crate::Record::new("orders", "v")).await.unwrap();
txn.commit().await.unwrap();
assert_eq!(txn.metrics().producer.records_sent, 1);
let share = kafka.share_consumer("s").build().await.unwrap();
let _ = share.poll(Duration::from_millis(20)).await;
assert_eq!(share.metrics().consumer.polls, 1);
assert_eq!(share.metrics().producer.records_sent, 0);
let admin = kafka.admin();
assert!(admin.metrics().connections.connections_created >= 1);
assert_eq!(admin.metrics().producer, Default::default());
assert_eq!(kafka.metrics().producer.records_sent, 1);
txn.close().await.unwrap();
drop(txn);
assert_eq!(kafka.metrics().producer.records_sent, 0);
share.close().await.unwrap();
}