diff --git a/constellation/src/bin/main.rs b/constellation/src/bin/main.rs index dcac2b2..9fe78f4 100644 --- a/constellation/src/bin/main.rs +++ b/constellation/src/bin/main.rs @@ -1,5 +1,6 @@ use anyhow::{bail, Result}; use clap::{Parser, ValueEnum}; +use metrics::{describe_counter, describe_gauge, describe_histogram, Unit}; use metrics_exporter_prometheus::PrometheusBuilder; use std::net::SocketAddr; use std::num::NonZero; @@ -244,16 +245,6 @@ fn run( let process_collector = metrics_process::Collector::default(); process_collector.describe(); - metrics::describe_gauge!( - "storage_available", - metrics::Unit::Bytes, - "available to be allocated" - ); - metrics::describe_gauge!( - "storage_free", - metrics::Unit::Bytes, - "unused bytes in filesystem" - ); if let Some(ref p) = data_dir { if let Err(e) = fs4::available_space(p) { eprintln!("fs4 failed to get available space. may not be supported here? space metrics may be absent. e: {e:?}"); @@ -301,7 +292,10 @@ fn run( fn install_metrics_server(metrics_bind: SocketAddr) -> Result<()> { println!("installing metrics server..."); - #[expect(deprecated, reason = "would change counters to _total suffix, needs dash updates")] + #[expect( + deprecated, + reason = "would change counters to _total suffix, needs dash updates" + )] PrometheusBuilder::new() .idle_timeout( metrics_util::MetricKindMask::ALL, @@ -313,10 +307,108 @@ fn install_metrics_server(metrics_bind: SocketAddr) -> Result<()> { .set_enable_unit_suffix(true) .with_http_listener(metrics_bind) .install()?; + describe_metrics(); println!("metrics server installed! listening at {metrics_bind:?}"); Ok(()) } +fn describe_metrics() { + describe_gauge!( + "storage_available", + Unit::Bytes, + "available to be allocated" + ); + describe_gauge!("storage_free", Unit::Bytes, "unused bytes in filesystem"); + describe_counter!( + "jetstream_connnect", + Unit::Count, + "attempts to connect to a jetstream server" + ); + describe_counter!( + "jetstream_read", + Unit::Count, + "attempts to read an event from jetstream" + ); + describe_counter!( + "jetstream_read_fail", + Unit::Count, + "failures to read events from jetstream" + ); + describe_counter!( + "jetstream_read_bytes", + Unit::Bytes, + "total received message bytes from jetstream" + ); + describe_counter!( + "jetstream_read_bytes_decompressed", + Unit::Bytes, + "total decompressed message bytes from jetstream" + ); + describe_histogram!( + "jetstream_read_bytes_decompressed", + Unit::Bytes, + "decompressed size of jetstream messages" + ); + describe_counter!( + "jetstream_events", + Unit::Count, + "valid json messages received" + ); + describe_histogram!( + "jetstream_events_queued", + Unit::Count, + "event messages waiting in queue" + ); + describe_gauge!( + "jetstream_cursor_age", + Unit::Microseconds, + "microseconds between our clock and the jetstream event's time_us" + ); + describe_counter!( + "consumer_events_non_actionable", + Unit::Count, + "count of non-actionable events" + ); + describe_counter!( + "consumer_events_actionable", + Unit::Count, + "count of action by type. *all* atproto record delete events are included" + ); + describe_counter!( + "consumer_events_actionable_links", + Unit::Count, + "total links encountered" + ); + describe_histogram!( + "consumer_events_actionable_links", + Unit::Count, + "number of links per message" + ); + #[cfg(feature = "rocks")] + { + describe_histogram!( + "storage_rocksdb_read_seconds", + Unit::Seconds, + "duration of the read stage of actions" + ); + describe_histogram!( + "storage_rocksdb_action_seconds", + Unit::Seconds, + "duration of read + write of actions" + ); + describe_counter!( + "storage_rocksdb_batch_ops_total", + Unit::Count, + "total batched operations from actions" + ); + describe_histogram!( + "storage_rocksdb_delete_account_ops", + Unit::Count, + "total batched ops for account deletions" + ); + } +} + #[cfg(test)] mod tests { use constellation::consumer::get_actionable; diff --git a/constellation/src/consumer/jetstream.rs b/constellation/src/consumer/jetstream.rs index 3808b48..7e58155 100644 --- a/constellation/src/consumer/jetstream.rs +++ b/constellation/src/consumer/jetstream.rs @@ -1,7 +1,5 @@ use anyhow::{bail, Result}; -use metrics::{ - counter, describe_counter, describe_gauge, describe_histogram, gauge, histogram, Unit, -}; +use metrics::{counter, gauge, histogram}; use std::io::{Cursor, ErrorKind, Read}; use std::net::ToSocketAddrs; use std::thread; @@ -19,52 +17,6 @@ pub fn consume_jetstream( stream: String, staying_alive: CancellationToken, ) -> Result<()> { - describe_counter!( - "jetstream_connnect", - Unit::Count, - "attempts to connect to a jetstream server" - ); - describe_counter!( - "jetstream_read", - Unit::Count, - "attempts to read an event from jetstream" - ); - describe_counter!( - "jetstream_read_fail", - Unit::Count, - "failures to read events from jetstream" - ); - describe_counter!( - "jetstream_read_bytes", - Unit::Bytes, - "total received message bytes from jetstream" - ); - describe_counter!( - "jetstream_read_bytes_decompressed", - Unit::Bytes, - "total decompressed message bytes from jetstream" - ); - describe_histogram!( - "jetstream_read_bytes_decompressed", - Unit::Bytes, - "decompressed size of jetstream messages" - ); - describe_counter!( - "jetstream_events", - Unit::Count, - "valid json messages received" - ); - describe_histogram!( - "jetstream_events_queued", - Unit::Count, - "event messages waiting in queue" - ); - describe_gauge!( - "jetstream_cursor_age", - Unit::Microseconds, - "microseconds between our clock and the jetstream event's time_us" - ); - let dict = DecoderDictionary::copy(JETSTREAM_ZSTD_DICTIONARY); let mut connect_retries = 0; let mut latest_cursor = cursor; diff --git a/constellation/src/consumer/mod.rs b/constellation/src/consumer/mod.rs index 9ba9c14..a5dabd9 100644 --- a/constellation/src/consumer/mod.rs +++ b/constellation/src/consumer/mod.rs @@ -7,7 +7,7 @@ use anyhow::Result; use jetstream::consume_jetstream; use jsonl_file::consume_jsonl_file; use links::{parse_any_link, record::walk_record, CollectedLink}; -use metrics::{counter, describe_counter, describe_histogram, histogram, Unit}; +use metrics::{counter, histogram}; use std::path::PathBuf; use std::sync::atomic::{AtomicU32, Ordering}; use std::sync::Arc; @@ -23,27 +23,6 @@ pub fn consume( stream: String, staying_alive: CancellationToken, ) -> Result<()> { - describe_counter!( - "consumer_events_non_actionable", - Unit::Count, - "count of non-actionable events" - ); - describe_counter!( - "consumer_events_actionable", - Unit::Count, - "count of action by type. *all* atproto record delete events are included" - ); - describe_counter!( - "consumer_events_actionable_links", - Unit::Count, - "total links encountered" - ); - describe_histogram!( - "consumer_events_actionable_links", - Unit::Count, - "number of links per message" - ); - let mut fixture_cursor = None; let (receiver, consumer_handle) = if let Some(f) = fixture { let (sender, receiver) = flume::bounded(21); diff --git a/constellation/src/storage/rocks_store.rs b/constellation/src/storage/rocks_store.rs index 3974636..993931a 100644 --- a/constellation/src/storage/rocks_store.rs +++ b/constellation/src/storage/rocks_store.rs @@ -7,7 +7,7 @@ use crate::{CountsByCount, Did, ManyToManyItem, RecordId}; use anyhow::{anyhow, bail, Result}; use bincode::Options as BincodeOptions; use links::CollectedLink; -use metrics::{counter, describe_counter, describe_histogram, histogram, Unit}; +use metrics::{counter, histogram}; use ratelimit::Ratelimiter; use rocksdb::backup::{BackupEngine, BackupEngineOptions}; use rocksdb::{ @@ -256,7 +256,6 @@ fn now() -> u64 { impl RocksStorage { pub fn new(path: impl AsRef) -> Result { - Self::describe_metrics(); let me = RocksStorage::open_readmode(path, false)?; me.global_init()?; Ok(me) @@ -420,29 +419,6 @@ impl RocksStorage { Ok(()) } - fn describe_metrics() { - describe_histogram!( - "storage_rocksdb_read_seconds", - Unit::Seconds, - "duration of the read stage of actions" - ); - describe_histogram!( - "storage_rocksdb_action_seconds", - Unit::Seconds, - "duration of read + write of actions" - ); - describe_counter!( - "storage_rocksdb_batch_ops_total", - Unit::Count, - "total batched operations from actions" - ); - describe_histogram!( - "storage_rocksdb_delete_account_ops", - Unit::Count, - "total batched ops for account deletions" - ); - } - fn merge_op_extend_did_ids( key: &[u8], existing: Option<&[u8]>,