Skip to main content

a3s_runtime/state/
record.rs

1use crate::contract::{
2    RuntimeActionRequest, RuntimeApplyRequest, RuntimeExecRequest, RuntimeExecResult,
3    RuntimeObservation, RuntimeRemoval, RuntimeUnitSpec,
4};
5use serde::{Deserialize, Serialize};
6
7#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
8#[serde(rename_all = "snake_case")]
9pub enum RuntimeActionKind {
10    Stop,
11    Remove,
12}
13
14#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
15#[serde(rename_all = "snake_case")]
16pub enum RuntimeRequestKind {
17    Apply,
18    Stop,
19    Remove,
20    Exec,
21}
22
23impl From<RuntimeActionKind> for RuntimeRequestKind {
24    fn from(value: RuntimeActionKind) -> Self {
25        match value {
26            RuntimeActionKind::Stop => Self::Stop,
27            RuntimeActionKind::Remove => Self::Remove,
28        }
29    }
30}
31
32#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
33#[serde(rename_all = "snake_case")]
34pub enum RuntimeRequestState {
35    Pending,
36    Completed,
37}
38
39#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
40#[serde(deny_unknown_fields)]
41pub struct RuntimeRequestReceipt {
42    pub schema: String,
43    pub request_id: String,
44    pub unit_id: String,
45    pub generation: u64,
46    pub kind: RuntimeRequestKind,
47    pub request_digest: String,
48    /// Effective absolute deadline captured when the request was first reserved.
49    ///
50    /// Exec always has a deadline because its relative timeout is converted to
51    /// an absolute value. Persisting it prevents a pending replay from extending
52    /// the original execution budget.
53    pub deadline_at_ms: Option<u64>,
54    pub state: RuntimeRequestState,
55    pub observation: Option<RuntimeObservation>,
56    pub removal: Option<RuntimeRemoval>,
57    pub exec_result: Option<RuntimeExecResult>,
58}
59
60impl RuntimeRequestReceipt {
61    pub const SCHEMA: &'static str = "a3s.runtime.request-receipt.v2";
62
63    pub(crate) fn pending_apply(request: &RuntimeApplyRequest) -> Result<Self, String> {
64        Ok(Self::pending(
65            request.request_id.clone(),
66            request.spec.unit_id.clone(),
67            request.spec.generation,
68            RuntimeRequestKind::Apply,
69            request.digest()?,
70            request.deadline_at_ms,
71        ))
72    }
73
74    pub(crate) fn pending_action(
75        kind: RuntimeActionKind,
76        request: &RuntimeActionRequest,
77    ) -> Result<Self, String> {
78        Ok(Self::pending(
79            request.request_id.clone(),
80            request.unit_id.clone(),
81            request.generation,
82            kind.into(),
83            request.digest()?,
84            request.deadline_at_ms,
85        ))
86    }
87
88    pub(crate) fn pending_exec(
89        request: &RuntimeExecRequest,
90        started_at_ms: u64,
91    ) -> Result<Self, String> {
92        let relative_deadline = started_at_ms.saturating_add(request.timeout_ms);
93        let deadline_at_ms = request
94            .deadline_at_ms
95            .map_or(relative_deadline, |absolute| {
96                absolute.min(relative_deadline)
97            });
98        Ok(Self::pending(
99            request.request_id.clone(),
100            request.unit_id.clone(),
101            request.generation,
102            RuntimeRequestKind::Exec,
103            request.digest()?,
104            Some(deadline_at_ms),
105        ))
106    }
107
108    fn pending(
109        request_id: String,
110        unit_id: String,
111        generation: u64,
112        kind: RuntimeRequestKind,
113        request_digest: String,
114        deadline_at_ms: Option<u64>,
115    ) -> Self {
116        Self {
117            schema: Self::SCHEMA.into(),
118            request_id,
119            unit_id,
120            generation,
121            kind,
122            request_digest,
123            deadline_at_ms,
124            state: RuntimeRequestState::Pending,
125            observation: None,
126            removal: None,
127            exec_result: None,
128        }
129    }
130
131    pub(crate) fn complete_with_observation(&mut self, observation: RuntimeObservation) {
132        self.state = RuntimeRequestState::Completed;
133        self.observation = Some(observation);
134        self.removal = None;
135        self.exec_result = None;
136    }
137
138    pub(crate) fn complete_with_removal(&mut self, removal: RuntimeRemoval) {
139        self.state = RuntimeRequestState::Completed;
140        self.observation = None;
141        self.removal = Some(removal);
142        self.exec_result = None;
143    }
144
145    pub(crate) fn complete_with_exec_result(&mut self, result: RuntimeExecResult) {
146        self.state = RuntimeRequestState::Completed;
147        self.observation = None;
148        self.removal = None;
149        self.exec_result = Some(result);
150    }
151
152    pub fn validate(&self) -> Result<(), String> {
153        if self.schema != Self::SCHEMA {
154            return Err(format!(
155                "unsupported Runtime request receipt schema {:?}",
156                self.schema
157            ));
158        }
159        crate::contract::validate_id("request_id", &self.request_id, 512)?;
160        crate::contract::validate_id("unit_id", &self.unit_id, 512)?;
161        if self.generation == 0 {
162            return Err("Runtime request receipt generation must be positive".into());
163        }
164        if self.deadline_at_ms == Some(0)
165            || (self.kind == RuntimeRequestKind::Exec && self.deadline_at_ms.is_none())
166        {
167            return Err("Runtime request receipt deadline is invalid".into());
168        }
169        crate::contract::validate_digest(&self.request_digest)?;
170        match (
171            self.kind,
172            self.state,
173            &self.observation,
174            &self.removal,
175            &self.exec_result,
176        ) {
177            (_, RuntimeRequestState::Pending, None, None, None) => Ok(()),
178            (
179                RuntimeRequestKind::Apply | RuntimeRequestKind::Stop,
180                RuntimeRequestState::Completed,
181                Some(observation),
182                None,
183                None,
184            ) => {
185                observation.validate()?;
186                self.validate_unit_result(&observation.unit_id, observation.generation)
187            }
188            (
189                RuntimeRequestKind::Remove,
190                RuntimeRequestState::Completed,
191                None,
192                Some(removal),
193                None,
194            ) => {
195                removal.validate()?;
196                if removal.request_id != self.request_id {
197                    return Err("Runtime removal receipt request identity mismatch".into());
198                }
199                self.validate_unit_result(&removal.unit_id, removal.generation)
200            }
201            (
202                RuntimeRequestKind::Exec,
203                RuntimeRequestState::Completed,
204                None,
205                None,
206                Some(result),
207            ) => {
208                result.validate()?;
209                if result.request_id != self.request_id {
210                    return Err("Runtime exec receipt request identity mismatch".into());
211                }
212                self.validate_unit_result(
213                    &result.observation.unit_id,
214                    result.observation.generation,
215                )
216            }
217            _ => Err("Runtime request receipt result does not match its kind and state".into()),
218        }
219    }
220
221    fn validate_unit_result(&self, unit_id: &str, generation: u64) -> Result<(), String> {
222        if unit_id != self.unit_id || generation != self.generation {
223            return Err("Runtime request receipt result identity mismatch".into());
224        }
225        Ok(())
226    }
227}
228
229#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
230#[serde(deny_unknown_fields)]
231pub struct RuntimeUnitRecord {
232    pub schema: String,
233    pub spec: RuntimeUnitSpec,
234    pub observation: RuntimeObservation,
235    pub removed_at_ms: Option<u64>,
236}
237
238impl RuntimeUnitRecord {
239    pub const SCHEMA: &'static str = "a3s.runtime.unit-record.v2";
240
241    pub(crate) fn new(request: &RuntimeApplyRequest, now_ms: u64) -> Result<Self, String> {
242        let record = Self {
243            schema: Self::SCHEMA.into(),
244            spec: request.spec.clone(),
245            observation: RuntimeObservation::accepted(&request.spec, now_ms)?,
246            removed_at_ms: None,
247        };
248        record.validate()?;
249        Ok(record)
250    }
251
252    pub fn validate(&self) -> Result<(), String> {
253        if self.schema != Self::SCHEMA {
254            return Err(format!(
255                "unsupported Runtime unit record schema {:?}",
256                self.schema
257            ));
258        }
259        self.spec.validate()?;
260        self.observation.validate_against(&self.spec)
261    }
262}
263
264#[derive(Debug, Clone, PartialEq, Eq)]
265pub struct RuntimeStateReservation {
266    pub dispatch: bool,
267    pub record: RuntimeUnitRecord,
268    pub receipt: RuntimeRequestReceipt,
269}