diff --git a/crates/common/src/lib.rs b/crates/common/src/lib.rs index 521fa9e..b0b2b26 100644 --- a/crates/common/src/lib.rs +++ b/crates/common/src/lib.rs @@ -70,6 +70,7 @@ pub mod cache; pub mod audit; // inbuxa: the audit log (audit-hold-lock spec, AU) pub mod hold; // inbuxa: legal holds (audit-hold-lock spec, LH) pub mod privacy; // inbuxa: the personal-data catalog, evaluated +pub mod reachability; // inbuxa: whether the outside world reaches each node's ports pub mod config; pub mod expr; pub mod i18n; @@ -129,6 +130,8 @@ pub const KV_LOCK_QUEUE_MESSAGE: u8 = 21; pub const KV_LOCK_TASK: u8 = 23; pub const KV_LOCK_DAV: u8 = 25; pub const KV_SIEVE_ID: u8 = 26; +// inbuxa: far above upstream's prefixes, so a new one of theirs never collides +pub const KV_PORT_REACHABILITY: u8 = 200; #[derive(Clone)] pub struct Server { diff --git a/crates/common/src/reachability.rs b/crates/common/src/reachability.rs new file mode 100644 index 0000000..df72f80 --- /dev/null +++ b/crates/common/src/reachability.rs @@ -0,0 +1,293 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! Whether the outside world can reach each node's ports (settings-reorg, +//! Ports: the reachability check). +//! +//! A server can't answer this about itself: a connection to its own public +//! address never leaves the machine, so it passes whatever the firewall in +//! front says. In a cluster the other nodes are outside that machine. Every +//! ten minutes each node resolves every other active node's hostname, as a +//! sender would, and tries a TCP connection to each listener port on each +//! address. What it saw goes in the shared in-memory store for an hour, under +//! (target, prober), so whichever node the admin asks can report it all. +//! +//! A single server has no one outside to ask. It reports only whether each +//! port is listening, and says so. +//! +//! A connection is all that's tried: nothing is sent, so no protocol logs a +//! session and no rate limit counts it. + +use crate::{KV_PORT_REACHABILITY, Server}; +use registry::schema::{enums::ClusterNodeStatus, structs::NetworkListener}; +use serde::{Deserialize, Serialize}; +use serde_json::{Value, json}; +use std::{ + collections::BTreeSet, + net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr}, + time::{Duration, Instant}, +}; +use store::{dispatch::lookup::KeyValue, write::now}; + +/// How often each node probes the others. +pub const PROBE_INTERVAL: Duration = Duration::from_secs(600); +/// How long one node's view of another is kept: long enough to span a missed round. +const KEEP_FOR: u64 = 3600; +const CONNECT_TIMEOUT: Duration = Duration::from_secs(5); + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct Probe { + pub port: u16, + pub address: String, + pub ok: bool, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct Report { + /// Unix seconds. + pub checked_at: u64, + pub probes: Vec, + /// The hostname didn't resolve, so nothing could be tried. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +/// The ports a sender or client could reach: every listener's port, leaving +/// out listeners bound only to loopback, which are private by design. +pub fn public_ports<'x>(listeners: impl IntoIterator) -> Vec { + listeners + .into_iter() + .flat_map(|l| l.bind.iter()) + .map(|addr| addr.0) + .filter(|addr| !addr.ip().is_loopback()) + .map(|addr| addr.port()) + .collect::>() + .into_iter() + .collect() +} + +fn key(target: &str, prober: &str) -> Vec { + format!("{target}\n{prober}").into_bytes() +} + +async fn connect(address: SocketAddr) -> Result<(), String> { + match tokio::time::timeout(CONNECT_TIMEOUT, tokio::net::TcpStream::connect(address)).await { + Ok(Ok(_)) => Ok(()), + Ok(Err(err)) => Err(err.to_string()), + Err(_) => Err("no answer within 5 seconds".into()), + } +} + +/// Tries each port on each address `hostname` resolves to. +pub async fn probe_host(hostname: &str, ports: &[u16]) -> Report { + let checked_at = now(); + let addresses = match tokio::net::lookup_host((hostname, 0)).await { + Ok(found) => found.map(|a| a.ip()).collect::>(), + Err(err) => { + return Report { + checked_at, + probes: vec![], + error: Some(format!("{hostname} doesn't resolve: {err}")), + }; + } + }; + let tries = addresses.iter().flat_map(|ip| { + ports.iter().map(move |port| { + let address = SocketAddr::new(*ip, *port); + async move { + let result = connect(address).await; + Probe { + port: *port, + address: ip.to_string(), + ok: result.is_ok(), + error: result.err(), + } + } + }) + }); + Report { + checked_at, + probes: futures::future::join_all(tries).await, + error: None, + } +} + +async fn listeners(server: &Server) -> trc::Result> { + Ok(server + .registry() + .list::() + .await? + .into_iter() + .map(|l| l.object) + .collect()) +} + +/// Where to knock to see a port listening on this machine: the bound +/// address, or loopback of the same family for a wildcard bind. +pub fn local_targets<'x>( + listeners: impl IntoIterator, +) -> Vec { + listeners + .into_iter() + .flat_map(|l| l.bind.iter()) + .map(|addr| addr.0) + .filter(|addr| !addr.ip().is_loopback()) + .map(|addr| match addr.ip() { + IpAddr::V4(ip) if ip.is_unspecified() => { + SocketAddr::new(Ipv4Addr::LOCALHOST.into(), addr.port()) + } + IpAddr::V6(ip) if ip.is_unspecified() => { + SocketAddr::new(Ipv6Addr::LOCALHOST.into(), addr.port()) + } + _ => addr, + }) + .collect::>() + .into_iter() + .collect() +} + +/// One round: this node probes every other active node and records what it saw. +pub async fn probe_peers(server: &Server) -> trc::Result<()> { + let nodes = server.registry().cluster_node_list().await?; + let me = server.registry().node_id() as u64; + let Some(prober) = nodes + .iter() + .find(|n| n.node_id == me) + .map(|n| n.hostname.clone()) + else { + return Ok(()); + }; + let ports = public_ports(&listeners(server).await?); + for target in nodes.iter().filter(|n| { + n.node_id != me && n.status == ClusterNodeStatus::Active && n.hostname != prober + }) { + let report = probe_host(&target.hostname, &ports).await; + server + .in_memory_store() + .key_set( + KeyValue::with_prefix( + KV_PORT_REACHABILITY, + key(&target.hostname, &prober), + serde_json::to_vec(&report).unwrap_or_default(), + ) + .expires(KEEP_FOR), + ) + .await?; + } + Ok(()) +} + +/// What `GET /api/ports/check` answers. +pub async fn report(server: &Server) -> trc::Result { + let listeners = listeners(server).await?; + let ports = public_ports(&listeners); + let nodes = if server.core.storage.coordinator.is_enabled() { + server.registry().cluster_node_list().await? + } else { + vec![] + }; + let active = nodes + .iter() + .filter(|n| n.status == ClusterNodeStatus::Active) + .collect::>(); + + if active.len() < 2 { + // No one outside to ask: only whether each port is listening here. + let started = Instant::now(); + let listening = futures::future::join_all(local_targets(&listeners).into_iter().map( + |address| async move { + let result = connect(address).await; + json!({ "port": address.port(), "address": address.ip().to_string(), "listening": result.is_ok() }) + }, + )) + .await; + return Ok(json!({ + "mode": "local", + "ports": ports, + "listening": listening, + "ms": started.elapsed().as_millis() as u64, + })); + } + + let mut out = Vec::new(); + for target in &active { + let mut seen_by = Vec::new(); + for prober in active.iter().filter(|p| p.node_id != target.node_id) { + let stored = server + .in_memory_store() + .key_get::(KeyValue::<()>::build_key( + KV_PORT_REACHABILITY, + key(&target.hostname, &prober.hostname), + )) + .await?; + let report = stored.and_then(|raw| serde_json::from_str::(&raw).ok()); + seen_by.push(json!({ "prober": prober.hostname, "report": report })); + } + out.push(json!({ "hostname": target.hostname, "seenBy": seen_by })); + } + Ok(json!({ + "mode": "cluster", + "ports": ports, + "intervalSeconds": PROBE_INTERVAL.as_secs(), + "nodes": out, + })) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn listener(binds: &[&str]) -> NetworkListener { + NetworkListener { + bind: registry::schema::prelude::Map::new( + binds.iter().map(|b| b.parse().unwrap()).collect(), + ), + ..Default::default() + } + } + + #[test] + fn public_ports_leave_out_loopback_only_listeners() { + let listeners = [ + listener(&["[::]:25"]), + listener(&["0.0.0.0:993", "[::]:993"]), + listener(&["127.0.0.1:8080"]), + listener(&["203.0.113.5:465"]), + ]; + assert_eq!(public_ports(listeners.iter()), vec![25, 465, 993]); + assert_eq!( + local_targets(listeners.iter()) + .iter() + .map(ToString::to_string) + .collect::>(), + vec!["127.0.0.1:993", "203.0.113.5:465", "[::1]:25", "[::1]:993"] + ); + } + + #[tokio::test] + async fn probe_host_reports_open_and_closed_ports() { + let open = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let open_port = open.local_addr().unwrap().port(); + let closed_port = { + let l = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + l.local_addr().unwrap().port() + }; + let report = probe_host("127.0.0.1", &[open_port, closed_port]).await; + assert_eq!(report.error, None); + let ok = |port| report.probes.iter().find(|p| p.port == port).unwrap().ok; + assert!(ok(open_port)); + assert!(!ok(closed_port)); + } + + #[tokio::test] + async fn probe_host_says_when_a_name_does_not_resolve() { + let report = probe_host("does-not-exist.invalid", &[25]).await; + assert!(report.probes.is_empty()); + assert!(report.error.unwrap().contains("doesn't resolve")); + } +} diff --git a/crates/http/src/api/mod.rs b/crates/http/src/api/mod.rs index b35a721..df1f611 100644 --- a/crates/http/src/api/mod.rs +++ b/crates/http/src/api/mod.rs @@ -131,6 +131,18 @@ impl ManagementApi for Server { let answer = jmap::inbuxa::directory_test::test(self, &request).await?; Ok(JsonResponse::new(answer).no_cache().into_http_response()) } + // inbuxa: whether the outside world reaches each node's ports + "ports" if path.get(1).copied() == Some("check") => { + let (_in_flight, access_token) = self.authenticate_headers(req, session).await?; + if access_token.tenant_id().is_some() { + return Err(trc::JmapEvent::Forbidden + .into_err() + .details("Port checks are for server-level administrators.")); + } + access_token.enforce_permission(Permission::SysNetworkListenerGet)?; + let answer = common::reachability::report(self).await?; + Ok(JsonResponse::new(answer).no_cache().into_http_response()) + } "account" => { // Authenticate request let (_in_flight, access_token) = self.authenticate_headers(req, session).await?; diff --git a/crates/services/src/task_manager/scheduler.rs b/crates/services/src/task_manager/scheduler.rs index d283ef6..a32a2d6 100644 --- a/crates/services/src/task_manager/scheduler.rs +++ b/crates/services/src/task_manager/scheduler.rs @@ -46,6 +46,8 @@ enum Event { StoreMetrics, // inbuxa: MON-25: alert evaluation EvaluateAlerts, + // inbuxa: settings-reorg: probe the other nodes' ports + ProbePeerPorts, } /// When the next metric-history tick is due (MON-4), read from the registry @@ -95,6 +97,11 @@ pub fn spawn_task_scheduler(inner: Arc) { Instant::now() + server.registry().refresh_node_id_interval(), Event::RenewNodeIdLease, ); + // inbuxa: first round a minute after start, once the others have a lease + queue.schedule( + Instant::now() + Duration::from_secs(60), + Event::ProbePeerPorts, + ); } // Spam classifier training @@ -229,6 +236,19 @@ pub fn spawn_task_scheduler(inner: Arc) { } }); } + Event::ProbePeerPorts => { + queue.schedule( + Instant::now() + common::reachability::PROBE_INTERVAL, + Event::ProbePeerPorts, + ); + + let server = server.clone(); + tokio::spawn(async move { + if let Err(err) = common::reachability::probe_peers(&server).await { + trc::error!(err.details("Failed to probe the other nodes' ports")); + } + }); + } Event::OtelMetrics => { if let Some(otel) = &server.core.metrics.otel { queue.schedule(Instant::now() + otel.interval, Event::OtelMetrics); @@ -476,6 +496,7 @@ impl Event { Event::RenewNodeIdLease => "renewNodeIdLease", Event::StoreMetrics => "storeMetrics", Event::EvaluateAlerts => "evaluateAlerts", + Event::ProbePeerPorts => "probePeerPorts", } } }