Ports: each node checks the others' ports from outside
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct Report {
|
||||
/// Unix seconds.
|
||||
pub checked_at: u64,
|
||||
pub probes: Vec<Probe>,
|
||||
/// The hostname didn't resolve, so nothing could be tried.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub error: Option<String>,
|
||||
}
|
||||
|
||||
/// 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<Item = &'x NetworkListener>) -> Vec<u16> {
|
||||
listeners
|
||||
.into_iter()
|
||||
.flat_map(|l| l.bind.iter())
|
||||
.map(|addr| addr.0)
|
||||
.filter(|addr| !addr.ip().is_loopback())
|
||||
.map(|addr| addr.port())
|
||||
.collect::<BTreeSet<_>>()
|
||||
.into_iter()
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn key(target: &str, prober: &str) -> Vec<u8> {
|
||||
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::<BTreeSet<_>>(),
|
||||
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<Vec<NetworkListener>> {
|
||||
Ok(server
|
||||
.registry()
|
||||
.list::<NetworkListener>()
|
||||
.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<Item = &'x NetworkListener>,
|
||||
) -> Vec<SocketAddr> {
|
||||
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::<BTreeSet<_>>()
|
||||
.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<Value> {
|
||||
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::<Vec<_>>();
|
||||
|
||||
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::<String>(KeyValue::<()>::build_key(
|
||||
KV_PORT_REACHABILITY,
|
||||
key(&target.hostname, &prober.hostname),
|
||||
))
|
||||
.await?;
|
||||
let report = stored.and_then(|raw| serde_json::from_str::<Report>(&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<_>>(),
|
||||
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"));
|
||||
}
|
||||
}
|
||||
@@ -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?;
|
||||
|
||||
@@ -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<Inner>) {
|
||||
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<Inner>) {
|
||||
}
|
||||
});
|
||||
}
|
||||
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",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user