use super::{
decisions::{
LiveRemoteAllowedInputs, RedispatchInputs, drag_held_effective, live_remote_may_run,
pose_settle_due, served_by_label, should_redispatch,
},
dispatch::{connection_is_stale, start_remote_render},
state::{Orchestrator, current_drag_held, current_pose, lock},
update::{continue_lane, redraw_from_epoch},
};
use crate::{
MainWindow, RemoteWorkerModel,
bridge::{
frame_cache::guide_pass::GuideCache,
remote::{
LiveDispatch,
handoff::{HandoffAction, HandoffEvent, HandoffMachine, HandoffState},
live_lane::LiveLane,
live_remote_dispatch,
remote_render::RemoteConnectionHandle,
},
render_thread::{RedrawGate, RenderContext},
},
gui::show_toast,
settings::{LiveComputeTarget, RemoteEndpoint, SettingsPersister},
};
use slint::ComponentHandle;
use std::{
sync::{Arc, Mutex, PoisonError, atomic::Ordering},
time::{Duration, Instant},
};
const SETTLE_DEBOUNCE: Duration = Duration::from_millis(600);
const POSE_SETTLE_DEBOUNCE: Duration = Duration::from_millis(250);
const DRAG_HELD_WATCHDOG: Duration = Duration::from_secs(20);
const POLL_INTERVAL_MS: u64 = 100;
#[derive(Clone)]
pub struct RemoteOrchestratorHandle(Arc<Mutex<Orchestrator>>);
impl RemoteOrchestratorHandle {
pub fn request_redisplay(&self, ui: &MainWindow, render_ctx: &Arc<Mutex<RenderContext>>) {
let Some((width, height)) = lock(&self.0)
.live_lane
.as_ref()
.map(|lane| lane.epoch().dimensions())
else {
return;
};
redraw_from_epoch(ui, render_ctx, &self.0, width, height);
}
}
#[must_use]
pub fn setup_remote_rendering(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
settings_store: &Arc<SettingsPersister>,
) -> (slint::Timer, RemoteOrchestratorHandle) {
let state = Arc::new(Mutex::new(Orchestrator {
handoff: HandoffMachine::new(),
remote_handle: None,
live_lane: None,
lane_scene: None,
current_request_id: None,
next_request_id: 1,
refusal_noted: false,
worker_label: String::new(),
scene_seen: (0, Instant::now()),
gave_up_for: None,
remote_connection: None,
display_only_refused: false,
remote_hdr: None,
last_pose: current_pose(render_ctx),
last_change_at: Instant::now(),
drag_was_held: false,
guide_cache: GuideCache::new(),
pending_guide_gen: None,
pending_denoise_gen: None,
last_denoised: None,
redraw_gate: RedrawGate::new(),
last_redraw_at: None,
live_remote_was_allowed: true,
}));
let handle = RemoteOrchestratorHandle(state.clone());
let timer = slint::Timer::default();
let ui_weak = ui.as_weak();
let render_ctx_poll = render_ctx.clone();
let settings_store_poll = settings_store.clone();
timer.start(
slint::TimerMode::Repeated,
Duration::from_millis(POLL_INTERVAL_MS),
move || {
let Some(ui) = ui_weak.upgrade() else {
return;
};
poll_tick(&ui, &render_ctx_poll, &settings_store_poll, &state);
},
);
(timer, handle)
}
fn poll_tick(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
settings_store: &Arc<SettingsPersister>,
state: &Arc<Mutex<Orchestrator>>,
) {
let remote_allowed = update_live_remote_suspension(render_ctx, state);
let pose = current_pose(render_ctx);
let held_raw = observe_drag_edge(render_ctx, state);
let changed = pose != lock(state).last_pose;
if changed {
lock(state).last_pose = pose;
lock(state).last_change_at = Instant::now();
let actions = lock(state).handoff.handle(HandoffEvent::OrientationChanged);
apply_actions(&actions, render_ctx, state);
sync_served_by_to_ui(ui, render_ctx, state);
sync_camera_moving_to_ctx(render_ctx, state);
return;
}
let (previewing, quiet_for) = {
let s = lock(state);
(
matches!(s.handoff.state(), HandoffState::Previewing),
s.last_change_at.elapsed(),
)
};
let held = drag_held_effective(held_raw, quiet_for, DRAG_HELD_WATCHDOG);
if pose_settle_due(previewing, held, quiet_for, POSE_SETTLE_DEBOUNCE) && remote_allowed {
let endpoint = settle_endpoint(ui, render_ctx, settings_store, state);
drop_stale_connection(state, endpoint.as_ref());
let actions = lock(state).handoff.handle(HandoffEvent::SettleElapsed {
worker_available: endpoint.is_some(),
});
apply_actions(&actions, render_ctx, state);
if let Some(endpoint) = endpoint {
start_remote_render(ui, render_ctx, endpoint, state);
}
}
sync_camera_moving_to_ctx(render_ctx, state);
reconcile_lane_with_ctx(render_ctx, state);
if remote_allowed {
resume_idle_lane(ui, render_ctx, state);
}
maybe_redispatch(ui, render_ctx, settings_store, state, held);
reconcile_served_by_after_release(ui, render_ctx, state);
sync_served_by_to_ui(ui, render_ctx, state);
}
fn observe_drag_edge(
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) -> bool {
let held_raw = current_drag_held(render_ctx);
let was_held = std::mem::replace(&mut lock(state).drag_was_held, held_raw);
if was_held && !held_raw {
lock(state).last_change_at = Instant::now();
}
held_raw
}
#[must_use]
pub(super) fn live_remote_allowed(render_ctx: &Arc<Mutex<RenderContext>>) -> bool {
let ctx = render_ctx.lock().unwrap_or_else(PoisonError::into_inner);
live_remote_may_run(LiveRemoteAllowedInputs {
tab_visible: ctx.tab_visible,
paused: ctx.paused,
export_active: ctx.export_active(),
})
}
#[must_use]
fn update_live_remote_suspension(
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) -> bool {
let allowed = live_remote_allowed(render_ctx);
let was_allowed = std::mem::replace(&mut lock(state).live_remote_was_allowed, allowed);
if was_allowed && !allowed {
suspend_live_remote(render_ctx, state);
}
allowed
}
fn suspend_live_remote(render_ctx: &Arc<Mutex<RenderContext>>, state: &Arc<Mutex<Orchestrator>>) {
let Some(display_only) = lock(state)
.live_lane
.as_ref()
.map(LiveLane::is_display_only)
else {
return;
};
if display_only {
apply_actions(
&[
HandoffAction::SendCancelToWorker,
HandoffAction::DiscardRemotePartial,
],
render_ctx,
state,
);
if matches!(
lock(state).handoff.state(),
HandoffState::Settling | HandoffState::RemoteRendering
) {
let actions = lock(state).handoff.handle(HandoffEvent::RemoteFailed);
apply_actions(&actions, render_ctx, state);
}
} else {
apply_actions(&[HandoffAction::SendCancelToWorker], render_ctx, state);
}
}
fn maybe_redispatch(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
settings_store: &Arc<SettingsPersister>,
state: &Arc<Mutex<Orchestrator>>,
drag_held: bool,
) {
let (generation, tab_visible, suspended, epoch_in_ctx, mode) = {
let mut ctx = render_ctx.lock().unwrap_or_else(PoisonError::into_inner);
(
ctx.scene_generation(),
ctx.tab_visible,
ctx.paused || ctx.export_active(),
ctx.live_epoch.is_some(),
ctx.live_compute_target,
)
};
if matches!(mode, LiveComputeTarget::LocalOnly) {
return;
}
let now = Instant::now();
let inputs = {
let mut s = lock(state);
if s.scene_seen.0 != generation {
s.scene_seen = (generation, now);
}
RedispatchInputs {
handoff_idle: matches!(s.handoff.state(), HandoffState::Idle),
tab_visible,
epoch_live: epoch_in_ctx || s.live_lane.is_some(),
scene_stable: now.saturating_duration_since(s.scene_seen.1) >= SETTLE_DEBOUNCE,
gave_up_for_this_scene: s.gave_up_for == Some((generation, mode)),
suspended,
drag_held,
}
};
if !should_redispatch(inputs) {
return;
}
let Some(endpoint) = settle_endpoint(ui, render_ctx, settings_store, state) else {
return;
};
drop_stale_connection(state, Some(&endpoint));
let actions = lock(state).handoff.handle(HandoffEvent::Redispatch);
apply_actions(&actions, render_ctx, state);
start_remote_render(ui, render_ctx, endpoint, state);
}
fn drop_stale_connection(state: &Arc<Mutex<Orchestrator>>, endpoint: Option<&RemoteEndpoint>) {
let mut s = lock(state);
if connection_is_stale(
s.remote_connection
.as_ref()
.map(RemoteConnectionHandle::worker),
endpoint.map(|e| &e.connection),
) {
s.remote_connection = None;
s.display_only_refused = false;
s.remote_hdr = None;
}
}
fn settle_endpoint(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
settings_store: &Arc<SettingsPersister>,
state: &Arc<Mutex<Orchestrator>>,
) -> Option<RemoteEndpoint> {
let endpoint = settings_store.snapshot().settings.remote;
let remote_hdr = lock(state).remote_hdr;
let decision = {
let ctx = render_ctx.lock().unwrap_or_else(PoisonError::into_inner);
live_remote_dispatch(
ctx.live_compute_target,
endpoint.is_some(),
ctx.env_map.as_ref(),
remote_hdr,
)
};
match decision {
LiveDispatch::Remote => {
lock(state).refusal_noted = false;
endpoint
}
LiveDispatch::Local => None,
LiveDispatch::Refused(refusal) => {
let first_time = !std::mem::replace(&mut lock(state).refusal_noted, true);
if first_time {
show_toast(ui, refusal.live_note(), "info");
}
None
}
}
}
fn reconcile_lane_with_ctx(
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) {
let ctx_epoch = render_ctx
.lock()
.unwrap_or_else(PoisonError::into_inner)
.live_epoch
.clone();
let mut s = lock(state);
let released = s.live_lane.as_ref().is_some_and(|lane| {
ctx_epoch
.as_ref()
.is_none_or(|epoch| !Arc::ptr_eq(epoch, lane.epoch()))
});
if !released {
return;
}
if let Some(handle) = s.remote_handle.take() {
handle.cancel();
}
if let Some(mut lane) = s.live_lane.take() {
lane.abandon();
}
s.lane_scene = None;
if matches!(
s.handoff.state(),
HandoffState::Settling | HandoffState::RemoteRendering
) {
s.handoff.handle(HandoffEvent::RemoteFailed);
}
}
fn resume_idle_lane(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) {
let idle = lock(state)
.live_lane
.as_ref()
.is_some_and(LiveLane::is_idle);
if idle {
continue_lane(ui, &ui.as_weak(), render_ctx, state);
}
}
fn reconcile_served_by_after_release(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) {
let still_remote_active = render_ctx
.lock()
.unwrap_or_else(PoisonError::into_inner)
.remote_active;
if still_remote_active {
return;
}
let served_by_remote = matches!(
lock(state).handoff.served_by(),
crate::bridge::remote::handoff::ImageSource::Remote
);
if served_by_remote {
lock(state).handoff.handle(HandoffEvent::SceneInvalidated);
sync_served_by_to_ui(ui, render_ctx, state);
}
}
fn sync_camera_moving_to_ctx(
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) {
let moving = matches!(lock(state).handoff.state(), HandoffState::Previewing);
render_ctx
.lock()
.unwrap_or_else(PoisonError::into_inner)
.camera_moving = moving;
}
pub(super) fn sync_served_by_to_ui(
ui: &MainWindow,
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) {
let (mode, epoch) = {
let ctx = render_ctx.lock().unwrap_or_else(PoisonError::into_inner);
let epoch = ctx.remote_active.then(|| ctx.live_epoch.clone()).flatten();
(ctx.effective_live_target(), epoch)
};
let combined_remote_samples = epoch.map_or(0, |epoch| epoch.remote_done());
let (handoff_says_remote, worker_label) = {
let s = lock(state);
(
matches!(
s.handoff.served_by(),
crate::bridge::remote::handoff::ImageSource::Remote
),
s.worker_label.clone(),
)
};
let label = served_by_label(
mode,
handoff_says_remote,
combined_remote_samples,
&worker_label,
);
let model = ui.global::<RemoteWorkerModel>();
if model.get_served_by_remote() != label.is_some() {
model.set_served_by_remote(label.is_some());
}
if let Some(label) = label
&& model.get_served_by_worker_name().as_str() != label
{
model.set_served_by_worker_name(label.into());
}
}
pub(super) fn apply_actions(
actions: &[HandoffAction],
render_ctx: &Arc<Mutex<RenderContext>>,
state: &Arc<Mutex<Orchestrator>>,
) {
for action in actions {
match action {
HandoffAction::DiscardLocalPreview | HandoffAction::SendRenderRequestToWorker => {}
HandoffAction::SendCancelToWorker => {
if let Some(handle) = &lock(state).remote_handle {
handle.cancel();
}
}
HandoffAction::DiscardRemotePartial => {
let mut s = lock(state);
s.remote_handle = None;
if let Some(mut lane) = s.live_lane.take() {
lane.abandon();
}
s.lane_scene = None;
if let Some(pending) = s.pending_guide_gen.take() {
pending.cancel.store(true, Ordering::Relaxed);
}
s.pending_denoise_gen = None;
s.last_denoised = None;
drop(s);
render_ctx
.lock()
.unwrap_or_else(PoisonError::into_inner)
.release_remote();
}
}
}
}