use super::engine::{
BatchContext, LaneShared, PREVIEW_MAX_BOUNCES, PreviewItem, RecordRevision, ResolvedDesign,
finish_item, finish_item_at_revision, panic_message, push_progress, resolve_design,
update_remote,
};
use crate::{
MainWindow,
bridge::{
preview_render::{self, PreviewJob, PreviewView},
preview_wait::RemoteShortfall,
},
gui::batch::{
remote_backoff::{REMOTE_WAIT_SLICE, RemoteBackoff, wait_unless},
remote_dispatch::{DispatcherGroup, RemoteStatus},
},
settings::WorkerSettings,
};
use slint::Weak;
use std::{
panic::{self, AssertUnwindSafe},
sync::atomic::{AtomicBool, Ordering},
time::Instant,
};
use tracing::{info, warn};
const SIT_OUT_TOAST: &str = "Remote worker is failing every preview request -- remote lane \
paused, rendering locally (see indicatrix-cut.log)";
enum Attempt {
Rendered(Vec<u8>),
Failed(RemoteShortfall),
NotAttempted,
}
#[derive(Default)]
struct LaneHealth {
backoff: RemoteBackoff,
}
fn catch_render_checked(
f: impl FnOnce() -> Result<Vec<u8>, RemoteShortfall>,
) -> Result<Vec<u8>, RemoteShortfall> {
panic::catch_unwind(AssertUnwindSafe(f)).unwrap_or_else(|payload| {
Err(RemoteShortfall::Failed(format!(
"render panicked: {}",
panic_message(&*payload)
)))
})
}
fn render_item_remote(
ctx: &BatchContext<'_>,
worker: &WorkerSettings,
resolved: &ResolvedDesign,
view: PreviewView,
cancel: &AtomicBool,
) -> Result<Vec<u8>, RemoteShortfall> {
let job = PreviewJob {
planes: &resolved.planes,
material: &resolved.material,
size: ctx.preview_size,
spp: ctx.preview_spp,
max_bounces: PREVIEW_MAX_BOUNCES,
};
catch_render_checked(|| preview_render::render_view_remote_checked(&job, view, worker, cancel))
}
fn wait_out_failures(shared: &LaneShared<'_>, ui_weak: &Weak<MainWindow>, backoff: &RemoteBackoff) {
let Some(delay) = backoff.wait_needed(Instant::now()) else {
return;
};
let title = format!(
"Remote paused after {} failure(s) -- retrying in {delay:.0?}",
backoff.consecutive_failures()
);
update_remote(shared.progress, |remote| remote.note(&title));
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
wait_unless(delay, REMOTE_WAIT_SLICE, || {
shared.cancel.load(Ordering::Relaxed) || shared.queue.shared_is_empty()
});
}
fn attempt_remote(
shared: &LaneShared<'_>,
worker: Option<&WorkerSettings>,
resolved: Option<&ResolvedDesign>,
view: PreviewView,
) -> Attempt {
match (resolved, worker) {
(Some(r), Some(w)) => match render_item_remote(shared.ctx, w, r, view, shared.cancel) {
Ok(bytes) => Attempt::Rendered(bytes),
Err(shortfall) => Attempt::Failed(shortfall),
},
_ => Attempt::NotAttempted,
}
}
fn note_success(health: &mut LaneHealth) {
let failures = health.backoff.consecutive_failures();
if failures > 0 {
info!(recovered_after = failures, "preview remote lane recovered");
}
health.backoff.record_success();
}
fn note_failure(
shared: &LaneShared<'_>,
health: &mut LaneHealth,
ui_weak: &Weak<MainWindow>,
item: PreviewItem,
shortfall: &RemoteShortfall,
fallback_to_local: bool,
) {
if matches!(shortfall, RemoteShortfall::Cancelled) {
return;
}
let retry_in = health.backoff.record_failure(Instant::now());
warn!(
entry_id = item.entry_id,
view = ?item.view,
%shortfall,
consecutive = health.backoff.consecutive_failures(),
?retry_in,
"preview remote lane: item failed; requeued locally"
);
if fallback_to_local
&& health.backoff.sitting_out()
&& !shared.sit_out_toasted.swap(true, Ordering::Relaxed)
{
let _ = ui_weak.upgrade_in_event_loop(|ui| {
crate::gui::show_toast(&ui, SIT_OUT_TOAST, "warning");
});
}
}
fn dispose_unrendered(
shared: &LaneShared<'_>,
ui_weak: &Weak<MainWindow>,
item: PreviewItem,
fallback_to_local: bool,
) {
if fallback_to_local {
shared.queue.return_to_local(item);
} else {
finish_item(shared, ui_weak, item.entry_id, item.view, None);
}
}
fn serve_item(
shared: &LaneShared<'_>,
ui_weak: &Weak<MainWindow>,
worker: Option<&WorkerSettings>,
fallback_to_local: bool,
health: &mut LaneHealth,
item: PreviewItem,
) {
let resolved = resolve_design(shared.ctx, item.entry_id);
let title = resolved
.as_ref()
.map_or_else(|| format!("Design #{}", item.entry_id), |r| r.title.clone());
update_remote(shared.progress, |remote| remote.item_started(&title));
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
let attempt = attempt_remote(shared, worker, resolved.as_ref(), item.view);
update_remote(shared.progress, RemoteStatus::item_ended);
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
match attempt {
Attempt::Rendered(bytes) => {
note_success(health);
let revision = resolved
.as_ref()
.map_or(RecordRevision::Unknown, |r| r.revision);
finish_item_at_revision(
shared,
ui_weak,
item.entry_id,
item.view,
Some(bytes),
revision,
);
}
Attempt::Failed(shortfall) => {
note_failure(shared, health, ui_weak, item, &shortfall, fallback_to_local);
dispose_unrendered(shared, ui_weak, item, fallback_to_local);
}
Attempt::NotAttempted => dispose_unrendered(shared, ui_weak, item, fallback_to_local),
}
}
pub(super) fn run_remote_lane(
shared: &LaneShared<'_>,
ui_weak: &Weak<MainWindow>,
worker: Option<&WorkerSettings>,
fallback_to_local: bool,
dispatchers: &DispatcherGroup<'_>,
) {
let _counted_out_on_drop = dispatchers.guard();
update_remote(shared.progress, RemoteStatus::dispatcher_started);
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
let mut health = LaneHealth::default();
loop {
if fallback_to_local {
wait_out_failures(shared, ui_weak, &health.backoff);
}
if shared.cancel.load(Ordering::Relaxed) {
break;
}
let Some(item) = shared.queue.claim_shared() else {
break;
};
serve_item(
shared,
ui_weak,
worker,
fallback_to_local,
&mut health,
item,
);
}
update_remote(shared.progress, RemoteStatus::dispatcher_ended);
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
}