pub struct AdminClient { /* private fields */ }Expand description
Kafka admin client for cluster administration.
Implementations§
Source§impl AdminClient
impl AdminClient
Sourcepub async fn describe_acls(
&self,
filter: AclFilter,
) -> Result<DescribeAclsResult>
pub async fn describe_acls( &self, filter: AclFilter, ) -> Result<DescribeAclsResult>
Sourcepub async fn create_acls(
&self,
acls: Vec<AclBinding>,
) -> Result<CreateAclsResult>
pub async fn create_acls( &self, acls: Vec<AclBinding>, ) -> Result<CreateAclsResult>
Create ACLs.
Returns Ok(result) when the RPC succeeds. An Ok return does not
mean every ACL was created — inspect each element of
CreateAclsResult::results for per-ACL failures.
§Arguments
acls- List of ACL bindings to create
§Example
let acl = AclBinding::allow_read_topic("my-topic", "User:alice");
admin.create_acls(vec![acl]).await?;Sourcepub async fn delete_acls(
&self,
filters: Vec<AclBindingFilter>,
) -> Result<DeleteAclsResult>
pub async fn delete_acls( &self, filters: Vec<AclBindingFilter>, ) -> Result<DeleteAclsResult>
Delete ACLs matching the specified filters.
Returns Ok(result) when the RPC succeeds. An Ok return does not
mean every filter matched or every ACL was deleted — inspect each
element of DeleteAclsResult for per-filter failures.
§Arguments
filters- List of ACL binding filters to match for deletion
§Example
// Delete all ACLs for a specific topic
let filter = AclBindingFilter {
resource_type: AclResourceType::Topic,
resource_name: Some("my-topic".to_string()),
pattern_type: AclPatternType::Literal,
principal: None,
host: None,
operation: AclOperation::Any,
permission_type: AclPermissionType::Any,
};
admin.delete_acls(vec![filter]).await?;Source§impl AdminClient
impl AdminClient
Sourcepub async fn describe_configs_per_resource(
&self,
request: DescribeConfigsRequest,
) -> Result<Vec<DescribeConfigsResourceResult>>
pub async fn describe_configs_per_resource( &self, request: DescribeConfigsRequest, ) -> Result<Vec<DescribeConfigsResourceResult>>
Describe configuration for one or more resources (topics, brokers, etc.), preserving per-resource errors and attribution.
Uses DescribeConfigs (API Key 32). Build a DescribeConfigsRequest
via its convenience constructors (for_topic, for_broker) or manually
populate the resources field for multi-resource queries.
Each returned DescribeConfigsResourceResult carries the resource it
describes and that resource’s error, if any. Use this rather than
describe_configs whenever more than one
resource is requested, or whenever an authorization failure must be
distinguishable from “this resource has no configs”.
Sourcepub async fn describe_configs(
&self,
request: DescribeConfigsRequest,
) -> Result<Vec<ConfigEntry>>
pub async fn describe_configs( &self, request: DescribeConfigsRequest, ) -> Result<Vec<ConfigEntry>>
Describe configuration for one or more resources, flattening every resource’s entries into a single list.
§Errors
Unlike the previous behaviour, a per-resource error is not silently
swallowed: if any requested resource failed (for example with
TOPIC_AUTHORIZATION_FAILED), this returns that error rather than an
empty or partial list that is indistinguishable from “no configs”.
Use describe_configs_per_resource
to inspect partial results across a multi-resource request.
Sourcepub async fn alter_topic_config(
&self,
topic: &str,
configs: HashMap<String, String>,
) -> Result<AlterConfigResult>
pub async fn alter_topic_config( &self, topic: &str, configs: HashMap<String, String>, ) -> Result<AlterConfigResult>
Alter configuration for a topic.
Uses IncrementalAlterConfigs (API Key 44) to set individual config keys without replacing the entire config. Each key-value pair is applied as a SET operation.
IncrementalAlterConfigs is a controller-only API: the request is
routed to the current controller and re-issued against a freshly
resolved controller on NOT_CONTROLLER.
Sourcepub async fn list_topics(&self) -> Result<Vec<String>>
pub async fn list_topics(&self) -> Result<Vec<String>>
List all topics.
Forces a metadata refresh first, so this is safe to call immediately
after create_topics: the refresh waits out the
retry.backoff.ms rate limiter rather than returning a stale cache as
though it were current.
Sourcepub async fn topic_config(
&self,
topic: &str,
keys: &[&str],
) -> Result<HashMap<String, ConfigValue>>
pub async fn topic_config( &self, topic: &str, keys: &[&str], ) -> Result<HashMap<String, ConfigValue>>
Fetch specific config keys for a topic.
A convenience wrapper around describe_configs
for the common case of reading a small set of well-known topic-level keys.
Pass an empty keys slice to fetch all config entries for the topic.
Returns a map of config key → ConfigValue.
§Examples
let cfg = admin.topic_config("my-topic", &["retention.ms", "cleanup.policy"]).await?;
if let Some(krafka::admin::ConfigValue::Value(v)) = cfg.get("retention.ms") {
println!("retention.ms = {v}");
}Sourcepub async fn describe_topics<S: AsRef<str>>(
&self,
topics: &[S],
) -> Result<HashMap<String, Arc<TopicInfo>>>
pub async fn describe_topics<S: AsRef<str>>( &self, topics: &[S], ) -> Result<HashMap<String, Arc<TopicInfo>>>
Describe topics by name, returning a map of topic name → [TopicInfo].
Topics not found in cluster metadata are absent from the returned map; callers can detect missing topics by comparing request and response key sets.
Sourcepub async fn describe_topic(
&self,
topic: &str,
) -> Result<Option<Arc<TopicInfo>>>
pub async fn describe_topic( &self, topic: &str, ) -> Result<Option<Arc<TopicInfo>>>
Describe a single topic by name.
Returns None if the topic does not exist or is not visible in cluster metadata.
Sourcepub async fn describe_cluster(&self) -> Result<DescribeClusterResult>
pub async fn describe_cluster(&self) -> Result<DescribeClusterResult>
Describe the cluster using the DescribeCluster API (Key 60).
Returns cluster metadata including cluster ID, controller, brokers, and authorized operations.
Sourcepub async fn partition_count(&self, topic: &str) -> Result<Option<usize>>
pub async fn partition_count(&self, topic: &str) -> Result<Option<usize>>
Get partition count for a topic.
Sourcepub fn request_timeout(&self) -> Duration
pub fn request_timeout(&self) -> Duration
Get the request timeout.
Source§impl AdminClient
impl AdminClient
Sourcepub async fn describe_features(&self) -> Result<DescribeFeaturesResult>
pub async fn describe_features(&self) -> Result<DescribeFeaturesResult>
Describe broker-supported and cluster-finalized features (KIP-584).
Sends an ApiVersions request (v3+) to any broker and extracts the
feature information from the tagged fields. The response includes:
- Features supported by the responding broker (per-broker)
- Cluster-wide finalized features and their epoch (cluster-wide)
§Example
let features = admin.describe_features().await?;
for f in &features.supported_features {
println!("{}: v{}–v{}", f.name, f.min_version, f.max_version);
}
for f in &features.finalized_features {
println!("{}: v{}–v{} (finalized)", f.name, f.min_version_level, f.max_version_level);
}Sourcepub async fn update_features(
&self,
feature_updates: Vec<FeatureUpdateKey>,
validate_only: bool,
) -> Result<UpdateFeaturesResult>
pub async fn update_features( &self, feature_updates: Vec<FeatureUpdateKey>, validate_only: bool, ) -> Result<UpdateFeaturesResult>
Update cluster-wide finalized feature version levels (KIP-584).
This is a destructive operation — downgrades and deletions can be
data-lossy. Only the controller serves this request, so the client
resolves the controller from cluster metadata and sends it there
directly, re-resolving and retrying on NOT_CONTROLLER.
Requires ALTER permission on the cluster.
§Per-feature results are version-dependent
UpdateFeatures v2 (Kafka 4.0+) removed the per-feature Results array
from the response: the controller now answers with a single top-level
outcome. Against a v2 broker,
UpdateFeaturesResult::results is therefore empty, and a
successful return means every requested update was applied. Against a
v0/v1 broker the per-feature breakdown is still populated. Treat
Ok(_) — not a non-empty results — as the success signal, and use
results only as extra diagnostics when it happens to be present.
§Example
use krafka::protocol::FeatureUpdateKey;
// `Ok` means the update was applied. `results` is a v0/v1-only detail.
let results = admin.update_features(
vec![FeatureUpdateKey::upgrade("metadata.version", 17)],
false, // validate_only
).await?;
for result in &results.results {
if let Some(e) = &result.error {
eprintln!("Failed to update {}: {e}", result.feature);
}
}Source§impl AdminClient
impl AdminClient
Sourcepub async fn delete_consumer_group_offsets(
&self,
group_id: &str,
topic_partitions: &[(&str, &[i32])],
) -> Result<OffsetDeleteResult>
pub async fn delete_consumer_group_offsets( &self, group_id: &str, topic_partitions: &[(&str, &[i32])], ) -> Result<OffsetDeleteResult>
Delete committed offsets for a consumer group.
This is a destructive operation — deleted offsets cannot be
recovered. The consumer group must be in the Empty state.
The request is sent to the group coordinator.
§Example
let results = admin.delete_offsets(
"my-group",
&[("my-topic", &[0, 1, 2])],
).await?;Sourcepub async fn describe_consumer_group_offsets(
&self,
group_id: &str,
topic_partitions: Option<&[(&str, &[i32])]>,
visibility: OffsetVisibility,
) -> Result<Vec<GroupOffsetEntry>>
pub async fn describe_consumer_group_offsets( &self, group_id: &str, topic_partitions: Option<&[(&str, &[i32])]>, visibility: OffsetVisibility, ) -> Result<Vec<GroupOffsetEntry>>
Fetch committed offsets for a consumer group.
Pass topic_partitions to fetch offsets for specific partitions,
or None to fetch all committed offsets for the group.
The request is sent to the group coordinator.
visibility decides what to do about an offset a transaction has
staged but not yet committed — see OffsetVisibility. It is an
explicit argument rather than a default because the two answers differ
on exactly the pipelines where the difference matters, and a lag
dashboard silently reading the unstable value reports progress that an
abort can take back.
§Example
let offsets = admin
.describe_consumer_group_offsets("my-group", None, OffsetVisibility::IncludeUnstable)
.await?;
for entry in &offsets {
println!("{}/{}: {}", entry.topic, entry.partition, entry.committed_offset);
}§Errors
With OffsetVisibility::StableOnly, a partition whose offset is
staged inside an unresolved transaction is reported by the broker as
UNSTABLE_OFFSET_COMMIT; it is surfaced rather than omitted, since an
omitted partition is indistinguishable from one the group never
committed.
Sourcepub async fn alter_consumer_group_offsets(
&self,
group_id: &str,
topic_offsets: &[(&str, &[(i32, i64)])],
) -> Result<Vec<AlterGroupOffsetResult>>
pub async fn alter_consumer_group_offsets( &self, group_id: &str, topic_offsets: &[(&str, &[(i32, i64)])], ) -> Result<Vec<AlterGroupOffsetResult>>
Alter committed offsets for a consumer group.
Sets each specified partition’s committed offset. The consumer group
must be in the Empty state (no active members).
The request is sent to the group coordinator.
§Example
admin
.alter_consumer_group_offsets(
"my-group",
&[("my-topic", &[(0, 100), (1, 200)])],
)
.await?;Source§impl AdminClient
impl AdminClient
Sourcepub async fn describe_consumer_groups(
&self,
group_ids: Vec<String>,
) -> Result<Vec<ConsumerGroupDescription>>
pub async fn describe_consumer_groups( &self, group_ids: Vec<String>, ) -> Result<Vec<ConsumerGroupDescription>>
Describe consumer groups.
Automatically detects whether each group uses the classic protocol or the new consumer protocol (KIP-848) and dispatches to the appropriate API:
- Classic groups → DescribeGroups (Key 15)
- Consumer groups → ConsumerGroupDescribe (Key 69)
The returned ConsumerGroupDescription is a unified type.
Fields specific to one protocol variant are Option-wrapped.
§Example
let groups = admin
.describe_consumer_groups(vec!["my-group".to_string()])
.await?;
for group in &groups {
println!("{}: type={}, state={}, members={}",
group.group_id, group.group_type, group.state, group.members.len());
}Sourcepub async fn list_consumer_groups(
&self,
filter: &GroupListing,
) -> Result<Vec<ConsumerGroupListing>>
pub async fn list_consumer_groups( &self, filter: &GroupListing, ) -> Result<Vec<ConsumerGroupListing>>
List all consumer groups on the cluster.
Returns a list of all consumer groups with their protocol types.
§Example
let groups = admin.list_consumer_groups(&GroupListing::all()).await?;
for group in &groups {
println!("{} ({})", group.group_id, group.protocol_type);
}filter is applied by the broker. On a cluster with thousands of
groups that is the difference between transferring the whole registry
and transferring the handful you asked about — see GroupListing.
Sourcepub async fn delete_records(
&self,
offsets: HashMap<(String, i32), i64>,
timeout: Duration,
) -> Result<Vec<DeleteRecordResult>>
pub async fn delete_records( &self, offsets: HashMap<(String, i32), i64>, timeout: Duration, ) -> Result<Vec<DeleteRecordResult>>
Delete records from topic partitions before the specified offsets.
Records with offsets less than the specified offset for each partition will be marked for deletion. This adjusts the log start offset.
§Arguments
offsets- Map of (topic, partition) to the offset before which to deletetimeout- Operation timeout
§Example
use std::collections::HashMap;
let mut offsets = HashMap::new();
offsets.insert(("my-topic".to_string(), 0), 100i64);
let results = admin.delete_records(offsets, Duration::from_secs(30)).await?;Sourcepub async fn offset_for_leader_epoch(
&self,
partitions: Vec<(String, i32, i32)>,
) -> Result<Vec<LeaderEpochResult>>
pub async fn offset_for_leader_epoch( &self, partitions: Vec<(String, i32, i32)>, ) -> Result<Vec<LeaderEpochResult>>
Get the end offset for each partition at the given leader epoch.
This is used to detect log truncation after a leader change. For each topic-partition, the broker returns the end offset for the requested leader epoch. If the epoch is no longer valid, the broker returns the epoch and offset where the log was truncated.
§Arguments
partitions- List of (topic, partition, leader_epoch) tuples
§Example
let results = admin.offset_for_leader_epoch(
vec![("my-topic".to_string(), 0, 5)]
).await?;
for r in &results {
println!("{}:{} epoch={} end_offset={}", r.topic, r.partition, r.leader_epoch, r.end_offset);
}Source§impl AdminClient
impl AdminClient
Sourcepub async fn list_offsets(
&self,
topic_partitions: &[(&str, &[i32])],
spec: OffsetSpec,
) -> Result<Vec<ListOffsetResult>>
pub async fn list_offsets( &self, topic_partitions: &[(&str, &[i32])], spec: OffsetSpec, ) -> Result<Vec<ListOffsetResult>>
List offsets for one or more topic-partitions.
Each request is routed to the partition’s current leader. Metadata
is refreshed once on NotLeaderForPartition errors before retrying.
§Arguments
topic_partitions— slice of(topic_name, partition_ids)pairs.spec— which offset to fetch (Earliest,Latest, orTimestamp).
§Example
use krafka::admin::{AdminClient, OffsetSpec};
let results = admin
.list_offsets(&[("my-topic", &[0, 1, 2])], OffsetSpec::Latest)
.await?;
for r in &results {
println!("{}/{}: offset={}", r.topic, r.partition, r.offset);
}Sourcepub async fn consumer_group_lag(
&self,
group_id: &str,
topic_partitions: Option<&[(&str, &[i32])]>,
) -> Result<Vec<ConsumerGroupLag>>
pub async fn consumer_group_lag( &self, group_id: &str, topic_partitions: Option<&[(&str, &[i32])]>, ) -> Result<Vec<ConsumerGroupLag>>
Compute consumer group lag for the specified topics.
Lag is defined as end_offset − committed_offset for each
topic-partition. Partitions with no committed offset have
committed_offset = None and lag = None.
Partitions whose end offset could not be fetched report
end_offset = None, lag = None, and the reason in
ConsumerGroupLag::end_offset_error — never lag = 0, which would
make a stalled consumer look healthy.
This method issues two parallel-ish requests:
describe_consumer_group_offsetsfor the committed positions.list_offsetswithOffsetSpec::Latestfor the end offsets.
The consumer group does not need to be stopped.
§Arguments
group_id— consumer group ID.topic_partitions— which partitions to measure; passNoneto measure all partitions that the group has committed offsets for.
§Example
let lag = admin
.consumer_group_lag("my-group", Some(&[("my-topic", &[0, 1, 2])]))
.await?;
for entry in &lag {
println!(
"{}/{}: lag={:?}",
entry.topic, entry.partition, entry.lag
);
}Source§impl AdminClient
impl AdminClient
Sourcepub async fn describe_log_dirs(
&self,
topics: Option<Vec<DescribableLogDirTopic>>,
) -> Result<Vec<LogDirInfo>>
pub async fn describe_log_dirs( &self, topics: Option<Vec<DescribableLogDirTopic>>, ) -> Result<Vec<LogDirInfo>>
Describe log directories on all known brokers.
Each broker maintains one or more log directories; this method queries every broker and returns per-directory information including sizes, partition assignments, and (v4+) volume capacity.
Pass None for topics to describe all partitions on every
broker, or pass a list of [DescribableLogDirTopic] to filter.
§Example
// Describe all log dirs on every broker
let dirs = admin.describe_log_dirs(None).await?;
for dir in &dirs {
println!("broker {} dir {}: {:?}", dir.broker_id, dir.log_dir, dir.error);
}
// Describe specific topic partitions
use krafka::protocol::DescribableLogDirTopic;
let filter = vec![DescribableLogDirTopic {
topic: "my-topic".into(),
partitions: vec![0, 1, 2],
}];
let dirs = admin.describe_log_dirs(Some(filter)).await?;Sourcepub async fn elect_leaders(
&self,
election_type: ElectionType,
topic_partitions: Option<Vec<ElectLeadersTopicPartitions>>,
timeout: Duration,
) -> Result<Vec<ElectLeadersResult>>
pub async fn elect_leaders( &self, election_type: ElectionType, topic_partitions: Option<Vec<ElectLeadersTopicPartitions>>, timeout: Duration, ) -> Result<Vec<ElectLeadersResult>>
Trigger leader election for the specified partitions.
When topic_partitions is None, leaders for all partitions are
elected. The election_type controls whether to perform a preferred
or unclean leader election (requires broker v1+; v0 always does
preferred election).
Returns per-partition results — individual partitions may fail even when the RPC succeeds.
§Example
use krafka::protocol::ElectionType;
use std::time::Duration;
// Preferred election for all partitions
let results = admin
.elect_leaders(ElectionType::Preferred, None, Duration::from_secs(60))
.await?;Sourcepub async fn alter_partition_reassignments(
&self,
topics: Vec<ReassignableTopic>,
timeout: Duration,
) -> Result<AlterReassignmentsResult>
pub async fn alter_partition_reassignments( &self, topics: Vec<ReassignableTopic>, timeout: Duration, ) -> Result<AlterReassignmentsResult>
Alter partition reassignments.
Initiates or cancels partition reassignments. To cancel a pending
reassignment, set replicas to None for that partition.
This is a destructive operation — reassigning partitions moves data between brokers and can significantly impact cluster load.
Returns per-partition results — individual partitions may fail even when the RPC succeeds.
§Example
use krafka::protocol::{ReassignableTopic, ReassignablePartition};
use std::time::Duration;
let results = admin.alter_partition_reassignments(
vec![ReassignableTopic {
name: "my-topic".into(),
partitions: vec![ReassignablePartition {
partition_index: 0,
replicas: Some(vec![1, 2, 3]),
}],
}],
Duration::from_secs(60),
).await?;Sourcepub async fn alter_partition_reassignments_opts(
&self,
topics: Vec<ReassignableTopic>,
timeout: Duration,
allow_replication_factor_change: bool,
) -> Result<AlterReassignmentsResult>
pub async fn alter_partition_reassignments_opts( &self, topics: Vec<ReassignableTopic>, timeout: Duration, allow_replication_factor_change: bool, ) -> Result<AlterReassignmentsResult>
alter_partition_reassignments
with control over whether the move may change a partition’s replication
factor (AlterPartitionReassignments v1, Kafka 4.1).
Passing false asks the broker to reject any target replica set whose
size differs from the partition’s current replica count, answering with
INVALID_REPLICATION_FACTOR instead of applying the change. That is the
safer default for generated reassignment plans, where an off-by-one
would otherwise silently reduce durability.
§Version requirement
The field only exists from v1. Against a v0 broker
allow_replication_factor_change = false cannot be honoured, so
this returns ProtocolErrorKind::UnknownApiVersion rather than
quietly sending a request that permits the change. Passing true — the
broker default and the v0 behaviour — works against every version.
Sourcepub async fn list_partition_reassignments(
&self,
topics: Option<Vec<ListPartitionReassignmentsTopic>>,
timeout: Duration,
) -> Result<Vec<PartitionReassignmentInfo>>
pub async fn list_partition_reassignments( &self, topics: Option<Vec<ListPartitionReassignmentsTopic>>, timeout: Duration, ) -> Result<Vec<PartitionReassignmentInfo>>
List ongoing partition reassignments.
When topics is None, all ongoing reassignments are listed.
Otherwise, only the specified topic-partitions are checked.
§Example
// List all ongoing reassignments
let reassignments = admin
.list_partition_reassignments(None, Duration::from_secs(60))
.await?;
for topic in &reassignments {
for p in &topic.partitions {
println!("{} p{}: adding {:?}, removing {:?}",
topic.name, p.partition_index, p.adding_replicas, p.removing_replicas);
}
}Sourcepub async fn alter_replica_log_dirs(
&self,
dirs: Vec<AlterReplicaLogDir>,
) -> Result<Vec<AlterReplicaLogDirsResult>>
pub async fn alter_replica_log_dirs( &self, dirs: Vec<AlterReplicaLogDir>, ) -> Result<Vec<AlterReplicaLogDirsResult>>
Move partition replicas to a different log directory on the broker.
This is a destructive operation — moving replicas between log directories triggers data copying and can impact broker I/O.
This is a per-broker operation. Each broker is sent only the topic-partitions it actually hosts a replica for, resolved from cached cluster metadata. Brokers that host none of the requested replicas are not contacted at all.
This matters: broadcasting the full request to every broker asks each
one to move replicas it does not own, producing a storm of
REPLICA_NOT_AVAILABLE / KAFKA_STORAGE_ERROR results that are
indistinguishable from genuine failures on the brokers that do own them.
If metadata knows no replicas for a requested partition, the request is sent to the partition’s leader when known, and reported as an error otherwise — it is never broadcast.
Returns per-partition results — individual partitions may fail even when the RPC succeeds.
§Example
use krafka::protocol::{AlterReplicaLogDir, AlterReplicaLogDirTopic};
let results = admin.alter_replica_log_dirs(vec![
AlterReplicaLogDir {
path: "/data/kafka-logs-2".into(),
topics: vec![AlterReplicaLogDirTopic {
name: "my-topic".into(),
partitions: vec![0, 1],
}],
},
]).await?;Source§impl AdminClient
impl AdminClient
Sourcepub async fn describe_client_quotas(
&self,
components: &[(&str, i8, Option<&str>)],
strict: bool,
) -> Result<DescribeClientQuotasResult>
pub async fn describe_client_quotas( &self, components: &[(&str, i8, Option<&str>)], strict: bool, ) -> Result<DescribeClientQuotasResult>
Describe client quotas matching the given filter.
§Arguments
components- Filter components. Each component specifies an entity type and match criteria. The broker returns entities matching all components.strict- Iftrue, exclude entities with unspecified entity types (i.e., only return entities that exactly match all given component types).
§Filter Match Types
Each component has a match_type:
0(exact): match the entity with the given name1(default): match the default entity for this type2(any specified): match any entity with a name (non-default)
§Example
// Describe all quotas for user "alice"
let results = admin.describe_client_quotas(
&[("user", 0, Some("alice"))],
false,
).await?;Sourcepub async fn alter_client_quotas(
&self,
entries: &[QuotaAlteration<'_>],
validate_only: bool,
) -> Result<Vec<AlterClientQuotaResult>>
pub async fn alter_client_quotas( &self, entries: &[QuotaAlteration<'_>], validate_only: bool, ) -> Result<Vec<AlterClientQuotaResult>>
Alter client quotas.
Each entry specifies an entity (user, client-id, ip) and a set of quota operations (set or remove). Results are returned per-entity.
§Arguments
entries- Quota alterations. Each entry has an entity and operations.validate_only- Iftrue, validate the request without applying changes.
AlterClientQuotas is a controller-only API: the request is routed
to the current controller and re-issued against a freshly resolved
controller on NOT_CONTROLLER.
§Example
use krafka::admin::QuotaAlteration;
// Set producer byte rate quota for user "alice"
let results = admin.alter_client_quotas(
&[QuotaAlteration {
entity: vec![("user", Some("alice"))],
ops: vec![("producer_byte_rate", Some(1_048_576.0))],
}],
false,
).await?;Source§impl AdminClient
impl AdminClient
Sourcepub async fn describe_user_scram_credentials(
&self,
users: Option<Vec<String>>,
) -> Result<DescribeUserScramCredentialsResult>
pub async fn describe_user_scram_credentials( &self, users: Option<Vec<String>>, ) -> Result<DescribeUserScramCredentialsResult>
Sourcepub async fn alter_user_scram_credentials(
&self,
deletions: Vec<ScramCredentialDeletion>,
upsertions: Vec<ScramCredentialUpsertion>,
) -> Result<Vec<AlterScramCredentialResult>>
pub async fn alter_user_scram_credentials( &self, deletions: Vec<ScramCredentialDeletion>, upsertions: Vec<ScramCredentialUpsertion>, ) -> Result<Vec<AlterScramCredentialResult>>
Alter (upsert or delete) SCRAM credentials for users.
This is a destructive operation — deleting a SCRAM credential removes the user’s ability to authenticate with that mechanism.
AlterUserScramCredentials is a controller-only API: the request is
routed to the current controller and re-issued against a freshly
resolved controller on NOT_CONTROLLER.
§Example
use krafka::protocol::{ScramCredentialDeletion, ScramCredentialUpsertion};
use krafka::auth::ScramMechanism;
use zeroize::Zeroizing;
let results = admin.alter_user_scram_credentials(
vec![ScramCredentialDeletion {
name: "alice".into(),
mechanism: ScramMechanism::Sha512,
}],
vec![ScramCredentialUpsertion {
name: "bob".into(),
mechanism: ScramMechanism::Sha256,
iterations: 8192,
salt: Zeroizing::new(vec![1, 2, 3]),
salted_password: Zeroizing::new(vec![4, 5, 6]),
}],
).await?;Source§impl AdminClient
impl AdminClient
Read a share group’s share-partition start offsets (KIP-932).
Pass topics = None to describe every topic-partition the group
holds state for. Passing Some(&[]) describes nothing — the wire
protocol distinguishes a null topics array from an empty one, and so
does this method.
lag is Some(_) only when the
coordinator supports DescribeShareGroupOffsets v1 (KIP-1226, Kafka
4.3+); against an older broker it is None rather than a misleading
zero.
The request goes to the group coordinator, which is where share-group state lives.
§Example
// Everything the group knows about.
let all = admin.describe_share_group_offsets("orders-share", None).await?;
for p in &all.partitions {
println!("{}-{} start={} lag={:?}", p.topic, p.partition, p.start_offset, p.lag);
}
// Just two partitions of one topic.
let some = admin
.describe_share_group_offsets("orders-share", Some(&[("orders", &[0, 1][..])]))
.await?;§Errors
Returns an error if the client is closed, a topic name is invalid, the
coordinator cannot be found, or the broker supports no compatible
DescribeShareGroupOffsets version.
Reset a share group’s share-partition start offsets (KIP-932).
This is a destructive operation. Moving the start offset backwards re-delivers records the group already processed; moving it forwards skips records permanently.
The group must be empty. A group with a live member is answered
with NON_EMPTY_GROUP, for the same reason
alter_consumer_group_offsets
requires it: rewriting the start offset under an active member would
hand it records it has already acquired.
§Example
// Rewind two partitions to the beginning of the log.
let results = admin
.alter_share_group_offsets("orders-share", &[("orders", &[(0, 0), (1, 0)][..])])
.await?;
for r in &results {
if let Some(e) = &r.error {
eprintln!("{}-{} failed: {e}", r.topic, r.partition);
}
}§Errors
Returns an error if the client is closed, a topic name is invalid, the
coordinator cannot be found, the broker supports no compatible
AlterShareGroupOffsets version, or the request fails at the top level
(including NON_EMPTY_GROUP).
Delete a share group’s offset state for whole topics (KIP-932).
This is a destructive operation — the deleted state cannot be recovered, and the group restarts those topics from its configured reset policy. The group must be empty.
Use it after retiring a topic, so the coordinator stops carrying state for partitions that no longer exist.
§Errors
Returns an error if the client is closed, a topic name is invalid, the
coordinator cannot be found, the broker supports no compatible
DeleteShareGroupOffsets version, or the request fails at the top level
(including NON_EMPTY_GROUP).
Source§impl AdminClient
impl AdminClient
Sourcepub async fn describe_streams_groups(
&self,
group_ids: &[&str],
) -> Result<Vec<DescribedStreamsGroup>>
pub async fn describe_streams_groups( &self, group_ids: &[&str], ) -> Result<Vec<DescribedStreamsGroup>>
Describe Streams groups (KIP-1071, Kafka 4.1+).
Returns each group’s topology, members, and per-member task assignments and changelog offsets. Like every group API, the request goes to each group’s coordinator, so groups spread across brokers are batched per coordinator and issued once per broker.
§What this is for
krafka has no Streams runtime and cannot join a Streams group — that is
StreamsGroupHeartbeat (key 88), whose request carries the application
topology. This is the observational half: it is what an operator needs
to answer “is this Streams application healthy?” without running one.
Two fields are the ones worth alerting on:
- A member whose
topology_epochis below the group’sStreamsTopology::epochis still running an older topology. - A member whose
assignmentdiffers from itstarget_assignmenthas not finished rebalancing. Persistently so usually means restoration is not keeping up — comparetask_offsetsagainsttask_end_offsetsfor the lag that explains it.
§Errors
Per-group failures are reported in each
DescribedStreamsGroup::error_code
rather than failing the whole call, so one unknown group does not hide
the others. A broker that does not support the API at all — anything
before Kafka 4.1 — fails the call with
ProtocolErrorKind::UnknownApiVersion.
§Example
let groups = admin.describe_streams_groups(&["my-streams-app"]).await?;
for group in &groups {
if !group.error_code.is_ok() {
eprintln!("{}: {:?}", group.group_id, group.error_code);
continue;
}
let topology_epoch = group.topology.as_ref().map_or(-1, |t| t.epoch);
for member in &group.members {
let lagging = member.topology_epoch < topology_epoch;
let rebalancing = member.assignment != member.target_assignment;
println!(
"{} active={} lagging_topology={lagging} rebalancing={rebalancing}",
member.member_id,
member.assignment.active_tasks.len(),
);
}
}Source§impl AdminClient
impl AdminClient
Sourcepub async fn create_delegation_token(
&self,
owner: Option<(&str, &str)>,
renewers: &[(&str, &str)],
max_lifetime: Option<Duration>,
) -> Result<CreateDelegationTokenResult>
pub async fn create_delegation_token( &self, owner: Option<(&str, &str)>, renewers: &[(&str, &str)], max_lifetime: Option<Duration>, ) -> Result<CreateDelegationTokenResult>
Create a delegation token.
Delegation tokens allow a principal to delegate authentication to another principal without sharing credentials (KIP-48). The token HMAC can be used for SASL/SCRAM authentication.
CreateDelegationToken is a controller-only API: the request is
routed to the current controller and re-issued against a freshly
resolved controller on NOT_CONTROLLER.
§Arguments
owner- Principal the token is issued for, as a(type, name)pair, orNonefor the authenticated caller. Requesting a token on behalf of another principal is KIP-373 and needsCreateDelegationTokenv3+ plusCreateTokensauthorisation on that principal — it is how a superuser provisions a token for a service account that never authenticates interactively. The requester is recorded alongside the owner and comes back onDelegationToken::token_requester_principal_name, which is what an audit trail needs.renewers- Principals authorized to renew the token (type, name pairs). Pass an empty slice to allow only the token owner to renew.max_lifetime- Maximum token lifetime. UseNonefor the server default (typically 7 days).
Sourcepub async fn renew_delegation_token(
&self,
hmac: &[u8],
renew_period: Duration,
) -> Result<RenewDelegationTokenResult>
pub async fn renew_delegation_token( &self, hmac: &[u8], renew_period: Duration, ) -> Result<RenewDelegationTokenResult>
Renew a delegation token, extending its expiry time.
§Arguments
hmac- HMAC of the token to renew (fromDelegationToken::hmac).renew_period- How long to extend the token’s lifetime.
Sourcepub async fn expire_delegation_token(
&self,
hmac: &[u8],
expiry_period: Option<Duration>,
) -> Result<ExpireDelegationTokenResult>
pub async fn expire_delegation_token( &self, hmac: &[u8], expiry_period: Option<Duration>, ) -> Result<ExpireDelegationTokenResult>
Expire a delegation token, revoking it before its natural expiry.
§Arguments
hmac- HMAC of the token to expire (fromDelegationToken::hmac).expiry_period- How long until the token expires. PassNoneto expire the token immediately (sends-1to the broker).
Sourcepub async fn describe_delegation_token(
&self,
owners: Option<&[(&str, &str)]>,
) -> Result<Vec<DelegationToken>>
pub async fn describe_delegation_token( &self, owners: Option<&[(&str, &str)]>, ) -> Result<Vec<DelegationToken>>
Describe delegation tokens visible to the caller.
§Arguments
owners- Filter by token owners (type, name pairs). PassNoneto return all tokens visible to the caller.
Source§impl AdminClient
impl AdminClient
Sourcepub async fn create_topics(
&self,
topics: Vec<NewTopic>,
timeout: Duration,
validate_only: bool,
) -> Result<Vec<CreateTopicResult>>
pub async fn create_topics( &self, topics: Vec<NewTopic>, timeout: Duration, validate_only: bool, ) -> Result<Vec<CreateTopicResult>>
Create topics.
CreateTopics is a controller-only API. The request is routed to the
current controller and re-issued against a freshly resolved controller if
the broker answers NOT_CONTROLLER — see
get_controller_connection.
A NOT_CONTROLLER that survives every retry is returned as an Err, not
as an Ok carrying a per-topic error string.
Returns Ok(results) when the RPC succeeds. An Ok return does not
mean every topic was created — inspect each
CreateTopicResult::error for per-topic failures. Every requested
topic is guaranteed to appear in the result: a topic the broker omits
from its response is reported with an explicit error rather than
silently vanishing.
§Parameters
topics— Descriptions of the topics to create.timeout— How long the broker should wait for the creation to complete.validate_only— Whentrue, the broker validates the request but does not create any topics. Useful for pre-flight checks. Requires CreateTopics v2+ (Kafka 0.11+); all modern brokers support this.
Sourcepub async fn delete_topics(
&self,
topics: Vec<String>,
timeout: Duration,
) -> Result<Vec<DeleteTopicResult>>
pub async fn delete_topics( &self, topics: Vec<String>, timeout: Duration, ) -> Result<Vec<DeleteTopicResult>>
Delete topics.
DeleteTopics is a controller-only API; see
create_topics for how controller routing and
NOT_CONTROLLER retries work.
Returns Ok(results) when the RPC succeeds. An Ok return does not
mean every topic was deleted — inspect each
DeleteTopicResult::error for per-topic failures. Topics omitted by
the broker are reported with an explicit error.
Sourcepub async fn create_partitions(
&self,
topic: impl Into<String>,
new_total_count: i32,
timeout: Duration,
validate_only: bool,
) -> Result<CreatePartitionsResult>
pub async fn create_partitions( &self, topic: impl Into<String>, new_total_count: i32, timeout: Duration, validate_only: bool, ) -> Result<CreatePartitionsResult>
Increase the number of partitions for a topic.
Note: Partition count can only be increased, never decreased.
CreatePartitions is a controller-only API; see
create_topics for how controller routing and
NOT_CONTROLLER retries work.
§Parameters
topic— Topic to expand.new_total_count— The total partition count after the increase.timeout— How long the broker should wait for the change to complete.validate_only— Whentrue, the broker validates the request but does not create any partitions.
Source§impl AdminClient
impl AdminClient
Sourcepub async fn describe_producers(
&self,
topic_partitions: &[(&str, &[i32])],
) -> Result<Vec<DescribeProducersTopicResult>>
pub async fn describe_producers( &self, topic_partitions: &[(&str, &[i32])], ) -> Result<Vec<DescribeProducersTopicResult>>
Describe active producers on the given topic-partitions.
Routes each topic-partition to its leader broker via cached metadata for optimal performance. Falls back to any broker if the leader is unknown.
Returns per-partition producer state useful for debugging transactional and idempotent producers.
§Example
let results = admin
.describe_producers(&[("my-topic", &[0, 1])])
.await?;Sourcepub async fn describe_transactions(
&self,
transactional_ids: &[&str],
) -> Result<Vec<TransactionDescription>>
pub async fn describe_transactions( &self, transactional_ids: &[&str], ) -> Result<Vec<TransactionDescription>>
Sourcepub async fn list_transactions(
&self,
state_filters: &[&str],
producer_id_filters: &[i64],
duration_filter: i64,
transactional_id_pattern: Option<&str>,
) -> Result<ListTransactionsResult>
pub async fn list_transactions( &self, state_filters: &[&str], producer_id_filters: &[i64], duration_filter: i64, transactional_id_pattern: Option<&str>, ) -> Result<ListTransactionsResult>
List transactions matching the given filters.
Queries all brokers and merges results, because each broker only knows about transactions it coordinates.
Pass empty slices for state_filters and producer_id_filters, -1
for duration_filter and None for transactional_id_pattern to list
all transactions.
transactional_id_pattern requires ListTransactions v2 (KIP-1152,
Kafka 4.1+). Against an older broker the negotiated version is lower and
the pattern cannot be sent, so this method rejects the call rather than
silently returning unfiltered results.
§Example
// List all ongoing transactions
let txns = admin
.list_transactions(&["Ongoing"], &[], -1, None)
.await?;
// Only transactions whose transactional ID starts with "orders-"
let txns = admin
.list_transactions(&[], &[], -1, Some("orders-.*"))
.await?;Sourcepub async fn list_config_resources(
&self,
resource_types: &[ConfigResourceType],
) -> Result<Vec<ListedConfigResource>>
pub async fn list_config_resources( &self, resource_types: &[ConfigResourceType], ) -> Result<Vec<ListedConfigResource>>
List config resources known to the cluster (KIP-1142).
Kafka 4.1 renamed API key 74 from ListClientMetricsResources to
ListConfigResources and generalised it: v0 could only enumerate
client-metrics subscriptions (KIP-714), while v1 enumerates any config
resource type — topics, brokers, broker loggers, groups and client
metrics.
Pass an empty resource_types slice to get whichever types the broker
lists by default.
§Version behaviour
Against a pre-4.1 broker only v0 is available, which ignores
resource_types and returns client-metrics subscriptions. Requesting a
type other than ConfigResourceType::ClientMetrics there would
silently return the wrong set, so this method rejects the call instead.
§Example
use krafka::protocol::ConfigResourceType;
// Everything the broker lists by default.
let all = admin.list_config_resources(&[]).await?;
// Just the topics.
let topics = admin
.list_config_resources(&[ConfigResourceType::Topic])
.await?;Sourcepub async fn list_client_metrics_resources(&self) -> Result<Vec<String>>
pub async fn list_client_metrics_resources(&self) -> Result<Vec<String>>
List client metrics subscription names (KIP-714).
Convenience wrapper over list_config_resources
that asks for ConfigResourceType::ClientMetrics and returns only the
names. Works against any broker: this is exactly what API key 74 did
before KIP-1142 renamed it.
§Example
let names = admin.list_client_metrics_resources().await?;
for name in &names {
println!("subscription: {name}");
}Sourcepub async fn write_txn_markers(
&self,
markers: &[WritableTxnMarker],
) -> Result<Vec<WriteTxnMarkersResult>>
pub async fn write_txn_markers( &self, markers: &[WritableTxnMarker], ) -> Result<Vec<WriteTxnMarkersResult>>
Write transaction markers (COMMIT or ABORT) to the given topic-partitions.
This is an inter-broker API used to finalize transactions.
The admin client exposes it primarily for aborting stuck transactions
(abort_transaction).
§Routing
WriteTxnMarkers must reach the leader of each partition: only the
leader can append the marker to the log. This method groups each
marker’s topic-partitions by their current leader and sends one
WriteTxnMarkers request per leader, then merges the per-partition
results back together.
Sending the whole marker set to a single broker instead would leave a
transaction that spans several brokers only partially finalised: the
other leaders answer NOT_LEADER_OR_FOLLOWER, the transaction stays
open, and read_committed consumers block on the last stable offset
indefinitely.
§Errors
Returns Err if a partition’s leader is unknown, or if the request to a
leader fails outright — a marker that was never written must not be
reported as success. Per-partition broker errors are surfaced in
WriteTxnMarkersPartitionResult::error.
§Example
use krafka::protocol::{WritableTxnMarker, WritableTxnMarkerTopic};
let results = admin
.write_txn_markers(&[WritableTxnMarker {
producer_id: 42,
producer_epoch: 5,
transaction_result: false, // ABORT
topics: vec![WritableTxnMarkerTopic {
name: "my-topic".into(),
partition_indexes: vec![0, 1],
}],
coordinator_epoch: 10,
transaction_version: 0,
}])
.await?;Sourcepub async fn abort_transaction(
&self,
transactional_id: &str,
) -> Result<Vec<WriteTxnMarkersResult>>
pub async fn abort_transaction( &self, transactional_id: &str, ) -> Result<Vec<WriteTxnMarkersResult>>
Abort a stuck transaction by writing an ABORT marker.
This is the admin-friendly wrapper around
write_txn_markers that discovers the
affected partitions via
describe_transactions, determines the
owning coordinator epoch via
describe_producers, and writes an ABORT
marker to each partition’s leader.
§The coordinator epoch is discovered, never assumed
The partition leader validates the marker’s coordinator_epoch against
the epoch it has cached for that producer, and accepts any epoch when
it has none cached. Passing a fabricated 0 therefore succeeds
exactly on the partitions where validation cannot protect you — aborting
a transaction that a live, newer coordinator still owns and discarding
committed exactly-once data.
This method instead reads the real epoch from the active producer state
on the target partitions, mirroring Java’s AbortTransactionSpec.
abort_transaction_with_epoch
takes the epoch explicitly when the caller already knows it.
§Errors
Returns an error when the transaction cannot be described, when no active producer state exposes a coordinator epoch (so no safe epoch can be determined), or when the discovered producer state disagrees across partitions.
§Example
admin.abort_transaction("my-transactional-id").await?;Sourcepub async fn abort_transaction_with_epoch(
&self,
transactional_id: &str,
coordinator_epoch: i32,
) -> Result<Vec<WriteTxnMarkersResult>>
pub async fn abort_transaction_with_epoch( &self, transactional_id: &str, coordinator_epoch: i32, ) -> Result<Vec<WriteTxnMarkersResult>>
Abort a stuck transaction using a caller-supplied coordinator epoch.
Use this when the epoch is already known (for example from a previous
describe_producers call); otherwise prefer
abort_transaction, which discovers it.
§Warning
Supplying an epoch older than the one the partition leader has cached is rejected by the broker, but supplying an arbitrary epoch to a leader with no cached epoch is accepted. Never invent a value here.
Sourcepub async fn describe_metadata_quorum(
&self,
topic_partitions: &[(&str, &[i32])],
) -> Result<DescribeQuorumResult>
pub async fn describe_metadata_quorum( &self, topic_partitions: &[(&str, &[i32])], ) -> Result<DescribeQuorumResult>
Describe the KRaft quorum for the given topic-partitions.
In a KRaft-mode cluster this returns the current voters, observers, leader, leader epoch, and high watermark for each quorum partition.
The primary use case is inspecting __cluster_metadata partition 0.
§Example
let result = admin
.describe_metadata_quorum(&[("__cluster_metadata", &[0])])
.await?;Source§impl AdminClient
impl AdminClient
Sourcepub fn builder() -> AdminClientBuilder
pub fn builder() -> AdminClientBuilder
Create a new admin client builder.
Sourcepub async fn get_controller_connection(&self) -> Result<Arc<BrokerConnection>>
pub async fn get_controller_connection(&self) -> Result<Arc<BrokerConnection>>
Get a connection to the cluster controller.
Controller-only APIs must be routed here rather than to an arbitrary
broker. A non-controller broker does forward such requests, but during a
controller failover it answers NOT_CONTROLLER (41) instead — which the
admin APIs surface only as a per-item error string, so a caller checking
nothing but the Result concludes the operation succeeded.
§Errors
Returns crate::error::ErrorCode::UnknownControllerId if the cluster
reports no controller even after a metadata refresh.
Sourcepub async fn delete_consumer_groups(
&self,
group_ids: Vec<String>,
) -> Result<Vec<DeleteGroupResult>>
pub async fn delete_consumer_groups( &self, group_ids: Vec<String>, ) -> Result<Vec<DeleteGroupResult>>
Delete consumer groups by ID.
Returns one DeleteGroupResult per group. Each result may contain
an error if that particular group could not be deleted (e.g., it has
active members).
Sourcepub async fn describe_topic_partitions(
&self,
topics: Vec<String>,
) -> Result<DescribeTopicPartitionsResult>
pub async fn describe_topic_partitions( &self, topics: Vec<String>, ) -> Result<DescribeTopicPartitionsResult>
Describe topic partitions using the DescribeTopicPartitions API (Key 75).
Returns detailed per-partition information including leader, replicas, ISR, eligible leader replicas (ELR), and offline replicas. Supports pagination for topics with many partitions.
§Example
let result = admin
.describe_topic_partitions(vec!["my-topic".to_string()])
.await?;
for topic in &result.topics {
println!("{}: {} partitions", topic.name.as_deref().unwrap_or("?"), topic.partitions.len());
for p in &topic.partitions {
println!(" partition {}: leader={}, isr={:?}", p.partition_index, p.leader_id, p.isr_nodes);
}
}Sourcepub fn update_seed_brokers(&self, servers: Vec<String>) -> Result<()>
pub fn update_seed_brokers(&self, servers: Vec<String>) -> Result<()>
Replace the bootstrap server list at runtime (KIP-899).
The new addresses are used on the next metadata refresh that falls back to bootstrap servers. Does not close existing connections.
§Errors
Returns an error if servers is empty.
Sourcepub async fn refresh_tls(&self) -> Result<()>
pub async fn refresh_tls(&self) -> Result<()>
Re-read TLS certificate and key files from disk and atomically install the new material for all future connections (KIP-1288).
Existing TLS sessions are unaffected: they keep the connector they handshaked with and are replaced naturally as connections cycle. On error the previously loaded certificates stay active, so a call made mid-rotation against a half-written PEM is safe to retry.
No-op when TLS is not configured.
Use this for event-driven rotation (an inotify watch, a sidecar
signal). For unattended rotation set
TransportConfig::tls_reload_interval
instead and krafka reloads on a timer.
§Errors
Returns an error if the certificate or key files cannot be read or parsed.
Sourcepub async fn rebootstrap(&self)
pub async fn rebootstrap(&self)
Force a rebootstrap: close all connections, clear the metadata cache, and fall back to bootstrap servers (KIP-899).
Sourcepub async fn close(&self)
pub async fn close(&self)
Close the admin client.
Sets the closed flag so that subsequent operations fail fast.
§Connection teardown depends on pool ownership
- Own pool (the client was built from
bootstrap_servers): all broker connections are torn down. In-flight admin RPCs that have not yet received a response will fail with a network error, so callers should let long-running admin operations finish first. - Shared pool (built via
AdminClientBuilder::with_client): connections are left untouched. The pool belongs to theKrafkaClient, and closing it here would kill every producer and consumer connection on that client and fail all in-flight Produce/Fetch requests. Close theKrafkaClientto release those connections.
Calling close() more than once is a no-op.
Sourcepub fn owns_pool(&self) -> bool
pub fn owns_pool(&self) -> bool
Whether this admin client owns its connection pool.
false when the pool was borrowed from a
KrafkaClient via
AdminClientBuilder::with_client; in that case
close does not tear down connections.
Sourcepub fn connection_metrics(&self) -> Arc<ConnectionMetrics> ⓘ
pub fn connection_metrics(&self) -> Arc<ConnectionMetrics> ⓘ
Get the shared connection metrics handle used by this admin client’s broker pool.
Trait Implementations§
Source§impl Drop for AdminClient
impl Drop for AdminClient
Auto Trait Implementations§
impl !Freeze for AdminClient
impl !RefUnwindSafe for AdminClient
impl !UnwindSafe for AdminClient
impl Send for AdminClient
impl Sync for AdminClient
impl Unpin for AdminClient
impl UnsafeUnpin for AdminClient
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more