1use meerkat_core::lifecycle::{InputId, RunId};
6use serde::{Deserialize, Serialize};
7
8use crate::accept::AcceptOutcome;
9use crate::identifiers::LogicalRuntimeId;
10use crate::input::Input;
11use crate::input_state::{InputLifecycleState, InputState, StoredInputState};
12use crate::runtime_event::RuntimeEventEnvelope;
13use crate::runtime_state::RuntimeState;
14
15#[derive(Debug, Clone, thiserror::Error)]
17#[non_exhaustive]
18pub enum RuntimeDriverError {
19 #[error("Runtime not ready: {state}")]
21 NotReady { state: RuntimeState },
22
23 #[error("Runtime not found: {runtime_id}")]
30 NotFound { runtime_id: LogicalRuntimeId },
31
32 #[error("Input validation failed: {reason}")]
34 ValidationFailed { reason: String },
35
36 #[error("Runtime destroyed")]
38 Destroyed,
39
40 #[error("Recovery corruption: {reason}")]
42 RecoveryCorruption { reason: String },
43
44 #[error("Internal error: {0}")]
46 Internal(String),
47}
48
49#[derive(Debug, Clone, thiserror::Error)]
51#[non_exhaustive]
52pub enum RuntimeControlPlaneError {
53 #[error("Runtime not found: {0}")]
55 NotFound(LogicalRuntimeId),
56
57 #[error("Invalid state for operation: {state}")]
59 InvalidState { state: RuntimeState },
60
61 #[error("Store error: {0}")]
63 StoreError(String),
64
65 #[error("Internal error: {0}")]
67 Internal(String),
68}
69
70#[derive(Debug, Clone, Serialize, Deserialize)]
72pub struct RecoveryReport {
73 pub inputs_recovered: usize,
75 pub inputs_abandoned: usize,
77 pub inputs_requeued: usize,
79 #[serde(default, skip_serializing_if = "Vec::is_empty")]
81 pub details: Vec<String>,
82}
83
84#[derive(Debug, Clone, Serialize, Deserialize)]
86pub struct RetireReport {
87 pub inputs_abandoned: usize,
89 #[serde(default)]
91 pub inputs_pending_drain: usize,
92}
93
94#[derive(Debug, Clone, Serialize, Deserialize)]
96pub struct ResetReport {
97 pub inputs_abandoned: usize,
99}
100
101#[derive(Debug, Clone, Serialize, Deserialize)]
103pub struct RecycleReport {
104 pub inputs_transferred: usize,
106}
107
108#[derive(Debug, Clone, Serialize, Deserialize)]
110pub struct DestroyReport {
111 pub inputs_abandoned: usize,
113}
114
115#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
120#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
121pub trait RuntimeDriver: Send + Sync {
122 async fn accept_input(&mut self, input: Input) -> Result<AcceptOutcome, RuntimeDriverError>;
124
125 async fn on_runtime_event(
127 &mut self,
128 event: RuntimeEventEnvelope,
129 ) -> Result<(), RuntimeDriverError>;
130
131 async fn recover(&mut self) -> Result<RecoveryReport, RuntimeDriverError>;
133
134 fn runtime_state(&self) -> RuntimeState;
136
137 fn input_state(&self, input_id: &InputId) -> Option<&InputState>;
139
140 fn input_phase(&self, input_id: &InputId) -> Option<InputLifecycleState>;
142
143 fn input_last_run_id(&self, input_id: &InputId) -> Option<RunId>;
145
146 fn input_last_boundary_sequence(&self, input_id: &InputId) -> Option<u64>;
148
149 fn stored_input_state(&self, input_id: &InputId) -> Option<StoredInputState>;
151
152 fn stored_input_states_snapshot(&self) -> Result<Vec<StoredInputState>, RuntimeDriverError>;
158
159 fn input_id_for_idempotency_key(&self, idempotency_key: &str) -> Option<InputId>;
165
166 fn active_input_ids(&self) -> Vec<InputId>;
168}
169
170#[cfg_attr(not(target_arch = "wasm32"), async_trait::async_trait)]
172#[cfg_attr(target_arch = "wasm32", async_trait::async_trait(?Send))]
173pub trait RuntimeControlPlane: Send + Sync {
174 async fn ingest(
176 &self,
177 runtime_id: &LogicalRuntimeId,
178 input: Input,
179 ) -> Result<AcceptOutcome, RuntimeControlPlaneError>;
180
181 async fn publish_event(
183 &self,
184 event: RuntimeEventEnvelope,
185 ) -> Result<(), RuntimeControlPlaneError>;
186
187 async fn retire(
189 &self,
190 runtime_id: &LogicalRuntimeId,
191 ) -> Result<RetireReport, RuntimeControlPlaneError>;
192
193 async fn recycle(
195 &self,
196 runtime_id: &LogicalRuntimeId,
197 ) -> Result<RecycleReport, RuntimeControlPlaneError>;
198
199 async fn reset(
201 &self,
202 runtime_id: &LogicalRuntimeId,
203 ) -> Result<ResetReport, RuntimeControlPlaneError>;
204
205 async fn recover(
207 &self,
208 runtime_id: &LogicalRuntimeId,
209 ) -> Result<RecoveryReport, RuntimeControlPlaneError>;
210
211 async fn runtime_state(
213 &self,
214 runtime_id: &LogicalRuntimeId,
215 ) -> Result<RuntimeState, RuntimeControlPlaneError>;
216
217 async fn destroy(
219 &self,
220 runtime_id: &LogicalRuntimeId,
221 ) -> Result<DestroyReport, RuntimeControlPlaneError>;
222
223 async fn load_boundary_receipt(
225 &self,
226 runtime_id: &LogicalRuntimeId,
227 run_id: &RunId,
228 sequence: u64,
229 ) -> Result<Option<meerkat_core::lifecycle::RunBoundaryReceipt>, RuntimeControlPlaneError>;
230}
231
232#[cfg(test)]
233#[allow(clippy::unwrap_used)]
234mod tests {
235 use super::*;
236
237 fn _assert_driver_object_safe(_: &dyn RuntimeDriver) {}
239 fn _assert_control_plane_object_safe(_: &dyn RuntimeControlPlane) {}
240
241 #[test]
242 fn runtime_driver_error_display() {
243 let err = RuntimeDriverError::NotReady {
244 state: RuntimeState::Initializing,
245 };
246 assert!(err.to_string().contains("initializing"));
247
248 let err = RuntimeDriverError::ValidationFailed {
249 reason: "bad input".into(),
250 };
251 assert!(err.to_string().contains("bad input"));
252 }
253
254 #[test]
255 fn runtime_control_plane_error_display() {
256 let err = RuntimeControlPlaneError::NotFound(LogicalRuntimeId::new("missing"));
257 assert!(err.to_string().contains("missing"));
258 }
259
260 #[test]
261 fn recovery_report_serde() {
262 let report = RecoveryReport {
263 inputs_recovered: 5,
264 inputs_abandoned: 1,
265 inputs_requeued: 3,
266 details: vec!["requeued 3 staged inputs".into()],
267 };
268 let json = serde_json::to_value(&report).unwrap();
269 let parsed: RecoveryReport = serde_json::from_value(json).unwrap();
270 assert_eq!(parsed.inputs_recovered, 5);
271 }
272}