Compare commits

...
Author SHA1 Message Date
molluskandClaude Opus 4.8 319d0c5e29 S11: make presence/discovery state honest on apply failure
SetPresenceMode and the Discoverable time-box auto-revert both committed
the new presence_mode to local state *before* apply_discovery and only
log_msg'd on failure, so a failed off-transition could leave the n0 DNS
PkarrPublisher running while the UI showed not-discoverable (privacy /
reality mismatch — security-open-handoff S11, from the W7 P7 review).

Fix (Codex, senior-reviewed):
- discovery.rs: pure resolve_presence_transition(prev, requested, apply_ok)
  -> (mode, Option<error>) seam — on failure keep the previous (truthful)
  mode and surface a message. +4 unit tests.
- apply_discovery now builds the replacement resolver/publisher services
  BEFORE clearing the service set, so a builder failure leaves the old
  posture fully intact (no partial state) — "keep previous mode" is then
  provably truthful.
- Both SetPresenceMode and the time-box revert apply discovery first, route
  through the seam, commit only the truthful mode, and surface failures via
  the existing PresenceModeReverted (corrects the picker) + UiEvent::Error.
  No new wire/event variant.
- A failed off-transition stays Discoverable and arms a 60s retry
  (DISCOVERY_REVERT_RETRY) so the beacon never stands stuck.
- P3 notes documented: relay-resolve exposes n0 query metadata (by design);
  no explicit iroh unpublish API exists, so the bounded ~30s pkarr TTL
  linger is documented, not behavior-changed; DirectOnly stays no-n0.

306 lib tests / clippy --all-targets / release all green (re-run by senior).
Runtime publish-stop behavior still wants a 2-machine / packet-capture check.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-18 03:32:38 -04:00
mollusk d56c2c90b2 Merge codex-security-hardening: S10 log redaction + T1/T2/T5/T6/T7 trust-boundary fixes 2026-06-18 03:05:56 -04:00
molluskandClaude Opus 4.8 5086e86bd2 Security hardening: log redaction + 5 trust-boundary fixes (S10, T1/T2/T5/T6/T7)
Codex (gpt-5.5) implementer branch, senior-reviewed.

- S10 (High): redact capabilities/chat from logs; create log 0600 + chmod
  existing; rotate at 5 MiB. New short_id/short_bytes_hex/redact_for_log seams.
- T1 (P2): friends-ALPN authorizes (handler(from)) before reading any peer
  bytes; unauthorized conns closed pre-read (DoS relief).
- T2 (P2): per-(author,kind) replay gate on state-changing gossip (Announce/
  Leave) only; Chat bypasses it, preserving the S2 no-monotonic-ts decision.
- T5 (P2): cap inbound Opus datagrams at 4 + MAX_OPUS_PAYLOAD (4000).
- T6 (P2): bind friend-Pong room ticket host to the authenticated responder
  (interpret_pong/probe now thread the remote id) — blocks Join-button
  redirect/phishing. Non-regressive given the W7 P3 restamp design.
- T7 (P3): sanitize_ticket caps/validates PeerState.sharing at gossip ingest
  so invalid offers never render a Watch button.

302 lib tests pass (was 291), clippy --all-targets clean, release builds.
Tests-green only; DoS relief + 2-machine replay/redirect behavior want a
field test. W7 P7 (n0 DNS privacy) reviewed read-only — findings to triage.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-18 03:05:56 -04:00
molluskandClaude Opus 4.8 54780fa73b Remove Codex task-report.md from repo root
Transient implementer handoff note; its content is preserved in the
handoff docs. Not repo content.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-18 01:51:26 -04:00
molluskandClaude Opus 4.8 b1aa751a84 Merge codex-home-empty-state-ui: focused empty-state Home layout
Fresh/empty Home keeps Create/Join dominant via a tested home_layout_mode
seam (FocusedEmpty / ThreeColumn / Stacked); once Recents or Friends has
content the normal three-card layout returns. Quieter empty-state cards.

Conflict resolution:
- HomeLayoutMode enum/fn coexists with the SettingsCategory enum (separate
  derives); both unit tests kept.
- Top bar: the wishlist Hotkeys-info button is always shown; the room-layout
  button is hidden on Home (home-empty's intent) and shown in Room. The
  auto-merge had wedged the info tooltip into the conditional as a stray
  expression — split into separate info_button / layout_button bindings.

291 lib tests pass, clippy --all-targets clean, bin builds.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-18 01:36:19 -04:00
molluskandClaude Opus 4.8 9efab491c7 Merge codex-settings-category-nav: category-navigated Settings
Settings is split into navigable categories (left sidebar ≥820px wide,
pick_list dropdown below) instead of one long scroll. Integrated with the
wishlist branch's hotkey editor by giving it its own "Hotkeys" category
(7 categories total: Audio, Hotkeys, Recording, Profile, Appearance,
Network, Notifications).

Conflict resolution: the wishlist branch had inserted a Hotkeys section
into the old long-scroll between Microphone and Recording; relocated it
into a dedicated SettingsCategory::Hotkeys arm and updated the category
stability test.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-18 01:32:17 -04:00
molluskandClaude Opus 4.8 f3f399a748 Merge codex-wishlist-audio-hotkeys: per-peer EQ (W2), spatial pan + stereo bus (W1), focused hotkeys (W5), + A21/A22/A14 fixes
W2: src/audio/eq.rs 3-band RBJ biquad EQ, per-peer, flat=bypass.
W1: src/audio/pan.rs constant-power pan; mixer/playback/recorder converted
    to a stereo bus, bit-for-bit dual-mono at pan=0.
W5: src/hotkeys.rs config-backed focused hotkey map + Settings editor + info popup.
A21: jitter resets on large seq discontinuities (sender restart / far jump).
A22: WAV writer guards RIFF/data size overflow.
A14: orderly window-close shutdown (finalize recordings, leave room, close net).
W3 (PipeWire routing) intentionally left as a design note.

No new deps; no wire/serialization changes. Tests-green only; audio + 2-machine
field verification pending.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-18 01:29:45 -04:00
molluskandClaude Opus 4.8 1afdccbefe Merge codex-security-s9: bind gossip Announce addr to authenticated author (S9)
Reject signed Announce(PeerState) whose embedded state.addr.id does not
match the authenticated payload.author, closing the residual S2 gap where
a valid signer could advertise another node's EndpointAddr.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-18 01:29:34 -04:00
mollusk 7724da73b8 Refine home empty-state layout 2026-06-17 17:06:45 -04:00
mollusk 8982df364e Fix gossip announce address binding 2026-06-17 16:17:03 -04:00
mollusk 33e3998e7c Add orderly shutdown on window close 2026-06-16 17:37:00 -04:00
mollusk 44bad7b70b Fix jitter restart and WAV size overflow 2026-06-16 17:28:41 -04:00
mollusk 20643a24de Add audio controls and focused hotkeys 2026-06-16 17:23:38 -04:00
18 changed files with 2337 additions and 278 deletions
+622 -120
View File
File diff suppressed because it is too large Load Diff
+316
View File
@@ -0,0 +1,316 @@
//! Per-peer listener-side voice EQ.
//!
//! The EQ is deliberately small and local: three RBJ cookbook biquads at fixed
//! voice-oriented frequencies, with only gain exposed to the UI. State lives per
//! peer in the playout mixer so filter delay registers are continuous across 20ms
//! Opus frames; flat settings are treated as bypass so the default path is cheap
//! and sample-exact.
use serde::{Deserialize, Serialize};
const DEFAULT_SAMPLE_RATE: f32 = 48_000.0;
const LOW_SHELF_HZ: f32 = 160.0;
const MID_PEAK_HZ: f32 = 2_400.0;
const HIGH_SHELF_HZ: f32 = 6_500.0;
const MID_Q: f32 = 1.0;
const SHELF_Q: f32 = std::f32::consts::FRAC_1_SQRT_2;
const FLAT_EPSILON_DB: f32 = 0.001;
/// UI and config clamp for each band. Wide enough to be useful for voice, narrow
/// enough that a peer cannot accidentally make the listener-side limiter do all
/// the work.
pub const EQ_GAIN_DB_MIN: f32 = -12.0;
pub const EQ_GAIN_DB_MAX: f32 = 12.0;
/// Persisted per-peer EQ gains, in decibels. `Default` is flat/bypassed.
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq)]
pub struct EqSettings {
#[serde(default)]
pub low_gain_db: f32,
#[serde(default)]
pub mid_gain_db: f32,
#[serde(default)]
pub high_gain_db: f32,
}
impl Default for EqSettings {
fn default() -> Self {
Self {
low_gain_db: 0.0,
mid_gain_db: 0.0,
high_gain_db: 0.0,
}
}
}
impl EqSettings {
pub fn flat() -> Self {
Self::default()
}
/// Clamp all public gains to the supported UI/DSP range.
pub fn clamped(self) -> Self {
Self {
low_gain_db: self.low_gain_db.clamp(EQ_GAIN_DB_MIN, EQ_GAIN_DB_MAX),
mid_gain_db: self.mid_gain_db.clamp(EQ_GAIN_DB_MIN, EQ_GAIN_DB_MAX),
high_gain_db: self.high_gain_db.clamp(EQ_GAIN_DB_MIN, EQ_GAIN_DB_MAX),
}
}
/// True when the EQ should be bypassed entirely.
pub fn is_flat(self) -> bool {
self.low_gain_db.abs() <= FLAT_EPSILON_DB
&& self.mid_gain_db.abs() <= FLAT_EPSILON_DB
&& self.high_gain_db.abs() <= FLAT_EPSILON_DB
}
}
/// A stateful three-band EQ. One instance belongs to one decoded peer stream.
pub struct Eq {
settings: EqSettings,
low: Biquad,
mid: Biquad,
high: Biquad,
}
impl Eq {
/// Build an EQ at the application's audio rate (48 kHz).
pub fn new(settings: EqSettings) -> Self {
Self::with_sample_rate(settings, DEFAULT_SAMPLE_RATE)
}
fn with_sample_rate(settings: EqSettings, sample_rate: f32) -> Self {
let settings = settings.clamped();
Self {
settings,
low: Biquad::low_shelf(sample_rate, LOW_SHELF_HZ, settings.low_gain_db, SHELF_Q),
mid: Biquad::peaking(sample_rate, MID_PEAK_HZ, settings.mid_gain_db, MID_Q),
high: Biquad::high_shelf(sample_rate, HIGH_SHELF_HZ, settings.high_gain_db, SHELF_Q),
}
}
pub fn settings(&self) -> EqSettings {
self.settings
}
/// Process one mono PCM frame in place. Flat settings are sample-exact bypass.
pub fn process_frame(&mut self, frame: &mut [i16]) {
if self.settings.is_flat() {
return;
}
for sample in frame {
let x = *sample as f32;
let y = self.high.process(self.mid.process(self.low.process(x)));
*sample = y.round().clamp(i16::MIN as f32, i16::MAX as f32) as i16;
}
}
}
#[derive(Debug, Clone, Copy)]
struct Coeffs {
b0: f32,
b1: f32,
b2: f32,
a1: f32,
a2: f32,
}
impl Coeffs {
fn normalized(b0: f32, b1: f32, b2: f32, a0: f32, a1: f32, a2: f32) -> Self {
let inv_a0 = 1.0 / a0;
Self {
b0: b0 * inv_a0,
b1: b1 * inv_a0,
b2: b2 * inv_a0,
a1: a1 * inv_a0,
a2: a2 * inv_a0,
}
}
fn all_finite(self) -> bool {
self.b0.is_finite()
&& self.b1.is_finite()
&& self.b2.is_finite()
&& self.a1.is_finite()
&& self.a2.is_finite()
}
}
/// Direct Form II transposed biquad. The two delay registers are the state that
/// must survive across frames.
struct Biquad {
coeffs: Coeffs,
z1: f32,
z2: f32,
}
impl Biquad {
fn new(coeffs: Coeffs) -> Self {
debug_assert!(coeffs.all_finite());
Self {
coeffs,
z1: 0.0,
z2: 0.0,
}
}
fn low_shelf(sample_rate: f32, freq: f32, gain_db: f32, q: f32) -> Self {
let (a, cos_w0, alpha) = rbj_terms(sample_rate, freq, gain_db, q);
let sqrt_a = a.sqrt();
let b0 = a * ((a + 1.0) - (a - 1.0) * cos_w0 + 2.0 * sqrt_a * alpha);
let b1 = 2.0 * a * ((a - 1.0) - (a + 1.0) * cos_w0);
let b2 = a * ((a + 1.0) - (a - 1.0) * cos_w0 - 2.0 * sqrt_a * alpha);
let a0 = (a + 1.0) + (a - 1.0) * cos_w0 + 2.0 * sqrt_a * alpha;
let a1 = -2.0 * ((a - 1.0) + (a + 1.0) * cos_w0);
let a2 = (a + 1.0) + (a - 1.0) * cos_w0 - 2.0 * sqrt_a * alpha;
Self::new(Coeffs::normalized(b0, b1, b2, a0, a1, a2))
}
fn peaking(sample_rate: f32, freq: f32, gain_db: f32, q: f32) -> Self {
let (a, cos_w0, alpha) = rbj_terms(sample_rate, freq, gain_db, q);
let b0 = 1.0 + alpha * a;
let b1 = -2.0 * cos_w0;
let b2 = 1.0 - alpha * a;
let a0 = 1.0 + alpha / a;
let a1 = -2.0 * cos_w0;
let a2 = 1.0 - alpha / a;
Self::new(Coeffs::normalized(b0, b1, b2, a0, a1, a2))
}
fn high_shelf(sample_rate: f32, freq: f32, gain_db: f32, q: f32) -> Self {
let (a, cos_w0, alpha) = rbj_terms(sample_rate, freq, gain_db, q);
let sqrt_a = a.sqrt();
let b0 = a * ((a + 1.0) + (a - 1.0) * cos_w0 + 2.0 * sqrt_a * alpha);
let b1 = -2.0 * a * ((a - 1.0) + (a + 1.0) * cos_w0);
let b2 = a * ((a + 1.0) + (a - 1.0) * cos_w0 - 2.0 * sqrt_a * alpha);
let a0 = (a + 1.0) - (a - 1.0) * cos_w0 + 2.0 * sqrt_a * alpha;
let a1 = 2.0 * ((a - 1.0) - (a + 1.0) * cos_w0);
let a2 = (a + 1.0) - (a - 1.0) * cos_w0 - 2.0 * sqrt_a * alpha;
Self::new(Coeffs::normalized(b0, b1, b2, a0, a1, a2))
}
fn process(&mut self, x: f32) -> f32 {
let y = self.coeffs.b0 * x + self.z1;
self.z1 = self.coeffs.b1 * x - self.coeffs.a1 * y + self.z2;
self.z2 = self.coeffs.b2 * x - self.coeffs.a2 * y;
// Avoid carrying denormal-sized state forever on long quiet tails.
if self.z1.abs() < 1.0e-20 {
self.z1 = 0.0;
}
if self.z2.abs() < 1.0e-20 {
self.z2 = 0.0;
}
y
}
}
fn rbj_terms(sample_rate: f32, freq: f32, gain_db: f32, q: f32) -> (f32, f32, f32) {
let sr = sample_rate.max(1.0);
let f = freq.clamp(1.0, sr * 0.49);
let w0 = 2.0 * std::f32::consts::PI * f / sr;
let a = 10.0f32.powf(gain_db / 40.0);
let alpha = w0.sin() / (2.0 * q.max(0.001));
(a, w0.cos(), alpha)
}
#[cfg(test)]
mod tests {
use super::*;
fn sine(freq: f32, len: usize, amp: f32) -> Vec<i16> {
(0..len)
.map(|n| {
let t = n as f32 / DEFAULT_SAMPLE_RATE;
(amp * (2.0 * std::f32::consts::PI * freq * t).sin()).round() as i16
})
.collect()
}
fn rms(frame: &[i16]) -> f32 {
let sum: f32 = frame.iter().map(|&s| (s as f32).powi(2)).sum();
(sum / frame.len().max(1) as f32).sqrt()
}
#[test]
fn flat_eq_is_sample_exact_identity() {
let mut eq = Eq::new(EqSettings::flat());
let mut frame: Vec<i16> = (-480..480).map(|n| (n * 31) as i16).collect();
let original = frame.clone();
eq.process_frame(&mut frame);
assert_eq!(frame, original);
}
#[test]
fn low_shelf_boost_raises_low_frequency_energy() {
let mut eq = Eq::new(EqSettings {
low_gain_db: 9.0,
..EqSettings::flat()
});
let mut low = sine(100.0, 48_000, 3_000.0);
let before = rms(&low);
eq.process_frame(&mut low);
let after = rms(&low);
assert!(after > before * 1.6, "low shelf should boost low RMS: {before} -> {after}");
}
#[test]
fn high_shelf_boost_raises_high_frequency_energy() {
let mut eq = Eq::new(EqSettings {
high_gain_db: 9.0,
..EqSettings::flat()
});
let mut high = sine(8_000.0, 48_000, 3_000.0);
let before = rms(&high);
eq.process_frame(&mut high);
let after = rms(&high);
assert!(after > before * 1.6, "high shelf should boost high RMS: {before} -> {after}");
}
#[test]
fn coefficients_are_finite_across_supported_gain_range() {
for gain in [EQ_GAIN_DB_MIN, -6.0, 0.0, 6.0, EQ_GAIN_DB_MAX] {
for b in [
Biquad::low_shelf(DEFAULT_SAMPLE_RATE, LOW_SHELF_HZ, gain, SHELF_Q),
Biquad::peaking(DEFAULT_SAMPLE_RATE, MID_PEAK_HZ, gain, MID_Q),
Biquad::high_shelf(DEFAULT_SAMPLE_RATE, HIGH_SHELF_HZ, gain, SHELF_Q),
] {
assert!(b.coeffs.all_finite(), "coefficients must be finite at {gain} dB");
}
}
}
#[test]
fn hot_signal_does_not_nan_or_wrap() {
let mut eq = Eq::new(EqSettings {
low_gain_db: 12.0,
mid_gain_db: 12.0,
high_gain_db: 12.0,
});
let mut frame = sine(1_000.0, 48_000, 30_000.0);
eq.process_frame(&mut frame);
let peak = frame
.iter()
.map(|&s| i32::from(s).abs())
.max()
.unwrap_or(0);
assert!(peak > 1_000, "processed signal should retain audible energy");
assert!(
frame.iter().any(|&s| s > 0) && frame.iter().any(|&s| s < 0),
"a boosted sine should retain both polarities"
);
}
#[test]
fn settings_are_clamped() {
let s = EqSettings {
low_gain_db: -99.0,
mid_gain_db: 2.0,
high_gain_db: 99.0,
}
.clamped();
assert_eq!(s.low_gain_db, EQ_GAIN_DB_MIN);
assert_eq!(s.mid_gain_db, 2.0);
assert_eq!(s.high_gain_db, EQ_GAIN_DB_MAX);
}
}
+11 -4
View File
@@ -1,17 +1,22 @@
use std::sync::mpsc::{Sender, Receiver};
use std::sync::mpsc::{Receiver, Sender};
use std::sync::Arc;
use std::sync::atomic::AtomicUsize;
use thiserror::Error;
/// Target depth of the playback ring buffer, in samples (48kHz mono).
/// Playback output channel count. Capture/encode/network remain mono; only the
/// listener-side playout bus is stereo.
pub const PLAYBACK_CHANNELS: usize = 2;
/// Target depth of the playback ring buffer, in interleaved samples (48kHz
/// stereo).
///
/// The playout chain is paced to keep the ring near this level: production is
/// driven by how fast PipeWire actually drains the ring (the hardware clock),
/// not by a fixed software timer — which is what eliminates the producer/
/// consumer beat that otherwise churns ~20% of audio into drops + silence.
/// 2880 = 60ms = 3×20ms frames, comfortably above the 2048-sample max quantum
/// 5760 = 60ms = 3×20ms stereo frames, comfortably above the 2048-frame max quantum
/// so a single hardware pull can never empty the ring before the mixer refills.
pub const PLAYBACK_TARGET_SAMPLES: usize = 2880;
pub const PLAYBACK_TARGET_SAMPLES: usize = 2880 * PLAYBACK_CHANNELS;
#[derive(Error, Debug)]
pub enum AudioError {
@@ -52,9 +57,11 @@ pub trait AudioBackend: Send + Sync {
}
pub mod echo_cancel;
pub mod eq;
pub mod gate;
pub mod limiter;
pub mod multitrack;
pub mod pan;
pub mod pipewire_impl;
pub mod pw_cli;
pub mod recorder;
+77
View File
@@ -0,0 +1,77 @@
//! Listener-side stereo pan law.
//!
//! Capture, Opus, and the network stay mono. These helpers are used only after a
//! peer has been decoded locally, just before the playout mix is written to the
//! stereo playback bus.
/// Clamp and compute constant-power pan gains for `pan` in `[-1.0, 1.0]`.
///
/// - `-1.0` is hard left `(1, 0)`
/// - `0.0` is center `(sqrt(1/2), sqrt(1/2))`
/// - `1.0` is hard right `(0, 1)`
pub fn pan_gains(pan: f32) -> (f32, f32) {
let pan = pan.clamp(-1.0, 1.0);
let theta = (pan + 1.0) * std::f32::consts::FRAC_PI_4;
(theta.cos(), theta.sin())
}
/// Gains used by the legacy-compatible playback mixer.
///
/// The pure law above is constant-power. The existing application, however, was
/// mono and users heard the full old mono signal in both ears. Scaling by sqrt(2)
/// makes `pan = 0` exactly dual-mono `(1, 1)`, preserving the default sound while
/// still following the same equal-power curve as a peer is moved away from center.
pub fn playback_pan_gains(pan: f32) -> (f32, f32) {
let (left, right) = pan_gains(pan);
(left * std::f32::consts::SQRT_2, right * std::f32::consts::SQRT_2)
}
#[cfg(test)]
mod tests {
use super::*;
const EPS: f32 = 1.0e-6;
#[test]
fn hard_left_and_right_are_endpoints() {
assert_eq!(pan_gains(-1.0), (1.0, 0.0));
let (l, r) = pan_gains(1.0);
assert!(l.abs() < EPS, "left at hard-right should be zero-ish, got {l}");
assert!((r - 1.0).abs() < EPS, "right at hard-right should be one, got {r}");
}
#[test]
fn center_is_equal_and_power_preserving() {
let (l, r) = pan_gains(0.0);
assert!((l - r).abs() < EPS);
assert!((l - std::f32::consts::FRAC_1_SQRT_2).abs() < EPS);
assert!(((l * l + r * r) - 1.0).abs() < EPS);
}
#[test]
fn gains_move_monotonically() {
let pans = [-1.0, -0.5, 0.0, 0.5, 1.0];
let mut prev_l = f32::INFINITY;
let mut prev_r = f32::NEG_INFINITY;
for pan in pans {
let (l, r) = pan_gains(pan);
assert!(l <= prev_l + EPS, "left gain must not rise as pan moves right");
assert!(r >= prev_r - EPS, "right gain must not fall as pan moves right");
prev_l = l;
prev_r = r;
}
}
#[test]
fn playback_center_preserves_legacy_dual_mono() {
let (l, r) = playback_pan_gains(0.0);
assert!((l - 1.0).abs() < EPS);
assert!((r - 1.0).abs() < EPS);
}
#[test]
fn input_is_clamped() {
assert_eq!(pan_gains(-9.0), pan_gains(-1.0));
assert_eq!(pan_gains(9.0), pan_gains(1.0));
}
}
+23 -18
View File
@@ -283,8 +283,9 @@ fn run_playback(
let core = context.connect_rc(None)
.map_err(|e| AudioError::Init(e.to_string()))?;
// Ring buffer setup: 9600 samples (200ms capacity for mono 48kHz).
const RING_CAPACITY: usize = 9600;
// Ring buffer setup: 19200 interleaved samples (200ms capacity for stereo
// 48kHz).
const RING_CAPACITY: usize = 9600 * crate::audio::PLAYBACK_CHANNELS;
let rb = HeapRb::<i16>::new(RING_CAPACITY);
let (mut producer, consumer) = rb.split();
@@ -371,7 +372,7 @@ fn run_playback(
let data = &mut datas[0];
let mut total_size = 0;
if let Some(slice) = data.data() {
let stride = 2; // S16LE Mono = 2 bytes per frame
let stride = 2 * crate::audio::PLAYBACK_CHANNELS; // S16LE stereo
// Fill exactly what the graph asked for this cycle (with
// a safe fallback), never the whole mapped slice — that
// over-pull past the ring depth was the original crackle.
@@ -383,17 +384,20 @@ fn run_playback(
user_data.callback_count.fetch_add(1, Ordering::Relaxed);
let mut starved = 0u64;
for i in 0..n_frames {
let val = match user_data.consumer.try_pop() {
Some(v) => v,
None => {
starved += 1;
0
}
};
let bytes = val.to_le_bytes();
let start = i * stride;
slice[start] = bytes[0];
slice[start + 1] = bytes[1];
for ch in 0..crate::audio::PLAYBACK_CHANNELS {
let val = match user_data.consumer.try_pop() {
Some(v) => v,
None => {
starved += 1;
0
}
};
let bytes = val.to_le_bytes();
let offset = start + ch * 2;
slice[offset] = bytes[0];
slice[offset + 1] = bytes[1];
}
}
if starved > 0 {
// One wait-free atomic add per quantum — RT-safe.
@@ -403,7 +407,8 @@ fn run_playback(
// actually pulled (excluding underruns, which removed
// nothing) so the mixer paces against true ring depth.
// Wait-free fetch_sub, RT-safe.
let popped = n_frames - starved as usize;
let requested_samples = n_frames * crate::audio::PLAYBACK_CHANNELS;
let popped = requested_samples - starved as usize;
if popped > 0 {
user_data.fill_gauge.fetch_sub(popped, Ordering::Relaxed);
}
@@ -411,7 +416,7 @@ fn run_playback(
}
let chunk = data.chunk_mut();
*chunk.offset_mut() = 0;
*chunk.stride_mut() = 2;
*chunk.stride_mut() = (2 * crate::audio::PLAYBACK_CHANNELS) as _;
*chunk.size_mut() = total_size as _;
}
}
@@ -422,7 +427,7 @@ fn run_playback(
let mut audio_info = spa::param::audio::AudioInfoRaw::new();
audio_info.set_format(spa::param::audio::AudioFormat::S16LE);
audio_info.set_rate(48000);
audio_info.set_channels(1); // Mono
audio_info.set_channels(crate::audio::PLAYBACK_CHANNELS as u32); // Stereo playback
let obj = pw::spa::pod::Object {
type_: pw::spa::utils::SpaTypes::ObjectParamFormat.as_raw(),
@@ -450,7 +455,7 @@ fn run_playback(
// `frames_to_produce`). `requested()`, not the buffer size, now governs
// per-cycle output, so this is a generous max rather than a hard pin.
const MAX_QUANTUM_FRAMES: i32 = 8192;
const STRIDE: i32 = 2; // S16LE mono = 2 bytes/frame
const STRIDE: i32 = 2 * crate::audio::PLAYBACK_CHANNELS as i32; // S16LE stereo
let buffers_obj = pw::spa::pod::Object {
type_: pw::spa::utils::SpaTypes::ObjectParamBuffers.as_raw(),
id: pw::spa::param::ParamType::Buffers.as_raw(),
@@ -555,7 +560,7 @@ fn run_playback(
if verbose || du > 0 || dd > 0 {
crate::log_msg(&format!(
"playout-health: fill={fill} samples (~{}ms) | underrun +{du} samples/s (total {u}) | dropped +{dd} frames/s (total {d}) | quantum={q} frames, {dc} callbacks/s",
fill / 48,
fill / (48 * crate::audio::PLAYBACK_CHANNELS),
));
}
}
+55 -8
View File
@@ -22,6 +22,8 @@ use std::path::{Path, PathBuf};
const SAMPLE_RATE: u32 = 48_000;
const BITS_PER_SAMPLE: u16 = 16;
const CHANNELS: u16 = 1;
const RIFF_DATA_OVERHEAD: u64 = 36;
const MAX_RIFF_DATA_BYTES: u64 = u32::MAX as u64 - RIFF_DATA_OVERHEAD;
/// Cap on buffered mic samples (~200ms). Bounds how far recording lag can drift
/// if the capture clock runs persistently faster than playout — past this we drop
@@ -34,7 +36,7 @@ const MAX_MIC_FIFO: usize = SAMPLE_RATE as usize / 5;
pub struct WavWriter {
file: File,
/// Bytes of PCM data written so far (for the size fields).
data_bytes: u32,
data_bytes: u64,
}
impl WavWriter {
@@ -42,7 +44,10 @@ impl WavWriter {
pub fn new(path: &Path) -> io::Result<Self> {
let mut file = File::create(path)?;
file.write_all(&Self::header(0))?;
Ok(Self { file, data_bytes: 0 })
Ok(Self {
file,
data_bytes: 0,
})
}
/// The 44-byte canonical WAV/PCM header for the given data length in bytes.
@@ -68,21 +73,40 @@ impl WavWriter {
/// Append PCM samples to the data chunk.
pub fn write_samples(&mut self, samples: &[i16]) -> io::Result<()> {
let added_bytes = u64::try_from(samples.len())
.ok()
.and_then(|len| len.checked_mul(2))
.ok_or_else(|| io::Error::other("WAV sample buffer too large"))?;
let new_data_bytes = self
.data_bytes
.checked_add(added_bytes)
.ok_or_else(|| io::Error::other("WAV data size overflow"))?;
if new_data_bytes > MAX_RIFF_DATA_BYTES {
return Err(io::Error::other("WAV too large for RIFF"));
}
let mut buf = Vec::with_capacity(samples.len() * 2);
for &s in samples {
buf.extend_from_slice(&s.to_le_bytes());
}
self.file.write_all(&buf)?;
self.data_bytes += (samples.len() * 2) as u32;
self.data_bytes = new_data_bytes;
Ok(())
}
/// Patch the RIFF + data size fields and flush. Consumes the writer.
pub fn finalize(mut self) -> io::Result<()> {
let data_bytes = u32::try_from(self.data_bytes)
.map_err(|_| io::Error::other("WAV too large for RIFF"))?;
let riff_size = self
.data_bytes
.checked_add(RIFF_DATA_OVERHEAD)
.and_then(|size| u32::try_from(size).ok())
.ok_or_else(|| io::Error::other("WAV too large for RIFF"))?;
self.file.seek(SeekFrom::Start(4))?;
self.file.write_all(&(36 + self.data_bytes).to_le_bytes())?;
self.file.write_all(&riff_size.to_le_bytes())?;
self.file.seek(SeekFrom::Start(40))?;
self.file.write_all(&self.data_bytes.to_le_bytes())?;
self.file.write_all(&data_bytes.to_le_bytes())?;
self.file.flush()?;
Ok(())
}
@@ -209,11 +233,29 @@ mod tests {
let _ = std::fs::remove_file(&path);
}
#[test]
fn wav_writer_rejects_data_that_would_overflow_riff_header() {
let dir = std::env::temp_dir();
let path = dir.join(format!("peerspeak-overflow-{}.wav", std::process::id()));
let mut w = WavWriter::new(&path).unwrap();
w.data_bytes = MAX_RIFF_DATA_BYTES - 1;
let before_len = std::fs::metadata(&path).unwrap().len();
let err = w.write_samples(&[0]).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::Other);
assert_eq!(w.data_bytes, MAX_RIFF_DATA_BYTES - 1);
assert_eq!(std::fs::metadata(&path).unwrap().len(), before_len);
drop(w);
let _ = std::fs::remove_file(&path);
}
#[test]
fn mic_is_summed_with_mix_when_present() {
let dir = std::env::temp_dir();
let mut r = Recorder {
writer: WavWriter::new(&dir.join(format!("ps-sum-{}.wav", std::process::id()))).unwrap(),
writer: WavWriter::new(&dir.join(format!("ps-sum-{}.wav", std::process::id())))
.unwrap(),
mic_fifo: VecDeque::new(),
path: PathBuf::new(),
};
@@ -223,7 +265,11 @@ mod tests {
r.write_frame(&[10, 20]).unwrap();
assert_eq!(r.mic_fifo.len(), 1, "two samples consumed, one mic left");
r.write_frame(&[0, 0]).unwrap();
assert_eq!(r.mic_fifo.len(), 0, "remaining mic sample consumed; rest is silence");
assert_eq!(
r.mic_fifo.len(),
0,
"remaining mic sample consumed; rest is silence"
);
let _ = r.finalize();
}
@@ -231,7 +277,8 @@ mod tests {
fn mic_fifo_is_capped() {
let dir = std::env::temp_dir();
let mut r = Recorder {
writer: WavWriter::new(&dir.join(format!("ps-cap-{}.wav", std::process::id()))).unwrap(),
writer: WavWriter::new(&dir.join(format!("ps-cap-{}.wav", std::process::id())))
.unwrap(),
mic_fifo: VecDeque::new(),
path: PathBuf::new(),
};
+4 -2
View File
@@ -26,7 +26,7 @@ use std::time::Duration;
use peerspeak::audio::AudioBackend;
use peerspeak::audio::pipewire_impl::PipeWireBackend;
use peerspeak::core::jitter::FRAME_SAMPLES; // 960 samples = 20ms @ 48kHz mono
use peerspeak::core::jitter::FRAME_SAMPLES; // 960 mono frames = 20ms @ 48kHz
const SAMPLE_RATE: f32 = 48_000.0;
@@ -69,11 +69,13 @@ async fn main() {
tokio::time::sleep(Duration::from_millis(2)).await;
continue;
}
let mut frame = Vec::with_capacity(FRAME_SAMPLES);
let mut frame = Vec::with_capacity(FRAME_SAMPLES * peerspeak::audio::PLAYBACK_CHANNELS);
for _ in 0..FRAME_SAMPLES {
let t = n as f32 / SAMPLE_RATE;
// 0.25 amplitude: clearly audible but not harsh.
let sample = (0.25 * i16::MAX as f32 * (2.0 * std::f32::consts::PI * freq * t).sin()) as i16;
// Stereo playback bus: duplicate the probe tone to L/R.
frame.push(sample);
frame.push(sample);
n += 1;
}
+27 -1
View File
@@ -1,6 +1,7 @@
use crate::notify::Sound;
use crate::theme::AppTheme;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::fs;
use std::path::PathBuf;
@@ -232,6 +233,17 @@ pub struct AppConfig {
/// capped (see `recents`). Defaulted empty so older configs upgrade cleanly.
#[serde(default)]
pub recents: Vec<crate::recents::Recent>,
/// Per-peer listener-side EQ settings, keyed by peer node id string. Local
/// preference only; never sent to peers.
#[serde(default)]
pub peer_eq: HashMap<String, crate::audio::eq::EqSettings>,
/// Per-peer listener-side pan (`-1.0` left, `0.0` center, `1.0` right),
/// keyed by peer node id string. Local preference only.
#[serde(default)]
pub peer_pan: HashMap<String, f32>,
/// Focused app-local keyboard shortcuts.
#[serde(default)]
pub hotkeys: crate::hotkeys::HotkeyMap,
/// Last window size (px), restored as the initial size on next launch.
/// Saved on close.
#[serde(default = "default_window_width")]
@@ -287,6 +299,9 @@ impl Default for AppConfig {
sound_reconnect_failed_enabled: true,
pixelpass_path: None,
recents: Vec::new(),
peer_eq: HashMap::new(),
peer_pan: HashMap::new(),
hotkeys: crate::hotkeys::HotkeyMap::default(),
window_width: default_window_width(),
window_height: default_window_height(),
window_x: None,
@@ -411,6 +426,18 @@ mod tests {
assert_eq!(deserialized.window_height, 760.0);
// Configs predating the recents list load an empty list.
assert!(deserialized.recents.is_empty());
// Configs predating per-peer listener shaping load flat/center/default
// shortcut settings.
assert!(deserialized.peer_eq.is_empty());
assert!(deserialized.peer_pan.is_empty());
assert_eq!(
crate::hotkeys::format_binding(
deserialized
.hotkeys
.binding(crate::hotkeys::HotkeyAction::PushToTalk)
),
"Space"
);
}
#[test]
@@ -596,4 +623,3 @@ mod tests {
assert_eq!(config.noise_gate_threshold, 0.01);
}
}
+62 -3
View File
@@ -59,6 +59,10 @@ const PRIME_TIMEOUT_TICKS: usize = 25;
/// badly behind, so we drop the oldest and resync rather than grow unbounded.
const MAX_BUFFERED_FRAMES: usize = 32;
/// Sequence discontinuities larger than this (~10s at 20ms/frame) are treated
/// as a restarted/new stream, not ordinary packet loss or reordering.
const MAX_REASONABLE_SEQ_GAP: u32 = 500;
pub struct JitterBuffer {
decoder: OpusDecoder,
/// Reorder window: sequence number -> encoded Opus payload.
@@ -116,6 +120,14 @@ impl JitterBuffer {
}
}
fn reset_to_stream(&mut self, seq: u32, payload: Vec<u8>) {
self.packets.clear();
self.packets.insert(seq, payload);
self.next_seq = None;
self.clean_run = 0;
self.buffering_ticks = 0;
}
/// Store a received packet, dropping ones we've already played past and
/// bounding total depth.
pub fn insert(&mut self, seq: u32, payload: Vec<u8>) {
@@ -124,9 +136,19 @@ impl JitterBuffer {
if let Some(next) = self.next_seq
&& seq_before(seq, next)
{
if next.wrapping_sub(seq) > MAX_REASONABLE_SEQ_GAP {
self.reset_to_stream(seq, payload);
return;
}
self.note_disruption();
return;
}
if let Some(next) = self.next_seq
&& seq.wrapping_sub(next) > MAX_REASONABLE_SEQ_GAP
{
self.reset_to_stream(seq, payload);
return;
}
self.packets.insert(seq, payload);
while self.packets.len() > MAX_BUFFERED_FRAMES {
@@ -283,6 +305,44 @@ mod tests {
assert_eq!(jb.packets.len(), 1); // only seq 7 remains buffered
}
#[test]
fn far_behind_sequence_resets_as_restarted_stream() {
let mut jb = JitterBuffer::new().unwrap();
jb.next_seq = Some(5_000);
jb.packets.insert(5_000, vec![9]);
jb.clean_run = 12;
jb.buffering_ticks = 4;
jb.insert(0, vec![1]);
assert_eq!(jb.next_seq, None);
assert_eq!(jb.packets.len(), 1);
assert_eq!(jb.packets.get(&0).map(Vec::as_slice), Some(&[1][..]));
assert_eq!(jb.clean_run, 0);
assert_eq!(jb.buffering_ticks, 0);
}
#[test]
fn far_ahead_sequence_resets_to_bound_plc_run() {
let mut jb = JitterBuffer::new().unwrap();
jb.next_seq = Some(10);
jb.packets.insert(10, vec![9]);
jb.clean_run = 12;
jb.buffering_ticks = 4;
let jumped_seq = 10 + MAX_REASONABLE_SEQ_GAP + 1;
jb.insert(jumped_seq, vec![2]);
assert_eq!(jb.next_seq, None);
assert_eq!(jb.packets.len(), 1);
assert_eq!(
jb.packets.get(&jumped_seq).map(Vec::as_slice),
Some(&[2][..])
);
assert_eq!(jb.clean_run, 0);
assert_eq!(jb.buffering_ticks, 0);
}
#[test]
fn test_seq_before_ordering() {
// Basic ordering
@@ -358,7 +418,7 @@ mod tests {
fn is_idle_reflects_buffer_state() {
let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap();
let mut jb = JitterBuffer::new().unwrap();
// Fresh buffer
assert!(jb.is_idle());
@@ -369,7 +429,7 @@ mod tests {
// Prime (3 frames)
jb.insert(1, frame(&mut enc, 1000));
jb.insert(2, frame(&mut enc, 1000));
// Drain past the end so it underruns
assert!(jb.pop_frame().is_some());
assert!(jb.pop_frame().is_some());
@@ -635,4 +695,3 @@ mod tests {
assert_eq!(jb.clean_run, 0);
}
}
+14 -5
View File
@@ -11,6 +11,10 @@ pub enum CoreCommand {
/// for a NEW room; it's ignored when joining (the label rides in the ticket).
Join { name: String, ticket: String, room_name: String, input_device: Option<String>, output_device: Option<String>, echo_cancellation: bool, avatar: crate::avatar::Avatar },
Leave,
/// Orderly app shutdown: finalize recordings, leave any active room, stop local
/// audio/screen-share work, close the persistent network stack, then ack with
/// [`UiEvent::ShutdownComplete`].
Shutdown,
ToggleMute,
/// Change our avatar (W4) and re-announce it to the room over presence.
SetAvatar(crate::avatar::Avatar),
@@ -18,6 +22,10 @@ pub enum CoreCommand {
SetPttMode(bool),
SetPttActive(bool),
SetPeerVolume(EndpointId, f32),
/// Listener-side per-peer EQ. Local only; never leaves this app instance.
SetPeerEq(EndpointId, crate::audio::eq::EqSettings),
/// Listener-side per-peer pan. Local only; never leaves this app instance.
SetPeerPan(EndpointId, f32),
/// Locally mute/unmute a peer: when muted, their audio is decoded (so levels
/// still show) but not mixed into our output.
SetPeerMuted(EndpointId, bool),
@@ -116,11 +124,12 @@ pub enum UiEvent {
/// joinable gathering (with a one-click ticket). Emitted by the outbound ping
/// scheduler; absence of a recent event = treat as offline.
FriendPresence { id: EndpointId, presence: FriendPresence },
/// The Discoverable time-box elapsed (W7 P6): the core auto-reverted our presence
/// posture to the carried `mode` (always `Normal`) and stopped publishing. The
/// GUI must mirror + persist this so its presence picker stops showing
/// Discoverable. Distinct from a user-driven change so the GUI knows to update
/// without having issued the command itself.
/// Core corrected the committed presence posture. Usually the Discoverable
/// time-box elapsed and the core auto-reverted to `Normal`; on discovery apply
/// failure, this carries the previous truthful mode. The GUI must mirror +
/// persist this so its presence picker matches the endpoint's discovery state.
PresenceModeReverted { mode: PresenceMode },
/// Core finished orderly app shutdown and the GUI can exit.
ShutdownComplete,
Error(String),
}
+347 -51
View File
@@ -2,6 +2,7 @@ pub mod messages;
pub mod jitter;
use crate::audio::{AudioBackend, pipewire_impl::PipeWireBackend};
use crate::audio::eq::{Eq, EqSettings};
use crate::codec::{AudioEncoder, opus_impl::OpusEncoder};
use crate::core::jitter::{JitterBuffer, FRAME_SAMPLES};
use crate::network::{
@@ -12,6 +13,7 @@ use crate::network::{
use crate::core::messages::{CoreCommand, UiEvent};
use crate::config::{NetworkMode, RecordingMode};
use crate::presence::PresenceMode;
use crate::audio::multitrack::MultitrackRecorder;
use iroh::{Endpoint, EndpointAddr, EndpointId, RelayMode, SecretKey, endpoint::presets, protocol::Router};
use iroh_gossip::net::Gossip;
@@ -49,6 +51,12 @@ impl CoreController {
pub fn send(&self, cmd: CoreCommand) -> bool {
self.cmd_tx.try_send(cmd).is_ok()
}
/// Clone the command sender for asynchronous one-shot sends that should wait
/// for channel capacity instead of failing immediately on a full queue.
pub fn command_sender(&self) -> mpsc::Sender<CoreCommand> {
self.cmd_tx.clone()
}
}
/// How long a peer may stay "reconnecting" after a transient drop before we give
@@ -57,6 +65,32 @@ impl CoreController {
/// clears from the room promptly.
const RECONNECT_GRACE: Duration = Duration::from_secs(45);
/// Opus frames sent by our encoder are one 20 ms mono frame, normally far below
/// this. 4000 bytes still leaves room for large valid Opus packets (well above a
/// 48 kHz / 60 ms frame) while bounding malicious datagram copy/decode churn.
const MAX_OPUS_PAYLOAD: usize = 4000;
/// If the Discoverable time-box tries to revert but discovery service reconfiguration
/// fails, retry soon while keeping the UI in the still-possible publishing state.
const DISCOVERY_REVERT_RETRY: Duration = Duration::from_secs(60);
fn audio_datagram_len_ok(len: usize) -> bool {
(4..=4 + MAX_OPUS_PAYLOAD).contains(&len)
}
fn arm_discovery_retry(
discovery_deadline: &mut Option<tokio::time::Instant>,
now: tokio::time::Instant,
) {
let retry_deadline = now + DISCOVERY_REVERT_RETRY;
if discovery_deadline
.map(|current| current > retry_deadline)
.unwrap_or(true)
{
*discovery_deadline = Some(retry_deadline);
}
}
/// Per-peer reconnect grace timers (see [`RECONNECT_GRACE`]). Shared between the
/// room-event task (which arms one on a transient drop and cancels it on a
/// gossip rejoin) and the conn-event task (which cancels it when the audio link
@@ -214,6 +248,7 @@ fn stop_mic_monitor(backend: &PipeWireBackend, monitor: Option<MicMonitor>) {
/// limiter (see [`crate::audio::limiter`]) can ride it down to the ceiling instead
/// of the old hard clip shattering loud moments. Peers shorter than `frame_len`
/// contribute 0 past their end; an empty peer set yields a silent bus.
#[cfg(test)]
fn mix_frames(peer_frames: &[Vec<i16>], frame_len: usize) -> Vec<i32> {
let mut mixed = vec![0i32; frame_len];
for frame in peer_frames {
@@ -224,6 +259,44 @@ fn mix_frames(peer_frames: &[Vec<i16>], frame_len: usize) -> Vec<i32> {
mixed
}
/// Sum per-peer mono frames into one interleaved stereo `i32` bus. Center pan is
/// a special exact dual-mono path so the default listener mix is bit-for-bit the
/// old mono sum duplicated to both ears.
fn mix_stereo_frames(peer_frames: &[(Vec<i16>, f32)], frame_len: usize) -> Vec<i32> {
let mut mixed = vec![0i32; frame_len * crate::audio::PLAYBACK_CHANNELS];
for (frame, pan) in peer_frames {
if pan.abs() <= f32::EPSILON {
for (i, &sample) in frame.iter().take(frame_len).enumerate() {
let idx = i * crate::audio::PLAYBACK_CHANNELS;
let s = sample as i32;
mixed[idx] += s;
mixed[idx + 1] += s;
}
continue;
}
let (left_gain, right_gain) = crate::audio::pan::playback_pan_gains(*pan);
for (i, &sample) in frame.iter().take(frame_len).enumerate() {
let idx = i * crate::audio::PLAYBACK_CHANNELS;
let x = sample as f32;
mixed[idx] += (x * left_gain).round() as i32;
mixed[idx + 1] += (x * right_gain).round() as i32;
}
}
mixed
}
/// Fold an interleaved stereo frame to mono for the existing mixed WAV writers.
/// Center/default pan folds back to the exact old mono mix.
fn stereo_to_mono(stereo: &[i16]) -> Vec<i16> {
let mut mono = Vec::with_capacity(stereo.len() / crate::audio::PLAYBACK_CHANNELS);
for pair in stereo.chunks_exact(crate::audio::PLAYBACK_CHANNELS) {
let sum = pair[0] as i32 + pair[1] as i32;
mono.push((sum / 2).clamp(i16::MIN as i32, i16::MAX as i32) as i16);
}
mono
}
/// Handles the transport's per-peer link-state stream (`ConnEvent`): arms/cancels
/// reconnect grace timers, tracks which peers we've linked with, and forwards
/// link state to the UI. Pulled out of the conn-event task as a unit so the
@@ -411,12 +484,11 @@ impl NetStack {
/// `DnsAddressLookup`, mirroring the `N0` preset) is added when `plan.resolver`; the
/// n0 DNS *publisher* (`PkarrPublisher`) when `plan.publisher`.
///
/// Idempotent and reversible: it clears the whole service set and reinstalls exactly
/// what the plan wants, so flipping `publisher` off simply drops the publisher (its
/// republish task ends when the last clone is dropped, and the already-published
/// record TTL-expires within ~30s) without an endpoint rebuild and without disturbing
/// resolution. The brief clear→re-add window is a few synchronous calls; presence
/// toggles are rare, so a concurrent dial racing it is not a practical concern.
/// Idempotent and reversible: it builds the replacement services first, then clears
/// the service set and reinstalls exactly what the plan wants. Flipping `publisher`
/// off drops the publisher (its republish task ends when the last clone is dropped,
/// and the already-published record TTL-expires within ~30s) without an endpoint
/// rebuild and without disturbing resolution.
fn apply_discovery(
endpoint: &Endpoint,
memory_lookup: &iroh::address_lookup::memory::MemoryLookup,
@@ -427,16 +499,34 @@ fn apply_discovery(
pkarr::{PkarrPublisher, PkarrResolver},
};
let services = endpoint.address_lookup()?;
let pkarr_resolver = if plan.resolver {
Some(PkarrResolver::n0_dns().into_address_lookup(endpoint)?)
} else {
None
};
let dns_resolver = if plan.resolver {
Some(DnsAddressLookup::n0_dns().into_address_lookup(endpoint)?)
} else {
None
};
let publisher = if plan.publisher {
Some(PkarrPublisher::n0_dns().into_address_lookup(endpoint)?)
} else {
None
};
services.clear();
// Always keep the local, server-free lookup (this is what ticket/gossip dialing
// depends on — it must survive every posture, including DirectOnly).
services.add(memory_lookup.clone());
if plan.resolver {
services.add(PkarrResolver::n0_dns().into_address_lookup(endpoint)?);
services.add(DnsAddressLookup::n0_dns().into_address_lookup(endpoint)?);
if let Some(pkarr_resolver) = pkarr_resolver {
services.add(pkarr_resolver);
}
if plan.publisher {
services.add(PkarrPublisher::n0_dns().into_address_lookup(endpoint)?);
if let Some(dns_resolver) = dns_resolver {
services.add(dns_resolver);
}
if let Some(publisher) = publisher {
services.add(publisher);
}
Ok(())
}
@@ -588,7 +678,7 @@ async fn probe_friends_once(
let ep = endpoint.clone();
set.spawn(async move {
match crate::presence_net::probe(&ep, addr).await {
Ok(reply) => crate::presence::interpret_pong(&reply).map(|p| (id, p)),
Ok((from, reply)) => crate::presence::interpret_pong(&reply, from).map(|p| (id, p)),
Err(_) => None,
}
});
@@ -665,6 +755,8 @@ async fn run_core_loop(
let is_multitrack = Arc::new(AtomicBool::new(false));
let mut recording_mode = RecordingMode::default();
let peer_volumes = Arc::new(Mutex::new(HashMap::<EndpointId, f32>::new()));
let peer_eq = Arc::new(Mutex::new(HashMap::<EndpointId, EqSettings>::new()));
let peer_pan = Arc::new(Mutex::new(HashMap::<EndpointId, f32>::new()));
// Peers locally muted by us: decoded for level metering but not mixed.
let locally_muted = Arc::new(Mutex::new(HashSet::<EndpointId>::new()));
let mut current_name = "Anonymous".to_string();
@@ -794,26 +886,85 @@ async fn run_core_loop(
// W7 P6 time-box: Discoverable auto-reverts to Normal after DISCOVERY_TIMEBOX
// so a publish beacon never stands indefinitely. The branch is disabled
// (`if` guard) unless a deadline is armed; `unwrap_or_else` is unreachable
// belt-and-braces. On fire: stop publishing, drop to Normal, tell the GUI.
// belt-and-braces. On fire: stop publishing first, then commit Normal only
// if the endpoint's discovery services accepted the non-publishing plan.
_ = tokio::time::sleep_until(
discovery_deadline.unwrap_or_else(tokio::time::Instant::now),
), if discovery_deadline.is_some() => {
discovery_deadline = None;
*presence_mode.lock().unwrap() = crate::presence::PresenceMode::Normal;
let plan = crate::discovery::lookup_plan(network_mode, false);
if let Err(e) = apply_discovery(&net.endpoint, &net.memory_lookup, plan) {
crate::log_msg(&format!("discovery: time-box revert failed: {e:#}"));
let previous_mode = *presence_mode.lock().unwrap();
if previous_mode != PresenceMode::Discoverable {
discovery_deadline = None;
continue;
}
let requested_mode = PresenceMode::Normal;
let now = tokio::time::Instant::now();
let plan = crate::discovery::lookup_plan(
network_mode,
requested_mode.publishes_to_discovery(),
);
let apply_result = apply_discovery(&net.endpoint, &net.memory_lookup, plan);
let (committed_mode, transition_error) =
crate::discovery::resolve_presence_transition(
previous_mode,
requested_mode,
apply_result.is_ok(),
);
*presence_mode.lock().unwrap() = committed_mode;
discovery_deadline = if committed_mode == PresenceMode::Discoverable {
Some(now + DISCOVERY_REVERT_RETRY)
} else {
None
};
match apply_result {
Ok(()) => {
crate::log_msg(
"discovery: Discoverable time-box elapsed → reverting to Normal",
);
let _ = ui_tx
.send(UiEvent::PresenceModeReverted {
mode: PresenceMode::Normal,
})
.await;
}
Err(e) => {
crate::log_msg(&format!("discovery: time-box revert failed: {e:#}"));
if committed_mode != requested_mode {
let _ = ui_tx
.send(UiEvent::PresenceModeReverted {
mode: committed_mode,
})
.await;
}
if let Some(message) = transition_error {
let _ = ui_tx
.send(UiEvent::Error(format!("{message} ({e:#})")))
.await;
}
}
}
crate::log_msg("discovery: Discoverable time-box elapsed → reverting to Normal");
let _ = ui_tx
.send(UiEvent::PresenceModeReverted {
mode: crate::presence::PresenceMode::Normal,
})
.await;
continue;
}
};
match cmd {
CoreCommand::Shutdown => {
crate::log_msg("Core shutdown requested");
// Finalize recordings while capture/mixer feeders are still alive.
stop_recording(&recorder, &is_recording, &multitrack, &is_multitrack, &ui_tx).await;
stop_mic_monitor(&audio_backend, mic_monitor.take());
if let Some(session) = active_session.take() {
session.shutdown(audio_backend.clone()).await;
net.audio_router.clear();
}
*current_room.lock().unwrap() = None;
net.shutdown().await;
let _ = ui_tx.send(UiEvent::ShutdownComplete).await;
break;
}
CoreCommand::Join { name, ticket, room_name, input_device, output_device, echo_cancellation, avatar } => {
current_name = name.clone();
current_avatar = avatar;
@@ -858,7 +1009,12 @@ async fn run_core_loop(
let ticket_str = if ticket.trim().is_empty() || ticket == "create" {
let topic_id: [u8; 32] = rand::random();
let host_addr = endpoint.addr();
crate::log_msg(&format!("Creating room. host_addr={:?}, topic_id={:?}", host_addr, topic_id));
crate::log_msg(&format!(
"Creating room. host_id={}, host_addrs={}, topic={}",
crate::short_id(&host_addr.id.to_string()),
host_addr.addrs.len(),
crate::short_bytes_hex(&topic_id)
));
// The creator's chosen cosmetic label rides in the ticket so
// every joiner inherits it; sanitize it before it leaves here.
let label = crate::sanitize::sanitize_name(&room_name);
@@ -866,7 +1022,10 @@ async fn run_core_loop(
ticket.to_string()
} else {
let ticket_str = ticket.trim().to_string();
crate::log_msg(&format!("Joining room with existing ticket={}", ticket_str));
crate::log_msg(&format!(
"Joining room with existing ticket={}",
crate::redact_for_log(&ticket_str)
));
ticket_str
};
@@ -905,7 +1064,17 @@ async fn run_core_loop(
.map(|peers| peers.values().cloned().collect())
.unwrap_or_default();
crate::log_msg(&format!("Attempting room_state.join with self_state={:?}, extra_bootstrap={:?}", self_state, extra_bootstrap.iter().map(|a| a.id).collect::<Vec<_>>()));
let extra_bootstrap_ids = extra_bootstrap
.iter()
.map(|a| crate::short_id(&a.id.to_string()))
.collect::<Vec<_>>();
crate::log_msg(&format!(
"Attempting room_state.join self_id={}, self_name={:?}, sharing={}, extra_bootstrap={:?}",
crate::short_id(&self_state.addr.id.to_string()),
self_state.name,
self_state.sharing.is_some(),
extra_bootstrap_ids
));
if let Err(e) = room_state.join(&ticket_str, self_state.clone(), extra_bootstrap).await {
crate::log_msg(&format!("Error room_state.join failed: {:?}", e));
let _ = ui_tx.send(UiEvent::Error(format!("Failed to join room: {}", e))).await;
@@ -1070,8 +1239,9 @@ async fn run_core_loop(
};
while let Some((from_peer, bytes)) = datagram_rx.recv().await {
if bytes.len() < 4 {
continue; // malformed: missing sequence header
if !audio_datagram_len_ok(bytes.len()) {
// Malformed (< sequence header) or oversized Opus payload.
continue;
}
let seq = u32::from_le_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
let payload = bytes[4..].to_vec();
@@ -1104,6 +1274,8 @@ async fn run_core_loop(
let jitter_mixer = jitter.clone();
let is_deafened_clone = is_deafened.clone();
let peer_volumes_mixer = peer_volumes.clone();
let peer_eq_mixer = peer_eq.clone();
let peer_pan_mixer = peer_pan.clone();
let locally_muted_mixer = locally_muted.clone();
let output_gain_mixer = output_gain.clone();
let ui_tx_mixer = ui_tx.clone();
@@ -1117,6 +1289,9 @@ async fn run_core_loop(
// the ceiling instead of hard-clipping. State carries across
// frames (see audio::limiter).
let mut limiter = crate::audio::limiter::SoftLimiter::new(48_000);
// Per-peer EQ filter state. Settings are live-cloned each
// cycle; state is rebuilt only when a peer's EQ changes.
let mut peer_eqs: HashMap<EndpointId, Eq> = HashMap::new();
// When the ring is at/above target we have nothing to do; nap
// briefly and re-check. Short enough (relative to the ~60ms
// target and ~21ms device quantum) that we always refill well
@@ -1140,8 +1315,11 @@ async fn run_core_loop(
}
let current_volumes = peer_volumes_mixer.lock().await.clone();
let current_eq = peer_eq_mixer.lock().await.clone();
let current_pans = peer_pan_mixer.lock().await.clone();
let muted_peers = locally_muted_mixer.lock().await.clone();
let mut peer_frames = Vec::new();
let mut peer_frames: Vec<(Vec<i16>, f32)> = Vec::new();
let mut peers_seen = HashSet::new();
// Multitrack stem capture: tap each peer's RAW decoded frame
// (pre-volume, pre-mute, pre-limiter) so the stems are clean
@@ -1167,10 +1345,31 @@ async fn run_core_loop(
let vol = current_volumes.get(&peer_id).copied().unwrap_or(1.0);
apply_volume(&mut frame, vol);
let eq_settings = current_eq
.get(&peer_id)
.copied()
.unwrap_or_default()
.clamped();
if eq_settings.is_flat() {
peer_eqs.remove(&peer_id);
} else {
let needs_rebuild = peer_eqs
.get(&peer_id)
.map(|eq| eq.settings() != eq_settings)
.unwrap_or(true);
if needs_rebuild {
peer_eqs.insert(peer_id, Eq::new(eq_settings));
}
if let Some(eq) = peer_eqs.get_mut(&peer_id) {
eq.process_frame(&mut frame);
}
}
// Level is recorded even for locally-muted peers so
// the UI still shows that they're speaking.
let peak = level_peaks.entry(peer_id).or_insert(0.0);
*peak = peak.max(frame_level(&frame));
peers_seen.insert(peer_id);
// Locally muted: decoded above (jitter buffer advances,
// level shown) but not mixed into our output.
@@ -1178,16 +1377,23 @@ async fn run_core_loop(
continue;
}
peer_frames.push(frame);
let pan = current_pans
.get(&peer_id)
.copied()
.unwrap_or(0.0)
.clamp(-1.0, 1.0);
peer_frames.push((frame, pan));
}
}
peer_eqs.retain(|id, _| peers_seen.contains(id) || current_eq.contains_key(id));
// Lossless i32 sum, then the limiter applies the master
// output gain (in f32, so a boost past the ceiling is
// limited too) and rides peaks down to the ceiling.
let mixed_sum = mix_frames(&peer_frames, FRAME_SAMPLES);
let mixed_sum = mix_stereo_frames(&peer_frames, FRAME_SAMPLES);
let out_gain = f32::from_bits(output_gain_mixer.load(Ordering::Relaxed));
let mixed = limiter.process(&mixed_sum, out_gain);
let record_mix = stereo_to_mono(&mixed);
// Record the true call audio, independent of local deafen —
// deafen only silences our own monitor, not what the call
@@ -1200,7 +1406,7 @@ async fn run_core_loop(
for (id, f) in &stems {
mt.write_peer(*id, f)?;
}
mt.write_mix(&mixed)?;
mt.write_mix(&record_mix)?;
mt.end_cycle()
})();
if let Err(e) = res {
@@ -1209,13 +1415,13 @@ async fn run_core_loop(
}
} else if is_recording_mixer.load(Ordering::Relaxed)
&& let Some(rec) = recorder_mixer.lock().unwrap().as_mut()
&& let Err(e) = rec.write_frame(&mixed)
&& let Err(e) = rec.write_frame(&record_mix)
{
crate::log_msg(&format!("Recording write failed: {e}"));
}
let frame_to_send = if is_deafened_clone.load(Ordering::Relaxed) {
vec![0i16; FRAME_SAMPLES]
vec![0i16; mixed.len()]
} else {
mixed
};
@@ -1522,6 +1728,26 @@ async fn run_core_loop(
guard.insert(peer_id, vol);
}
CoreCommand::SetPeerEq(peer_id, settings) => {
let settings = settings.clamped();
let mut guard = peer_eq.lock().await;
if settings.is_flat() {
guard.remove(&peer_id);
} else {
guard.insert(peer_id, settings);
}
}
CoreCommand::SetPeerPan(peer_id, pan) => {
let pan = pan.clamp(-1.0, 1.0);
let mut guard = peer_pan.lock().await;
if pan.abs() <= 0.001 {
guard.remove(&peer_id);
} else {
guard.insert(peer_id, pan);
}
}
CoreCommand::SetPeerMuted(peer_id, muted) => {
let mut guard = locally_muted.lock().await;
if muted {
@@ -1649,22 +1875,60 @@ async fn run_core_loop(
}
CoreCommand::SetPresenceMode(mode) => {
*presence_mode.lock().unwrap() = mode;
// W7 P6: re-apply n0 DNS discovery for the new posture (publish on iff
// Discoverable). Runtime — no endpoint rebuild; clears + reinstalls the
// address-lookup services. The resolver stays on regardless so we can
// still look up moved friends.
let plan = crate::discovery::lookup_plan(network_mode, mode.publishes_to_discovery());
if let Err(e) = apply_discovery(&net.endpoint, &net.memory_lookup, plan) {
crate::log_msg(&format!("discovery: apply failed: {e:#}"));
let previous_mode = *presence_mode.lock().unwrap();
let now = tokio::time::Instant::now();
if previous_mode == mode {
// Same-mode requests are no-ops for discovery wiring, but keep the
// existing UX: re-selecting Discoverable restarts the clock.
discovery_deadline = if mode == PresenceMode::Discoverable {
Some(now + crate::discovery::DISCOVERY_TIMEBOX)
} else {
None
};
continue;
}
// Arm (Discoverable) or cancel (any other posture) the auto-revert
// time-box. Re-selecting Discoverable restarts the clock.
discovery_deadline = if mode == crate::presence::PresenceMode::Discoverable {
Some(tokio::time::Instant::now() + crate::discovery::DISCOVERY_TIMEBOX)
// W7 P6/S11: re-apply n0 DNS discovery for the requested posture
// first, then commit the presence mode only if the endpoint accepted
// that discovery plan. This keeps the UI truthful when dropping the
// publisher fails.
let plan =
crate::discovery::lookup_plan(network_mode, mode.publishes_to_discovery());
let apply_result = apply_discovery(&net.endpoint, &net.memory_lookup, plan);
let (committed_mode, transition_error) =
crate::discovery::resolve_presence_transition(
previous_mode,
mode,
apply_result.is_ok(),
);
*presence_mode.lock().unwrap() = committed_mode;
if committed_mode == PresenceMode::Discoverable {
if apply_result.is_ok() && mode == PresenceMode::Discoverable {
discovery_deadline = Some(now + crate::discovery::DISCOVERY_TIMEBOX);
} else {
arm_discovery_retry(&mut discovery_deadline, now);
}
} else {
None
};
discovery_deadline = None;
}
if let Err(e) = apply_result {
crate::log_msg(&format!("discovery: apply failed: {e:#}"));
if committed_mode != mode {
let _ = ui_tx
.send(UiEvent::PresenceModeReverted {
mode: committed_mode,
})
.await;
}
if let Some(message) = transition_error {
let _ = ui_tx
.send(UiEvent::Error(format!("{message} ({e:#})")))
.await;
}
}
}
CoreCommand::SetRecordingMode(mode) => {
@@ -1865,7 +2129,10 @@ async fn run_core_loop(
#[cfg(test)]
mod tests {
use super::{apply_volume, frame_level, mix_frames, MicLevelMeter, MIC_LEVEL_REPORT_SAMPLES};
use super::{
apply_volume, audio_datagram_len_ok, frame_level, mix_frames, mix_stereo_frames,
stereo_to_mono, MicLevelMeter, MAX_OPUS_PAYLOAD, MIC_LEVEL_REPORT_SAMPLES,
};
/// A frame of constant amplitude with the given sample count.
fn frame(amp: i16, len: usize) -> Vec<i16> {
@@ -1881,6 +2148,15 @@ mod tests {
assert!(m.push(&frame(1000, MIC_LEVEL_REPORT_SAMPLES)).is_some());
}
#[test]
fn audio_datagram_length_gate_preserves_header_and_caps_payload() {
assert!(!audio_datagram_len_ok(0));
assert!(!audio_datagram_len_ok(3));
assert!(audio_datagram_len_ok(4));
assert!(audio_datagram_len_ok(4 + MAX_OPUS_PAYLOAD));
assert!(!audio_datagram_len_ok(5 + MAX_OPUS_PAYLOAD));
}
#[test]
fn mic_meter_holds_the_peak_across_the_window() {
let mut m = MicLevelMeter::new();
@@ -1924,6 +2200,27 @@ mod tests {
assert_eq!(mixed, vec![100i32, -200, 300, -400]);
}
#[test]
fn centered_stereo_mix_is_exact_dual_mono() {
let a = vec![100, -200, 300, -400];
let b = vec![50, 200, -100, 400];
let mixed = mix_stereo_frames(&[(a, 0.0), (b, 0.0)], 4);
assert_eq!(mixed, vec![150, 150, 0, 0, 200, 200, 0, 0]);
}
#[test]
fn hard_left_pan_only_contributes_left_channel() {
let frame = vec![100, 200];
let mixed = mix_stereo_frames(&[(frame, -1.0)], 2);
assert_eq!(mixed, vec![141, 0, 283, 0]);
}
#[test]
fn stereo_fold_down_averages_pairs() {
let mono = stereo_to_mono(&[100, 100, 200, 0, i16::MAX, i16::MAX]);
assert_eq!(mono, vec![100, 100, i16::MAX]);
}
#[test]
fn two_peers_sum_sample_by_sample() {
let a = vec![100, -200, 300, -400];
@@ -2038,4 +2335,3 @@ mod tests {
assert!((level - 0.5).abs() < 1e-3, "mid-range level was {level}");
}
}
+96 -10
View File
@@ -8,14 +8,19 @@
//! The model (from `docs/contacts-plan.md` P6, decided 2026-06-16):
//! - **Resolving is always allowed on relay-capable modes** — a stationary friend
//! (typically in `Normal`) must be able to look up a friend who moved networks. A
//! resolve is a DNS query to n0 that publishes nothing; it only fires when a saved
//! address is stale and the dial falls through to discovery.
//! resolve is a DNS query to n0 that publishes nothing, but still exposes query
//! timing/source metadata to n0; it only fires when a saved address is stale and
//! the dial falls through to discovery.
//! - **Publishing is gated on `Discoverable`** and asymmetric: only the mover
//! publishes their address to n0 DNS; everyone else just looks it up.
//! - **Stopping publishing removes the local publisher service**; iroh does not
//! expose an explicit unpublish call here, so already-published pkarr records can
//! linger until their default ~30s TTL expires.
//! - **`DirectOnly` is the explicit no-server posture** — neither resolve nor publish
//! ever touches n0 there, regardless of the Discoverable toggle.
use crate::config::NetworkMode;
use crate::presence::PresenceMode;
use std::time::Duration;
/// How long `Discoverable` stays on before auto-reverting to `Normal`. Discovery is
@@ -46,12 +51,43 @@ pub fn lookup_plan(network_mode: NetworkMode, want_publish: bool) -> LookupPlan
match network_mode {
// The explicit serverless posture: no n0 contact at all, even to resolve.
// A Discoverable toggle here is intentionally inert.
NetworkMode::DirectOnly => LookupPlan { resolver: false, publisher: false },
NetworkMode::DirectOnly => LookupPlan {
resolver: false,
publisher: false,
},
// Relay-capable: always resolve (so a stationary friend can find a mover);
// publish only when the user opted into Discoverable.
NetworkMode::RelayNoDiscovery | NetworkMode::N0Full => {
LookupPlan { resolver: true, publisher: want_publish }
}
NetworkMode::RelayNoDiscovery | NetworkMode::N0Full => LookupPlan {
resolver: true,
publisher: want_publish,
},
}
}
/// Decide which presence mode may be committed after attempting to apply discovery
/// services for `requested`.
///
/// On failure, keep the previous mode: it is the only locally truthful state because
/// the endpoint's discovery services may still reflect the old posture. Same-mode
/// requests are no-ops from a presence-truth perspective and do not surface an error.
pub fn resolve_presence_transition(
previous: PresenceMode,
requested: PresenceMode,
apply_ok: bool,
) -> (PresenceMode, Option<String>) {
if previous == requested {
return (previous, None);
}
if apply_ok {
(requested, None)
} else {
(
previous,
Some(format!(
"Couldn't update discovery mode; keeping {previous}."
)),
)
}
}
@@ -64,12 +100,18 @@ mod tests {
for mode in [NetworkMode::RelayNoDiscovery, NetworkMode::N0Full] {
assert_eq!(
lookup_plan(mode, false),
LookupPlan { resolver: true, publisher: false },
LookupPlan {
resolver: true,
publisher: false
},
"{mode:?}: resolve always on, no publish when not Discoverable"
);
assert_eq!(
lookup_plan(mode, true),
LookupPlan { resolver: true, publisher: true },
LookupPlan {
resolver: true,
publisher: true
},
"{mode:?}: Discoverable adds publish on top of resolve"
);
}
@@ -79,12 +121,18 @@ mod tests {
fn direct_only_never_touches_n0_even_when_discoverable() {
assert_eq!(
lookup_plan(NetworkMode::DirectOnly, false),
LookupPlan { resolver: false, publisher: false }
LookupPlan {
resolver: false,
publisher: false
}
);
// The serverless posture overrides the Discoverable request entirely.
assert_eq!(
lookup_plan(NetworkMode::DirectOnly, true),
LookupPlan { resolver: false, publisher: false }
LookupPlan {
resolver: false,
publisher: false
}
);
}
@@ -92,4 +140,42 @@ mod tests {
fn timebox_is_thirty_minutes() {
assert_eq!(DISCOVERY_TIMEBOX, Duration::from_secs(1800));
}
#[test]
fn presence_transition_commits_requested_mode_after_successful_apply() {
assert_eq!(
resolve_presence_transition(PresenceMode::Normal, PresenceMode::Discoverable, true),
(PresenceMode::Discoverable, None)
);
}
#[test]
fn presence_transition_keeps_previous_mode_when_apply_fails() {
let (mode, err) =
resolve_presence_transition(PresenceMode::Normal, PresenceMode::Discoverable, false);
assert_eq!(mode, PresenceMode::Normal);
assert!(err.unwrap().contains("keeping Normal"));
}
#[test]
fn presence_transition_keeps_discoverable_when_off_transition_fails() {
let (mode, err) =
resolve_presence_transition(PresenceMode::Discoverable, PresenceMode::Normal, false);
assert_eq!(mode, PresenceMode::Discoverable);
assert!(err.unwrap().contains("keeping Discoverable"));
}
#[test]
fn presence_transition_same_mode_is_noop_without_error() {
assert_eq!(
resolve_presence_transition(
PresenceMode::Discoverable,
PresenceMode::Discoverable,
false
),
(PresenceMode::Discoverable, None)
);
}
}
+284
View File
@@ -0,0 +1,284 @@
//! Focused, app-local keyboard shortcuts.
//!
//! These helpers are intentionally pure: key serialization, formatting, lookup,
//! and conflict detection live here, while iced event handling stays at the app
//! edge. There are no OS-global shortcuts.
use iced::keyboard;
use serde::{Deserialize, Serialize};
/// A serializable key identity. Modifiers are deliberately out of scope for this
/// first pass; iced delivers the focused app key and we compare that exact key.
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum KeyBinding {
Named(String),
Character(String),
}
impl KeyBinding {
pub fn from_key(key: &keyboard::Key) -> Option<Self> {
match key {
keyboard::Key::Named(named) => Some(Self::Named(format!("{named:?}"))),
keyboard::Key::Character(ch) => {
let s = ch.to_string();
if s.is_empty() {
None
} else {
Some(Self::Character(s.to_lowercase()))
}
}
keyboard::Key::Unidentified => None,
}
}
pub fn label(&self) -> String {
match self {
KeyBinding::Named(name) => name.clone(),
KeyBinding::Character(ch) => ch.to_uppercase(),
}
}
}
/// Parse a hand-editable binding string from config/docs/tests. Empty and
/// `"unset"` are unbound.
pub fn parse_binding(input: &str) -> Option<KeyBinding> {
let trimmed = input.trim();
if trimmed.is_empty() || trimmed.eq_ignore_ascii_case("unset") {
return None;
}
if trimmed.chars().count() == 1 {
Some(KeyBinding::Character(trimmed.to_lowercase()))
} else {
Some(KeyBinding::Named(trimmed.to_string()))
}
}
pub fn format_binding(binding: Option<&KeyBinding>) -> String {
binding
.map(KeyBinding::label)
.unwrap_or_else(|| "unset".to_string())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum HotkeyAction {
ToggleMute,
ToggleDeafen,
OpenSettings,
PushToTalk,
LeaveRoom,
}
impl HotkeyAction {
pub const ALL: [HotkeyAction; 5] = [
HotkeyAction::ToggleMute,
HotkeyAction::ToggleDeafen,
HotkeyAction::OpenSettings,
HotkeyAction::PushToTalk,
HotkeyAction::LeaveRoom,
];
pub fn label(self) -> &'static str {
match self {
HotkeyAction::ToggleMute => "Toggle mute",
HotkeyAction::ToggleDeafen => "Toggle deafen",
HotkeyAction::OpenSettings => "Open Settings",
HotkeyAction::PushToTalk => "Push-to-talk",
HotkeyAction::LeaveRoom => "Leave room",
}
}
pub fn tier(self) -> HotkeyTier {
match self {
HotkeyAction::ToggleMute
| HotkeyAction::ToggleDeafen
| HotkeyAction::OpenSettings => HotkeyTier::AppWide,
HotkeyAction::PushToTalk | HotkeyAction::LeaveRoom => HotkeyTier::RoomOnly,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum HotkeyTier {
AppWide,
RoomOnly,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct HotkeyContext {
pub in_call: bool,
}
impl HotkeyContext {
fn allows(self, action: HotkeyAction) -> bool {
matches!(action.tier(), HotkeyTier::AppWide) || self.in_call
}
}
/// Persisted shortcut map. Defaults preserve the old Space push-to-talk binding
/// and add a few function-key app shortcuts that do not collide with typing.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct HotkeyMap {
#[serde(default = "default_mute")]
pub toggle_mute: Option<KeyBinding>,
#[serde(default = "default_deafen")]
pub toggle_deafen: Option<KeyBinding>,
#[serde(default = "default_settings")]
pub open_settings: Option<KeyBinding>,
#[serde(default = "default_ptt")]
pub push_to_talk: Option<KeyBinding>,
#[serde(default)]
pub leave_room: Option<KeyBinding>,
}
impl Default for HotkeyMap {
fn default() -> Self {
Self {
toggle_mute: default_mute(),
toggle_deafen: default_deafen(),
open_settings: default_settings(),
push_to_talk: default_ptt(),
leave_room: None,
}
}
}
fn named(name: &str) -> Option<KeyBinding> {
Some(KeyBinding::Named(name.to_string()))
}
fn default_mute() -> Option<KeyBinding> {
named("F9")
}
fn default_deafen() -> Option<KeyBinding> {
named("F10")
}
fn default_settings() -> Option<KeyBinding> {
named("F2")
}
fn default_ptt() -> Option<KeyBinding> {
named("Space")
}
impl HotkeyMap {
pub fn binding(&self, action: HotkeyAction) -> Option<&KeyBinding> {
match action {
HotkeyAction::ToggleMute => self.toggle_mute.as_ref(),
HotkeyAction::ToggleDeafen => self.toggle_deafen.as_ref(),
HotkeyAction::OpenSettings => self.open_settings.as_ref(),
HotkeyAction::PushToTalk => self.push_to_talk.as_ref(),
HotkeyAction::LeaveRoom => self.leave_room.as_ref(),
}
}
pub fn set_binding(&mut self, action: HotkeyAction, binding: Option<KeyBinding>) {
match action {
HotkeyAction::ToggleMute => self.toggle_mute = binding,
HotkeyAction::ToggleDeafen => self.toggle_deafen = binding,
HotkeyAction::OpenSettings => self.open_settings = binding,
HotkeyAction::PushToTalk => self.push_to_talk = binding,
HotkeyAction::LeaveRoom => self.leave_room = binding,
}
}
pub fn lookup_key(&self, key: &keyboard::Key, context: HotkeyContext) -> Option<HotkeyAction> {
let pressed = KeyBinding::from_key(key)?;
HotkeyAction::ALL
.into_iter()
.find(|&action| context.allows(action) && self.binding(action) == Some(&pressed))
}
pub fn lookup_binding(
&self,
binding: &KeyBinding,
context: HotkeyContext,
) -> Option<HotkeyAction> {
HotkeyAction::ALL
.into_iter()
.find(|&action| context.allows(action) && self.binding(action) == Some(binding))
}
pub fn conflicts(&self) -> Vec<HotkeyConflict> {
let mut conflicts = Vec::new();
let actions = HotkeyAction::ALL;
for i in 0..actions.len() {
for j in (i + 1)..actions.len() {
let a = actions[i];
let b = actions[j];
if let (Some(ab), Some(bb)) = (self.binding(a), self.binding(b))
&& ab == bb
{
conflicts.push(HotkeyConflict {
binding: ab.clone(),
first: a,
second: b,
});
}
}
}
conflicts
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct HotkeyConflict {
pub binding: KeyBinding,
pub first: HotkeyAction,
pub second: HotkeyAction,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn unset_actions_format_as_unset() {
assert_eq!(format_binding(None), "unset");
assert_eq!(parse_binding("unset"), None);
assert_eq!(parse_binding(""), None);
}
#[test]
fn duplicate_binding_is_detected() {
let mut map = HotkeyMap::default();
map.set_binding(HotkeyAction::ToggleMute, parse_binding("M"));
map.set_binding(HotkeyAction::ToggleDeafen, parse_binding("m"));
let conflicts = map.conflicts();
assert_eq!(conflicts.len(), 1);
assert_eq!(conflicts[0].first, HotkeyAction::ToggleMute);
assert_eq!(conflicts[0].second, HotkeyAction::ToggleDeafen);
}
#[test]
fn lookup_respects_room_tier() {
let mut map = HotkeyMap::default();
map.set_binding(HotkeyAction::LeaveRoom, parse_binding("Escape"));
let binding = parse_binding("Escape").unwrap();
assert_eq!(
map.lookup_binding(&binding, HotkeyContext { in_call: false }),
None,
"room-only shortcuts should not fire outside a call"
);
assert_eq!(
map.lookup_binding(&binding, HotkeyContext { in_call: true }),
Some(HotkeyAction::LeaveRoom)
);
}
#[test]
fn default_ptt_is_space() {
let map = HotkeyMap::default();
assert_eq!(
format_binding(map.binding(HotkeyAction::PushToTalk)),
"Space"
);
}
#[test]
fn parse_single_character_case_folds() {
assert_eq!(parse_binding("M"), Some(KeyBinding::Character("m".to_string())));
assert_eq!(format_binding(parse_binding("m").as_ref()), "M");
}
}
+118 -6
View File
@@ -16,10 +16,15 @@ pub mod sanitize;
pub mod avatar;
pub mod recents;
pub mod discovery;
pub mod hotkeys;
use std::path::PathBuf;
use std::fs::File;
use std::path::{Path, PathBuf};
use std::sync::OnceLock;
const LOG_MAX_BYTES: u64 = 5 * 1024 * 1024;
const LOG_MODE: u32 = 0o600;
/// Resolves the log file path once: `$XDG_STATE_HOME/peerspeak/peerspeak.log`
/// (via `dirs::state_dir`), falling back to the system temp dir. Computed lazily
/// so we never hardcode a per-user path.
@@ -42,6 +47,65 @@ pub fn log_file_path() -> PathBuf {
log_path().clone()
}
/// Short, human-matchable id prefix for diagnostics. Never use this where the
/// full value is needed for protocol behavior.
pub fn short_id(id: &str) -> String {
id.chars().take(8).collect()
}
/// Redact a capability-bearing value for logs while keeping a tiny prefix for
/// support correlation. Tickets and endpoint addresses are bearer capabilities:
/// logging the full string is equivalent to leaking the room/share.
pub fn redact_for_log(value: &str) -> String {
let value = value.trim();
if value.is_empty() {
"<redacted:empty>".to_string()
} else {
format!("<redacted:{}...>", short_id(value))
}
}
pub fn short_bytes_hex(bytes: &[u8]) -> String {
bytes.iter()
.take(6)
.map(|b| format!("{b:02x}"))
.collect::<Vec<_>>()
.join("")
}
fn rotated_log_path(path: &Path) -> PathBuf {
let file_name = path.file_name().and_then(|n| n.to_str()).unwrap_or("peerspeak.log");
path.with_file_name(format!("{file_name}.1"))
}
fn prepare_log_file(path: &Path) -> std::io::Result<File> {
prepare_log_file_with_limit(path, LOG_MAX_BYTES)
}
fn prepare_log_file_with_limit(path: &Path, max_bytes: u64) -> std::io::Result<File> {
use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
if let Some(parent) = path.parent() {
let _ = std::fs::create_dir_all(parent);
}
if std::fs::metadata(path).is_ok_and(|m| m.len() > max_bytes) {
let rotated = rotated_log_path(path);
let _ = std::fs::remove_file(&rotated);
if std::fs::rename(path, &rotated).is_err() {
let _ = std::fs::OpenOptions::new().write(true).truncate(true).open(path);
}
}
let file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.mode(LOG_MODE)
.open(path)?;
let _ = std::fs::set_permissions(path, std::fs::Permissions::from_mode(LOG_MODE));
Ok(file)
}
pub fn log_msg(msg: &str) {
// Format the whole line into one buffer first, then emit it with a single
// `write_all`. The file is opened with `O_APPEND`, so a lone `write()` is
@@ -51,13 +115,61 @@ pub fn log_msg(msg: &str) {
Ok(time) => format!("[{}.{:03}] {}\n", time.as_secs(), time.subsec_millis(), msg),
Err(_) => format!("{}\n", msg),
};
if let Ok(mut file) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(log_path())
{
if let Ok(mut file) = prepare_log_file(log_path()) {
use std::io::Write;
let _ = file.write_all(line.as_bytes());
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write;
use std::os::unix::fs::PermissionsExt;
fn temp_log_dir() -> PathBuf {
let stamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
std::env::temp_dir().join(format!("peerspeak-log-test-{}-{stamp}", std::process::id()))
}
#[test]
fn redaction_keeps_only_a_short_prefix() {
let secret = "abcdefghijklmnopqrstuvwxyz";
let redacted = redact_for_log(secret);
assert!(redacted.contains("abcdefgh"));
assert!(!redacted.contains("ijklmnopqrstuvwxyz"));
assert_eq!(redact_for_log(" "), "<redacted:empty>");
}
#[test]
fn log_file_is_created_private() {
let dir = temp_log_dir();
let path = dir.join("peerspeak.log");
let _file = prepare_log_file(&path).unwrap();
let mode = std::fs::metadata(&path).unwrap().permissions().mode() & 0o777;
assert_eq!(mode, LOG_MODE);
let _ = std::fs::remove_dir_all(dir);
}
#[test]
fn oversized_log_is_rotated_on_open() {
let dir = temp_log_dir();
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("peerspeak.log");
{
let mut file = std::fs::File::create(&path).unwrap();
file.write_all(b"oversized").unwrap();
}
let _file = prepare_log_file_with_limit(&path, 4).unwrap();
let rotated = rotated_log_path(&path);
assert_eq!(std::fs::read_to_string(rotated).unwrap(), "oversized");
assert_eq!(std::fs::metadata(&path).unwrap().len(), 0);
let _ = std::fs::remove_dir_all(dir);
}
}
+168 -13
View File
@@ -43,7 +43,7 @@ impl std::fmt::Debug for GossipPayload {
f.debug_struct("GossipPayload")
.field("author", &self.author)
.field("ts", &self.ts)
.field("msg", &self.msg)
.field("msg_kind", &gossip_message_kind(&self.msg))
.finish_non_exhaustive()
}
}
@@ -79,6 +79,58 @@ enum GossipReject {
BadSignature,
/// Timestamp outside the freshness window — stale (replay) or implausibly future.
OutOfWindow,
/// A signed Announce advertised an address for a different node id.
AnnounceAddressMismatch,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
enum StateMutationKind {
Announce,
Leave,
}
fn gossip_message_kind(msg: &GossipMessage) -> &'static str {
match msg {
GossipMessage::Announce(_) => "Announce",
GossipMessage::Leave => "Leave",
GossipMessage::Chat { .. } => "Chat",
}
}
fn state_mutation_kind(msg: &GossipMessage) -> Option<StateMutationKind> {
match msg {
GossipMessage::Announce(_) => Some(StateMutationKind::Announce),
GossipMessage::Leave => Some(StateMutationKind::Leave),
GossipMessage::Chat { .. } => None,
}
}
fn admit_state_mutation(
seen: &mut HashMap<(EndpointId, StateMutationKind), u64>,
author: EndpointId,
msg: &GossipMessage,
ts: u64,
) -> bool {
let Some(kind) = state_mutation_kind(msg) else {
return true;
};
let key = (author, kind);
if seen.get(&key).is_some_and(|last_ts| ts <= *last_ts) {
return false;
}
seen.insert(key, ts);
true
}
fn peer_state_for_log(state: &PeerState) -> String {
format!(
"name={:?}, muted={}, addr_id={}, addrs={}, sharing={}",
state.name,
state.is_muted,
crate::short_id(&state.addr.id.to_string()),
state.addr.addrs.len(),
state.sharing.is_some()
)
}
/// Authenticate a received payload against the room topic and local clock. The
@@ -99,6 +151,10 @@ fn verify_gossip(
if now_ms.abs_diff(payload.ts) > window_ms {
return Err(GossipReject::OutOfWindow);
}
if let GossipMessage::Announce(state) = &payload.msg
&& state.addr.id != payload.author {
return Err(GossipReject::AnnounceAddressMismatch);
}
Ok(())
}
@@ -185,11 +241,21 @@ impl RoomState for IrohGossipState {
self_state: PeerState,
extra_bootstrap: Vec<EndpointAddr>,
) -> Result<(), NetError> {
crate::log_msg(&format!("RoomState::join: self_id={:?}, self_name={:?}, ticket={}", self_state.addr.id, self_state.name, ticket_str));
crate::log_msg(&format!(
"RoomState::join: self_id={}, self_name={:?}, ticket={}",
crate::short_id(&self_state.addr.id.to_string()),
self_state.name,
crate::redact_for_log(ticket_str)
));
let ticket = ticket_str.parse::<PeerSpeakTicket>()?;
let topic_id = TopicId::from_bytes(ticket.topic_id);
crate::log_msg(&format!("Parsed ticket. host_id={:?}, host_addrs={:?}, topic={:?}", ticket.host_addr.id, ticket.host_addr.addrs, topic_id));
crate::log_msg(&format!(
"Parsed ticket. host_id={}, host_addrs={}, topic={}",
crate::short_id(&ticket.host_addr.id.to_string()),
ticket.host_addr.addrs.len(),
crate::short_bytes_hex(&ticket.topic_id)
));
// Stop any currently running topic
let _ = self.leave().await;
@@ -236,6 +302,7 @@ impl RoomState for IrohGossipState {
let handle = tokio::spawn(async move {
crate::log_msg(&format!("Spawned gossip topic loop for self_id={:?}", self_id));
let mut state_mutations_seen = HashMap::new();
// Broadcast initial state
let initial_payload = {
@@ -285,7 +352,26 @@ impl RoomState for IrohGossipState {
continue;
}
crate::log_msg(&format!("Gossip Event::Received from author={:?}, payload={:?}", payload.author, payload.msg));
if !admit_state_mutation(
&mut state_mutations_seen,
payload.author,
&payload.msg,
payload.ts,
) {
crate::log_msg(&format!(
"Gossip dropped replayed state mutation author={}, kind={}, ts={}",
crate::short_id(&payload.author.to_string()),
gossip_message_kind(&payload.msg),
payload.ts
));
continue;
}
crate::log_msg(&format!(
"Gossip Event::Received author={}, kind={}",
crate::short_id(&payload.author.to_string()),
gossip_message_kind(&payload.msg)
));
match payload.msg {
GossipMessage::Announce(mut state) => {
@@ -299,6 +385,10 @@ impl RoomState for IrohGossipState {
// monogram, so a malformed/oversized/bomb
// image can't crash or exhaust us (W4).
state.avatar = state.avatar.sanitize_incoming();
// Screen-share tickets are capabilities and
// peer-supplied: cap/validate once at ingest
// so invalid offers never render a Watch button.
state.sharing = state.sharing.and_then(crate::screenshare::sanitize_ticket);
let (is_new, state_changed) = {
let mut peer_map = peers.lock().unwrap();
let is_new = !peer_map.contains_key(&payload.author);
@@ -310,11 +400,19 @@ impl RoomState for IrohGossipState {
};
if is_new {
crate::log_msg(&format!("Gossip new peer joined: {:?}, state: {:?}", payload.author, state));
crate::log_msg(&format!(
"Gossip new peer joined: {}, state: {}",
crate::short_id(&payload.author.to_string()),
peer_state_for_log(&state)
));
address_lookup.add_endpoint_info(state.addr.clone());
let _ = event_tx.send(RoomEvent::PeerJoined(payload.author, state)).await;
} else if state_changed {
crate::log_msg(&format!("Gossip peer state updated: {:?}, state: {:?}", payload.author, state));
crate::log_msg(&format!(
"Gossip peer state updated: {}, state: {}",
crate::short_id(&payload.author.to_string()),
peer_state_for_log(&state)
));
let _ = event_tx.send(RoomEvent::PeerUpdated(payload.author, state)).await;
}
}
@@ -390,7 +488,10 @@ impl RoomState for IrohGossipState {
}
async fn update_self_state(&self, self_state: PeerState) -> Result<(), NetError> {
crate::log_msg(&format!("RoomState::update_self_state: state={:?}", self_state));
crate::log_msg(&format!(
"RoomState::update_self_state: state: {}",
peer_state_for_log(&self_state)
));
*self.self_state.lock().unwrap() = Some(self_state.clone());
let sender_opt = self.active_sender.lock().unwrap().clone();
@@ -489,11 +590,10 @@ mod tests {
use super::*;
use crate::network::PeerState;
use iroh::SecretKey;
use std::collections::HashMap;
fn sample_peer_state() -> PeerState {
let secret = SecretKey::generate();
let public = secret.public();
let addr = iroh::EndpointAddr::from(public);
fn sample_peer_state_for(id: EndpointId) -> PeerState {
let addr = iroh::EndpointAddr::from(id);
PeerState {
name: "TestPeerGossip".to_string(),
is_muted: true,
@@ -563,7 +663,7 @@ mod tests {
fn test_gossip_payload_announce_round_trip() {
let secret = SecretKey::generate();
let topic = [9u8; 32];
let peer_state = sample_peer_state();
let peer_state = sample_peer_state_for(secret.public());
let payload = sign_gossip(&secret, &topic, 1000, GossipMessage::Announce(peer_state.clone()));
let serialized = serde_json::to_string(&payload).unwrap();
@@ -732,5 +832,60 @@ mod tests {
// Within the window (clock skew tolerance) → accepted.
assert!(verify_gossip(&p, &topic, 1_000_000 + GOSSIP_FRESHNESS_MS - 1, GOSSIP_FRESHNESS_MS).is_ok());
}
}
#[test]
fn verify_rejects_announce_with_address_for_another_identity() {
let signer = SecretKey::generate();
let advertised = SecretKey::generate();
let topic = [6u8; 32];
let state = sample_peer_state_for(advertised.public());
let p = sign_gossip(&signer, &topic, 5_000, GossipMessage::Announce(state));
assert_eq!(
verify_gossip(&p, &topic, 5_000, GOSSIP_FRESHNESS_MS),
Err(GossipReject::AnnounceAddressMismatch)
);
}
#[test]
fn state_mutation_replay_gate_drops_replayed_leave_and_announce() {
let author = fresh_id();
let mut seen = HashMap::new();
assert!(admit_state_mutation(&mut seen, author, &GossipMessage::Leave, 10));
assert!(!admit_state_mutation(&mut seen, author, &GossipMessage::Leave, 10));
assert!(!admit_state_mutation(&mut seen, author, &GossipMessage::Leave, 9));
assert!(admit_state_mutation(&mut seen, author, &GossipMessage::Leave, 11));
let announce = GossipMessage::Announce(sample_peer_state_for(author));
assert!(admit_state_mutation(&mut seen, author, &announce, 10));
assert!(!admit_state_mutation(&mut seen, author, &announce, 10));
assert!(!admit_state_mutation(&mut seen, author, &announce, 9));
assert!(admit_state_mutation(&mut seen, author, &announce, 12));
}
#[test]
fn state_mutation_replay_gate_leaves_chat_ordering_untouched() {
let author = fresh_id();
let mut seen = HashMap::new();
let later_chat = GossipMessage::Chat { name: "A".into(), text: "later".into(), ts: 200 };
let earlier_chat = GossipMessage::Chat { name: "A".into(), text: "earlier".into(), ts: 100 };
assert!(admit_state_mutation(&mut seen, author, &later_chat, 200));
assert!(admit_state_mutation(&mut seen, author, &earlier_chat, 100));
assert!(admit_state_mutation(&mut seen, author, &later_chat, 200));
assert!(seen.is_empty(), "chat must not populate the state-mutation replay map");
}
#[test]
fn state_mutation_replay_gate_is_per_author_and_kind() {
let author = fresh_id();
let other = fresh_id();
let mut seen = HashMap::new();
let announce = GossipMessage::Announce(sample_peer_state_for(author));
assert!(admit_state_mutation(&mut seen, author, &GossipMessage::Leave, 5));
assert!(admit_state_mutation(&mut seen, author, &announce, 5));
assert!(admit_state_mutation(&mut seen, other, &GossipMessage::Leave, 5));
}
}
+40 -21
View File
@@ -110,26 +110,32 @@ pub enum FriendPresence {
InRoom { name: String, ticket: String },
}
/// Interpret a peer's reply defensively. Only a `Pong` is a reply (a `Ping` is
/// not, so it yields `None`). When the peer reports a room, we **sanitize the
/// peer-supplied name** and **only surface it as joinable if the ticket actually
/// parses** as a [`crate::network::PeerSpeakTicket`] — a garbage or hostile
/// ticket downgrades the friend to plain `Online` rather than offering a dead /
/// dangerous Join button. (We still never auto-join; the user clicks.)
pub fn interpret_pong(msg: &ControlMsg) -> Option<FriendPresence> {
/// Interpret a peer's reply defensively. `from` must be the connection's
/// authenticated remote id, not any value carried in the payload. Only a `Pong`
/// is a reply (a `Ping` is not, so it yields `None`). When the peer reports a
/// room, we **sanitize the peer-supplied name** and **only surface it as joinable
/// if the ticket actually parses** as a [`crate::network::PeerSpeakTicket`] and
/// points back at the replying friend. A garbage/redirect ticket downgrades the
/// friend to plain `Online` rather than offering a dead or attacker-controlled
/// Join button. (We still never auto-join; the user clicks.)
pub fn interpret_pong(msg: &ControlMsg, from: EndpointId) -> Option<FriendPresence> {
match msg {
ControlMsg::Ping => None,
ControlMsg::Pong { room: None } => Some(FriendPresence::Online),
ControlMsg::Pong { room: Some(r) } => {
if r.ticket.parse::<crate::network::PeerSpeakTicket>().is_ok() {
Some(FriendPresence::InRoom {
name: crate::sanitize::sanitize_name(&r.name),
ticket: r.ticket.clone(),
})
} else {
let Ok(ticket) = r.ticket.parse::<crate::network::PeerSpeakTicket>() else {
// Online, but the advertised room is unusable — don't offer Join.
Some(FriendPresence::Online)
return Some(FriendPresence::Online);
};
if ticket.host_addr.id != from {
// Online, but the advertised room redirects away from the friend
// who authenticated this Pong — don't offer a phishing Join.
return Some(FriendPresence::Online);
}
Some(FriendPresence::InRoom {
name: crate::sanitize::sanitize_name(&r.name),
ticket: r.ticket.clone(),
})
}
}
}
@@ -206,21 +212,22 @@ mod tests {
#[test]
fn interpret_ping_is_not_a_reply() {
assert_eq!(interpret_pong(&ControlMsg::Ping), None);
assert_eq!(interpret_pong(&ControlMsg::Ping, id()), None);
}
#[test]
fn interpret_pong_online_and_inroom() {
let friend = id();
// No room -> Online.
assert_eq!(
interpret_pong(&ControlMsg::Pong { room: None }),
interpret_pong(&ControlMsg::Pong { room: None }, friend),
Some(FriendPresence::Online)
);
// Valid ticket -> InRoom with a sanitized name.
let t = valid_ticket(id());
let t = valid_ticket(friend);
let got = interpret_pong(&ControlMsg::Pong {
room: Some(RoomPresence { name: "HangOut".into(), ticket: t.clone() }),
});
}, friend);
assert_eq!(got, Some(FriendPresence::InRoom { name: "HangOut".into(), ticket: t }));
}
@@ -230,17 +237,29 @@ mod tests {
// Online — no dead/hostile Join button is surfaced.
let got = interpret_pong(&ControlMsg::Pong {
room: Some(RoomPresence { name: "Trap".into(), ticket: "not-a-ticket".into() }),
});
}, id());
assert_eq!(got, Some(FriendPresence::Online));
}
#[test]
fn interpret_pong_rejects_ticket_for_a_different_host() {
let friend = id();
let attacker = id();
let t = valid_ticket(attacker);
let got = interpret_pong(&ControlMsg::Pong {
room: Some(RoomPresence { name: "Redirect".into(), ticket: t }),
}, friend);
assert_eq!(got, Some(FriendPresence::Online));
}
#[test]
fn interpret_pong_sanitizes_a_hostile_room_name() {
// Control/bidi characters in a peer-supplied name are stripped.
let t = valid_ticket(id());
let friend = id();
let t = valid_ticket(friend);
let got = interpret_pong(&ControlMsg::Pong {
room: Some(RoomPresence { name: "Hang\u{202e}Out\u{0007}".into(), ticket: t.clone() }),
});
}, friend);
match got {
Some(FriendPresence::InRoom { name, .. }) => {
assert!(!name.contains('\u{202e}'), "bidi override must be stripped");
+17 -15
View File
@@ -49,16 +49,17 @@ fn decode(bytes: &[u8]) -> Result<ControlMsg> {
serde_json::from_slice(bytes).context("failed to decode control message")
}
/// Probe `peer` for presence: send a `Ping`, return their `Pong`. An error means
/// no usable reply (offline / unreachable / refused / malformed) — the caller
/// treats that as "appears offline". `peer` is usually a bare [`EndpointId`]
/// (friends store the stable id); a full [`EndpointAddr`] is also accepted (and
/// used by hermetic tests).
pub async fn probe(endpoint: &Endpoint, peer: impl Into<EndpointAddr>) -> Result<ControlMsg> {
/// Probe `peer` for presence: send a `Ping`, return their authenticated id and
/// `Pong`. An error means no usable reply (offline / unreachable / refused /
/// malformed) — the caller treats that as "appears offline". `peer` is usually a
/// bare [`EndpointId`] (friends store the stable id); a full [`EndpointAddr`] is
/// also accepted (and used by hermetic tests).
pub async fn probe(endpoint: &Endpoint, peer: impl Into<EndpointAddr>) -> Result<(EndpointId, ControlMsg)> {
let conn = tokio::time::timeout(IO_TIMEOUT, endpoint.connect(peer, FRIENDS_ALPN))
.await
.context("timed out connecting to peer")?
.context("failed to connect to peer")?;
let from = conn.remote_id();
let io = async {
let (mut send, mut recv) = conn.open_bi().await.context("failed to open control stream")?;
@@ -77,7 +78,7 @@ pub async fn probe(endpoint: &Endpoint, peer: impl Into<EndpointAddr>) -> Result
.await
.context("timed out awaiting pong")?;
conn.close(VarInt::from_u32(0), b"done");
result
result.map(|msg| (from, msg))
}
/// A reply policy: given the *authenticated* remote id, decide whether and how to
@@ -110,6 +111,10 @@ async fn handle(incoming: Incoming, handler: Handler) -> Result<()> {
async fn exchange(conn: &iroh::endpoint::Connection, handler: &Handler) -> Result<()> {
// The authenticated remote id — NOT anything the peer puts in the payload.
let from = conn.remote_id();
let Some(reply) = handler(from) else {
conn.close(VarInt::from_u32(0), b"not authorized");
return Ok(());
};
let io = async {
let (mut send, mut recv) = conn.accept_bi().await.context("failed to accept stream")?;
@@ -118,13 +123,9 @@ async fn exchange(conn: &iroh::endpoint::Connection, handler: &Handler) -> Resul
ControlMsg::Ping => {}
other => bail!("expected a ping, got {other:?}"),
}
// Ask the policy what to send. None -> answer nothing (stranger / invisible):
// finish the stream with no bytes so the prober sees an empty (unusable) reply.
if let Some(reply) = handler(from) {
send.write_all(&encode(&reply)?)
.await
.context("failed to write pong")?;
}
send.write_all(&encode(&reply)?)
.await
.context("failed to write pong")?;
send.finish().context("failed to finish reply stream")?;
Ok::<_, anyhow::Error>(())
};
@@ -220,10 +221,11 @@ mod tests {
let serve_task = tokio::spawn(async move { serve(server_ep, handler).await });
// The allowed prober gets a Pong with the room.
let pong = tokio::time::timeout(Duration::from_secs(15), probe(&prober, server_addr.clone()))
let (from, pong) = tokio::time::timeout(Duration::from_secs(15), probe(&prober, server_addr.clone()))
.await
.expect("probe timed out")
.expect("probe failed");
assert_eq!(from, server_addr.id);
match pong {
ControlMsg::Pong { room: Some(r) } => assert_eq!(r.name, "HangOut"),
other => panic!("expected Pong with a room, got {other:?}"),
+56 -1
View File
@@ -25,6 +25,10 @@ use tokio::process::{Child, Command};
/// points elsewhere.
const PIXELPASS_BIN: &str = "pixelpass";
/// Pixelpass endpoint tickets are normally ~140 chars. Leave headroom for format
/// growth, but reject unbounded gossip payloads before the UI offers "Watch".
const MAX_TICKET_LEN: usize = 512;
/// How long to wait for the host to emit its ticket / the viewer to connect
/// before giving up and killing the child. Startup is normally sub-second; this
/// is only a safety net so a hung pixelpass can't wedge the caller forever.
@@ -108,6 +112,19 @@ pub fn viewer_args(ticket: &str) -> Vec<String> {
]
}
/// Sanitize a peer-advertised pixelpass ticket at the gossip boundary. Peerspeak
/// intentionally does not depend on pixelpass/iroh-tickets, so this validates the
/// stable CLI ticket envelope we consume: bounded ASCII endpoint tickets beginning
/// with `endpoint`. Invalid input becomes `None`, which removes the Watch button.
pub fn sanitize_ticket(ticket: String) -> Option<String> {
let ticket = ticket.trim();
let valid_len = !ticket.is_empty() && ticket.len() <= MAX_TICKET_LEN;
let valid_shape = ticket.starts_with("endpoint")
&& ticket.len() > "endpoint".len()
&& ticket.bytes().all(|b| b.is_ascii_alphanumeric());
(valid_len && valid_shape).then(|| ticket.to_string())
}
/// Resolve the pixelpass binary: an explicit config override (used only if it
/// points at an existing file), otherwise the first `pixelpass` found on
/// `$PATH`. `None` means it isn't installed — a normal, handled state. An
@@ -267,12 +284,29 @@ where
tokio::spawn(async move {
while let Ok(Some(line)) = lines.next_line().await {
if let Some(ev) = parse_pixelpass_event(&line) {
crate::log_msg(&format!("pixelpass {role}: {ev:?}"));
crate::log_msg(&format!("pixelpass {role}: {}", event_for_log(&ev)));
}
}
});
}
fn event_for_log(ev: &PixelpassEvent) -> String {
match ev {
PixelpassEvent::Ticket(ticket) => format!("ticket {}", crate::redact_for_log(ticket)),
PixelpassEvent::Connected(_) => "connected".to_string(),
PixelpassEvent::ViewerJoined { active, max } => {
format!("viewer_joined active={active} max={max}")
}
PixelpassEvent::ViewerLeft { active, max } => {
format!("viewer_left active={active} max={max}")
}
PixelpassEvent::Refused(reason) => format!("viewer_refused reason={reason:?}"),
PixelpassEvent::CaptureStarted => "capture_started".to_string(),
PixelpassEvent::CaptureStopped => "capture_stopped".to_string(),
PixelpassEvent::Other => "other".to_string(),
}
}
/// Open the viewer stream URL in a media player. Mirrors pixelpass's own
/// low-latency mpv invocation; falls back to vlc. The player is reaped in a
/// background task so it doesn't linger as a zombie when its window closes.
@@ -340,6 +374,27 @@ mod tests {
);
}
#[test]
fn sanitize_ticket_accepts_pixelpass_endpoint_ticket_shape() {
let ticket = "endpointaabwxjexzensznfvuudiapn5tyzws3angd2merarm";
assert_eq!(sanitize_ticket(format!(" {ticket}\n")), Some(ticket.to_string()));
}
#[test]
fn sanitize_ticket_rejects_oversized_or_garbage_ticket() {
assert_eq!(sanitize_ticket("not-a-ticket".into()), None);
assert_eq!(sanitize_ticket(format!("endpoint{}", "a".repeat(MAX_TICKET_LEN))), None);
assert_eq!(sanitize_ticket("endpointabc-def".into()), None);
}
#[test]
fn event_log_redacts_ticket_values() {
let ticket = "endpointaabwxjexzensznfvuudiapn5tyzws3angd2merarm".to_string();
let log = event_for_log(&PixelpassEvent::Ticket(ticket.clone()));
assert!(log.contains("endpoint"));
assert!(!log.contains(&ticket["endpoint".len() + 8..]));
}
#[test]
fn parses_ticket() {
assert_eq!(