use super::{
capability::RemoteCapability,
rate::{remote_chunk_samples, remote_marginal_rate, shortfall},
};
use crate::bridge::{
export_thread::{sample_cursor::SampleCursor, scene_snapshot::SceneSnapshot},
remote::remote_render::{self, RemoteRenderRequest, RemoteUpdate},
};
use glam::Vec3;
use indicatrix_net::{SceneState, client::Accumulator};
use std::{
collections::VecDeque,
sync::{
Arc, Mutex, PoisonError,
atomic::{AtomicBool, AtomicU32, Ordering},
mpsc::{self, RecvTimeoutError},
},
time::{Duration, Instant},
};
const REQUEST_ID: u32 = 1;
const LIVENESS_TIMEOUT: Duration = Duration::from_secs(8);
const FIRST_EVENT_TIMEOUT: Duration = Duration::from_secs(30);
const fn liveness_deadline(seen_first_update: bool) -> Duration {
if seen_first_update {
LIVENESS_TIMEOUT
} else {
FIRST_EVENT_TIMEOUT
}
}
const CANCEL_WAIT_TIMEOUT: Duration = Duration::from_secs(10);
#[must_use]
pub(in crate::bridge::export_thread) fn scene_state_from_snapshot(
snapshot: &SceneSnapshot,
width: u32,
height: u32,
) -> SceneState {
SceneState {
width,
height,
yaw: snapshot.yaw,
pitch: snapshot.pitch,
distance: snapshot.distance,
light_yaw: snapshot.light_yaw,
light_pitch: snapshot.light_pitch,
exposure: snapshot.exposure,
max_bounces: snapshot.max_bounces,
lighting_preset: snapshot.lighting_preset,
material: snapshot.material.clone(),
planes: snapshot.active_planes.clone(),
girdle_frosted: !snapshot.facet_finishes.is_empty(),
}
}
#[expect(
clippy::too_many_arguments,
reason = "every argument is a distinct piece of one RenderRequest's own identity \
(worker capability, scene, sample range, resolution, the shared \
accumulator, cancellation) -- bundling them into a struct would just move \
the same count into field access, not reduce it, and `RemoteRenderRequest` \
already plays that role one level down in `bridge::remote::remote_render`"
)]
pub(in crate::bridge::export_thread) fn run_remote_batch(
capability: &RemoteCapability,
scene: SceneState,
first_sample: u32,
samples: u32,
width: u32,
height: u32,
accumulator: &Arc<Mutex<Accumulator>>,
cancel: &AtomicBool,
) -> (u32, bool, Option<String>, Option<f64>) {
let (tx, rx) = mpsc::channel::<RemoteUpdate>();
let handle = remote_render::spawn_remote_render(
RemoteRenderRequest {
worker: capability.worker.clone(),
request_id: REQUEST_ID,
scene,
first_sample,
samples,
width,
height,
},
Arc::clone(accumulator),
move |update| {
let _ = tx.send(update);
},
);
let mut first_progress: Option<(Instant, u32)> = None;
let mut cancel_sent = false;
let mut cancel_sent_at: Option<Instant> = None;
let mut last_update = Instant::now();
let mut first_update_seen = false;
loop {
if !cancel_sent && cancel.load(Ordering::Relaxed) {
handle.cancel();
cancel_sent = true;
cancel_sent_at = Some(Instant::now());
}
match rx.recv_timeout(Duration::from_millis(100)) {
Ok(update) => {
last_update = Instant::now();
first_update_seen = true;
match update {
RemoteUpdate::Done { cancelled, .. } => {
let samples_done = accumulator
.lock()
.unwrap_or_else(PoisonError::into_inner)
.samples_done();
let rate = first_progress.and_then(|(first_instant, first_samples)| {
remote_marginal_rate(
samples_done.saturating_sub(first_samples),
Instant::now().saturating_duration_since(first_instant),
)
});
return (samples_done, cancelled, None, rate);
}
RemoteUpdate::Failed { message, .. } => {
let acc = accumulator.lock().unwrap_or_else(PoisonError::into_inner);
return (acc.samples_done(), false, Some(message), None);
}
RemoteUpdate::Frame { samples_done, .. }
| RemoteUpdate::Progress { samples_done, .. } => {
if first_progress.is_none() && samples_done > 0 {
first_progress = Some((Instant::now(), samples_done));
}
}
RemoteUpdate::Connected { .. } | RemoteUpdate::Preview { .. } => {}
}
}
Err(RecvTimeoutError::Timeout) => match cancel_sent_at {
Some(sent_at) if sent_at.elapsed() > CANCEL_WAIT_TIMEOUT => {
tracing::warn!(
"remote render: worker never confirmed cancellation within \
{CANCEL_WAIT_TIMEOUT:?} -- returning with samples completed so \
far"
);
let acc = accumulator.lock().unwrap_or_else(PoisonError::into_inner);
return (acc.samples_done(), true, None, None);
}
None if last_update.elapsed() > liveness_deadline(first_update_seen) => {
let message = format!(
"worker silent for {:.0?} -- treating as failed",
last_update.elapsed()
);
tracing::warn!("remote render: {message}");
let acc = accumulator.lock().unwrap_or_else(PoisonError::into_inner);
return (acc.samples_done(), false, Some(message), None);
}
Some(_) | None => {}
},
Err(RecvTimeoutError::Disconnected) => {
let acc = accumulator.lock().unwrap_or_else(PoisonError::into_inner);
return (
acc.samples_done(),
false,
Some("remote render worker thread ended unexpectedly".to_string()),
None,
);
}
}
}
}
pub(in crate::bridge::export_thread) struct RemoteProgress {
accum: Mutex<Vec<Vec3>>,
traced: AtomicU32,
in_flight: Mutex<Option<Arc<Mutex<Accumulator>>>>,
notes: Mutex<VecDeque<String>>,
}
impl RemoteProgress {
pub(in crate::bridge::export_thread) fn new(pixel_count: usize) -> Self {
Self {
accum: Mutex::new(vec![Vec3::ZERO; pixel_count]),
traced: AtomicU32::new(0),
in_flight: Mutex::new(None),
notes: Mutex::new(VecDeque::new()),
}
}
pub(in crate::bridge::export_thread) fn samples_done(&self) -> u32 {
let in_flight = self
.in_flight
.lock()
.unwrap_or_else(PoisonError::into_inner);
let live = in_flight.as_ref().map_or(0, |acc| {
acc.lock()
.unwrap_or_else(PoisonError::into_inner)
.samples_done()
});
drop(in_flight);
self.traced.load(Ordering::Relaxed) + live
}
pub(in crate::bridge::export_thread) fn preview_buffer(&self) -> Vec<Vec3> {
let mut buf = self
.accum
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clone();
let in_flight = self
.in_flight
.lock()
.unwrap_or_else(PoisonError::into_inner);
if let Some(acc) = in_flight.as_ref() {
let acc = acc.lock().unwrap_or_else(PoisonError::into_inner);
for (dst, src) in buf.iter_mut().zip(acc.buffer()) {
*dst += *src;
}
}
drop(in_flight);
buf
}
pub(in crate::bridge::export_thread) fn take_note(&self) -> Option<String> {
self.notes
.lock()
.unwrap_or_else(PoisonError::into_inner)
.pop_front()
}
pub(in crate::bridge::export_thread) fn into_buffer(self) -> Vec<Vec3> {
self.accum
.into_inner()
.unwrap_or_else(PoisonError::into_inner)
}
}
const MAX_CONSECUTIVE_REMOTE_FAILURES: u32 = 2;
pub(in crate::bridge::export_thread) struct RemoteLaneOutcome {
pub(in crate::bridge::export_thread) fatal: Option<String>,
}
#[expect(
clippy::too_many_arguments,
reason = "every argument is a distinct piece of one remote lane's own identity \
(the shared cursor, worker capability, scene, resolution, the initial \
rate estimate, the cross-thread progress/notes state, whether a failed \
chunk may fall back to local, and cancellation/completion signalling) -- \
bundling them into a struct would just move the same count into field \
access, not reduce it"
)]
pub(in crate::bridge::export_thread) fn run_remote_lane(
cursor: &SampleCursor,
capability: &RemoteCapability,
scene_state: &SceneState,
width: u32,
height: u32,
initial_rate: f64,
progress: &RemoteProgress,
fallback_to_local: bool,
cancel: &AtomicBool,
remote_lane_done: &AtomicBool,
) -> RemoteLaneOutcome {
let outcome = run_remote_lane_claim_loop(
cursor,
capability,
scene_state,
width,
height,
initial_rate,
progress,
fallback_to_local,
cancel,
);
remote_lane_done.store(true, Ordering::Release);
outcome
}
#[expect(
clippy::too_many_arguments,
reason = "see `run_remote_lane`'s own identical `#[expect]` -- this is that \
function's claim loop, split out only so the `remote_lane_done` \
store-on-every-exit-path guarantee lives in exactly one place (the \
wrapper) rather than being duplicated at every `return` this loop has"
)]
fn run_remote_lane_claim_loop(
cursor: &SampleCursor,
capability: &RemoteCapability,
scene_state: &SceneState,
width: u32,
height: u32,
initial_rate: f64,
progress: &RemoteProgress,
fallback_to_local: bool,
cancel: &AtomicBool,
) -> RemoteLaneOutcome {
let mut rate = initial_rate;
let mut consecutive_failures = 0u32;
loop {
if cancel.load(Ordering::Relaxed) {
return RemoteLaneOutcome { fatal: None };
}
let Some((start, count)) = cursor.claim(remote_chunk_samples(rate)) else {
return RemoteLaneOutcome { fatal: None };
};
let chunk_accumulator = Arc::new(Mutex::new(Accumulator::new(width, height)));
*progress
.in_flight
.lock()
.unwrap_or_else(PoisonError::into_inner) = Some(Arc::clone(&chunk_accumulator));
let (done, cancelled, error, measured_rate) = run_remote_batch(
capability,
scene_state.clone(),
start,
count,
width,
height,
&chunk_accumulator,
cancel,
);
*progress
.in_flight
.lock()
.unwrap_or_else(PoisonError::into_inner) = None;
if done > 0 {
let chunk = chunk_accumulator
.lock()
.unwrap_or_else(PoisonError::into_inner);
let mut merged = progress
.accum
.lock()
.unwrap_or_else(PoisonError::into_inner);
for (dst, src) in merged.iter_mut().zip(chunk.buffer()) {
*dst += *src;
}
drop(merged);
drop(chunk);
progress.traced.fetch_add(done, Ordering::Relaxed);
}
if let Some(measured) = measured_rate {
rate = measured;
}
if cancelled {
return RemoteLaneOutcome { fatal: None };
}
let missing = shortfall(count, done);
if missing == 0 {
consecutive_failures = 0;
continue;
}
let message = error.as_deref().unwrap_or("connection ended early");
if !fallback_to_local {
return RemoteLaneOutcome {
fatal: Some(format!(
"Remote worker failed ({message}) -- {missing} of {count} samples in \
this chunk were never traced, and this export has no local fallback \
(Compute: Remote)."
)),
};
}
cursor.return_to_local(start + done, missing);
progress
.notes
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push_back(format!(
"Remote worker failed ({message}) partway through a chunk -- the \
remaining {missing} samples will finish locally. Every sample it \
completed first is still included."
));
consecutive_failures += 1;
if consecutive_failures >= MAX_CONSECUTIVE_REMOTE_FAILURES {
progress
.notes
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push_back(
"Remote worker failed repeatedly -- no longer sending it further \
work; finishing the rest of the export locally."
.to_string(),
);
return RemoteLaneOutcome { fatal: None };
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bridge::export_thread::scene_snapshot::SceneSnapshot;
use indicatrix::{geometry::cuts::StandardGemCuts, optics::materials::GemMaterial};
#[test]
fn scene_state_from_snapshot_carries_the_snapshots_own_bounce_cap_not_a_viewport_default() {
let snapshot = SceneSnapshot {
yaw: 0.6,
pitch: 0.45,
distance: 2.4,
light_yaw: 0.85,
light_pitch: 0.95,
material: GemMaterial::diamond(),
lighting_preset: indicatrix::optics::raytracer::LightingPreset::RingLights,
max_bounces: 64, exposure: 1.0,
active_planes: StandardGemCuts::standard_round_brilliant(),
facet_finishes: Vec::new(),
env_map: None,
};
let state = scene_state_from_snapshot(&snapshot, 1920, 1080);
assert_eq!(
state.max_bounces, 64,
"the remote RenderRequest's scene must carry the export's OWN bounce cap"
);
}
#[test]
fn liveness_deadline_grants_the_first_event_grace_before_any_update_has_arrived() {
assert_eq!(
liveness_deadline(false),
FIRST_EVENT_TIMEOUT,
"the wait for a dispatch's very first RemoteUpdate (even Connected) must \
use the longer grace, not the steady-state deadline"
);
}
#[test]
fn liveness_deadline_switches_to_the_tighter_steady_state_timeout_once_seen() {
assert_eq!(
liveness_deadline(true),
LIVENESS_TIMEOUT,
"once a dispatch has produced at least one update, every wait after it must \
use the steady-state liveness timeout, not the first-event grace"
);
}
#[test]
fn first_event_timeout_is_strictly_longer_than_liveness_timeout() {
assert!(FIRST_EVENT_TIMEOUT > LIVENESS_TIMEOUT);
}
}