feat: harden the RNS carrier: idle-link reaper, refusal logging, patched microReticulum fork

This commit is contained in:
randogoth 2026-09-29 10:27:44 +03:00
parent 07f064460f
commit 8ea6865e7e
6 changed files with 110 additions and 16 deletions

View file

@ -12,6 +12,7 @@
use std::collections::{HashMap, HashSet};
use std::ffi::CString;
use std::sync::{Mutex, OnceLock};
use std::time::{Duration, Instant};
use crate::bind::TransportBindValues;
use crate::proto::{AUTH_FULL_TOKENS, INTERNAL_ERROR, MALFORMED, RATE_LIMITED, TOO_LARGE};
@ -59,6 +60,7 @@ pub struct RnsArgs {
pub udp_listen_port: u16,
pub udp_forward_host: Option<String>,
pub udp_forward_port: u16,
pub link_idle_secs: u64,
pub max_tokens: u16,
pub main_quota: i64,
pub requests_quota: i64,
@ -71,6 +73,10 @@ struct RnsState {
destination: [u8; 16],
sessions: Mutex<HashMap<[u8; 16], Session<'static>>>,
links: Mutex<HashSet<[u8; 16]>>,
/// Last-seen instant per link, the reaper's clock. Updated before the
/// limiters so a client still sending — even one being refused — counts
/// as alive; only a silent link is reapable.
last_activity: Mutex<HashMap<[u8; 16], Instant>>,
request_limiter: RateLimiter,
byte_limiter: ByteRateLimiter,
max_links: usize,
@ -117,6 +123,7 @@ pub fn start(args: RnsArgs) -> anyhow::Result<String> {
destination,
sessions: Mutex::new(HashMap::new()),
links: Mutex::new(HashSet::new()),
last_activity: Mutex::new(HashMap::new()),
request_limiter: RateLimiter::new(args.rate_link_requests),
byte_limiter: ByteRateLimiter::new(args.rate_link_bytes),
max_links: args.max_links,
@ -124,6 +131,13 @@ pub fn start(args: RnsArgs) -> anyhow::Result<String> {
};
STATE.get_or_init(|| state);
// microReticulum never times a link out and the shim has no server-side
// close, so a client that dies mid-link would hold its Store handle and
// maxLinks slot forever; 0 keeps the old never-expire behavior.
if args.link_idle_secs > 0 {
std::thread::spawn(move || reap_loop(Duration::from_secs(args.link_idle_secs)));
}
let storage_dir = CString::new(args.storage_dir).unwrap();
let listen_host = CString::new(args.udp_listen_host).unwrap();
let forward_host = args.udp_forward_host.map(|h| CString::new(h).unwrap());
@ -185,6 +199,9 @@ fn handle_request(request: &[u8], link_id: &[u8; 16]) -> Vec<u8> {
// an AUTH carrying a full token set, whichever is larger.
let max_request = (state.config.max_envelope + 34)
.max(AUTH_FULL_TOKENS + 32 * state.config.max_tokens as usize);
// Any request, even one about to be refused, proves the link is alive.
state.last_activity.lock().unwrap().insert(*link_id, Instant::now());
if request.len() > max_request {
log::warn!(
"request of {} bytes over {} byte cap",
@ -194,6 +211,9 @@ fn handle_request(request: &[u8], link_id: &[u8; 16]) -> Vec<u8> {
return response(TOO_LARGE);
}
if !state.request_limiter.allow(&link) || !state.byte_limiter.allow(&link, request.len()) {
// Pre-dispatch refusals never reach the session metrics, so they get
// their own line or the journal undercounts a spinning client.
log::debug!("link {link} rate limited: request refused");
return response(RATE_LIMITED);
}
@ -264,6 +284,45 @@ extern "C" fn smolmail_rns_take_response(out: *mut u8, cap: usize) -> usize {
}
}
/// Links whose last activity is at least `idle` old, with their age.
fn expired_links(
last_activity: &HashMap<[u8; 16], Instant>,
now: Instant,
idle: Duration,
) -> Vec<([u8; 16], Duration)> {
last_activity
.iter()
.filter(|(_, seen)| now.duration_since(**seen) >= idle)
.map(|(id, seen)| (*id, now.duration_since(*seen)))
.collect()
}
/// Evicts links that have gone silent: microReticulum never times a link
/// out and the shim has no server-side close, so a client that dies
/// mid-link would otherwise hold its Store handle and maxLinks slot
/// forever. A reaped link that somehow speaks again is treated as a fresh
/// session (unauthenticated until the next AUTH), which is safe.
fn reap_loop(idle: Duration) {
let interval = (idle / 4).clamp(Duration::from_secs(1), Duration::from_secs(60));
loop {
std::thread::sleep(interval);
let Some(state) = STATE.get() else {
return;
};
let expired = expired_links(&state.last_activity.lock().unwrap(), Instant::now(), idle);
for (link, age) in expired {
state.links.lock().unwrap().remove(&link);
state.sessions.lock().unwrap().remove(&link);
state.last_activity.lock().unwrap().remove(&link);
log::debug!(
"link {} reaped after {:?} idle",
data_encoding::HEXLOWER.encode(&link),
age
);
}
}
}
#[no_mangle]
extern "C" fn smolmail_rns_on_link_opened(link_id: *const u8) -> i32 {
let Some(state) = STATE.get() else {
@ -280,6 +339,9 @@ extern "C" fn smolmail_rns_on_link_opened(link_id: *const u8) -> i32 {
return 1;
}
links.insert(link);
if let Some(state) = STATE.get() {
state.last_activity.lock().unwrap().insert(link, Instant::now());
}
log::debug!("link {} opened ({} open)", data_encoding::HEXLOWER.encode(&link), links.len());
0
}
@ -290,6 +352,7 @@ extern "C" fn smolmail_rns_on_link_closed(link_id: *const u8) {
let link: [u8; 16] = unsafe { std::slice::from_raw_parts(link_id, 16).try_into().unwrap() };
state.links.lock().unwrap().remove(&link);
state.sessions.lock().unwrap().remove(&link);
state.last_activity.lock().unwrap().remove(&link);
log::debug!("link {} closed", data_encoding::HEXLOWER.encode(&link));
}
}
@ -311,4 +374,18 @@ mod tests {
"799855f4955f1b09fd20a13cd84f4e71"
);
}
#[test]
fn expired_links_selects_only_silent_ones() {
let mut last_activity = HashMap::new();
let now = Instant::now();
let idle = Duration::from_secs(300);
last_activity.insert([1u8; 16], now - Duration::from_secs(301)); // silent
last_activity.insert([2u8; 16], now - Duration::from_secs(299)); // active
last_activity.insert([3u8; 16], now); // just seen
let expired = expired_links(&last_activity, now, idle);
assert_eq!(expired.len(), 1);
assert_eq!(expired[0].0, [1u8; 16]);
assert!(expired[0].1 >= idle);
}
}