Skip to main content

regain_core/
session.rs

1use crate::timing::{
2    ACKNOWLEDGED_CONTROLS, CLOSE_SECONDS, MAX_READ_RETRY_OVERHEAD_SECONDS, PERSISTENT_CONTROLS,
3    USB_BIND_SECONDS, USB_REBIND_PAUSE_SECONDS, USB_RESET_SECONDS,
4};
5use crate::*;
6use anyhow::{Context, Result, ensure};
7use chrono::Utc;
8use serde_json::{Value, json};
9use std::{
10    collections::BTreeMap,
11    sync::{Arc, Mutex},
12    time::Duration,
13};
14use tokio::time::Instant;
15
16#[derive(Clone, Copy, PartialEq, Eq)]
17enum OpenPurpose {
18    Acquisition,
19    ThermalShutdown,
20}
21/// One owner per camera. All public entry points are serialized by the frontend.
22/// Status and queued controls remain readable while capture owns the worker.
23pub struct Session {
24    pub status: SharedStatus,
25    pub selection: Selection,
26    runtime: Runtime,
27    log: Diagnostic,
28    worker: Option<Worker>,
29    direct: bool,
30    ever_opened: bool,
31    applied: BTreeMap<i32, i64>,
32    recovery_temperature: Option<f64>,
33    recovery_power: Option<i64>,
34    recovery_target: Option<i64>,
35    settle_required: bool,
36    cooling_seeded: bool,
37    usb_target: Option<String>,
38    cooling: cooling::Mailbox,
39}
40impl Session {
41    pub fn new(selection: Selection, runtime: Runtime, log: Diagnostic) -> Result<Self> {
42        selection.recovery.validate()?;
43        let status = Arc::new(Mutex::new(Status::default()));
44        let observed = status.clone();
45        let log: Diagnostic = Arc::new(move |level, event, message| {
46            // The direct worker reports transfer retries while its status call
47            // is still pending. Surface those records without parsing prose or
48            // changing any recovery decisions. Final frame metadata reconciles
49            // successful reads if stderr delivery lagged behind the reply.
50            if matches!(event, "transfer.retry" | "transfer.exhausted") {
51                let mut state = observed.lock().unwrap();
52                if matches!(
53                    state.phase.as_str(),
54                    "Starting exposure" | "Exposing" | "Downloading"
55                ) {
56                    if event == "transfer.retry" {
57                        state.retry.usb_reads = state.retry.usb_reads.saturating_add(1);
58                    }
59                    state.retry.last_failure = Some(message.into());
60                }
61            }
62            log(level, event, message);
63        });
64        Ok(Self {
65            status,
66            direct: selection.direct,
67            selection,
68            runtime,
69            log,
70            worker: None,
71            ever_opened: false,
72            applied: BTreeMap::new(),
73            recovery_temperature: None,
74            recovery_power: None,
75            recovery_target: None,
76            settle_required: false,
77            cooling_seeded: false,
78            usb_target: None,
79            cooling: cooling::Mailbox::default(),
80        })
81    }
82    fn emit(&self, level: &str, event: &str, message: impl AsRef<str>) {
83        let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
84            (self.log)(level, event, message.as_ref())
85        }));
86    }
87    fn phase(&self, phase: impl Into<String>) {
88        let phase = phase.into();
89        self.status.lock().unwrap().phase = phase.clone();
90        self.emit(
91            if matches!(
92                phase.as_str(),
93                "Idle" | "Starting exposure" | "Exposing" | "Downloading"
94            ) {
95                "debug"
96            } else {
97                "info"
98            },
99            "session.phase",
100            phase,
101        );
102    }
103    pub fn snapshot(&self) -> Status {
104        self.status.lock().unwrap().clone()
105    }
106    pub fn cooling(&self) -> cooling::CoolingHandle {
107        cooling::CoolingHandle {
108            mailbox: self.cooling.clone(),
109            status: self.status.clone(),
110        }
111    }
112    /// Drive one reserved command while idle. Capture drives the same mailbox
113    /// at safe worker checkpoints. Frontends retain/serialize this operation.
114    pub async fn service_cooling(&mut self, token: &CancellationToken) -> Result<()> {
115        self.service_cooling_with(&mut self.settings(), token).await
116    }
117    async fn service_cooling_with(
118        &mut self,
119        settings: &mut BTreeMap<i32, i64>,
120        token: &CancellationToken,
121    ) -> Result<()> {
122        let Some(request) = self.cooling.claim() else {
123            return Ok(());
124        };
125        let (kind, value, deadline) = {
126            let request = request.lock().unwrap();
127            (request.kind, request.value, request.deadline)
128        };
129        let deadline = deadline.min(
130            Instant::now()
131                + Duration::from_secs_f64(self.selection.recovery.command_timeout_seconds),
132        );
133        let mailbox = self.cooling.clone();
134        if let Err(error) = cooling::validate(&self.status, kind, value) {
135            let _ = mailbox.finish(&request, Err(error), || {});
136            return Ok(());
137        }
138        if token.is_cancelled() {
139            let _ = mailbox.finish(&request, Err(cooling::CoolingError::Cancelled), || {});
140            return Err(Failure::Cancelled.into());
141        }
142        let remaining = deadline
143            .saturating_duration_since(Instant::now())
144            .as_secs_f64();
145        if remaining == 0. {
146            let _ = mailbox.finish(&request, Err(cooling::CoolingError::Expired), || {});
147            return Ok(());
148        }
149        if !mailbox.dispatch(&request) {
150            return Ok(());
151        }
152        let result = self
153            .write_control(kind, value, deadline, token, false)
154            .await;
155        let observation = result.as_ref().ok().copied();
156        let result = mailbox.finish(&request, result.map(|observed| observed.value), || {
157            self.applied.insert(kind, value);
158            settings.insert(kind, value);
159            let mut state = self.status.lock().unwrap();
160            state.values.insert(kind, value);
161            if let Some(observation) = observation {
162                state.observations.insert(kind, observation);
163            }
164        });
165        match result {
166            Err(error @ cooling::CoolingError::Uncertain { .. }) => {
167                self.control_result(Err(error)).await.map(|_| ())
168            }
169            Err(cooling::CoolingError::Cancelled) => Err(Failure::Cancelled.into()),
170            _ => Ok(()), // Unsent expiry/failure belongs to the command receipt.
171        }
172    }
173    /// Acknowledge gain/offset while idle, using the same serialized worker as
174    /// capture. The caller must retain this future through dispatch/cleanup.
175    /// Unlike queue_control, success means write plus readback, not desired intent.
176    /// The absolute deadline includes time spent waiting for owner admission.
177    pub async fn set_imaging_control(
178        &mut self,
179        kind: i32,
180        value: i64,
181        deadline: Instant,
182        token: &CancellationToken,
183    ) -> Result<i64> {
184        ensure!(
185            matches!(kind, 0 | 5),
186            invalid("Only gain and offset are supported")
187        );
188        let state = self.snapshot();
189        ensure!(
190            state.connected && state.control_connection_available && self.worker.is_some(),
191            invalid("Camera controls are unavailable")
192        );
193        ensure!(
194            !self.cooling().pending(),
195            invalid("A cooler command is pending")
196        );
197        let cap = state
198            .controls
199            .get(&kind)
200            .ok_or_else(|| invalid("Control is unavailable"))?;
201        ensure!(
202            cap.kind == kind
203                && cap.writable
204                && cap.min <= cap.max
205                && (cap.min..=cap.max).contains(&value),
206            invalid("Control value is outside writable capabilities")
207        );
208        if token.is_cancelled() {
209            return Err(Failure::Cancelled.into());
210        }
211        if Instant::now() >= deadline {
212            return Err(cooling::CoolingError::Expired.into());
213        }
214        let deadline = deadline.min(
215            Instant::now()
216                + Duration::from_secs_f64(self.selection.recovery.command_timeout_seconds),
217        );
218        let result = self.write_control(kind, value, deadline, token, true).await;
219        let observation = self.control_result(result).await?;
220        self.applied.insert(kind, observation.value);
221        let mut state = self.status.lock().unwrap();
222        state.values.insert(kind, observation.value);
223        state.observations.insert(kind, observation);
224        Ok(observation.value)
225    }
226    async fn write_control(
227        &mut self,
228        kind: i32,
229        value: i64,
230        deadline: Instant,
231        token: &CancellationToken,
232        allow_sdk_offset_clamp: bool,
233    ) -> std::result::Result<ControlObservation, cooling::CoolingError> {
234        let mut set_acknowledged = false;
235        let result = async {
236            self.worker
237                .as_mut()
238                .context("Camera worker is disconnected")?
239                .control_call(
240                    "set",
241                    json!({"control":kind,"value":value}),
242                    deadline,
243                    token,
244                )
245                .await?;
246            set_acknowledged = true;
247            self.emit(
248                "debug",
249                "control.write_acknowledged",
250                format!("Control {kind}; readback pending"),
251            );
252            let remaining = deadline
253                .saturating_duration_since(Instant::now())
254                .as_secs_f64();
255            ensure!(remaining > 0., "Control deadline expired before readback");
256            let request_started = Instant::now();
257            let reply = self
258                .worker
259                .as_mut()
260                .context("Camera worker is disconnected")?
261                .control_call("get-observation", json!({"control":kind}), deadline, token)
262                .await?
263                .0;
264            let observation = serde_json::from_value::<ControlObservationReply>(reply)
265                .map_err(|_| invalid("Invalid control readback"))?
266                .normalize(request_started)?;
267            let actual = observation.value;
268            if actual != value {
269                let cap = self.snapshot().controls.get(&kind).cloned();
270                ensure!(
271                    allow_sdk_offset_clamp
272                        && !self.direct
273                        && kind == 5
274                        && cap.is_some_and(|cap| (cap.min..=cap.max).contains(&actual)),
275                    "Control readback differs from request"
276                );
277                self.emit(
278                    "info",
279                    "control.clamped",
280                    format!("SDK applied offset {actual} instead of {value}"),
281                );
282            }
283            ensure!(
284                Instant::now() < deadline,
285                "Control readback exceeded its deadline"
286            );
287            Ok::<_, anyhow::Error>(observation)
288        }
289        .await;
290        result.map_err(|error| {
291            if !set_acknowledged
292                && matches!(
293                    error.downcast_ref::<cooling::CoolingError>(),
294                    Some(cooling::CoolingError::Expired)
295                )
296            {
297                cooling::CoolingError::Expired
298            } else if !set_acknowledged
299                && matches!(error.downcast_ref::<Failure>(), Some(Failure::Cancelled))
300            {
301                cooling::CoolingError::Cancelled
302            } else {
303                let code = match error.downcast_ref::<Failure>() {
304                    Some(Failure::Worker { code, .. } | Failure::UncertainControl { code, .. }) => {
305                        *code
306                    }
307                    _ => None,
308                };
309                cooling::CoolingError::Uncertain {
310                    message: format!("{error:#}"),
311                    code,
312                }
313            }
314        })
315    }
316    async fn control_result(
317        &mut self,
318        result: std::result::Result<ControlObservation, cooling::CoolingError>,
319    ) -> Result<ControlObservation> {
320        match result {
321            Err(cooling::CoolingError::Uncertain { message, code }) => {
322                self.invalidate().await;
323                Err(Failure::UncertainControl { message, code }.into())
324            }
325            Err(cooling::CoolingError::Cancelled) => {
326                // Worker cancellation can retire an otherwise untouched process.
327                // Withdraw its core availability before another owner operation.
328                self.invalidate().await;
329                Err(Failure::Cancelled.into())
330            }
331            Err(error) => Err(error.into()),
332            Ok(actual) => Ok(actual),
333        }
334    }
335    /// Opt into Regain-owned WB. The caller serializes this with capture, like
336    /// all Session operations. Settings and effective AWB gains survive recovery.
337    pub async fn set_white_balance(
338        &mut self,
339        settings: crate::white_balance::Settings,
340        token: &CancellationToken,
341    ) -> Result<()> {
342        settings.gains.validate()?;
343        let result = self
344            .call(
345                "white-balance",
346                serde_json::to_value(settings)?,
347                None,
348                token,
349            )
350            .await?
351            .0;
352        self.status.lock().unwrap().white_balance =
353            Some(serde_json::from_value(result["settings"].clone())?);
354        // Legacy WB/flip controls must not be replayed over managed WB.
355        for kind in [3, 4, 9] {
356            self.applied.remove(&kind);
357        }
358        Ok(())
359    }
360    pub async fn simulate_read_failures(
361        &mut self,
362        count: u32,
363        token: &CancellationToken,
364    ) -> Result<()> {
365        ensure!(
366            self.runtime.simulate && self.direct,
367            "Read fault injection requires a simulated direct camera"
368        );
369        self.call(
370            "simulate-read-failures",
371            json!({"count":count}),
372            None,
373            token,
374        )
375        .await?;
376        Ok(())
377    }
378    pub fn seed_recovery(&mut self, temperature: Option<f64>, power: Option<i64>) {
379        self.recovery_temperature = temperature.filter(|t| t.is_finite());
380        self.recovery_power = power;
381        self.settle_required = true;
382    }
383    pub fn seed_recovery_target(&mut self, target: Option<i64>) {
384        self.recovery_target = target.filter(|v| (-40..=30).contains(v));
385    }
386    pub fn queue_control(status: &SharedStatus, kind: i32, value: i64) -> Result<()> {
387        let mut state = status.lock().unwrap();
388        ensure!(
389            state.white_balance.is_none() || !matches!(kind, 3 | 4 | 9),
390            Failure::Invalid("WB and flip controls are owned by managed white balance".into())
391        );
392        let cap = state
393            .controls
394            .get(&kind)
395            .ok_or_else(|| invalid(format!("Control {kind} unavailable")))?;
396        ensure!(
397            cap.writable && value >= cap.min && value <= cap.max,
398            Failure::Invalid(format!(
399                "Control {kind} must be writable and between {} and {}",
400                cap.min, cap.max
401            ))
402        );
403        ensure!(
404            PERSISTENT_CONTROLS.contains(&kind),
405            Failure::Invalid("Control is not a persistent imaging setting".into())
406        );
407        state.values.insert(kind, value);
408        Ok(())
409    }
410    async fn call(
411        &mut self,
412        method: &str,
413        params: Value,
414        seconds: Option<f64>,
415        token: &CancellationToken,
416    ) -> Result<(Value, Vec<u8>)> {
417        self.worker
418            .as_mut()
419            .context("Camera worker is disconnected")?
420            .call(
421                method,
422                params,
423                seconds.unwrap_or(self.selection.recovery.command_timeout_seconds),
424                token,
425            )
426            .await
427    }
428    async fn delay(&self, seconds: f64, token: &CancellationToken) -> Result<()> {
429        tokio::select! {biased;_=token.cancelled()=>Err(Failure::Cancelled.into()),_=tokio::time::sleep(Duration::from_secs_f64(seconds))=>Ok(())}
430    }
431    async fn invalidate(&mut self) {
432        self.cooling.retire();
433        if self.ever_opened && !self.settle_required {
434            let state = self.snapshot();
435            self.recovery_temperature = state.values.get(&8).map(|v| *v as f64 / 10.);
436            self.recovery_power = state.values.get(&15).copied();
437            self.recovery_target = self
438                .applied
439                .get(&16)
440                .copied()
441                .or_else(|| state.values.get(&16).copied());
442            self.settle_required = true;
443        }
444        self.status.lock().unwrap().control_connection_available = false;
445        if let Some(mut worker) = self.worker.take() {
446            worker.kill().await;
447        }
448        self.status.lock().unwrap().process_id = None;
449    }
450    async fn fallback(&mut self, reason: &str, token: &CancellationToken) -> Result<()> {
451        self.status.lock().unwrap().retry.last_failure = Some(reason.into());
452        self.emit(
453            "warning",
454            "backend.fallback",
455            format!("Switching to SDK fallback: {reason}"),
456        );
457        self.invalidate().await;
458        self.direct = false;
459        self.phase("SDK fallback reconnect delay");
460        self.delay(self.selection.recovery.reconnect_delay_seconds, token)
461            .await
462    }
463    pub async fn connect(&mut self, token: &CancellationToken) -> Result<()> {
464        let result = async {
465            match self.open(token).await {
466                Err(error) if self.direct && self.selection.sdk_fallback && retryable(&error) => {
467                    self.fallback(&error.to_string(), token).await?;
468                    self.open(token).await
469                }
470                other => other,
471            }?;
472            if self.selection.recovery.usb_reset_after_failures > 0 && self.usb_target.is_none() {
473                // First connection can learn the serial. Release its handle before
474                // binding the physical device, then reopen that exact serial.
475                self.invalidate().await;
476                self.bind_usb(token).await?;
477                self.open(token).await?;
478            }
479            self.ever_opened = true;
480            self.status.lock().unwrap().connected = true;
481            self.phase("Idle");
482            Ok(())
483        }
484        .await;
485        if let Err(error) = &result {
486            if !token.is_cancelled() {
487                self.status.lock().unwrap().retry.last_failure = Some(format!("{error:#}"));
488            }
489            self.emit("warning", "connection.failed", format!("{error:#}"));
490            self.invalidate().await;
491        }
492        result
493    }
494    async fn bind_usb(&mut self, token: &CancellationToken) -> Result<()> {
495        if token.is_cancelled() {
496            return Err(Failure::Cancelled.into());
497        }
498        if self.runtime.simulate {
499            self.usb_target = Some("simulation".into());
500            return Ok(());
501        }
502        let serial = self
503            .selection
504            .serial
505            .as_deref()
506            .context("USB recovery requires a camera serial")?;
507        let args = [
508            "zwo",
509            "camera-direct",
510            "--usb-target",
511            &self.selection.name,
512            serial,
513        ];
514        let command = self.runtime.usb_command(&args, USB_BIND_SECONDS);
515        let encoded = tokio::select! { biased; _=token.cancelled()=>return Err(Failure::Cancelled.into()), r=command=>r? };
516        let target = regain_transport::usb::Target::decode(&encoded)?;
517        ensure!(
518            target.serial.eq_ignore_ascii_case(serial),
519            "USB recovery serial mismatch"
520        );
521        self.usb_target = Some(encoded);
522        self.emit(
523            "info",
524            "usb.bound",
525            "USB recovery bound to the selected camera's serial and physical location",
526        );
527        Ok(())
528    }
529    async fn reset_usb(&mut self, token: &CancellationToken) -> Result<()> {
530        if token.is_cancelled() {
531            return Err(Failure::Cancelled.into());
532        }
533        let target = self
534            .usb_target
535            .as_deref()
536            .context("No verified USB recovery target")?;
537        ensure!(
538            self.worker.is_none(),
539            "Close the camera worker before USB recovery"
540        );
541        self.phase("USB recovery");
542        self.emit(
543            "warning",
544            "usb.reset",
545            "Resetting the selected camera; the retained frame is abandoned",
546        );
547        if !self.runtime.simulate {
548            // Once dispatched, finish the bounded helper before honoring abort so
549            // a Linux port cycle can always re-enable the port.
550            self.runtime
551                .usb_command(
552                    &[
553                        "usb",
554                        if self.selection.recovery.usb_port_cycle {
555                            "cycle"
556                        } else {
557                            "reset"
558                        },
559                        target,
560                    ],
561                    USB_RESET_SECONDS,
562                )
563                .await?;
564        }
565        self.usb_target = None;
566        if token.is_cancelled() {
567            return Err(Failure::Cancelled.into());
568        }
569        self.phase("Waiting for USB camera");
570        let deadline = Instant::now() + Duration::from_secs(USB_BIND_SECONDS);
571        loop {
572            match self.bind_usb(token).await {
573                Ok(()) => break,
574                Err(e) if token.is_cancelled() || Instant::now() >= deadline => return Err(e),
575                Err(_) => self.delay(USB_REBIND_PAUSE_SECONDS, token).await?,
576            }
577        }
578        self.emit(
579            "info",
580            "usb.returned",
581            "The same camera serial is available after USB recovery",
582        );
583        Ok(())
584    }
585    async fn open(&mut self, token: &CancellationToken) -> Result<()> {
586        self.open_worker(token, OpenPurpose::Acquisition).await
587    }
588    async fn open_worker(&mut self, token: &CancellationToken, purpose: OpenPurpose) -> Result<()> {
589        if self.ever_opened && self.selection.serial.is_none() {
590            return Err(invalid(
591                "Automatic recovery requires a camera serial number",
592            ));
593        }
594        self.phase(if purpose == OpenPurpose::ThermalShutdown {
595            "Closing camera"
596        } else {
597            "Opening"
598        });
599        self.worker = Some(self.runtime.spawn(self.direct, self.log.clone()).await?);
600        self.applied.clear();
601        // Values retain desired recovery settings; old evidence cannot describe
602        // a replacement worker or its newly negotiated capabilities.
603        self.status.lock().unwrap().observations.clear();
604        self.cooling_seeded = false;
605        let (result, _) = self
606            .call(
607                "open",
608                json!({"name":self.selection.name,"serial":self.selection.serial}),
609                None,
610                token,
611            )
612            .await?;
613        let serial = result["serial"].as_str().map(str::to_owned);
614        ensure!(
615            self.selection.serial.is_none() || self.selection.serial == serial,
616            Failure::Invalid("Camera identity changed".into())
617        );
618        let info = &result["info"];
619        ensure!(
620            info["name"] == self.selection.name
621                && info["formats"]
622                    .as_array()
623                    .is_some_and(|a| a.contains(&json!(2))),
624            Failure::Invalid("Camera identity or RAW16 support changed".into())
625        );
626        let previous = self.snapshot();
627        ensure!(
628            !self.ever_opened
629                || (previous.info["width"] == info["width"]
630                    && previous.info["height"] == info["height"]),
631            Failure::Invalid("Camera geometry changed".into())
632        );
633        let mut controls: Vec<Control> = serde_json::from_value(result["controls"].clone())
634            .map_err(|_| invalid("Invalid camera controls"))?;
635        if self.direct
636            && self.selection.sdk_fallback
637            && matches!(
638                self.selection.name.as_str(),
639                "ZWO ASI2600MM Duo" | "ZWO ASI220MM Mini"
640            )
641            && let Some(exp) = controls.iter_mut().find(|c| c.kind == 1)
642        {
643            exp.max = 2_000_000_000;
644        }
645        ensure!(
646            controls.iter().any(|c| c.kind == 1),
647            Failure::Invalid("Exposure control unavailable".into())
648        );
649        if info["cooled"] == true {
650            ensure!(
651                [8, 16, 17]
652                    .iter()
653                    .all(|k| controls.iter().any(|c| c.kind == *k)),
654                Failure::Invalid("Required cooling controls unavailable".into())
655            );
656        }
657        self.selection.serial = serial.clone();
658        {
659            let mut state = self.status.lock().unwrap();
660            state.info = info.clone();
661            state.serial = serial;
662            state.sdk_version = result["sdkVersion"].as_str().unwrap_or("unknown").into();
663            state.backend = if self.direct { "direct" } else { "sdk" }.into();
664            state.sdk_fallback = self.selection.direct && !self.direct;
665            state.controls = controls.into_iter().map(|c| (c.kind, c)).collect();
666            for c in state.controls.values().cloned().collect::<Vec<_>>() {
667                if (c.writable || c.kind == 6) && PERSISTENT_CONTROLS.contains(&c.kind) {
668                    state.values.entry(c.kind).or_insert(c.value);
669                } else {
670                    state.values.insert(c.kind, c.value);
671                }
672            }
673            state.control_connection_available = purpose == OpenPurpose::Acquisition;
674            state.process_id = self.worker.as_ref().and_then(Worker::pid);
675            state.white_balance_capabilities = result["whiteBalance"].clone();
676        }
677        if purpose == OpenPurpose::Acquisition
678            && let Some(settings) = previous.white_balance
679        {
680            // A new worker has no previous estimate for Locked to freeze.
681            if settings.mode == crate::white_balance::Mode::Locked {
682                self.set_white_balance(
683                    crate::white_balance::Settings {
684                        mode: crate::white_balance::Mode::Manual,
685                        ..settings
686                    },
687                    token,
688                )
689                .await?;
690            }
691            self.set_white_balance(settings, token).await?;
692        }
693        self.emit(
694            "info",
695            if purpose == OpenPurpose::ThermalShutdown {
696                "camera.cleanup_opened"
697            } else {
698                "connection.opened"
699            },
700            format!(
701                "Camera opened using {}{}; serial {}{}",
702                if self.direct { "direct" } else { "SDK" },
703                if self.selection.direct && !self.direct {
704                    " fallback"
705                } else {
706                    ""
707                },
708                self.selection.serial.as_deref().unwrap_or("unavailable"),
709                if purpose == OpenPurpose::ThermalShutdown {
710                    "; thermal shutdown only"
711                } else {
712                    ""
713                }
714            ),
715        );
716        if purpose == OpenPurpose::Acquisition {
717            self.cooling.activate();
718        }
719        Ok(())
720    }
721    fn settings(&self) -> BTreeMap<i32, i64> {
722        let state = self.snapshot();
723        state
724            .values
725            .into_iter()
726            .filter(|(k, _)| {
727                (state.white_balance.is_none() || !matches!(k, 3 | 4 | 9))
728                    && state.controls.get(k).is_some_and(|c| c.writable || *k == 6)
729                    && PERSISTENT_CONTROLS.contains(k)
730            })
731            .collect()
732    }
733    async fn apply(
734        &mut self,
735        values: &mut BTreeMap<i32, i64>,
736        token: &CancellationToken,
737    ) -> Result<()> {
738        let mut ordered: Vec<_> = values.iter().map(|(k, v)| (*k, *v)).collect();
739        ordered.sort_by_key(|(k, _)| if *k == 17 { 100 } else { *k });
740        for (kind, value) in ordered {
741            if self.snapshot().white_balance.is_some() && matches!(kind, 3 | 4 | 9) {
742                continue;
743            }
744            if self.applied.get(&kind) == Some(&value) {
745                continue;
746            }
747            let state = self.snapshot();
748            let cap = state
749                .controls
750                .get(&kind)
751                .ok_or_else(|| invalid(format!("Control {kind} disappeared after reconnect")))?;
752            if !cap.writable {
753                ensure!(
754                    cap.value == value,
755                    Failure::Invalid(format!("Read-only control {kind} changed"))
756                );
757                self.applied.insert(kind, value);
758                continue;
759            }
760            let observation = if ACKNOWLEDGED_CONTROLS.contains(&kind) {
761                let deadline = Instant::now()
762                    + Duration::from_secs_f64(self.selection.recovery.command_timeout_seconds);
763                let result = self.write_control(kind, value, deadline, token, true).await;
764                Some(self.control_result(result).await?)
765            } else {
766                self.call("set", json!({"control":kind,"value":value}), None, token)
767                    .await?;
768                None
769            };
770            let actual = if let Some(observation) = observation {
771                observation.value
772            } else {
773                self.call("get", json!({"control":kind}), None, token)
774                    .await?
775                    .0
776                    .as_i64()
777                    .ok_or_else(|| invalid("Invalid control readback"))?
778            };
779            if actual != value {
780                ensure!(
781                    !self.direct && kind == 5 && actual >= cap.min && actual <= cap.max,
782                    "Control {kind} readback {actual} differs from requested {value}"
783                );
784                // Acknowledged controls already emit their clamp diagnostic in
785                // the shared helper; other controls cannot use this policy.
786                values.insert(kind, actual);
787                let mut shared = self.status.lock().unwrap();
788                if shared.values.get(&kind) == Some(&value) {
789                    shared.values.insert(kind, actual);
790                }
791            }
792            self.applied.insert(kind, actual);
793            if let Some(observation) = observation {
794                self.status
795                    .lock()
796                    .unwrap()
797                    .observations
798                    .insert(kind, observation);
799            }
800        }
801        if self.direct
802            && self.settle_required
803            && !self.cooling_seeded
804            && values.get(&17).is_some_and(|v| *v != 0)
805            && let (Some(power), Some(temperature)) =
806                (self.recovery_power, self.recovery_temperature)
807        {
808            self.call(
809                "resume-cooling",
810                json!({"power":power,"temperature":temperature,"previousTarget":self.recovery_target.or_else(|| values.get(&16).copied()).unwrap_or(0)}),
811                None,
812                token,
813            )
814            .await?;
815            self.cooling_seeded = true;
816            self.emit(
817                "info",
818                "cooling.resumed",
819                format!(
820                    "Resumed prior cooler demand {power}% at restored target {} C",
821                    values.get(&16).unwrap_or(&0)
822                ),
823            );
824        }
825        Ok(())
826    }
827    async fn observe_control(
828        &mut self,
829        kind: i32,
830        token: &CancellationToken,
831    ) -> Result<ControlObservation> {
832        let request_started = Instant::now();
833        let (reply, pixels) = self
834            .call("get-observation", json!({"control":kind}), None, token)
835            .await?;
836        ensure!(
837            pixels.is_empty(),
838            invalid("Unexpected observation image payload")
839        );
840        serde_json::from_value::<ControlObservationReply>(reply)
841            .map_err(|_| invalid("Invalid control observation"))?
842            .normalize(request_started)
843    }
844    async fn read_environment(
845        &mut self,
846        token: &CancellationToken,
847    ) -> Result<(Option<f64>, Option<i64>)> {
848        for kind in [8, 15] {
849            if self.snapshot().controls.contains_key(&kind) {
850                let observation = self.observe_control(kind, token).await?;
851                let mut state = self.status.lock().unwrap();
852                state.values.insert(kind, observation.value);
853                state.observations.insert(kind, observation);
854            }
855        }
856        let state = self.snapshot();
857        Ok((
858            state.values.get(&8).map(|v| *v as f64 / 10.),
859            state.values.get(&15).copied(),
860        ))
861    }
862    /// Read-only idle telemetry on the existing worker. Unlike refresh, this
863    /// never opens a replacement worker or applies desired configuration.
864    pub async fn refresh_environment(&mut self, token: &CancellationToken) -> Result<()> {
865        let state = self.snapshot();
866        ensure!(
867            state.connected && state.control_connection_available && self.worker.is_some(),
868            invalid("Camera controls are unavailable")
869        );
870        let result = self.read_environment(token).await.map(|_| ());
871        if let Err(error) = &result {
872            self.emit("warning", "environment.failed", format!("{error:#}"));
873            self.invalidate().await;
874        }
875        result
876    }
877    pub async fn refresh(&mut self, token: &CancellationToken) -> Result<()> {
878        let result = async {
879            if self.worker.is_none() {
880                self.phase("Restoring camera controls");
881                self.delay(self.selection.recovery.reconnect_delay_seconds, token)
882                    .await?;
883                self.open(token).await?;
884            }
885            self.apply(&mut self.settings(), token).await?;
886            self.read_environment(token).await?;
887            Ok(())
888        }
889        .await;
890        if let Err(error) = &result {
891            if !token.is_cancelled() {
892                self.status.lock().unwrap().retry.last_failure = Some(format!("{error:#}"));
893            }
894            self.emit("warning", "controls.failed", format!("{error:#}"));
895            self.invalidate().await;
896        }
897        result
898    }
899    pub fn ready_timeout(&self, seconds: f64) -> f64 {
900        let o = &self.selection.recovery;
901        seconds
902            + o.exposure_grace_seconds
903            + if self.direct {
904                (1 + if self.retained() || seconds <= o.maximum_retry_exposure_seconds {
905                    o.direct_read_retries
906                } else {
907                    0
908                }) as f64
909                    * o.download_timeout_seconds
910                    + o.direct_read_retries as f64
911                        * self.snapshot().info["readRetryOverheadSeconds"]
912                            .as_f64()
913                            .unwrap_or(0.)
914                            .clamp(0., MAX_READ_RETRY_OVERHEAD_SECONDS)
915            } else {
916                0.
917            }
918    }
919    fn retained(&self) -> bool {
920        self.direct && self.snapshot().info["retainedFrameReads"] == true
921    }
922    pub async fn capture(&mut self, e: Exposure, token: &CancellationToken) -> Result<Frame> {
923        if self.snapshot().white_balance.is_some() {
924            let state = self.snapshot();
925            crate::white_balance::WhiteBalance::validate_geometry(
926                state.info["color"] == true,
927                state.info["bayer"].as_u64(),
928                e.bin,
929            )
930            .map_err(|error| invalid(error.to_string()))?;
931        }
932        let state = self.snapshot();
933        validate_capture(
934            &state.info,
935            &state.controls,
936            &e,
937            self.direct && self.selection.sdk_fallback,
938        )?;
939        let options = self.selection.recovery.clone();
940        let seconds = e.microseconds as f64 / 1e6;
941        let retries = options.replacement_exposures(e.microseconds);
942        let mut settings = self.settings();
943        let mut prior = self
944            .recovery_temperature
945            .or_else(|| state.values.get(&8).map(|v| *v as f64 / 10.));
946        let mut power = self
947            .recovery_power
948            .or_else(|| state.values.get(&15).copied());
949        let mut last = None;
950        let mut usb_resets = 0;
951        {
952            let mut status = self.status.lock().unwrap();
953            status.retry.recaptures = 0;
954            status.retry.downloads = 0;
955            status.retry.usb_reads = 0;
956        }
957        for attempt in 0..=retries {
958            if attempt > 0 && !token.is_cancelled() {
959                self.status.lock().unwrap().retry.recaptures = attempt;
960            }
961            let usb_reads_before = self.snapshot().retry.usb_reads;
962            let result = self
963                .attempt(&e, &mut settings, &mut prior, &mut power, attempt, token)
964                .await;
965            match result {
966                Ok(mut frame) => {
967                    {
968                        let reads = frame.metadata["readRecoveries"]
969                            .as_u64()
970                            .unwrap_or(0)
971                            .min(u32::MAX as u64) as u32;
972                        let mut status = self.status.lock().unwrap();
973                        // Close the diagnostic-counting window atomically with
974                        // reconciliation: late stderr records must not count
975                        // the same successful retry a second time.
976                        status.retry.usb_reads = usb_reads_before.saturating_add(reads);
977                        status.phase = "Idle".into();
978                    }
979                    frame.metadata["recoveries"] = json!(attempt);
980                    frame.metadata["usbResets"] = json!(usb_resets);
981                    let retries = self.snapshot().retry;
982                    frame.metadata["retainedReadRetries"] = json!(retries.usb_reads);
983                    frame.metadata["downloadRetriesTotal"] = json!(retries.downloads);
984                    if attempt > 0 || retries.downloads > 0 || retries.usb_reads > 0 {
985                        self.emit("info", "capture.recovered", format!(
986                            "Returning {}x{} image after {attempt} replacement exposures, {} SDK read retries and {} retained-frame retries across all attempts; delivered frame used {} retained-frame retries and {} handle reopens",
987                            e.width, e.height, retries.downloads, retries.usb_reads,
988                            frame.metadata["readRecoveries"].as_u64().unwrap_or(0),
989                            frame.metadata["handleReopens"].as_u64().unwrap_or(0)));
990                    }
991                    self.phase("Idle");
992                    return Ok(frame);
993                }
994                Err(error) => {
995                    if token.is_cancelled() {
996                        // Finish framed replies before reaching this boundary.
997                        // Acknowledged stop preserves the worker and cooler;
998                        // an uncertain stop still retires the isolated worker.
999                        let stopped = self
1000                            .call("stop", Value::Null, None, &CancellationToken::new())
1001                            .await;
1002                        if let Err(stop_error) = stopped {
1003                            let message = format!(
1004                                "Stop was not acknowledged; retiring worker: {stop_error:#}"
1005                            );
1006                            {
1007                                let mut state = self.status.lock().unwrap();
1008                                state.error = Some(message.clone());
1009                                state.retry.last_failure = Some(message.clone());
1010                            }
1011                            self.emit("warning", "capture.abort_failed", message);
1012                            self.invalidate().await;
1013                        } else {
1014                            self.status.lock().unwrap().error = None;
1015                            self.emit("info", "capture.aborted", "Client cancelled capture; camera stopped without reopening or replacing exposure");
1016                        }
1017                        self.phase("Aborted");
1018                        return Err(Failure::Cancelled.into());
1019                    }
1020                    {
1021                        let mut state = self.status.lock().unwrap();
1022                        state.error = Some(format!("{error:#}"));
1023                        if !token.is_cancelled() {
1024                            state.retry.last_failure = Some(format!("{error:#}"));
1025                        }
1026                        state.sdk_error_code = match error.downcast_ref::<Failure>() {
1027                            Some(
1028                                Failure::Worker { code, .. }
1029                                | Failure::UncertainControl { code, .. },
1030                            ) => *code,
1031                            _ => None,
1032                        };
1033                    }
1034                    let can_retry = retryable(&error) && attempt < retries && !token.is_cancelled();
1035                    self.emit(
1036                        "warning",
1037                        "capture.failed",
1038                        format!("Attempt {}/{}: {error:#}", attempt + 1, retries + 1),
1039                    );
1040                    self.invalidate().await;
1041                    if !can_retry {
1042                        last = Some(error);
1043                        break;
1044                    }
1045                    if usb_resets == 0
1046                        && options.usb_reset_after_failures > 0
1047                        && attempt + 1 >= options.usb_reset_after_failures
1048                    {
1049                        usb_resets += 1;
1050                        if let Err(e) = self.reset_usb(token).await {
1051                            self.emit("error", "usb.failed", format!("{e:#}"));
1052                            last = Some(e);
1053                            break;
1054                        }
1055                    }
1056                    if self.direct && self.selection.sdk_fallback {
1057                        self.direct = false;
1058                        self.emit(
1059                            "warning",
1060                            "backend.fallback",
1061                            "Next permitted exposure retry will use SDK fallback",
1062                        );
1063                    }
1064                    self.emit("warning","capture.retry",format!("Scheduling replacement exposure {}/{}: {seconds} s, {}x{}, bin {}; reconnect delay {} s",attempt+1,retries,e.width,e.height,e.bin,options.reconnect_delay_seconds));
1065                    last = Some(error);
1066                }
1067            }
1068        }
1069        self.phase(if token.is_cancelled() {
1070            "Aborted"
1071        } else {
1072            "Error"
1073        });
1074        Err(last.unwrap_or_else(|| invalid("Capture failed")))
1075    }
1076    async fn attempt(
1077        &mut self,
1078        e: &Exposure,
1079        settings: &mut BTreeMap<i32, i64>,
1080        prior: &mut Option<f64>,
1081        power: &mut Option<i64>,
1082        attempt: u32,
1083        token: &CancellationToken,
1084    ) -> Result<Frame> {
1085        let options = self.selection.recovery.clone();
1086        if token.is_cancelled() {
1087            return Err(Failure::Cancelled.into());
1088        }
1089        if attempt > 0 || self.worker.is_none() {
1090            self.phase("Reconnect delay");
1091            self.invalidate().await;
1092            self.delay(options.reconnect_delay_seconds, token).await?;
1093            self.open(token).await?;
1094        }
1095        self.service_cooling_with(settings, token).await?;
1096        self.apply(settings, token).await?;
1097        if self.settle_required
1098            && settings.get(&17).copied().unwrap_or(0) != 0
1099            && let Some(t) = *prior
1100        {
1101            self.settle(t, *power, settings, token).await?;
1102        }
1103        let observed = self.read_environment(token).await?;
1104        *prior = observed.0.or(*prior);
1105        *power = observed.1.or(*power);
1106        if self.direct && self.selection.sdk_fallback {
1107            let result = self
1108                .call("validate", serde_json::to_value(e)?, None, token)
1109                .await;
1110            if let Err(error) = result {
1111                if matches!(
1112                    error.downcast_ref::<Failure>(),
1113                    Some(Failure::Worker { code: Some(8), .. })
1114                ) {
1115                    self.fallback(&error.to_string(), token).await?;
1116                    self.open(token).await?;
1117                    self.apply(settings, token).await?;
1118                    if settings.get(&17).copied().unwrap_or(0) != 0
1119                        && let Some(t) = *prior
1120                    {
1121                        self.settle(t, *power, settings, token).await?;
1122                    }
1123                } else {
1124                    return Err(error);
1125                }
1126            }
1127        }
1128        self.settle_required = false;
1129        self.recovery_temperature = None;
1130        self.recovery_power = None;
1131        self.recovery_target = None;
1132        self.phase("Starting exposure");
1133        let started = Utc::now();
1134        let seconds = e.microseconds as f64 / 1e6;
1135        let mut params = serde_json::to_value(e)?;
1136        if self.direct {
1137            params["readRetries"] = json!(if self.retained()
1138                || seconds <= options.maximum_retry_exposure_seconds
1139            {
1140                options.direct_read_retries
1141            } else {
1142                0
1143            });
1144            params["captureTimeoutSeconds"] =
1145                json!(self.ready_timeout(seconds) + options.command_timeout_seconds);
1146            params["transferTimeoutSeconds"] = json!(options.download_timeout_seconds);
1147            params["readChunkKiB"] = json!(options.direct_read_chunk_kib);
1148        }
1149        // Never cancel a framed pipe exchange halfway through a reply. Check
1150        // cancellation at owner checkpoints, then send a bounded stop command.
1151        let exchange_token = CancellationToken::new();
1152        if token.is_cancelled() {
1153            return Err(Failure::Cancelled.into());
1154        }
1155        self.call("start", params, None, &exchange_token).await?;
1156        self.phase("Exposing");
1157        let clock = Instant::now();
1158        let mut environment_sample = Instant::now();
1159        loop {
1160            if token.is_cancelled() {
1161                return Err(Failure::Cancelled.into());
1162            }
1163            self.service_cooling_with(settings, token).await?;
1164            let state = self
1165                .call("status", Value::Null, None, &exchange_token)
1166                .await?
1167                .0
1168                .as_i64()
1169                .ok_or_else(|| invalid("Invalid exposure state"))?;
1170            self.status.lock().unwrap().sdk_exposure_state = Some(state);
1171            if state == 2 {
1172                break;
1173            }
1174            ensure!(state == 1, "Exposure ended in camera state {state}");
1175            ensure!(
1176                clock.elapsed().as_secs_f64() <= self.ready_timeout(seconds),
1177                "Exposure readiness timed out"
1178            );
1179            if environment_sample.elapsed() >= Duration::from_secs(2) {
1180                self.read_environment(&exchange_token).await?;
1181                environment_sample = Instant::now();
1182            }
1183            self.delay(0.025, token).await?;
1184        }
1185        self.service_cooling_with(settings, token).await?;
1186        self.phase("Downloading");
1187        let mut reads = 0;
1188        let (mut metadata, pixels) = loop {
1189            match self
1190                .worker
1191                .as_mut()
1192                .context("Camera worker is disconnected")?
1193                .call_image(
1194                    "download",
1195                    Value::Null,
1196                    options.download_timeout_seconds,
1197                    &exchange_token,
1198                    e.bytes()?,
1199                )
1200                .await
1201            {
1202                Ok(frame) => break frame,
1203                Err(error) if !self.retained() && retryable(&error) && !token.is_cancelled() => {
1204                    self.status.lock().unwrap().retry.last_failure = Some(format!("{error:#}"));
1205                    self.emit("warning", "transfer.failed", format!("{error:#}"));
1206                    let state = self
1207                        .call("status", Value::Null, None, token)
1208                        .await?
1209                        .0
1210                        .as_i64()
1211                        .ok_or_else(|| invalid("Invalid post-transfer state"))?;
1212                    self.status.lock().unwrap().sdk_exposure_state = Some(state);
1213                    if reads >= options.ready_frame_download_retries || state != 2 {
1214                        return Err(error);
1215                    }
1216                    reads += 1;
1217                    self.status.lock().unwrap().retry.downloads += 1;
1218                    self.phase(format!(
1219                        "Rereading ready frame ({reads}/{})",
1220                        options.ready_frame_download_retries
1221                    ));
1222                    self.delay(options.reconnect_delay_seconds, token).await?;
1223                }
1224                Err(error) => return Err(error),
1225            }
1226        };
1227        if token.is_cancelled() {
1228            return Err(Failure::Cancelled.into());
1229        }
1230        {
1231            let mut state = self.status.lock().unwrap();
1232            let recovered = metadata["readRecoveries"]
1233                .as_u64()
1234                .unwrap_or(0)
1235                .min(u32::MAX as u64) as u32;
1236            if let Some(failure) = metadata["readErrors"]
1237                .as_array()
1238                .or_else(|| metadata["readoutErrors"].as_array())
1239                .and_then(|errors| errors.last())
1240                .and_then(Value::as_str)
1241            {
1242                state.retry.last_failure = Some(failure.into());
1243            } else if recovered > 0 && state.retry.usb_reads == 0 {
1244                state.retry.last_failure = Some(
1245                    "USB frame read failed; frame recovered (worker supplied no failure detail)"
1246                        .into(),
1247                );
1248            }
1249            if let Some(failure) = metadata["cleanupError"].as_str() {
1250                state.retry.last_failure = Some(failure.into());
1251            }
1252        }
1253        ensure!(
1254            pixels.len() == e.bytes()?,
1255            Failure::Invalid("Image length differs from requested ROI".into())
1256        );
1257        for (key, expected) in [("width", e.width), ("height", e.height)] {
1258            ensure!(
1259                metadata[key].as_u64() == Some(expected as u64),
1260                Failure::Invalid(format!("Frame {key} differs from request"))
1261            );
1262        }
1263        if let Some(error) = metadata["cleanupError"].as_str() {
1264            self.emit(
1265                "warning",
1266                "capture.cleanup_failed",
1267                format!("Frame preserved; reconnect required: {error}"),
1268            );
1269            self.invalidate().await;
1270        } else {
1271            self.service_cooling_with(settings, token).await?;
1272        }
1273        metadata["startedUtc"] = json!(started);
1274        if self.snapshot().white_balance.is_some() {
1275            let settings = serde_json::from_value(metadata["whiteBalance"]["settings"].clone())
1276                .map_err(|_| invalid("Missing or invalid managed white balance frame metadata"))?;
1277            self.status.lock().unwrap().white_balance = Some(settings);
1278        }
1279        metadata["endedUtc"] = json!(Utc::now());
1280        metadata["exposure"] = serde_json::to_value(e)?;
1281        metadata["controls"] = serde_json::to_value(settings)?;
1282        metadata["backend"] = json!(if self.direct { "direct" } else { "sdk" });
1283        metadata["sdkFallback"] = json!(self.selection.direct && !self.direct);
1284        metadata["downloadRetries"] = json!(reads);
1285        Ok(Frame {
1286            exposure: e.clone(),
1287            metadata,
1288            pixels: pixels.into(),
1289        })
1290    }
1291    async fn settle(
1292        &mut self,
1293        prior: f64,
1294        prior_power: Option<i64>,
1295        settings: &mut BTreeMap<i32, i64>,
1296        token: &CancellationToken,
1297    ) -> Result<()> {
1298        let o = self.selection.recovery.clone();
1299        let clock = Instant::now();
1300        let mut stable = 0;
1301        let mut hold = CoolingHold::default();
1302        self.phase(format!("Restoring cooling near {prior:.1} C"));
1303        while clock.elapsed().as_secs_f64() < o.cooling_timeout_seconds {
1304            self.service_cooling_with(settings, token).await?;
1305            if settings.get(&17).copied().unwrap_or(0) == 0 {
1306                return Ok(());
1307            }
1308            let target = *settings.get(&16).unwrap_or(&0) as f64;
1309            let (temperature, power) = self.read_environment(token).await?;
1310            let output = prior_power.is_none_or(|p| p <= 10 || power.is_some_and(|v| v >= p - 10));
1311            let near = temperature.is_some_and(|t| {
1312                t >= prior.min(target) - o.temperature_tolerance_c
1313                    && t <= prior + o.temperature_tolerance_c
1314            });
1315            let target_held = hold.observe(
1316                temperature,
1317                power,
1318                target,
1319                o.temperature_tolerance_c,
1320                o.cooling_sample_seconds,
1321                clock.elapsed().as_secs_f64(),
1322            );
1323            stable = if (near && output) || target_held {
1324                stable + 1
1325            } else {
1326                0
1327            };
1328            self.emit("info","cooling.wait",format!("Temperature {temperature:?} C (prior {prior}), power {power:?}% (prior {prior_power:?}%), stable {stable}/{}",o.cooling_stable_samples));
1329            if stable >= o.cooling_stable_samples {
1330                self.emit(
1331                    "info",
1332                    "cooling.recovered",
1333                    format!("Cooling recovered at {temperature:?} C; restored setpoint {target} C"),
1334                );
1335                return Ok(());
1336            }
1337            self.delay(o.cooling_sample_seconds, token).await?;
1338        }
1339        anyhow::bail!("Camera did not recover its prior cooling temperature and output")
1340    }
1341    pub async fn close(&mut self) {
1342        self.cooling.retire();
1343        self.status.lock().unwrap().connected = false;
1344        let token = CancellationToken::new();
1345        let closed = if self.worker.is_some() {
1346            match self
1347                .call("close", Value::Null, Some(CLOSE_SECONDS), &token)
1348                .await
1349            {
1350                Ok(_) => true,
1351                Err(error) => {
1352                    let message = format!("Camera close or settings restoration failed: {error:#}");
1353                    self.emit("warning", "camera.cleanup_failed", &message);
1354                    self.status.lock().unwrap().error = Some(message);
1355                    false
1356                }
1357            }
1358        } else {
1359            false
1360        };
1361        self.invalidate().await;
1362        if !closed
1363            && self.ever_opened
1364            && self
1365                .snapshot()
1366                .controls
1367                .values()
1368                .any(|c| c.writable && matches!(c.kind, 17 | 21))
1369            && self.selection.serial.is_some()
1370        {
1371            let cleanup = async {
1372                self.delay(self.selection.recovery.reconnect_delay_seconds, &token)
1373                    .await?;
1374                self.open_worker(&token, OpenPurpose::ThermalShutdown)
1375                    .await?;
1376                // Backend close disables each supported thermal actuator and
1377                // attempts both even if one write fails. Never restore settings.
1378                self.call("close", Value::Null, Some(CLOSE_SECONDS), &token)
1379                    .await?;
1380                Result::<()>::Ok(())
1381            }
1382            .await;
1383            if let Err(error) = cleanup {
1384                self.emit(
1385                    "warning",
1386                    "cooling.cleanup_failed",
1387                    format!("Could not disable cooler/dew heater: {error:#}"),
1388                );
1389            }
1390            self.invalidate().await;
1391        }
1392        self.status.lock().unwrap().connected = false;
1393        self.cooling.retire();
1394        self.usb_target = None;
1395        self.phase("Disconnected");
1396    }
1397}
1398
1399impl Drop for Session {
1400    fn drop(&mut self) {
1401        self.cooling.retire();
1402    }
1403}
1404
1405#[derive(Default)]
1406pub struct CoolingHold {
1407    since: Option<f64>,
1408    first: f64,
1409    last: Option<f64>,
1410}
1411impl CoolingHold {
1412    pub fn observe(
1413        &mut self,
1414        temperature: Option<f64>,
1415        power: Option<i64>,
1416        target: f64,
1417        tolerance: f64,
1418        sample: f64,
1419        elapsed: f64,
1420    ) -> bool {
1421        let valid = temperature
1422            .is_some_and(|t| t.is_finite() && (t - target).abs() <= tolerance.min(1.))
1423            && power.is_some_and(|p| p > 0);
1424        if !valid {
1425            self.since = None;
1426            self.last = None;
1427            return false;
1428        }
1429        let t = temperature.unwrap();
1430        if self.since.is_none()
1431            || self
1432                .last
1433                .is_some_and(|last| elapsed < last || elapsed - last > 5f64.max(sample * 2.))
1434            || t > self.first + 0.2
1435        {
1436            self.since = Some(elapsed);
1437            self.first = t;
1438        }
1439        self.last = Some(elapsed);
1440        elapsed - self.since.unwrap() >= 30.
1441    }
1442}
1443
1444#[cfg(test)]
1445mod tests {
1446    use super::*;
1447    fn runtime() -> Runtime {
1448        let directory = std::env::var_os("REGAIN_TEST_WORKERS")
1449            .map(std::path::PathBuf::from)
1450            .unwrap_or_else(|| {
1451                std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../../target/debug")
1452            });
1453        Runtime {
1454            directory,
1455            sdk: "unused".into(),
1456            simulate: true,
1457            sdk_simulation: Some(json!({"instant":true})),
1458        }
1459    }
1460    fn selection(direct: bool) -> Selection {
1461        Selection {
1462            name: if direct {
1463                "ZWO ASI676MC"
1464            } else {
1465                "ZWO Simulated"
1466            }
1467            .into(),
1468            serial: None,
1469            direct,
1470            sdk_fallback: false,
1471            recovery: RecoveryOptions {
1472                reconnect_delay_seconds: 0.01,
1473                cooling_sample_seconds: 0.01,
1474                cooling_stable_samples: 1,
1475                ..RecoveryOptions::default()
1476            },
1477        }
1478    }
1479    fn exposure() -> Exposure {
1480        Exposure {
1481            width: 64,
1482            height: 64,
1483            bin: 1,
1484            x: 0,
1485            y: 0,
1486            microseconds: 10000,
1487            dark: true,
1488        }
1489    }
1490    fn log() -> Diagnostic {
1491        Arc::new(|_, _, _| {})
1492    }
1493    async fn exposing(status: &SharedStatus) {
1494        let deadline = Instant::now() + Duration::from_secs(10);
1495        loop {
1496            if status.lock().unwrap().phase == "Exposing" {
1497                return;
1498            }
1499            assert!(Instant::now() < deadline, "Capture never reached Exposing");
1500            tokio::time::sleep(Duration::from_millis(1)).await;
1501        }
1502    }
1503    #[tokio::test]
1504    async fn worker_evidence_age_survives_core_refresh_and_invalid_readback_is_uncertain() {
1505        let token = CancellationToken::new();
1506        let mut session = Session::new(selection(false), runtime(), log()).unwrap();
1507        session.connect(&token).await.unwrap();
1508        session.refresh(&token).await.unwrap();
1509        session
1510            .call(
1511                "simulation",
1512                json!({"observationControl":8,
1513            "observationReply":{"value":-100,"ageSeconds":120.25}}),
1514                None,
1515                &token,
1516            )
1517            .await
1518            .unwrap();
1519        let requested = Instant::now();
1520        session.read_environment(&token).await.unwrap();
1521        let received = Instant::now();
1522        let observed = session.snapshot().observations[&8];
1523        assert_eq!(observed.value, -100);
1524        // The internal request admission lies between these two boundaries.
1525        let age = Duration::from_secs_f64(120.25);
1526        assert!(observed.observed_at >= requested - age);
1527        assert!(observed.observed_at <= received - age);
1528        let previous = session.snapshot().observations[&0];
1529        session
1530            .call(
1531                "simulation",
1532                json!({"observationControl":0,
1533            "observationReply":{"value":123,"ageSeconds":-1}}),
1534                None,
1535                &token,
1536            )
1537            .await
1538            .unwrap();
1539        let error = session
1540            .set_imaging_control(0, 123, Instant::now() + Duration::from_secs(5), &token)
1541            .await
1542            .unwrap_err();
1543        assert!(matches!(
1544            error.downcast_ref::<Failure>(),
1545            Some(Failure::UncertainControl { .. })
1546        ));
1547        assert!(!retryable(&error));
1548        assert_eq!(session.snapshot().observations[&0], previous);
1549        assert!(!session.snapshot().control_connection_available);
1550        assert!(session.worker.is_none());
1551        session.close().await;
1552    }
1553    #[tokio::test]
1554    async fn queued_persistent_control_invalid_observation_never_replays_or_starts_exposure() {
1555        for (kind, value) in [(0, 123), (5, 20), (16, -15), (17, 0)] {
1556            let token = CancellationToken::new();
1557            let events = Arc::new(Mutex::new(Vec::new()));
1558            let sink = events.clone();
1559            let mut session = Session::new(
1560                selection(false),
1561                runtime(),
1562                Arc::new(move |_, event, _| sink.lock().unwrap().push(event.to_owned())),
1563            )
1564            .unwrap();
1565            session.connect(&token).await.unwrap();
1566            session.refresh(&token).await.unwrap();
1567            let previous = session.snapshot().observations[&kind];
1568            assert_ne!(
1569                previous.value, value,
1570                "The test must dispatch a changed setting"
1571            );
1572            let writes = events
1573                .lock()
1574                .unwrap()
1575                .iter()
1576                .filter(|event| event.as_str() == "control.write_acknowledged")
1577                .count();
1578            session
1579                .call(
1580                    "simulation",
1581                    json!({"observationControl":kind,
1582                "observationReply":{"value":value,"ageSeconds":-1}}),
1583                    None,
1584                    &token,
1585                )
1586                .await
1587                .unwrap();
1588            Session::queue_control(&session.status, kind, value).unwrap();
1589            let error = session.capture(exposure(), &token).await.err().unwrap();
1590            assert!(matches!(
1591                error.downcast_ref::<Failure>(),
1592                Some(Failure::UncertainControl { .. })
1593            ));
1594            assert!(!retryable(&error));
1595            assert_eq!(session.snapshot().observations[&kind], previous);
1596            assert!(!session.snapshot().control_connection_available);
1597            assert!(session.worker.is_none());
1598            let observed = events.lock().unwrap().clone();
1599            assert_eq!(
1600                observed
1601                    .iter()
1602                    .filter(|event| event.as_str() == "control.write_acknowledged")
1603                    .count(),
1604                writes + 1
1605            );
1606            assert_eq!(
1607                observed
1608                    .iter()
1609                    .filter(|event| event.as_str() == "connection.opened")
1610                    .count(),
1611                1
1612            );
1613            assert!(!observed.iter().any(|event| event == "capture.retry"));
1614            assert_ne!(session.snapshot().phase, "Exposing");
1615            session.close().await;
1616        }
1617    }
1618    #[tokio::test]
1619    async fn idle_environment_read_never_applies_queued_settings_or_reopens_lost_worker() {
1620        for direct in [false, true] {
1621            let token = CancellationToken::new();
1622            let mut selected = selection(direct);
1623            if direct {
1624                selected.name = "ZWO ASI585MM Pro".into();
1625            }
1626            let mut session = Session::new(selected, runtime(), log()).unwrap();
1627            session.connect(&token).await.unwrap();
1628            session.refresh(&token).await.unwrap();
1629            let before = session.snapshot();
1630            Session::queue_control(&session.status, 0, 234).unwrap();
1631            Session::queue_control(&session.status, 16, -20).unwrap();
1632            session.refresh_environment(&token).await.unwrap();
1633            let after = session.snapshot();
1634            assert_eq!(after.process_id, before.process_id);
1635            for kind in [0, 5, 16, 17] {
1636                assert_eq!(after.observations[&kind], before.observations[&kind]);
1637                assert_eq!(
1638                    session
1639                        .call("get", json!({"control":kind}), None, &token)
1640                        .await
1641                        .unwrap()
1642                        .0,
1643                    before.observations[&kind].value
1644                );
1645            }
1646            assert_eq!(after.values[&0], 234);
1647            assert_eq!(after.values[&16], -20);
1648            for kind in [8, 15] {
1649                assert!(
1650                    after.observations[&kind].observed_at > before.observations[&kind].observed_at
1651                );
1652            }
1653            session.invalidate().await;
1654            assert!(session.refresh_environment(&token).await.is_err());
1655            assert!(session.worker.is_none());
1656            assert!(session.snapshot().process_id.is_none());
1657            assert_eq!(session.snapshot().observations, after.observations);
1658            session.close().await;
1659        }
1660    }
1661    #[tokio::test]
1662    async fn failed_idle_environment_read_retires_worker_without_reopening_or_applying_settings() {
1663        let token = CancellationToken::new();
1664        let mut session = Session::new(selection(false), runtime(), log()).unwrap();
1665        session.connect(&token).await.unwrap();
1666        session.refresh(&token).await.unwrap();
1667        let before = session.snapshot();
1668        Session::queue_control(&session.status, 0, 234).unwrap();
1669        session
1670            .call(
1671                "simulation",
1672                json!({"observationControl":8,
1673            "observationReply":{"value":100,"ageSeconds":-1}}),
1674                None,
1675                &token,
1676            )
1677            .await
1678            .unwrap();
1679        assert!(session.refresh_environment(&token).await.is_err());
1680        assert!(!session.snapshot().control_connection_available);
1681        assert!(session.worker.is_none());
1682        assert_eq!(session.snapshot().observations, before.observations);
1683        assert_eq!(session.snapshot().values[&0], 234);
1684        assert!(session.refresh_environment(&token).await.is_err());
1685        assert!(session.worker.is_none());
1686        assert_eq!(session.snapshot().observations, before.observations);
1687        session.close().await;
1688    }
1689    #[tokio::test]
1690    async fn observations_are_acknowledged_evidence_not_desired_settings_or_new_worker_defaults() {
1691        for direct in [false, true] {
1692            let token = CancellationToken::new();
1693            let mut selected = selection(direct);
1694            if direct {
1695                selected.name = "ZWO ASI585MM Pro".into();
1696            }
1697            let mut session = Session::new(selected, runtime(), log()).unwrap();
1698            session.connect(&token).await.unwrap();
1699            assert!(session.snapshot().observations.is_empty());
1700            Session::queue_control(&session.status, 0, 123).unwrap();
1701            assert!(session.snapshot().observations.is_empty());
1702            session.refresh(&token).await.unwrap();
1703            let observed = session.snapshot().observations;
1704            for kind in [0, 5, 8, 15, 16, 17] {
1705                assert!(observed[&kind].observed_at <= Instant::now());
1706            }
1707            assert_eq!(observed[&0].value, 123);
1708            Session::queue_control(&session.status, 0, 234).unwrap();
1709            let queued = session.snapshot();
1710            assert_eq!(queued.values[&0], 234);
1711            assert_eq!(queued.observations[&0], observed[&0]);
1712            let json = serde_json::to_value(&queued).unwrap();
1713            assert!(json.get("observations").is_none());
1714            assert!(json.to_string().find("observedAt").is_none());
1715            session.refresh(&token).await.unwrap();
1716            let acknowledged = session.snapshot().observations[&0];
1717            assert_eq!(acknowledged.value, 234);
1718            assert!(acknowledged.observed_at >= observed[&0].observed_at);
1719            session.invalidate().await;
1720            // Diagnostics can retain known aged evidence while unavailable.
1721            assert_eq!(session.snapshot().observations[&0], acknowledged);
1722            session.open(&token).await.unwrap();
1723            assert!(session.snapshot().observations.is_empty());
1724            assert_eq!(session.snapshot().values[&0], 234);
1725            session.refresh(&token).await.unwrap();
1726            assert_eq!(session.snapshot().observations[&0].value, 234);
1727            session.close().await;
1728        }
1729    }
1730    #[tokio::test]
1731    async fn acknowledged_imaging_controls_apply_before_success_and_survive_worker_recovery() {
1732        for direct in [false, true] {
1733            let token = CancellationToken::new();
1734            let mut session = Session::new(selection(direct), runtime(), log()).unwrap();
1735            session.connect(&token).await.unwrap();
1736            session.refresh(&token).await.unwrap();
1737            for (kind, value) in [(0, 123), (5, 20)] {
1738                assert_eq!(
1739                    session
1740                        .set_imaging_control(
1741                            kind,
1742                            value,
1743                            Instant::now() + Duration::from_secs(5),
1744                            &token
1745                        )
1746                        .await
1747                        .unwrap(),
1748                    value
1749                );
1750                assert_eq!(session.snapshot().values[&kind], value);
1751                let observation = session.snapshot().observations[&kind];
1752                assert_eq!(observation.value, value);
1753                assert!(observation.observed_at <= Instant::now());
1754                assert_eq!(session.applied[&kind], value);
1755                assert_eq!(
1756                    session
1757                        .call("get", json!({"control":kind}), None, &token)
1758                        .await
1759                        .unwrap()
1760                        .0,
1761                    value
1762                );
1763            }
1764            session.invalidate().await;
1765            let frame = session.capture(exposure(), &token).await.unwrap();
1766            assert_eq!(frame.metadata["controls"]["0"], 123);
1767            assert_eq!(frame.metadata["controls"]["5"], 20);
1768            for (kind, value) in [(0, 123), (5, 20)] {
1769                assert_eq!(
1770                    session
1771                        .call("get", json!({"control":kind}), None, &token)
1772                        .await
1773                        .unwrap()
1774                        .0,
1775                    value
1776                );
1777            }
1778            session.close().await;
1779        }
1780    }
1781    #[tokio::test]
1782    async fn imaging_control_mismatch_retires_worker_without_publishing_or_retrying() {
1783        let token = CancellationToken::new();
1784        let mut rt = runtime();
1785        rt.sdk_simulation = Some(json!({"instant":true,"clampControl":0,"clampMinimum":200}));
1786        let mut session = Session::new(selection(false), rt, log()).unwrap();
1787        session.connect(&token).await.unwrap();
1788        let before = session.snapshot().values[&0];
1789        let error = session
1790            .set_imaging_control(0, 100, Instant::now() + Duration::from_secs(5), &token)
1791            .await
1792            .unwrap_err();
1793        assert!(matches!(
1794            error.downcast_ref::<Failure>(),
1795            Some(Failure::UncertainControl { .. })
1796        ));
1797        assert!(!retryable(&error));
1798        assert_eq!(session.snapshot().values[&0], before);
1799        assert!(!session.snapshot().control_connection_available);
1800        assert!(session.worker.is_none());
1801        assert!(session.applied.is_empty());
1802        session.close().await;
1803    }
1804    #[tokio::test]
1805    async fn acknowledged_sdk_offset_preserves_existing_bounded_clamp_policy() {
1806        let token = CancellationToken::new();
1807        let mut rt = runtime();
1808        rt.sdk_simulation = Some(json!({"instant":true,"clampControl":5,"clampMinimum":20}));
1809        let mut session = Session::new(selection(false), rt, log()).unwrap();
1810        session.connect(&token).await.unwrap();
1811        assert_eq!(
1812            session
1813                .set_imaging_control(5, 0, Instant::now() + Duration::from_secs(5), &token)
1814                .await
1815                .unwrap(),
1816            20
1817        );
1818        assert_eq!(session.snapshot().values[&5], 20);
1819        assert_eq!(session.applied[&5], 20);
1820        let frame = session.capture(exposure(), &token).await.unwrap();
1821        assert_eq!(frame.metadata["controls"]["5"], 20);
1822        let maximum = session.snapshot().controls[&5].max;
1823        session
1824            .call(
1825                "simulation",
1826                json!({"clampControl":5,"clampMinimum":maximum+1}),
1827                None,
1828                &token,
1829            )
1830            .await
1831            .unwrap();
1832        let error = session
1833            .set_imaging_control(5, 0, Instant::now() + Duration::from_secs(5), &token)
1834            .await
1835            .unwrap_err();
1836        assert!(matches!(
1837            error.downcast_ref::<Failure>(),
1838            Some(Failure::UncertainControl { .. })
1839        ));
1840        assert_eq!(session.snapshot().values[&5], 20);
1841        assert!(session.worker.is_none());
1842        session.close().await;
1843    }
1844    #[tokio::test]
1845    async fn imaging_control_preflight_rejects_without_io_or_desired_state_changes() {
1846        let token = CancellationToken::new();
1847        let mut session = Session::new(selection(false), runtime(), log()).unwrap();
1848        session.connect(&token).await.unwrap();
1849        session.refresh(&token).await.unwrap();
1850        let before = session.snapshot();
1851        let applied = session.applied.clone();
1852        for (kind, value) in [(1, 1000), (16, -10), (0, before.controls[&0].max + 1)] {
1853            let error = session
1854                .set_imaging_control(kind, value, Instant::now() + Duration::from_secs(5), &token)
1855                .await
1856                .unwrap_err();
1857            assert!(matches!(
1858                error.downcast_ref::<Failure>(),
1859                Some(Failure::Invalid(_))
1860            ));
1861        }
1862        for malformed in 0..4 {
1863            {
1864                let mut state = session.status.lock().unwrap();
1865                let cap = state.controls.get_mut(&0).unwrap();
1866                match malformed {
1867                    0 => cap.kind = 5,
1868                    1 => cap.min = cap.max + 1,
1869                    2 => cap.writable = false,
1870                    _ => {
1871                        state.controls.remove(&0);
1872                    }
1873                }
1874            }
1875            assert!(
1876                session
1877                    .set_imaging_control(0, 123, Instant::now() + Duration::from_secs(5), &token)
1878                    .await
1879                    .is_err()
1880            );
1881            session
1882                .status
1883                .lock()
1884                .unwrap()
1885                .controls
1886                .insert(0, before.controls[&0].clone());
1887        }
1888        let cancelled = CancellationToken::new();
1889        cancelled.cancel();
1890        let error = session
1891            .set_imaging_control(0, 123, Instant::now() + Duration::from_secs(5), &cancelled)
1892            .await
1893            .unwrap_err();
1894        assert!(matches!(
1895            error.downcast_ref::<Failure>(),
1896            Some(Failure::Cancelled)
1897        ));
1898        let error = session
1899            .set_imaging_control(0, 123, Instant::now(), &token)
1900            .await
1901            .unwrap_err();
1902        assert!(matches!(
1903            error.downcast_ref::<cooling::CoolingError>(),
1904            Some(cooling::CoolingError::Expired)
1905        ));
1906        let cooler = session
1907            .cooling()
1908            .submit(16, -15, Duration::from_secs(5))
1909            .unwrap();
1910        assert!(
1911            session
1912                .set_imaging_control(0, 123, Instant::now() + Duration::from_secs(5), &token)
1913                .await
1914                .is_err()
1915        );
1916        drop(cooler);
1917        assert_eq!(session.snapshot().values, before.values);
1918        assert_eq!(session.applied, applied);
1919        assert_eq!(session.snapshot().process_id, before.process_id);
1920        for kind in [0, 5] {
1921            assert_eq!(
1922                session
1923                    .call("get", json!({"control":kind}), None, &token)
1924                    .await
1925                    .unwrap()
1926                    .0,
1927                before.values[&kind]
1928            );
1929        }
1930        session.close().await;
1931    }
1932    async fn acknowledged_cooling(direct: bool) {
1933        let token = CancellationToken::new();
1934        let mut sel = selection(direct);
1935        if direct {
1936            sel.name = "ZWO ASI585MM Pro".into();
1937        }
1938        sel.recovery.ready_frame_download_retries = 0;
1939        let mut rt = runtime();
1940        rt.sdk_simulation = Some(json!({"instant":false}));
1941        let mut session = Session::new(sel, rt, log()).unwrap();
1942        session.connect(&token).await.unwrap();
1943        let handle = session.cooling();
1944        let shared = session.status.clone();
1945        let original = shared.lock().unwrap().values[&16];
1946        Session::queue_control(&shared, 0, 123).unwrap();
1947        if !direct {
1948            session
1949                .call("fault", json!({"kind":"download"}), None, &token)
1950                .await
1951                .unwrap();
1952        }
1953        let capture = session.capture(
1954            Exposure {
1955                microseconds: 6_000_000,
1956                ..exposure()
1957            },
1958            &token,
1959        );
1960        let commands = async {
1961            exposing(&shared).await;
1962            let first = handle.submit(16, -10, Duration::from_secs(5)).unwrap();
1963            assert_eq!(shared.lock().unwrap().values[&16], original);
1964            assert_eq!(first.wait().await, Ok(-10));
1965            for (kind, value) in [(17, 1), (16, -15)] {
1966                assert_eq!(
1967                    handle
1968                        .submit(kind, value, Duration::from_secs(5))
1969                        .unwrap()
1970                        .wait()
1971                        .await,
1972                    Ok(value)
1973                );
1974            }
1975            drop(handle.submit(16, -20, Duration::from_secs(5)).unwrap());
1976            // Legacy deferred imaging intent must not change frozen capture or
1977            // replacement settings. Only acknowledged cooler keys are live.
1978            Session::queue_control(&shared, 0, 200).unwrap();
1979        };
1980        let (frame, ()) = tokio::join!(capture, commands);
1981        let frame = frame.unwrap();
1982        assert_eq!(frame.metadata["controls"]["16"], -15);
1983        assert_eq!(frame.metadata["controls"]["17"], 1);
1984        assert_eq!(frame.metadata["controls"]["0"], 123);
1985        assert_eq!(frame.metadata["recoveries"], if direct { 0 } else { 1 });
1986        assert_eq!(
1987            frame.pixels.as_ref(),
1988            (0..4096u16).flat_map(u16::to_le_bytes).collect::<Vec<_>>()
1989        );
1990        assert_eq!(session.snapshot().values[&16], -15);
1991        assert_eq!(session.snapshot().values[&17], 1);
1992        assert!(!handle.pending());
1993        session.invalidate().await;
1994        // A later capture restores the acknowledged target after worker loss.
1995        let next = session.capture(exposure(), &token).await.unwrap();
1996        assert_eq!(next.metadata["controls"]["16"], -15);
1997        assert_eq!(next.metadata["controls"]["17"], 1);
1998        assert_eq!(next.metadata["controls"]["0"], 200);
1999        session.close().await;
2000    }
2001    #[tokio::test]
2002    async fn acknowledged_sdk_cooling_survives_replacement_and_preserves_imaging_settings() {
2003        acknowledged_cooling(false).await;
2004    }
2005    #[tokio::test]
2006    async fn acknowledged_direct_cooling_survives_worker_recovery() {
2007        acknowledged_cooling(true).await;
2008    }
2009    #[tokio::test]
2010    async fn cooler_readback_mismatch_stops_capture_without_retry_or_target_publication() {
2011        let token = CancellationToken::new();
2012        let mut rt = runtime();
2013        rt.sdk_simulation = Some(json!({"instant":false}));
2014        let events = Arc::new(Mutex::new(Vec::new()));
2015        let sink = events.clone();
2016        let mut session = Session::new(
2017            selection(false),
2018            rt,
2019            Arc::new(move |_, event, _| sink.lock().unwrap().push(event.to_owned())),
2020        )
2021        .unwrap();
2022        session.connect(&token).await.unwrap();
2023        session.refresh(&token).await.unwrap();
2024        session
2025            .call(
2026                "simulation",
2027                json!({"clampControl":16,"clampMinimum":5}),
2028                None,
2029                &token,
2030            )
2031            .await
2032            .unwrap();
2033        let handle = session.cooling();
2034        let shared = session.status.clone();
2035        let original = shared.lock().unwrap().values[&16];
2036        let capture = session.capture(
2037            Exposure {
2038                microseconds: 6_000_000,
2039                ..exposure()
2040            },
2041            &token,
2042        );
2043        let commands = async {
2044            exposing(&shared).await;
2045            assert!(matches!(
2046                handle
2047                    .submit(16, -10, Duration::from_secs(5))
2048                    .unwrap()
2049                    .wait()
2050                    .await,
2051                Err(cooling::CoolingError::Uncertain { .. })
2052            ));
2053        };
2054        let (result, ()) = tokio::join!(capture, commands);
2055        let error = result.err().unwrap();
2056        assert!(matches!(
2057            error.downcast_ref::<Failure>(),
2058            Some(Failure::UncertainControl { .. })
2059        ));
2060        assert!(!retryable(&error));
2061        assert_eq!(session.snapshot().values[&16], original);
2062        assert_eq!(session.snapshot().observations[&16].value, original);
2063        assert!(!session.snapshot().control_connection_available);
2064        assert!(
2065            !events
2066                .lock()
2067                .unwrap()
2068                .iter()
2069                .any(|event| event == "capture.retry")
2070        );
2071        assert!(matches!(
2072            handle.submit(17, 1, Duration::from_secs(1)),
2073            Err(cooling::CoolingError::Unavailable)
2074        ));
2075        session.close().await;
2076    }
2077    #[tokio::test]
2078    async fn idle_cooling_acknowledges_and_teardown_rejects_queued_requests() {
2079        for direct in [false, true] {
2080            let token = CancellationToken::new();
2081            let mut sel = selection(direct);
2082            if direct {
2083                sel.name = "ZWO ASI585MM Pro".into();
2084            }
2085            let mut session = Session::new(sel, runtime(), log()).unwrap();
2086            session.connect(&token).await.unwrap();
2087            let handle = session.cooling();
2088            let receipt = handle.submit(16, -10, Duration::from_secs(5)).unwrap();
2089            session.service_cooling(&token).await.unwrap();
2090            assert_eq!(receipt.wait().await, Ok(-10));
2091            assert_eq!(session.snapshot().values[&16], -10);
2092            let receipt = handle.submit(16, -20, Duration::from_secs(5)).unwrap();
2093            session.close().await;
2094            assert_eq!(
2095                receipt.wait().await,
2096                Err(cooling::CoolingError::Unavailable)
2097            );
2098            assert!(!handle.pending());
2099        }
2100    }
2101    #[tokio::test]
2102    async fn live_cooler_target_and_disable_are_acknowledged_during_recovery_settle() {
2103        for direct in [false, true] {
2104            let token = CancellationToken::new();
2105            let mut sel = selection(direct);
2106            if direct {
2107                sel.name = "ZWO ASI585MM Pro".into();
2108            }
2109            let mut session = Session::new(sel, runtime(), log()).unwrap();
2110            session.connect(&token).await.unwrap();
2111            let handle = session.cooling();
2112            let shared = session.status.clone();
2113            let enabled = handle.submit(17, 1, Duration::from_secs(5)).unwrap();
2114            session.service_cooling(&token).await.unwrap();
2115            assert_eq!(enabled.wait().await, Ok(1));
2116            // Neither simulator meets this prior temperature/output, so settle
2117            // must remain active until the caller explicitly disables cooling.
2118            session.seed_recovery(Some(-30.), Some(80));
2119            let capture = session.capture(exposure(), &token);
2120            let commands = async {
2121                let deadline = Instant::now() + Duration::from_secs(10);
2122                loop {
2123                    if shared
2124                        .lock()
2125                        .unwrap()
2126                        .phase
2127                        .starts_with("Restoring cooling")
2128                    {
2129                        break;
2130                    }
2131                    assert!(Instant::now() < deadline);
2132                    tokio::time::sleep(Duration::from_millis(1)).await;
2133                }
2134                for (kind, value) in [(16, -15), (17, 0)] {
2135                    assert_eq!(
2136                        handle
2137                            .submit(kind, value, Duration::from_secs(5))
2138                            .unwrap()
2139                            .wait()
2140                            .await,
2141                        Ok(value)
2142                    );
2143                }
2144            };
2145            let (frame, ()) = tokio::join!(capture, commands);
2146            let frame = frame.unwrap();
2147            assert_eq!(frame.metadata["controls"]["16"], -15);
2148            assert_eq!(frame.metadata["controls"]["17"], 0);
2149            assert_eq!(frame.metadata["recoveries"], 0);
2150            assert!(!handle.pending());
2151            session.close().await;
2152        }
2153    }
2154    #[tokio::test]
2155    async fn direct_usb_read_size_survives_worker_replacement_and_leaves_sdk_bandwidth_unchanged() {
2156        let token = CancellationToken::new();
2157        let mut sel = selection(true);
2158        sel.recovery.direct_read_chunk_kib = 64;
2159        let mut s = Session::new(sel, runtime(), log()).unwrap();
2160        s.connect(&token).await.unwrap();
2161        assert!(!s.snapshot().controls[&6].writable);
2162        assert_eq!(s.snapshot().values[&6], 40);
2163        assert_eq!(
2164            s.capture(exposure(), &token).await.unwrap().metadata["readChunkKiB"],
2165            64
2166        );
2167        s.call("simulate-read-failures", json!({"count":3}), None, &token)
2168            .await
2169            .unwrap();
2170        let frame = s.capture(exposure(), &token).await.unwrap();
2171        assert_eq!(frame.metadata["recoveries"], 1);
2172        assert_eq!(frame.metadata["readChunkKiB"], 64);
2173        s.close().await;
2174    }
2175    #[test]
2176    fn direct_usb_read_size_configuration_preserves_key_default_and_bounds() {
2177        let old: RecoveryOptions = serde_json::from_value(json!({"directReadRetries":1})).unwrap();
2178        assert_eq!(old.direct_read_chunk_kib, 1024);
2179        let current: RecoveryOptions =
2180            serde_json::from_value(json!({"directReadChunkKiB":64})).unwrap();
2181        assert_eq!(current.direct_read_chunk_kib, 64);
2182        assert_eq!(
2183            serde_json::to_value(current).unwrap()["directReadChunkKiB"],
2184            64
2185        );
2186        for value in [0, 3, 512, 1024, 2048, u32::MAX] {
2187            let options = RecoveryOptions {
2188                direct_read_chunk_kib: value,
2189                ..RecoveryOptions::default()
2190            };
2191            assert_eq!(options.validate().is_ok(), matches!(value, 512 | 1024));
2192        }
2193    }
2194    #[tokio::test]
2195    async fn cooled_worker_recovery_restores_output_once_and_honors_warming_or_disable() {
2196        let token = CancellationToken::new();
2197        let mut sel = selection(true);
2198        sel.name = "ZWO ASI585MM Pro".into();
2199        let mut s = Session::new(sel, runtime(), log()).unwrap();
2200        s.connect(&token).await.unwrap();
2201        Session::queue_control(&s.status, 16, 10).unwrap();
2202        Session::queue_control(&s.status, 17, 1).unwrap();
2203        s.refresh(&token).await.unwrap();
2204        s.call(
2205            "resume-cooling",
2206            json!({"power":40,"temperature":25.0,"previousTarget":10}),
2207            None,
2208            &token,
2209        )
2210        .await
2211        .unwrap();
2212        s.refresh(&token).await.unwrap();
2213        assert_eq!(s.snapshot().values[&15], 40);
2214        s.call("simulate-read-failures", json!({"count":3}), None, &token)
2215            .await
2216            .unwrap();
2217        let frame = s.capture(exposure(), &token).await.unwrap();
2218        assert_eq!(frame.metadata["recoveries"], 1);
2219        assert_eq!(s.snapshot().values[&15], 40);
2220        assert!(s.cooling_seeded);
2221        // Ordinary refreshes must not repeatedly reseed the regulator.
2222        s.call(
2223            "resume-cooling",
2224            json!({"power":30,"temperature":25.0,"previousTarget":10}),
2225            None,
2226            &token,
2227        )
2228        .await
2229        .unwrap();
2230        s.refresh(&token).await.unwrap();
2231        assert_eq!(s.snapshot().values[&15], 30);
2232        s.invalidate().await;
2233        Session::queue_control(&s.status, 16, 11).unwrap();
2234        s.refresh(&token).await.unwrap();
2235        assert_eq!(s.snapshot().values[&15], 0); // A warmer requested target wins.
2236        s.invalidate().await;
2237        Session::queue_control(&s.status, 17, 0).unwrap();
2238        s.refresh(&token).await.unwrap();
2239        assert_eq!(s.snapshot().values[&17], 0);
2240        assert!(!s.cooling_seeded);
2241        s.close().await;
2242    }
2243
2244    #[tokio::test]
2245    async fn managed_white_balance_survives_worker_recovery_and_retains_locked_gains() {
2246        use crate::white_balance::{Gains, Mode, Output, Settings};
2247        for direct in [false, true] {
2248            let token = CancellationToken::new();
2249            let mut s = Session::new(selection(direct), runtime(), log()).unwrap();
2250            s.connect(&token).await.unwrap();
2251            assert_eq!(s.snapshot().white_balance_capabilities["supported"], true);
2252            let settings = Settings {
2253                mode: Mode::Manual,
2254                gains: Gains { red: 2., blue: 0.5 },
2255                output: Output::Corrected,
2256            };
2257            s.set_white_balance(settings, &token).await.unwrap();
2258            s.set_white_balance(
2259                Settings {
2260                    mode: Mode::Locked,
2261                    ..settings
2262                },
2263                &token,
2264            )
2265            .await
2266            .unwrap();
2267            assert!(Session::queue_control(&s.status, 3, 50).is_err());
2268            let e = Exposure {
2269                dark: false,
2270                ..exposure()
2271            };
2272            let first = s.capture(e.clone(), &token).await.unwrap();
2273            assert_eq!(first.metadata["whiteBalance"]["applied"], true);
2274            s.invalidate().await;
2275            let recovered = s.capture(e, &token).await.unwrap();
2276            assert_eq!(first.pixels, recovered.pixels);
2277            assert_eq!(
2278                s.snapshot().white_balance.unwrap(),
2279                Settings {
2280                    mode: Mode::Locked,
2281                    ..settings
2282                }
2283            );
2284            s.set_white_balance(
2285                Settings {
2286                    mode: Mode::Once,
2287                    ..Settings::default()
2288                },
2289                &token,
2290            )
2291            .await
2292            .unwrap();
2293            s.capture(
2294                Exposure {
2295                    dark: false,
2296                    ..exposure()
2297                },
2298                &token,
2299            )
2300            .await
2301            .unwrap();
2302            let effective = s.snapshot().white_balance.unwrap();
2303            assert_eq!(effective.mode, Mode::Locked);
2304            s.invalidate().await;
2305            s.capture(exposure(), &token).await.unwrap();
2306            assert_eq!(s.snapshot().white_balance.unwrap(), effective);
2307            s.invalidate().await;
2308        }
2309    }
2310    #[test]
2311    fn cooling_requires_output_or_sustained_setpoint_and_resets_on_warming_or_gaps() {
2312        let mut hold = CoolingHold::default();
2313        for i in 0..=15 {
2314            assert_eq!(
2315                hold.observe(Some(-10.), Some(20), -10., 2., 2., i as f64 * 2.),
2316                i == 15
2317            );
2318        }
2319        assert!(!hold.observe(Some(-9.7), Some(20), -10., 2., 2., 32.));
2320        assert!(!hold.observe(Some(-10.), Some(20), -10., 2., 2., 100.));
2321        assert!(!hold.observe(Some(-10.), Some(0), -10., 2., 2., 102.));
2322    }
2323    #[tokio::test]
2324    async fn direct_worker_acknowledges_capture_cooling_over_production_framing() {
2325        // Transport primitive only: Session's common acknowledged control queue
2326        // is a separate requirement. Never discover or activate physical cameras.
2327        for mode in ["still", "video"] {
2328            let token = CancellationToken::new();
2329            let mut worker = runtime().spawn(true, log()).await.unwrap();
2330            worker
2331                .call("open", json!({"name":"ZWO ASI585MM Pro"}), 15., &token)
2332                .await
2333                .unwrap();
2334            worker
2335                .call(
2336                    "start",
2337                    json!({"mode":mode,"maxFps":120.0,"width":64,"height":64,
2338                "x":0,"y":0,"bin":1,"microseconds":6_000_000,"dark":false}),
2339                    15.,
2340                    &token,
2341                )
2342                .await
2343                .unwrap();
2344            for (control, value) in [(16, -10), (17, 1), (16, -15)] {
2345                worker
2346                    .call("set", json!({"control":control,"value":value}), 15., &token)
2347                    .await
2348                    .unwrap();
2349                assert_eq!(
2350                    worker
2351                        .call("get", json!({"control":control}), 15., &token)
2352                        .await
2353                        .unwrap()
2354                        .0,
2355                    value
2356                );
2357                assert_eq!(
2358                    worker
2359                        .call("status", Value::Null, 15., &token)
2360                        .await
2361                        .unwrap()
2362                        .0,
2363                    1
2364                );
2365            }
2366            let deadline = Instant::now() + Duration::from_secs(15);
2367            while worker
2368                .call("status", Value::Null, 15., &token)
2369                .await
2370                .unwrap()
2371                .0
2372                == 1
2373            {
2374                assert!(Instant::now() < deadline);
2375                tokio::time::sleep(Duration::from_millis(10)).await;
2376            }
2377            let (metadata, pixels) = worker
2378                .call_image("download", Value::Null, 15., &token, 8192)
2379                .await
2380                .unwrap();
2381            assert_eq!(metadata["mode"], mode);
2382            assert_eq!(
2383                pixels,
2384                (0..4096u16).flat_map(u16::to_le_bytes).collect::<Vec<_>>()
2385            );
2386            worker
2387                .call("close", Value::Null, 15., &token)
2388                .await
2389                .unwrap();
2390            worker.kill().await;
2391        }
2392    }
2393
2394    #[tokio::test]
2395    async fn environment_refreshes_before_sdk_and_direct_exposures_finish() {
2396        for direct in [false, true] {
2397            let token = CancellationToken::new();
2398            let mut sel = selection(direct);
2399            if direct {
2400                sel.name = "ZWO ASI6200MM Pro".into();
2401            }
2402            let mut rt = runtime();
2403            rt.sdk_simulation = Some(json!({"instant":false}));
2404            let mut s = Session::new(sel, rt, log()).unwrap();
2405            s.connect(&token).await.unwrap();
2406            let shared = s.status.clone();
2407            let capture = async {
2408                s.capture(
2409                    Exposure {
2410                        microseconds: 6_000_000,
2411                        ..exposure()
2412                    },
2413                    &token,
2414                )
2415                .await
2416            };
2417            let observe = async {
2418                let deadline = Instant::now() + Duration::from_secs(5);
2419                while shared.lock().unwrap().phase != "Exposing" {
2420                    assert!(Instant::now() < deadline);
2421                    tokio::time::sleep(Duration::from_millis(10)).await;
2422                }
2423                // Make stale frontend values distinguishable from worker telemetry.
2424                {
2425                    let mut state = shared.lock().unwrap();
2426                    state.values.insert(8, 999);
2427                    state.values.insert(15, 99);
2428                }
2429                loop {
2430                    let state = shared.lock().unwrap().clone();
2431                    // The two worker reads update status separately. Wait for
2432                    // both before checking their values, including on slow CI.
2433                    if state.values[&8] != 999 && state.values[&15] != 99 {
2434                        assert_eq!(state.phase, "Exposing");
2435                        assert_eq!(state.values[&8], if direct { 250 } else { -100 });
2436                        assert_eq!(state.values[&15], if direct { 0 } else { 30 });
2437                        break;
2438                    }
2439                    assert!(
2440                        Instant::now() < deadline,
2441                        "Telemetry stayed stale during capture"
2442                    );
2443                    tokio::time::sleep(Duration::from_millis(20)).await;
2444                }
2445            };
2446            let (frame, ()) = tokio::join!(capture, observe);
2447            assert_eq!(frame.unwrap().pixels.len(), 8192);
2448            s.close().await;
2449        }
2450    }
2451    #[tokio::test]
2452    async fn usb_escalation_is_opt_in_bounded_and_obeys_exposure_limits() {
2453        for (threshold, retries, seconds, resets) in [
2454            (0, 3, 0.01, 0),
2455            (2, 3, 0.01, 1),
2456            (1, 0, 0.01, 0),
2457            (1, 3, 31., 0),
2458            (4, 3, 0.01, 0),
2459        ] {
2460            let events = Arc::new(Mutex::new(Vec::<String>::new()));
2461            let sink = events.clone();
2462            let mut sel = selection(false);
2463            sel.recovery.usb_reset_after_failures = threshold;
2464            sel.recovery.max_retries = retries;
2465            sel.recovery.ready_frame_download_retries = 0;
2466            let mut rt = runtime();
2467            rt.sdk_simulation = Some(json!({"instant":true,"fault":"download"}));
2468            let mut session = Session::new(
2469                sel,
2470                rt,
2471                Arc::new(move |_, event, _| sink.lock().unwrap().push(event.into())),
2472            )
2473            .unwrap();
2474            let token = CancellationToken::new();
2475            session.connect(&token).await.unwrap();
2476            assert!(
2477                session
2478                    .capture(
2479                        Exposure {
2480                            microseconds: (seconds * 1e6) as u64,
2481                            ..exposure()
2482                        },
2483                        &token
2484                    )
2485                    .await
2486                    .is_err()
2487            );
2488            assert_eq!(
2489                events
2490                    .lock()
2491                    .unwrap()
2492                    .iter()
2493                    .filter(|e| e.as_str() == "usb.reset")
2494                    .count(),
2495                resets
2496            );
2497            session.close().await;
2498        }
2499    }
2500    #[tokio::test]
2501    async fn usb_recovery_restores_controls_and_cancellation_never_resets() {
2502        let token = CancellationToken::new();
2503        let mut sel = selection(false);
2504        sel.recovery.usb_reset_after_failures = 1;
2505        sel.recovery.ready_frame_download_retries = 0;
2506        let mut session = Session::new(sel, runtime(), log()).unwrap();
2507        session.connect(&token).await.unwrap();
2508        Session::queue_control(&session.status, 0, 230).unwrap();
2509        session
2510            .call("fault", json!({"kind":"download"}), None, &token)
2511            .await
2512            .unwrap();
2513        let frame = session.capture(exposure(), &token).await.unwrap();
2514        assert_eq!(frame.metadata["usbResets"], 1);
2515        assert_eq!(frame.metadata["recoveries"], 1);
2516        assert_eq!(frame.metadata["controls"]["0"], 230);
2517        let cancelled = CancellationToken::new();
2518        cancelled.cancel();
2519        assert!(
2520            session
2521                .capture(exposure(), &cancelled)
2522                .await
2523                .err()
2524                .unwrap()
2525                .is::<Failure>()
2526        );
2527        session.close().await;
2528    }
2529    #[tokio::test]
2530    async fn sdk_ready_frame_read_retry_does_not_replace_exposure() {
2531        let token = CancellationToken::new();
2532        let mut s = Session::new(selection(false), runtime(), log()).unwrap();
2533        s.connect(&token).await.unwrap();
2534        s.call("fault", json!({"kind":"download"}), None, &token)
2535            .await
2536            .unwrap();
2537        let frame = s.capture(exposure(), &token).await.unwrap();
2538        assert_eq!(frame.metadata["recoveries"], 0);
2539        assert_eq!(frame.metadata["downloadRetries"], 1);
2540        assert_eq!(frame.pixels.len(), 8192);
2541        let status = s.snapshot();
2542        assert_eq!(status.retry.downloads, 1);
2543        assert_eq!(status.retry.recaptures, 0);
2544        assert!(status.recovery_info().contains("state: Idle; retries: 1"));
2545        assert!(status.retry.last_failure.as_deref().unwrap().contains("11"));
2546        s.capture(exposure(), &token).await.unwrap();
2547        assert_eq!(s.snapshot().retry.downloads, 0);
2548        assert_eq!(s.snapshot().retry.last_failure, status.retry.last_failure);
2549        s.close().await;
2550    }
2551    #[tokio::test]
2552    async fn sdk_replacement_restores_controls_and_respects_long_exposure_limit() {
2553        for (seconds, allowed) in [(30., true), (30.000001, false), (1200., false)] {
2554            let token = CancellationToken::new();
2555            let mut sel = selection(false);
2556            sel.recovery.ready_frame_download_retries = 0;
2557            let mut s = Session::new(sel, runtime(), log()).unwrap();
2558            s.connect(&token).await.unwrap();
2559            Session::queue_control(&s.status, 0, 250).unwrap();
2560            Session::queue_control(&s.status, 16, -10).unwrap();
2561            s.call("fault", json!({"kind":"download"}), None, &token)
2562                .await
2563                .unwrap();
2564            let result = s
2565                .capture(
2566                    Exposure {
2567                        microseconds: (seconds * 1e6) as u64,
2568                        ..exposure()
2569                    },
2570                    &token,
2571                )
2572                .await;
2573            assert_eq!(result.is_ok(), allowed);
2574            if let Ok(frame) = result {
2575                assert_eq!(frame.metadata["recoveries"], 1);
2576                assert_eq!(frame.metadata["controls"]["0"], 250);
2577                assert_eq!(s.snapshot().retry.recaptures, 1);
2578                assert!(
2579                    s.snapshot()
2580                        .recovery_info()
2581                        .contains("state: Idle; retries: 1")
2582                );
2583            } else {
2584                assert_eq!(s.snapshot().retry.recaptures, 0);
2585                assert!(
2586                    s.snapshot()
2587                        .recovery_info()
2588                        .contains("state: Error; retries: 0")
2589                );
2590            }
2591            s.close().await;
2592        }
2593    }
2594    #[tokio::test]
2595    async fn direct_retained_reads_and_abort_share_the_recovery_session() {
2596        let token = CancellationToken::new();
2597        let mut s = Session::new(selection(true), runtime(), log()).unwrap();
2598        s.connect(&token).await.unwrap();
2599        s.call("simulate-read-failures", json!({"count":2}), None, &token)
2600            .await
2601            .unwrap();
2602        let frame = s.capture(exposure(), &token).await.unwrap();
2603        assert_eq!(frame.metadata["readRecoveries"], 2);
2604        assert_eq!(frame.metadata["recoveries"], 0);
2605        assert_eq!(s.snapshot().retry.usb_reads, 2);
2606        assert!(
2607            s.snapshot()
2608                .recovery_info()
2609                .contains("state: Idle; retries: 2")
2610        );
2611        let abort = CancellationToken::new();
2612        let cancel = abort.clone();
2613        tokio::spawn(async move {
2614            tokio::time::sleep(Duration::from_millis(20)).await;
2615            cancel.cancel();
2616        });
2617        assert!(
2618            s.capture(
2619                Exposure {
2620                    microseconds: 2000000,
2621                    ..exposure()
2622                },
2623                &abort
2624            )
2625            .await
2626            .is_err()
2627        );
2628        assert_eq!(
2629            s.capture(exposure(), &token).await.unwrap().pixels.len(),
2630            8192
2631        );
2632        s.close().await;
2633    }
2634    #[tokio::test]
2635    async fn intentional_abort_keeps_worker_cooling_and_does_not_retry() {
2636        for direct in [false, true] {
2637            let token = CancellationToken::new();
2638            let mut sel = selection(direct);
2639            if direct {
2640                sel.name = "ZWO ASI6200MM Pro".into();
2641            }
2642            let mut rt = runtime();
2643            rt.sdk_simulation = Some(json!({"instant":false,"temperature":250,"coolerPower":70}));
2644            let events = Arc::new(Mutex::new(Vec::new()));
2645            let recorded = events.clone();
2646            let mut s = Session::new(
2647                sel,
2648                rt,
2649                Arc::new(move |_, event, _| recorded.lock().unwrap().push(event.to_owned())),
2650            )
2651            .unwrap();
2652            s.connect(&token).await.unwrap();
2653            Session::queue_control(&s.status, 16, 25).unwrap();
2654            Session::queue_control(&s.status, 17, 1).unwrap();
2655            s.refresh(&token).await.unwrap();
2656            if direct {
2657                s.call(
2658                    "resume-cooling",
2659                    json!({"power":70,"temperature":25.,"previousTarget":25}),
2660                    None,
2661                    &token,
2662                )
2663                .await
2664                .unwrap();
2665                s.refresh(&token).await.unwrap();
2666            }
2667            let pid = s.snapshot().process_id;
2668            let abort = CancellationToken::new();
2669            let cancel = abort.clone();
2670            let status = s.status.clone();
2671            let cancellation = tokio::spawn(async move {
2672                exposing(&status).await;
2673                cancel.cancel();
2674            });
2675            let result = tokio::time::timeout(
2676                Duration::from_secs(5),
2677                s.capture(
2678                    Exposure {
2679                        microseconds: 600_000_000,
2680                        ..exposure()
2681                    },
2682                    &abort,
2683                ),
2684            )
2685            .await
2686            .unwrap();
2687            cancellation.await.unwrap();
2688            assert!(matches!(
2689                result.err().unwrap().downcast_ref::<Failure>(),
2690                Some(Failure::Cancelled)
2691            ));
2692            assert_eq!(s.snapshot().process_id, pid);
2693            assert!(s.snapshot().control_connection_available);
2694            assert_eq!(s.snapshot().values[&17], 1);
2695            assert_eq!(s.snapshot().values[&15], 70);
2696            assert!(s.snapshot().retry.last_failure.is_none());
2697            assert_eq!(
2698                s.capture(exposure(), &token).await.unwrap().metadata["recoveries"],
2699                0
2700            );
2701            assert_eq!(s.snapshot().process_id, pid);
2702            assert!(
2703                !events
2704                    .lock()
2705                    .unwrap()
2706                    .iter()
2707                    .any(|e| e == "capture.failed" || e == "capture.retry")
2708            );
2709            s.close().await;
2710        }
2711    }
2712    #[tokio::test]
2713    async fn thermal_cleanup_does_not_restore_white_balance_or_reactivate_controls() {
2714        let token = CancellationToken::new();
2715        let mut s = Session::new(selection(false), runtime(), log()).unwrap();
2716        s.connect(&token).await.unwrap();
2717        s.set_white_balance(
2718            crate::white_balance::Settings {
2719                mode: crate::white_balance::Mode::Manual,
2720                ..Default::default()
2721            },
2722            &token,
2723        )
2724        .await
2725        .unwrap();
2726        s.invalidate().await;
2727        s.status.lock().unwrap().connected = false;
2728        s.open_worker(&token, OpenPurpose::ThermalShutdown)
2729            .await
2730            .unwrap();
2731        assert!(!s.snapshot().connected);
2732        assert!(!s.snapshot().control_connection_available);
2733        assert_eq!(s.snapshot().phase, "Closing camera");
2734        assert!(s.cooling().submit(17, 1, Duration::from_secs(1)).is_err());
2735        let actual = s
2736            .call("white-balance", Value::Null, None, &token)
2737            .await
2738            .unwrap()
2739            .0;
2740        assert_eq!(actual["managed"], false);
2741        s.close().await;
2742    }
2743    #[tokio::test]
2744    async fn failed_abort_retires_worker_before_another_capture() {
2745        let token = CancellationToken::new();
2746        let mut s = Session::new(selection(true), runtime(), log()).unwrap();
2747        s.connect(&token).await.unwrap();
2748        let pid = s.snapshot().process_id;
2749        s.call("simulation", json!({"cleanupFailure":true}), None, &token)
2750            .await
2751            .unwrap();
2752        let abort = CancellationToken::new();
2753        let cancel = abort.clone();
2754        let status = s.status.clone();
2755        let cancellation = tokio::spawn(async move {
2756            exposing(&status).await;
2757            cancel.cancel();
2758        });
2759        let result = s
2760            .capture(
2761                Exposure {
2762                    microseconds: 600_000_000,
2763                    ..exposure()
2764                },
2765                &abort,
2766            )
2767            .await;
2768        cancellation.await.unwrap();
2769        assert!(matches!(
2770            result.err().unwrap().downcast_ref::<Failure>(),
2771            Some(Failure::Cancelled)
2772        ));
2773        assert!(!s.snapshot().control_connection_available);
2774        assert_eq!(s.snapshot().process_id, None);
2775        assert!(
2776            s.snapshot()
2777                .retry
2778                .last_failure
2779                .unwrap()
2780                .contains("Stop was not acknowledged")
2781        );
2782        s.capture(exposure(), &token).await.unwrap();
2783        assert_ne!(s.snapshot().process_id, pid);
2784        s.close().await;
2785    }
2786    #[tokio::test]
2787    async fn direct_fallback_revalidates_same_identity_and_uses_one_budget() {
2788        let token = CancellationToken::new();
2789        let mut sel = selection(true);
2790        sel.sdk_fallback = true;
2791        let mut runtime = runtime();
2792        runtime.sdk_simulation = Some(
2793            json!({"name":"ZWO ASI676MC","width":3552,"height":3552,"bins":[1],"serial":"direct-simulator","cooled":false,"instant":true}),
2794        );
2795        let mut s = Session::new(sel, runtime, log()).unwrap();
2796        s.connect(&token).await.unwrap();
2797        s.call("simulate-read-failures", json!({"count":3}), None, &token)
2798            .await
2799            .unwrap();
2800        let frame = s.capture(exposure(), &token).await.unwrap();
2801        assert_eq!(frame.metadata["backend"], "sdk");
2802        assert_eq!(frame.metadata["sdkFallback"], true);
2803        assert_eq!(frame.metadata["recoveries"], 1);
2804        assert_eq!(frame.metadata["retainedReadRetries"], 2);
2805        assert_eq!(frame.metadata["readRecoveries"].as_u64().unwrap_or(0), 0);
2806        s.close().await;
2807    }
2808    #[test]
2809    fn invalid_limits_are_rejected() {
2810        assert!(
2811            RecoveryOptions {
2812                maximum_retry_exposure_seconds: f64::NAN,
2813                ..RecoveryOptions::default()
2814            }
2815            .validate()
2816            .is_err()
2817        );
2818        assert!(
2819            RecoveryOptions {
2820                direct_read_retries: 6,
2821                ..RecoveryOptions::default()
2822            }
2823            .validate()
2824            .is_err()
2825        );
2826    }
2827}