Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
82e1740d3c | ||
|
|
4dc1bcd546 |
@@ -225,6 +225,20 @@ state change; rate-limit pings), tickets from friends (validate defensively, no
|
||||
auto-join), the discovery publish (only when toggled, ideally auto-expiring).
|
||||
`cargo audit` (JSON store → no new deps expected). Field test on dopedart.
|
||||
|
||||
**Local hardening DONE 2026-06-27:** inbound friend-presence replies are now
|
||||
rate-limited per authenticated friend id (`PresenceRateLimiter`: burst 4, refill
|
||||
1/15s) and wired into the live friends listener before it builds a `Pong`; denied
|
||||
probes get the same silent no-data close as unauthorized probes. Existing
|
||||
defensive reply handling still validates room tickets against the authenticated
|
||||
friend id and never auto-joins. Verified with `cargo test presence`,
|
||||
`cargo test --lib`, `cargo clippy --all-targets -- -D warnings`, and
|
||||
`cargo audit --no-fetch --stale` (local DB; reports only the two already-allowed
|
||||
unmaintained advisories in `deny.toml`). A fresh advisory fetch was blocked in
|
||||
this sandbox by network restrictions.
|
||||
|
||||
**Remaining:** live 2-machine field test on dopedart, a fresh online
|
||||
`cargo audit`, and any follow-up findings from that test.
|
||||
|
||||
## The connect flow (the user's scenario, end to end)
|
||||
1. Friend X, at a coffee shop, opens peerspeak and starts a gathering labeled
|
||||
"HangOut."
|
||||
|
||||
+15
-3
@@ -1074,22 +1074,34 @@ async fn run_core_loop(
|
||||
// Join, cleared on Leave.
|
||||
let current_room: Arc<std::sync::Mutex<Option<crate::presence::RoomPresence>>> =
|
||||
Arc::new(std::sync::Mutex::new(None));
|
||||
let presence_rate_limiter =
|
||||
Arc::new(std::sync::Mutex::new(crate::presence::PresenceRateLimiter::default()));
|
||||
|
||||
// Reply policy for the idle friends listener (B2): answer friends only, never
|
||||
// while invisible (`should_answer`), and report our current gathering so a friend
|
||||
// can one-click join. Reads the shared snapshots, so it stays correct as they
|
||||
// change and survives a network-stack rebuild. Pure-sync (no awaits, no lock held
|
||||
// across one). Built once and handed to every `build_net_stack`.
|
||||
// can one-click join. Rate-limits allowed friends before building a reply, so a
|
||||
// spammy saved peer gets the same silent close as an unauthorized peer. Reads the
|
||||
// shared snapshots, so it stays correct as they change and survives a network-stack
|
||||
// rebuild. Pure-sync (no awaits, no lock held across one). Built once and handed
|
||||
// to every `build_net_stack`.
|
||||
let friends_handler: crate::presence_net::Handler = {
|
||||
let friends = friends.clone();
|
||||
let presence_mode = presence_mode.clone();
|
||||
let current_room = current_room.clone();
|
||||
let presence_rate_limiter = presence_rate_limiter.clone();
|
||||
Arc::new(move |from| {
|
||||
let mode = *presence_mode.lock().unwrap();
|
||||
let allowed = crate::presence::should_answer(&from, &friends.lock().unwrap(), mode);
|
||||
if !allowed {
|
||||
return None;
|
||||
}
|
||||
if !presence_rate_limiter
|
||||
.lock()
|
||||
.unwrap()
|
||||
.allow(from, std::time::Instant::now())
|
||||
{
|
||||
return None;
|
||||
}
|
||||
let room = current_room.lock().unwrap().clone();
|
||||
Some(crate::presence::ControlMsg::Pong { room })
|
||||
})
|
||||
|
||||
@@ -15,6 +15,16 @@
|
||||
use crate::friends::FriendStore;
|
||||
use iroh::EndpointId;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::HashMap;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
/// Maximum immediate presence replies to one friend before throttling. Normal
|
||||
/// presence polling is once per minute, so this only catches repeated/manual or
|
||||
/// abusive probes while still allowing a short burst after app startup.
|
||||
pub const PRESENCE_RATE_LIMIT_BURST: u32 = 4;
|
||||
|
||||
/// Refill one presence-reply token per friend at this cadence.
|
||||
pub const PRESENCE_RATE_LIMIT_REFILL: Duration = Duration::from_secs(15);
|
||||
|
||||
/// The user's presence posture — how reachable they are to friends while idle.
|
||||
/// Persisted in `AppConfig`; the default keeps you privately reachable to friends
|
||||
@@ -100,6 +110,48 @@ pub fn should_answer(from: &EndpointId, friends: &FriendStore, mode: PresenceMod
|
||||
mode.answers_pings() && friends.contains(from)
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
struct RateBucket {
|
||||
tokens: u32,
|
||||
last_refill: Instant,
|
||||
}
|
||||
|
||||
/// Per-friend limiter for inbound presence pings. It is intentionally keyed by
|
||||
/// the authenticated connection id, not payload data. Callers should only invoke
|
||||
/// it after [`should_answer`] passes, so strangers do not consume memory here.
|
||||
#[derive(Debug, Default, Clone)]
|
||||
pub struct PresenceRateLimiter {
|
||||
buckets: HashMap<EndpointId, RateBucket>,
|
||||
}
|
||||
|
||||
impl PresenceRateLimiter {
|
||||
/// Return whether `from` may receive a presence reply at `now`.
|
||||
///
|
||||
/// This is a token bucket: each friend starts with a small burst and regains
|
||||
/// one token every [`PRESENCE_RATE_LIMIT_REFILL`]. A denied probe should be
|
||||
/// answered with no data, matching the listener's "reveal nothing" policy.
|
||||
pub fn allow(&mut self, from: EndpointId, now: Instant) -> bool {
|
||||
let bucket = self.buckets.entry(from).or_insert(RateBucket {
|
||||
tokens: PRESENCE_RATE_LIMIT_BURST,
|
||||
last_refill: now,
|
||||
});
|
||||
|
||||
let elapsed = now.saturating_duration_since(bucket.last_refill);
|
||||
let refill = elapsed.as_secs() / PRESENCE_RATE_LIMIT_REFILL.as_secs();
|
||||
if refill > 0 {
|
||||
let refill = refill.min(u32::MAX as u64) as u32;
|
||||
bucket.tokens = PRESENCE_RATE_LIMIT_BURST.min(bucket.tokens.saturating_add(refill));
|
||||
bucket.last_refill = now;
|
||||
}
|
||||
|
||||
if bucket.tokens == 0 {
|
||||
return false;
|
||||
}
|
||||
bucket.tokens -= 1;
|
||||
true
|
||||
}
|
||||
}
|
||||
|
||||
/// What we learned about a friend from a successful ping reply.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum FriendPresence {
|
||||
@@ -177,6 +229,36 @@ mod tests {
|
||||
assert!(!should_answer(&stranger, &friends, PresenceMode::Invisible));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn presence_rate_limiter_allows_a_small_burst_then_refills() {
|
||||
let mut limiter = PresenceRateLimiter::default();
|
||||
let friend = id();
|
||||
let now = Instant::now();
|
||||
|
||||
for _ in 0..PRESENCE_RATE_LIMIT_BURST {
|
||||
assert!(limiter.allow(friend, now));
|
||||
}
|
||||
assert!(!limiter.allow(friend, now));
|
||||
assert!(!limiter.allow(friend, now + PRESENCE_RATE_LIMIT_REFILL - Duration::from_millis(1)));
|
||||
|
||||
assert!(limiter.allow(friend, now + PRESENCE_RATE_LIMIT_REFILL));
|
||||
assert!(!limiter.allow(friend, now + PRESENCE_RATE_LIMIT_REFILL));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn presence_rate_limiter_is_per_peer() {
|
||||
let mut limiter = PresenceRateLimiter::default();
|
||||
let a = id();
|
||||
let b = id();
|
||||
let now = Instant::now();
|
||||
|
||||
for _ in 0..PRESENCE_RATE_LIMIT_BURST {
|
||||
assert!(limiter.allow(a, now));
|
||||
}
|
||||
assert!(!limiter.allow(a, now));
|
||||
assert!(limiter.allow(b, now));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn presence_mode_flags() {
|
||||
assert!(PresenceMode::Discoverable.publishes_to_discovery());
|
||||
|
||||
+109
-4
@@ -332,7 +332,13 @@ pub async fn spawn_host(
|
||||
.args(host_args(audio_app))
|
||||
.stdin(Stdio::null())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::null())
|
||||
// Capture stderr (not null): pixelpass prints its startup precondition
|
||||
// failures there — a missing GStreamer plugin / `pactl`, each with an
|
||||
// actionable "Install hint: sudo apt install ..." line. If the host dies
|
||||
// before its ticket we fold that tail into our error so the user sees
|
||||
// *what to install* instead of a dead-end "exited before a ticket". On
|
||||
// the success path we drain it in the background so the pipe can't fill.
|
||||
.stderr(Stdio::piped())
|
||||
.kill_on_drop(true)
|
||||
.spawn()?;
|
||||
|
||||
@@ -340,6 +346,7 @@ pub async fn spawn_host(
|
||||
.stdout
|
||||
.take()
|
||||
.ok_or_else(|| std::io::Error::other("pixelpass host stdout missing"))?;
|
||||
let stderr = child.stderr.take();
|
||||
let mut lines = BufReader::new(stdout).lines();
|
||||
|
||||
let ticket = match read_until(&mut lines, |e| match e {
|
||||
@@ -351,9 +358,10 @@ pub async fn spawn_host(
|
||||
Ok(Some(t)) => t,
|
||||
Ok(None) => {
|
||||
let _ = child.kill().await;
|
||||
return Err(std::io::Error::other(
|
||||
"pixelpass host exited before emitting a ticket",
|
||||
));
|
||||
let detail = read_stderr_tail(stderr).await;
|
||||
return Err(std::io::Error::other(format!(
|
||||
"pixelpass host exited before emitting a ticket{detail}"
|
||||
)));
|
||||
}
|
||||
Err(e) => {
|
||||
let _ = child.kill().await;
|
||||
@@ -361,10 +369,65 @@ pub async fn spawn_host(
|
||||
}
|
||||
};
|
||||
|
||||
if let Some(stderr) = stderr {
|
||||
drain_stderr_in_background(stderr);
|
||||
}
|
||||
drain_in_background(lines, "host", notices);
|
||||
Ok((child, ticket))
|
||||
}
|
||||
|
||||
/// Read a killed pixelpass child's stderr to EOF and reduce it to a short,
|
||||
/// user-facing diagnostic tail via [`pixelpass_failure_detail`]. Bounded: the
|
||||
/// caller kills the child first, so the pipe EOFs promptly. Returns an empty
|
||||
/// string when stderr was already taken or carried nothing useful.
|
||||
async fn read_stderr_tail(stderr: Option<tokio::process::ChildStderr>) -> String {
|
||||
use tokio::io::AsyncReadExt;
|
||||
let Some(mut stderr) = stderr else {
|
||||
return String::new();
|
||||
};
|
||||
let mut buf = Vec::new();
|
||||
let _ = stderr.read_to_end(&mut buf).await;
|
||||
pixelpass_failure_detail(&String::from_utf8_lossy(&buf))
|
||||
}
|
||||
|
||||
/// Discard a running pixelpass child's stderr in the background so its pipe
|
||||
/// can't fill and stall the host (mirrors [`drain_in_background`] for stdout).
|
||||
fn drain_stderr_in_background(mut stderr: tokio::process::ChildStderr) {
|
||||
use tokio::io::AsyncReadExt;
|
||||
tokio::spawn(async move {
|
||||
let mut buf = [0u8; 4096];
|
||||
while let Ok(n) = stderr.read(&mut buf).await {
|
||||
if n == 0 {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/// Extract a human-useful tail from a failed pixelpass child's stderr to append
|
||||
/// to our error. pixelpass writes actionable startup errors there (a missing
|
||||
/// GStreamer element / `pactl` plus an `Install hint: sudo apt install ...`
|
||||
/// line), which is exactly what a freshly-installed host needs to see. The
|
||||
/// decorative host banner (box-drawing) is dropped — it only prints on the
|
||||
/// success path, but we filter it defensively. Pure: no I/O. Returns an empty
|
||||
/// string when there's nothing worth surfacing (so callers can append blindly).
|
||||
pub fn pixelpass_failure_detail(stderr: &str) -> String {
|
||||
let useful: Vec<&str> = stderr
|
||||
.lines()
|
||||
.map(str::trim_end)
|
||||
.filter(|l| !l.trim().is_empty())
|
||||
.filter(|l| !l.trim_start().starts_with(['│', '┌', '└', '├']))
|
||||
.collect();
|
||||
if useful.is_empty() {
|
||||
return String::new();
|
||||
}
|
||||
// The anyhow error and its install hint are the *last* lines printed, so
|
||||
// keep the tail rather than the head.
|
||||
const MAX_LINES: usize = 12;
|
||||
let start = useful.len().saturating_sub(MAX_LINES);
|
||||
format!("\n\npixelpass reported:\n{}", useful[start..].join("\n"))
|
||||
}
|
||||
|
||||
/// Spawn a pixelpass viewer for `ticket`, wait for it to connect, and open the
|
||||
/// stream in a local player (mpv, falling back to vlc). Returns the live viewer
|
||||
/// child so the caller can kill it on room-leave; it also self-exits when the
|
||||
@@ -633,6 +696,48 @@ mod tests {
|
||||
assert_eq!(parse_audio_apps(stdout.as_bytes()), vec!["mpv".to_string()]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn failure_detail_surfaces_install_hint_and_drops_banner() {
|
||||
// The real shape of a fresh-host failure: anyhow error + install hint on
|
||||
// stderr. We must keep those (so the user knows what to apt install) and
|
||||
// drop the decorative banner box-drawing lines.
|
||||
let stderr = "\
|
||||
┌─ PixelPass · host ─────────────────────────────────────────
|
||||
│ display server : Wayland
|
||||
└────────────────────────────────────────────────────────────
|
||||
Error: GStreamer element `vah264enc` not available.
|
||||
Install hint: sudo apt install gstreamer1.0-plugins-bad
|
||||
";
|
||||
let detail = pixelpass_failure_detail(stderr);
|
||||
assert!(detail.starts_with("\n\npixelpass reported:\n"));
|
||||
assert!(detail.contains("vah264enc` not available"));
|
||||
assert!(detail.contains("sudo apt install gstreamer1.0-plugins-bad"));
|
||||
assert!(!detail.contains('│'), "banner box-drawing must be dropped");
|
||||
assert!(!detail.contains('┌'));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn failure_detail_empty_when_nothing_useful() {
|
||||
// Blank / banner-only stderr yields an empty string so the caller can
|
||||
// append it to the base message unconditionally without trailing noise.
|
||||
assert_eq!(pixelpass_failure_detail(""), "");
|
||||
assert_eq!(pixelpass_failure_detail(" \n \n"), "");
|
||||
assert_eq!(
|
||||
pixelpass_failure_detail("│ display server : Wayland\n│ capture : x\n"),
|
||||
""
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn failure_detail_keeps_only_the_tail() {
|
||||
// A long stderr is truncated to its last lines (where the real error
|
||||
// and hint live), not its head.
|
||||
let body: String = (0..30).map(|i| format!("line {i}\n")).collect();
|
||||
let detail = pixelpass_failure_detail(&body);
|
||||
assert!(detail.contains("line 29"));
|
||||
assert!(!detail.contains("line 0\n"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn help_probe_detects_strict_audio_flag() {
|
||||
// A new pixelpass advertises the flag; an old one doesn't. The probe must
|
||||
|
||||
Reference in New Issue
Block a user