use kacrab_protocol::{
KafkaUuid,
generated::{
CreatableReplicaAssignment, CreatableTopic, CreatableTopicConfig,
CreatePartitionsAssignment, CreatePartitionsTopic, DescribeConfigsResource, ErrorCode,
alter_configs_request::{AlterConfigsResource, AlterableConfig},
incremental_alter_configs_request::{
AlterConfigsResource as IncrementalAlterConfigsResource,
AlterableConfig as IncrementalAlterableConfig,
},
},
};
use crate::common::{OffsetAndMetadata, TopicPartition};
#[derive(Debug, Clone)]
pub struct NewTopic {
name: String,
num_partitions: i32,
replication_factor: i16,
replica_assignments: Vec<(i32, Vec<i32>)>,
configs: Vec<(String, Option<String>)>,
}
impl NewTopic {
#[must_use]
pub fn new(name: impl Into<String>, num_partitions: i32, replication_factor: i16) -> Self {
Self {
name: name.into(),
num_partitions,
replication_factor,
replica_assignments: Vec::new(),
configs: Vec::new(),
}
}
#[must_use]
pub fn with_replica_assignments(
name: impl Into<String>,
assignments: Vec<(i32, Vec<i32>)>,
) -> Self {
Self {
name: name.into(),
num_partitions: -1,
replication_factor: -1,
replica_assignments: assignments,
configs: Vec::new(),
}
}
#[must_use]
pub fn config(mut self, name: impl Into<String>, value: Option<String>) -> Self {
self.configs.push((name.into(), value));
self
}
#[must_use]
pub fn name(&self) -> &str {
&self.name
}
pub(super) fn into_creatable(self) -> CreatableTopic {
CreatableTopic {
name: self.name.into(),
num_partitions: self.num_partitions,
replication_factor: self.replication_factor,
assignments: self
.replica_assignments
.into_iter()
.map(|(partition_index, broker_ids)| CreatableReplicaAssignment {
partition_index,
broker_ids,
_unknown_tagged_fields: Vec::new(),
})
.collect(),
configs: self
.configs
.into_iter()
.map(|(name, value)| CreatableTopicConfig {
name: name.into(),
value: value.map(Into::into),
_unknown_tagged_fields: Vec::new(),
})
.collect(),
_unknown_tagged_fields: Vec::new(),
}
}
}
#[derive(Debug, Clone)]
pub struct NewPartitions {
topic: String,
total_count: i32,
new_assignments: Vec<Vec<i32>>,
}
impl NewPartitions {
#[must_use]
pub fn increase_to(topic: impl Into<String>, total_count: i32) -> Self {
Self {
topic: topic.into(),
total_count,
new_assignments: Vec::new(),
}
}
#[must_use]
pub fn assigning(mut self, new_assignments: Vec<Vec<i32>>) -> Self {
self.new_assignments = new_assignments;
self
}
#[must_use]
pub fn topic(&self) -> &str {
&self.topic
}
pub(super) fn into_topic(self) -> CreatePartitionsTopic {
let assignments = if self.new_assignments.is_empty() {
None
} else {
Some(
self.new_assignments
.into_iter()
.map(|broker_ids| CreatePartitionsAssignment {
broker_ids,
_unknown_tagged_fields: Vec::new(),
})
.collect(),
)
};
CreatePartitionsTopic {
name: self.topic.into(),
count: self.total_count,
assignments,
_unknown_tagged_fields: Vec::new(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum ResourceType {
Unknown,
Topic,
Broker,
BrokerLogger,
Group,
ClientMetrics,
}
impl ResourceType {
pub(super) const fn to_wire(self) -> i8 {
match self {
Self::Unknown => 0,
Self::Topic => 2,
Self::Broker => 4,
Self::BrokerLogger => 8,
Self::ClientMetrics => 16,
Self::Group => 32,
}
}
pub(super) const fn from_wire(value: i8) -> Self {
match value {
2 => Self::Topic,
4 => Self::Broker,
8 => Self::BrokerLogger,
16 => Self::ClientMetrics,
32 => Self::Group,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct ConfigResource {
pub resource_type: ResourceType,
pub name: String,
}
impl ConfigResource {
#[must_use]
pub fn topic(name: impl Into<String>) -> Self {
Self {
resource_type: ResourceType::Topic,
name: name.into(),
}
}
#[must_use]
pub fn broker(broker_id: i32) -> Self {
Self {
resource_type: ResourceType::Broker,
name: broker_id.to_string(),
}
}
pub(super) fn to_describe(&self) -> DescribeConfigsResource {
DescribeConfigsResource {
resource_type: self.resource_type.to_wire(),
resource_name: self.name.clone().into(),
configuration_keys: None,
_unknown_tagged_fields: Vec::new(),
}
}
pub(super) fn to_alter(&self, entries: Vec<ConfigEntry>) -> AlterConfigsResource {
AlterConfigsResource {
resource_type: self.resource_type.to_wire(),
resource_name: self.name.clone().into(),
configs: entries
.into_iter()
.map(|entry| AlterableConfig {
name: entry.name.into(),
value: entry.value.map(Into::into),
_unknown_tagged_fields: Vec::new(),
})
.collect(),
_unknown_tagged_fields: Vec::new(),
}
}
pub(super) fn to_incremental(
&self,
ops: Vec<AlterConfigOp>,
) -> IncrementalAlterConfigsResource {
IncrementalAlterConfigsResource {
resource_type: self.resource_type.to_wire(),
resource_name: self.name.clone().into(),
configs: ops
.into_iter()
.map(|op| IncrementalAlterableConfig {
name: op.name.into(),
config_operation: op.op_type.to_wire(),
value: op.value.map(Into::into),
_unknown_tagged_fields: Vec::new(),
})
.collect(),
_unknown_tagged_fields: Vec::new(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AlterConfigOpType {
Set,
Delete,
Append,
Subtract,
}
impl AlterConfigOpType {
pub(super) const fn to_wire(self) -> i8 {
match self {
Self::Set => 0,
Self::Delete => 1,
Self::Append => 2,
Self::Subtract => 3,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AlterConfigOp {
pub name: String,
pub value: Option<String>,
pub op_type: AlterConfigOpType,
}
impl AlterConfigOp {
#[must_use]
pub fn set(name: impl Into<String>, value: impl Into<String>) -> Self {
Self {
name: name.into(),
value: Some(value.into()),
op_type: AlterConfigOpType::Set,
}
}
#[must_use]
pub fn delete(name: impl Into<String>) -> Self {
Self {
name: name.into(),
value: None,
op_type: AlterConfigOpType::Delete,
}
}
#[must_use]
pub fn append(name: impl Into<String>, value: impl Into<String>) -> Self {
Self {
name: name.into(),
value: Some(value.into()),
op_type: AlterConfigOpType::Append,
}
}
#[must_use]
pub fn subtract(name: impl Into<String>, value: impl Into<String>) -> Self {
Self {
name: name.into(),
value: Some(value.into()),
op_type: AlterConfigOpType::Subtract,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConfigEntry {
pub name: String,
pub value: Option<String>,
pub read_only: bool,
pub is_sensitive: bool,
pub source: ConfigSource,
}
impl ConfigEntry {
#[must_use]
pub fn set(name: impl Into<String>, value: Option<String>) -> Self {
Self {
name: name.into(),
value,
read_only: false,
is_sensitive: false,
source: ConfigSource::Unknown,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ConfigSource {
Unknown,
TopicConfig,
DynamicBrokerConfig,
DynamicDefaultBrokerConfig,
StaticBrokerConfig,
DefaultConfig,
DynamicBrokerLoggerConfig,
}
impl ConfigSource {
pub(super) const fn from_wire(value: i8) -> Self {
match value {
1 => Self::TopicConfig,
2 => Self::DynamicBrokerConfig,
3 => Self::DynamicDefaultBrokerConfig,
4 => Self::StaticBrokerConfig,
5 => Self::DefaultConfig,
6 => Self::DynamicBrokerLoggerConfig,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone)]
pub struct ResourceConfig {
pub resource: ConfigResource,
pub entries: Vec<ConfigEntry>,
}
pub use crate::common::Node;
#[derive(Debug, Clone)]
pub struct ClusterDescription {
pub cluster_id: Option<String>,
pub controller: Option<Node>,
pub nodes: Vec<Node>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TopicListing {
pub name: String,
pub topic_id: KafkaUuid,
pub is_internal: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TopicPartitionInfo {
pub partition: i32,
pub leader: Option<Node>,
pub replicas: Vec<Node>,
pub isr: Vec<Node>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TopicDescription {
pub name: String,
pub topic_id: KafkaUuid,
pub is_internal: bool,
pub partitions: Vec<TopicPartitionInfo>,
}
#[derive(Debug, Clone, Default)]
pub struct ListTopicsOptions {
pub list_internal: bool,
}
#[derive(Debug, Clone, Default)]
pub struct DescribeTopicsOptions;
#[derive(Debug, Clone, Default)]
pub struct CreateTopicsOptions {
pub validate_only: bool,
}
#[derive(Debug, Clone, Default)]
pub struct CreatePartitionsOptions {
pub validate_only: bool,
}
#[derive(Debug, Clone, Default)]
pub struct AlterConfigsOptions {
pub validate_only: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerGroupListing {
pub group_id: String,
pub is_simple_consumer_group: bool,
pub state: Option<GroupState>,
pub group_type: Option<GroupType>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[expect(missing_docs, reason = "Variants mirror Kafka's GroupState names 1:1.")]
pub enum GroupState {
Unknown,
PreparingRebalance,
CompletingRebalance,
Stable,
Dead,
Empty,
Assigning,
Reconciling,
}
impl GroupState {
#[must_use]
pub fn from_broker(state: &str) -> Self {
match state {
"PreparingRebalance" => Self::PreparingRebalance,
"CompletingRebalance" => Self::CompletingRebalance,
"Stable" => Self::Stable,
"Dead" => Self::Dead,
"Empty" => Self::Empty,
"Assigning" => Self::Assigning,
"Reconciling" => Self::Reconciling,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[expect(missing_docs, reason = "Variants mirror Kafka's GroupType names 1:1.")]
pub enum GroupType {
Unknown,
Classic,
Consumer,
Share,
Streams,
}
impl GroupType {
#[must_use]
pub fn from_broker(group_type: &str) -> Self {
match group_type.to_ascii_lowercase().as_str() {
"classic" => Self::Classic,
"consumer" => Self::Consumer,
"share" => Self::Share,
"streams" => Self::Streams,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MemberDescription {
pub member_id: String,
pub group_instance_id: Option<String>,
pub rack_id: Option<String>,
pub client_id: String,
pub host: String,
pub assignment: Vec<TopicPartition>,
pub target_assignment: Vec<TopicPartition>,
pub member_epoch: Option<i32>,
pub upgraded: Option<bool>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerGroupDescription {
pub group_id: String,
pub is_simple_consumer_group: bool,
pub members: Vec<MemberDescription>,
pub partition_assignor: String,
pub state: GroupState,
pub group_type: GroupType,
pub coordinator: Node,
pub authorized_operations: Vec<AclOperation>,
pub group_epoch: Option<i32>,
pub target_assignment_epoch: Option<i32>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct GroupOffset {
pub partition: TopicPartition,
pub offset: OffsetAndMetadata,
}
#[derive(Debug, Clone, Default)]
pub struct ListConsumerGroupsOptions {
pub states_filter: Vec<String>,
pub types_filter: Vec<String>,
}
#[derive(Debug, Clone, Default)]
pub struct DescribeConsumerGroupsOptions {
pub include_authorized_operations: bool,
}
#[derive(Debug, Clone, Default)]
pub struct ListConsumerGroupOffsetsOptions {
pub partitions: Vec<TopicPartition>,
pub require_stable: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ElectionType {
Preferred,
Unclean,
}
impl ElectionType {
pub(super) const fn to_wire(self) -> i8 {
match self {
Self::Preferred => 0,
Self::Unclean => 1,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OffsetSpec {
Earliest,
Latest,
MaxTimestamp,
Timestamp(i64),
}
impl OffsetSpec {
pub(super) const fn to_wire(self) -> i64 {
match self {
Self::Earliest => -2,
Self::Latest => -1,
Self::MaxTimestamp => -3,
Self::Timestamp(timestamp) => timestamp,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ListOffsetsResult {
pub partition: TopicPartition,
pub offset: i64,
pub timestamp: i64,
pub leader_epoch: Option<i32>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DeletedRecords {
pub partition: TopicPartition,
pub low_watermark: i64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProducerState {
pub producer_id: i64,
pub producer_epoch: i32,
pub last_sequence: i32,
pub last_timestamp: i64,
pub coordinator_epoch: i32,
pub current_transaction_start_offset: i64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PartitionProducerState {
pub partition: TopicPartition,
pub active_producers: Vec<ProducerState>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TransactionDescription {
pub transactional_id: String,
pub state: String,
pub producer_id: i64,
pub producer_epoch: i16,
pub transaction_timeout_ms: i32,
pub transaction_start_time_ms: i64,
pub topic_partitions: Vec<TopicPartition>,
pub coordinator: Node,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TransactionListing {
pub transactional_id: String,
pub producer_id: i64,
pub state: String,
}
#[derive(Debug, Clone, Default)]
pub struct ListTransactionsOptions {
pub state_filters: Vec<String>,
pub producer_id_filters: Vec<i64>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LogDirReplicaInfo {
pub size: i64,
pub offset_lag: i64,
pub is_future: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LogDirDescription {
pub log_dir: String,
pub error: Option<ErrorCode>,
pub total_bytes: i64,
pub usable_bytes: i64,
pub replicas: Vec<(TopicPartition, LogDirReplicaInfo)>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BrokerLogDirs {
pub broker_id: i32,
pub log_dirs: Vec<LogDirDescription>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NewPartitionReassignment {
pub topic_partition: TopicPartition,
pub replicas: Option<Vec<i32>>,
}
impl NewPartitionReassignment {
#[must_use]
pub const fn assigning(topic_partition: TopicPartition, replicas: Vec<i32>) -> Self {
Self {
topic_partition,
replicas: Some(replicas),
}
}
#[must_use]
pub const fn cancel(topic_partition: TopicPartition) -> Self {
Self {
topic_partition,
replicas: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PartitionReassignment {
pub topic_partition: TopicPartition,
pub replicas: Vec<i32>,
pub adding_replicas: Vec<i32>,
pub removing_replicas: Vec<i32>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FeatureUpdateUpgradeType {
Upgrade,
SafeDowngrade,
UnsafeDowngrade,
}
impl FeatureUpdateUpgradeType {
pub(super) const fn to_wire(self) -> i8 {
match self {
Self::Upgrade => 1,
Self::SafeDowngrade => 2,
Self::UnsafeDowngrade => 3,
}
}
pub(super) const fn allows_downgrade(self) -> bool {
!matches!(self, Self::Upgrade)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FeatureUpdate {
pub feature: String,
pub max_version_level: i16,
pub upgrade_type: FeatureUpdateUpgradeType,
}
impl FeatureUpdate {
#[must_use]
pub fn new(
feature: impl Into<String>,
max_version_level: i16,
upgrade_type: FeatureUpdateUpgradeType,
) -> Self {
Self {
feature: feature.into(),
max_version_level,
upgrade_type,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AclResourceType {
Unknown,
Any,
Topic,
Group,
Cluster,
TransactionalId,
DelegationToken,
User,
}
impl AclResourceType {
pub(super) const fn to_wire(self) -> i8 {
match self {
Self::Unknown => 0,
Self::Any => 1,
Self::Topic => 2,
Self::Group => 3,
Self::Cluster => 4,
Self::TransactionalId => 5,
Self::DelegationToken => 6,
Self::User => 7,
}
}
pub(super) const fn from_wire(value: i8) -> Self {
match value {
1 => Self::Any,
2 => Self::Topic,
3 => Self::Group,
4 => Self::Cluster,
5 => Self::TransactionalId,
6 => Self::DelegationToken,
7 => Self::User,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AclPatternType {
Unknown,
Any,
Match,
Literal,
Prefixed,
}
impl AclPatternType {
pub(super) const fn to_wire(self) -> i8 {
match self {
Self::Unknown => 0,
Self::Any => 1,
Self::Match => 2,
Self::Literal => 3,
Self::Prefixed => 4,
}
}
pub(super) const fn from_wire(value: i8) -> Self {
match value {
1 => Self::Any,
2 => Self::Match,
3 => Self::Literal,
4 => Self::Prefixed,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[expect(
missing_docs,
reason = "Variants mirror Kafka's AclOperation constants 1:1."
)]
pub enum AclOperation {
Unknown,
Any,
All,
Read,
Write,
Create,
Delete,
Alter,
Describe,
ClusterAction,
DescribeConfigs,
AlterConfigs,
IdempotentWrite,
CreateTokens,
DescribeTokens,
}
impl AclOperation {
pub(super) const fn to_wire(self) -> i8 {
match self {
Self::Unknown => 0,
Self::Any => 1,
Self::All => 2,
Self::Read => 3,
Self::Write => 4,
Self::Create => 5,
Self::Delete => 6,
Self::Alter => 7,
Self::Describe => 8,
Self::ClusterAction => 9,
Self::DescribeConfigs => 10,
Self::AlterConfigs => 11,
Self::IdempotentWrite => 12,
Self::CreateTokens => 13,
Self::DescribeTokens => 14,
}
}
pub(super) const fn from_wire(value: i8) -> Self {
match value {
1 => Self::Any,
2 => Self::All,
3 => Self::Read,
4 => Self::Write,
5 => Self::Create,
6 => Self::Delete,
7 => Self::Alter,
8 => Self::Describe,
9 => Self::ClusterAction,
10 => Self::DescribeConfigs,
11 => Self::AlterConfigs,
12 => Self::IdempotentWrite,
13 => Self::CreateTokens,
14 => Self::DescribeTokens,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AclPermissionType {
Unknown,
Any,
Deny,
Allow,
}
impl AclPermissionType {
pub(super) const fn to_wire(self) -> i8 {
match self {
Self::Unknown => 0,
Self::Any => 1,
Self::Deny => 2,
Self::Allow => 3,
}
}
pub(super) const fn from_wire(value: i8) -> Self {
match value {
1 => Self::Any,
2 => Self::Deny,
3 => Self::Allow,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AclBinding {
pub resource_type: AclResourceType,
pub resource_name: String,
pub pattern_type: AclPatternType,
pub principal: String,
pub host: String,
pub operation: AclOperation,
pub permission_type: AclPermissionType,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AclBindingFilter {
pub resource_type: AclResourceType,
pub resource_name: Option<String>,
pub pattern_type: AclPatternType,
pub principal: Option<String>,
pub host: Option<String>,
pub operation: AclOperation,
pub permission_type: AclPermissionType,
}
impl AclBindingFilter {
#[must_use]
pub const fn any() -> Self {
Self {
resource_type: AclResourceType::Any,
resource_name: None,
pattern_type: AclPatternType::Any,
principal: None,
host: None,
operation: AclOperation::Any,
permission_type: AclPermissionType::Any,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct ClientQuotaEntity {
pub entries: Vec<(String, Option<String>)>,
}
impl ClientQuotaEntity {
pub const USER: &'static str = "user";
pub const CLIENT_ID: &'static str = "client-id";
pub const IP: &'static str = "ip";
#[must_use]
pub const fn new(entries: Vec<(String, Option<String>)>) -> Self {
Self { entries }
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ClientQuotaMatch {
Exact(String),
Default,
Any,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ClientQuotaFilterComponent {
pub entity_type: String,
pub match_type: ClientQuotaMatch,
}
#[derive(Debug, Clone, PartialEq)]
pub struct ClientQuotaEntry {
pub entity: ClientQuotaEntity,
pub quotas: Vec<(String, f64)>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct ClientQuotaOp {
pub key: String,
pub value: Option<f64>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct ClientQuotaAlteration {
pub entity: ClientQuotaEntity,
pub ops: Vec<ClientQuotaOp>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ScramMechanism {
Unknown,
ScramSha256,
ScramSha512,
}
impl ScramMechanism {
pub(super) const fn to_wire(self) -> i8 {
match self {
Self::Unknown => 0,
Self::ScramSha256 => 1,
Self::ScramSha512 => 2,
}
}
pub(super) const fn from_wire(value: i8) -> Self {
match value {
1 => Self::ScramSha256,
2 => Self::ScramSha512,
_ => Self::Unknown,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ScramCredentialInfo {
pub mechanism: ScramMechanism,
pub iterations: i32,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UserScramCredentials {
pub user: String,
pub credentials: Vec<ScramCredentialInfo>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ScramCredentialUpsertion {
pub user: String,
pub mechanism: ScramMechanism,
pub iterations: i32,
pub salt: Vec<u8>,
pub salted_password: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ScramCredentialDeletion {
pub user: String,
pub mechanism: ScramMechanism,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DelegationToken {
pub token_id: String,
pub owner_principal_type: String,
pub owner_principal_name: String,
pub issue_timestamp_ms: i64,
pub expiry_timestamp_ms: i64,
pub max_timestamp_ms: i64,
pub hmac: Vec<u8>,
pub renewers: Vec<(String, String)>,
}
#[derive(Debug, Clone, Default)]
pub struct CreateDelegationTokenOptions {
pub owner: Option<(String, String)>,
pub renewers: Vec<(String, String)>,
pub max_lifetime_ms: i64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ReplicaLogDirAssignment {
pub topic_partition: TopicPartition,
pub broker_id: i32,
pub log_dir: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FencedProducer {
pub transactional_id: String,
pub producer_id: i64,
pub producer_epoch: i16,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AbortTransactionSpec {
pub topic_partition: TopicPartition,
pub producer_id: i64,
pub producer_epoch: i16,
pub coordinator_epoch: i32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct SupportedVersionRange {
pub min_version: i16,
pub max_version: i16,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct FinalizedVersionRange {
pub min_version_level: i16,
pub max_version_level: i16,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FeatureMetadata {
pub finalized_features_epoch: Option<i64>,
pub supported_features: Vec<(String, SupportedVersionRange)>,
pub finalized_features: Vec<(String, FinalizedVersionRange)>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct MemberToRemove {
pub member_id: Option<String>,
pub group_instance_id: Option<String>,
}
impl MemberToRemove {
#[must_use]
pub fn static_member(group_instance_id: impl Into<String>) -> Self {
Self {
member_id: None,
group_instance_id: Some(group_instance_id.into()),
}
}
#[must_use]
pub fn dynamic_member(member_id: impl Into<String>) -> Self {
Self {
member_id: Some(member_id.into()),
group_instance_id: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ReplicaLogDirInfo {
pub topic_partition: TopicPartition,
pub broker_id: i32,
pub current_log_dir: Option<String>,
pub future_log_dir: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct QuorumReplicaState {
pub replica_id: i32,
pub log_end_offset: i64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QuorumInfo {
pub leader_id: i32,
pub leader_epoch: i32,
pub high_watermark: i64,
pub voters: Vec<QuorumReplicaState>,
pub observers: Vec<QuorumReplicaState>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RaftVoterEndpoint {
pub name: String,
pub host: String,
pub port: u16,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ShareGroupDescription {
pub group_id: String,
pub state: GroupState,
pub group_epoch: i32,
pub assignor_name: String,
pub members: Vec<MemberDescription>,
pub coordinator: Node,
pub authorized_operations: Vec<AclOperation>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StreamsGroupDescription {
pub group_id: String,
pub state: GroupState,
pub group_epoch: i32,
pub members: Vec<MemberDescription>,
pub coordinator: Node,
pub authorized_operations: Vec<AclOperation>,
}
#[cfg(test)]
mod tests {
use kacrab_protocol::KafkaString;
use super::*;
#[test]
fn resource_type_wire_round_trips() {
for ty in [
ResourceType::Topic,
ResourceType::Broker,
ResourceType::BrokerLogger,
ResourceType::ClientMetrics,
ResourceType::Group,
ResourceType::Unknown,
] {
assert_eq!(ResourceType::from_wire(ty.to_wire()), ty);
}
assert_eq!(ResourceType::from_wire(99), ResourceType::Unknown);
}
#[test]
fn config_source_maps_known_and_unknown() {
assert_eq!(ConfigSource::from_wire(1), ConfigSource::TopicConfig);
assert_eq!(ConfigSource::from_wire(4), ConfigSource::StaticBrokerConfig);
assert_eq!(ConfigSource::from_wire(5), ConfigSource::DefaultConfig);
assert_eq!(ConfigSource::from_wire(-1), ConfigSource::Unknown);
assert_eq!(ConfigSource::from_wire(42), ConfigSource::Unknown);
}
#[test]
fn new_topic_partition_count_and_configs() {
let topic = NewTopic::new("orders", 6, 3)
.config("retention.ms", Some("60000".to_owned()))
.config("cleanup.policy", None);
assert_eq!(topic.name(), "orders");
let wire = topic.into_creatable();
assert_eq!(wire.name.as_str(), "orders");
assert_eq!(wire.num_partitions, 6);
assert_eq!(wire.replication_factor, 3);
assert!(wire.assignments.is_empty());
assert_eq!(wire.configs.len(), 2);
assert_eq!(wire.configs[0].name.as_str(), "retention.ms");
assert_eq!(
wire.configs[0].value.as_ref().map(KafkaString::as_str),
Some("60000")
);
assert_eq!(wire.configs[1].value, None);
}
#[test]
fn new_topic_with_replica_assignments_sends_negative_counts() {
let topic =
NewTopic::with_replica_assignments("orders", vec![(0, vec![1, 2]), (1, vec![2, 3])]);
let wire = topic.into_creatable();
assert_eq!(wire.num_partitions, -1);
assert_eq!(wire.replication_factor, -1);
assert_eq!(wire.assignments.len(), 2);
assert_eq!(wire.assignments[0].partition_index, 0);
assert_eq!(wire.assignments[0].broker_ids, vec![1, 2]);
}
#[test]
fn new_partitions_increase_and_assign() {
let plain = NewPartitions::increase_to("orders", 8).into_topic();
assert_eq!(plain.name.as_str(), "orders");
assert_eq!(plain.count, 8);
assert!(plain.assignments.is_none());
let assigned = NewPartitions::increase_to("orders", 8)
.assigning(vec![vec![1, 2], vec![3, 4]])
.into_topic();
let assignments = assigned.assignments.expect("assignments present");
assert_eq!(assignments.len(), 2);
assert_eq!(assignments[1].broker_ids, vec![3, 4]);
}
#[test]
fn config_resource_builders_set_type_and_name() {
let topic = ConfigResource::topic("orders");
assert_eq!(topic.resource_type, ResourceType::Topic);
assert_eq!(topic.name, "orders");
let broker = ConfigResource::broker(7);
assert_eq!(broker.resource_type, ResourceType::Broker);
assert_eq!(broker.name, "7");
let describe = broker.to_describe();
assert_eq!(describe.resource_type, ResourceType::Broker.to_wire());
assert_eq!(describe.resource_name.as_str(), "7");
assert!(describe.configuration_keys.is_none());
}
#[test]
fn alter_config_op_type_wire_values_match_kafka() {
assert_eq!(AlterConfigOpType::Set.to_wire(), 0);
assert_eq!(AlterConfigOpType::Delete.to_wire(), 1);
assert_eq!(AlterConfigOpType::Append.to_wire(), 2);
assert_eq!(AlterConfigOpType::Subtract.to_wire(), 3);
}
#[test]
fn election_and_offset_spec_wire_values_match_kafka() {
assert_eq!(ElectionType::Preferred.to_wire(), 0);
assert_eq!(ElectionType::Unclean.to_wire(), 1);
assert_eq!(OffsetSpec::Earliest.to_wire(), -2);
assert_eq!(OffsetSpec::Latest.to_wire(), -1);
assert_eq!(OffsetSpec::MaxTimestamp.to_wire(), -3);
assert_eq!(OffsetSpec::Timestamp(1234).to_wire(), 1234);
}
#[test]
fn config_resource_to_incremental_maps_ops() {
let resource = ConfigResource::topic("orders");
let incremental = resource.to_incremental(vec![
AlterConfigOp::set("retention.ms", "60000"),
AlterConfigOp::delete("cleanup.policy"),
AlterConfigOp::append("follower.replication.throttled.replicas", "1:2"),
]);
assert_eq!(incremental.resource_name.as_str(), "orders");
assert_eq!(incremental.configs.len(), 3);
assert_eq!(incremental.configs[0].config_operation, 0);
assert_eq!(
incremental.configs[0]
.value
.as_ref()
.map(KafkaString::as_str),
Some("60000")
);
assert_eq!(incremental.configs[1].config_operation, 1);
assert_eq!(incremental.configs[1].value, None);
assert_eq!(incremental.configs[2].config_operation, 2);
}
#[test]
fn config_resource_to_alter_carries_entries() {
let resource = ConfigResource::topic("orders");
let alter = resource.to_alter(vec![
ConfigEntry::set("retention.ms", Some("60000".to_owned())),
ConfigEntry::set("cleanup.policy", None),
]);
assert_eq!(alter.resource_name.as_str(), "orders");
assert_eq!(alter.configs.len(), 2);
assert_eq!(alter.configs[0].name.as_str(), "retention.ms");
assert_eq!(
alter.configs[0].value.as_ref().map(KafkaString::as_str),
Some("60000")
);
assert_eq!(alter.configs[1].value, None);
}
#[test]
fn acl_enums_round_trip_through_wire_values() {
for ty in [
AclResourceType::Any,
AclResourceType::Topic,
AclResourceType::Group,
AclResourceType::Cluster,
AclResourceType::TransactionalId,
AclResourceType::DelegationToken,
AclResourceType::User,
AclResourceType::Unknown,
] {
assert_eq!(AclResourceType::from_wire(ty.to_wire()), ty);
}
assert_eq!(AclResourceType::from_wire(99), AclResourceType::Unknown);
for pattern in [
AclPatternType::Any,
AclPatternType::Match,
AclPatternType::Literal,
AclPatternType::Prefixed,
AclPatternType::Unknown,
] {
assert_eq!(AclPatternType::from_wire(pattern.to_wire()), pattern);
}
for permission in [
AclPermissionType::Any,
AclPermissionType::Deny,
AclPermissionType::Allow,
AclPermissionType::Unknown,
] {
assert_eq!(
AclPermissionType::from_wire(permission.to_wire()),
permission
);
}
assert_eq!(AclOperation::Read.to_wire(), 3);
assert_eq!(AclOperation::DescribeTokens.to_wire(), 14);
for op in [
AclOperation::All,
AclOperation::Read,
AclOperation::Write,
AclOperation::Create,
AclOperation::Delete,
AclOperation::Alter,
AclOperation::Describe,
AclOperation::IdempotentWrite,
AclOperation::Unknown,
] {
assert_eq!(AclOperation::from_wire(op.to_wire()), op);
}
}
#[test]
fn acl_binding_filter_any_matches_everything() {
let filter = AclBindingFilter::any();
assert_eq!(filter.resource_type, AclResourceType::Any);
assert_eq!(filter.pattern_type, AclPatternType::Any);
assert_eq!(filter.operation, AclOperation::Any);
assert_eq!(filter.permission_type, AclPermissionType::Any);
assert_eq!(filter.resource_name, None);
assert_eq!(filter.principal, None);
assert_eq!(filter.host, None);
}
#[test]
fn member_to_remove_constructors_set_one_identifier() {
let static_member = MemberToRemove::static_member("instance-1");
assert_eq!(
static_member.group_instance_id.as_deref(),
Some("instance-1")
);
assert_eq!(static_member.member_id, None);
let dynamic_member = MemberToRemove::dynamic_member("member-1");
assert_eq!(dynamic_member.member_id.as_deref(), Some("member-1"));
assert_eq!(dynamic_member.group_instance_id, None);
}
#[test]
fn scram_mechanism_round_trips_and_feature_upgrade_maps() {
for mechanism in [
ScramMechanism::ScramSha256,
ScramMechanism::ScramSha512,
ScramMechanism::Unknown,
] {
assert_eq!(ScramMechanism::from_wire(mechanism.to_wire()), mechanism);
}
assert_eq!(ScramMechanism::from_wire(9), ScramMechanism::Unknown);
assert_eq!(FeatureUpdateUpgradeType::Upgrade.to_wire(), 1);
assert!(!FeatureUpdateUpgradeType::Upgrade.allows_downgrade());
assert_eq!(FeatureUpdateUpgradeType::SafeDowngrade.to_wire(), 2);
assert!(FeatureUpdateUpgradeType::SafeDowngrade.allows_downgrade());
assert_eq!(FeatureUpdateUpgradeType::UnsafeDowngrade.to_wire(), 3);
assert!(FeatureUpdateUpgradeType::UnsafeDowngrade.allows_downgrade());
}
}