//! Live-edge catch-up for the screen-share viewer. //! //! PixelPass carries the share as MPEG-TS over a reliable, ordered transport. On //! a lossy link (satellite handovers are the pathological case) every loss burst //! becomes retransmission plus head-of-line blocking, and the viewer absorbs the //! stall as buffered latency. Nothing in the chain ever trims that buffer back, //! so the picture ends up seconds behind the host and stays there. //! //! Measured on a `tc netem` rig that simulates a satellite link (40 ms +/- 20 ms //! jitter, 0.5% loss, a 250 ms/30%-loss handover burst every 15 s): a viewer with //! ordinary timestamp pacing settles ~1.24 s behind. mpv's `--untimed` does NOT //! help (~1.38 s, marginally worse) because it only removes pacing at //! *presentation* while audio still drains at 1x the DAC rate, so an accumulated //! buffer never shrinks. Returning to the live edge requires consuming the //! backlog faster than it arrives. //! //! So we nudge playback slightly faster than realtime while the buffer is deep, //! and drop back to 1x once it has drained. mpv's default pitch correction //! (`scaletempo2`) keeps a 5% speedup inaudible, and because audio and video are //! sped up together A/V sync is preserved — unlike `--untimed`. //! //! The control law and the JSON-IPC message handling are pure functions with //! tests; the only I/O is [`drive`], which talks to mpv's `--input-ipc-server` //! socket. use std::path::{Path, PathBuf}; use std::time::Duration; /// Buffer depth (seconds) above which we start draining. pub const CACHE_HIGH_S: f64 = 1.0; /// Buffer depth (seconds) below which we return to realtime. pub const CACHE_LOW_S: f64 = 0.4; /// The buffer depth we aim to sit at; the drain rate is proportional to how far /// above this the buffer actually is. pub const CACHE_TARGET_S: f64 = 0.5; /// Extra playback rate per second of excess buffer. pub const CATCHUP_GAIN: f64 = 0.05; /// Hard ceiling on the drain rate. Beyond this the speedup stops being /// unnoticeable, and a share that far behind is better served by the operator /// restarting it than by a chipmunk impression. pub const MAX_CATCHUP_SPEED: f64 = 1.15; /// Normal realtime playback. pub const NORMAL_SPEED: f64 = 1.0; /// How often we sample the buffer depth. pub const POLL_INTERVAL: Duration = Duration::from_millis(500); /// Smallest rate change worth sending to the player. pub const SPEED_EPSILON: f64 = 0.005; /// The property we watch on the viewer. const CACHE_PROPERTY: &str = "demuxer-cache-duration"; /// Decide the playback rate for the next interval. /// /// Proportional, because a fixed small speedup cannot recover a large backlog in /// any reasonable time: draining 6 s at 1.05x takes two minutes, which a viewer /// experiences as "still broken". The drain rate instead scales with how deep /// the buffer is, so a bad handover is cleared in tens of seconds while a small /// excursion still gets only a gentle, inaudible nudge. /// /// Deliberately hysteretic: between [`CACHE_LOW_S`] and [`CACHE_HIGH_S`] the /// current rate is held, so a buffer hovering near a single threshold cannot /// oscillate the speed (and with it the audio pitch) every poll. Pure. /// /// A non-finite reading (mpv reports `null` before playback starts, and the /// caller maps that to NaN) holds the current rate rather than guessing. pub fn catchup_speed(cache_s: f64, current: f64) -> f64 { if !cache_s.is_finite() { return current; } if cache_s < CACHE_LOW_S { return NORMAL_SPEED; } if cache_s <= CACHE_HIGH_S { return current; } let excess = cache_s - CACHE_TARGET_S; (NORMAL_SPEED + CATCHUP_GAIN * excess).clamp(NORMAL_SPEED, MAX_CATCHUP_SPEED) } /// Where mpv should create its IPC socket. Kept separate from the runtime /// lookup so tests can pin a directory. Pure. pub fn socket_path(dir: &Path, token: u64) -> PathBuf { dir.join(format!("peerspeak-mpv-{token}.sock")) } /// The directory for the IPC socket: the XDG runtime dir when the session /// provides one (tmpfs, user-private, cleaned at logout), else the temp dir. pub fn socket_dir() -> PathBuf { std::env::var_os("XDG_RUNTIME_DIR") .map(PathBuf::from) .unwrap_or_else(std::env::temp_dir) } /// A `get_property` request for the buffer depth. Pure. pub fn get_cache_request(request_id: u64) -> String { format!(r#"{{"command":["get_property","{CACHE_PROPERTY}"],"request_id":{request_id}}}"#) } /// A `set_property` request for the playback rate. Pure. pub fn set_speed_request(request_id: u64, speed: f64) -> String { format!(r#"{{"command":["set_property","speed",{speed}],"request_id":{request_id}}}"#) } /// Extract the buffer depth from one line of mpv's IPC output. /// /// mpv interleaves unsolicited event lines with command replies, so a line is /// only ours when it carries the matching `request_id`. Returns: /// - `Some(Some(secs))` — our reply, with a usable number, /// - `Some(None)` — our reply, but no number (mpv sends `"data":null` before /// playback starts, and reports `error` while the demuxer has no cache yet), /// - `None` — not our reply (an event, or another command's response). /// /// Pure. pub fn parse_cache_response(line: &str, request_id: u64) -> Option> { let value: serde_json::Value = serde_json::from_str(line.trim()).ok()?; let id = value.get("request_id")?.as_u64()?; if id != request_id { return None; } if value.get("error").and_then(|e| e.as_str()) != Some("success") { return Some(None); } Some(value.get("data").and_then(|d| d.as_f64())) } /// Drive one mpv viewer's playback rate over its JSON IPC socket. /// /// Runs until mpv exits (the socket dies), so it is spawned detached alongside /// the player and needs no shutdown signal. Every failure path just ends the /// task: catch-up is an optimization, and a viewer that never gets it still /// plays, exactly as before this existed. #[cfg(unix)] pub async fn drive(socket: PathBuf) { use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::UnixStream; // mpv creates the socket a moment after exec, so the first connects race it. let mut stream = None; for _ in 0..40 { match UnixStream::connect(&socket).await { Ok(s) => { stream = Some(s); break; } Err(_) => tokio::time::sleep(Duration::from_millis(250)).await, } } let Some(stream) = stream else { crate::log_msg("livesync: mpv IPC socket never appeared; catch-up disabled"); return; }; let (read_half, mut write_half) = stream.into_split(); let mut lines = BufReader::new(read_half).lines(); let mut request_id: u64 = 0; let mut speed = NORMAL_SPEED; loop { tokio::time::sleep(POLL_INTERVAL).await; request_id += 1; let query = format!("{}\n", get_cache_request(request_id)); if write_half.write_all(query.as_bytes()).await.is_err() { break; } // Skip event lines until our reply arrives. let cache = loop { match lines.next_line().await { Ok(Some(line)) => { if let Some(value) = parse_cache_response(&line, request_id) { break value; } } // Socket closed or unreadable: mpv is gone. _ => return, } }; let cache = cache.unwrap_or(f64::NAN); let next = catchup_speed(cache, speed); // A proportional law would otherwise re-send on every wobble of the // reading; only a change worth hearing is worth a round trip. if (next - speed).abs() > SPEED_EPSILON { speed = next; request_id += 1; let set = format!("{}\n", set_speed_request(request_id, speed)); if write_half.write_all(set.as_bytes()).await.is_err() { break; } crate::log_msg(&format!( "livesync: cache {cache:.2}s -> playback speed {speed}x" )); } } } #[cfg(test)] mod tests { use super::*; #[test] fn deep_buffer_speeds_up_and_drained_buffer_returns_to_realtime() { assert!(catchup_speed(1.5, NORMAL_SPEED) > NORMAL_SPEED); assert_eq!(catchup_speed(0.1, MAX_CATCHUP_SPEED), NORMAL_SPEED); } #[test] fn drain_rate_scales_with_how_far_behind_we_are() { // The point of the proportional law: a small excursion gets a gentle // nudge, a deep backlog gets real recovery. let small = catchup_speed(1.5, NORMAL_SPEED); let large = catchup_speed(4.0, NORMAL_SPEED); assert!( large > small, "deeper buffer must drain faster: {small} vs {large}" ); assert!( (small - 1.05).abs() < 1e-9, "1.5s buffer -> 1.05x, got {small}" ); } #[test] fn drain_rate_is_capped_so_it_never_sounds_absurd() { // The ~6 s standing buffer measured on the netem rig, and far worse. assert_eq!(catchup_speed(6.0, NORMAL_SPEED), MAX_CATCHUP_SPEED); assert_eq!(catchup_speed(600.0, NORMAL_SPEED), MAX_CATCHUP_SPEED); } #[test] fn hysteresis_band_holds_the_current_speed() { // Between the marks nothing changes, whichever side we came from — // this is what stops the rate (and audio pitch) oscillating. for cache in [CACHE_LOW_S, 0.7, CACHE_HIGH_S] { assert_eq!(catchup_speed(cache, NORMAL_SPEED), NORMAL_SPEED); assert_eq!(catchup_speed(cache, MAX_CATCHUP_SPEED), MAX_CATCHUP_SPEED); } } #[test] fn unknown_cache_holds_the_current_speed() { assert_eq!( catchup_speed(f64::NAN, MAX_CATCHUP_SPEED), MAX_CATCHUP_SPEED ); assert_eq!(catchup_speed(f64::INFINITY, NORMAL_SPEED), NORMAL_SPEED); } #[test] fn a_full_handover_cycle_drains_then_settles() { // Buffer grows through a loss burst, then drains as we play faster. let mut speed = NORMAL_SPEED; for cache in [0.2, 0.5, 1.2, 3.4, 1.4, 0.9, 0.6, 0.3, 0.2] { speed = catchup_speed(cache, speed); } assert_eq!( speed, NORMAL_SPEED, "should be back at realtime once drained" ); } #[test] fn requests_are_valid_json_with_their_ids() { let get: serde_json::Value = serde_json::from_str(&get_cache_request(7)).unwrap(); assert_eq!(get["request_id"], 7); assert_eq!(get["command"][0], "get_property"); assert_eq!(get["command"][1], CACHE_PROPERTY); let set: serde_json::Value = serde_json::from_str(&set_speed_request(8, 1.05)).unwrap(); assert_eq!(set["request_id"], 8); assert_eq!(set["command"][0], "set_property"); assert_eq!(set["command"][1], "speed"); assert_eq!(set["command"][2], 1.05); } #[test] fn parses_our_reply_only() { assert_eq!( parse_cache_response(r#"{"error":"success","data":1.25,"request_id":3}"#, 3), Some(Some(1.25)) ); // Another command's reply, and an unsolicited event, are not ours. assert_eq!( parse_cache_response(r#"{"error":"success","data":1.25,"request_id":4}"#, 3), None ); assert_eq!( parse_cache_response(r#"{"event":"playback-restart"}"#, 3), None ); assert_eq!(parse_cache_response("not json", 3), None); } #[test] fn reply_without_a_usable_number_is_ours_but_empty() { // mpv before playback starts, and while the demuxer has no cache. assert_eq!( parse_cache_response(r#"{"error":"success","data":null,"request_id":1}"#, 1), Some(None) ); assert_eq!( parse_cache_response(r#"{"error":"property unavailable","request_id":1}"#, 1), Some(None) ); } #[test] fn socket_path_is_scoped_to_its_token() { let a = socket_path(Path::new("/run/user/1000"), 42); assert_eq!(a, Path::new("/run/user/1000/peerspeak-mpv-42.sock")); assert_ne!(a, socket_path(Path::new("/run/user/1000"), 43)); } }