datum-mq 0.10.9

Kafka sources and sinks for Datum streams, with native and rdkafka backends
Documentation
#![cfg_attr(not(feature = "rdkafka"), allow(dead_code))]

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 produced_batch(&self, count: u64) {
        self.inner.produced.fetch_add(count, Ordering::Relaxed);
    }

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

    pub(crate) fn delivery_failed_batch(&self, count: u64) {
        self.inner
            .delivery_failures
            .fetch_add(count, Ordering::Relaxed);
    }

    pub(crate) fn queue_full(&self) {
        self.inner.queue_full.fetch_add(1, Ordering::Relaxed);
    }
    pub(crate) fn replace_from(&self, snapshot: KafkaMetricsSnapshot) {
        self.inner
            .assigned_partitions
            .store(snapshot.assigned_partitions, Ordering::Relaxed);
        self.inner
            .revoked_partitions
            .store(snapshot.revoked_partitions, Ordering::Relaxed);
        self.inner
            .lost_partitions
            .store(snapshot.lost_partitions, Ordering::Relaxed);
        self.inner
            .rebalances
            .store(snapshot.rebalances, Ordering::Relaxed);
        self.inner
            .emitted
            .store(snapshot.emitted, Ordering::Relaxed);
        self.inner
            .committed_offsets
            .store(snapshot.committed_offsets, Ordering::Relaxed);
        self.inner
            .commit_failures
            .store(snapshot.commit_failures, Ordering::Relaxed);
        self.inner.paused.store(snapshot.paused, Ordering::Relaxed);
        self.inner
            .outstanding
            .store(snapshot.outstanding, Ordering::Relaxed);
        self.inner
            .producer_in_flight
            .store(snapshot.producer_in_flight, Ordering::Relaxed);
        self.inner
            .produced
            .store(snapshot.produced, Ordering::Relaxed);
        self.inner
            .delivery_failures
            .store(snapshot.delivery_failures, Ordering::Relaxed);
        self.inner
            .queue_full
            .store(snapshot.queue_full, Ordering::Relaxed);
        self.inner
            .high_watermark
            .store(snapshot.high_watermark, Ordering::Relaxed);
        self.inner
            .committed_watermark
            .store(snapshot.committed_watermark, Ordering::Relaxed);
    }
}