use kafka_conn::ApiKey;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum CoordinatorKind {
Group,
Transaction,
}
impl CoordinatorKind {
pub const fn key_type(self) -> i8 {
match self {
CoordinatorKind::Group => 0,
CoordinatorKind::Transaction => 1,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BrokerSelector {
Caller,
PartitionLeader,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Routing {
Any,
Controller,
Coordinator(CoordinatorKind),
Specific(BrokerSelector),
}
pub const fn routing(api_key: ApiKey) -> Routing {
match api_key {
ApiKey::CreateTopics
| ApiKey::DeleteTopics
| ApiKey::CreatePartitions
| ApiKey::AlterPartitionReassignments
| ApiKey::ListPartitionReassignments
| ApiKey::ElectLeaders
| ApiKey::UpdateFeatures => Routing::Controller,
ApiKey::OffsetCommit
| ApiKey::OffsetFetch
| ApiKey::OffsetDelete
| ApiKey::JoinGroup
| ApiKey::Heartbeat
| ApiKey::LeaveGroup
| ApiKey::SyncGroup
| ApiKey::DescribeGroups
| ApiKey::DeleteGroups
| ApiKey::ConsumerGroupDescribe
| ApiKey::ConsumerGroupHeartbeat
| ApiKey::ShareGroupDescribe
| ApiKey::ShareGroupHeartbeat
| ApiKey::DescribeShareGroupOffsets
| ApiKey::AlterShareGroupOffsets
| ApiKey::DeleteShareGroupOffsets
| ApiKey::TxnOffsetCommit => Routing::Coordinator(CoordinatorKind::Group),
ApiKey::InitProducerId
| ApiKey::AddPartitionsToTxn
| ApiKey::AddOffsetsToTxn
| ApiKey::EndTxn
| ApiKey::DescribeTransactions => Routing::Coordinator(CoordinatorKind::Transaction),
ApiKey::DescribeLogDirs | ApiKey::AlterReplicaLogDirs | ApiKey::DescribeProducers => {
Routing::Specific(BrokerSelector::Caller)
}
ApiKey::Produce | ApiKey::Fetch | ApiKey::ListOffsets | ApiKey::OffsetForLeaderEpoch => {
Routing::Specific(BrokerSelector::PartitionLeader)
}
_ => Routing::Any,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_four_classes_from_claude_md() {
for key in [
ApiKey::CreateTopics,
ApiKey::DeleteTopics,
ApiKey::CreatePartitions,
ApiKey::AlterPartitionReassignments,
ApiKey::ElectLeaders,
ApiKey::UpdateFeatures,
] {
assert_eq!(routing(key), Routing::Controller, "{key}");
}
assert_eq!(
routing(ApiKey::OffsetFetch),
Routing::Coordinator(CoordinatorKind::Group)
);
assert_eq!(
routing(ApiKey::DescribeTransactions),
Routing::Coordinator(CoordinatorKind::Transaction)
);
assert_eq!(
routing(ApiKey::DescribeLogDirs),
Routing::Specific(BrokerSelector::Caller)
);
assert_eq!(
routing(ApiKey::DescribeProducers),
Routing::Specific(BrokerSelector::Caller)
);
for key in [
ApiKey::DescribeConfigs,
ApiKey::DescribeAcls,
ApiKey::ListGroups,
ApiKey::Metadata,
] {
assert_eq!(routing(key), Routing::Any, "{key}");
}
}
#[test]
fn the_read_path_goes_to_the_leader() {
for key in [ApiKey::Fetch, ApiKey::ListOffsets, ApiKey::Produce] {
assert_eq!(
routing(key),
Routing::Specific(BrokerSelector::PartitionLeader),
"{key}"
);
}
}
#[test]
fn share_and_consumer_group_apis_route_like_classic_groups() {
for key in [
ApiKey::ConsumerGroupDescribe,
ApiKey::ShareGroupDescribe,
ApiKey::DescribeShareGroupOffsets,
] {
assert_eq!(
routing(key),
Routing::Coordinator(CoordinatorKind::Group),
"{key}"
);
}
}
#[test]
fn an_api_key_this_build_cannot_name_routes_anywhere() {
assert_eq!(routing(ApiKey::Unknown(9_999)), Routing::Any);
}
#[test]
fn find_coordinator_key_types_match_the_protocol() {
assert_eq!(CoordinatorKind::Group.key_type(), 0);
assert_eq!(CoordinatorKind::Transaction.key_type(), 1);
}
}