pixel-change-check-client 0.1.2

Replicates your screen exactly and sends only the pixels that changed, over QUIC or a relay, to a native or browser viewer.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
//! The native viewer: connect, authenticate, reconstruct, present.
//!
//! The state machine is `pcc::Compositor`, so this file holds only the
//! parts that are genuinely about presenting: choosing a transport,
//! telling the user what is happening, and drawing. Two rules the UI
//! follows that the old code got wrong:
//!
//! * a failed session ends the window loop instead of leaving it spinning
//!   on a screen that will never change;
//! * a resize re-reads the surface dimensions, so a sharer that changes
//!   geometry does not leave a viewer rendering into a stale buffer.

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 {
    /// `host:port` of a sharer listening for direct connections.
    pub connect: Option<String>,
    /// `host:port` of a relay, for NAT'd peers.
    pub relay: Option<String>,
    /// Relay session code, required with `relay`.
    pub session: Option<String>,
    /// SHA-256 fingerprint of the sharer's certificate. Required with
    /// `connect`: this is what makes the connection authenticated.
    pub pin: String,
    pub token: SessionToken,
    /// Open a native window; otherwise print periodic status.
    pub show_window: bool,
    /// Reconnect after a failure instead of exiting.
    pub reconnect: bool,
}

/// How many video frames the playout clock may hold while waiting for
/// audio time to reach them.
///
/// Receipt: at 30 fps this is a second of video, far more than any
/// healthy session needs, and still bounded -- past it the oldest frame is
/// dropped, because a stale frame is worth less than a fresh one.
const PLAYOUT_DEPTH: usize = 30;

/// What the presentation layer needs to know, shared with the receive
/// task.
/// Why a session stopped, so a failed connection exits non-zero rather
/// than looking like a clean run to whatever is driving the CLI.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Stop {
    /// The sharer ended the session normally.
    Clean,
    /// The session could not be established or was lost.
    Failed,
}

#[derive(Default)]
struct Surface {
    frame: Option<Arc<Frame>>,
    terminal: Option<(String, Stop)>,
}

/// Entry point for `pcc view`. The network loop runs on a background
/// tokio task; the window loop stays on the calling thread because GUI
/// toolkits generally require the main thread on macOS.
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 {
        // Reconnect as a new synchronisation epoch: every attempt starts
        // from a clean compositor, so a half-received surface from the
        // previous attempt can never be mistaken for a valid one.
        loop {
            {
                let mut s = bg_surface.lock();
                s.frame = None;
                s.terminal = None;
            }
            // The connection is rebuilt per attempt, so it is captured
            // here rather than hoisted out of the reconnect loop.
            let outcome = receive_once(&view_args, &bg_surface, &metrics_bg).await;
            *bg_attempts.lock() += 1;
            // Every exit path must release the presentation loop, or a
            // cleanly-ended session leaves it spinning forever.
            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) => {
                    // The whole chain: a bare "handshake failed" is not
                    // something anyone can act on.
                    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;
        }
        // Wake the presentation loop even if it was waiting on a frame
        // that will never arrive.
        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());
    }

    // A session that never connected is a failure, and a script driving
    // this must be able to tell.
    let outcome = surface.lock().terminal.clone();
    match outcome {
        Some((message, Stop::Failed)) if !args.reconnect => Err(anyhow::anyhow!(message)),
        _ => Ok(()),
    }
}

/// A sealed session presented as one transport.
///
/// It cannot be split a second time: the two directions share nothing,
/// but each already owns its own counter, and splitting would hand the
/// caller halves that are individually useless.
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"
    );

    // Kept only so the audio stream can be accepted; the transport owns
    // its own handle.
    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?;
        // QUIC streams are not visible to the peer until data flows, so
        // the handshake kick is what makes `accept_bi` return. Sending
        // Hello does double duty: it authenticates and it opens the
        // stream.
        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>"
        );
    };

    // End-to-end encryption, before anything else is exchanged. The token
    // is the pre-shared key, so there is no second credential to manage
    // and a wrong token fails here rather than after a snapshot has been
    // decoded.
    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) => {
            // A dropped stream loses the sharer's reason, because a
            // QUIC send stream that is not finished is reset rather
            // than flushed. The viewer knows what it just sent, so it
            // says the useful thing: the token was refused.
            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,
        )),
        // The viewer *reads* what the host seals, so the source opens
        // with the host-to-viewer direction. Getting this backwards fails
        // authentication on the first frame, which is at least a loud
        // failure rather than a silent downgrade.
        source: Box::new(crate::network::SealedSource::new(
            source,
            session.host_to_viewer,
        )),
    });

    // Audio arrives on its own unidirectional stream, so a multi-megabyte
    // snapshot can never delay a 20 ms frame. It is accepted only when the
    // sharer opened one; the magic header is what identifies it, because
    // quinn 0.10 has no typed streams.
    let mut audio_stream: Option<quinn::RecvStream> = None;
    // The two clocks, and the bridge between them. Until the estimator
    // converges, video plays as soon as it arrives rather than being
    // scheduled against a plausible-but-wrong offset.
    let mut sync = crate::audio::sync::Estimator::new();
    let mut sharer_origin: Option<std::time::Instant> = None;
    // The output device lives on its own thread, so nothing here holds a
    // `!Send` stream across an await.
    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 {
            // quinn 0.10 has no typed streams, so the header written by
            // `open_stream` is what marks this one as audio. It is read
            // from the stream itself: a unidirectional stream is
            // read-only and cannot be split.
            let mut probe = stream;
            match crate::audio::AudioReceiver::read_header(&mut probe).await {
                Ok(true) => {
                    match crate::audio::AudioReceiver::new() {
                        Ok(receiver) => {
                            // Video is released against the audio clock,
                            // never the other way round: the audio device
                            // has a buffer and a latency nobody controls,
                            // so its time defines "now".
                            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;
    // With audio, a frame waits for the audio clock rather than being
    // shown the instant it decodes. Without audio there is no clock, and
    // the frame is shown immediately.
    let mut pending: Option<Frame> = None;

    loop {
        // Drain whatever audio has arrived, then take one visual message.
        // A blocked visual stream must not stop audio, and vice versa.
        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))) => {
                    // The first audio frame pins the sharer's clock to
                    // this process's clock; every later one refines it.
                    let origin = *sharer_origin.get_or_insert_with(std::time::Instant::now);
                    sync.observe(crate::audio::sync::OffsetSample::new(
                        (origin + frame.pts()).elapsed(),
                        // No RTT measurement of our own here; the
                        // estimator's clamp keeps a bad sample from
                        // poisoning the playout.
                        std::time::Duration::from_millis(0),
                    ));
                    out.push(&frame.pcm);
                }
                Ok(Ok(None)) => {}
                Ok(Err(e)) => warn!("Audio frame refused: {e}"),
                Err(_) => {}
            }
        }
        // Release whatever the audio clock says is due.
        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() {
                        // Hold it: the audio clock decides when it is due.
                        pending = Some(frame);
                    } else {
                        surface.lock().frame = Some(Arc::new(frame));
                    }
                    // The gap between joining and this frame is the single
                    // number a new user cares about most.
                    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() {
                                // Hold it: the audio clock decides when
                                // it is due, never the arrival.
                                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(_)) => {
                        // The compositor knows precisely what is missing.
                        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 {
                    // We have a snapshot but the sharer says it is
                    // further ahead; the next update will carry it.
                }
            }
            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 { .. } => {
                // Viewer-to-sharer messages; a sharer sending one back is
                // a protocol error, not something to act on.
                warn!("Ignoring a viewer-only message from the sharer");
            }
        }
    }
}

fn server_name_of(target: &str) -> &str {
    // The certificate is validated by fingerprint, so the SNI only has to
    // be syntactically valid. The host part is the natural choice.
    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);

    // The window can be resized; the surface is re-read every time its
    // dimensions change, so a sharer that changes geometry does not leave
    // this loop drawing into a stale buffer.
    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…"),
        }
    }
}

/// Wait for a frame, optionally giving up after a deadline. The old code
/// looped here forever, so a failed connection left the app hung with no
/// way out and no explanation.
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));
    }
}

/// The address a viewer is pointed at, resolved once so the headless log
/// and error messages can name something concrete.
pub fn describe(target: &str) -> String {
    match target.parse::<SocketAddr>() {
        Ok(a) => a.to_string(),
        Err(_) => target.to_string(),
    }
}