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}
21pub 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 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 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(()), }
172 }
173 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 self.invalidate().await;
329 Err(Failure::Cancelled.into())
330 }
331 Err(error) => Err(error.into()),
332 Ok(actual) => Ok(actual),
333 }
334 }
335 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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); 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 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 {
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 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}