Skip to main content

regain_core/
model.rs

1use anyhow::{Result, ensure};
2use serde::{Deserialize, Serialize};
3use serde_json::Value;
4use std::{
5    collections::BTreeMap,
6    sync::{Arc, Mutex},
7};
8
9#[derive(Debug)]
10pub enum Failure {
11    Invalid(String),
12    Cancelled,
13    Worker {
14        message: String,
15        code: Option<i32>,
16    },
17    /// A dispatched control has no known acknowledgement. Retire its worker;
18    /// retrying the command or capture could repeat an already applied write.
19    UncertainControl {
20        message: String,
21        code: Option<i32>,
22    },
23}
24impl std::fmt::Display for Failure {
25    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
26        match self {
27            Self::Invalid(m)
28            | Self::Worker { message: m, .. }
29            | Self::UncertainControl { message: m, .. } => f.write_str(m),
30            Self::Cancelled => f.write_str("Exposure aborted"),
31        }
32    }
33}
34impl std::error::Error for Failure {}
35pub fn invalid(message: impl Into<String>) -> anyhow::Error {
36    Failure::Invalid(message.into()).into()
37}
38pub fn retryable(error: &anyhow::Error) -> bool {
39    match error.downcast_ref::<Failure>() {
40        Some(Failure::Invalid(_) | Failure::Cancelled | Failure::UncertainControl { .. }) => false,
41        Some(Failure::Worker { code: Some(c), .. }) => {
42            matches!(c, 1 | 2 | 4 | 5 | 11 | 12 | 15 | 16)
43        }
44        _ => true,
45    }
46}
47#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
48pub struct Exposure {
49    pub width: u32,
50    pub height: u32,
51    pub bin: u32,
52    pub x: u32,
53    pub y: u32,
54    pub microseconds: u64,
55    pub dark: bool,
56}
57impl Exposure {
58    pub fn bytes(&self) -> Result<usize> {
59        let count = u64::from(self.width)
60            .checked_mul(u64::from(self.height))
61            .and_then(|pixels| pixels.checked_mul(2))
62            .ok_or_else(|| invalid("Invalid frame size"))?;
63        ensure!(
64            count > 0 && count <= 512 * 1024 * 1024,
65            Failure::Invalid("Invalid frame size".into())
66        );
67        Ok(count as usize)
68    }
69}
70#[derive(Clone, Debug, Serialize, Deserialize)]
71pub struct Control {
72    #[serde(rename = "type")]
73    pub kind: i32,
74    pub min: i64,
75    pub max: i64,
76    pub value: i64,
77    pub writable: bool,
78}
79pub use crate::recovery::RecoveryOptions;
80#[derive(Clone, Debug, Serialize, Deserialize)]
81#[serde(rename_all = "camelCase")]
82pub struct Selection {
83    pub name: String,
84    pub serial: Option<String>,
85    pub direct: bool,
86    pub sdk_fallback: bool,
87    pub recovery: RecoveryOptions,
88}
89#[derive(Clone, Debug, Default, Serialize)]
90#[serde(rename_all = "camelCase")]
91pub struct Status {
92    pub info: Value,
93    pub serial: Option<String>,
94    pub sdk_version: String,
95    pub backend: String,
96    pub sdk_fallback: bool,
97    pub phase: String,
98    pub error: Option<String>,
99    pub retry: RetryStatus,
100    pub controls: BTreeMap<i32, Control>,
101    pub values: BTreeMap<i32, i64>,
102    /// Acknowledged values, separate from queued desired settings. Monotonic
103    /// timestamps belong to this process and never enter the status JSON.
104    #[serde(skip)]
105    pub observations: BTreeMap<i32, ControlObservation>,
106    pub connected: bool,
107    pub control_connection_available: bool,
108    pub sdk_exposure_state: Option<i64>,
109    pub sdk_error_code: Option<i32>,
110    pub process_id: Option<u32>,
111    pub white_balance_capabilities: Value,
112    pub white_balance: Option<crate::white_balance::Settings>,
113}
114/// Attempts made for the current/latest capture; the last failure survives
115/// successful recovery and subsequent healthy captures in this connection.
116#[derive(Clone, Debug, Default, Serialize)]
117#[serde(rename_all = "camelCase")]
118pub struct RetryStatus {
119    pub recaptures: u32,
120    pub downloads: u32,
121    pub usb_reads: u32,
122    pub last_failure: Option<String>,
123}
124impl Status {
125    pub fn recovery_info(&self) -> String {
126        let retry = &self.retry;
127        let failure = retry.last_failure.as_deref().unwrap_or("none");
128        let failure = failure.split_whitespace().collect::<Vec<_>>().join(" ");
129        format!(
130            "state: {}; retries: {} (recaptures: {}, downloads: {}, USB reads: {}); last failure: {}",
131            if self.phase.is_empty() {
132                "Disconnected"
133            } else {
134                &self.phase
135            },
136            u64::from(retry.recaptures) + u64::from(retry.downloads) + u64::from(retry.usb_reads),
137            retry.recaptures,
138            retry.downloads,
139            retry.usb_reads,
140            failure
141        )
142    }
143}
144pub type SharedStatus = Arc<Mutex<Status>>;
145pub type Diagnostic = Arc<dyn Fn(&str, &str, &str) + Send + Sync>;
146#[derive(Clone, Copy, Debug, PartialEq, Eq)]
147pub struct ControlObservation {
148    pub value: i64,
149    pub observed_at: tokio::time::Instant,
150}
151/// Worker-relative age; never transfer a process-local monotonic timestamp.
152/// Receivers must also account for the request/response transit time.
153#[derive(Clone, Debug, Serialize, Deserialize)]
154#[serde(rename_all = "camelCase", deny_unknown_fields)]
155pub struct ControlObservationReply {
156    pub value: i64,
157    pub age_seconds: f64,
158}
159impl ControlObservationReply {
160    pub fn normalize(&self, request_started: tokio::time::Instant) -> Result<ControlObservation> {
161        Ok(ControlObservation {
162            value: self.value,
163            observed_at: self.observed_at(request_started)?,
164        })
165    }
166    /// Conservatively include all IPC/worker time by subtracting the reported
167    /// age from request admission, never from response receipt. Process-local
168    /// monotonic clock epochs do not need to agree across the pipe.
169    pub fn observed_at(
170        &self,
171        request_started: tokio::time::Instant,
172    ) -> Result<tokio::time::Instant> {
173        let age = std::time::Duration::try_from_secs_f64(self.age_seconds)
174            .map_err(|_| invalid("Invalid control observation age"))?;
175        request_started
176            .checked_sub(age)
177            .ok_or_else(|| invalid("Unrepresentable control observation time"))
178    }
179}
180#[derive(Clone)]
181pub struct Frame {
182    pub exposure: Exposure,
183    pub pixels: Arc<[u8]>,
184    pub metadata: Value,
185}
186
187pub fn validate_exposure(
188    info: &Value,
189    controls: &BTreeMap<i32, Control>,
190    e: &Exposure,
191) -> Result<()> {
192    e.bytes()?;
193    let bins = info["bins"]
194        .as_array()
195        .ok_or_else(|| invalid("Camera binning is unavailable"))?;
196    let w = info["width"].as_u64().unwrap_or(0);
197    let h = info["height"].as_u64().unwrap_or(0);
198    let exp = controls
199        .get(&1)
200        .ok_or_else(|| invalid("Exposure control unavailable"))?;
201    let alignment = info["originAlignment"]
202        .as_u64()
203        .or_else(|| info["originAlignmentX"].as_u64())
204        .unwrap_or(1)
205        .max(1);
206    let align_y = info["originAlignmentY"]
207        .as_u64()
208        .unwrap_or(if info["originAlignment"].is_number() {
209            2
210        } else {
211            1
212        })
213        .max(1);
214    ensure!(e.bin>0 && bins.iter().any(|b|b.as_u64()==Some(e.bin as u64)) && e.width.is_multiple_of(8) && e.height.is_multiple_of(2) &&
215        (u64::from(e.x)+u64::from(e.width))*u64::from(e.bin)<=w && (u64::from(e.y)+u64::from(e.height))*u64::from(e.bin)<=h &&
216        u64::from(e.width)*u64::from(e.bin)>=info["minimumWidth"].as_u64().unwrap_or(8) && u64::from(e.height)*u64::from(e.bin)>=info["minimumHeight"].as_u64().unwrap_or(2) &&
217        u64::from(e.x)*u64::from(e.bin)%alignment==0 && u64::from(e.y)*u64::from(e.bin)%align_y==0 &&
218        e.microseconds>0 && e.microseconds>=exp.min.max(0) as u64 && e.microseconds<=exp.max.max(0) as u64,
219        Failure::Invalid("Exposure or ROI is outside camera capabilities (width multiple of 8, height multiple of 2; see camera alignment)".into()));
220    Ok(())
221}
222
223/// Validate against the SDK's ROI rules when fallback is explicitly permitted.
224/// The worker's validate command decides whether it must switch before exposing.
225pub fn validate_capture(
226    info: &Value,
227    controls: &BTreeMap<i32, Control>,
228    e: &Exposure,
229    sdk_fallback: bool,
230) -> Result<()> {
231    if !sdk_fallback {
232        return validate_exposure(info, controls, e);
233    }
234    let mut info = info.clone();
235    info["minimumWidth"] = serde_json::json!(8);
236    info["minimumHeight"] = serde_json::json!(2);
237    info["originAlignment"] = serde_json::json!(1);
238    info["originAlignmentX"] = serde_json::json!(1);
239    info["originAlignmentY"] = serde_json::json!(1);
240    validate_exposure(&info, controls, e)
241}
242
243#[cfg(test)]
244mod observation_tests {
245    use super::*;
246    use std::time::Duration;
247    use tokio::time::Instant;
248
249    #[test]
250    fn worker_age_normalization_includes_transit_and_rejects_invalid_times() {
251        let request = Instant::now();
252        let reply = ControlObservationReply {
253            value: -100,
254            age_seconds: 20.25,
255        };
256        let observed = reply.observed_at(request).unwrap();
257        let received = request + Duration::from_secs(5);
258        assert_eq!(
259            received.duration_since(observed),
260            Duration::from_secs_f64(25.25)
261        );
262        for invalid_age in [-1.0, f64::NAN, f64::INFINITY, f64::NEG_INFINITY, f64::MAX] {
263            assert!(
264                ControlObservationReply {
265                    value: 0,
266                    age_seconds: invalid_age
267                }
268                .observed_at(request)
269                .is_err()
270            );
271        }
272        assert_eq!(
273            ControlObservationReply {
274                value: 0,
275                age_seconds: 0.0
276            }
277            .observed_at(request)
278            .unwrap(),
279            request
280        );
281        for invalid in [
282            serde_json::json!({"value":true,"ageSeconds":0}),
283            serde_json::json!({"value":0,"ageSeconds":"0"}),
284            serde_json::json!({"value":0}),
285            serde_json::json!({"value":0,"ageSeconds":0,"timestamp":"private clock epoch"}),
286        ] {
287            assert!(serde_json::from_value::<ControlObservationReply>(invalid).is_err());
288        }
289    }
290}