use std::sync::{Arc, Condvar, Mutex};
use std::time::Duration;
use frust_core::anim::FrameTime;
#[cfg(feature = "perf-trace")]
use crate::perf::MarkerQueue;
use crate::perf::UiSpans;
pub const NO_RENDER_THREAD_VAR: &str = "FRUST_NO_RENDER_THREAD";
pub fn render_thread_enabled() -> bool {
!render_thread_kill_switch(
option_env!("FRUST_NO_RENDER_THREAD"),
std::env::var(NO_RENDER_THREAD_VAR).ok().as_deref(),
)
}
fn render_thread_kill_switch(compile_time: Option<&str>, runtime: Option<&str>) -> bool {
fn is_set_non_zero(value: Option<&str>) -> bool {
matches!(value, Some(v) if v != "0")
}
is_set_non_zero(compile_time) || is_set_non_zero(runtime)
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct SurfaceSize {
pub width: u32,
pub height: u32,
pub scale: f64,
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct FrameMeta {
pub frame_time: FrameTime,
pub size: SurfaceSize,
pub frame_id: u64,
}
#[derive(Debug, Clone)]
pub struct SceneFrame<S> {
pub scene: S,
pub meta: FrameMeta,
pub ui_spans: UiSpans,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RenderEvent {
SurfaceCreated,
SurfaceChanged,
SurfaceDestroyed,
Pause,
Resume,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RenderPhase {
NoSurface,
Active,
Paused,
}
impl RenderPhase {
pub fn can_render(self) -> bool {
matches!(self, RenderPhase::Active)
}
}
pub fn next_render_phase(current: RenderPhase, event: RenderEvent) -> RenderPhase {
match event {
RenderEvent::SurfaceCreated => RenderPhase::Active,
RenderEvent::SurfaceDestroyed => RenderPhase::NoSurface,
RenderEvent::SurfaceChanged => current,
RenderEvent::Pause => match current {
RenderPhase::NoSurface => RenderPhase::NoSurface,
_ => RenderPhase::Paused,
},
RenderEvent::Resume => match current {
RenderPhase::NoSurface => RenderPhase::NoSurface,
_ => RenderPhase::Active,
},
}
}
#[derive(Debug)]
pub struct Ack {
shared: Arc<AckShared>,
}
#[derive(Debug)]
pub struct AckWaiter {
shared: Arc<AckShared>,
}
#[derive(Debug)]
struct AckShared {
done: Mutex<bool>,
signal: Condvar,
}
pub fn ack_pair() -> (AckWaiter, Ack) {
let shared = Arc::new(AckShared {
done: Mutex::new(false),
signal: Condvar::new(),
});
(
AckWaiter {
shared: shared.clone(),
},
Ack { shared },
)
}
impl Ack {
pub fn acknowledge(self) {
}
}
impl Drop for Ack {
fn drop(&mut self) {
let mut done = self.shared.done.lock().unwrap();
*done = true;
drop(done);
self.shared.signal.notify_all();
}
}
impl AckWaiter {
pub fn wait(self) {
let mut done = self.shared.done.lock().unwrap();
while !*done {
done = self.shared.signal.wait(done).unwrap();
}
}
#[must_use = "the caller must handle a timeout (proceed degraded) rather than assume the barrier was honored"]
pub fn wait_timeout(self, timeout: Duration) -> bool {
let done = self.shared.done.lock().unwrap();
let (done, _timeout_result) = self
.shared
.signal
.wait_timeout_while(done, timeout, |done| !*done)
.unwrap();
*done
}
pub fn completed(&self) -> bool {
*self.shared.done.lock().unwrap()
}
}
#[derive(Debug)]
pub enum RenderCommand<W> {
SurfaceCreated {
window: W,
size: SurfaceSize,
},
SurfaceChanged {
size: SurfaceSize,
},
SurfaceDestroyed {
ack: Ack,
},
Pause {
ack: Ack,
},
Resume,
}
impl<W> RenderCommand<W> {
pub fn event(&self) -> RenderEvent {
match self {
RenderCommand::SurfaceCreated { .. } => RenderEvent::SurfaceCreated,
RenderCommand::SurfaceChanged { .. } => RenderEvent::SurfaceChanged,
RenderCommand::SurfaceDestroyed { .. } => RenderEvent::SurfaceDestroyed,
RenderCommand::Pause { .. } => RenderEvent::Pause,
RenderCommand::Resume => RenderEvent::Resume,
}
}
pub fn requires_ack(&self) -> bool {
matches!(
self,
RenderCommand::SurfaceDestroyed { .. } | RenderCommand::Pause { .. }
)
}
}
#[derive(Debug)]
struct Inbox<S, W> {
latest: Option<SceneFrame<S>>,
dropped: u64,
commands: Vec<RenderCommand<W>>,
#[cfg(feature = "perf-trace")]
pending_markers: MarkerQueue,
sender_alive: bool,
receiver_alive: bool,
}
#[derive(Debug)]
struct Channel<S, W> {
inbox: Mutex<Inbox<S, W>>,
signal: Condvar,
}
#[derive(Debug)]
pub struct RenderSender<S, W> {
channel: Arc<Channel<S, W>>,
}
#[derive(Debug)]
pub struct RenderReceiver<S, W> {
channel: Arc<Channel<S, W>>,
}
#[derive(Debug)]
pub struct RenderBatch<S, W> {
pub commands: Vec<RenderCommand<W>>,
pub scene: Option<SceneFrame<S>>,
pub disconnected: bool,
}
pub fn render_channel<S, W>() -> (RenderSender<S, W>, RenderReceiver<S, W>) {
let channel = Arc::new(Channel {
inbox: Mutex::new(Inbox {
latest: None,
dropped: 0,
commands: Vec::new(),
#[cfg(feature = "perf-trace")]
pending_markers: MarkerQueue::new(),
sender_alive: true,
receiver_alive: true,
}),
signal: Condvar::new(),
});
(
RenderSender {
channel: channel.clone(),
},
RenderReceiver { channel },
)
}
impl<S, W> RenderSender<S, W> {
pub fn send_scene(&self, frame: SceneFrame<S>) -> Option<SceneFrame<S>> {
let mut inbox = self.channel.inbox.lock().unwrap();
if !inbox.receiver_alive {
#[cfg(feature = "perf-trace")]
drop(crate::perf::take_pending_markers());
return Some(frame);
}
#[cfg(feature = "perf-trace")]
inbox
.pending_markers
.absorb(crate::perf::take_pending_markers());
let stale = inbox.latest.replace(frame);
if stale.is_some() {
inbox.dropped += 1;
}
drop(inbox);
self.channel.signal.notify_one();
stale
}
pub fn send_command(&self, command: RenderCommand<W>) {
let mut inbox = self.channel.inbox.lock().unwrap();
if !inbox.receiver_alive {
drop(inbox);
drop(command);
return;
}
inbox.commands.push(command);
drop(inbox);
self.channel.signal.notify_one();
}
#[must_use = "the caller must wait() on the returned AckWaiter before backgrounding"]
pub fn pause(&self) -> AckWaiter {
let (waiter, ack) = ack_pair();
self.send_command(RenderCommand::Pause { ack });
waiter
}
#[must_use = "the caller must wait() on the returned AckWaiter before releasing the window"]
pub fn destroy_surface(&self) -> AckWaiter {
let (waiter, ack) = ack_pair();
self.send_command(RenderCommand::SurfaceDestroyed { ack });
waiter
}
}
impl<S, W> Drop for RenderSender<S, W> {
fn drop(&mut self) {
let mut inbox = self.channel.inbox.lock().unwrap();
inbox.sender_alive = false;
drop(inbox);
self.channel.signal.notify_all();
}
}
impl<S, W> Drop for RenderReceiver<S, W> {
fn drop(&mut self) {
let mut inbox = self.channel.inbox.lock().unwrap();
inbox.receiver_alive = false;
let commands = std::mem::take(&mut inbox.commands);
let latest = inbox.latest.take();
#[cfg(feature = "perf-trace")]
inbox.pending_markers.clear();
drop(inbox);
drop(commands);
drop(latest);
self.channel.signal.notify_all();
}
}
impl<S, W> RenderReceiver<S, W> {
pub fn wait_next(&self) -> RenderBatch<S, W> {
let mut inbox = self.channel.inbox.lock().unwrap();
inbox = self
.channel
.signal
.wait_while(inbox, |i| {
i.commands.is_empty() && i.latest.is_none() && i.sender_alive
})
.unwrap();
drain(&mut inbox)
}
pub fn try_next(&self) -> RenderBatch<S, W> {
let mut inbox = self.channel.inbox.lock().unwrap();
drain(&mut inbox)
}
pub fn dropped_frames(&self) -> u64 {
self.channel.inbox.lock().unwrap().dropped
}
}
fn drain<S, W>(inbox: &mut Inbox<S, W>) -> RenderBatch<S, W> {
#[cfg(feature = "perf-trace")]
if inbox.latest.is_some() && !inbox.pending_markers.is_empty() {
crate::perf::stage_markers(inbox.pending_markers.take());
}
RenderBatch {
commands: std::mem::take(&mut inbox.commands),
scene: inbox.latest.take(),
disconnected: !inbox.sender_alive,
}
}
#[derive(Debug)]
pub struct SceneReturnSender<S> {
slot: Arc<Mutex<Option<S>>>,
}
#[derive(Debug)]
pub struct SceneReturnReceiver<S> {
slot: Arc<Mutex<Option<S>>>,
}
pub fn scene_return_channel<S>() -> (SceneReturnSender<S>, SceneReturnReceiver<S>) {
let slot = Arc::new(Mutex::new(None));
(
SceneReturnSender { slot: slot.clone() },
SceneReturnReceiver { slot },
)
}
impl<S> SceneReturnSender<S> {
pub fn give_back(&self, scene: S) {
let mut slot = self.slot.lock().unwrap();
*slot = Some(scene);
}
}
impl<S> SceneReturnReceiver<S> {
pub fn try_recv(&self) -> Option<S> {
self.slot.lock().unwrap().take()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicBool, Ordering};
#[cfg(feature = "perf-trace")]
use std::sync::mpsc;
use std::thread;
use std::time::{Duration, Instant};
#[test]
fn render_thread_kill_switch_off_when_neither_set() {
assert!(!render_thread_kill_switch(None, None));
}
#[test]
fn render_thread_kill_switch_on_when_compile_time_set_non_zero() {
assert!(render_thread_kill_switch(Some("1"), None));
}
#[test]
fn render_thread_kill_switch_on_when_runtime_set_non_zero() {
assert!(render_thread_kill_switch(None, Some("1")));
}
#[test]
fn render_thread_kill_switch_off_when_either_is_literal_zero_and_other_unset() {
assert!(!render_thread_kill_switch(Some("0"), None));
assert!(!render_thread_kill_switch(None, Some("0")));
}
#[test]
fn render_thread_kill_switch_on_when_either_source_wins() {
assert!(render_thread_kill_switch(Some("0"), Some("1")));
assert!(render_thread_kill_switch(Some("1"), Some("0")));
}
#[test]
fn render_thread_enabled_default_process_env_is_on() {
assert!(render_thread_enabled());
}
const ALL_PHASES: [RenderPhase; 3] = [
RenderPhase::NoSurface,
RenderPhase::Active,
RenderPhase::Paused,
];
#[test]
fn created_always_lands_in_active() {
for phase in ALL_PHASES {
assert_eq!(
next_render_phase(phase, RenderEvent::SurfaceCreated),
RenderPhase::Active,
"creating (or recreating) a surface must reach Active from {phase:?}"
);
}
}
#[test]
fn destroyed_always_lands_in_no_surface() {
for phase in ALL_PHASES {
assert_eq!(
next_render_phase(phase, RenderEvent::SurfaceDestroyed),
RenderPhase::NoSurface,
"destroying a surface must reach NoSurface from {phase:?}"
);
}
}
#[test]
fn changed_keeps_the_current_phase() {
for phase in ALL_PHASES {
assert_eq!(
next_render_phase(phase, RenderEvent::SurfaceChanged),
phase,
"a resize must not change whether we can render, from {phase:?}"
);
}
}
#[test]
fn pause_barriers_from_active_but_not_from_no_surface() {
assert_eq!(
next_render_phase(RenderPhase::Active, RenderEvent::Pause),
RenderPhase::Paused
);
assert_eq!(
next_render_phase(RenderPhase::Paused, RenderEvent::Pause),
RenderPhase::Paused
);
assert_eq!(
next_render_phase(RenderPhase::NoSurface, RenderEvent::Pause),
RenderPhase::NoSurface
);
}
#[test]
fn resume_reactivates_unless_the_surface_is_gone() {
assert_eq!(
next_render_phase(RenderPhase::Paused, RenderEvent::Resume),
RenderPhase::Active
);
assert_eq!(
next_render_phase(RenderPhase::Active, RenderEvent::Resume),
RenderPhase::Active
);
assert_eq!(
next_render_phase(RenderPhase::NoSurface, RenderEvent::Resume),
RenderPhase::NoSurface
);
}
#[test]
fn only_active_can_render() {
assert!(RenderPhase::Active.can_render());
assert!(!RenderPhase::NoSurface.can_render());
assert!(!RenderPhase::Paused.can_render());
}
#[test]
fn command_event_projection_matches_each_variant() {
let created: RenderCommand<()> = RenderCommand::SurfaceCreated {
window: (),
size: test_size(),
};
assert_eq!(created.event(), RenderEvent::SurfaceCreated);
assert!(!created.requires_ack());
let changed: RenderCommand<()> = RenderCommand::SurfaceChanged { size: test_size() };
assert_eq!(changed.event(), RenderEvent::SurfaceChanged);
assert!(!changed.requires_ack());
let resume: RenderCommand<()> = RenderCommand::Resume;
assert_eq!(resume.event(), RenderEvent::Resume);
assert!(!resume.requires_ack());
let (_w, ack) = ack_pair();
let paused: RenderCommand<()> = RenderCommand::Pause { ack };
assert_eq!(paused.event(), RenderEvent::Pause);
assert!(paused.requires_ack());
let (_w, ack) = ack_pair();
let destroyed: RenderCommand<()> = RenderCommand::SurfaceDestroyed { ack };
assert_eq!(destroyed.event(), RenderEvent::SurfaceDestroyed);
assert!(destroyed.requires_ack());
}
fn test_size() -> SurfaceSize {
SurfaceSize {
width: 1080,
height: 1920,
scale: 2.0,
}
}
fn frame(id: u32) -> SceneFrame<u32> {
SceneFrame {
scene: id,
meta: FrameMeta {
frame_time: FrameTime::from_nanos(u64::from(id)),
size: test_size(),
frame_id: u64::from(id),
},
ui_spans: UiSpans::default(),
}
}
#[test]
fn latest_wins_replaces_an_untaken_scene() {
let (tx, rx) = render_channel::<u32, ()>();
tx.send_scene(frame(1));
tx.send_scene(frame(2));
tx.send_scene(frame(3));
let batch = rx.try_next();
assert_eq!(batch.scene.expect("a scene is pending").scene, 3);
assert_eq!(rx.dropped_frames(), 2, "two stale frames were replaced");
let batch = rx.try_next();
assert!(batch.scene.is_none(), "the slot is drained after a take");
}
#[test]
fn taking_between_sends_drops_nothing() {
let (tx, rx) = render_channel::<u32, ()>();
tx.send_scene(frame(1));
assert_eq!(rx.try_next().scene.unwrap().scene, 1);
tx.send_scene(frame(2));
assert_eq!(rx.try_next().scene.unwrap().scene, 2);
assert_eq!(
rx.dropped_frames(),
0,
"each scene was taken before the next"
);
}
#[test]
fn send_scene_returns_none_when_the_slot_was_empty() {
let (tx, _rx) = render_channel::<u32, ()>();
assert!(
tx.send_scene(frame(1)).is_none(),
"the first send has nothing stale to hand back"
);
}
#[test]
fn send_scene_returns_the_overwritten_stale_frame() {
let (tx, rx) = render_channel::<u32, ()>();
assert!(tx.send_scene(frame(1)).is_none());
let stale = tx
.send_scene(frame(2))
.expect("the untaken frame(1) must be returned");
assert_eq!(stale.scene, 1, "the returned frame is the one replaced");
let batch = rx.try_next();
assert_eq!(batch.scene.expect("a scene is pending").scene, 2);
assert_eq!(rx.dropped_frames(), 1);
}
#[test]
fn send_scene_after_receiver_death_hands_the_frame_back() {
let (tx, rx) = render_channel::<u32, ()>();
drop(rx);
let returned = tx
.send_scene(frame(1))
.expect("a scene sent after receiver death must still be handed back");
assert_eq!(returned.scene, 1);
}
#[test]
fn scene_return_try_recv_is_none_before_any_give_back() {
let (_tx, rx) = scene_return_channel::<u32>();
assert!(rx.try_recv().is_none());
}
#[test]
fn scene_return_round_trips_a_single_scene() {
let (tx, rx) = scene_return_channel::<u32>();
tx.give_back(7);
assert_eq!(rx.try_recv(), Some(7), "the given-back scene is polled out");
assert!(
rx.try_recv().is_none(),
"the slot is drained after a take, like the forward channel"
);
}
#[test]
fn scene_return_is_depth_1_latest_wins() {
let (tx, rx) = scene_return_channel::<u32>();
tx.give_back(1);
tx.give_back(2);
assert_eq!(
rx.try_recv(),
Some(2),
"only the most recently given-back scene survives"
);
}
#[test]
fn scene_return_never_blocks_the_ui_thread() {
let (_tx, rx) = scene_return_channel::<u32>();
let start = Instant::now();
assert!(rx.try_recv().is_none());
assert!(
start.elapsed() < Duration::from_millis(50),
"try_recv must return immediately, never park"
);
}
#[test]
fn wait_next_blocks_until_a_scene_arrives() {
let (tx, rx) = render_channel::<u32, ()>();
let handle = thread::spawn(move || rx.wait_next().scene.map(|f| f.scene));
thread::sleep(Duration::from_millis(20));
tx.send_scene(frame(7));
assert_eq!(handle.join().unwrap(), Some(7));
}
#[test]
fn commands_drain_in_fifo_order_with_the_latest_scene() {
let (tx, rx) = render_channel::<u32, ()>();
tx.send_command(RenderCommand::SurfaceCreated {
window: (),
size: test_size(),
});
tx.send_scene(frame(1));
tx.send_command(RenderCommand::SurfaceChanged { size: test_size() });
tx.send_scene(frame(2));
let batch = rx.try_next();
let events: Vec<RenderEvent> = batch.commands.iter().map(RenderCommand::event).collect();
assert_eq!(
events,
vec![RenderEvent::SurfaceCreated, RenderEvent::SurfaceChanged],
"commands preserve FIFO order"
);
assert_eq!(
batch.scene.unwrap().scene,
2,
"only the freshest scene rides along"
);
assert!(!batch.disconnected);
}
#[test]
fn wait_next_wakes_and_reports_disconnection_when_sender_dropped() {
let (tx, rx) = render_channel::<u32, ()>();
let handle = thread::spawn(move || rx.wait_next().disconnected);
thread::sleep(Duration::from_millis(20));
drop(tx); assert!(
handle.join().unwrap(),
"dropping the sender wakes wait_next with a disconnection"
);
}
#[test]
fn pending_work_drains_before_disconnection_is_reported() {
let (tx, rx) = render_channel::<u32, ()>();
tx.send_scene(frame(5));
drop(tx);
let batch = rx.try_next();
assert_eq!(batch.scene.unwrap().scene, 5);
assert!(batch.disconnected);
}
#[test]
fn ack_unblocks_the_waiter_on_acknowledge() {
let (waiter, ack) = ack_pair();
assert!(!waiter.completed());
ack.acknowledge();
assert!(waiter.completed());
waiter.wait(); }
#[test]
fn ack_unblocks_the_waiter_on_drop_as_a_safety_net() {
let (waiter, ack) = ack_pair();
drop(ack); assert!(waiter.completed());
}
#[test]
fn pause_barrier_blocks_the_ui_thread_until_the_render_thread_quiesces() {
let (tx, rx) = render_channel::<u32, ()>();
let quiesced = Arc::new(AtomicBool::new(false));
let quiesced_render = quiesced.clone();
let render = thread::spawn(move || {
loop {
let batch = rx.wait_next();
for command in batch.commands {
if let RenderCommand::Pause { ack } = command {
thread::sleep(Duration::from_millis(30));
quiesced_render.store(true, Ordering::SeqCst);
ack.acknowledge();
return;
}
}
if batch.disconnected {
return;
}
}
});
let waiter = tx.pause();
waiter.wait();
assert!(
quiesced.load(Ordering::SeqCst),
"the render thread must have quiesced before the UI thread proceeded"
);
render.join().unwrap();
}
#[test]
fn destroy_surface_barrier_orders_resource_teardown_before_window_release() {
let (tx, rx) = render_channel::<u32, ()>();
let resources_dropped = Arc::new(AtomicBool::new(false));
let resources_dropped_render = resources_dropped.clone();
let render = thread::spawn(move || {
loop {
let batch = rx.wait_next();
for command in batch.commands {
if let RenderCommand::SurfaceDestroyed { ack } = command {
thread::sleep(Duration::from_millis(30));
resources_dropped_render.store(true, Ordering::SeqCst);
ack.acknowledge();
return;
}
}
if batch.disconnected {
return;
}
}
});
let waiter = tx.destroy_surface();
waiter.wait();
assert!(
resources_dropped.load(Ordering::SeqCst),
"surface resources must be dropped before the UI thread releases the window"
);
render.join().unwrap();
}
#[test]
fn a_leftover_scene_is_not_rendered_after_a_pause() {
let (tx, rx) = render_channel::<u32, ()>();
tx.send_scene(frame(1)); let _waiter = tx.pause();
let mut phase = RenderPhase::Active;
let mut submitted: Option<u32> = None;
let batch = rx.try_next();
for command in batch.commands {
phase = next_render_phase(phase, command.event());
if let RenderCommand::Pause { ack } = command {
ack.acknowledge();
}
}
if phase.can_render() {
submitted = batch.scene.map(|f| f.scene);
}
assert_eq!(phase, RenderPhase::Paused);
assert!(
submitted.is_none(),
"a leftover scene must not be submitted after a Pause"
);
}
const REGRESSION_DEADLINE: Duration = Duration::from_secs(5);
#[test]
fn receiver_drop_drains_a_queued_orphaned_ack() {
let (tx, rx) = render_channel::<u32, ()>();
let waiter = tx.pause();
drop(rx);
assert!(
waiter.wait_timeout(REGRESSION_DEADLINE),
"dropping the receiver must drain the queued Pause's Ack so the UI waiter unblocks"
);
}
#[test]
fn command_sent_after_receiver_death_fires_its_ack() {
let (tx, rx) = render_channel::<u32, ()>();
drop(rx);
let waiter = tx.destroy_surface();
assert!(
waiter.wait_timeout(REGRESSION_DEADLINE),
"an ack-carrying command sent after receiver death must fire its Ack safety net"
);
}
#[test]
fn receiver_death_mid_flight_unblocks_a_waiting_ui_thread() {
let (tx, rx) = render_channel::<u32, ()>();
let (report_tx, report_rx) = std::sync::mpsc::channel();
let ui = thread::spawn(move || {
let honored = tx.destroy_surface().wait_timeout(REGRESSION_DEADLINE);
report_tx.send(honored).unwrap();
});
thread::sleep(Duration::from_millis(20));
drop(rx);
let honored = report_rx
.recv_timeout(Duration::from_secs(10))
.expect("the UI thread must report back, not stay deadlocked");
assert!(
honored,
"receiver drop must unblock a UI thread already waiting on the barrier"
);
ui.join().unwrap();
}
#[test]
fn a_scene_sent_after_receiver_death_is_dropped_not_queued() {
let (tx, rx) = render_channel::<u32, ()>();
drop(rx);
tx.send_scene(frame(1)); }
#[test]
fn wait_timeout_returns_true_when_acknowledged() {
let (waiter, ack) = ack_pair();
ack.acknowledge();
assert!(
waiter.wait_timeout(REGRESSION_DEADLINE),
"an acknowledged barrier must report honored"
);
}
#[test]
fn wait_timeout_expires_false_when_never_acknowledged() {
let (waiter, _ack) = ack_pair();
assert!(
!waiter.wait_timeout(Duration::from_millis(20)),
"wait_timeout must return false when the ack never fires"
);
}
#[cfg(feature = "perf-trace")]
fn bench_stats() -> crate::perf::FrameStats {
crate::perf::FrameStats::with_capacity_enabled_and_raw(8, true, true)
}
#[cfg(feature = "perf-trace")]
fn record_one(stats: &mut crate::perf::FrameStats) {
stats.record(crate::perf::FramePasses::from_split(
UiSpans::default(),
crate::perf::RenderSpans::default(),
));
}
#[cfg(feature = "perf-trace")]
struct Lockstep {
to_peer: mpsc::Sender<()>,
from_peer: mpsc::Receiver<()>,
}
#[cfg(feature = "perf-trace")]
impl Lockstep {
fn signal(&self) {
self.to_peer
.send(())
.expect("the peer thread must still be running");
}
fn wait(&self) {
self.from_peer
.recv_timeout(Duration::from_secs(5))
.expect("the peer thread must reach its next step");
}
}
#[cfg(feature = "perf-trace")]
fn lockstep_pair() -> (Lockstep, Lockstep) {
let (ui_tx, ui_rx) = mpsc::channel();
let (render_tx, render_rx) = mpsc::channel();
(
Lockstep {
to_peer: ui_tx,
from_peer: render_rx,
},
Lockstep {
to_peer: render_tx,
from_peer: ui_rx,
},
)
}
#[cfg(feature = "perf-trace")]
#[test]
fn markers_land_on_the_frames_that_carried_them_across_the_split() {
let _guard = crate::perf::marker_test_guard(true);
let (tx, rx) = render_channel::<u32, ()>();
let mut stats = bench_stats();
let (ui_step, render_step) = lockstep_pair();
thread::scope(|scope| {
scope.spawn(move || {
let _ui_guard = crate::perf::marker_test_guard(true);
crate::perf::mark_scenario_start("s3-create1k");
tx.send_scene(frame(1));
ui_step.signal();
ui_step.wait();
crate::perf::mark_scenario_end("s3-create1k");
tx.send_scene(frame(2));
ui_step.signal();
});
render_step.wait();
let batch1 = rx.try_next();
assert!(batch1.scene.is_some(), "the render thread took build k");
render_step.signal();
render_step.wait();
record_one(&mut stats);
let batch2 = rx.try_next();
assert!(batch2.scene.is_some(), "the render thread took build k+1");
record_one(&mut stats);
});
assert_eq!(
stats.marker_log(),
[
"bench-scenario-start n=1 s3-create1k",
"bench-scenario-end n=2 s3-create1k",
],
"each edge names the recorded frame that actually carried it"
);
}
#[cfg(feature = "perf-trace")]
#[test]
fn a_window_whose_scene_is_replaced_collapses_onto_the_first_recorded_frame() {
let _guard = crate::perf::marker_test_guard(true);
let (tx, rx) = render_channel::<u32, ()>();
let mut stats = bench_stats();
let (ui_step, render_step) = lockstep_pair();
thread::scope(|scope| {
scope.spawn(move || {
let _ui_guard = crate::perf::marker_test_guard(true);
crate::perf::mark_scenario_start("s3-update");
tx.send_scene(frame(1));
crate::perf::mark_scenario_end("s3-update");
let stale = tx.send_scene(frame(2));
assert!(stale.is_some(), "build k's scene was replaced, not taken");
ui_step.signal();
});
render_step.wait();
let batch = rx.try_next();
assert_eq!(batch.scene.map(|f| f.scene), Some(2));
record_one(&mut stats);
});
assert_eq!(
stats.marker_log(),
[
"bench-scenario-start n=1 s3-update",
"bench-scenario-end n=1 s3-update",
]
);
assert_eq!(rx.dropped_frames(), 1);
}
#[cfg(feature = "perf-trace")]
#[test]
fn a_marker_raised_after_the_handoff_never_lands_on_the_frame_in_flight() {
let _guard = crate::perf::marker_test_guard(true);
let (tx, rx) = render_channel::<u32, ()>();
let mut stats = bench_stats();
let (ui_step, render_step) = lockstep_pair();
thread::scope(|scope| {
scope.spawn(move || {
let _ui_guard = crate::perf::marker_test_guard(true);
tx.send_scene(frame(1));
ui_step.signal();
ui_step.wait();
crate::perf::mark_scenario_start("s3-update");
ui_step.signal();
ui_step.wait();
tx.send_scene(frame(2));
ui_step.signal();
});
render_step.wait();
assert!(rx.try_next().scene.is_some(), "build k crossed");
render_step.signal();
render_step.wait();
record_one(&mut stats);
assert!(
stats.marker_log().is_empty(),
"a marker still queued on the UI thread cannot reach frame 1"
);
render_step.signal();
render_step.wait();
assert!(rx.try_next().scene.is_some(), "build k+1 crossed");
record_one(&mut stats);
});
assert_eq!(stats.marker_log(), ["bench-scenario-start n=2 s3-update"]);
}
#[cfg(feature = "perf-trace")]
#[test]
fn markers_wait_in_the_inbox_until_a_scene_is_actually_taken() {
let _guard = crate::perf::marker_test_guard(true);
let (tx, rx) = render_channel::<u32, ()>();
let mut stats = bench_stats();
let (ui_step, render_step) = lockstep_pair();
thread::scope(|scope| {
scope.spawn(move || {
let _ui_guard = crate::perf::marker_test_guard(true);
tx.send_scene(frame(1));
ui_step.signal();
ui_step.wait();
crate::perf::mark_scenario_start("s3-swap");
tx.send_command(RenderCommand::SurfaceChanged { size: test_size() });
ui_step.signal();
ui_step.wait();
tx.send_scene(frame(2));
ui_step.signal();
});
render_step.wait();
assert!(rx.try_next().scene.is_some());
record_one(&mut stats);
assert!(stats.marker_log().is_empty(), "frame 1 carried no marker");
render_step.signal();
render_step.wait();
let batch = rx.try_next();
assert!(
batch.scene.is_none(),
"a command-only wakeup takes no scene"
);
record_one(&mut stats);
assert!(
stats.marker_log().is_empty(),
"no scene was taken, so nothing may be attributed to frame 2"
);
render_step.signal();
render_step.wait();
assert!(rx.try_next().scene.is_some());
record_one(&mut stats);
});
assert_eq!(stats.marker_log(), ["bench-scenario-start n=3 s3-swap"]);
}
#[cfg(feature = "perf-trace")]
#[test]
fn a_frame_carrying_no_marker_emits_no_marker_line() {
let _guard = crate::perf::marker_test_guard(true);
let (tx, rx) = render_channel::<u32, ()>();
let mut stats = bench_stats();
tx.send_scene(frame(1));
assert!(rx.try_next().scene.is_some());
record_one(&mut stats);
assert!(stats.marker_log().is_empty());
assert_eq!(
stats.total_frames(),
1,
"the frame itself is still recorded"
);
}
#[cfg(feature = "perf-trace")]
#[test]
fn receiver_death_discards_markers_no_frame_will_ever_name() {
let _guard = crate::perf::marker_test_guard(true);
let (tx, rx) = render_channel::<u32, ()>();
let (next_tx, next_rx) = render_channel::<u32, ()>();
let mut stats = bench_stats();
let (ui_step, render_step) = lockstep_pair();
thread::scope(|scope| {
scope.spawn(move || {
let _ui_guard = crate::perf::marker_test_guard(true);
crate::perf::mark_scenario_start("s3-clear");
tx.send_scene(frame(1));
assert!(
crate::perf::take_pending_markers().is_empty(),
"send_scene moved the marker out of this thread's queue"
);
ui_step.signal();
ui_step.wait();
crate::perf::mark_scenario_end("s3-clear");
assert!(
tx.send_scene(frame(2)).is_some(),
"a send to a dead receiver hands the frame straight back"
);
assert!(
crate::perf::take_pending_markers().is_empty(),
"the post-mortem marker was discarded, not left queued"
);
next_tx.send_scene(frame(3));
ui_step.signal();
});
render_step.wait();
drop(rx);
render_step.signal();
render_step.wait();
assert!(next_rx.try_next().scene.is_some());
record_one(&mut stats);
});
assert!(
stats.marker_log().is_empty(),
"nothing leaked from the dead channel into the next one"
);
}
#[cfg(feature = "perf-trace")]
#[test]
fn a_marker_raised_on_a_thread_that_hands_no_frame_off_is_never_emitted() {
let _guard = crate::perf::marker_test_guard(true);
let (tx, rx) = render_channel::<u32, ()>();
let mut stats = bench_stats();
let (ui_step, render_step) = lockstep_pair();
thread::scope(|scope| {
scope.spawn(move || {
let _ui_guard = crate::perf::marker_test_guard(true);
thread::scope(|inner| {
inner.spawn(|| {
let _worker_guard = crate::perf::marker_test_guard(true);
crate::perf::mark_scenario_start("s4-parse");
});
});
tx.send_scene(frame(1));
ui_step.signal();
});
render_step.wait();
assert!(rx.try_next().scene.is_some());
record_one(&mut stats);
});
assert!(
stats.marker_log().is_empty(),
"the worker thread's marker belongs to no handed-off frame"
);
}
#[cfg(feature = "perf-trace")]
#[test]
fn an_undrained_inbox_stops_growing_at_the_queue_cap() {
let _guard = crate::perf::marker_test_guard(true);
let (tx, rx) = render_channel::<u32, ()>();
let mut stats = bench_stats();
let (ui_step, render_step) = lockstep_pair();
let per_send = crate::perf::MARKER_QUEUE_CAP * 3 / 4;
thread::scope(|scope| {
scope.spawn(move || {
let _ui_guard = crate::perf::marker_test_guard(true);
for send in 0..2u32 {
for i in 0..per_send {
crate::perf::mark_scenario_start(&format!("s3-op{send}-{i}"));
}
tx.send_scene(frame(send + 1));
}
ui_step.signal();
});
render_step.wait();
assert!(rx.try_next().scene.is_some());
record_one(&mut stats);
});
let dropped = per_send * 2 - crate::perf::MARKER_QUEUE_CAP;
let log = stats.marker_log();
assert_eq!(
log.len(),
crate::perf::MARKER_QUEUE_CAP + 1,
"a capped inbox's worth of markers, plus one overflow notice"
);
assert_eq!(
log[0],
format!("bench-scenario-start n=1 s3-op0-{dropped}"),
"the oldest edges were the ones dropped"
);
assert_eq!(
log[crate::perf::MARKER_QUEUE_CAP],
format!("frust-perf marker-overflow n=1 dropped={dropped}")
);
}
}