use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
use car_browser::{
ChromiumBackend, FrameReceiver, ScreencastFrame, ScreencastPump, FRAME_CHANNEL_CAP,
};
use tokio::sync::{mpsc, watch, Mutex};
const SUPERVISOR_STOP_GRACE: std::time::Duration = std::time::Duration::from_secs(5);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FrameAudience {
Viewer,
Model,
}
struct Consumer {
tx: mpsc::Sender<ScreencastFrame>,
audience: FrameAudience,
epoch: Instant,
}
fn publish(consumers: &mut Vec<Consumer>, frame: &ScreencastFrame, blackout: bool) -> usize {
let mut delivered = 0usize;
consumers.retain(|c| {
if blackout && c.audience == FrameAudience::Model {
return true;
}
let stamped = ScreencastFrame {
jpeg: frame.jpeg.clone(),
viewport: frame.viewport,
captured_at: c.epoch.elapsed().as_secs_f64(),
};
match c.tx.try_send(stamped) {
Ok(()) => {
delivered += 1;
true
}
Err(mpsc::error::TrySendError::Full(_)) => true,
Err(mpsc::error::TrySendError::Closed(_)) => false,
}
});
delivered
}
pub struct FrameFanout {
consumers: Arc<Mutex<Vec<Consumer>>>,
blackout: Arc<AtomicBool>,
backend: Mutex<Option<Arc<ChromiumBackend>>>,
supervisor: Arc<Mutex<Option<tokio::task::JoinHandle<()>>>>,
supervisor_generation: AtomicU64,
quality: AtomicI64,
idle: watch::Sender<u64>,
cancel: watch::Sender<u64>,
teardown: Mutex<()>,
width: u32,
height: u32,
}
impl FrameFanout {
pub fn new(width: u32, height: u32, default_quality: i64) -> Self {
Self {
consumers: Arc::new(Mutex::new(Vec::new())),
blackout: Arc::new(AtomicBool::new(false)),
backend: Mutex::new(None),
supervisor: Arc::new(Mutex::new(None)),
supervisor_generation: AtomicU64::new(0),
quality: AtomicI64::new(default_quality.clamp(1, 100)),
idle: watch::channel(0).0,
cancel: watch::channel(0).0,
teardown: Mutex::new(()),
width,
height,
}
}
pub fn set_blackout(&self, active: bool) {
self.blackout.store(active, Ordering::SeqCst);
}
pub fn blackout_active(&self) -> bool {
self.blackout.load(Ordering::SeqCst)
}
pub async fn bind(&self, backend: Arc<ChromiumBackend>) {
*self.backend.lock().await = Some(backend);
self.ensure_supervisor().await;
}
pub async fn subscribe(&self, audience: FrameAudience) -> (FrameReceiver, Instant) {
let (tx, rx) = mpsc::channel(FRAME_CHANNEL_CAP);
let epoch = Instant::now();
self.consumers.lock().await.push(Consumer {
tx,
audience,
epoch,
});
self.ensure_supervisor().await;
(rx, epoch)
}
pub fn quality(&self) -> i64 {
self.quality.load(Ordering::SeqCst)
}
pub async fn set_quality(&self, quality: i64) {
let quality = quality.clamp(1, 100);
if self.quality.swap(quality, Ordering::SeqCst) == quality {
return;
}
self.stop_supervisor().await;
self.ensure_supervisor().await;
}
pub async fn prune(&self) {
let empty = {
let mut consumers = self.consumers.lock().await;
consumers.retain(|c| !c.tx.is_closed());
consumers.is_empty()
};
if empty {
self.idle.send_modify(|n| *n = n.wrapping_add(1));
}
}
#[cfg(test)]
pub(crate) fn idle_watch_for_test(&self) -> watch::Receiver<u64> {
self.idle.subscribe()
}
#[cfg(test)]
pub async fn consumer_count_for_test(&self) -> usize {
self.consumers.lock().await.len()
}
async fn ensure_supervisor(&self) {
let mut running = self.supervisor.lock().await;
if running.as_ref().is_some_and(|t| !t.is_finished()) {
return;
}
let Some(backend) = self.backend.lock().await.clone() else {
return;
};
if self.consumers.lock().await.is_empty() {
return;
}
let consumers = Arc::clone(&self.consumers);
let blackout = Arc::clone(&self.blackout);
let idle = self.idle.subscribe();
let cancel = self.cancel.subscribe();
let quality = self.quality.load(Ordering::SeqCst);
let (width, height) = (self.width, self.height);
let slot = Arc::clone(&self.supervisor);
self.supervisor_generation.fetch_add(1, Ordering::AcqRel);
*running = Some(tokio::spawn(async move {
supervise(
backend, consumers, blackout, idle, cancel, quality, width, height, slot,
)
.await;
}));
}
async fn stop_supervisor(&self) {
let _teardown = self.teardown.lock().await;
if self
.supervisor
.lock()
.await
.as_ref()
.is_none_or(|t| t.is_finished())
{
return;
}
let generation = self.supervisor_generation.load(Ordering::Acquire);
self.cancel.send_modify(|n| *n = n.wrapping_add(1));
let deadline = Instant::now() + SUPERVISOR_STOP_GRACE;
loop {
{
let mut slot = self.supervisor.lock().await;
if self.supervisor_generation.load(Ordering::Acquire) != generation {
return;
}
match slot.as_ref() {
None => return,
Some(task) if task.is_finished() => {
*slot = None;
return;
}
Some(_) => {}
}
}
if Instant::now() >= deadline {
tracing::debug!(
"browser stream: supervisor did not stop within the grace period; aborting"
);
let mut slot = self.supervisor.lock().await;
if self.supervisor_generation.load(Ordering::Acquire) != generation {
return;
}
if let Some(task) = slot.take() {
drop(slot);
task.abort();
let _ = task.await;
}
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
}
async fn publish_exit(
slot: &Mutex<Option<tokio::task::JoinHandle<()>>>,
consumers: &Mutex<Vec<Consumer>>,
) -> bool {
let mut slot = slot.lock().await;
if !consumers.lock().await.is_empty() {
return false;
}
*slot = None;
true
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum RetryWake {
Retry,
Idle,
Cancelled,
}
async fn wait_for_retry(
tabs: &mut watch::Receiver<car_browser::TabsSnapshot>,
idle: &mut watch::Receiver<u64>,
cancel: &mut watch::Receiver<u64>,
) -> RetryWake {
tokio::select! {
changed = tabs.changed() => {
if changed.is_ok() { RetryWake::Retry } else { RetryWake::Idle }
}
_ = idle.changed() => RetryWake::Idle,
_ = cancel.changed() => RetryWake::Cancelled,
}
}
async fn supervise(
backend: Arc<ChromiumBackend>,
consumers: Arc<Mutex<Vec<Consumer>>>,
blackout: Arc<AtomicBool>,
mut idle: watch::Receiver<u64>,
mut cancel: watch::Receiver<u64>,
quality: i64,
width: u32,
height: u32,
slot: Arc<Mutex<Option<tokio::task::JoinHandle<()>>>>,
) {
let mut tabs = backend.subscribe_tabs();
loop {
if publish_exit(&slot, &consumers).await {
return;
}
let Ok(page) = backend.page_handle().await else {
match wait_for_retry(&mut tabs, &mut idle, &mut cancel).await {
RetryWake::Retry => continue,
RetryWake::Cancelled => return,
RetryWake::Idle => {
if publish_exit(&slot, &consumers).await {
return;
}
continue;
}
}
};
let attached_to = backend.active_tab_id();
let Ok(pump) = ScreencastPump::attach(&page, quality, 1, width, height).await else {
match wait_for_retry(&mut tabs, &mut idle, &mut cancel).await {
RetryWake::Retry => continue,
RetryWake::Cancelled => return,
RetryWake::Idle => {
if publish_exit(&slot, &consumers).await {
return;
}
continue;
}
}
};
let Ok((mut incoming, _started)) = pump.subscribe().await else {
match wait_for_retry(&mut tabs, &mut idle, &mut cancel).await {
RetryWake::Retry => continue,
RetryWake::Cancelled => return,
RetryWake::Idle => {
if publish_exit(&slot, &consumers).await {
return;
}
continue;
}
}
};
if consumers.lock().await.is_empty() {
pump.stop().await;
if publish_exit(&slot, &consumers).await {
return;
}
continue;
}
loop {
tokio::select! {
frame = incoming.recv() => {
let Some(frame) = frame else { break };
let blacked = blackout.load(Ordering::SeqCst);
let mut list = consumers.lock().await;
publish(&mut list, &frame, blacked);
if list.is_empty() {
drop(list);
pump.stop().await;
if publish_exit(&slot, &consumers).await {
return;
}
break;
}
}
changed = tabs.changed() => {
if changed.is_err() {
break;
}
if backend.active_tab_id() != attached_to {
break;
}
}
_ = idle.changed() => {
if consumers.lock().await.is_empty() {
pump.stop().await;
if publish_exit(&slot, &consumers).await {
return;
}
break;
}
}
_ = cancel.changed() => {
pump.stop().await;
return;
}
}
}
pump.stop().await;
}
}
#[cfg(test)]
mod tests {
use super::*;
use car_browser::models::Viewport;
fn frame(byte: u8) -> ScreencastFrame {
ScreencastFrame {
jpeg: vec![byte].into(),
viewport: Viewport {
width: 4,
height: 3,
device_pixel_ratio: 1.0,
},
captured_at: 0.0,
}
}
fn consumer(audience: FrameAudience) -> (Consumer, mpsc::Receiver<ScreencastFrame>) {
let (tx, rx) = mpsc::channel(FRAME_CHANNEL_CAP);
(
Consumer {
tx,
audience,
epoch: Instant::now(),
},
rx,
)
}
#[tokio::test]
async fn a_prune_that_lands_while_the_supervisor_is_attaching_is_not_lost() {
let fanout = FrameFanout::new(4, 3, 80);
let mut idle = fanout.idle_watch_for_test();
let (rx, _epoch) = fanout.subscribe(FrameAudience::Viewer).await;
drop(rx);
fanout.prune().await;
tokio::time::timeout(std::time::Duration::from_millis(100), idle.changed())
.await
.expect("the idle signal must still be there when the supervisor finally parks on it")
.expect("the fan-out is still alive");
assert_eq!(fanout.consumer_count_for_test().await, 0);
}
#[tokio::test]
async fn a_supervisor_tearing_down_does_not_exit_once_a_consumer_has_arrived() {
let slot: Mutex<Option<tokio::task::JoinHandle<()>>> =
Mutex::new(Some(tokio::spawn(async {
std::future::pending::<()>().await
})));
let consumers: Mutex<Vec<Consumer>> = Mutex::new(Vec::new());
let (c, _rx) = consumer(FrameAudience::Viewer);
consumers.lock().await.push(c);
assert!(
!publish_exit(&slot, &consumers).await,
"a registered consumer must keep the supervisor alive"
);
assert!(
slot.lock().await.is_some(),
"and the slot must still read as live, or nothing would restart it"
);
consumers.lock().await.clear();
assert!(publish_exit(&slot, &consumers).await);
assert!(
slot.lock().await.is_none(),
"the exiting task must publish its own absence, not leave a stale handle"
);
}
#[test]
fn every_consumer_receives_a_frame_when_no_blackout_is_active() {
let (viewer, mut viewer_rx) = consumer(FrameAudience::Viewer);
let (model, mut model_rx) = consumer(FrameAudience::Model);
let mut list = vec![viewer, model];
assert_eq!(publish(&mut list, &frame(7), false), 2);
assert_eq!(viewer_rx.try_recv().unwrap().jpeg.as_ref(), &[7]);
assert_eq!(model_rx.try_recv().unwrap().jpeg.as_ref(), &[7]);
}
#[test]
fn blackout_suppresses_the_recording_consumer_and_only_that_one() {
let (viewer, mut viewer_rx) = consumer(FrameAudience::Viewer);
let (model, mut model_rx) = consumer(FrameAudience::Model);
let mut list = vec![viewer, model];
assert_eq!(publish(&mut list, &frame(1), true), 1);
assert_eq!(
viewer_rx.try_recv().unwrap().jpeg.as_ref(),
&[1],
"the human driving must keep seeing the page"
);
assert!(
model_rx.try_recv().is_err(),
"no frame captured during a blackout may reach disk"
);
}
#[test]
fn a_gated_consumer_resumes_once_the_blackout_lifts() {
let (viewer, _viewer_rx) = consumer(FrameAudience::Viewer);
let (model, mut model_rx) = consumer(FrameAudience::Model);
let mut list = vec![viewer, model];
publish(&mut list, &frame(1), true);
publish(&mut list, &frame(2), true);
assert_eq!(list.len(), 2, "the gated consumer stays registered");
publish(&mut list, &frame(3), false);
assert_eq!(model_rx.try_recv().unwrap().jpeg.as_ref(), &[3]);
assert!(
model_rx.try_recv().is_err(),
"only the post-blackout frame lands"
);
}
#[test]
fn a_dropped_consumer_is_pruned_and_the_survivor_keeps_streaming() {
let (viewer, viewer_rx) = consumer(FrameAudience::Viewer);
let (other, mut other_rx) = consumer(FrameAudience::Viewer);
let mut list = vec![viewer, other];
drop(viewer_rx);
assert_eq!(publish(&mut list, &frame(5), false), 1);
assert_eq!(list.len(), 1);
assert_eq!(other_rx.try_recv().unwrap().jpeg.as_ref(), &[5]);
}
#[tokio::test]
async fn captured_at_is_stamped_per_consumer_epoch() {
const GAP: std::time::Duration = std::time::Duration::from_millis(30);
let (early, mut early_rx) = consumer(FrameAudience::Viewer);
let mut list = vec![early];
tokio::time::sleep(GAP).await;
let (late, mut late_rx) = consumer(FrameAudience::Viewer);
list.push(late);
publish(&mut list, &frame(1), false);
let early_at = early_rx.try_recv().unwrap().captured_at;
let late_at = late_rx.try_recv().unwrap().captured_at;
assert!(
early_at > late_at,
"the earlier subscriber sees a larger elapsed time ({early_at} vs {late_at})"
);
let gap = early_at - late_at;
assert!(
gap >= GAP.as_secs_f64() * 0.9,
"the two stamps must differ by the subscription gap, so the late \
subscriber is measuring its own epoch and not a shared one \
(early {early_at}, late {late_at}, difference {gap})"
);
}
#[tokio::test]
async fn subscribing_without_a_browser_yields_no_frames_and_no_supervisor() {
let fanout = FrameFanout::new(1920, 1080, 80);
let (mut rx, _epoch) = fanout.subscribe(FrameAudience::Viewer).await;
assert!(
fanout.supervisor.lock().await.is_none(),
"no browser bound yet, so nothing to capture"
);
assert!(rx.try_recv().is_err());
}
#[tokio::test]
async fn the_blackout_flag_round_trips() {
let fanout = FrameFanout::new(1920, 1080, 80);
assert!(!fanout.blackout_active());
fanout.set_blackout(true);
assert!(fanout.blackout_active());
fanout.set_blackout(false);
assert!(!fanout.blackout_active());
}
}