use super::dispatch::{CANCEL_WAIT_TIMEOUT, liveness_deadline};
use crate::{
bridge::{
export_thread::{params::ComputeTarget, tonemap_png::tonemap_accumulation},
remote::remote_render::{self, RemoteFinalImageRequest, RemoteUpdate},
},
settings::{ExportTransfer, WorkerSettings},
};
use glam::Vec3;
use indicatrix::color::ColorSpace;
use indicatrix_net::{
SceneState,
client::{Accumulator, PreviewSnapshot},
};
use slint::{Rgba8Pixel, SharedPixelBuffer};
use std::{
collections::{BTreeMap, BTreeSet},
sync::{
Arc, Mutex, PoisonError,
atomic::{AtomicBool, AtomicU32, Ordering},
mpsc::{self, RecvTimeoutError},
},
time::{Duration, Instant},
};
const REQUEST_ID: u32 = 1;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(in crate::bridge::export_thread) enum TransferPlan {
FinalPicture,
FullData,
FullDataRefusedBefore,
}
#[must_use]
pub(in crate::bridge::export_thread) const fn plan_export_transfer(
requested: ExportTransfer,
compute_target: ComputeTarget,
remote_configured: bool,
hdr_refused: bool,
refused_before: bool,
) -> TransferPlan {
let wanted = matches!(requested, ExportTransfer::FinalPicture)
&& !matches!(compute_target, ComputeTarget::LocalOnly)
&& remote_configured
&& !hdr_refused;
match (wanted, refused_before) {
(false, _) => TransferPlan::FullData,
(true, true) => TransferPlan::FullDataRefusedBefore,
(true, false) => TransferPlan::FinalPicture,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(in crate::bridge::export_thread) enum FinalPictureOutcome {
Completed {
rgba: Vec<u8>,
reclaimed_samples: u32,
},
Cancelled,
Unsupported(String),
Failed(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(in crate::bridge::export_thread) enum FinalPictureFollowUp {
Use(Vec<u8>),
Cancelled,
FallBackToFullData { note: String, remember: bool },
Fail(String),
}
#[must_use]
pub(in crate::bridge::export_thread) fn final_picture_follow_up(
outcome: FinalPictureOutcome,
compute_target: ComputeTarget,
) -> FinalPictureFollowUp {
match outcome {
FinalPictureOutcome::Completed { rgba, .. } => FinalPictureFollowUp::Use(rgba),
FinalPictureOutcome::Cancelled => FinalPictureFollowUp::Cancelled,
FinalPictureOutcome::Unsupported(message) => FinalPictureFollowUp::FallBackToFullData {
note: format!(
"The remote does not support final-picture transfer ({message}); \
using full data instead."
),
remember: true,
},
FinalPictureOutcome::Failed(message) => {
if matches!(compute_target, ComputeTarget::RemoteOnly) {
FinalPictureFollowUp::Fail(format!("Remote final-picture render failed: {message}"))
} else {
FinalPictureFollowUp::FallBackToFullData {
note: format!(
"Remote final-picture render failed ({message}); continuing with \
full data."
),
remember: false,
}
}
}
}
}
static FINAL_PICTURE_REFUSED: Mutex<BTreeSet<String>> = Mutex::new(BTreeSet::new());
pub(in crate::bridge::export_thread) fn remember_final_picture_refused(worker: &WorkerSettings) {
FINAL_PICTURE_REFUSED
.lock()
.unwrap_or_else(PoisonError::into_inner)
.insert(worker.address.clone());
}
#[must_use]
pub(in crate::bridge::export_thread) fn final_picture_refused(worker: &WorkerSettings) -> bool {
FINAL_PICTURE_REFUSED
.lock()
.unwrap_or_else(PoisonError::into_inner)
.contains(&worker.address)
}
pub fn forget_final_picture_refusals() {
FINAL_PICTURE_REFUSED
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clear();
FINAL_PICTURE_RATES
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clear();
}
pub(in crate::bridge::export_thread) struct LocalShare<'a> {
pub samples: u32,
pub done: &'a AtomicU32,
pub result: mpsc::Receiver<(Vec<Vec3>, Duration)>,
pub stop: &'a AtomicBool,
}
pub(in crate::bridge::export_thread) const DEFAULT_VIEWER_SHARE: f64 = 0.10;
#[must_use]
pub(in crate::bridge::export_thread) fn viewer_share(
samples: u32,
local_rate: Option<f64>,
remote_rate: Option<f64>,
) -> u32 {
let half = samples / 2;
let fraction = match (local_rate, remote_rate) {
(Some(local), Some(remote))
if local.is_finite() && remote.is_finite() && local > 0.0 && remote > 0.0 =>
{
local / (local + remote)
}
_ => DEFAULT_VIEWER_SHARE,
};
let share = (f64::from(samples) * fraction).round().max(0.0) as u32;
share.min(half)
}
#[derive(Debug, Clone, Copy, Default)]
struct SplitRates {
local_px_rate: Option<f64>,
remote_px_rate: Option<f64>,
}
static FINAL_PICTURE_RATES: Mutex<BTreeMap<String, SplitRates>> = Mutex::new(BTreeMap::new());
#[must_use]
pub(in crate::bridge::export_thread) fn split_rates(
worker: &WorkerSettings,
pixels: u32,
) -> (Option<f64>, Option<f64>) {
let entry = FINAL_PICTURE_RATES
.lock()
.unwrap_or_else(PoisonError::into_inner)
.get(&worker.address)
.copied();
let Some(entry) = entry else {
return (None, None);
};
let denorm = |px_rate: Option<f64>| px_rate.map(|rate| rate / f64::from(pixels.max(1)));
(denorm(entry.local_px_rate), denorm(entry.remote_px_rate))
}
pub(in crate::bridge::export_thread) fn record_split_rates(
worker: &WorkerSettings,
pixels: u32,
local: Option<f64>,
remote: Option<f64>,
) {
if local.is_none() && remote.is_none() {
return;
}
let pixels = f64::from(pixels.max(1));
let local = local.map(|rate| rate * pixels);
let remote = remote.map(|rate| rate * pixels);
FINAL_PICTURE_RATES
.lock()
.unwrap_or_else(PoisonError::into_inner)
.entry(worker.address.clone())
.and_modify(|entry| {
if let Some(local) = local {
entry.local_px_rate = Some(local);
}
if let Some(remote) = remote {
entry.remote_px_rate = Some(remote);
}
})
.or_insert(SplitRates {
local_px_rate: local,
remote_px_rate: remote,
});
}
pub(in crate::bridge::export_thread) fn decode_final_picture(
accumulator: &Accumulator,
width: u32,
height: u32,
) -> Result<Vec<u8>, String> {
let picture = accumulator
.final_image()
.ok_or_else(|| "the remote finished without sending a picture".to_string())?;
if (picture.width, picture.height) != (width, height) {
return Err(format!(
"the remote sent a {}x{} picture for a {width}x{height} request",
picture.width, picture.height
));
}
indicatrix_net::display::decode_rgba8(picture.encoding, width, height, &picture.bytes)
.map_err(|e| format!("the remote's picture did not decode: {e}"))
}
fn preview_image(snapshot: &PreviewSnapshot) -> SharedPixelBuffer<Rgba8Pixel> {
let rgba = tonemap_accumulation(
snapshot.width,
snapshot.height,
snapshot.samples_done.max(1),
&snapshot.buffer,
ColorSpace::Srgb,
);
let mut buffer = SharedPixelBuffer::<Rgba8Pixel>::new(snapshot.width, snapshot.height);
let dst = buffer.make_mut_slice();
let src: &[Rgba8Pixel] = bytemuck::cast_slice(&rgba);
dst.copy_from_slice(src);
buffer
}
pub(in crate::bridge::export_thread) fn run_final_image_request(
worker: &WorkerSettings,
scene: SceneState,
samples: u32,
color_space: ColorSpace,
cancel: &AtomicBool,
local: Option<&LocalShare<'_>>,
mut on_progress: impl FnMut(u32, u32, Option<SharedPixelBuffer<Rgba8Pixel>>),
) -> FinalPictureOutcome {
let (width, height) = (scene.width, scene.height);
let accumulator = Arc::new(Mutex::new(Accumulator::new(width, height)));
let (tx, rx) = mpsc::channel::<RemoteUpdate>();
let viewer_samples = local.map_or(0, |share| share.samples);
let handle = remote_render::spawn_final_image_request(
RemoteFinalImageRequest {
worker: worker.clone(),
request_id: REQUEST_ID,
scene,
first_sample: 0,
samples,
color_space,
viewer_samples,
},
Arc::clone(&accumulator),
move |update| {
let _ = tx.send(update);
},
);
let picture_bytes = u64::from(width) * u64::from(height) * 4;
let contribution_bytes =
u64::from(width) * u64::from(height) * indicatrix_net::radiance::BYTES_PER_PIXEL as u64;
let mut cancel_sent_at: Option<Instant> = None;
let mut last_update = Instant::now();
let mut seen_first = false;
let mut contribution_sent = false;
loop {
if cancel_sent_at.is_none() && cancel.load(Ordering::Relaxed) {
handle.cancel();
cancel_sent_at = Some(Instant::now());
stop_local(local);
}
if !contribution_sent
&& let Some(share) = local
&& let Ok((sum, _elapsed)) = share.result.try_recv()
&& handle.contribute(sum)
{
contribution_sent = true;
last_update = Instant::now();
}
let local_done = local.map_or(0, |share| share.done.load(Ordering::Relaxed));
match rx.recv_timeout(Duration::from_millis(100)) {
Ok(update) => {
last_update = Instant::now();
seen_first = true;
if let Some(outcome) = on_update(
update,
&accumulator,
width,
height,
local_done,
&mut on_progress,
) {
stop_local(local);
return outcome;
}
}
Err(RecvTimeoutError::Timeout) => {
let in_flight_bytes = if contribution_sent {
contribution_bytes
} else {
picture_bytes
};
match cancel_sent_at {
Some(sent_at) if sent_at.elapsed() > CANCEL_WAIT_TIMEOUT => {
tracing::warn!("final picture: the remote never confirmed the cancel");
stop_local(local);
return FinalPictureOutcome::Cancelled;
}
None if last_update.elapsed()
> liveness_deadline(seen_first, in_flight_bytes) =>
{
stop_local(local);
return FinalPictureOutcome::Failed(format!(
"remote silent for {:.0?}",
last_update.elapsed()
));
}
Some(_) | None => {}
}
}
Err(RecvTimeoutError::Disconnected) => {
stop_local(local);
return FinalPictureOutcome::Failed(
"the remote connection thread ended unexpectedly".to_string(),
);
}
}
}
}
fn stop_local(local: Option<&LocalShare<'_>>) {
if let Some(share) = local {
share.stop.store(true, Ordering::Relaxed);
}
}
fn on_update(
update: RemoteUpdate,
accumulator: &Mutex<Accumulator>,
width: u32,
height: u32,
local_done: u32,
on_progress: &mut impl FnMut(u32, u32, Option<SharedPixelBuffer<Rgba8Pixel>>),
) -> Option<FinalPictureOutcome> {
match update {
RemoteUpdate::Progress { samples_done, .. } => {
on_progress(samples_done, local_done, None);
None
}
RemoteUpdate::Preview { .. } => {
let guard = accumulator.lock().unwrap_or_else(PoisonError::into_inner);
if let Some(snapshot) = guard.last_preview() {
let samples_done = snapshot.samples_done;
let image = preview_image(snapshot);
drop(guard);
on_progress(samples_done, local_done, Some(image));
}
None
}
RemoteUpdate::Done {
cancelled: true, ..
} => Some(FinalPictureOutcome::Cancelled),
RemoteUpdate::Done {
cancelled: false, ..
} => {
let guard = accumulator.lock().unwrap_or_else(PoisonError::into_inner);
let reclaimed_samples = guard
.done_stats()
.map_or(0, |stats| stats.reclaimed_samples);
let decoded = decode_final_picture(&guard, width, height);
drop(guard);
Some(match decoded {
Ok(rgba) => FinalPictureOutcome::Completed {
rgba,
reclaimed_samples,
},
Err(message) => FinalPictureOutcome::Failed(message),
})
}
RemoteUpdate::Unsupported { message, .. } => {
Some(FinalPictureOutcome::Unsupported(message))
}
RemoteUpdate::Failed { message, .. } => Some(FinalPictureOutcome::Failed(message)),
RemoteUpdate::Connected { .. }
| RemoteUpdate::FinalImage { .. }
| RemoteUpdate::Frame { .. }
| RemoteUpdate::DisplayFrame { .. }
| RemoteUpdate::CapabilityChanged { .. } => None,
}
}
#[cfg(test)]
mod tests;