#![allow(clippy::let_underscore_untyped)]
use krafka::network::TransportConfig;
#[test]
fn every_client_builder_accepts_a_transport_config() {
let _ = |t: TransportConfig| krafka::consumer::Consumer::builder().transport(t);
let _ = |t: TransportConfig| krafka::producer::Producer::builder().transport(t);
let _ = |t: TransportConfig| krafka::admin::AdminClient::builder().transport(t);
let _ = |t: TransportConfig| krafka::producer::TransactionalProducer::builder().transport(t);
let _ =
|t: TransportConfig| krafka::client::KrafkaClient::builder("localhost:9092").transport(t);
#[cfg(feature = "unstable-protocol")]
let _ = |t: TransportConfig| krafka::share_consumer::ShareConsumer::builder().transport(t);
}
#[test]
fn client_builders_expose_a_synchronous_validation_terminal() {
let _: fn() -> krafka::error::Result<krafka::consumer::ConsumerConfig> =
|| krafka::consumer::Consumer::builder().build_config();
let _: fn() -> krafka::error::Result<krafka::producer::ProducerConfig> =
|| krafka::producer::Producer::builder().build_config();
let _: fn() -> krafka::error::Result<krafka::admin::AdminConfig> =
|| krafka::admin::AdminClient::builder().build_config();
}
#[test]
fn tls_rotation_is_reachable_from_every_client() {
async fn _consumer(c: &krafka::consumer::Consumer) -> krafka::error::Result<()> {
c.refresh_tls().await
}
async fn _producer(p: &krafka::producer::Producer) -> krafka::error::Result<()> {
p.refresh_tls().await
}
async fn _admin(a: &krafka::admin::AdminClient) -> krafka::error::Result<()> {
a.refresh_tls().await
}
async fn _client(k: &krafka::client::KrafkaClient) -> krafka::error::Result<()> {
k.refresh_tls().await
}
}
#[test]
fn metrics_are_readable_without_an_async_context() {
fn _producer(p: &krafka::producer::Producer) {
let _snapshot: krafka::producer::ProducerMetricsSnapshot = p.metrics();
let _conn: std::sync::Arc<krafka::metrics::ConnectionMetrics> = p.connection_metrics();
}
fn _consumer(c: &krafka::consumer::Consumer) {
let _m: &std::sync::Arc<krafka::metrics::ConsumerMetrics> = c.metrics();
let _conn: std::sync::Arc<krafka::metrics::ConnectionMetrics> = c.connection_metrics();
}
fn _admin(a: &krafka::admin::AdminClient) {
let _conn: std::sync::Arc<krafka::metrics::ConnectionMetrics> = a.connection_metrics();
}
}
#[test]
fn api_version_negotiation_is_synchronous() {
fn _negotiate(conn: &krafka::network::BrokerConnection) {
let _: Option<i16> = conn.negotiate_api_version(krafka::protocol::ApiKey::Fetch, 18, 4);
let _: Option<i16> = conn.negotiate_api_version_max(krafka::protocol::ApiKey::Produce, 13);
let _: Option<krafka::protocol::ApiVersionRange> =
conn.get_api_version(krafka::protocol::ApiKey::Metadata);
}
}
#[test]
fn client_builders_share_one_configuration_surface() {
use krafka::metadata::MetadataRecoveryStrategy;
use std::time::Duration;
macro_rules! assert_common_setters {
($make:expr) => {{
let _ = |v: &str| $make.client_id(v);
let _ = |d: Duration| $make.request_timeout(d);
let _ = |d: Duration| $make.connect_timeout(d);
let _ = |d: Duration| $make.metadata_max_age(d);
let _ = |s: MetadataRecoveryStrategy| $make.metadata_recovery_strategy(s);
let _ = |t: TransportConfig| $make.transport(t);
let _ = |a: krafka::auth::AuthConfig| $make.auth(a);
let _ = |u: &str, p: &str| $make.sasl_plain(u, p);
let _ = |u: &str, p: &str| $make.sasl_scram_sha256(u, p);
let _ = |u: &str, p: &str| $make.sasl_scram_sha512(u, p);
let _ = |t: &str| $make.sasl_oauthbearer(t);
}};
}
assert_common_setters!(krafka::consumer::Consumer::builder());
assert_common_setters!(krafka::producer::Producer::builder());
assert_common_setters!(krafka::admin::AdminClient::builder());
assert_common_setters!(krafka::producer::TransactionalProducer::builder());
assert_common_setters!(krafka::client::KrafkaClient::builder("localhost:9092"));
#[cfg(feature = "unstable-protocol")]
assert_common_setters!(krafka::share_consumer::ShareConsumer::builder());
}
#[test]
fn every_long_lived_client_shares_one_operational_surface() {
macro_rules! assert_lifecycle {
($ty:ty) => {{
async fn _refresh(c: &$ty) -> krafka::error::Result<()> {
c.refresh_tls().await
}
async fn _rebootstrap(c: &$ty) {
c.rebootstrap().await;
}
fn _seeds(c: &$ty, s: Vec<String>) -> krafka::error::Result<()> {
c.update_seed_brokers(s)
}
fn _conn_metrics(c: &$ty) -> std::sync::Arc<krafka::metrics::ConnectionMetrics> {
c.connection_metrics()
}
fn _closed(c: &$ty) -> bool {
c.is_closed()
}
}};
}
assert_lifecycle!(krafka::producer::Producer);
assert_lifecycle!(krafka::consumer::Consumer);
assert_lifecycle!(krafka::admin::AdminClient);
assert_lifecycle!(krafka::producer::TransactionalProducer);
#[cfg(feature = "unstable-protocol")]
assert_lifecycle!(krafka::share_consumer::ShareConsumer);
fn _consumer_wakeup(c: &krafka::consumer::Consumer) {
c.wakeup();
}
#[cfg(feature = "unstable-protocol")]
fn _share_wakeup(c: &krafka::share_consumer::ShareConsumer) {
c.wakeup();
}
fn _txn(p: &krafka::producer::TransactionalProducer) {
let _: krafka::producer::ProducerMetricsSnapshot = p.metrics();
}
#[cfg(feature = "unstable-protocol")]
fn _share(c: &krafka::share_consumer::ShareConsumer) {
let _: std::sync::Arc<krafka::metrics::ConsumerMetrics> = c.metrics();
}
}