use std::collections::HashMap;
use kafka_protocol::messages::MetadataResponse;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BrokerAddr {
pub node_id: i32,
pub host: String,
pub port: i32,
}
impl BrokerAddr {
#[must_use]
pub fn addr(&self) -> String {
format!("{}:{}", self.host, self.port)
}
}
pub type PartitionKey = (String, i32);
#[derive(Debug, Default)]
pub struct Metadata {
brokers: HashMap<i32, BrokerAddr>,
leaders: HashMap<PartitionKey, i32>,
partition_counts: HashMap<String, i32>,
controller: Option<i32>,
}
impl Metadata {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn update(&mut self, resp: &MetadataResponse) {
if resp.controller_id.0 >= 0 {
self.controller = Some(resp.controller_id.0);
}
for broker in &resp.brokers {
self.brokers.insert(
broker.node_id.0,
BrokerAddr {
node_id: broker.node_id.0,
host: broker.host.to_string(),
port: broker.port,
},
);
}
for topic in &resp.topics {
if topic.error_code != 0 {
continue;
}
let Some(name) = topic.name.as_ref() else {
continue;
};
self.partition_counts.insert(
name.0.to_string(),
i32::try_from(topic.partitions.len()).unwrap_or(i32::MAX),
);
for partition in &topic.partitions {
if partition.error_code != 0 || partition.leader_id.0 < 0 {
continue;
}
self.leaders.insert(
(name.0.to_string(), partition.partition_index),
partition.leader_id.0,
);
}
}
}
#[must_use]
pub fn leader_for(&self, topic: &str, partition: i32) -> Option<&BrokerAddr> {
let node = self.leaders.get(&(topic.to_owned(), partition))?;
self.brokers.get(node)
}
#[must_use]
pub fn broker(&self, node_id: i32) -> Option<&BrokerAddr> {
self.brokers.get(&node_id)
}
#[must_use]
pub fn controller(&self) -> Option<&BrokerAddr> {
self.brokers.get(&self.controller?)
}
pub fn invalidate_controller(&mut self) {
self.controller = None;
}
pub fn invalidate_partition(&mut self, topic: &str, partition: i32) {
self.leaders.remove(&(topic.to_owned(), partition));
}
#[must_use]
pub fn partition_count(&self, topic: &str) -> i32 {
self.partition_counts.get(topic).copied().unwrap_or(0)
}
pub fn forget_topic(&mut self, topic: &str) {
self.partition_counts.remove(topic);
self.leaders.retain(|(name, _), _| name != topic);
}
pub fn brokers(&self) -> impl Iterator<Item = &BrokerAddr> {
self.brokers.values()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.brokers.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::*;
use kafka_protocol::messages::metadata_response::{
MetadataResponseBroker, MetadataResponsePartition, MetadataResponseTopic,
};
use kafka_protocol::messages::{BrokerId, TopicName};
use kafka_protocol::protocol::StrBytes;
fn broker(node_id: i32, host: &str, port: i32) -> MetadataResponseBroker {
let mut b = MetadataResponseBroker::default();
b.node_id = BrokerId(node_id);
b.host = StrBytes::from_string(host.to_owned());
b.port = port;
b
}
fn topic(
name: &str,
error_code: i16,
partitions: Vec<MetadataResponsePartition>,
) -> MetadataResponseTopic {
let mut t = MetadataResponseTopic::default();
t.name = Some(TopicName(StrBytes::from_string(name.to_owned())));
t.error_code = error_code;
t.partitions = partitions;
t
}
fn partition(index: i32, leader: i32, error_code: i16) -> MetadataResponsePartition {
let mut p = MetadataResponsePartition::default();
p.partition_index = index;
p.leader_id = BrokerId(leader);
p.error_code = error_code;
p
}
fn response(
brokers: Vec<MetadataResponseBroker>,
topics: Vec<MetadataResponseTopic>,
) -> MetadataResponse {
let mut r = MetadataResponse::default();
r.brokers = brokers;
r.topics = topics;
r
}
#[test]
fn a_leader_is_resolved_to_an_address() {
let mut md = Metadata::new();
md.update(&response(
vec![broker(1, "kafka-1", 9092), broker(2, "kafka-2", 9092)],
vec![topic("t", 0, vec![partition(0, 2, 0)])],
));
assert_eq!(md.leader_for("t", 0).unwrap().addr(), "kafka-2:9092");
assert!(md.leader_for("t", 1).is_none());
assert!(md.leader_for("other", 0).is_none());
}
#[test]
fn an_update_for_one_topic_leaves_the_others_alone() {
let mut md = Metadata::new();
md.update(&response(
vec![broker(1, "kafka-1", 9092)],
vec![topic("a", 0, vec![partition(0, 1, 0)])],
));
md.update(&response(
vec![broker(1, "kafka-1", 9092)],
vec![topic("b", 0, vec![partition(0, 1, 0)])],
));
assert!(md.leader_for("a", 0).is_some(), "topic a was forgotten");
assert!(md.leader_for("b", 0).is_some());
}
#[test]
fn a_moved_leader_is_picked_up() {
let mut md = Metadata::new();
md.update(&response(
vec![broker(1, "kafka-1", 9092), broker(2, "kafka-2", 9092)],
vec![topic("t", 0, vec![partition(0, 1, 0)])],
));
assert_eq!(md.leader_for("t", 0).unwrap().node_id, 1);
md.update(&response(
vec![broker(1, "kafka-1", 9092), broker(2, "kafka-2", 9092)],
vec![topic("t", 0, vec![partition(0, 2, 0)])],
));
assert_eq!(md.leader_for("t", 0).unwrap().node_id, 2);
}
#[test]
fn an_errored_partition_is_not_recorded() {
let mut md = Metadata::new();
md.update(&response(
vec![broker(1, "kafka-1", 9092)],
vec![topic("t", 0, vec![partition(0, 1, 9)])],
));
assert!(md.leader_for("t", 0).is_none());
}
#[test]
fn a_leaderless_partition_is_not_recorded() {
let mut md = Metadata::new();
md.update(&response(
vec![broker(1, "kafka-1", 9092)],
vec![topic("t", 0, vec![partition(0, -1, 0)])],
));
assert!(md.leader_for("t", 0).is_none());
}
#[test]
fn an_errored_topic_still_yields_its_brokers() {
let mut md = Metadata::new();
md.update(&response(
vec![broker(7, "kafka-7", 9092)],
vec![topic("t", 3, vec![partition(0, 7, 0)])],
));
assert!(md.leader_for("t", 0).is_none());
assert_eq!(md.broker(7).unwrap().addr(), "kafka-7:9092");
}
#[test]
fn the_partition_count_follows_metadata() {
let mut md = Metadata::new();
assert_eq!(md.partition_count("t"), 0);
md.update(&response(
vec![broker(1, "kafka-1", 9092)],
vec![topic(
"t",
0,
vec![partition(0, 1, 0), partition(1, 1, 0), partition(2, 1, 0)],
)],
));
assert_eq!(md.partition_count("t"), 3);
assert_eq!(md.partition_count("other"), 0);
}
#[test]
fn invalidating_one_partition_keeps_the_others() {
let mut md = Metadata::new();
md.update(&response(
vec![broker(1, "kafka-1", 9092)],
vec![topic("t", 0, vec![partition(0, 1, 0), partition(1, 1, 0)])],
));
md.invalidate_partition("t", 0);
assert!(md.leader_for("t", 0).is_none());
assert!(md.leader_for("t", 1).is_some());
assert!(
md.broker(1).is_some(),
"the broker itself is still reachable and its connection stays open"
);
}
}