Release 2026.9.28.6 #107
@@ -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"));
|
||||
}
|
||||
}
|
||||
@@ -156,15 +156,7 @@ async fn post_webhook_events(
|
||||
|
||||
// Add HMAC-SHA256 signature
|
||||
let mut headers = settings.headers.clone();
|
||||
if !settings.key.is_empty() {
|
||||
let key = hmac::Key::new(hmac::HMAC_SHA256, settings.key.as_bytes());
|
||||
let tag = hmac::sign(&key, body.as_bytes());
|
||||
|
||||
headers.insert(
|
||||
"X-Signature",
|
||||
STANDARD.encode(tag.as_ref()).parse().unwrap(),
|
||||
);
|
||||
}
|
||||
sign(&mut headers, &settings.key, &body);
|
||||
|
||||
// Send request
|
||||
let response = settings
|
||||
@@ -188,3 +180,150 @@ async fn post_webhook_events(
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
/// Adds the HMAC-SHA256 `X-Signature` a receiver checks, when the webhook has a key.
|
||||
fn sign(headers: &mut hyper::HeaderMap, key: &str, body: &str) {
|
||||
if !key.is_empty() {
|
||||
let key = hmac::Key::new(hmac::HMAC_SHA256, key.as_bytes());
|
||||
let tag = hmac::sign(&key, body.as_bytes());
|
||||
|
||||
headers.insert(
|
||||
"X-Signature",
|
||||
STANDARD.encode(tag.as_ref()).parse().unwrap(),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// inbuxa: "Send test" for a saved webhook (settings-reorg, Webhooks). One
|
||||
/// sample event, sent the way a real batch is: the same URL, headers, sign-in,
|
||||
/// signature, timeout and certificate checks. The event's type,
|
||||
/// `webhook.test`, is none the server raises, and an `X-Inbuxa-Test` header
|
||||
/// marks it, so a receiver can tell it apart. Answers the HTTP status, or why
|
||||
/// nothing came back.
|
||||
pub async fn send_test(hook: ®istry::schema::structs::WebHook) -> Result<u16, String> {
|
||||
let mut headers = hook
|
||||
.http_auth
|
||||
.build_headers(hook.http_headers.clone(), "application/json".into())
|
||||
.await
|
||||
.map_err(|err| format!("Unable to build HTTP headers: {err}"))?;
|
||||
let key = hook
|
||||
.signature_key
|
||||
.secret()
|
||||
.await
|
||||
.map_err(|err| format!("Unable to retrieve signature key: {err}"))?
|
||||
.unwrap_or_default()
|
||||
.into_owned();
|
||||
|
||||
let created = now();
|
||||
let body = serde_json::json!({
|
||||
"events": [{
|
||||
"id": format!("test-{created}"),
|
||||
"createdAt": mail_parser::DateTime::from_timestamp(created as i64).to_rfc3339(),
|
||||
"type": "webhook.test",
|
||||
"data": { "details": "A test from inbuxa Admin. Nothing happened on the server." },
|
||||
}]
|
||||
})
|
||||
.to_string();
|
||||
sign(&mut headers, &key, &body);
|
||||
headers.insert("X-Inbuxa-Test", "true".parse().unwrap());
|
||||
|
||||
let response = utils::http::http_client_builder(hook.allow_invalid_certs)
|
||||
.build()
|
||||
.map_err(|err| format!("Unable to build an HTTP client: {err}"))?
|
||||
.post(&hook.url)
|
||||
.timeout(hook.timeout.into_inner())
|
||||
.headers(headers)
|
||||
.body(body)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|err| format!("Webhook request to {} failed: {err}", hook.url))?;
|
||||
Ok(response.status().as_u16())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use registry::schema::structs::{SecretKeyOptional, SecretKeyValue, WebHook};
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
|
||||
/// One request in, the given status out; hands back what was received.
|
||||
async fn receiver(status: &'static str) -> (String, tokio::task::JoinHandle<String>) {
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let url = format!("http://{}/hook", listener.local_addr().unwrap());
|
||||
let task = tokio::spawn(async move {
|
||||
let (mut socket, _) = listener.accept().await.unwrap();
|
||||
let mut buf = Vec::new();
|
||||
let mut chunk = [0u8; 4096];
|
||||
loop {
|
||||
let n = socket.read(&mut chunk).await.unwrap();
|
||||
buf.extend_from_slice(&chunk[..n]);
|
||||
let text = String::from_utf8_lossy(&buf);
|
||||
if let Some(end) = text.find("\r\n\r\n") {
|
||||
let length = text[..end]
|
||||
.lines()
|
||||
.find_map(|l| {
|
||||
l.to_ascii_lowercase()
|
||||
.strip_prefix("content-length:")
|
||||
.map(|v| v.trim().parse::<usize>().unwrap())
|
||||
})
|
||||
.unwrap_or(0);
|
||||
if buf.len() >= end + 4 + length || n == 0 {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
socket
|
||||
.write_all(
|
||||
format!("HTTP/1.1 {status}\r\ncontent-length: 0\r\nconnection: close\r\n\r\n")
|
||||
.as_bytes(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
String::from_utf8_lossy(&buf).into_owned()
|
||||
});
|
||||
(url, task)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn send_test_signs_and_marks_the_sample() {
|
||||
let (url, task) = receiver("204 No Content").await;
|
||||
let hook = WebHook {
|
||||
url,
|
||||
enable: false,
|
||||
signature_key: SecretKeyOptional::Value(SecretKeyValue { secret: "k".into() }),
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(send_test(&hook).await, Ok(204));
|
||||
|
||||
let request = task.await.unwrap();
|
||||
let (head, body) = request.split_once("\r\n\r\n").unwrap();
|
||||
let head = head.to_ascii_lowercase();
|
||||
assert!(head.contains("x-inbuxa-test: true"), "{head}");
|
||||
let parsed: serde_json::Value = serde_json::from_str(body).unwrap();
|
||||
assert_eq!(parsed["events"][0]["type"], "webhook.test");
|
||||
let tag = hmac::sign(&hmac::Key::new(hmac::HMAC_SHA256, b"k"), body.as_bytes());
|
||||
assert!(
|
||||
head.contains(&format!(
|
||||
"x-signature: {}",
|
||||
STANDARD.encode(tag.as_ref()).to_ascii_lowercase()
|
||||
)),
|
||||
"{head}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn send_test_reports_what_came_back() {
|
||||
let (url, _task) = receiver("403 Forbidden").await;
|
||||
let hook = WebHook {
|
||||
url,
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(send_test(&hook).await, Ok(403));
|
||||
|
||||
let hook = WebHook {
|
||||
url: "http://127.0.0.1:9/hook".into(),
|
||||
..Default::default()
|
||||
};
|
||||
assert!(send_test(&hook).await.unwrap_err().contains("failed"));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -131,6 +131,29 @@ impl ManagementApi for Server {
|
||||
let answer = jmap::inbuxa::directory_test::test(self, &request).await?;
|
||||
Ok(JsonResponse::new(answer).no_cache().into_http_response())
|
||||
}
|
||||
// inbuxa: send one sample event to a saved webhook
|
||||
"webhook" if is_post && path.get(1).copied() == Some("test") => {
|
||||
let (_in_flight, access_token) = self.authenticate_headers(req, session).await?;
|
||||
jmap::inbuxa::webhook_test::assert_allowed(&access_token)?;
|
||||
let request = body
|
||||
.as_deref()
|
||||
.and_then(|body| serde_json::from_slice::<serde_json::Value>(body).ok())
|
||||
.unwrap_or_default();
|
||||
let answer = jmap::inbuxa::webhook_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?;
|
||||
|
||||
@@ -798,6 +798,10 @@ mod tests {
|
||||
assert!(delivery_facts(&mut Facts::default(), &message, "[email protected]").is_err());
|
||||
}
|
||||
|
||||
fn is_timestamp(value: &str) -> bool {
|
||||
chrono::DateTime::parse_from_rfc3339(value).is_ok()
|
||||
}
|
||||
|
||||
/// The settings questions a release prepares answers for (EX-26): every
|
||||
/// non-secret property of every settings object, at the object's own
|
||||
/// default, built exactly as a live question is.
|
||||
@@ -842,6 +846,12 @@ mod tests {
|
||||
if info.secret {
|
||||
continue;
|
||||
}
|
||||
// A date's default is the moment the object is built, so its
|
||||
// question changes every run and no live question ever
|
||||
// matches it: nothing worth preparing.
|
||||
if matches!(map[&property].as_str(), Some(v) if is_timestamp(v)) {
|
||||
continue;
|
||||
}
|
||||
let mut facts = Facts::default();
|
||||
push_setting(&mut facts, &object, &property, &info, &map[&property]);
|
||||
out.push((object.clone(), property, facts));
|
||||
@@ -856,6 +866,7 @@ mod tests {
|
||||
assert!(questions.len() > 500, "found {}", questions.len());
|
||||
assert!(questions.iter().any(|(o, p, _)| o == "x:Domain" && p == "dnsManagement"));
|
||||
assert!(!questions.iter().any(|(o, p, _)| o == "x:AiModel" && p == "httpAuth"));
|
||||
assert!(!questions.iter().any(|(o, p, _)| o == "x:Account" && p == "createdAt"));
|
||||
}
|
||||
|
||||
/// Writes `resources/explain/settings.json.gz` (EX-26). Run before a
|
||||
|
||||
@@ -18,6 +18,7 @@ pub mod ai_limits;
|
||||
pub mod log_settings;
|
||||
pub mod data_inventory;
|
||||
pub mod directory_test;
|
||||
pub mod webhook_test;
|
||||
pub mod explanation;
|
||||
pub mod protocol_policy;
|
||||
pub mod tenant_protocol_policy;
|
||||
|
||||
@@ -0,0 +1,61 @@
|
||||
/*
|
||||
* SPDX-FileCopyrightText: 2026 Coffey Labs
|
||||
*
|
||||
* SPDX-License-Identifier: AGPL-3.0-only
|
||||
*/
|
||||
|
||||
//! `POST /api/webhook/test`: send one sample event to a saved webhook
|
||||
//! (settings-reorg, Webhooks "Send test").
|
||||
//!
|
||||
//! ```json
|
||||
//! {"webhookId": "b"}
|
||||
//! ```
|
||||
//!
|
||||
//! The answer is `{"sent": true, "status": 200, "ms": 84}` when the receiver
|
||||
//! answered 2xx, `{"sent": false, "status": 403, …}` when it answered
|
||||
//! otherwise, and `{"sent": false, "error": "…"}` when nothing came back. The
|
||||
//! webhook is used as saved, even when it's off, so it can be tried before
|
||||
//! it's switched on. The request goes where the saved webhook already sends,
|
||||
//! so this gives nobody a reach they didn't have.
|
||||
//!
|
||||
//! For server-level administrators who may change webhooks.
|
||||
|
||||
use common::{Server, auth::AccessToken};
|
||||
use registry::schema::{enums::Permission, structs::WebHook};
|
||||
use serde_json::{Value, json};
|
||||
use std::{str::FromStr, time::Instant};
|
||||
use types::id::Id;
|
||||
|
||||
pub fn assert_allowed(access_token: &AccessToken) -> trc::Result<()> {
|
||||
if access_token.tenant_id().is_some() {
|
||||
return Err(trc::JmapEvent::Forbidden
|
||||
.into_err()
|
||||
.details("Webhook tests are for server-level administrators."));
|
||||
}
|
||||
access_token.enforce_permission(Permission::SysWebHookUpdate)
|
||||
}
|
||||
|
||||
pub async fn test(server: &Server, body: &Value) -> trc::Result<Value> {
|
||||
let webhook_id = body
|
||||
.get("webhookId")
|
||||
.and_then(Value::as_str)
|
||||
.and_then(|id| Id::from_str(id).ok())
|
||||
.ok_or_else(|| {
|
||||
trc::ResourceEvent::BadParameters
|
||||
.into_err()
|
||||
.details("Expected {\"webhookId\": …}")
|
||||
})?;
|
||||
let Some(hook) = server.registry().object::<WebHook>(webhook_id).await? else {
|
||||
return Ok(json!({ "sent": false, "error": "There's no such webhook. Save it first." }));
|
||||
};
|
||||
|
||||
let started = Instant::now();
|
||||
Ok(match common::telemetry::webhooks::send_test(&hook).await {
|
||||
Ok(status) => json!({
|
||||
"sent": (200..300).contains(&status),
|
||||
"status": status,
|
||||
"ms": started.elapsed().as_millis() as u64,
|
||||
}),
|
||||
Err(error) => json!({ "sent": false, "error": error }),
|
||||
})
|
||||
}
|
||||
@@ -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",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -81,7 +81,7 @@ fn legacy_setting(name: &str, is_set: impl Fn(&str) -> bool) -> Option<String> {
|
||||
#[macro_export]
|
||||
macro_rules! brand_version {
|
||||
() => {
|
||||
"2026.9.28.5"
|
||||
"2026.9.28.6"
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
Binary file not shown.
Reference in New Issue
Block a user