meerkat_mobkit/unified_runtime/
module_ops.rs1use std::time::Duration;
4
5use crate::runtime::{
6 DeliveryHistoryRequest, DeliveryHistoryResponse, DeliveryRecord, DeliverySendError,
7 DeliverySendRequest, GatingAuditEntry, GatingDecideError, GatingDecideRequest,
8 GatingDecisionResult, GatingEvaluateRequest, GatingEvaluateResult, GatingPendingEntry,
9 LifecycleEvent, MemoryIndexError, MemoryIndexRequest, MemoryIndexResult, MemoryQueryRequest,
10 MemoryQueryResult, MemoryStoreInfo, ModuleHealthTransition, RoutingResolution,
11 RoutingResolveError, RoutingResolveRequest, RuntimeMutationError, RuntimeRoute,
12 RuntimeRouteMutationError, ScheduleDefinition, ScheduleEvaluation, ScheduleValidationError,
13 SubscribeRequest, SubscribeResponse,
14};
15use crate::types::{EventEnvelope, UnifiedEvent};
16use crate::{ModuleRouteError, ModuleRouteRequest, ModuleRouteResponse, route_module_call};
17
18use super::UnifiedRuntime;
19use super::types::UnifiedRuntimeError;
20
21fn run_blocking<F, R>(f: F) -> R
28where
29 F: FnOnce() -> R + Send,
30 R: Send,
31{
32 std::thread::scope(|scope| {
33 scope
34 .spawn(f)
35 .join()
36 .unwrap_or_else(|e| std::panic::resume_unwind(e))
37 })
38}
39
40impl UnifiedRuntime {
41 pub async fn module_is_running(&self) -> bool {
42 self.module_runtime.lock().await.is_running()
43 }
44
45 pub async fn loaded_modules(&self) -> Vec<String> {
46 self.module_runtime.lock().await.loaded_modules()
47 }
48
49 pub async fn reconcile_modules(
51 &self,
52 modules: Vec<String>,
53 timeout: Duration,
54 ) -> Result<usize, RuntimeMutationError> {
55 let mut rt = self.module_runtime.lock().await;
56 run_blocking(|| rt.reconcile_modules(modules, timeout))
57 }
58
59 pub async fn resolve_routing(
61 &self,
62 request: RoutingResolveRequest,
63 ) -> Result<RoutingResolution, RoutingResolveError> {
64 let mut rt = self.module_runtime.lock().await;
65 run_blocking(|| rt.resolve_routing(request))
66 }
67
68 pub async fn send_delivery(
70 &self,
71 request: DeliverySendRequest,
72 ) -> Result<DeliveryRecord, DeliverySendError> {
73 let mut rt = self.module_runtime.lock().await;
74 run_blocking(|| rt.send_delivery(request))
75 }
76
77 pub async fn evaluate_schedule_tick(
78 &self,
79 schedules: &[ScheduleDefinition],
80 tick_ms: u64,
81 ) -> Result<ScheduleEvaluation, ScheduleValidationError> {
82 self.module_runtime
83 .lock()
84 .await
85 .evaluate_schedule_tick(schedules, tick_ms)
86 }
87
88 pub async fn list_runtime_routes(&self) -> Vec<RuntimeRoute> {
89 self.module_runtime.lock().await.list_runtime_routes()
90 }
91
92 pub async fn add_runtime_route(
93 &self,
94 route: RuntimeRoute,
95 ) -> Result<RuntimeRoute, RuntimeRouteMutationError> {
96 self.module_runtime.lock().await.add_runtime_route(route)
97 }
98
99 pub async fn delete_runtime_route(
100 &self,
101 route_key: &str,
102 ) -> Result<RuntimeRoute, RuntimeRouteMutationError> {
103 self.module_runtime
104 .lock()
105 .await
106 .delete_runtime_route(route_key)
107 }
108
109 pub async fn delivery_history(
110 &self,
111 request: DeliveryHistoryRequest,
112 ) -> DeliveryHistoryResponse {
113 self.module_runtime.lock().await.delivery_history(request)
114 }
115
116 pub async fn memory_stores(&self) -> Vec<MemoryStoreInfo> {
117 self.module_runtime.lock().await.memory_stores()
118 }
119
120 pub async fn memory_index(
129 &self,
130 request: MemoryIndexRequest,
131 ) -> Result<MemoryIndexResult, MemoryIndexError> {
132 let mut rt = self.module_runtime.lock().await;
133 run_blocking(|| rt.memory_index(request))
134 }
135
136 pub async fn memory_query(&self, request: MemoryQueryRequest) -> MemoryQueryResult {
137 self.module_runtime.lock().await.memory_query(request)
139 }
140
141 pub async fn evaluate_gating_action(
149 &self,
150 request: GatingEvaluateRequest,
151 ) -> GatingEvaluateResult {
152 let mut rt = self.module_runtime.lock().await;
153 run_blocking(|| rt.evaluate_gating_action(request))
154 }
155
156 pub async fn list_gating_pending(&self) -> Vec<GatingPendingEntry> {
157 self.module_runtime.lock().await.list_gating_pending()
158 }
159
160 pub async fn decide_gating_action(
161 &self,
162 request: GatingDecideRequest,
163 ) -> Result<GatingDecisionResult, GatingDecideError> {
164 self.module_runtime
165 .lock()
166 .await
167 .decide_gating_action(request)
168 }
169
170 pub async fn gating_audit_entries(&self, limit: usize) -> Vec<GatingAuditEntry> {
171 self.module_runtime.lock().await.gating_audit_entries(limit)
172 }
173
174 pub async fn spawn_member(
176 &self,
177 module_id: &str,
178 timeout: Duration,
179 ) -> Result<(), RuntimeMutationError> {
180 let mut rt = self.module_runtime.lock().await;
181 run_blocking(|| rt.spawn_member(module_id, timeout))
182 }
183
184 pub async fn route_module_call(
186 &self,
187 request: &ModuleRouteRequest,
188 timeout: Duration,
189 ) -> Result<ModuleRouteResponse, ModuleRouteError> {
190 let rt = self.module_runtime.lock().await;
191 run_blocking(|| route_module_call(&rt, request, timeout))
192 }
193
194 pub async fn module_lifecycle_events(&self) -> Vec<LifecycleEvent> {
195 self.module_runtime.lock().await.lifecycle_events.clone()
196 }
197
198 pub async fn module_health_transitions(&self) -> Vec<ModuleHealthTransition> {
199 self.module_runtime
200 .lock()
201 .await
202 .supervisor_report
203 .transitions
204 .clone()
205 }
206
207 pub async fn module_events(&self) -> Vec<EventEnvelope<UnifiedEvent>> {
208 self.module_runtime.lock().await.merged_events().to_vec()
209 }
210
211 pub async fn subscribe_events(
212 &self,
213 request: SubscribeRequest,
214 ) -> Result<SubscribeResponse, UnifiedRuntimeError> {
215 self.drain_mob_agent_events().await?;
216 self.module_runtime
217 .lock()
218 .await
219 .subscribe_events(request)
220 .map_err(UnifiedRuntimeError::Subscribe)
221 }
222}