96a35ca84c
ci / rust (push) Failing after 29s
ci / web (push) Failing after 35s
docker / build-push (., web/Dockerfile, punktfunk-web) (push) Successful in 5s
docker / build-push (ci, ci/rust-ci.Dockerfile, punktfunk-rust-ci) (push) Successful in 3s
docker / build-push (docs-site, docs-site/Dockerfile, punktfunk-docs) (push) Successful in 18s
ci / docs-site (push) Failing after 38s
apple / swift (push) Successful in 1m15s
docker / deploy-docs (push) Successful in 17s
New crate crates/punktfunk-client-linux (binary punktfunk-client), the native Linux client on the Option A architecture (2026-06-12 research): - GTK4/libadwaita shell linking punktfunk-core directly (no C ABI): mDNS host list, TOFU fingerprint prompt, SPAKE2 PIN pairing dialog, preferences (mode/bitrate/gamepad/shortcut capture), stats overlay, --connect host[:port] for scripting. - Video: FFmpeg software HEVC decode (LOW_DELAY, slice threads) -> RGBA -> GdkMemoryTexture inside GtkGraphicsOffload (the dmabuf subsurface path lights up when VAAPI lands; black-background keeps fullscreen scanout-eligible). - Audio: Opus -> PipeWire playback stream, the host virtual-mic's adaptive jitter ring inverted. - Input: keyboard as the exact inverse of the host VK table (evdev keycodes, layout-independent; unit-tested), absolute mouse through the Contain-fit transform, WHEEL_DELTA(120) scroll, compositor shortcut inhibition while streaming, Ctrl+Alt+Shift+Q release chord, F11 fullscreen. SDL3 gamepad capture (single pad-0 model) + rumble and DualSense lightbar feedback on the same thread. - Session pump owns video+audio pulls; the gamepad thread owns rumble+hidout — possible because NativeClient's plane receivers are now mutexed, making it Sync (Arc-shared, compiler-verified per-plane contract instead of the ABI's manual assertion). - Linux-gated deps + a stub main keep cargo build --workspace green on macOS. Validated live against serve --native on this box: 1920x1080@60, locked 60 fps, capture->decoded p50 ~6.4 ms (software decode, debug build). Teardown keys off AdwNavigationPage::hidden — NavigationView push fires a transient unmap/map cycle that must not end the session. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
227 lines
8.1 KiB
Rust
227 lines
8.1 KiB
Rust
//! Session controller: one worker thread runs connect → pump (video pull + decode, audio
|
|
//! pull + Opus decode, stats), feeding the GTK main loop over channels. The UI keeps the
|
|
//! `Arc<NativeClient>` from the `Connected` event for direct input sends (no extra hop on
|
|
//! the input path) — `NativeClient` is `Sync`, planes stay one-consumer-per-thread:
|
|
//! video+audio here, rumble+hidout on the gamepad thread.
|
|
|
|
use crate::video::{DecodedFrame, Decoder};
|
|
use crate::{audio, gamepad};
|
|
use punktfunk_core::client::NativeClient;
|
|
use punktfunk_core::config::{CompositorPref, GamepadPref, Mode};
|
|
use punktfunk_core::PunktfunkError;
|
|
use std::sync::atomic::{AtomicBool, Ordering};
|
|
use std::sync::Arc;
|
|
use std::time::{Duration, Instant};
|
|
|
|
pub struct SessionParams {
|
|
pub host: String,
|
|
pub port: u16,
|
|
pub mode: Mode,
|
|
pub gamepad: GamepadPref,
|
|
pub bitrate_kbps: u32,
|
|
/// Pinned host fingerprint; `None` = trust on first use (caller persists the observed one).
|
|
pub pin: Option<[u8; 32]>,
|
|
pub identity: (String, String),
|
|
}
|
|
|
|
#[derive(Clone, Copy, Default)]
|
|
pub struct Stats {
|
|
pub fps: f32,
|
|
pub mbps: f32,
|
|
pub decode_ms: f32,
|
|
/// Median capture→decoded latency over the last window (host-clock corrected).
|
|
pub latency_ms: f32,
|
|
}
|
|
|
|
pub enum SessionEvent {
|
|
Connected {
|
|
connector: Arc<NativeClient>,
|
|
mode: Mode,
|
|
fingerprint: [u8; 32],
|
|
},
|
|
Failed(String),
|
|
Ended(Option<String>),
|
|
Stats(Stats),
|
|
}
|
|
|
|
pub struct SessionHandle {
|
|
pub events: async_channel::Receiver<SessionEvent>,
|
|
pub frames: async_channel::Receiver<DecodedFrame>,
|
|
pub stop: Arc<AtomicBool>,
|
|
}
|
|
|
|
pub fn start(params: SessionParams) -> SessionHandle {
|
|
let (ev_tx, ev_rx) = async_channel::unbounded();
|
|
// Tiny frame queue, newest wins: force_send displaces the oldest when the UI lags.
|
|
let (frame_tx, frame_rx) = async_channel::bounded(2);
|
|
let stop = Arc::new(AtomicBool::new(false));
|
|
let stop_w = stop.clone();
|
|
std::thread::Builder::new()
|
|
.name("punktfunk-session".into())
|
|
.spawn(move || pump(params, ev_tx, frame_tx, stop_w))
|
|
.expect("spawn session thread");
|
|
SessionHandle {
|
|
events: ev_rx,
|
|
frames: frame_rx,
|
|
stop,
|
|
}
|
|
}
|
|
|
|
fn now_ns() -> u64 {
|
|
std::time::SystemTime::now()
|
|
.duration_since(std::time::UNIX_EPOCH)
|
|
.map(|d| d.as_nanos() as u64)
|
|
.unwrap_or(0)
|
|
}
|
|
|
|
fn pump(
|
|
params: SessionParams,
|
|
ev_tx: async_channel::Sender<SessionEvent>,
|
|
frame_tx: async_channel::Sender<DecodedFrame>,
|
|
stop: Arc<AtomicBool>,
|
|
) {
|
|
let connector = match NativeClient::connect(
|
|
¶ms.host,
|
|
params.port,
|
|
params.mode,
|
|
CompositorPref::Auto,
|
|
params.gamepad,
|
|
params.bitrate_kbps,
|
|
params.pin,
|
|
Some(params.identity),
|
|
Duration::from_secs(15),
|
|
) {
|
|
Ok(c) => Arc::new(c),
|
|
Err(e) => {
|
|
let msg = match e {
|
|
PunktfunkError::Crypto => {
|
|
"Host identity rejected — wrong fingerprint, or the host requires pairing"
|
|
.to_string()
|
|
}
|
|
PunktfunkError::Timeout => "Connection timed out".to_string(),
|
|
other => format!("Connect failed: {other:?}"),
|
|
};
|
|
let _ = ev_tx.send_blocking(SessionEvent::Failed(msg));
|
|
return;
|
|
}
|
|
};
|
|
let _ = ev_tx.send_blocking(SessionEvent::Connected {
|
|
connector: connector.clone(),
|
|
mode: connector.mode(),
|
|
fingerprint: connector.host_fingerprint,
|
|
});
|
|
|
|
let mut decoder = match Decoder::new() {
|
|
Ok(d) => d,
|
|
Err(e) => {
|
|
let _ = ev_tx.send_blocking(SessionEvent::Ended(Some(format!("video decoder: {e}"))));
|
|
return;
|
|
}
|
|
};
|
|
// Audio and gamepads are best-effort: a session without them still streams.
|
|
let player = audio::AudioPlayer::spawn()
|
|
.map_err(|e| tracing::warn!(error = %e, "audio disabled"))
|
|
.ok();
|
|
let mut opus_dec = opus::Decoder::new(48_000, opus::Channels::Stereo)
|
|
.map_err(|e| tracing::warn!(error = %e, "opus decoder failed — audio disabled"))
|
|
.ok();
|
|
let gamepad_thread = gamepad::spawn(connector.clone(), stop.clone());
|
|
|
|
let clock_offset = connector.clock_offset_ns;
|
|
let mut total_frames = 0u64;
|
|
let mut window_start = Instant::now();
|
|
let mut frames_n = 0u32;
|
|
let mut bytes_n = 0u64;
|
|
let mut decode_us_sum = 0u64;
|
|
let mut lat_us: Vec<u64> = Vec::with_capacity(256);
|
|
let mut pcm = vec![0f32; 5760 * 2]; // decode scratch: max Opus frame (120 ms stereo)
|
|
|
|
let end: Option<String> = loop {
|
|
if stop.load(Ordering::SeqCst) {
|
|
break None;
|
|
}
|
|
match connector.next_frame(Duration::from_millis(4)) {
|
|
Ok(frame) => {
|
|
let t0 = Instant::now();
|
|
match decoder.decode(&frame.data) {
|
|
Ok(Some(decoded)) => {
|
|
total_frames += 1;
|
|
if total_frames == 1 {
|
|
tracing::info!(
|
|
width = decoded.width,
|
|
height = decoded.height,
|
|
"first frame decoded"
|
|
);
|
|
}
|
|
// Latency: our wall clock expressed in the host's capture clock,
|
|
// minus the host-stamped capture pts (same math as client-rs).
|
|
let lat = (now_ns() as i128 + clock_offset as i128 - frame.pts_ns as i128)
|
|
.max(0) as u64;
|
|
if lat > 0 && lat < 10_000_000_000 {
|
|
lat_us.push(lat / 1000);
|
|
}
|
|
decode_us_sum += t0.elapsed().as_micros() as u64;
|
|
frames_n += 1;
|
|
bytes_n += frame.data.len() as u64;
|
|
let _ = frame_tx.force_send(decoded);
|
|
}
|
|
Ok(None) => {}
|
|
// Survivable (loss until the next IDR/RFI recovery) — keep feeding.
|
|
Err(e) => tracing::debug!(error = %e, "decode error (recovering)"),
|
|
}
|
|
}
|
|
Err(PunktfunkError::NoFrame) => {}
|
|
Err(PunktfunkError::Closed) => break Some("Host ended the session".to_string()),
|
|
Err(e) => break Some(format!("session: {e:?}")),
|
|
}
|
|
|
|
// Drain audio between frames (packets land every 5 ms; the queue holds 320 ms).
|
|
while let Ok(pkt) = connector.next_audio(Duration::ZERO) {
|
|
if let (Some(player), Some(dec)) = (&player, opus_dec.as_mut()) {
|
|
match dec.decode_float(&pkt.data, &mut pcm, false) {
|
|
Ok(samples) => player.push(pcm[..samples * 2].to_vec()),
|
|
Err(e) => tracing::debug!(error = %e, "opus decode"),
|
|
}
|
|
}
|
|
}
|
|
|
|
if window_start.elapsed() >= Duration::from_secs(1) {
|
|
let secs = window_start.elapsed().as_secs_f32();
|
|
lat_us.sort_unstable();
|
|
let p50 = lat_us.get(lat_us.len() / 2).copied().unwrap_or(0);
|
|
tracing::debug!(
|
|
fps = frames_n,
|
|
lat_p50_us = p50,
|
|
total_frames,
|
|
"stream window"
|
|
);
|
|
let _ = ev_tx.try_send(SessionEvent::Stats(Stats {
|
|
fps: frames_n as f32 / secs,
|
|
mbps: bytes_n as f32 * 8.0 / 1e6 / secs,
|
|
decode_ms: if frames_n > 0 {
|
|
decode_us_sum as f32 / frames_n as f32 / 1000.0
|
|
} else {
|
|
0.0
|
|
},
|
|
latency_ms: p50 as f32 / 1000.0,
|
|
}));
|
|
window_start = Instant::now();
|
|
frames_n = 0;
|
|
bytes_n = 0;
|
|
decode_us_sum = 0;
|
|
lat_us.clear();
|
|
}
|
|
};
|
|
|
|
tracing::info!(
|
|
total_frames,
|
|
reason = end.as_deref().unwrap_or("user"),
|
|
"session ended"
|
|
);
|
|
stop.store(true, Ordering::SeqCst); // take the gamepad thread down with us
|
|
if let Some(t) = gamepad_thread {
|
|
let _ = t.join();
|
|
}
|
|
let _ = ev_tx.send_blocking(SessionEvent::Ended(end));
|
|
}
|