use crate::{
BatchModel,
bridge::{
preview_render::{PREVIEW_LIGHT_PITCH, PREVIEW_LIGHT_YAW},
remote::remote_render,
},
gui::{
batch::{
batch_queue::WorkQueue,
preview::{RI_MATCH_TOLERANCE, seeded_random_unit, target_ri_for_design},
},
library::detail::reconstruct_planes,
},
settings::WorkerSettings,
};
use indicatrix::{
color::metrics::{PROFILE_AZIMUTHS_DEG, evaluate_full_axis_profile_at_azimuth},
geometry::plane::GpuFacetPlane,
optics::{materials::GemMaterial, raytracer::LightingPreset},
};
use indicatrix_net::{
SceneState,
messages::{AxisTiltCurves as WireAxisTiltCurves, TiltCurvesRequest, TiltCurvesResponse},
};
use indicatrix_vault::model::tilt_curves::{
AxisTiltCurves as StorageAxisTiltCurves, TILT_CURVE_AXIS_COUNT, TILT_CURVE_POINTS_PER_AXIS,
TiltPerformanceCurves,
};
use slint::{ComponentHandle, Weak};
use std::{
any::Any,
sync::{
Mutex, PoisonError,
atomic::{AtomicBool, AtomicU32, Ordering},
},
thread,
time::{Duration, SystemTime, UNIX_EPOCH},
};
use tracing::warn;
const REMOTE_SCENE_MAX_BOUNCES: u32 = 12;
const TILT_REQUEST_ID: u32 = 1;
const LOCAL_IDLE_POLL: Duration = Duration::from_millis(15);
pub(super) struct BatchContext<'a> {
pub(super) db: &'a Mutex<indicatrix_vault::db::sqlite::Database>,
pub(super) material_candidates:
&'a [indicatrix_vault::model::material_match::RiPresetCandidate],
}
struct ResolvedDesign {
title: String,
planes: Vec<GpuFacetPlane>,
material: GemMaterial,
}
fn resolve_design(ctx: &BatchContext<'_>, entry_id: i64) -> Option<ResolvedDesign> {
let full = {
let guard = ctx.db.lock().unwrap_or_else(PoisonError::into_inner);
guard.get_diagram_full(entry_id)
};
let Ok(Some(full)) = full else {
return None;
};
let facet_specs: Vec<indicatrix::geometry::cuts::FacetSpec> = full
.angle_settings
.iter()
.map(|a| indicatrix::geometry::cuts::FacetSpec {
facet: a.facet.clone(),
angle: a.angle.clone(),
index: a.index.clone(),
notes: a.notes.clone(),
})
.collect();
let planes = reconstruct_planes(
full.shape.as_deref(),
full.index_gear.as_deref(),
&facet_specs,
);
if planes.is_empty() {
return None;
}
let target_ri = target_ri_for_design(&full);
let material_name = {
let guard = ctx.db.lock().unwrap_or_else(PoisonError::into_inner);
let mut rng = seeded_random_unit(entry_id);
guard.ensure_preview_material(
entry_id,
target_ri,
ctx.material_candidates,
RI_MATCH_TOLERANCE,
&mut rng,
)
};
let Ok(Some(material_name)) = material_name else {
return None;
};
let material = GemMaterial::by_name(&material_name)?;
Some(ResolvedDesign {
title: full.title,
planes,
material,
})
}
fn wire_axis_to_storage(axis: WireAxisTiltCurves) -> Result<StorageAxisTiltCurves, String> {
let to_array = |v: Vec<f32>, name: &str| -> Result<[f32; TILT_CURVE_POINTS_PER_AXIS], String> {
v.try_into().map_err(|v: Vec<f32>| {
format!(
"{name} has {} points, expected {TILT_CURVE_POINTS_PER_AXIS}",
v.len()
)
})
};
Ok(StorageAxisTiltCurves {
brilliance_pct: to_array(axis.brilliance_pct, "brilliance_pct")?,
extinction_pct: to_array(axis.extinction_pct, "extinction_pct")?,
windowing_pct: to_array(axis.windowing_pct, "windowing_pct")?,
})
}
fn fetch_tilt_curves_remote(
worker: &WorkerSettings,
planes: &[GpuFacetPlane],
material: &GemMaterial,
cancel: &AtomicBool,
) -> Option<TiltPerformanceCurves> {
if cancel.load(Ordering::Relaxed) {
return None;
}
let scene = SceneState {
width: 1,
height: 1,
yaw: 0.0,
pitch: 0.0,
distance: 1.0,
light_yaw: PREVIEW_LIGHT_YAW,
light_pitch: PREVIEW_LIGHT_PITCH,
exposure: 1.0,
max_bounces: REMOTE_SCENE_MAX_BOUNCES,
lighting_preset: LightingPreset::RingLights,
material: material.clone(),
planes: planes.to_vec(),
girdle_frosted: false,
};
let (mut stream, welcome) = remote_render::connect_and_handshake(worker).ok()?;
if !welcome.tilt_curves {
return None;
}
let request = TiltCurvesRequest {
request_id: TILT_REQUEST_ID,
scene,
};
indicatrix_net::client::send_tilt_curves_request(&mut stream, &request).ok()?;
match indicatrix_net::client::recv_tilt_curves_response(&mut stream) {
Ok(TiltCurvesResponse::Curves(result)) => {
let mut axes_vec = Vec::with_capacity(result.axes.len());
for axis in result.axes {
match wire_axis_to_storage(axis) {
Ok(axis) => axes_vec.push(axis),
Err(e) => {
warn!("Remote tilt-curve reply had a malformed axis: {e}");
return None;
}
}
}
let axes: [StorageAxisTiltCurves; TILT_CURVE_AXIS_COUNT] = axes_vec.try_into().ok()?;
Some(TiltPerformanceCurves { axes })
}
Ok(TiltCurvesResponse::Cancelled { .. } | TiltCurvesResponse::Error(_)) | Err(_) => None,
}
}
fn compute_tilt_curves_locally(
planes: &[GpuFacetPlane],
material: &GemMaterial,
cancel: &AtomicBool,
) -> Option<TiltPerformanceCurves> {
let mut axes = Vec::with_capacity(PROFILE_AZIMUTHS_DEG.len());
for &azimuth_deg in &PROFILE_AZIMUTHS_DEG {
if cancel.load(Ordering::Relaxed) {
return None;
}
let (brilliance_pct, extinction_pct, windowing_pct) = evaluate_full_axis_profile_at_azimuth(
planes,
material,
azimuth_deg,
PREVIEW_LIGHT_YAW,
PREVIEW_LIGHT_PITCH,
);
axes.push(StorageAxisTiltCurves {
brilliance_pct,
extinction_pct,
windowing_pct,
});
}
let axes: [StorageAxisTiltCurves; TILT_CURVE_AXIS_COUNT] = axes.try_into().ok()?;
Some(TiltPerformanceCurves { axes })
}
fn save_curves(ctx: &BatchContext<'_>, entry_id: i64, curves: &TiltPerformanceCurves) -> bool {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |d| d.as_secs() as i64);
let guard = ctx.db.lock().unwrap_or_else(PoisonError::into_inner);
guard.save_tilt_curves(entry_id, curves, None, now).is_ok()
}
fn process_local_entry(ctx: &BatchContext<'_>, entry_id: i64, cancel: &AtomicBool) -> bool {
let Some(resolved) = resolve_design(ctx, entry_id) else {
return false;
};
let Some(curves) = compute_tilt_curves_locally(&resolved.planes, &resolved.material, cancel)
else {
return false;
};
save_curves(ctx, entry_id, &curves)
}
fn process_remote_entry(
ctx: &BatchContext<'_>,
worker: Option<&WorkerSettings>,
entry_id: i64,
cancel: &AtomicBool,
mut on_title: impl FnMut(&str),
) -> bool {
let Some(resolved) = resolve_design(ctx, entry_id) else {
return false;
};
on_title(&resolved.title);
let Some(worker) = worker else {
return false;
};
let Some(curves) =
fetch_tilt_curves_remote(worker, &resolved.planes, &resolved.material, cancel)
else {
return false;
};
save_curves(ctx, entry_id, &curves)
}
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())
}
#[derive(Default)]
pub(super) struct Tally {
pub(super) computed: AtomicU32,
pub(super) failed: AtomicU32,
}
#[derive(Default, Clone)]
pub(super) struct LiveProgress {
completed: u32,
local_active: u32,
remote_title: String,
remote_active: bool,
}
fn push_progress(
ui_weak: &Weak<crate::MainWindow>,
design_total: u32,
local_lane_total: u32,
progress: &Mutex<LiveProgress>,
) {
let snapshot = progress
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clone();
let ui_weak = ui_weak.clone();
let _ = ui_weak.upgrade_in_event_loop(move |ui| {
ui.global::<BatchModel>()
.set_tilt_design_index(snapshot.completed as i32);
ui.global::<BatchModel>()
.set_tilt_design_total(design_total as i32);
ui.global::<BatchModel>()
.set_tilt_local_active(snapshot.local_active as i32);
ui.global::<BatchModel>()
.set_tilt_local_lane_total(local_lane_total as i32);
ui.global::<BatchModel>()
.set_tilt_remote_title(snapshot.remote_title.into());
ui.global::<BatchModel>()
.set_tilt_remote_active(snapshot.remote_active);
});
}
fn set_remote_status(progress: &Mutex<LiveProgress>, active: bool, title: &str) {
let mut p = progress.lock().unwrap_or_else(PoisonError::into_inner);
p.remote_active = active;
p.remote_title = title.to_string();
}
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<i64>,
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) fn run_local_lane(
shared: &LaneShared<'_>,
ui_weak: &Weak<crate::MainWindow>,
remote_lane_done: &AtomicBool,
) {
loop {
if shared.cancel.load(Ordering::Relaxed) {
break;
}
let Some(entry_id) = shared.queue.claim_local() else {
if remote_lane_done.load(Ordering::Acquire) {
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 result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
process_local_entry(shared.ctx, entry_id, shared.cancel)
}));
let saved = match result {
Ok(saved) => saved,
Err(payload) => {
warn!(
"Tilt-curve computation panicked for entry {entry_id}: {}",
panic_message(&*payload)
);
false
}
};
if saved {
shared.tally.computed.fetch_add(1, Ordering::Relaxed);
} else {
shared.tally.failed.fetch_add(1, Ordering::Relaxed);
}
adjust_local_active(shared.progress, -1);
increment_completed(shared.progress);
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
}
}
pub(super) fn run_remote_lane(
shared: &LaneShared<'_>,
ui_weak: &Weak<crate::MainWindow>,
worker: Option<&WorkerSettings>,
fallback_to_local: bool,
remote_lane_done: &AtomicBool,
) {
loop {
if shared.cancel.load(Ordering::Relaxed) {
break;
}
let Some(entry_id) = shared.queue.claim_shared() else {
break;
};
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
process_remote_entry(shared.ctx, worker, entry_id, shared.cancel, |title| {
set_remote_status(shared.progress, true, title);
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
})
}));
let saved = match result {
Ok(saved) => saved,
Err(payload) => {
warn!(
"Remote tilt-curve dispatch panicked for entry {entry_id}: {}",
panic_message(&*payload)
);
false
}
};
if saved {
shared.tally.computed.fetch_add(1, Ordering::Relaxed);
increment_completed(shared.progress);
} else if fallback_to_local {
shared.queue.return_to_local(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,
);
}
set_remote_status(shared.progress, false, "");
push_progress(
ui_weak,
shared.design_total,
shared.local_lane_total,
shared.progress,
);
remote_lane_done.store(true, Ordering::Release);
}
pub(super) fn run_batch_lanes(
shared: &LaneShared<'_>,
plan: super::super::batch_queue::LanePlan,
remote_worker: Option<&WorkerSettings>,
local_lane_total: u32,
ui_weak: &Weak<crate::MainWindow>,
remote_lane_done: &AtomicBool,
) {
std::thread::scope(|scope| {
if plan.run_remote {
let remote_ui_weak = ui_weak.clone();
scope.spawn(move || {
run_remote_lane(
shared,
&remote_ui_weak,
remote_worker,
plan.fallback_to_local,
remote_lane_done,
);
});
}
if plan.run_local {
for _ in 0..local_lane_total {
let local_ui_weak = ui_weak.clone();
scope.spawn(move || {
run_local_lane(shared, &local_ui_weak, remote_lane_done);
});
}
}
});
}