use crate::network::{
connect_direct, Message, MessageSink, MessageSource, MessageTransport, NetworkConfig, Rev,
SessionToken,
};
use crate::pcc::types::{Frame, BYTES_PER_PIXEL};
use crate::pcc::{ApplyError, Compositor};
use crate::relay::{RelayRole, RelayTransport};
use anyhow::{Context, Result};
use parking_lot::Mutex;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use tracing::{info, warn};
#[derive(Clone)]
pub struct ViewArgs {
pub connect: Option<String>,
pub relay: Option<String>,
pub session: Option<String>,
pub pin: String,
pub token: SessionToken,
pub show_window: bool,
pub reconnect: bool,
}
const PLAYOUT_DEPTH: usize = 30;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Stop {
Clean,
Failed,
}
#[derive(Default)]
struct Surface {
frame: Option<Arc<Frame>>,
terminal: Option<(String, Stop)>,
}
pub fn run_view(args: ViewArgs, metrics: crate::telemetry::SharedMetrics) -> Result<()> {
let rt = tokio::runtime::Runtime::new()?;
let surface: Arc<Mutex<Surface>> = Arc::new(Mutex::new(Surface::default()));
let bg_surface = surface.clone();
let metrics_bg = metrics.clone();
let attempts = Arc::new(Mutex::new(0usize));
let bg_attempts = attempts.clone();
let view_args = args.clone();
rt.spawn(async move {
loop {
{
let mut s = bg_surface.lock();
s.frame = None;
s.terminal = None;
}
let outcome = receive_once(&view_args, &bg_surface, &metrics_bg).await;
*bg_attempts.lock() += 1;
let (message, stop) = match outcome {
Ok(()) => {
let n = bg_attempts.lock();
let m = format!("the sharer ended the session after {n} attempt(s)");
info!("Viewer session ended: {m}");
(m, Stop::Clean)
}
Err(e) => {
let message = format!("{e:#}");
warn!("Viewer session ended: {message}");
(message, Stop::Failed)
}
};
bg_surface.lock().terminal = Some((message, stop));
if !view_args.reconnect {
break;
}
tokio::time::sleep(Duration::from_secs(2)).await;
}
if bg_surface.lock().terminal.is_none() {
bg_surface.lock().terminal = Some(("session ended".into(), Stop::Clean));
}
});
if args.show_window {
if let Err(e) = run_window_loop(surface.clone()) {
warn!("Falling back to headless mode (no window available): {e}");
run_headless_loop(surface.clone());
}
} else {
run_headless_loop(surface.clone());
}
let outcome = surface.lock().terminal.clone();
match outcome {
Some((message, Stop::Failed)) if !args.reconnect => Err(anyhow::anyhow!(message)),
_ => Ok(()),
}
}
struct SealedSession {
sink: Box<dyn MessageSink>,
source: Box<dyn MessageSource>,
}
#[async_trait::async_trait]
impl MessageTransport for SealedSession {
async fn send(&mut self, msg: &Message) -> Result<()> {
self.sink.send(msg).await
}
async fn send_encoded(&mut self, bytes: &[u8]) -> Result<()> {
self.sink.send_encoded(bytes).await
}
async fn recv(&mut self) -> Result<Message> {
self.source.recv().await
}
fn split(self: Box<Self>) -> (Box<dyn MessageSink>, Box<dyn MessageSource>) {
unreachable!("a sealed session is used whole")
}
}
#[async_trait::async_trait]
impl MessageSink for Box<dyn MessageSink> {
async fn send(&mut self, msg: &Message) -> Result<()> {
(**self).send(msg).await
}
async fn send_encoded(&mut self, bytes: &[u8]) -> Result<()> {
(**self).send_encoded(bytes).await
}
}
#[async_trait::async_trait]
impl MessageSource for Box<dyn MessageSource> {
async fn recv(&mut self) -> Result<Message> {
(**self).recv().await
}
async fn recv_raw(&mut self) -> Result<Vec<u8>> {
(**self).recv_raw().await
}
}
async fn receive_once(
args: &ViewArgs,
surface: &Arc<Mutex<Surface>>,
metrics: &crate::telemetry::SharedMetrics,
) -> Result<()> {
anyhow::ensure!(
!(args.connect.is_some() && args.relay.is_some()),
"specify either --connect or --relay, not both"
);
let quic: Option<quinn::Connection> = None;
let transport: Box<dyn MessageTransport> = if let Some(target) = &args.connect {
let addr = crate::network::resolve(target)
.await
.with_context(|| format!("Could not resolve --connect {target}"))?;
let pin = crate::network::hex_to_der(&args.pin)?;
info!("Connecting directly to {addr}");
let mut transport = connect_direct(
&NetworkConfig::default(),
&pin,
target,
server_name_of(target),
)
.await?;
transport
.send(&Message::Hello {
token: args.token.as_str().to_string(),
})
.await?;
Box::new(transport)
} else if let Some(target) = &args.relay {
let session = args
.session
.clone()
.context("--session <CODE> is required when using --relay")?;
let addr = crate::network::resolve(target)
.await
.with_context(|| format!("Could not resolve --relay {target}"))?;
let pin = crate::network::hex_to_der(&args.pin)?;
info!("Connecting via relay {addr}, session '{session}'");
let mut transport = RelayTransport::connect(
addr,
server_name_of(target),
&pin,
session,
args.token.clone(),
RelayRole::Viewer,
)
.await?;
transport
.send(&Message::Hello {
token: args.token.as_str().to_string(),
})
.await?;
Box::new(transport)
} else {
anyhow::bail!(
"Specify either --connect <host:port> or --relay <host:port> --session <code>"
);
};
let keys = crate::network::e2e::KeyPair::generate();
let (mut sink, mut source) = transport.split();
let session = match crate::network::e2e::viewer_handshake(
&mut sink,
&mut source,
keys,
args.token.as_str(),
)
.await
{
Ok(s) => s,
Err(e) => {
let detail = e.to_string();
anyhow::bail!(
"the sharer refused this session (token, or a different build?): {detail}"
);
}
};
let mut transport: Box<dyn MessageTransport> = Box::new(SealedSession {
sink: Box::new(crate::network::SealedSink::new(
sink,
session.viewer_to_host,
)),
source: Box::new(crate::network::SealedSource::new(
source,
session.host_to_viewer,
)),
});
let mut audio_stream: Option<quinn::RecvStream> = None;
let mut sync = crate::audio::sync::Estimator::new();
let mut sharer_origin: Option<std::time::Instant> = None;
let mut audio_out: Option<crate::audio::AudioOutput> = None;
let mut audio_played: Option<crate::audio::output::PlayedReceiver<u64>> = None;
let mut audio: Option<crate::audio::AudioReceiver> = None;
if let Some(connection) = &quic {
if let Ok(stream) = connection.accept_uni().await {
let mut probe = stream;
match crate::audio::AudioReceiver::read_header(&mut probe).await {
Ok(true) => {
match crate::audio::AudioReceiver::new() {
Ok(receiver) => {
match crate::audio::AudioOutput::start::<u64>(PLAYOUT_DEPTH) {
Ok((out, played)) => {
info!("Audio playout started");
audio_out = Some(out);
audio_played = Some(played);
}
Err(e) => warn!("Audio output unavailable: {e}"),
}
audio = Some(receiver);
audio_stream = Some(probe);
}
Err(e) => warn!("Audio decoder unavailable: {e}"),
}
}
Ok(false) => info!("Unidirectional stream was not audio; ignoring it"),
Err(e) => warn!("Audio stream header unreadable: {e}"),
}
}
}
info!("Connected. Waiting for frames...");
let joined_at = std::time::Instant::now();
let mut compositor = Compositor::new();
let mut last_ack: Rev = 0;
let mut pending: Option<Frame> = None;
loop {
if let (Some(receiver), Some(stream), Some(out)) =
(audio.as_mut(), audio_stream.as_mut(), audio_out.as_ref())
{
match tokio::time::timeout(
std::time::Duration::from_millis(1),
receiver.read_frame(stream),
)
.await
{
Ok(Ok(Some(frame))) => {
let origin = *sharer_origin.get_or_insert_with(std::time::Instant::now);
sync.observe(crate::audio::sync::OffsetSample::new(
(origin + frame.pts()).elapsed(),
std::time::Duration::from_millis(0),
));
out.push(&frame.pcm);
}
Ok(Ok(None)) => {}
Ok(Err(e)) => warn!("Audio frame refused: {e}"),
Err(_) => {}
}
}
if let Some(rx) = audio_played.as_ref() {
while let Ok(played) = rx.try_recv() {
if played
.revisions
.contains(&pending.as_ref().map(|f| f.id).unwrap_or(u64::MAX))
&& !played.revisions.is_empty()
{
if let Some(frame) = pending.take() {
surface.lock().frame = Some(Arc::new(frame));
}
}
}
}
let msg = transport.recv().await?;
match msg {
Message::SnapshotBegin {
rev: _,
pts_us: _,
epoch,
width,
height,
format,
total_len,
chunks,
} => {
compositor.begin_snapshot(epoch, width, height, format, total_len, chunks)?;
}
Message::SnapshotChunk {
rev: _,
index,
data,
} => {
compositor.push_snapshot_chunk(index, &data)?;
}
Message::SnapshotCommit { rev, pts_us, epoch } => {
let apply_start = std::time::Instant::now();
compositor.commit_snapshot(rev, epoch)?;
metrics.apply.record_duration(apply_start.elapsed());
if let Some((w, h)) = compositor.dimensions() {
let frame = Frame::with_pts(rev, w, h, compositor.buffer().to_vec(), pts_us)?;
if sync.convergence() && audio_played.is_some() {
pending = Some(frame);
} else {
surface.lock().frame = Some(Arc::new(frame));
}
metrics
.first_exact_image
.record_duration(joined_at.elapsed());
metrics.first_paint.record_duration(joined_at.elapsed());
info!("Snapshot installed: {w}x{h} at revision {rev}");
}
if rev != last_ack {
last_ack = rev;
transport.send(&Message::Ack { rev }).await?;
}
}
Message::PartialUpdate {
rev,
pts_us,
epoch,
ops,
} => {
let apply_start = std::time::Instant::now();
match compositor.apply_ops(rev, epoch, &ops) {
Ok(()) => {
metrics.apply.record_duration(apply_start.elapsed());
if let Some((w, h)) = compositor.dimensions() {
let frame =
Frame::with_pts(rev, w, h, compositor.buffer().to_vec(), pts_us)?;
if sync.convergence() && audio_played.is_some() {
pending = Some(frame);
} else {
surface.lock().frame = Some(Arc::new(frame));
}
}
if rev != last_ack {
last_ack = rev;
transport.send(&Message::Ack { rev }).await?;
}
}
Err(ApplyError::Rejected(_)) => {
metrics.repairs.incr();
transport.send(&Message::RequestKeyframe).await?;
}
Err(ApplyError::Invalid(e)) => {
return Err(e.context("The sharer sent an update this viewer cannot apply"));
}
}
}
Message::KeepAlive { rev } => {
if rev > last_ack && last_ack == 0 {
}
}
Message::QualityConfig(cfg) => {
info!(
"Sharer adjusted quality: target_fps={} quality={:.1}",
cfg.target_fps, cfg.quality
);
}
Message::Error(e) => return Err(anyhow::anyhow!("{e}")),
Message::Bye => {
info!("Sharer ended the session");
return Ok(());
}
Message::Hello { .. }
| Message::Ack { .. }
| Message::RequestKeyframe
| Message::E2eOffer { .. }
| Message::E2eReply { .. } => {
warn!("Ignoring a viewer-only message from the sharer");
}
}
}
}
fn server_name_of(target: &str) -> &str {
match target.rsplit_once(':') {
Some((host, _port)) if !host.is_empty() => host,
_ => "pcc",
}
}
fn run_window_loop(surface: Arc<Mutex<Surface>>) -> Result<()> {
use minifb::{Key, Window, WindowOptions};
info!("Waiting for the first frame to size the window...");
let first = wait_for_frame(&surface, None)?;
let (width, height) = (first.width, first.height);
let mut window = Window::new(
"PixelChangeCheck Viewer",
width as usize,
height as usize,
WindowOptions::default(),
)
.context("Failed to open a window (no display available?)")?;
window.set_target_fps(60);
let mut current = (width, height);
let mut argb = vec![0u32; (width * height) as usize];
while window.is_open() && !window.is_key_down(Key::Escape) {
let (frame, terminal) = {
let s = surface.lock();
(s.frame.clone(), s.terminal.clone())
};
if let Some((reason, _)) = terminal {
warn!("Viewer stopped: {reason}");
break;
}
let Some(frame) = frame else { continue };
if (frame.width, frame.height) != current {
argb = vec![0u32; (frame.width * frame.height) as usize];
current = (frame.width, frame.height);
info!("Surface resized to {}x{}", current.0, current.1);
}
let expected = (current.0 * current.1) as usize * BYTES_PER_PIXEL;
if frame.data.len() != expected {
continue;
}
for (i, px) in frame
.data
.as_chunks::<BYTES_PER_PIXEL>()
.0
.iter()
.enumerate()
{
argb[i] = ((px[0] as u32) << 16) | ((px[1] as u32) << 8) | px[2] as u32;
}
window.update_with_buffer(&argb, current.0 as usize, current.1 as usize)?;
}
Ok(())
}
fn run_headless_loop(surface: Arc<Mutex<Surface>>) {
info!("Running headless (no window). Press Ctrl+C to quit.");
loop {
std::thread::sleep(Duration::from_secs(2));
let s = surface.lock();
if let Some((reason, _)) = &s.terminal {
info!("Viewer stopped: {reason}");
return;
}
match &s.frame {
Some(f) => info!(
"Receiving {}x{} ({} bytes)",
f.width,
f.height,
f.data.len()
),
None => info!("Connected, waiting for the first snapshot…"),
}
}
}
fn wait_for_frame(
surface: &Arc<Mutex<Surface>>,
_deadline: Option<Duration>,
) -> Result<std::sync::Arc<Frame>> {
loop {
let s = surface.lock();
if let Some((reason, _)) = &s.terminal {
anyhow::bail!("{reason}");
}
if let Some(frame) = &s.frame {
return Ok(frame.clone());
}
drop(s);
std::thread::sleep(Duration::from_millis(50));
}
}
pub fn describe(target: &str) -> String {
match target.parse::<SocketAddr>() {
Ok(a) => a.to_string(),
Err(_) => target.to_string(),
}
}