From d4dadf33e6c57c849887247faa0f8b5c37835f11 Mon Sep 17 00:00:00 2001 From: reverse Date: Tue, 18 Aug 2026 22:34:36 +0000 Subject: [PATCH] =?UTF-8?q?metrics:=20optional=20OpenMetrics/Prometheus=20?= =?UTF-8?q?endpoint=20(metrics=5Fbind,=20off=20by=20default)=20=E2=80=94?= =?UTF-8?q?=20commands/messages/connects=20counters=20bumped=20inline=20vi?= =?UTF-8?q?a=20shared=20atomics,=20users/channels/servers/links=20gauges?= =?UTF-8?q?=20republished=20each=20tick;=20no=20event=20round-trip=20on=20?= =?UTF-8?q?the=20hot=20path?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/ircd.rs | 15 +++++++ src/main.rs | 3 ++ src/modules/metrics.rs | 89 ++++++++++++++++++++++++++++++++++++++++++ src/modules/mod.rs | 1 + src/server.rs | 3 ++ src/users.rs | 3 ++ 6 files changed, 114 insertions(+) create mode 100644 src/modules/metrics.rs diff --git a/src/ircd.rs b/src/ircd.rs index c763974..efe9d5a 100644 --- a/src/ircd.rs +++ b/src/ircd.rs @@ -475,6 +475,11 @@ impl Ircd { ); return; } + use std::sync::atomic::Ordering::Relaxed; + self.server.metrics.commands.fetch_add(1, Relaxed); + if matches!(cmd, "PRIVMSG" | "NOTICE") { + self.server.metrics.messages.fetch_add(1, Relaxed); + } let _ = handler.handle(&mut self.server, uid, &msg.params); for m in &mut self.modules { @@ -669,6 +674,16 @@ impl Ircd { .send(uid, format!("ERROR :Closing link: ({reason})")); self.quit_user(uid, reason); } + // republish gauges (the core owns this state; the scrape thread only reads) + use std::sync::atomic::Ordering::Relaxed; + let m = &self.server.metrics; + m.users.store( + self.server.users.values().filter(|u| u.registered).count() as u64, + Relaxed, + ); + m.channels.store(self.server.channels.len() as u64, Relaxed); + m.servers.store(self.server.servers.len() as u64, Relaxed); + m.links.store(self.server.links.len() as u64, Relaxed); } /// Fire queued notify-hooks. Draining a queue (not iterating in place) lets a diff --git a/src/main.rs b/src/main.rs index 88cc497..33124dd 100644 --- a/src/main.rs +++ b/src/main.rs @@ -249,6 +249,9 @@ fn main() { // optional JSON-RPC-over-HTTP control interface (see crate::modules::rpc) echoircd::modules::rpc::maybe_start(&cfg, tx.clone()); + // optional OpenMetrics/Prometheus scrape endpoint (metrics_bind = host:port) + echoircd::modules::metrics::maybe_start(&cfg); + // optional WebSocket transport for browser IRC clients (see crate::websocket) echoircd::websocket::maybe_start(&cfg, tx.clone(), counter.clone()); diff --git a/src/modules/metrics.rs b/src/modules/metrics.rs new file mode 100644 index 0000000..f953a0e --- /dev/null +++ b/src/modules/metrics.rs @@ -0,0 +1,89 @@ +//! metrics — an optional Prometheus/OpenMetrics endpoint. Enable with +//! `metrics_bind = 127.0.0.1:9100` in the config (off by default). +//! +//! Counters live in a process-wide `Arc` of atomics: the core bumps them +//! inline (a relaxed atomic add, no lock, no event round-trip), and a tiny HTTP +//! thread reads them on scrape. Gauges (current users/channels/servers) are +//! republished each tick by the core, which owns that state. + +use std::io::{Read, Write}; +use std::net::TcpListener; +use std::sync::atomic::{AtomicU64, Ordering::Relaxed}; +use std::sync::{Arc, OnceLock}; + +use crate::config::Config; + +/// Every metric echoircd exposes. Counters only ever increase; gauges are set to +/// the live count each tick. +#[derive(Default)] +pub struct Metrics { + // counters (monotonic) + pub commands: AtomicU64, + pub messages: AtomicU64, + pub connects: AtomicU64, + // gauges (republished each tick) + pub users: AtomicU64, + pub channels: AtomicU64, + pub servers: AtomicU64, + pub links: AtomicU64, +} + +static METRICS: OnceLock> = OnceLock::new(); + +/// The shared metrics handle (created on first use). The core and the HTTP scrape +/// thread both call this, so they see the same atomics. +pub fn handle() -> Arc { + METRICS.get_or_init(|| Arc::new(Metrics::default())).clone() +} + +/// Render the current values in OpenMetrics/Prometheus text exposition format. +fn render(m: &Metrics) -> String { + let mut o = String::new(); + let counter = |o: &mut String, name: &str, help: &str, v: u64| { + o.push_str(&format!("# HELP {name} {help}\n# TYPE {name} counter\n{name} {v}\n")); + }; + let gauge = |o: &mut String, name: &str, help: &str, v: u64| { + o.push_str(&format!("# HELP {name} {help}\n# TYPE {name} gauge\n{name} {v}\n")); + }; + counter(&mut o, "echoircd_commands_total", "Commands dispatched.", m.commands.load(Relaxed)); + counter(&mut o, "echoircd_messages_total", "PRIVMSG/NOTICE handled.", m.messages.load(Relaxed)); + counter(&mut o, "echoircd_connects_total", "Client registrations completed.", m.connects.load(Relaxed)); + gauge(&mut o, "echoircd_users", "Registered users online.", m.users.load(Relaxed)); + gauge(&mut o, "echoircd_channels", "Channels in existence.", m.channels.load(Relaxed)); + gauge(&mut o, "echoircd_servers", "Servers known on the network.", m.servers.load(Relaxed)); + gauge(&mut o, "echoircd_links", "Direct server links.", m.links.load(Relaxed)); + o +} + +/// Start the scrape endpoint if `metrics_bind` is configured. Serves any GET with +/// the exposition text; it carries no secrets, so bind it somewhere private. +pub fn maybe_start(cfg: &Config) { + let Some(bind) = cfg.raw.get("metrics_bind").and_then(|v| v.last()).filter(|s| !s.is_empty()) else { + return; + }; + let bind = bind.to_string(); + let metrics = handle(); + match TcpListener::bind(&bind) { + Ok(listener) => { + eprintln!("echoircd metrics (OpenMetrics) on {bind}"); + std::thread::spawn(move || serve(listener, metrics)); + } + Err(e) => eprintln!("echoircd: cannot bind metrics {bind}: {e}"), + } +} + +fn serve(listener: TcpListener, metrics: Arc) { + for stream in listener.incoming() { + let Ok(mut s) = stream else { continue }; + // read (and ignore) the request head, then reply — this is a scrape, no routing + let mut buf = [0u8; 1024]; + let _ = s.read(&mut buf); + let body = render(&metrics); + let resp = format!( + "HTTP/1.1 200 OK\r\nContent-Type: text/plain; version=0.0.4\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ); + let _ = s.write_all(resp.as_bytes()); + } +} diff --git a/src/modules/mod.rs b/src/modules/mod.rs index c12c218..beae6a5 100644 --- a/src/modules/mod.rs +++ b/src/modules/mod.rs @@ -47,6 +47,7 @@ pub mod log_json; pub mod maphide; pub mod markread; pub mod metadata; +pub mod metrics; pub mod multiline; pub mod namedmodes; pub mod network_icon; diff --git a/src/server.rs b/src/server.rs index d52fb4a..1490748 100644 --- a/src/server.rs +++ b/src/server.rs @@ -174,6 +174,8 @@ pub struct Server { /// Module-owned server state, keyed by type. Each `modules/*.rs` stores its /// own struct here so features live in their own file instead of this one. pub ext: Extensible, + /// Prometheus counters/gauges, shared with the scrape thread (modules::metrics). + pub metrics: Arc, } impl Server { @@ -220,6 +222,7 @@ impl Server { event_tx, conn_counter, ext: Extensible::default(), + metrics: crate::modules::metrics::handle(), } } diff --git a/src/users.rs b/src/users.rs index abf968c..09ffb82 100644 --- a/src/users.rs +++ b/src/users.rs @@ -439,6 +439,9 @@ impl Server { if let Some(u) = self.users.get_mut(&uid) { u.registered = true; } + self.metrics + .connects + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); let nick = self .users .get(&uid)