#![allow(dead_code)]
mod client;
mod connection;
mod error;
mod model;
mod producer;
pub(crate) mod profile;
pub(crate) mod protocol;
mod security;
mod source;
pub(crate) use client::{
AutoOffsetReset, FetchTuning, GroupAssignor, NativeCommitPolicy, NativeKafkaConsumerConfig,
};
pub(crate) use error::{KafkaClientError, KafkaClientResult};
pub(crate) use model::{KafkaPayloadBatch, KafkaTimestamp, StartOffset, TopicPartitionAssignment};
pub(crate) use producer::{NativeKafkaProducerControl, NativeKafkaProducerHandle};
pub(crate) use security::{NativeSecurityConfig, native_security_config};
pub(crate) use source::{NativeKafkaControl, NativeKafkaSource};
pub(crate) const VERSION: &str = env!("CARGO_PKG_VERSION");
#[cfg(feature = "bench-internals")]
#[doc(hidden)]
#[must_use]
pub fn bench_encode_produce_batch(records: &[crate::ProducerRecord], idempotent: bool) -> Vec<u8> {
let refs = records.iter().collect::<Vec<_>>();
let identity = idempotent.then_some(protocol::BatchIdentity {
producer_id: 4242,
producer_epoch: 0,
base_sequence: 0,
});
let mut zstd_context =
protocol::ZstdEncoderContext::new().expect("benchmark zstd encoder context");
protocol::encode_produce_record_batch(
&refs,
protocol::CompressionCodec::None,
&mut zstd_context,
identity,
)
.expect("benchmark record batch encodes")
}
#[cfg(all(test, feature = "rdkafka"))]
pub(crate) use client::NativeKafkaConsumer;
#[cfg(test)]
pub(crate) use model::TopicPartition;
#[cfg(all(test, feature = "rdkafka"))]
mod tests;