use super::super::{capability::RemoteCapability, rate::remote_request_rate};
use crate::bridge::remote::remote_render::{self, RemoteRenderRequest, RemoteUpdate};
use indicatrix_net::{SceneState, client::Accumulator};
use std::{
sync::{
Arc, Mutex, PoisonError,
atomic::{AtomicBool, Ordering},
mpsc::{self, RecvTimeoutError},
},
time::{Duration, Instant},
};
const REQUEST_ID: u32 = 1;
pub(super) const LIVENESS_TIMEOUT: Duration = Duration::from_secs(8);
pub(super) const FIRST_EVENT_TIMEOUT: Duration = Duration::from_secs(30);
const MIN_ASSUMED_LINK_BYTES_PER_SEC: u64 = 4 * 1024 * 1024;
pub(super) const fn transfer_allowance(frame_bytes: u64) -> Duration {
Duration::from_secs(frame_bytes.div_ceil(MIN_ASSUMED_LINK_BYTES_PER_SEC))
}
pub(in crate::bridge::export_thread) const fn liveness_deadline(
seen_first_update: bool,
frame_bytes: u64,
) -> Duration {
let base = if seen_first_update {
LIVENESS_TIMEOUT
} else {
FIRST_EVENT_TIMEOUT
};
base.saturating_add(transfer_allowance(frame_bytes))
}
pub(in crate::bridge::export_thread) const CANCEL_WAIT_TIMEOUT: Duration = Duration::from_secs(10);
#[derive(Debug, Default)]
pub(super) struct ProgressSpan {
first: Option<(Instant, u32)>,
last: Option<(Instant, u32)>,
}
impl ProgressSpan {
pub(super) fn note(&mut self, at: Instant, samples_done: u32) {
let latest = self.last.map_or(0, |(_, samples)| samples);
if samples_done <= latest {
return;
}
if self.first.is_none() {
self.first = Some((at, samples_done));
}
self.last = Some((at, samples_done));
}
pub(super) fn advance(&self) -> Option<(u32, Duration)> {
let (first_at, first_samples) = self.first?;
let (last_at, last_samples) = self.last?;
if last_samples <= first_samples {
return None;
}
Some((
last_samples - first_samples,
last_at.saturating_duration_since(first_at),
))
}
}
#[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 request_sent_at = Instant::now();
let handle = remote_render::spawn_remote_render(
RemoteRenderRequest {
worker: capability.worker.clone(),
request_id: REQUEST_ID,
scene,
first_sample,
samples,
width,
height,
intent: indicatrix_net::messages::RequestIntent::Batch,
display_only: false,
},
Arc::clone(accumulator),
move |update| {
let _ = tx.send(update);
},
);
let mut span = ProgressSpan::default();
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;
let frame_bytes =
u64::from(width) * u64::from(height) * indicatrix_net::radiance::BYTES_PER_PIXEL as u64;
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) => {
let received_at = Instant::now();
last_update = received_at;
first_update_seen = true;
match update {
RemoteUpdate::Done { cancelled, .. } => {
let samples_done = accumulator
.lock()
.unwrap_or_else(PoisonError::into_inner)
.samples_done();
let rate = remote_request_rate(
capability.coordinator,
samples_done,
received_at.saturating_duration_since(request_sent_at),
span.advance(),
);
return (samples_done, cancelled, None, rate);
}
RemoteUpdate::Failed { message, .. }
| RemoteUpdate::Unsupported { 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, .. } => {
span.note(received_at, samples_done);
}
_ => {}
}
}
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, frame_bytes) =>
{
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,
);
}
}
}
}