datum-mq 0.10.1

Kafka sources and sinks for Datum streams, backed by rdkafka
Documentation
use std::sync::{
    Arc,
    atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering},
};

/// Cheap cloneable metrics handle for Kafka sources/sinks.
#[derive(Debug, Clone, Default)]
pub struct KafkaMetrics {
    inner: Arc<KafkaMetricsInner>,
}

#[derive(Debug, Default)]
struct KafkaMetricsInner {
    assigned_partitions: AtomicU64,
    revoked_partitions: AtomicU64,
    lost_partitions: AtomicU64,
    rebalances: AtomicU64,
    emitted: AtomicU64,
    committed_offsets: AtomicU64,
    commit_failures: AtomicU64,
    paused: AtomicBool,
    outstanding: AtomicU64,
    producer_in_flight: AtomicU64,
    produced: AtomicU64,
    delivery_failures: AtomicU64,
    queue_full: AtomicU64,
    high_watermark: AtomicI64,
    committed_watermark: AtomicI64,
}

/// Point-in-time metrics snapshot.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct KafkaMetricsSnapshot {
    pub assigned_partitions: u64,
    pub revoked_partitions: u64,
    pub lost_partitions: u64,
    pub rebalances: u64,
    pub emitted: u64,
    pub committed_offsets: u64,
    pub commit_failures: u64,
    pub paused: bool,
    pub outstanding: u64,
    pub producer_in_flight: u64,
    pub produced: u64,
    pub delivery_failures: u64,
    pub queue_full: u64,
    pub high_watermark: i64,
    pub committed_watermark: i64,
}

impl KafkaMetrics {
    #[must_use]
    pub fn snapshot(&self) -> KafkaMetricsSnapshot {
        KafkaMetricsSnapshot {
            assigned_partitions: self.inner.assigned_partitions.load(Ordering::Relaxed),
            revoked_partitions: self.inner.revoked_partitions.load(Ordering::Relaxed),
            lost_partitions: self.inner.lost_partitions.load(Ordering::Relaxed),
            rebalances: self.inner.rebalances.load(Ordering::Relaxed),
            emitted: self.inner.emitted.load(Ordering::Relaxed),
            committed_offsets: self.inner.committed_offsets.load(Ordering::Relaxed),
            commit_failures: self.inner.commit_failures.load(Ordering::Relaxed),
            paused: self.inner.paused.load(Ordering::Relaxed),
            outstanding: self.inner.outstanding.load(Ordering::Relaxed),
            producer_in_flight: self.inner.producer_in_flight.load(Ordering::Relaxed),
            produced: self.inner.produced.load(Ordering::Relaxed),
            delivery_failures: self.inner.delivery_failures.load(Ordering::Relaxed),
            queue_full: self.inner.queue_full.load(Ordering::Relaxed),
            high_watermark: self.inner.high_watermark.load(Ordering::Relaxed),
            committed_watermark: self.inner.committed_watermark.load(Ordering::Relaxed),
        }
    }

    pub(crate) fn add_assigned(&self, count: u64) {
        self.inner
            .assigned_partitions
            .fetch_add(count, Ordering::Relaxed);
        self.inner.rebalances.fetch_add(1, Ordering::Relaxed);
    }

    pub(crate) fn add_revoked(&self, count: u64) {
        self.inner
            .revoked_partitions
            .fetch_add(count, Ordering::Relaxed);
        self.inner.rebalances.fetch_add(1, Ordering::Relaxed);
    }

    pub(crate) fn add_lost(&self, count: u64) {
        self.inner
            .lost_partitions
            .fetch_add(count, Ordering::Relaxed);
        self.inner.rebalances.fetch_add(1, Ordering::Relaxed);
    }

    pub(crate) fn emitted(&self, outstanding: u64, high_watermark: i64) {
        self.inner.emitted.fetch_add(1, Ordering::Relaxed);
        self.inner.outstanding.store(outstanding, Ordering::Relaxed);
        self.inner
            .high_watermark
            .store(high_watermark, Ordering::Relaxed);
    }

    pub(crate) fn emitted_batch(&self, count: u64, outstanding: u64, high_watermark: i64) {
        self.inner.emitted.fetch_add(count, Ordering::Relaxed);
        self.inner.outstanding.store(outstanding, Ordering::Relaxed);
        self.inner
            .high_watermark
            .store(high_watermark, Ordering::Relaxed);
    }

    pub(crate) fn committed(&self, outstanding: u64, committed_watermark: i64) {
        self.inner.committed_offsets.fetch_add(1, Ordering::Relaxed);
        self.inner.outstanding.store(outstanding, Ordering::Relaxed);
        self.inner
            .committed_watermark
            .store(committed_watermark, Ordering::Relaxed);
    }

    pub(crate) fn set_outstanding(&self, outstanding: u64) {
        self.inner.outstanding.store(outstanding, Ordering::Relaxed);
    }

    pub(crate) fn commit_failed(&self) {
        self.inner.commit_failures.fetch_add(1, Ordering::Relaxed);
    }

    pub(crate) fn set_paused(&self, paused: bool) {
        self.inner.paused.store(paused, Ordering::Relaxed);
    }

    pub(crate) fn producer_in_flight(&self, in_flight: u64) {
        self.inner
            .producer_in_flight
            .store(in_flight, Ordering::Relaxed);
    }

    pub(crate) fn produced(&self) {
        self.inner.produced.fetch_add(1, Ordering::Relaxed);
    }

    pub(crate) fn delivery_failed(&self) {
        self.inner.delivery_failures.fetch_add(1, Ordering::Relaxed);
    }

    pub(crate) fn queue_full(&self) {
        self.inner.queue_full.fetch_add(1, Ordering::Relaxed);
    }
}