use super::{
super::{
capability::RemoteCapability,
rate::{remote_request_samples, shortfall},
},
progress::RemoteProgress,
run_batch::run_remote_batch,
};
use crate::bridge::export_thread::sample_cursor::SampleCursor;
use indicatrix_net::{SceneState, client::Accumulator};
use std::{
sync::{
Arc, Mutex, PoisonError,
atomic::{AtomicBool, Ordering},
},
thread,
time::{Duration, Instant},
};
pub(super) const MAX_CONSECUTIVE_REMOTE_FAILURES: u32 = 2;
pub(super) const REMOTE_RETRY_BACKOFF_INITIAL: Duration = Duration::from_secs(15);
pub(super) const REMOTE_RETRY_BACKOFF_MAX: Duration = Duration::from_secs(120);
pub(super) fn remote_retry_backoff(consecutive_failures: u32) -> Duration {
let doublings = consecutive_failures.saturating_sub(MAX_CONSECUTIVE_REMOTE_FAILURES);
let scaled = REMOTE_RETRY_BACKOFF_INITIAL.saturating_mul(1u32 << doublings.min(8));
scaled.min(REMOTE_RETRY_BACKOFF_MAX)
}
pub(super) fn pause_remote_lane(
cursor: &SampleCursor,
cancel: &AtomicBool,
total: Duration,
) -> bool {
let started = Instant::now();
while started.elapsed() < total {
if cancel.load(Ordering::Relaxed) || cursor.shared_pool_exhausted() {
return false;
}
thread::sleep(Duration::from_millis(100));
}
true
}
pub(in crate::bridge::export_thread) struct RemoteLaneOutcome {
pub(in crate::bridge::export_thread) fatal: Option<String>,
pub(in crate::bridge::export_thread) final_rate: Option<f64>,
}
#[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
}
fn merge_chunk_into_progress(
progress: &RemoteProgress,
chunk_accumulator: &Mutex<Accumulator>,
done: u32,
) {
if done == 0 {
return;
}
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);
}
#[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,
final_rate: Some(rate),
};
}
let Some((start, count)) =
cursor.claim(remote_request_samples(rate, capability.coordinator))
else {
return RemoteLaneOutcome {
fatal: None,
final_rate: Some(rate),
};
};
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;
merge_chunk_into_progress(progress, &chunk_accumulator, done);
if let Some(measured) = measured_rate {
rate = measured;
}
if cancelled {
return RemoteLaneOutcome {
fatal: None,
final_rate: Some(rate),
};
}
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)."
)),
final_rate: Some(rate),
};
}
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 {
let pause = remote_retry_backoff(consecutive_failures);
progress
.notes
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push_back(format!(
"Remote worker failed {consecutive_failures} chunks in a row -- \
pausing it for {}s before offering it more work; rendering \
continues locally in the meantime.",
pause.as_secs()
));
if !pause_remote_lane(cursor, cancel, pause) {
return RemoteLaneOutcome {
fatal: None,
final_rate: Some(rate),
};
}
}
}
}