datum-mq 0.10.9

Kafka sources and sinks for Datum streams, with native and rdkafka backends
Documentation
use std::{
    fs,
    future::Future,
    process::Command,
    sync::{
        OnceLock,
        atomic::{AtomicBool, AtomicU64, Ordering},
    },
    time::Instant,
};

#[derive(Debug, Clone, Copy)]
pub enum ProfileBucket {
    SocketFetch,
    ResponseParse,
    RecordDecode,
    PayloadCopy,
    OffsetBookkeeping,
    ChannelHandoff,
    SourceEmit,
    DatumConsume,
    CommitBookkeeping,
}

#[derive(Debug, Clone, Copy, Default)]
pub struct ProfileBucketSnapshot {
    pub calls: u64,
    pub wall_ns: u64,
    pub cpu_ns: u64,
}

#[derive(Debug, Clone, Copy, Default)]
pub struct ProfileSnapshot {
    pub socket_fetch: ProfileBucketSnapshot,
    pub response_parse: ProfileBucketSnapshot,
    pub record_decode: ProfileBucketSnapshot,
    pub payload_copy: ProfileBucketSnapshot,
    pub offset_bookkeeping: ProfileBucketSnapshot,
    pub channel_handoff: ProfileBucketSnapshot,
    pub source_emit: ProfileBucketSnapshot,
    pub datum_consume: ProfileBucketSnapshot,
    pub commit_bookkeeping: ProfileBucketSnapshot,
}

#[derive(Debug, Default)]
struct BucketCounters {
    calls: AtomicU64,
    wall_ns: AtomicU64,
    cpu_ns: AtomicU64,
}

#[derive(Debug, Default)]
struct ProfileCounters {
    socket_fetch: BucketCounters,
    response_parse: BucketCounters,
    record_decode: BucketCounters,
    payload_copy: BucketCounters,
    offset_bookkeeping: BucketCounters,
    channel_handoff: BucketCounters,
    source_emit: BucketCounters,
    datum_consume: BucketCounters,
    commit_bookkeeping: BucketCounters,
}

static ENABLED: AtomicBool = AtomicBool::new(false);
static COUNTERS: ProfileCounters = ProfileCounters {
    socket_fetch: BucketCounters {
        calls: AtomicU64::new(0),
        wall_ns: AtomicU64::new(0),
        cpu_ns: AtomicU64::new(0),
    },
    response_parse: BucketCounters {
        calls: AtomicU64::new(0),
        wall_ns: AtomicU64::new(0),
        cpu_ns: AtomicU64::new(0),
    },
    record_decode: BucketCounters {
        calls: AtomicU64::new(0),
        wall_ns: AtomicU64::new(0),
        cpu_ns: AtomicU64::new(0),
    },
    payload_copy: BucketCounters {
        calls: AtomicU64::new(0),
        wall_ns: AtomicU64::new(0),
        cpu_ns: AtomicU64::new(0),
    },
    offset_bookkeeping: BucketCounters {
        calls: AtomicU64::new(0),
        wall_ns: AtomicU64::new(0),
        cpu_ns: AtomicU64::new(0),
    },
    channel_handoff: BucketCounters {
        calls: AtomicU64::new(0),
        wall_ns: AtomicU64::new(0),
        cpu_ns: AtomicU64::new(0),
    },
    source_emit: BucketCounters {
        calls: AtomicU64::new(0),
        wall_ns: AtomicU64::new(0),
        cpu_ns: AtomicU64::new(0),
    },
    datum_consume: BucketCounters {
        calls: AtomicU64::new(0),
        wall_ns: AtomicU64::new(0),
        cpu_ns: AtomicU64::new(0),
    },
    commit_bookkeeping: BucketCounters {
        calls: AtomicU64::new(0),
        wall_ns: AtomicU64::new(0),
        cpu_ns: AtomicU64::new(0),
    },
};

pub fn set_enabled(enabled: bool) {
    ENABLED.store(enabled, Ordering::Relaxed);
}

pub fn reset() {
    COUNTERS.reset();
}

#[must_use]
pub fn enabled() -> bool {
    ENABLED.load(Ordering::Relaxed)
}

pub fn measure<T, F>(bucket: ProfileBucket, f: F) -> T
where
    F: FnOnce() -> T,
{
    if !enabled() {
        return f();
    }
    let wall_start = Instant::now();
    let cpu_start = thread_cpu_ns();
    let result = f();
    record(
        bucket,
        wall_start.elapsed().as_nanos() as u64,
        thread_cpu_ns().saturating_sub(cpu_start),
    );
    result
}

pub async fn measure_async<T, F>(bucket: ProfileBucket, future: F) -> T
where
    F: Future<Output = T>,
{
    if !enabled() {
        return future.await;
    }
    let wall_start = Instant::now();
    let cpu_start = thread_cpu_ns();
    let result = future.await;
    record(
        bucket,
        wall_start.elapsed().as_nanos() as u64,
        thread_cpu_ns().saturating_sub(cpu_start),
    );
    result
}

pub fn snapshot() -> ProfileSnapshot {
    ProfileSnapshot {
        socket_fetch: snapshot_bucket(&COUNTERS.socket_fetch),
        response_parse: snapshot_bucket(&COUNTERS.response_parse),
        record_decode: snapshot_bucket(&COUNTERS.record_decode),
        payload_copy: snapshot_bucket(&COUNTERS.payload_copy),
        offset_bookkeeping: snapshot_bucket(&COUNTERS.offset_bookkeeping),
        channel_handoff: snapshot_bucket(&COUNTERS.channel_handoff),
        source_emit: snapshot_bucket(&COUNTERS.source_emit),
        datum_consume: snapshot_bucket(&COUNTERS.datum_consume),
        commit_bookkeeping: snapshot_bucket(&COUNTERS.commit_bookkeeping),
    }
}

fn record(bucket: ProfileBucket, wall_ns: u64, cpu_ns: u64) {
    let counters = match bucket {
        ProfileBucket::SocketFetch => &COUNTERS.socket_fetch,
        ProfileBucket::ResponseParse => &COUNTERS.response_parse,
        ProfileBucket::RecordDecode => &COUNTERS.record_decode,
        ProfileBucket::PayloadCopy => &COUNTERS.payload_copy,
        ProfileBucket::OffsetBookkeeping => &COUNTERS.offset_bookkeeping,
        ProfileBucket::ChannelHandoff => &COUNTERS.channel_handoff,
        ProfileBucket::SourceEmit => &COUNTERS.source_emit,
        ProfileBucket::DatumConsume => &COUNTERS.datum_consume,
        ProfileBucket::CommitBookkeeping => &COUNTERS.commit_bookkeeping,
    };
    counters.calls.fetch_add(1, Ordering::Relaxed);
    counters.wall_ns.fetch_add(wall_ns, Ordering::Relaxed);
    counters.cpu_ns.fetch_add(cpu_ns, Ordering::Relaxed);
}

fn snapshot_bucket(counters: &BucketCounters) -> ProfileBucketSnapshot {
    ProfileBucketSnapshot {
        calls: counters.calls.load(Ordering::Relaxed),
        wall_ns: counters.wall_ns.load(Ordering::Relaxed),
        cpu_ns: counters.cpu_ns.load(Ordering::Relaxed),
    }
}

impl ProfileCounters {
    fn reset(&self) {
        self.socket_fetch.reset();
        self.response_parse.reset();
        self.record_decode.reset();
        self.payload_copy.reset();
        self.offset_bookkeeping.reset();
        self.channel_handoff.reset();
        self.source_emit.reset();
        self.datum_consume.reset();
        self.commit_bookkeeping.reset();
    }
}

impl BucketCounters {
    fn reset(&self) {
        self.calls.store(0, Ordering::Relaxed);
        self.wall_ns.store(0, Ordering::Relaxed);
        self.cpu_ns.store(0, Ordering::Relaxed);
    }
}

fn thread_cpu_ns() -> u64 {
    let Ok(stat) = fs::read_to_string("/proc/thread-self/stat") else {
        return 0;
    };
    let Some(close) = stat.rfind(')') else {
        return 0;
    };
    let fields = stat[close + 1..].split_whitespace().collect::<Vec<_>>();
    if fields.len() <= 12 {
        return 0;
    }
    let utime = fields[11].parse::<u64>().unwrap_or(0);
    let stime = fields[12].parse::<u64>().unwrap_or(0);
    ((utime + stime) as f64 * 1_000_000_000.0 / clock_ticks_per_second()) as u64
}

fn clock_ticks_per_second() -> f64 {
    static TICKS: OnceLock<f64> = OnceLock::new();
    *TICKS.get_or_init(|| {
        Command::new("getconf")
            .arg("CLK_TCK")
            .output()
            .ok()
            .and_then(|output| String::from_utf8(output.stdout).ok())
            .and_then(|value| value.trim().parse::<f64>().ok())
            .unwrap_or(100.0)
    })
}