Skip to main content

hara_native/live_session/
interpreter.rs

1use serde_json::{json, Value as JsonValue};
2
3use super::{
4    required_text, LiveBackend, LiveReplacementPolicy, LiveSession, LiveSessionCapabilities,
5    LiveSessionCommand, LiveSessionError, LiveSessionOperation, LiveSessionState,
6    LiveSessionStatus, LiveSettlement, LiveSource,
7};
8
9pub struct InterpreterLiveSession {
10    session_id: String,
11    source: LiveSource,
12    handle: Option<u64>,
13    generation_base: u64,
14    backend_generation: u64,
15    sequence: u64,
16    status: LiveSessionStatus,
17    pending_source: Option<LiveSource>,
18}
19
20impl InterpreterLiveSession {
21    pub fn start(
22        session_id: impl Into<String>,
23        source: LiveSource,
24    ) -> Result<Self, LiveSessionError> {
25        let session_id = required_text(session_id.into(), "session id")?;
26        let (handle, info) = start_backend(&session_id, &source)?;
27        let mut session = Self {
28            session_id,
29            source,
30            handle: Some(handle),
31            generation_base: 0,
32            backend_generation: 0,
33            sequence: 0,
34            status: LiveSessionStatus::Ready,
35            pending_source: None,
36        };
37        session.sync_info(&info)?;
38        Ok(session)
39    }
40
41    pub fn pending_revision(&self) -> Option<&str> {
42        self.pending_source.as_ref().map(LiveSource::revision)
43    }
44
45    fn generation(&self) -> u64 {
46        self.generation_base.saturating_add(self.backend_generation)
47    }
48
49    fn handle(&self) -> Result<u64, LiveSessionError> {
50        self.handle.ok_or_else(|| {
51            LiveSessionError::new(
52                "live-session/disposed",
53                "interpreter live session has been disposed",
54            )
55        })
56    }
57
58    fn refresh(&mut self) -> Result<(), LiveSessionError> {
59        let handle = self.handle()?;
60        let info = invoke_legacy(json!({"op": "info", "handle": handle}))?;
61        self.sync_info(&info)
62    }
63
64    fn sync_info(&mut self, info: &JsonValue) -> Result<(), LiveSessionError> {
65        self.backend_generation = required_u64(info, "generation")?;
66        self.sequence = required_u64(info, "sequence")?;
67        self.status = parse_status(required_string(info, "status")?)?;
68        Ok(())
69    }
70
71    fn invoke_handle(&self, mut request: JsonValue) -> Result<JsonValue, LiveSessionError> {
72        let handle = self.handle()?;
73        let object = request.as_object_mut().ok_or_else(|| {
74            LiveSessionError::new(
75                "live-session/internal",
76                "interpreter backend request must be a JSON object",
77            )
78        })?;
79        object.insert("handle".into(), JsonValue::from(handle));
80        invoke_legacy(request)
81    }
82
83    fn invoke_and_refresh(&mut self, request: JsonValue) -> Result<JsonValue, LiveSessionError> {
84        let payload = self.invoke_handle(request)?;
85        self.refresh()?;
86        Ok(payload)
87    }
88
89    fn restart(&mut self, source: LiveSource) -> Result<JsonValue, LiveSessionError> {
90        let next_generation = self.generation().saturating_add(1);
91        let (new_handle, new_info) = start_backend(&self.session_id, &source)?;
92        if let Some(old_handle) = self.handle {
93            if let Err(error) = invoke_legacy(json!({"op": "dispose", "handle": old_handle})) {
94                let _ = invoke_legacy(json!({"op": "dispose", "handle": new_handle}));
95                return Err(error);
96            }
97        }
98        self.source = source;
99        self.handle = Some(new_handle);
100        self.generation_base = next_generation;
101        self.backend_generation = 0;
102        self.sequence = 0;
103        self.status = LiveSessionStatus::Ready;
104        self.pending_source = None;
105        self.sync_info(&new_info)?;
106        Ok(new_info)
107    }
108
109    fn reset(&mut self) -> Result<JsonValue, LiveSessionError> {
110        if let Some(source) = self.pending_source.take() {
111            return match self.restart(source.clone()) {
112                Ok(payload) => Ok(payload),
113                Err(error) => {
114                    self.pending_source = Some(source);
115                    Err(error)
116                }
117            };
118        }
119        self.invoke_and_refresh(json!({"op": "reset"}))
120    }
121
122    fn dispose(&mut self) -> Result<JsonValue, LiveSessionError> {
123        let Some(handle) = self.handle else {
124            self.status = LiveSessionStatus::Disposed;
125            return Ok(JsonValue::Bool(false));
126        };
127        let payload = invoke_legacy(json!({"op": "dispose", "handle": handle}))?;
128        self.handle = None;
129        self.status = LiveSessionStatus::Disposed;
130        self.pending_source = None;
131        Ok(payload)
132    }
133}
134
135impl LiveSession for InterpreterLiveSession {
136    fn state(&self) -> LiveSessionState {
137        LiveSessionState {
138            session_id: self.session_id.clone(),
139            source_id: self.source.source_id().to_owned(),
140            generation: self.generation(),
141            revision: self.source.revision().to_owned(),
142            sequence: self.sequence,
143            backend: LiveBackend::Interpreter,
144            status: self.status,
145        }
146    }
147
148    fn capabilities(&self) -> LiveSessionCapabilities {
149        LiveSessionCapabilities {
150            backend: LiveBackend::Interpreter,
151            operations: vec![
152                LiveSessionOperation::Snapshot,
153                LiveSessionOperation::Step,
154                LiveSessionOperation::Run,
155                LiveSessionOperation::Resume,
156                LiveSessionOperation::Resolve,
157                LiveSessionOperation::Reject,
158                LiveSessionOperation::Update,
159                LiveSessionOperation::Reset,
160                LiveSessionOperation::Cancel,
161                LiveSessionOperation::Dispose,
162            ],
163            replacement_policies: vec![
164                LiveReplacementPolicy::Restart,
165                LiveReplacementPolicy::ReplaceOnNextStart,
166            ],
167        }
168    }
169
170    fn dispatch_command(
171        &mut self,
172        command: LiveSessionCommand,
173    ) -> Result<JsonValue, LiveSessionError> {
174        match command {
175            LiveSessionCommand::Snapshot => self.invoke_and_refresh(json!({"op": "snapshot"})),
176            LiveSessionCommand::Step => self.invoke_and_refresh(json!({"op": "step"})),
177            LiveSessionCommand::Run { boundary_limit } => self.invoke_and_refresh(json!({
178                "op": "run",
179                "boundaryLimit": boundary_limit,
180            })),
181            LiveSessionCommand::Call { .. } => Err(LiveSessionError::new(
182                "live-session/unsupported-operation",
183                "interpreter backend does not support direct function calls",
184            )),
185            LiveSessionCommand::Pause => Err(LiveSessionError::new(
186                "live-session/unsupported-operation",
187                "interpreter backend does not support pause",
188            )),
189            LiveSessionCommand::Resume { settlement } => {
190                let mut request = json!({"op": "resume"});
191                if let Some(settlement) = settlement {
192                    request["settlement"] = settlement_json(settlement);
193                }
194                self.invoke_and_refresh(request)
195            }
196            LiveSessionCommand::Resolve { value } => self.invoke_and_refresh(json!({
197                "op": "resolve-suspension",
198                "value": value,
199            })),
200            LiveSessionCommand::Reject { error } => self.invoke_and_refresh(json!({
201                "op": "reject-suspension",
202                "error": error,
203            })),
204            LiveSessionCommand::Update { source, policy } => match policy {
205                LiveReplacementPolicy::Restart => self.restart(source),
206                LiveReplacementPolicy::ReplaceOnNextStart => {
207                    let revision = source.revision().to_owned();
208                    self.pending_source = Some(source);
209                    Ok(json!({
210                        "accepted": true,
211                        "activation": "next-start",
212                        "revision": revision,
213                    }))
214                }
215                LiveReplacementPolicy::PreserveRuntime => Err(LiveSessionError::new(
216                    "live-session/unsupported-replacement",
217                    "interpreter backend does not support preserve-runtime replacement",
218                )),
219            },
220            LiveSessionCommand::Reset => self.reset(),
221            LiveSessionCommand::Cancel => self.invoke_and_refresh(json!({"op": "cancel"})),
222            LiveSessionCommand::Dispose => self.dispose(),
223        }
224    }
225}
226
227impl Drop for InterpreterLiveSession {
228    fn drop(&mut self) {
229        if let Some(handle) = self.handle {
230            let _ = invoke_legacy(json!({"op": "dispose", "handle": handle}));
231            self.handle = None;
232        }
233    }
234}
235
236fn start_backend(
237    session_id: &str,
238    source: &LiveSource,
239) -> Result<(u64, JsonValue), LiveSessionError> {
240    let info = invoke_legacy(json!({
241        "op": "start",
242        "sessionId": session_id,
243        "sourceId": source.source_id(),
244        "source": source.source(),
245    }))?;
246    let handle = required_u64(&info, "handle")?;
247    Ok((handle, info))
248}
249
250fn settlement_json(settlement: LiveSettlement) -> JsonValue {
251    match settlement {
252        LiveSettlement::Fulfilled(value) => json!({
253            "status": "fulfilled",
254            "value": value,
255        }),
256        LiveSettlement::Rejected(error) => json!({
257            "status": "rejected",
258            "error": error,
259        }),
260    }
261}
262
263fn invoke_legacy(request: JsonValue) -> Result<JsonValue, LiveSessionError> {
264    let encoded = request.to_string();
265    let bytes = crate::interpreter_observation::invoke_json(&encoded);
266    let response: JsonValue = serde_json::from_slice(&bytes).map_err(|error| {
267        LiveSessionError::backend(format!(
268            "interpreter observation returned invalid JSON: {error}"
269        ))
270    })?;
271    if response.get("ok").and_then(JsonValue::as_bool) == Some(true) {
272        return Ok(response.get("value").cloned().unwrap_or(JsonValue::Null));
273    }
274    let error = response.get("error").and_then(JsonValue::as_object);
275    let code = error
276        .and_then(|value| value.get("code"))
277        .and_then(JsonValue::as_str)
278        .unwrap_or("interpreter-observation/error");
279    let message = error
280        .and_then(|value| value.get("message"))
281        .and_then(JsonValue::as_str)
282        .unwrap_or("interpreter observation request failed");
283    Err(LiveSessionError::new(
284        format!("live-session/backend/{code}"),
285        message,
286    ))
287}
288
289fn required_u64(value: &JsonValue, field: &str) -> Result<u64, LiveSessionError> {
290    value.get(field).and_then(JsonValue::as_u64).ok_or_else(|| {
291        LiveSessionError::backend(format!(
292            "interpreter observation response requires unsigned {field}"
293        ))
294    })
295}
296
297fn required_string<'a>(value: &'a JsonValue, field: &str) -> Result<&'a str, LiveSessionError> {
298    value.get(field).and_then(JsonValue::as_str).ok_or_else(|| {
299        LiveSessionError::backend(format!(
300            "interpreter observation response requires string {field}"
301        ))
302    })
303}
304
305fn parse_status(status: &str) -> Result<LiveSessionStatus, LiveSessionError> {
306    match status {
307        "ready" => Ok(LiveSessionStatus::Ready),
308        "running" => Ok(LiveSessionStatus::Running),
309        "suspended" => Ok(LiveSessionStatus::Suspended),
310        "returned" => Ok(LiveSessionStatus::Returned),
311        "failed" => Ok(LiveSessionStatus::Failed),
312        "cancelled" => Ok(LiveSessionStatus::Cancelled),
313        "disposed" => Ok(LiveSessionStatus::Disposed),
314        other => Err(LiveSessionError::backend(format!(
315            "unknown interpreter live-session status: {other}"
316        ))),
317    }
318}
319
320#[cfg(test)]
321mod tests {
322    use super::{invoke_legacy, InterpreterLiveSession};
323    use crate::live_session::LiveSource;
324    use serde_json::json;
325
326    #[test]
327    fn dropping_an_interpreter_adapter_releases_its_backend_handle() {
328        let handle = {
329            let session = InterpreterLiveSession::start(
330                "fixture/live-interpreter-drop",
331                LiveSource::new("drop.hal", "sha256:drop", "(+ 1 2)").unwrap(),
332            )
333            .unwrap();
334            session.handle.expect("started session must own a handle")
335        };
336
337        let error = invoke_legacy(json!({"op": "info", "handle": handle})).unwrap_err();
338        assert_eq!(
339            error.code(),
340            "live-session/backend/interpreter-observation/no-session"
341        );
342    }
343}