use super::{
dispatch::dispatch_next_chunk,
poll::{apply_actions, live_remote_allowed, sync_served_by_to_ui},
state::{Orchestrator, lock},
};
use crate::{
MainWindow, ViewportModel,
bridge::{
frame_cache::guide_pass::GuideCache,
remote::{
handoff::HandoffEvent,
live_lane::{ChunkVerdict, LiveLane},
remote_render::RemoteUpdate,
},
render_thread::{RenderContext, tonemap_running_average},
},
gui::{remote::worker_callbacks::backend_label, show_toast},
settings::LiveComputeTarget,
};
use indicatrix::geometry::tool::StoneGeometry;
use indicatrix_net::client::Accumulator;
use slint::{ComponentHandle, Weak};
use std::{
sync::{Arc, Mutex, PoisonError},
time::{Duration, Instant},
};
use super::super::generation::{
DenoiseGenerationJob, adopt_ready_denoise, adopt_ready_guides, spawn_denoise_generation,
};
const REMOTE_REDRAW_MIN_INTERVAL: Duration = Duration::from_millis(33);
#[must_use]
fn redraw_is_due(last_redraw_at: Option<Instant>, min_interval: Duration, now: Instant) -> bool {
last_redraw_at.is_none_or(|t| now.duration_since(t) >= min_interval)
}
fn is_combining(render_ctx: &Arc<Mutex<RenderContext>>) -> bool {
matches!(
render_ctx
.lock()
.unwrap_or_else(PoisonError::into_inner)
.effective_live_target(),
LiveComputeTarget::Both
)
}
pub(super) fn handle_remote_update(
ui_weak: &Weak<MainWindow>,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
chunk: (u32, u32, &Arc<Mutex<Accumulator>>),
update: RemoteUpdate,
) {
let (width, height, chunk_accumulator) = chunk;
let chunk_accumulator = Arc::clone(chunk_accumulator);
let render_ctx = render_ctx.clone();
let state = Arc::clone(state);
let wants_gated_redraw = matches!(
update,
RemoteUpdate::Frame { .. }
| RemoteUpdate::Preview { .. }
| RemoteUpdate::DisplayFrame { .. }
);
if wants_gated_redraw {
if is_combining(&render_ctx) {
if let RemoteUpdate::Frame {
request_id,
samples_done,
} = update
&& let Some(lane) = lock(&state).live_lane.as_mut()
{
lane.observe_progress(request_id, samples_done, Instant::now());
}
return;
}
let orch = lock(&state);
let due = redraw_is_due(
orch.last_redraw_at,
REMOTE_REDRAW_MIN_INTERVAL,
Instant::now(),
);
if !due || orch.redraw_gate.submit(()).is_none() {
return;
}
drop(orch);
}
let ui_weak_for_next = ui_weak.clone();
let _ = ui_weak.upgrade_in_event_loop(move |ui| {
if wants_gated_redraw {
lock(&state).redraw_gate.take();
}
if lock(&state).current_request_id != Some(update.request_id()) {
return;
}
apply_update(
&UpdateCtx {
ui: &ui,
ui_weak: &ui_weak_for_next,
render_ctx: &render_ctx,
state: &state,
chunk_accumulator: &chunk_accumulator,
width,
height,
},
update,
);
});
}
struct UpdateCtx<'a> {
ui: &'a MainWindow,
ui_weak: &'a Weak<MainWindow>,
render_ctx: &'a Arc<Mutex<RenderContext>>,
state: &'a Arc<Mutex<Orchestrator>>,
chunk_accumulator: &'a Mutex<Accumulator>,
width: u32,
height: u32,
}
fn apply_update(c: &UpdateCtx<'_>, update: RemoteUpdate) {
let (ui, render_ctx, state) = (c.ui, c.render_ctx, c.state);
match update {
RemoteUpdate::Connected { info, .. } => {
let actions = lock(state)
.handoff
.handle(HandoffEvent::RemoteStreamStarted);
apply_actions(&actions, render_ctx, state);
{
let mut s = lock(state);
s.worker_label = backend_label(info.render.as_ref());
s.remote_hdr = Some(info.render.as_ref().is_some_and(|r| r.hdr));
}
sync_served_by_to_ui(ui, render_ctx, state);
}
RemoteUpdate::Frame {
request_id,
samples_done,
} => {
tracing::trace!("remote chunk {request_id}: {samples_done} samples done");
if let Some(lane) = lock(state).live_lane.as_mut() {
lane.observe_progress(request_id, samples_done, Instant::now());
}
redraw_from_epoch(ui, render_ctx, state, c.width, c.height);
}
RemoteUpdate::Preview { .. } => {
redraw_from_epoch(ui, render_ctx, state, c.width, c.height);
}
RemoteUpdate::Progress {
request_id,
samples_done,
} => {
if let Some(lane) = lock(state).live_lane.as_mut() {
lane.observe_progress(request_id, samples_done, Instant::now());
}
}
RemoteUpdate::Done {
request_id,
cancelled,
} => {
if cancelled {
on_chunk_paused(state, request_id);
} else {
if lane_is_display_only(state) {
show_display_frame(ui, state, c.chunk_accumulator, c.width, c.height);
}
on_chunk_done(ui, c.ui_weak, render_ctx, state, request_id);
}
}
RemoteUpdate::Failed {
request_id,
message,
} => on_chunk_failed(ui, c.ui_weak, render_ctx, state, request_id, &message),
RemoteUpdate::Unsupported {
request_id,
message,
} => {
if lane_is_display_only(state) {
on_display_only_refused(ui, render_ctx, state, &message);
} else {
on_chunk_failed(ui, c.ui_weak, render_ctx, state, request_id, &message);
}
}
RemoteUpdate::DisplayFrame {
request_id,
samples_done,
} => {
tracing::trace!("remote display frame {request_id}: {samples_done} samples");
show_display_frame(ui, state, c.chunk_accumulator, c.width, c.height);
}
RemoteUpdate::CapabilityChanged { render, .. } => {
{
let mut s = lock(state);
s.worker_label = backend_label(render.as_ref());
s.remote_hdr = Some(render.as_ref().is_some_and(|r| r.hdr));
}
sync_served_by_to_ui(ui, render_ctx, state);
}
RemoteUpdate::FinalImage { .. } => {}
}
}
fn lane_is_display_only(state: &Arc<Mutex<Orchestrator>>) -> bool {
lock(state)
.live_lane
.as_ref()
.is_some_and(LiveLane::is_display_only)
}
fn display_frame_rgba(
accumulator: &Accumulator,
width: u32,
height: u32,
) -> Result<Option<Vec<u8>>, String> {
let Some(frame) = accumulator.last_display_frame() else {
return Ok(None);
};
if (frame.width, frame.height) != (width, height) {
return Err(format!(
"a {}x{} display frame for a {width}x{height} view",
frame.width, frame.height
));
}
indicatrix_net::display::decode_rgba8(frame.encoding, width, height, &frame.bytes)
.map(Some)
.map_err(|e| e.to_string())
}
fn show_display_frame(
ui: &MainWindow,
state: &Arc<Mutex<Orchestrator>>,
chunk_accumulator: &Mutex<Accumulator>,
width: u32,
height: u32,
) {
let decoded = display_frame_rgba(
&chunk_accumulator
.lock()
.unwrap_or_else(PoisonError::into_inner),
width,
height,
);
match decoded {
Ok(Some(rgba)) => {
lock(state).last_redraw_at = Some(Instant::now());
push_image(ui, width, height, &rgba);
}
Ok(None) => {}
Err(message) => tracing::warn!("dropping a remote display frame: {message}"),
}
}
fn on_display_only_refused(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
message: &str,
) {
{
let mut s = lock(state);
s.display_only_refused = true;
s.remote_handle = None;
if let Some(mut lane) = s.live_lane.take() {
lane.abandon();
}
s.lane_scene = None;
}
render_ctx
.lock()
.unwrap_or_else(PoisonError::into_inner)
.release_remote();
let actions = lock(state).handoff.handle(HandoffEvent::RemoteFailed);
apply_actions(&actions, render_ctx, state);
show_toast(
ui,
&format!(
"The remote cannot send finished live frames ({message}); using full data \
instead."
),
"info",
);
}
pub(super) fn continue_lane(
ui: &MainWindow,
ui_weak: &Weak<MainWindow>,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) {
if live_remote_allowed(render_ctx) && dispatch_next_chunk(ui_weak, render_ctx, state) {
return;
}
let finished = lock(state)
.live_lane
.as_ref()
.is_some_and(LiveLane::is_finished);
if finished {
finish_epoch(ui, render_ctx, state);
}
}
fn on_chunk_done(
ui: &MainWindow,
ui_weak: &Weak<MainWindow>,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
request_id: u32,
) {
let merged = lock(state)
.live_lane
.as_mut()
.and_then(|lane| lane.chunk_done(request_id, Instant::now()));
if merged.is_none() {
return;
}
lock(state).remote_handle = None;
continue_lane(ui, ui_weak, render_ctx, state);
}
fn on_chunk_paused(state: &Arc<Mutex<Orchestrator>>, request_id: u32) {
let merged = lock(state)
.live_lane
.as_mut()
.and_then(|lane| lane.chunk_paused(request_id));
if merged.is_none() {
return;
}
lock(state).remote_handle = None;
}
fn on_chunk_failed(
ui: &MainWindow,
ui_weak: &Weak<MainWindow>,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
request_id: u32,
message: &str,
) {
let verdict = lock(state)
.live_lane
.as_mut()
.and_then(|lane| lane.chunk_failed(request_id));
lock(state).remote_handle = None;
match verdict {
None => {}
Some(ChunkVerdict::Continue) => {
tracing::warn!("remote chunk {request_id} failed ({message}); retrying");
continue_lane(ui, ui_weak, render_ctx, state);
}
Some(ChunkVerdict::GaveUp) => {
let mode = render_ctx
.lock()
.unwrap_or_else(PoisonError::into_inner)
.live_compute_target;
{
let mut s = lock(state);
s.gave_up_for = s
.live_lane
.as_ref()
.map(|lane| (lane.epoch().scene_generation(), mode));
}
let actions = lock(state).handoff.handle(HandoffEvent::RemoteFailed);
apply_actions(&actions, render_ctx, state);
if is_combining(render_ctx) {
show_toast(
ui,
&format!(
"Remote worker failed twice in a row ({message}); finishing this \
image locally."
),
"info",
);
} else {
{
let mut s = lock(state);
s.live_lane = None;
s.lane_scene = None;
}
render_ctx
.lock()
.unwrap_or_else(PoisonError::into_inner)
.release_remote();
show_toast(ui, &format!("Remote render failed: {message}"), "error");
}
}
}
}
fn finish_epoch(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) {
let Some((width, height)) = lock(state)
.live_lane
.as_ref()
.map(|lane| lane.epoch().dimensions())
else {
return;
};
redraw_from_epoch(ui, render_ctx, state, width, height);
let actions = lock(state).handoff.handle(HandoffEvent::RemoteDone);
apply_actions(&actions, render_ctx, state);
sync_served_by_to_ui(ui, render_ctx, state);
lock(state).remote_handle = None;
}
pub(super) fn redraw_from_epoch(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
width: u32,
height: u32,
) {
if is_combining(render_ctx) {
return;
}
let Some(epoch) = lock(state)
.live_lane
.as_ref()
.filter(|lane| !lane.is_display_only())
.map(|lane| Arc::clone(lane.epoch()))
else {
return;
};
if epoch.dimensions() != (width, height) {
return;
}
let (buffer, remote_count) = epoch.remote_snapshot();
let samples_done = remote_count.max(1);
let (yaw, pitch, distance, planes, tools, denoise_enabled, n_d) = {
let ctx = render_ctx.lock().unwrap_or_else(PoisonError::into_inner);
(
ctx.yaw,
ctx.pitch,
ctx.distance,
ctx.active_planes.clone(),
ctx.active_tools.clone(),
ctx.denoise_enabled,
ctx.active_n_d(),
)
};
let stone = StoneGeometry {
planes: &planes,
tools: &tools,
};
let desired_key = GuideCache::key_for_geom(width, height, yaw, pitch, distance, stone, n_d);
let cached_denoised = {
let mut orch = lock(state);
if denoise_enabled {
if let Some(fresh) =
adopt_ready_denoise(&desired_key, orch.pending_denoise_gen.as_ref())
{
orch.pending_denoise_gen = None;
orch.last_denoised = Some((desired_key.clone(), fresh));
}
match orch.last_denoised.as_ref() {
Some((key, bytes)) if *key == desired_key => Some(bytes.clone()),
_ => {
orch.last_denoised = None;
None
}
}
} else {
orch.pending_denoise_gen = None;
orch.last_denoised = None;
None
}
};
let bytes = cached_denoised
.unwrap_or_else(|| tonemap_running_average(width, height, samples_done, &buffer));
{
let mut orch = lock(state);
if denoise_enabled {
let pending_matches = orch
.pending_denoise_gen
.as_ref()
.is_some_and(|p| p.key == desired_key);
if !pending_matches {
let Orchestrator {
guide_cache,
pending_guide_gen,
..
} = &mut *orch;
if adopt_ready_guides(&desired_key, guide_cache, pending_guide_gen.as_ref()) {
let guides = guide_cache
.ensure_geom(width, height, yaw, pitch, distance, stone, n_d)
.clone();
orch.pending_denoise_gen = Some(spawn_denoise_generation(
desired_key,
DenoiseGenerationJob {
width,
height,
samples_done,
buffer,
guides,
yaw,
pitch,
distance,
planes: planes.as_ref().clone(),
tools: tools.as_ref().clone(),
n_d,
},
));
}
}
}
orch.last_redraw_at = Some(Instant::now());
}
push_image(ui, width, height, &bytes);
}
fn push_image(ui: &MainWindow, width: u32, height: u32, rgba: &[u8]) {
let mut fb = crate::bridge::pixel_buffer::FramebufferTransfer::new(width, height);
let image = fb.copy_from_gpu_slice(rgba);
ui.global::<ViewportModel>()
.set_render_image(slint::Image::from_rgba8(image));
ui.global::<ViewportModel>().set_has_render(true);
}
#[cfg(test)]
mod tests {
use super::*;
use indicatrix_net::messages::{DisplayEncoding, DisplayFrameHeader, StreamEvent};
fn with_display_frame(width: u32, height: u32, rgba: &[u8]) -> Accumulator {
let png = indicatrix_net::display::encode_rgba8(DisplayEncoding::Png, width, height, rgba)
.expect("encodes");
let mut acc = Accumulator::new(width, height);
acc.begin_request(1);
let header = DisplayFrameHeader {
request_id: 1,
samples_done: 32,
width,
height,
encoding: DisplayEncoding::Png,
payload_len: png.len() as u32,
};
acc.apply(&StreamEvent::DisplayFrame(header), Some(&png))
.expect("applies");
acc
}
#[test]
fn a_display_frame_decodes_to_exactly_the_pixels_the_remote_sent() {
let rgba: Vec<u8> = (0..3 * 2 * 4).map(|i| (i * 7 % 256) as u8).collect();
let acc = with_display_frame(3, 2, &rgba);
assert_eq!(display_frame_rgba(&acc, 3, 2), Ok(Some(rgba)));
}
#[test]
fn no_display_frame_yet_shows_nothing_and_a_mis_sized_one_is_refused() {
assert_eq!(display_frame_rgba(&Accumulator::new(2, 2), 2, 2), Ok(None));
let acc = with_display_frame(2, 2, &[5_u8; 16]);
assert!(
display_frame_rgba(&acc, 4, 4).is_err(),
"a frame for another view size never reaches the viewport"
);
}
#[test]
fn no_prior_redraw_is_always_due() {
assert!(redraw_is_due(
None,
REMOTE_REDRAW_MIN_INTERVAL,
Instant::now()
));
}
#[test]
fn a_redraw_within_the_interval_is_not_due() {
let last = Instant::now();
let now = last + Duration::from_millis(10);
assert!(
!redraw_is_due(Some(last), REMOTE_REDRAW_MIN_INTERVAL, now),
"10ms after the last redraw is well inside the ~33ms interval"
);
}
#[test]
fn a_redraw_exactly_at_the_interval_boundary_is_due() {
let last = Instant::now();
let now = last + REMOTE_REDRAW_MIN_INTERVAL;
assert!(
redraw_is_due(Some(last), REMOTE_REDRAW_MIN_INTERVAL, now),
"the boundary itself (elapsed >= min_interval) must count as due"
);
}
#[test]
fn a_redraw_well_past_the_interval_is_due() {
let last = Instant::now();
let now = last + REMOTE_REDRAW_MIN_INTERVAL + Duration::from_secs(1);
assert!(redraw_is_due(Some(last), REMOTE_REDRAW_MIN_INTERVAL, now));
}
#[test]
fn a_burst_of_updates_within_one_interval_finds_at_most_one_due() {
let start = Instant::now();
let mut last_redraw_at: Option<Instant> = None;
let mut due_count = 0;
for ms in [0u64, 5, 10, 15, 20, 25, 30] {
let now = start + Duration::from_millis(ms);
if redraw_is_due(last_redraw_at, REMOTE_REDRAW_MIN_INTERVAL, now) {
due_count += 1;
last_redraw_at = Some(now);
}
}
assert_eq!(
due_count, 1,
"only the FIRST update of a burst inside one ~33ms interval may redraw"
);
let now = start + REMOTE_REDRAW_MIN_INTERVAL + Duration::from_millis(1);
assert!(redraw_is_due(
last_redraw_at,
REMOTE_REDRAW_MIN_INTERVAL,
now
));
}
}