kacrab 0.1.1

A Kafka client for Rust, built from the protocol up.
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
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
#![expect(
    clippy::cast_precision_loss,
    reason = "Metric counts/bytes are coarse observability samples; f64 mantissa loss is \
              acceptable."
)]
//! Kafka-compatible producer metrics, mirroring Kafka's `SenderMetricsRegistry`.
//!
//! The producer's primary metrics facade is the lock-free
//! [`super::super::metrics::ProducerMetrics`] atomic counter set, exposed under Rust-native names.
//! This registry additionally publishes the same measurements under Kafka's metric names and
//! semantics (windowed Rate, Avg, Max) on the `producer-metrics` group plus per-topic
//! instances on the `producer-topic-metrics` group, so applications that query by
//! Kafka metric name see the expected sensors.

use std::{
    collections::{BTreeMap, HashMap},
    sync::Mutex,
    time::{SystemTime, UNIX_EPOCH},
};

use super::registry::{MetricName, Metrics, SensorId};

const CLIENT_GROUP: &str = "producer-metrics";
const TOPIC_GROUP: &str = "producer-topic-metrics";

fn now_ms() -> u64 {
    u64::try_from(
        SystemTime::now()
            .duration_since(UNIX_EPOCH)
            .map_or(0, |elapsed| elapsed.as_millis()),
    )
    .unwrap_or(u64::MAX)
}

/// Kafka-named producer client-level and per-topic metrics.
#[derive(Debug)]
pub(crate) struct SenderMetricsRegistry {
    inner: Mutex<RegistryInner>,
}

#[derive(Debug)]
struct RegistryInner {
    metrics: Metrics,
    client: ClientSensors,
    topics: HashMap<String, TopicSensors>,
}

/// Client-level (`producer-metrics`) sensor handles.
#[derive(Debug)]
struct ClientSensors {
    records_sent: SensorId,
    record_errors: SensorId,
    record_retries: SensorId,
    batch_split: SensorId,
    bytes: SensorId,
    compression_rate: SensorId,
    batch_size: SensorId,
    records_per_request: SensorId,
    request_latency: SensorId,
    record_size: SensorId,
    produce_throttle_time: SensorId,
    record_queue_time: SensorId,
    requests_in_flight: SensorId,
    metadata_age: SensorId,
    buffer_available_bytes: SensorId,
    waiting_threads: SensorId,
    buffer_exhausted: SensorId,
    bufferpool_wait_time: SensorId,
}

/// Per-topic (`producer-topic-metrics`) sensor handles.
#[derive(Debug, Clone, Copy)]
struct TopicSensors {
    records_sent: SensorId,
    bytes: SensorId,
    compression_rate: SensorId,
    record_retries: SensorId,
    record_errors: SensorId,
}

impl Default for SenderMetricsRegistry {
    fn default() -> Self {
        let mut metrics = Metrics::new();
        let client = ClientSensors::register(&mut metrics);
        Self {
            inner: Mutex::new(RegistryInner {
                metrics,
                client,
                topics: HashMap::new(),
            }),
        }
    }
}

impl ClientSensors {
    #[expect(
        clippy::too_many_lines,
        reason = "Registers the full Kafka client-level sensor set in one place."
    )]
    fn register(metrics: &mut Metrics) -> Self {
        let records_sent = meter(
            metrics,
            "records-sent",
            "record-send-rate",
            "The average number of records sent per second.",
            "record-send-total",
            "The total number of records sent.",
        );
        let record_errors = meter(
            metrics,
            "record-errors",
            "record-error-rate",
            "The average per-second number of record sends that resulted in errors.",
            "record-error-total",
            "The total number of record sends that resulted in errors.",
        );
        let record_retries = meter(
            metrics,
            "record-retries",
            "record-retry-rate",
            "The average per-second number of retried record sends.",
            "record-retry-total",
            "The total number of retried record sends.",
        );
        let batch_split = meter(
            metrics,
            "batch-split",
            "batch-split-rate",
            "The average number of batch splits per second.",
            "batch-split-total",
            "The total number of batch splits.",
        );
        let bytes = meter(
            metrics,
            "bytes",
            "byte-rate",
            "The average number of bytes sent per second.",
            "byte-total",
            "The total number of bytes sent.",
        );
        let compression_rate = avg_only(
            metrics,
            "compression-rate",
            "compression-rate-avg",
            "The average compression rate of record batches.",
        );
        let batch_size = avg_max(
            metrics,
            "batch-size",
            "batch-size-avg",
            "The average number of bytes sent per partition per-request.",
            "batch-size-max",
            "The max number of bytes sent per partition per-request.",
        );
        let records_per_request = avg_only(
            metrics,
            "records-per-request",
            "records-per-request-avg",
            "The average number of records per request.",
        );
        let request_latency = avg_max(
            metrics,
            "request-latency",
            "request-latency-avg",
            "The average request latency in ms.",
            "request-latency-max",
            "The maximum request latency in ms.",
        );
        let record_size = avg_max(
            metrics,
            "record-size",
            "record-size-avg",
            "The average record size.",
            "record-size-max",
            "The maximum record size.",
        );
        let produce_throttle_time = avg_max(
            metrics,
            "produce-throttle-time",
            "produce-throttle-time-avg",
            "The average time in ms a request was throttled by a broker.",
            "produce-throttle-time-max",
            "The maximum time in ms a request was throttled by a broker.",
        );
        let record_queue_time = avg_max(
            metrics,
            "record-queue-time",
            "record-queue-time-avg",
            "The average time in ms record batches spent in the send buffer.",
            "record-queue-time-max",
            "The maximum time in ms record batches spent in the send buffer.",
        );
        let requests_in_flight = value_only(
            metrics,
            "requests-in-flight",
            "requests-in-flight",
            "The current number of in-flight requests awaiting a response.",
        );
        let metadata_age = value_only(
            metrics,
            "metadata-age",
            "metadata-age",
            "The age in seconds of the current producer metadata being used.",
        );
        let buffer_available_bytes = value_only(
            metrics,
            "buffer-available-bytes",
            "buffer-available-bytes",
            "The total amount of buffer memory that is not being used.",
        );
        let waiting_threads = value_only(
            metrics,
            "waiting-threads",
            "waiting-threads",
            "The number of user threads blocked waiting for buffer memory to enqueue their \
             records.",
        );
        let buffer_exhausted = meter(
            metrics,
            "buffer-exhausted",
            "buffer-exhausted-rate",
            "The average per-second number of record sends that are blocked on buffer memory \
             exhaustion.",
            "buffer-exhausted-total",
            "The total number of record sends that are blocked on buffer memory exhaustion.",
        );
        let bufferpool_wait_time = avg_max(
            metrics,
            "bufferpool-wait-time",
            "bufferpool-wait-time-avg",
            "The average time in ms an appender waits for space allocation.",
            "bufferpool-wait-time-max",
            "The maximum time in ms an appender waits for space allocation.",
        );
        Self {
            records_sent,
            record_errors,
            record_retries,
            batch_split,
            bytes,
            compression_rate,
            batch_size,
            records_per_request,
            request_latency,
            record_size,
            produce_throttle_time,
            record_queue_time,
            requests_in_flight,
            metadata_age,
            buffer_available_bytes,
            waiting_threads,
            buffer_exhausted,
            bufferpool_wait_time,
        }
    }
}

impl RegistryInner {
    fn topic_sensors(&mut self, topic: &str) -> TopicSensors {
        if let Some(sensors) = self.topics.get(topic) {
            return *sensors;
        }
        let sensors = TopicSensors::register(&mut self.metrics, topic);
        let _previous = self.topics.insert(topic.to_owned(), sensors);
        sensors
    }
}

impl TopicSensors {
    fn register(metrics: &mut Metrics, topic: &str) -> Self {
        let records_sent = topic_meter(
            metrics,
            topic,
            "records-sent",
            "record-send-rate",
            "The average number of records sent per second for a topic.",
            "record-send-total",
            "The total number of records sent for a topic.",
        );
        let bytes = topic_meter(
            metrics,
            topic,
            "bytes",
            "byte-rate",
            "The average number of bytes sent per second for a topic.",
            "byte-total",
            "The total number of bytes sent for a topic.",
        );
        let compression_rate = topic_avg(
            metrics,
            topic,
            "compression-rate",
            "compression-rate",
            "The average compression rate of record batches for a topic.",
        );
        let record_retries = topic_meter(
            metrics,
            topic,
            "record-retries",
            "record-retry-rate",
            "The average per-second number of retried record sends for a topic.",
            "record-retry-total",
            "The total number of retried record sends for a topic.",
        );
        let record_errors = topic_meter(
            metrics,
            topic,
            "record-errors",
            "record-error-rate",
            "The average per-second number of record sends that resulted in errors for a topic.",
            "record-error-total",
            "The total number of record sends that resulted in errors for a topic.",
        );
        Self {
            records_sent,
            bytes,
            compression_rate,
            record_retries,
            record_errors,
        }
    }
}

impl SenderMetricsRegistry {
    fn record(&self, select: impl Fn(&ClientSensors) -> SensorId, value: f64) {
        let mut inner = self
            .inner
            .lock()
            .unwrap_or_else(std::sync::PoisonError::into_inner);
        let sensor = select(&inner.client);
        let _ignored = inner.metrics.record_at_ms(sensor, value, now_ms());
    }

    /// Record a sent batch (records, bytes, compression) for the client and topic.
    pub(crate) fn record_batch(
        &self,
        topic: &str,
        records: u64,
        bytes: u64,
        compression_ratio: f64,
    ) {
        let now = now_ms();
        let mut inner = self
            .inner
            .lock()
            .unwrap_or_else(std::sync::PoisonError::into_inner);
        let records = records as f64;
        let bytes = bytes as f64;
        let client = [
            (inner.client.records_sent, records),
            (inner.client.bytes, bytes),
            (inner.client.batch_size, bytes),
            (inner.client.compression_rate, compression_ratio),
        ];
        for (sensor, value) in client {
            let _ignored = inner.metrics.record_at_ms(sensor, value, now);
        }
        let topic_sensors = inner.topic_sensors(topic);
        let topic_records = [
            (topic_sensors.records_sent, records),
            (topic_sensors.bytes, bytes),
            (topic_sensors.compression_rate, compression_ratio),
        ];
        for (sensor, value) in topic_records {
            let _ignored = inner.metrics.record_at_ms(sensor, value, now);
        }
    }

    /// Record the number of records carried by one produce request.
    pub(crate) fn record_records_per_request(&self, records: u64) {
        self.record(|client| client.records_per_request, records as f64);
    }

    /// Record a produce request round-trip latency in milliseconds.
    pub(crate) fn record_request_latency(&self, latency_ms: f64) {
        self.record(|client| client.request_latency, latency_ms);
    }

    /// Record many serialized record sizes (one produce batch) with a single
    /// lock acquisition and a single clock read, instead of paying both per
    /// record. This keeps the per-record metric semantics (avg/max see every
    /// value) while removing the per-record lock + `clock_gettime` overhead.
    pub(crate) fn record_record_sizes(&self, sizes: &[usize]) {
        if sizes.is_empty() {
            return;
        }
        let now = now_ms();
        let mut inner = self
            .inner
            .lock()
            .unwrap_or_else(std::sync::PoisonError::into_inner);
        let sensor = inner.client.record_size;
        for &size in sizes {
            let _ignored = inner.metrics.record_at_ms(sensor, size as f64, now);
        }
    }

    /// Record a broker-imposed throttle window in milliseconds.
    pub(crate) fn record_throttle_time(&self, throttle_ms: f64) {
        self.record(|client| client.produce_throttle_time, throttle_ms);
    }

    /// Record the time a batch spent buffered before being drained.
    pub(crate) fn record_queue_time(&self, queue_ms: f64) {
        self.record(|client| client.record_queue_time, queue_ms);
    }

    /// Update the current in-flight request gauge.
    pub(crate) fn set_requests_in_flight(&self, in_flight: usize) {
        self.record(|client| client.requests_in_flight, in_flight as f64);
    }

    /// Update the metadata-age gauge in seconds.
    pub(crate) fn set_metadata_age(&self, age_seconds: f64) {
        self.record(|client| client.metadata_age, age_seconds);
    }

    /// Update the available-buffer-memory gauge.
    pub(crate) fn set_buffer_available_bytes(&self, available: usize) {
        self.record(|client| client.buffer_available_bytes, available as f64);
    }

    /// Update the gauge of user threads blocked waiting for buffer memory.
    pub(crate) fn set_waiting_threads(&self, waiting: usize) {
        self.record(|client| client.waiting_threads, waiting as f64);
    }

    /// Record one buffer-memory exhaustion event (an append blocked on buffer
    /// memory) and the time it spent waiting for space allocation.
    pub(crate) fn record_buffer_exhausted(&self, wait_ms: f64) {
        let now = now_ms();
        let mut inner = self
            .inner
            .lock()
            .unwrap_or_else(std::sync::PoisonError::into_inner);
        let exhausted = inner.client.buffer_exhausted;
        let _ignored = inner.metrics.record_at_ms(exhausted, 1.0, now);
        let wait = inner.client.bufferpool_wait_time;
        let _ignored = inner.metrics.record_at_ms(wait, wait_ms, now);
    }

    /// Record a record-send error for the client and (optionally) a topic.
    pub(crate) fn record_error(&self, topic: Option<&str>) {
        let now = now_ms();
        let mut inner = self
            .inner
            .lock()
            .unwrap_or_else(std::sync::PoisonError::into_inner);
        let client = inner.client.record_errors;
        let _ignored = inner.metrics.record_at_ms(client, 1.0, now);
        if let Some(topic) = topic {
            let sensor = inner.topic_sensors(topic).record_errors;
            let _ignored = inner.metrics.record_at_ms(sensor, 1.0, now);
        }
    }

    /// Record a record-send retry for the client and (optionally) a topic.
    pub(crate) fn record_retry(&self, topic: Option<&str>) {
        let now = now_ms();
        let mut inner = self
            .inner
            .lock()
            .unwrap_or_else(std::sync::PoisonError::into_inner);
        let client = inner.client.record_retries;
        let _ignored = inner.metrics.record_at_ms(client, 1.0, now);
        if let Some(topic) = topic {
            let sensor = inner.topic_sensors(topic).record_retries;
            let _ignored = inner.metrics.record_at_ms(sensor, 1.0, now);
        }
    }

    /// Record a batch split.
    pub(crate) fn record_split(&self) {
        self.record(|client| client.batch_split, 1.0);
    }

    /// Snapshot all registered Kafka-named metrics as `"group:name[:tag=value]" -> value`.
    pub(crate) fn kafka_metrics(&self) -> BTreeMap<String, f64> {
        let inner = self
            .inner
            .lock()
            .unwrap_or_else(std::sync::PoisonError::into_inner);
        inner
            .metrics
            .registered_metrics()
            .map(|(name, metric)| (metric_key(name), metric.metric_value()))
            .collect()
    }
}

fn metric_key(name: &MetricName) -> String {
    use std::fmt::Write as _;
    let mut key = format!("{}:{}", name.group(), name.name());
    for (tag_key, tag_value) in name.tags() {
        let _ignored = write!(key, ":{tag_key}={tag_value}");
    }
    key
}

// `MetricName::tags()` returns a `&BTreeMap<String, String>`, already in sorted
// key order, so `metric_key` produces a stable string for each metric.

/// Assert — in debug/test builds only — that a metric registered successfully.
///
/// `sensor_add_*` is a no-op when the identical metric is already on the sensor,
/// and only errors on a genuine name collision with a *different* metric — a
/// setup bug that would otherwise leave a sensor silently reading zero. Release
/// builds keep the intentional-ignore behavior; debug and test builds surface it.
fn expect_metric_registered<E: core::fmt::Debug>(result: Result<(), E>) {
    debug_assert!(
        result.is_ok(),
        "metric registration failed (duplicate metric name?): {result:?}"
    );
    let _ignored = result;
}

#[expect(
    clippy::too_many_arguments,
    reason = "Sensor builder threads sensor + rate/total metric names and descriptions."
)]
fn meter(
    metrics: &mut Metrics,
    sensor_name: &str,
    rate_name: &str,
    rate_desc: &str,
    total_name: &str,
    total_desc: &str,
) -> SensorId {
    let sensor = metrics.sensor(sensor_name);
    let rate = metrics.metric_name(rate_name, CLIENT_GROUP, rate_desc);
    let total = metrics.metric_name(total_name, CLIENT_GROUP, total_desc);
    expect_metric_registered(metrics.sensor_add_meter(sensor, rate, total));
    sensor
}

#[expect(
    clippy::too_many_arguments,
    reason = "Sensor builder threads sensor + avg/max metric names and descriptions."
)]
fn avg_max(
    metrics: &mut Metrics,
    sensor_name: &str,
    avg_name: &str,
    avg_desc: &str,
    max_name: &str,
    max_desc: &str,
) -> SensorId {
    let sensor = metrics.sensor(sensor_name);
    let avg = metrics.metric_name(avg_name, CLIENT_GROUP, avg_desc);
    let max = metrics.metric_name(max_name, CLIENT_GROUP, max_desc);
    expect_metric_registered(metrics.sensor_add_avg(sensor, avg));
    expect_metric_registered(metrics.sensor_add_max(sensor, max));
    sensor
}

fn avg_only(metrics: &mut Metrics, sensor_name: &str, avg_name: &str, avg_desc: &str) -> SensorId {
    let sensor = metrics.sensor(sensor_name);
    let avg = metrics.metric_name(avg_name, CLIENT_GROUP, avg_desc);
    expect_metric_registered(metrics.sensor_add_avg(sensor, avg));
    sensor
}

fn value_only(
    metrics: &mut Metrics,
    sensor_name: &str,
    value_name: &str,
    value_desc: &str,
) -> SensorId {
    let sensor = metrics.sensor(sensor_name);
    let value = metrics.metric_name(value_name, CLIENT_GROUP, value_desc);
    expect_metric_registered(metrics.sensor_add_value(sensor, value));
    sensor
}

#[expect(
    clippy::too_many_arguments,
    reason = "Sensor builder threads topic + sensor + rate/total metric names and descriptions."
)]
fn topic_meter(
    metrics: &mut Metrics,
    topic: &str,
    sensor_suffix: &str,
    rate_name: &str,
    rate_desc: &str,
    total_name: &str,
    total_desc: &str,
) -> SensorId {
    let sensor = metrics.sensor(format!("topic.{topic}.{sensor_suffix}"));
    let rate = metrics.metric_name_with_tags(rate_name, TOPIC_GROUP, rate_desc, [("topic", topic)]);
    let total =
        metrics.metric_name_with_tags(total_name, TOPIC_GROUP, total_desc, [("topic", topic)]);
    expect_metric_registered(metrics.sensor_add_meter(sensor, rate, total));
    sensor
}

fn topic_avg(
    metrics: &mut Metrics,
    topic: &str,
    sensor_suffix: &str,
    avg_name: &str,
    avg_desc: &str,
) -> SensorId {
    let sensor = metrics.sensor(format!("topic.{topic}.{sensor_suffix}"));
    let avg = metrics.metric_name_with_tags(avg_name, TOPIC_GROUP, avg_desc, [("topic", topic)]);
    expect_metric_registered(metrics.sensor_add_avg(sensor, avg));
    sensor
}

#[cfg(test)]
mod tests {
    #![allow(clippy::float_cmp, reason = "Metric totals are exact integer sums.")]

    use super::SenderMetricsRegistry;

    #[test]
    fn exposes_java_named_client_and_topic_metrics() {
        let registry = SenderMetricsRegistry::default();
        registry.record_batch("orders", 3, 300, 0.5);
        registry.record_batch("orders", 2, 200, 0.5);
        registry.record_error(Some("orders"));
        registry.record_retry(Some("orders"));
        registry.record_split();
        registry.record_buffer_exhausted(4.0);
        registry.record_buffer_exhausted(6.0);
        registry.set_buffer_available_bytes(2048);
        registry.set_waiting_threads(2);

        let metrics = registry.kafka_metrics();

        // Kafka BufferPool metrics: exhaustion count, wait-time avg/max, and gauges.
        assert_eq!(
            metrics.get("producer-metrics:buffer-exhausted-total"),
            Some(&2.0)
        );
        assert_eq!(
            metrics.get("producer-metrics:bufferpool-wait-time-avg"),
            Some(&5.0)
        );
        assert_eq!(
            metrics.get("producer-metrics:bufferpool-wait-time-max"),
            Some(&6.0)
        );
        assert_eq!(
            metrics.get("producer-metrics:buffer-available-bytes"),
            Some(&2048.0)
        );
        assert_eq!(metrics.get("producer-metrics:waiting-threads"), Some(&2.0));

        // Client-level cumulative totals (Kafka Meter total = sum of recorded values).
        assert_eq!(
            metrics.get("producer-metrics:record-send-total"),
            Some(&5.0)
        );
        assert_eq!(metrics.get("producer-metrics:byte-total"), Some(&500.0));
        assert_eq!(
            metrics.get("producer-metrics:record-error-total"),
            Some(&1.0)
        );
        assert_eq!(
            metrics.get("producer-metrics:record-retry-total"),
            Some(&1.0)
        );
        assert_eq!(
            metrics.get("producer-metrics:batch-split-total"),
            Some(&1.0)
        );
        assert_eq!(
            metrics.get("producer-metrics:compression-rate-avg"),
            Some(&0.5)
        );

        // Per-topic instances under the producer-topic-metrics group.
        assert_eq!(
            metrics.get("producer-topic-metrics:record-send-total:topic=orders"),
            Some(&5.0)
        );
        assert_eq!(
            metrics.get("producer-topic-metrics:byte-total:topic=orders"),
            Some(&500.0)
        );
        assert_eq!(
            metrics.get("producer-topic-metrics:record-error-total:topic=orders"),
            Some(&1.0)
        );
        assert_eq!(
            metrics.get("producer-topic-metrics:record-retry-total:topic=orders"),
            Some(&1.0)
        );

        // Rate/avg/gauge sensors are registered under their Kafka names.
        for name in [
            "producer-metrics:record-send-rate",
            "producer-metrics:record-error-rate",
            "producer-metrics:records-per-request-avg",
            "producer-metrics:request-latency-avg",
            "producer-metrics:batch-size-avg",
            "producer-metrics:requests-in-flight",
            "producer-metrics:metadata-age",
        ] {
            assert!(metrics.contains_key(name), "missing metric {name}");
        }
    }
}