use std::sync::Arc;
use std::time::Duration;
use tracing::info;
use crate::client::Kafka;
use crate::error::Result;
use crate::metadata::ClusterMetadata;
use crate::metrics::{ClientInstanceId, Metrics, MetricsSource};
use crate::network::ConnectionPool;
use crate::telemetry::{ClientType, Telemetry};
pub use crate::consumer::TopicPartition;
pub use crate::protocol::{
AclBinding, AclBindingFilter, AclOperation, AclPatternType, AclPermissionType, AclResourceType,
ConfigResourceType, DescribedStreamsGroup, ElectionType, FeatureUpdateKey, FeatureUpgradeType,
FinalizedFeature, ListedConfigResource, ScramCredentialDeletion, ScramCredentialUpsertion,
StreamsAssignment, StreamsEndpoint, StreamsGroupMember, StreamsKeyValue, StreamsSubtopology,
StreamsTaskIds, StreamsTaskOffset, StreamsTopicInfo, StreamsTopology, SupportedFeature,
};
macro_rules! admin_options {
(
$(#[$meta:meta])*
$name:ident {
$( $(#[$field_meta:meta])* $field:ident : $ty:ty ),* $(,)?
}
$( optional {
$( $(#[$opt_meta:meta])* $opt:ident : $opt_ty:ty ),* $(,)?
} )?
) => {
$(#[$meta])*
#[derive(Debug, Clone, Default)]
#[must_use]
pub struct $name {
timeout: Option<std::time::Duration>,
$( $field: $ty, )*
$( $( $opt: Option<$opt_ty>, )* )?
}
impl $name {
pub fn timeout(mut self, timeout: std::time::Duration) -> Self {
self.timeout = Some(timeout);
self
}
$(
$(#[$field_meta])*
pub fn $field(mut self, value: $ty) -> Self {
self.$field = value;
self
}
)*
$( $(
$(#[$opt_meta])*
pub fn $opt(mut self, value: $opt_ty) -> Self {
self.$opt = Some(value);
self
}
)* )?
}
};
}
mod acls;
mod configs;
mod driver;
mod features;
mod group_offsets;
mod groups;
mod offsets;
mod partitions;
mod quotas;
mod scram;
mod share_group_offsets;
mod streams_groups;
mod tokens;
mod topics;
mod transactions;
pub use acls::{
AclFilter, CreateAclsOptions, DeleteAclsOptions, DeleteAclsResult, DescribeAclsOptions,
};
pub use configs::{
ClusterBroker, ClusterDescription, ConfigEntry, ConfigOp, ConfigParseError, ConfigResource,
ConfigSynonymEntry, ConfigValue, DescribeClusterOptions, DescribeConfigsOptions,
IncrementalAlterConfigsOptions, ListConfigResourcesOptions,
};
pub use features::{DescribeFeaturesOptions, FeatureMetadata, UpdateFeaturesOptions};
pub use group_offsets::{
AlterConsumerGroupOffsetsOptions, ConsumerGroupLag, ConsumerGroupLagOptions,
DeleteConsumerGroupOffsetsOptions, GroupOffset, ListConsumerGroupOffsetsOptions,
};
pub use groups::{
ConsumerGroupDescription, ConsumerGroupListing, ConsumerGroupMember,
DeleteConsumerGroupsOptions, DescribeConsumerGroupsOptions, GroupType,
ListConsumerGroupsOptions, TopicPartitionAssignment,
};
pub use offsets::{
DeleteRecordsOptions, EpochEndOffset, ListOffsetsOptions, ListedOffset,
OffsetForLeaderEpochOptions, OffsetSpec,
};
pub use partitions::{
AlterPartitionReassignmentsOptions, AlterReplicaLogDirsOptions, DescribeLogDirsOptions,
ElectLeadersOptions, ListPartitionReassignmentsOptions, LogDirDescription, LogDirReplica,
PartitionReassignment, TopicPartitionReplica,
};
pub use quotas::{
AlterClientQuotasOptions, ClientQuotaAlteration, ClientQuotaEntity, ClientQuotaFilter,
DescribeClientQuotasOptions, QuotaMatch,
};
pub use scram::{
AlterUserScramCredentialsOptions, DescribeUserScramCredentialsOptions, ScramCredentialInfo,
};
pub use share_group_offsets::{
AlterShareGroupOffsetsOptions, DeleteShareGroupOffsetsOptions,
DescribeShareGroupOffsetsOptions, SharePartitionOffset,
};
pub use streams_groups::DescribeStreamsGroupsOptions;
pub use tokens::{
CreateDelegationTokenOptions, DelegationToken, DelegationTokenPrincipal,
DescribeDelegationTokenOptions, ExpireDelegationTokenOptions, RenewDelegationTokenOptions,
};
pub use topics::{
CreatePartitionsOptions, CreateTopicsOptions, DeleteTopicsOptions, DescribeTopicsOptions,
ListTopicsOptions, NewTopic, PartitionDescription, TopicDescription,
};
pub use transactions::{
AbortTransactionOptions, DescribeMetadataQuorumOptions, DescribeProducersOptions,
DescribeTransactionsOptions, ListTransactionsOptions, ProducerState, QuorumInfo,
QuorumListener, QuorumNode, QuorumReplica, TransactionDescription, TransactionListing,
};
const DEFAULT_API_TIMEOUT: Duration = Duration::from_secs(60);
#[derive(Debug, Clone)]
pub(crate) struct AdminConfig {
pub(crate) default_api_timeout: Duration,
pub(crate) retry_backoff: crate::util::BackoffPolicy,
pub(crate) metrics_push: bool,
}
impl Default for AdminConfig {
fn default() -> Self {
Self {
default_api_timeout: DEFAULT_API_TIMEOUT,
retry_backoff: crate::util::BackoffPolicy {
initial_backoff: Duration::from_millis(100),
max_backoff: Duration::from_secs(1),
backoff_multiplier: 2.0,
jitter_factor: 0.2,
},
metrics_push: false,
}
}
}
pub struct AdminClient {
config: AdminConfig,
kafka: Kafka,
metadata: Arc<ClusterMetadata>,
pool: Arc<ConnectionPool>,
closed: std::sync::atomic::AtomicBool,
metrics_source: Arc<MetricsSource>,
telemetry: Telemetry,
}
impl std::fmt::Debug for AdminClient {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AdminClient")
.field("config", &self.config)
.field("closed", &self.is_closed())
.finish_non_exhaustive()
}
}
impl AdminClient {
pub(crate) fn new(kafka: Kafka) -> Self {
Self {
config: AdminConfig::default(),
metadata: Arc::clone(kafka.metadata()),
pool: Arc::clone(kafka.pool()),
closed: std::sync::atomic::AtomicBool::new(false),
metrics_source: MetricsSource::admin(&kafka),
telemetry: Telemetry::disabled(),
kafka,
}
}
pub fn default_api_timeout(mut self, timeout: Duration) -> Self {
self.config.default_api_timeout = timeout;
self
}
pub fn retry_backoff(mut self, backoff: Duration) -> Self {
self.config.retry_backoff.initial_backoff = backoff;
self
}
pub fn metrics_push(mut self, enable: bool) -> Self {
self.config.metrics_push = enable;
self.telemetry = Telemetry::start(
enable,
&self.kafka,
ClientType::Admin,
Arc::clone(&self.metrics_source),
);
self
}
pub async fn close(&self) -> Result<()> {
if !self.closed.swap(true, std::sync::atomic::Ordering::SeqCst) {
self.telemetry.close(self.config.default_api_timeout).await;
info!("admin client closed");
}
Ok(())
}
#[inline]
pub fn is_closed(&self) -> bool {
self.closed.load(std::sync::atomic::Ordering::SeqCst)
}
pub fn metrics(&self) -> Metrics {
self.metrics_source.snapshot()
}
pub async fn client_instance_id(&self, timeout: Duration) -> Result<Option<ClientInstanceId>> {
self.telemetry.client_instance_id(timeout).await
}
}
fn validate_topics<'a>(names: impl IntoIterator<Item = &'a str>) -> Result<()> {
crate::protocol::validate_topic_names(names)
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::error::KrafkaError;
#[test]
fn test_admin_client_is_send_sync() {
fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<AdminClient>();
}
#[test]
fn admin_retry_backoff_is_exponential_and_jittered() {
let policy = AdminConfig::default().retry_backoff;
assert_eq!(policy.initial_backoff, Duration::from_millis(100));
assert_eq!(policy.max_backoff, Duration::from_secs(1));
assert!(policy.backoff_multiplier > 1.0);
assert!(policy.jitter_factor() > 0.0);
}
#[tokio::test]
async fn admin_settings_reach_the_config() {
let admin = Kafka::detached()
.admin()
.default_api_timeout(Duration::from_secs(5))
.retry_backoff(Duration::from_millis(7));
assert_eq!(admin.config.default_api_timeout, Duration::from_secs(5));
assert_eq!(
admin.config.retry_backoff.initial_backoff,
Duration::from_millis(7)
);
}
#[tokio::test]
async fn close_is_idempotent_and_operations_fail_fast_after_it() {
let client = Kafka::detached().admin();
client.close().await.unwrap();
client.close().await.unwrap();
assert!(client.is_closed());
let err = client
.list_topics(ListTopicsOptions::default())
.await
.unwrap_err();
assert!(matches!(err, KrafkaError::Closed { .. }), "got: {err:?}");
}
}