1use kafka_conn::ApiKey;
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
23pub enum CoordinatorKind {
24 Group,
26 Transaction,
28}
29
30impl CoordinatorKind {
31 pub const fn key_type(self) -> i8 {
33 match self {
34 CoordinatorKind::Group => 0,
35 CoordinatorKind::Transaction => 1,
36 }
37 }
38}
39
40#[derive(Debug, Clone, Copy, PartialEq, Eq)]
42pub enum BrokerSelector {
43 Caller,
45 PartitionLeader,
47}
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51pub enum Routing {
52 Any,
54 Controller,
56 Coordinator(CoordinatorKind),
58 Specific(BrokerSelector),
60}
61
62pub const fn routing(api_key: ApiKey) -> Routing {
64 match api_key {
65 ApiKey::CreateTopics
69 | ApiKey::DeleteTopics
70 | ApiKey::CreatePartitions
71 | ApiKey::AlterPartitionReassignments
72 | ApiKey::ListPartitionReassignments
73 | ApiKey::ElectLeaders
74 | ApiKey::UpdateFeatures => Routing::Controller,
75
76 ApiKey::OffsetCommit
78 | ApiKey::OffsetFetch
79 | ApiKey::OffsetDelete
80 | ApiKey::JoinGroup
81 | ApiKey::Heartbeat
82 | ApiKey::LeaveGroup
83 | ApiKey::SyncGroup
84 | ApiKey::DescribeGroups
85 | ApiKey::DeleteGroups
86 | ApiKey::ConsumerGroupDescribe
87 | ApiKey::ConsumerGroupHeartbeat
88 | ApiKey::ShareGroupDescribe
89 | ApiKey::ShareGroupHeartbeat
90 | ApiKey::DescribeShareGroupOffsets
91 | ApiKey::AlterShareGroupOffsets
92 | ApiKey::DeleteShareGroupOffsets
93 | ApiKey::TxnOffsetCommit => Routing::Coordinator(CoordinatorKind::Group),
94
95 ApiKey::InitProducerId
97 | ApiKey::AddPartitionsToTxn
98 | ApiKey::AddOffsetsToTxn
99 | ApiKey::EndTxn
100 | ApiKey::DescribeTransactions => Routing::Coordinator(CoordinatorKind::Transaction),
101
102 ApiKey::DescribeLogDirs | ApiKey::AlterReplicaLogDirs | ApiKey::DescribeProducers => {
105 Routing::Specific(BrokerSelector::Caller)
106 }
107
108 ApiKey::Produce | ApiKey::Fetch | ApiKey::ListOffsets | ApiKey::OffsetForLeaderEpoch => {
110 Routing::Specific(BrokerSelector::PartitionLeader)
111 }
112
113 _ => Routing::Any,
116 }
117}
118
119#[cfg(test)]
120mod tests {
121 use super::*;
122
123 #[test]
124 fn the_four_classes_from_claude_md() {
125 for key in [
127 ApiKey::CreateTopics,
128 ApiKey::DeleteTopics,
129 ApiKey::CreatePartitions,
130 ApiKey::AlterPartitionReassignments,
131 ApiKey::ElectLeaders,
132 ApiKey::UpdateFeatures,
133 ] {
134 assert_eq!(routing(key), Routing::Controller, "{key}");
135 }
136
137 assert_eq!(
139 routing(ApiKey::OffsetFetch),
140 Routing::Coordinator(CoordinatorKind::Group)
141 );
142 assert_eq!(
143 routing(ApiKey::DescribeTransactions),
144 Routing::Coordinator(CoordinatorKind::Transaction)
145 );
146
147 assert_eq!(
149 routing(ApiKey::DescribeLogDirs),
150 Routing::Specific(BrokerSelector::Caller)
151 );
152 assert_eq!(
153 routing(ApiKey::DescribeProducers),
154 Routing::Specific(BrokerSelector::Caller)
155 );
156
157 for key in [
159 ApiKey::DescribeConfigs,
160 ApiKey::DescribeAcls,
161 ApiKey::ListGroups,
162 ApiKey::Metadata,
163 ] {
164 assert_eq!(routing(key), Routing::Any, "{key}");
165 }
166 }
167
168 #[test]
169 fn the_read_path_goes_to_the_leader() {
170 for key in [ApiKey::Fetch, ApiKey::ListOffsets, ApiKey::Produce] {
171 assert_eq!(
172 routing(key),
173 Routing::Specific(BrokerSelector::PartitionLeader),
174 "{key}"
175 );
176 }
177 }
178
179 #[test]
180 fn share_and_consumer_group_apis_route_like_classic_groups() {
181 for key in [
185 ApiKey::ConsumerGroupDescribe,
186 ApiKey::ShareGroupDescribe,
187 ApiKey::DescribeShareGroupOffsets,
188 ] {
189 assert_eq!(
190 routing(key),
191 Routing::Coordinator(CoordinatorKind::Group),
192 "{key}"
193 );
194 }
195 }
196
197 #[test]
198 fn an_api_key_this_build_cannot_name_routes_anywhere() {
199 assert_eq!(routing(ApiKey::Unknown(9_999)), Routing::Any);
202 }
203
204 #[test]
205 fn find_coordinator_key_types_match_the_protocol() {
206 assert_eq!(CoordinatorKind::Group.key_type(), 0);
207 assert_eq!(CoordinatorKind::Transaction.key_type(), 1);
208 }
209}