#![allow(clippy::let_underscore_untyped)]
#![allow(clippy::expect_used, clippy::unwrap_used)]
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();
let _: fn() -> krafka::error::Result<krafka::producer::TransactionalProducerConfig> =
|| krafka::producer::TransactionalProducer::builder().build_config();
}
#[test]
fn every_sasl_mechanism_is_reachable_over_both_transports() {
use krafka::auth::{AuthConfig, SaslMechanism, SecurityProtocol, TlsConfig};
fn cleartext_configs() -> Vec<(SaslMechanism, AuthConfig)> {
vec![
(
SaslMechanism::Plain,
AuthConfig::sasl_plain("user", "pass").expect("valid PLAIN credentials"),
),
(
SaslMechanism::ScramSha256,
AuthConfig::sasl_scram_sha256("user", "pass"),
),
(
SaslMechanism::ScramSha512,
AuthConfig::sasl_scram_sha512("user", "pass"),
),
(
SaslMechanism::OAuthBearer,
AuthConfig::sasl_oauthbearer("jwt"),
),
]
}
for (mechanism, cleartext) in cleartext_configs() {
assert_eq!(
cleartext.security_protocol(),
&SecurityProtocol::SaslPlaintext,
"{mechanism} must be constructible over SASL_PLAINTEXT"
);
assert_eq!(cleartext.sasl_mechanism(), Some(&mechanism));
assert!(cleartext.tls_config().is_none());
let encrypted = cleartext.with_tls(TlsConfig::new().with_ca_cert("/etc/kafka/ca.pem"));
assert_eq!(
encrypted.security_protocol(),
&SecurityProtocol::SaslSsl,
"{mechanism} must be constructible over SASL_SSL"
);
assert_eq!(
encrypted.sasl_mechanism(),
Some(&mechanism),
"{mechanism} must survive the TLS upgrade"
);
assert_eq!(
encrypted.tls_config().and_then(|t| t.ca_cert_path()),
Some("/etc/kafka/ca.pem"),
"{mechanism} must carry the caller's TLS settings, not defaults"
);
}
let msk = AuthConfig::aws_msk_iam("AKID", "secret", "us-east-1")
.with_tls(TlsConfig::new().with_sni_hostname("b-1.msk.example.com"));
assert_eq!(msk.security_protocol(), &SecurityProtocol::SaslSsl);
assert_eq!(msk.sasl_mechanism(), Some(&SaslMechanism::AwsMskIam));
assert_eq!(
msk.tls_config().and_then(|t| t.sni_hostname()),
Some("b-1.msk.example.com")
);
assert_eq!(
AuthConfig::sasl_scram_sha256_ssl("u", "p", TlsConfig::new()).security_protocol(),
AuthConfig::sasl_scram_sha256("u", "p")
.with_tls(TlsConfig::new())
.security_protocol()
);
assert_eq!(
AuthConfig::sasl_scram_sha512_ssl("u", "p", TlsConfig::new()).security_protocol(),
AuthConfig::sasl_scram_sha512("u", "p")
.with_tls(TlsConfig::new())
.security_protocol()
);
assert_eq!(
AuthConfig::plaintext()
.with_tls(TlsConfig::new())
.security_protocol(),
&SecurityProtocol::Ssl
);
}
#[test]
fn both_producer_builders_share_one_configuration_surface() {
use krafka::producer::{Producer, ProducerRecord, TransactionalProducer};
use krafka::protocol::Compression;
use std::sync::Arc;
use std::time::Duration;
#[derive(Debug)]
struct NoopInterceptor;
impl krafka::interceptor::ProducerInterceptor for NoopInterceptor {
fn on_send(&self, _record: &mut ProducerRecord) -> krafka::interceptor::InterceptorResult {
Ok(())
}
}
#[derive(Debug)]
struct NoopDlq;
impl krafka::dlq::DeadLetterQueue for NoopDlq {
fn send<'a>(
&'a self,
_record: ProducerRecord,
_error: String,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
Box::pin(async {})
}
}
struct NoopStore;
impl krafka::producer::ProducerStateStore for NoopStore {
async fn load(
&self,
) -> krafka::error::Result<Option<krafka::producer::ProducerIdentitySnapshot>> {
Ok(None)
}
async fn store(
&self,
_snapshot: &krafka::producer::ProducerIdentitySnapshot,
) -> krafka::error::Result<()> {
Ok(())
}
}
macro_rules! assert_producer_setters {
($make:expr) => {{
let _ = |d: Duration| $make.linger(d);
let _ = |n: usize| $make.batch_size(n);
let _ = |n: usize| $make.buffer_memory(n);
let _ = |d: Duration| $make.max_block(d);
let _ = |n: usize| $make.max_request_size(n);
let _ = |n: u32| $make.retries(n);
let _ = |d: Duration| $make.retry_backoff(d);
let _ = |d: Duration| $make.delivery_timeout(d);
let _ = |c: Compression| $make.compression(c);
let _ = |l: Option<i32>| $make.compression_level(l);
let _ = |t: &str, c: Compression| $make.topic_compression(t, c);
let _ = |d: Duration| $make.metadata_topic_cache_ttl(d);
let _ = || $make.disable_metadata_topic_cache_ttl();
let _ = |d: Duration| $make.metadata_recovery_rebootstrap_trigger(d);
let _ = |q: Arc<NoopDlq>| $make.dead_letter_queue(q);
let _ = |i: Arc<NoopInterceptor>| $make.interceptor(i);
let _ = |i: Arc<NoopInterceptor>| $make.add_interceptor(i);
let _ = || $make.state_store(NoopStore);
let _ = |c: &krafka::client::KrafkaClient| $make.with_client(c);
let _ = |p: krafka::producer::UniformStickyPartitioner| $make.partitioner(p);
let _ = |e: Arc<dyn krafka::serdes::Serializer>| $make.key_serializer(e);
let _ = |e: Arc<dyn krafka::serdes::Serializer>| $make.value_serializer(e);
let _ = || {
$make.sasl_oauthbearer_provider(|| async {
Ok(krafka::auth::OAuthBearerToken::new("jwt"))
})
};
#[cfg(feature = "socks5")]
let _ = |p: krafka::network::ProxyConfig| $make.proxy(p);
}};
}
assert_producer_setters!(Producer::builder());
assert_producer_setters!(TransactionalProducer::builder());
}
#[test]
fn both_producers_share_one_operational_surface() {
macro_rules! assert_producer_ops {
($ty:ty) => {{
async fn _flush(p: &$ty) -> krafka::error::Result<()> {
p.flush().await
}
async fn _close(p: &$ty) {
p.close().await;
}
async fn _close_timeout(p: &$ty, d: std::time::Duration) -> krafka::error::Result<()> {
p.close_with_timeout(d).await
}
fn _metrics(p: &$ty) -> krafka::producer::ProducerMetricsSnapshot {
p.metrics()
}
}};
}
assert_producer_ops!(krafka::producer::Producer);
assert_producer_ops!(krafka::producer::TransactionalProducer);
}
#[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);
}};
}
macro_rules! assert_shares_a_client {
($make:expr) => {{
let _ = |c: &krafka::client::KrafkaClient| $make.with_client(c);
}};
}
assert_shares_a_client!(krafka::consumer::Consumer::builder());
assert_shares_a_client!(krafka::producer::Producer::builder());
assert_shares_a_client!(krafka::admin::AdminClient::builder());
assert_shares_a_client!(krafka::producer::TransactionalProducer::builder());
#[cfg(feature = "unstable-protocol")]
assert_shares_a_client!(krafka::share_consumer::ShareConsumer::builder());
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());
}
#[cfg(feature = "unstable-protocol")]
#[test]
fn both_consumers_share_the_read_side_surface() {
use std::sync::Arc;
use std::time::Duration;
macro_rules! assert_consumer_read_surface {
($make:expr) => {{
let _ = |d: Duration| $make.fetch_max_wait(d);
let _ = |n: i32| $make.fetch_min_bytes(n);
let _ = |n: i32| $make.fetch_max_bytes(n);
let _ = |n: i32| $make.max_poll_records(n);
let _ = |n: i32| $make.max_buffered_records(n);
let _ = |d: Arc<dyn krafka::serdes::Deserializer>| $make.key_deserializer(d);
let _ = |d: Arc<dyn krafka::serdes::Deserializer>| $make.value_deserializer(d);
}};
}
assert_consumer_read_surface!(krafka::consumer::Consumer::builder());
assert_consumer_read_surface!(krafka::share_consumer::ShareConsumer::builder());
let _ = |n: i32| krafka::share_consumer::ShareConsumer::builder().max_records(n);
let _ = |n: i32| krafka::share_consumer::ShareConsumer::builder().batch_size(n);
}
#[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()
}
fn _owns_pool(c: &$ty) -> bool {
c.owns_pool()
}
}};
}
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();
}
}