use std::collections::{BTreeMap, BTreeSet, HashMap, VecDeque};
use bytes::Bytes;
use super::wire;
use crate::error::ErrorCode;
use crate::protocol::ApiKey;
pub(crate) const PRODUCER_STATE_BATCHES: usize = 5;
#[derive(Debug, Clone)]
pub struct BrokerNode {
pub node_id: i32,
pub host: String,
pub port: i32,
pub rack: Option<String>,
pub online: bool,
}
#[derive(Debug, Clone)]
pub struct PartitionState {
pub leader: i32,
pub leader_epoch: i32,
pub replicas: Vec<i32>,
pub isr: Vec<i32>,
pub log: Vec<Bytes>,
pub log_start_offset: i64,
pub next_offset: i64,
pub open_transactions: HashMap<i64, i64>,
pub producers: HashMap<i64, ProducerEntry>,
pub aborted_transactions: Vec<(i64, i64, i64)>,
}
impl PartitionState {
fn new(leader: i32) -> Self {
Self {
leader,
leader_epoch: 0,
replicas: vec![leader],
isr: vec![leader],
log: Vec::new(),
log_start_offset: 0,
next_offset: 0,
open_transactions: HashMap::new(),
producers: HashMap::new(),
aborted_transactions: Vec::new(),
}
}
pub(crate) fn last_stable_offset(&self) -> i64 {
self.open_transactions
.values()
.copied()
.min()
.unwrap_or(self.next_offset)
}
pub(crate) fn check_sequence(
&self,
producer_id: i64,
producer_epoch: i16,
first_sequence: i32,
record_count: i32,
pre_kip360: bool,
) -> SequenceCheck {
let last_sequence = last_sequence(first_sequence, record_count);
let Some(entry) = self.producers.get(&producer_id) else {
return if pre_kip360 && first_sequence != 0 {
SequenceCheck::Reject(ErrorCode::UnknownProducerId)
} else {
SequenceCheck::Append
};
};
if producer_epoch < entry.epoch {
return SequenceCheck::Reject(ErrorCode::InvalidProducerEpoch);
}
if producer_epoch > entry.epoch {
return if first_sequence == 0 {
SequenceCheck::Append
} else {
SequenceCheck::Reject(ErrorCode::OutOfOrderSequenceNumber)
};
}
if let Some(duplicate) = entry
.batches
.iter()
.find(|b| b.first_sequence == first_sequence && b.last_sequence == last_sequence)
{
return SequenceCheck::Duplicate(duplicate.base_offset);
}
let in_sequence = match entry.batches.back() {
Some(last) => next_sequence(last.last_sequence) == first_sequence,
None => first_sequence == 0,
};
if in_sequence {
SequenceCheck::Append
} else {
SequenceCheck::Reject(ErrorCode::OutOfOrderSequenceNumber)
}
}
pub(crate) fn record_batch(
&mut self,
producer_id: i64,
producer_epoch: i16,
first_sequence: i32,
record_count: i32,
base_offset: i64,
) {
let entry = self.producers.entry(producer_id).or_insert(ProducerEntry {
epoch: producer_epoch,
batches: VecDeque::new(),
});
if entry.epoch != producer_epoch {
entry.epoch = producer_epoch;
entry.batches.clear();
}
entry.batches.push_back(BatchMetadata {
first_sequence,
last_sequence: last_sequence(first_sequence, record_count),
base_offset,
});
while entry.batches.len() > PRODUCER_STATE_BATCHES {
entry.batches.pop_front();
}
}
pub(crate) fn append_marker(
&mut self,
committed: bool,
producer_id: i64,
producer_epoch: i16,
) -> i64 {
let marker = wire::control_batch(committed, producer_id, producer_epoch);
let marker_offset = self.append(&marker);
let first_offset = self.open_transactions.remove(&producer_id);
if !committed && let Some(first) = first_offset {
self.aborted_transactions
.push((producer_id, first, marker_offset));
}
let entry = self.producers.entry(producer_id).or_insert(ProducerEntry {
epoch: producer_epoch,
batches: VecDeque::new(),
});
if producer_epoch > entry.epoch {
entry.epoch = producer_epoch;
entry.batches.clear();
}
marker_offset
}
pub(crate) fn append(&mut self, batch: &Bytes) -> i64 {
let base_offset = self.next_offset;
let count = wire::batch_record_count(batch).unwrap_or(0);
self.log
.push(wire::stamp_batch(batch, base_offset, self.leader_epoch));
self.next_offset += count;
base_offset
}
pub(crate) fn read_from(&self, fetch_offset: i64) -> Bytes {
self.read_range(fetch_offset, i64::MAX)
}
pub(crate) fn read_range(&self, fetch_offset: i64, limit: i64) -> Bytes {
let mut out = Vec::new();
for batch in &self.log {
let base = wire::batch_base_offset(batch).unwrap_or(0);
let count = wire::batch_record_count(batch).unwrap_or(0);
if base + count > fetch_offset && base + count <= limit {
out.extend_from_slice(batch);
}
}
Bytes::from(out)
}
pub(crate) fn aborted_transactions_from(&self, fetch_offset: i64) -> Vec<(i64, i64)> {
self.aborted_transactions
.iter()
.filter(|(_, _, marker_offset)| *marker_offset >= fetch_offset)
.map(|(producer_id, first_offset, _)| (*producer_id, *first_offset))
.collect()
}
}
fn last_sequence(first_sequence: i32, record_count: i32) -> i32 {
let delta = record_count.max(1) - 1;
if first_sequence > i32::MAX - delta {
delta - (i32::MAX - first_sequence) - 1
} else {
first_sequence + delta
}
}
fn next_sequence(last_sequence: i32) -> i32 {
if last_sequence == i32::MAX {
0
} else {
last_sequence + 1
}
}
#[derive(Debug, Clone, Default)]
pub struct ProducerEntry {
pub epoch: i16,
pub batches: VecDeque<BatchMetadata>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct BatchMetadata {
pub first_sequence: i32,
pub last_sequence: i32,
pub base_offset: i64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum SequenceCheck {
Append,
Duplicate(i64),
Reject(ErrorCode),
}
#[derive(Debug, Clone)]
pub struct TopicState {
pub topic_id: [u8; 16],
pub partitions: Vec<PartitionState>,
}
#[derive(Debug, Clone, Default)]
pub struct BrokerTransaction {
pub producer_id: i64,
pub producer_epoch: i16,
pub last_producer_epoch: i16,
pub status: TxnStatus,
pub transaction_timeout_ms: i32,
pub partitions: Vec<(String, i32)>,
pub staged_offsets: HashMap<String, HashMap<(String, i32), CommittedOffset>>,
}
impl BrokerTransaction {
pub fn is_open(&self) -> bool {
self.status == TxnStatus::Ongoing
}
pub(crate) fn bump_epoch(&mut self) {
self.last_producer_epoch = self.producer_epoch;
self.producer_epoch = self.producer_epoch.saturating_add(1);
}
pub(crate) fn begin(&mut self) {
if !matches!(
self.status,
TxnStatus::PrepareCommit | TxnStatus::PrepareAbort
) {
self.status = TxnStatus::Ongoing;
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum TxnStatus {
#[default]
Empty,
Ongoing,
PrepareCommit,
PrepareAbort,
CompleteCommit,
CompleteAbort,
}
#[derive(Debug, Clone)]
pub struct GroupMember {
pub member_id: String,
pub group_instance_id: Option<String>,
pub metadata: Bytes,
pub client_id: String,
pub client_host: String,
}
#[derive(Debug, Clone)]
pub struct CommittedOffset {
pub offset: i64,
pub leader_epoch: i32,
pub metadata: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct GroupState {
pub generation_id: i32,
pub protocol_type: String,
pub protocol_name: Option<String>,
pub leader: String,
pub members: Vec<GroupMember>,
pub assignments: HashMap<String, Bytes>,
pub offsets: HashMap<(String, i32), CommittedOffset>,
pub member_seq: u32,
pub state: ClassicGroupState,
pub consumer_members: HashMap<String, ConsumerGroupMemberState>,
pub group_epoch: i32,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum ClassicGroupState {
#[default]
Empty,
CompletingRebalance,
Stable,
}
impl ClassicGroupState {
pub fn as_str(self) -> &'static str {
match self {
Self::Empty => "Empty",
Self::CompletingRebalance => "CompletingRebalance",
Self::Stable => "Stable",
}
}
}
#[derive(Debug, Clone, Default)]
pub struct ConsumerGroupMemberState {
pub member_epoch: i32,
pub instance_id: Option<String>,
pub subscribed_topics: Vec<String>,
pub assignment: HashMap<String, Vec<i32>>,
pub owned: HashMap<String, Vec<i32>>,
pub assignment_dirty: bool,
pub last_heartbeat: Option<tokio::time::Instant>,
}
#[derive(Debug, Clone, Default)]
pub struct ShareGroupState {
pub group_epoch: i32,
pub members: HashMap<String, ShareMemberState>,
pub partitions: HashMap<(String, i32), SharePartitionState>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct ShareSessionClose {
pub api_key: ApiKey,
pub node_id: i32,
pub group_id: String,
pub member_id: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct ListOffsetsLookup {
pub node_id: i32,
pub api_version: i16,
pub topic: String,
pub partition: i32,
pub timestamp: i64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct LeaveGroupMemberSeen {
pub group_id: String,
pub member_id: String,
pub group_instance_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct ConsumerGroupHeartbeatSeen {
pub group_id: String,
pub member_id: String,
pub member_epoch: i32,
pub instance_id: Option<String>,
pub server_assignor: Option<String>,
pub full: bool,
}
#[derive(Debug, Clone, Default)]
pub struct ShareSession {
pub epoch: i32,
pub partitions: BTreeSet<(String, i32)>,
}
impl ShareSession {
pub(crate) fn advance(&mut self) {
self.epoch = if self.epoch == i32::MAX {
1
} else {
self.epoch + 1
};
}
}
#[derive(Debug, Clone, Default)]
pub struct ShareMemberState {
pub member_epoch: i32,
pub subscribed_topics: Vec<String>,
pub assignment: HashMap<String, Vec<i32>>,
pub assignment_dirty: bool,
}
#[derive(Debug, Clone, Default)]
pub struct SharePartitionState {
pub start_offset: i64,
pub acquired: BTreeMap<i64, String>,
pub archived: BTreeSet<i64>,
pub delivery_counts: HashMap<i64, i16>,
pub acknowledgements: BTreeMap<i64, Vec<ShareAckType>>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum ShareAckType {
Gap,
Accept,
Release,
Reject,
Renew,
}
impl ShareAckType {
fn from_wire(code: i8) -> Option<Self> {
Some(match code {
0 => Self::Gap,
1 => Self::Accept,
2 => Self::Release,
3 => Self::Reject,
4 => Self::Renew,
_ => return None,
})
}
}
impl ShareGroupState {
pub(crate) fn release_member(&mut self, member_id: &str) {
for partition in self.partitions.values_mut() {
partition.acquired.retain(|_, holder| holder != member_id);
}
}
}
impl SharePartitionState {
pub(crate) fn is_available(&self, offset: i64) -> bool {
offset >= self.start_offset
&& !self.acquired.contains_key(&offset)
&& !self.archived.contains(&offset)
}
pub(crate) fn acquire(&mut self, offset: i64, member_id: &str) -> i16 {
self.acquired.insert(offset, member_id.to_string());
let count = self.delivery_counts.entry(offset).or_insert(0);
*count = count.saturating_add(1);
*count
}
pub(crate) fn held_by(&self, first: i64, last: i64, member_id: &str) -> bool {
(first..=last).all(|offset| {
self.acquired
.get(&offset)
.is_some_and(|holder| holder == member_id)
})
}
pub(crate) fn acknowledge(&mut self, offset: i64, acknowledge_type: i8) {
if let Some(ack) = ShareAckType::from_wire(acknowledge_type) {
self.acknowledgements.entry(offset).or_default().push(ack);
}
match acknowledge_type {
0 | 1 | 3 => {
self.acquired.remove(&offset);
self.delivery_counts.remove(&offset);
self.archived.insert(offset);
while self.archived.remove(&self.start_offset) {
self.start_offset += 1;
}
}
2 => {
self.acquired.remove(&offset);
}
_ => {}
}
}
}
#[derive(Debug, Clone, Default)]
pub struct StreamsGroupState {
pub group_state: String,
pub group_epoch: i32,
pub assignment_epoch: i32,
pub topology_epoch: Option<i32>,
pub subtopologies: Option<Vec<String>>,
pub members: Vec<StreamsMemberState>,
}
#[derive(Debug, Clone, Default)]
pub struct StreamsMemberState {
pub member_id: String,
pub member_epoch: i32,
pub topology_epoch: i32,
pub process_id: String,
pub user_endpoint: Option<(String, u16)>,
pub active_tasks: Vec<(String, Vec<i32>)>,
pub target_active_tasks: Vec<(String, Vec<i32>)>,
}
#[derive(Debug)]
pub struct ClusterState {
pub cluster_id: String,
pub brokers: Vec<BrokerNode>,
pub controller_id: i32,
pub topics: HashMap<String, TopicState>,
pub groups: HashMap<String, GroupState>,
pub share_groups: HashMap<String, ShareGroupState>,
pub streams_groups: HashMap<String, StreamsGroupState>,
pub group_coordinators: HashMap<String, i32>,
pub txn_coordinators: HashMap<String, i32>,
pub auto_create_topics: bool,
pub default_partitions: i32,
pub next_producer_id: i64,
pub transactions: HashMap<String, BrokerTransaction>,
pub transaction_max_timeout_ms: i32,
pub idempotence: bool,
pub hold_transaction_markers: bool,
pub throttle_time_ms: HashMap<ApiKey, i32>,
pub share_sessions: HashMap<(i32, String, String), ShareSession>,
pub finalized_features: HashMap<String, i16>,
pub finalized_features_epoch: i64,
pub api_version_overrides: HashMap<ApiKey, (i16, i16)>,
pub share_session_closes: Vec<ShareSessionClose>,
pub list_offsets_lookups: Vec<ListOffsetsLookup>,
pub leave_group_members: Vec<LeaveGroupMemberSeen>,
pub consumer_group_heartbeats: Vec<ConsumerGroupHeartbeatSeen>,
pub sasl_plain: Option<(String, String)>,
pub telemetry: Option<TelemetrySubscription>,
pub telemetry_pushes: Vec<TelemetryPush>,
topic_id_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct TelemetrySubscription {
pub requested_metrics: Vec<String>,
pub push_interval: std::time::Duration,
pub delta_temporality: bool,
pub accepted_compression_types: Vec<i8>,
pub client_instance_id: [u8; 16],
pub subscription_id: i32,
}
impl TelemetrySubscription {
pub fn new(
requested_metrics: impl IntoIterator<Item = impl Into<String>>,
push_interval: std::time::Duration,
) -> Self {
Self {
requested_metrics: requested_metrics.into_iter().map(Into::into).collect(),
push_interval,
delta_temporality: false,
accepted_compression_types: vec![0],
client_instance_id: *b"krafka-fake-inst",
subscription_id: 1,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct TelemetryPush {
pub client_instance_id: [u8; 16],
pub subscription_id: i32,
pub terminating: bool,
pub compression_type: i8,
pub metrics: bytes::Bytes,
}
impl ClusterState {
pub(crate) fn new(broker_count: usize) -> Self {
let brokers = (0..broker_count)
.map(|i| BrokerNode {
node_id: i as i32,
host: "127.0.0.1".to_string(),
port: 0,
rack: None,
online: true,
})
.collect();
Self {
cluster_id: "krafka-fake-cluster".to_string(),
brokers,
controller_id: 0,
topics: HashMap::new(),
groups: HashMap::new(),
share_groups: HashMap::new(),
streams_groups: HashMap::new(),
finalized_features: HashMap::new(),
finalized_features_epoch: 0,
api_version_overrides: HashMap::new(),
share_session_closes: Vec::new(),
list_offsets_lookups: Vec::new(),
leave_group_members: Vec::new(),
consumer_group_heartbeats: Vec::new(),
sasl_plain: None,
telemetry: None,
telemetry_pushes: Vec::new(),
group_coordinators: HashMap::new(),
txn_coordinators: HashMap::new(),
auto_create_topics: true,
default_partitions: 1,
next_producer_id: 1000,
transactions: HashMap::new(),
transaction_max_timeout_ms: 900_000,
idempotence: true,
hold_transaction_markers: false,
throttle_time_ms: HashMap::new(),
share_sessions: HashMap::new(),
topic_id_seq: 1,
}
}
pub fn default_coordinator(&self) -> i32 {
self.brokers
.iter()
.find(|b| b.online)
.map(|b| b.node_id)
.unwrap_or(-1)
}
pub fn group_coordinator(&self, group_id: &str) -> i32 {
self.group_coordinators
.get(group_id)
.copied()
.unwrap_or_else(|| self.default_coordinator())
}
pub fn txn_coordinator(&self, transactional_id: &str) -> i32 {
self.txn_coordinators
.get(transactional_id)
.copied()
.unwrap_or_else(|| self.default_coordinator())
}
pub fn broker(&self, node_id: i32) -> Option<&BrokerNode> {
self.brokers.iter().find(|b| b.node_id == node_id)
}
pub fn create_topic(&mut self, name: &str, partitions: i32) -> bool {
if self.topics.contains_key(name) {
return false;
}
let online: Vec<i32> = self
.brokers
.iter()
.filter(|b| b.online)
.map(|b| b.node_id)
.collect();
let partition_states = (0..partitions.max(1))
.map(|i| {
let leader = online
.get(i as usize % online.len().max(1))
.copied()
.unwrap_or(0);
PartitionState::new(leader)
})
.collect();
let mut topic_id = [0u8; 16];
topic_id[8..].copy_from_slice(&self.topic_id_seq.to_be_bytes());
self.topic_id_seq += 1;
self.topics.insert(
name.to_string(),
TopicState {
topic_id,
partitions: partition_states,
},
);
true
}
pub fn add_partitions(&mut self, name: &str, partitions: i32) -> usize {
let online: Vec<i32> = self
.brokers
.iter()
.filter(|b| b.online)
.map(|b| b.node_id)
.collect();
let Some(topic) = self.topics.get_mut(name) else {
return 0;
};
let existing = topic.partitions.len();
let target = partitions.max(0) as usize;
if target <= existing {
return 0;
}
for i in existing..target {
let leader = online.get(i % online.len().max(1)).copied().unwrap_or(0);
topic.partitions.push(PartitionState::new(leader));
}
target - existing
}
pub fn partition_mut(&mut self, topic: &str, partition: i32) -> Option<&mut PartitionState> {
self.topics
.get_mut(topic)
.and_then(|t| t.partitions.get_mut(usize::try_from(partition).ok()?))
}
pub fn partition(&self, topic: &str, partition: i32) -> Option<&PartitionState> {
self.topics
.get(topic)
.and_then(|t| t.partitions.get(usize::try_from(partition).ok()?))
}
pub fn allocate_producer_id(&mut self) -> (i64, i16) {
let id = self.next_producer_id;
self.next_producer_id += 1;
(id, 0)
}
pub fn delete_topic(&mut self, name: &str) -> bool {
if self.topics.remove(name).is_none() {
return false;
}
for group in self.share_groups.values_mut() {
group.partitions.retain(|(topic, _), _| topic != name);
}
for session in self.share_sessions.values_mut() {
session.partitions.retain(|(topic, _)| topic != name);
}
true
}
pub(crate) fn throttle(&self, api_key: ApiKey) -> i32 {
self.throttle_time_ms.get(&api_key).copied().unwrap_or(0)
}
pub(crate) fn pre_kip360(&self) -> bool {
self.api_version_overrides
.get(&ApiKey::InitProducerId)
.is_some_and(|&(_, max)| max < 3)
}
pub(crate) fn transaction_for_producer(&self, producer_id: i64) -> Option<String> {
self.transactions
.iter()
.find(|(_, t)| t.producer_id == producer_id)
.map(|(id, _)| id.clone())
}
pub(crate) fn end_transaction(
&mut self,
transactional_id: &str,
committed: bool,
bump_epoch: bool,
) {
let Some(txn) = self.transactions.get_mut(transactional_id) else {
return;
};
if bump_epoch {
txn.bump_epoch();
}
txn.status = if committed {
TxnStatus::PrepareCommit
} else {
TxnStatus::PrepareAbort
};
if !self.hold_transaction_markers {
self.write_transaction_markers(transactional_id);
}
}
pub(crate) fn fence_transaction(&mut self, transactional_id: &str) {
if let Some(txn) = self.transactions.get_mut(transactional_id)
&& txn.status == TxnStatus::Ongoing
{
txn.bump_epoch();
txn.last_producer_epoch = -1;
txn.status = TxnStatus::PrepareAbort;
if !self.hold_transaction_markers {
self.write_transaction_markers(transactional_id);
}
}
}
pub(crate) fn write_transaction_markers(&mut self, transactional_id: &str) {
let Some(txn) = self.transactions.get_mut(transactional_id) else {
return;
};
let committed = match txn.status {
TxnStatus::PrepareCommit => true,
TxnStatus::PrepareAbort => false,
_ => return,
};
txn.status = if committed {
TxnStatus::CompleteCommit
} else {
TxnStatus::CompleteAbort
};
let partitions = std::mem::take(&mut txn.partitions);
let staged = std::mem::take(&mut txn.staged_offsets);
let (producer_id, producer_epoch) = (txn.producer_id, txn.producer_epoch);
for (topic, partition) in &partitions {
if let Some(p) = self.partition_mut(topic, *partition) {
p.append_marker(committed, producer_id, producer_epoch);
}
}
if committed {
for (group_id, offsets) in staged {
let group = self.groups.entry(group_id).or_default();
group.offsets.extend(offsets);
}
}
}
pub fn next_member_id(&mut self, group_id: &str) -> String {
let group = self.groups.entry(group_id.to_string()).or_default();
group.member_seq += 1;
format!("krafka-fake-member-{}", group.member_seq)
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::protocol::{Record, RecordBatch};
fn batch(values: &[&str]) -> Bytes {
let mut b = RecordBatch::new();
b.records = values
.iter()
.enumerate()
.map(|(i, v)| {
Record::new(None, Some(Bytes::copy_from_slice(v.as_bytes())))
.with_offset_delta(i as i32)
})
.collect();
b.encode().unwrap()
}
#[test]
fn appending_assigns_consecutive_offsets() {
let mut p = PartitionState::new(0);
assert_eq!(p.append(&batch(&["a", "b"])), 0);
assert_eq!(p.next_offset, 2);
assert_eq!(p.append(&batch(&["c"])), 2);
assert_eq!(p.next_offset, 3);
}
#[test]
fn reading_returns_whole_batches_that_span_the_fetch_offset() {
let mut p = PartitionState::new(0);
p.append(&batch(&["a", "b"])); p.append(&batch(&["c"]));
assert!(p.read_from(0).len() > p.read_from(2).len());
assert!(!p.read_from(1).is_empty(), "offset 1 sits inside batch one");
assert!(p.read_from(3).is_empty(), "nothing at or beyond the end");
}
#[test]
fn coordinators_default_to_the_lowest_online_broker_and_follow_overrides() {
let mut state = ClusterState::new(3);
assert_eq!(state.group_coordinator("g"), 0);
state.brokers[0].online = false;
assert_eq!(state.group_coordinator("g"), 1);
state.group_coordinators.insert("g".to_string(), 2);
assert_eq!(state.group_coordinator("g"), 2);
}
#[test]
fn topic_creation_spreads_leadership_over_online_brokers() {
let mut state = ClusterState::new(3);
assert!(state.create_topic("t", 3));
assert!(!state.create_topic("t", 3), "re-creation is a no-op");
let leaders: Vec<i32> = state.topics["t"]
.partitions
.iter()
.map(|p| p.leader)
.collect();
assert_eq!(leaders, vec![0, 1, 2]);
}
}