use super::remote_lane::run_remote_lane;
use crate::{
BatchModel, MainWindow,
bridge::preview_render::{self, CacheKind, PreviewJob, PreviewView},
gui::{
batch::{
batch_queue::WorkQueue,
material_choice::ensure_balanced_material,
remote_dispatch::{DispatcherGroup, RemoteStatus},
},
progress_eta::{EtaEstimator, batch_eta_label},
},
settings::WorkerSettings,
};
use indicatrix::{
geometry::plane::GpuFacetPlane,
optics::{materials::GemMaterial, raytracer::DEFAULT_MAX_BOUNCES},
renderer::gpu_backend::GpuBackend,
};
use indicatrix_vault::db::sqlite::Database;
use slint::{ComponentHandle, Weak};
use std::{
any::Any,
collections::{BTreeSet, HashMap},
panic::{self, AssertUnwindSafe},
sync::{
Mutex, PoisonError,
atomic::{AtomicBool, AtomicU32, Ordering},
},
thread,
time::{Duration, SystemTime, UNIX_EPOCH},
};
use tracing::warn;
mod accumulator;
pub(super) use accumulator::{DesignAccum, RecordRevision};
use accumulator::{FinishedViews, record_item_result};
pub(super) const PREVIEW_MAX_BOUNCES: u32 = DEFAULT_MAX_BOUNCES;
pub const RI_MATCH_TOLERANCE: f64 = 0.02;
pub const FALLBACK_TARGET_RI: f64 = 2.417;
#[must_use]
pub fn target_ri_for_design(full: &indicatrix_vault::model::entry::FullDiagramRecord) -> f64 {
full.refractive_index
.as_deref()
.and_then(|s| s.trim().parse::<f64>().ok())
.filter(|ri| ri.is_finite() && *ri > 1.0)
.unwrap_or(FALLBACK_TARGET_RI)
}
const LOCAL_IDLE_POLL: Duration = Duration::from_millis(15);
pub(super) struct BatchContext<'a> {
pub(super) db: &'a Mutex<Database>,
pub(super) material_candidates:
&'a [indicatrix_vault::model::material_match::RiPresetCandidate],
pub(super) preview_size: u32,
pub(super) preview_spp: u32,
pub(super) solid: bool,
pub(super) gpu_retired: &'a AtomicBool,
pub(super) angle_table_entries: &'a Mutex<BTreeSet<i64>>,
}
#[derive(Debug, Clone, Copy)]
pub(super) struct PreviewItem {
pub(super) entry_id: i64,
pub(super) view: PreviewView,
}
pub(super) struct ResolvedDesign {
pub(super) title: String,
pub(super) planes: Vec<GpuFacetPlane>,
pub(super) material: GemMaterial,
pub(super) revision: RecordRevision,
}
pub(super) fn record_planes(
ctx: &BatchContext<'_>,
full: &indicatrix_vault::model::entry::FullDiagramRecord,
) -> Option<Vec<GpuFacetPlane>> {
super::super::record_planes_for_batch(full, ctx.angle_table_entries)
}
pub(super) fn resolve_design(ctx: &BatchContext<'_>, entry_id: i64) -> Option<ResolvedDesign> {
let (full, updated_at) = {
let guard = ctx.db.lock().unwrap_or_else(PoisonError::into_inner);
(
guard.get_diagram_full(entry_id),
guard.entry_updated_at(entry_id),
)
};
let Ok(Some(full)) = full else {
return None;
};
let Ok(updated_at) = updated_at else {
return None;
};
let planes = record_planes(ctx, &full)?;
if ctx.solid {
return Some(ResolvedDesign {
title: full.title,
planes,
material: GemMaterial::by_name("Diamond")?,
revision: RecordRevision::Stamp(updated_at),
});
}
let target_ri = target_ri_for_design(&full);
let material_name = {
let guard = ctx.db.lock().unwrap_or_else(PoisonError::into_inner);
ensure_balanced_material(
&guard,
entry_id,
target_ri,
ctx.material_candidates,
RI_MATCH_TOLERANCE,
&planes,
)
};
let Ok(Some(material_name)) = material_name else {
return None;
};
let material = GemMaterial::by_name(&material_name)?;
Some(ResolvedDesign {
title: full.title,
planes,
material,
revision: RecordRevision::Stamp(updated_at),
})
}
pub(super) fn panic_message(payload: &(dyn Any + Send)) -> String {
payload
.downcast_ref::<&str>()
.map(|s| (*s).to_string())
.or_else(|| payload.downcast_ref::<String>().cloned())
.unwrap_or_else(|| "unknown panic".to_string())
}
fn catch_local_render(
view: PreviewView,
gpu_retired: &AtomicBool,
f: impl FnOnce() -> Option<Vec<u8>>,
) -> Option<Vec<u8>> {
panic::catch_unwind(AssertUnwindSafe(f)).unwrap_or_else(|payload| {
let message = panic_message(&*payload);
if message.contains("wgpu") || message.contains("still mapped") {
tracing::error!(
"Preview render panicked for a {view:?} view with a GPU-fatal message; \
retiring the shared GPU backend for the rest of this batch: {message}"
);
gpu_retired.store(true, Ordering::Relaxed);
} else {
warn!("Preview render panicked for a {view:?} view: {message}");
}
None
})
}
fn render_item_local(
ctx: &BatchContext<'_>,
gpu: &GpuBackend,
resolved: &ResolvedDesign,
view: PreviewView,
) -> Option<Vec<u8>> {
if ctx.solid {
return catch_local_render(view, ctx.gpu_retired, || {
preview_render::render_view_solid(&resolved.planes, ctx.preview_size, view)
});
}
if ctx.gpu_retired.load(Ordering::Relaxed) {
return None;
}
let job = PreviewJob {
planes: &resolved.planes,
material: &resolved.material,
size: ctx.preview_size,
spp: ctx.preview_spp,
max_bounces: PREVIEW_MAX_BOUNCES,
};
catch_local_render(view, ctx.gpu_retired, || {
preview_render::render_view(&job, view, gpu)
})
}
#[derive(Default)]
pub(super) struct Tally {
pub(super) generated: AtomicU32,
pub(super) failed: AtomicU32,
}
#[derive(Default, Clone)]
pub(super) struct LiveProgress {
completed: u32,
local_active: u32,
remote: RemoteStatus,
eta: EtaEstimator,
}
pub(super) fn push_progress(
ui_weak: &Weak<MainWindow>,
design_total: u32,
local_lane_total: u32,
progress: &Mutex<LiveProgress>,
) {
let (snapshot, eta_text) = {
let mut p = progress.lock().unwrap_or_else(PoisonError::into_inner);
let now = std::time::Instant::now();
if design_total > 0 {
let fraction = f64::from(p.completed) / f64::from(design_total);
p.eta.observe(now, fraction);
}
let eta_text = batch_eta_label(p.completed, design_total, p.eta.eta(now));
(p.clone(), eta_text)
};
let ui_weak = ui_weak.clone();
let _ = ui_weak.upgrade_in_event_loop(move |ui| {
ui.global::<BatchModel>().set_preview_eta(eta_text.into());
ui.global::<BatchModel>()
.set_preview_design_index(snapshot.completed as i32);
ui.global::<BatchModel>()
.set_preview_design_total(design_total as i32);
ui.global::<BatchModel>()
.set_preview_local_active(snapshot.local_active as i32);
ui.global::<BatchModel>()
.set_preview_local_lane_total(local_lane_total as i32);
ui.global::<BatchModel>()
.set_preview_remote_title(snapshot.remote.title().into());
ui.global::<BatchModel>()
.set_preview_remote_active(snapshot.remote.is_active());
ui.global::<BatchModel>()
.set_preview_remote_in_flight(snapshot.remote.in_flight() as i32);
});
}
pub(super) fn update_remote(
progress: &Mutex<LiveProgress>,
change: impl FnOnce(&mut RemoteStatus),
) {
change(
&mut progress
.lock()
.unwrap_or_else(PoisonError::into_inner)
.remote,
);
}
fn increment_completed(progress: &Mutex<LiveProgress>) {
progress
.lock()
.unwrap_or_else(PoisonError::into_inner)
.completed += 1;
}
fn adjust_local_active(progress: &Mutex<LiveProgress>, delta: i32) {
let mut p = progress.lock().unwrap_or_else(PoisonError::into_inner);
p.local_active = p.local_active.saturating_add_signed(delta);
}
pub(super) struct LaneShared<'a> {
pub(super) ctx: &'a BatchContext<'a>,
pub(super) queue: &'a WorkQueue<PreviewItem>,
pub(super) design_state: &'a Mutex<HashMap<i64, DesignAccum>>,
pub(super) tally: &'a Tally,
pub(super) progress: &'a Mutex<LiveProgress>,
pub(super) cancel: &'a AtomicBool,
pub(super) design_total: u32,
pub(super) local_lane_total: u32,
pub(super) remote_lane_total: u32,
pub(super) sit_out_toasted: AtomicBool,
pub(super) thumbnail_cache: &'a crate::gui::batch::preview_cache::PreviewThumbnailCache,
}
fn save_finished_design(shared: &LaneShared<'_>, entry_id: i64, done: &FinishedViews) -> bool {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX));
let guard = shared.ctx.db.lock().unwrap_or_else(PoisonError::into_inner);
let expected_updated_at = match done.revision {
RecordRevision::Stamp(stamp) => stamp,
RecordRevision::Unknown => match guard.entry_updated_at(entry_id) {
Ok(stamp) => stamp,
Err(e) => {
warn!("Could not save the previews for entry {entry_id}: {e}");
return false;
}
},
RecordRevision::Conflicting => {
warn!(
"Discarded the previews rendered for entry {entry_id}: the design changed \
between its two views"
);
return false;
}
};
let material = guard.get_preview_material(entry_id).ok().flatten();
let kind = if shared.ctx.solid {
CacheKind::SolidDraft {
size: shared.ctx.preview_size,
}
} else {
CacheKind::Preview {
size: shared.ctx.preview_size,
spp: shared.ctx.preview_spp,
max_bounces: PREVIEW_MAX_BOUNCES,
}
};
let fingerprint = preview_render::cache_fingerprint(kind, material.as_deref());
let stored = guard.save_preview_images(
entry_id,
done.front.as_deref(),
done.top.as_deref(),
now,
&fingerprint,
expected_updated_at,
);
drop(guard);
match stored {
Ok(true) => true,
Ok(false) => {
warn!(
"Discarded the previews rendered for entry {entry_id}: the design changed \
while they were rendering"
);
false
}
Err(e) => {
warn!("Failed to save the previews for entry {entry_id}: {e}");
false
}
}
}
pub(super) fn finish_item(
shared: &LaneShared<'_>,
ui_weak: &Weak<MainWindow>,
entry_id: i64,
view: PreviewView,
bytes: Option<Vec<u8>>,
) {
finish_item_at_revision(
shared,
ui_weak,
entry_id,
view,
bytes,
RecordRevision::Unknown,
);
}
pub(super) fn finish_item_at_revision(
shared: &LaneShared<'_>,
ui_weak: &Weak<MainWindow>,
entry_id: i64,
view: PreviewView,
bytes: Option<Vec<u8>>,
revision: RecordRevision,
) {
let Some(done) = record_item_result(shared.design_state, entry_id, view, bytes, revision)
else {
return;
};
let saved = (done.front.is_some() || done.top.is_some())
&& save_finished_design(shared, entry_id, &done);
if saved {
shared.tally.generated.fetch_add(1, Ordering::Relaxed);
shared.thumbnail_cache.invalidate(ui_weak, entry_id);
} else {
shared.tally.failed.fetch_add(1, Ordering::Relaxed);
}
increment_completed(shared.progress);
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
}
pub(super) fn run_local_lane(
shared: &LaneShared<'_>,
gpu: &GpuBackend,
ui_weak: &Weak<MainWindow>,
remote_lane_done: &AtomicBool,
) {
loop {
if shared.cancel.load(Ordering::Relaxed) {
break;
}
let remote_finished = remote_lane_done.load(Ordering::Acquire);
let Some(item) = shared.queue.claim_local() else {
if remote_finished {
break;
}
thread::sleep(LOCAL_IDLE_POLL);
continue;
};
adjust_local_active(shared.progress, 1);
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
let resolved = resolve_design(shared.ctx, item.entry_id);
let revision = resolved
.as_ref()
.map_or(RecordRevision::Unknown, |r| r.revision);
let bytes = resolved.and_then(|r| render_item_local(shared.ctx, gpu, &r, item.view));
finish_item_at_revision(shared, ui_weak, item.entry_id, item.view, bytes, revision);
adjust_local_active(shared.progress, -1);
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
}
}
pub(super) fn build_items(entry_ids: &[i64]) -> Vec<PreviewItem> {
entry_ids
.iter()
.flat_map(|&entry_id| {
[
PreviewItem {
entry_id,
view: PreviewView::Front,
},
PreviewItem {
entry_id,
view: PreviewView::Top,
},
]
})
.collect()
}
pub(super) fn run_batch_lanes(
shared: &LaneShared<'_>,
gpu: &GpuBackend,
plan: super::super::batch_queue::LanePlan,
remote_worker: Option<&WorkerSettings>,
local_lane_total: u32,
ui_weak: &Weak<MainWindow>,
remote_lane_done: &AtomicBool,
) {
let dispatchers = DispatcherGroup::new(shared.remote_lane_total as usize, remote_lane_done);
let dispatchers = &dispatchers;
std::thread::scope(|scope| {
if plan.run_remote {
for _ in 0..shared.remote_lane_total {
let remote_ui_weak = ui_weak.clone();
scope.spawn(move || {
run_remote_lane(
shared,
&remote_ui_weak,
remote_worker,
plan.fallback_to_local,
dispatchers,
);
});
}
}
if plan.run_local {
for _ in 0..local_lane_total {
let local_ui_weak = ui_weak.clone();
scope.spawn(move || {
run_local_lane(shared, gpu, &local_ui_weak, remote_lane_done);
});
}
}
});
}