hara_native/live_session/
interpreter.rs1use 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}