krafka 0.14.0

A pure Rust, async-native Apache Kafka client
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
//! Mutable cluster state behind the fake broker.
//!
//! Everything a test can manipulate — brokers, topic leadership, group and
//! transaction coordinators, committed offsets, the in-memory logs — lives
//! here behind a single lock. Handlers take that lock for the duration of one
//! request, which is what makes request handling serialisable and the
//! resulting behaviour reproducible.

use std::collections::HashMap;

use bytes::Bytes;

use super::wire;

/// A broker in the fake cluster's metadata.
#[derive(Debug, Clone)]
pub struct BrokerNode {
    /// Broker ID as advertised in Metadata responses.
    pub node_id: i32,
    /// Advertised host.
    pub host: String,
    /// Advertised port.
    pub port: i32,
    /// Advertised rack, if any.
    pub rack: Option<String>,
    /// Whether the broker is presented as reachable.
    ///
    /// A broker marked down is still listed in Metadata (real Kafka keeps
    /// listing brokers it has lost contact with) but is never chosen as a
    /// leader or coordinator by the cluster-manipulation helpers.
    pub online: bool,
}

/// One partition's log and leadership.
#[derive(Debug, Clone)]
pub struct PartitionState {
    /// Broker ID currently leading this partition.
    pub leader: i32,
    /// Current leader epoch, bumped on every leadership change.
    pub leader_epoch: i32,
    /// Replica set.
    pub replicas: Vec<i32>,
    /// In-sync replica set.
    pub isr: Vec<i32>,
    /// Stored record batches, already stamped with their broker-assigned
    /// base offsets.
    pub log: Vec<Bytes>,
    /// Offset of the first record still retained.
    pub log_start_offset: i64,
    /// Offset that the next appended record will receive. Because the fake
    /// broker acknowledges writes immediately, this doubles as the high
    /// watermark.
    pub next_offset: 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,
        }
    }

    /// Append a producer's record batch, stamping it with the offset it was
    /// assigned, and return that base 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
    }

    /// Concatenate every stored batch whose base offset is at or after
    /// `fetch_offset`.
    ///
    /// Batches are returned whole: a fetch landing in the middle of a batch
    /// gets the entire batch, exactly as a real broker does, leaving the
    /// client to discard the records below its requested offset.
    pub(crate) fn read_from(&self, fetch_offset: 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 {
                out.extend_from_slice(batch);
            }
        }
        Bytes::from(out)
    }
}

/// A topic and its partitions.
#[derive(Debug, Clone)]
pub struct TopicState {
    /// Topic UUID. Only surfaced on API versions that carry it.
    pub topic_id: [u8; 16],
    /// Partitions, indexed by partition number.
    pub partitions: Vec<PartitionState>,
}

/// A member of a consumer group.
#[derive(Debug, Clone)]
pub struct GroupMember {
    /// Broker-assigned member ID.
    pub member_id: String,
    /// Static membership ID (KIP-345), if the member supplied one.
    pub group_instance_id: Option<String>,
    /// Subscription metadata the member sent in JoinGroup.
    pub metadata: Bytes,
}

/// A committed offset for one topic-partition in one group.
#[derive(Debug, Clone)]
pub struct CommittedOffset {
    /// The committed offset.
    pub offset: i64,
    /// Leader epoch recorded alongside the commit, or `-1`.
    pub leader_epoch: i32,
    /// Opaque metadata attached to the commit.
    pub metadata: Option<String>,
}

/// A consumer group.
#[derive(Debug, Clone, Default)]
pub struct GroupState {
    /// Current generation, incremented on each completed join.
    pub generation_id: i32,
    /// Protocol type, e.g. `consumer`.
    pub protocol_type: String,
    /// Protocol the broker selected for the generation.
    pub protocol_name: Option<String>,
    /// Member ID of the group leader.
    pub leader: String,
    /// Current members.
    pub members: Vec<GroupMember>,
    /// Assignments distributed by SyncGroup, keyed by member ID.
    pub assignments: HashMap<String, Bytes>,
    /// Committed offsets, keyed by `(topic, partition)`.
    pub offsets: HashMap<(String, i32), CommittedOffset>,
    /// Counter behind generated member IDs.
    pub member_seq: u32,
    /// KIP-848 members, keyed by client-generated member ID.
    ///
    /// Separate from [`Self::members`], which models the classic
    /// JoinGroup/SyncGroup protocol. The two protocols have different member
    /// identity and epoch rules, and conflating them in one map made it
    /// impossible to model either faithfully.
    pub consumer_members: HashMap<String, ConsumerGroupMemberState>,
    /// Epoch the whole group is on. Bumped whenever the set of members or
    /// their subscriptions changes, which is what forces reconciliation.
    pub group_epoch: i32,
}

/// One KIP-848 member's coordinator-side state.
#[derive(Debug, Clone, Default)]
pub struct ConsumerGroupMemberState {
    /// The epoch this member is currently on. `0` until its first assignment.
    pub member_epoch: i32,
    /// Static membership ID, if the member supplied one.
    pub instance_id: Option<String>,
    /// Topics the member last told the coordinator it subscribes to.
    pub subscribed_topics: Vec<String>,
    /// Partitions the coordinator has assigned, keyed by topic name.
    pub assignment: HashMap<String, Vec<i32>>,
    /// Partitions the member last *reported* owning.
    ///
    /// Distinct from [`Self::assignment`]: the coordinator may have granted
    /// partitions the member has not acknowledged yet, and may be waiting for
    /// the member to release partitions it still holds. Reconciliation is
    /// exactly the gap between these two fields.
    pub owned: HashMap<String, Vec<i32>>,
    /// Whether [`Self::assignment`] changed since it was last sent.
    ///
    /// The assignment field is only put on the wire when it moves; a null
    /// assignment means "keep what you have".
    pub assignment_dirty: bool,
}

/// The whole fake cluster.
#[derive(Debug)]
pub struct ClusterState {
    /// Cluster ID reported in Metadata.
    pub cluster_id: String,
    /// Brokers, in advertised order.
    pub brokers: Vec<BrokerNode>,
    /// Broker ID currently acting as controller, or `-1` for none.
    pub controller_id: i32,
    /// Topics, keyed by name.
    pub topics: HashMap<String, TopicState>,
    /// Consumer groups, keyed by group ID.
    pub groups: HashMap<String, GroupState>,
    /// Group coordinator overrides, keyed by group ID. Groups without an entry
    /// resolve to [`ClusterState::default_coordinator`].
    pub group_coordinators: HashMap<String, i32>,
    /// Transaction coordinator overrides, keyed by transactional ID.
    pub txn_coordinators: HashMap<String, i32>,
    /// Whether an unknown topic is created on first reference.
    pub auto_create_topics: bool,
    /// Partition count given to auto-created topics.
    pub default_partitions: i32,
    /// Counter behind allocated producer IDs.
    pub next_producer_id: i64,
    /// Producer epochs, keyed by producer ID.
    pub producer_epochs: HashMap<i64, i16>,
    /// Counter behind generated topic UUIDs.
    topic_id_seq: u64,
}

impl ClusterState {
    /// Build a cluster with `broker_count` brokers, none of them bound to a
    /// listener yet. [`super::FakeBroker`] fills in the real host and port once
    /// the sockets are open.
    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(),
            group_coordinators: HashMap::new(),
            txn_coordinators: HashMap::new(),
            auto_create_topics: true,
            default_partitions: 1,
            next_producer_id: 1000,
            producer_epochs: HashMap::new(),
            topic_id_seq: 1,
        }
    }

    /// The broker every group and transaction resolves to unless a test has
    /// moved it: the lowest-numbered online broker.
    pub fn default_coordinator(&self) -> i32 {
        self.brokers
            .iter()
            .find(|b| b.online)
            .map(|b| b.node_id)
            .unwrap_or(-1)
    }

    /// Resolve the coordinator for a consumer group.
    pub fn group_coordinator(&self, group_id: &str) -> i32 {
        self.group_coordinators
            .get(group_id)
            .copied()
            .unwrap_or_else(|| self.default_coordinator())
    }

    /// Resolve the coordinator for a transactional ID.
    pub fn txn_coordinator(&self, transactional_id: &str) -> i32 {
        self.txn_coordinators
            .get(transactional_id)
            .copied()
            .unwrap_or_else(|| self.default_coordinator())
    }

    /// Look up a broker by ID.
    pub fn broker(&self, node_id: i32) -> Option<&BrokerNode> {
        self.brokers.iter().find(|b| b.node_id == node_id)
    }

    /// Create a topic with `partitions` partitions, spreading leadership
    /// round-robin over the online brokers. Existing topics are left alone.
    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
    }

    /// Mutable access to one partition.
    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()?))
    }

    /// Read-only access to one partition.
    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()?))
    }

    /// Allocate a fresh producer ID with epoch 0.
    pub fn allocate_producer_id(&mut self) -> (i64, i16) {
        let id = self.next_producer_id;
        self.next_producer_id += 1;
        self.producer_epochs.insert(id, 0);
        (id, 0)
    }

    /// Generate the next member ID for a group.
    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);
    }

    /// A fetch that lands inside a batch must still receive the whole batch,
    /// matching real broker behaviour.
    #[test]
    fn reading_returns_whole_batches_that_span_the_fetch_offset() {
        let mut p = PartitionState::new(0);
        p.append(&batch(&["a", "b"])); // offsets 0..=1
        p.append(&batch(&["c"])); // offset 2

        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]);
    }
}