1use std::collections::VecDeque;
18use std::panic::{catch_unwind, AssertUnwindSafe};
19use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
20use std::sync::{Arc, Mutex};
21
22use crate::events::{emit, ControlEvent, EventSink};
23
24#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
27pub enum EffectAdmissionStage {
28 Validate,
29 Transform,
30 Normalize,
31 Policy,
32 Approval,
33 Resource,
34 Execute,
35 Verify,
36 Settle,
37}
38
39impl EffectAdmissionStage {
40 pub fn label(self) -> &'static str {
41 match self {
42 Self::Validate => "validate",
43 Self::Transform => "transform",
44 Self::Normalize => "normalize",
45 Self::Policy => "policy",
46 Self::Approval => "approval",
47 Self::Resource => "resource",
48 Self::Execute => "execute",
49 Self::Verify => "verify",
50 Self::Settle => "settle",
51 }
52 }
53}
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
57pub enum EffectKind {
58 Navigate,
60 Input,
62 Action,
64 Script,
66 Capture,
68 Process,
70 Command,
72}
73
74impl EffectKind {
75 pub fn label(self) -> &'static str {
76 match self {
77 Self::Navigate => "navigate",
78 Self::Input => "input",
79 Self::Action => "action",
80 Self::Script => "script",
81 Self::Capture => "capture",
82 Self::Process => "process",
83 Self::Command => "command",
84 }
85 }
86}
87
88#[derive(Debug, Clone, PartialEq, Eq)]
94pub struct EffectRequest {
95 pub id: u64,
97 pub kind: EffectKind,
98 pub method: String,
100 pub target: String,
102 pub params: String,
104}
105
106impl EffectRequest {
107 pub fn new(kind: EffectKind, method: impl Into<String>, target: impl Into<String>) -> Self {
108 Self {
109 id: 0,
110 kind,
111 method: method.into(),
112 target: target.into(),
113 params: String::new(),
114 }
115 }
116 pub fn with_params(mut self, params: impl Into<String>) -> Self {
117 self.params = params.into();
118 self
119 }
120}
121
122#[derive(Debug, Clone, PartialEq, Eq)]
124pub enum AdmissionDecision {
125 Allow,
126 Deny { reason: String },
127}
128
129impl AdmissionDecision {
130 pub fn deny(reason: impl Into<String>) -> Self {
131 Self::Deny {
132 reason: reason.into(),
133 }
134 }
135}
136
137pub trait AdmissionHook: Send + Sync {
142 fn validate(&self, _request: &EffectRequest) -> Result<(), String> {
144 Ok(())
145 }
146 fn approve(&self, request: &EffectRequest) -> AdmissionDecision;
148}
149
150impl<F> AdmissionHook for F
151where
152 F: Fn(&EffectRequest) -> AdmissionDecision + Send + Sync,
153{
154 fn approve(&self, request: &EffectRequest) -> AdmissionDecision {
155 self(request)
156 }
157}
158
159#[derive(Debug, Clone, Copy, Default)]
162pub struct AllowAll;
163
164impl AdmissionHook for AllowAll {
165 fn approve(&self, _: &EffectRequest) -> AdmissionDecision {
166 AdmissionDecision::Allow
167 }
168}
169
170#[derive(Debug, Clone, Default)]
172pub struct CancelToken(Arc<AtomicBool>);
173
174impl CancelToken {
175 pub fn new() -> Self {
176 Self::default()
177 }
178 pub fn cancel(&self) {
179 self.0.store(true, Ordering::SeqCst);
180 }
181 pub fn is_cancelled(&self) -> bool {
182 self.0.load(Ordering::SeqCst)
183 }
184}
185
186#[derive(Debug, Clone, PartialEq, Eq)]
188pub enum ExecError {
189 Failed(String),
190 Cancelled,
192}
193
194impl From<String> for ExecError {
195 fn from(s: String) -> Self {
196 Self::Failed(s)
197 }
198}
199impl From<&str> for ExecError {
200 fn from(s: &str) -> Self {
201 Self::Failed(s.into())
202 }
203}
204
205#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
207pub enum SettleOutcome {
208 Ok,
209 Denied,
210 Failed,
211 Cancelled,
212}
213
214impl SettleOutcome {
215 pub fn label(self) -> &'static str {
216 match self {
217 Self::Ok => "ok",
218 Self::Denied => "denied",
219 Self::Failed => "failed",
220 Self::Cancelled => "cancelled",
221 }
222 }
223}
224
225#[derive(Debug, Clone, PartialEq, Eq)]
227pub struct Settlement {
228 pub id: u64,
229 pub kind: EffectKind,
230 pub method: String,
231 pub target: String,
232 pub outcome: SettleOutcome,
233 pub stages: Vec<EffectAdmissionStage>,
235 pub executed: bool,
237 pub crashed: bool,
239 pub reason: Option<String>,
241}
242
243#[derive(Debug)]
246pub struct Effect<T> {
247 pub settlement: Settlement,
248 pub value: Option<T>,
249}
250
251impl<T> Effect<T> {
252 pub fn is_ok(&self) -> bool {
253 self.settlement.outcome == SettleOutcome::Ok
254 }
255 pub fn into_result(self) -> Result<T, String> {
257 match self.value {
258 Some(v) if self.settlement.outcome == SettleOutcome::Ok => Ok(v),
259 _ => Err(format!(
260 "{}: {}",
261 self.settlement.outcome.label(),
262 self.settlement.reason.unwrap_or_default()
263 )),
264 }
265 }
266}
267
268pub const DEFAULT_JOURNAL_CAPACITY: usize = 1024;
270
271pub struct EffectGate {
274 hook: Arc<dyn AdmissionHook>,
275 events: Option<EventSink>,
276 next_id: AtomicU64,
277 journal: Mutex<VecDeque<Settlement>>,
278 capacity: usize,
279}
280
281impl std::fmt::Debug for EffectGate {
282 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
283 f.debug_struct("EffectGate")
284 .field("next_id", &self.next_id.load(Ordering::SeqCst))
285 .field("capacity", &self.capacity)
286 .finish_non_exhaustive()
287 }
288}
289
290impl Default for EffectGate {
291 fn default() -> Self {
292 Self::allow_all()
293 }
294}
295
296impl EffectGate {
297 pub fn new(hook: impl AdmissionHook + 'static) -> Self {
298 Self::from_arc(Arc::new(hook))
299 }
300 pub fn from_arc(hook: Arc<dyn AdmissionHook>) -> Self {
301 Self {
302 hook,
303 events: None,
304 next_id: AtomicU64::new(1),
305 journal: Mutex::new(VecDeque::new()),
306 capacity: DEFAULT_JOURNAL_CAPACITY,
307 }
308 }
309 pub fn allow_all() -> Self {
311 Self::new(AllowAll)
312 }
313 pub fn with_events(mut self, sink: EventSink) -> Self {
314 self.events = Some(sink);
315 self
316 }
317 pub fn with_journal_capacity(mut self, capacity: usize) -> Self {
319 self.capacity = capacity;
320 self
321 }
322
323 pub fn settlements(&self) -> Vec<Settlement> {
325 lock(&self.journal).iter().cloned().collect()
326 }
327
328 pub fn admit<T>(
335 &self,
336 mut request: EffectRequest,
337 cancel: &CancelToken,
338 execute: impl FnOnce(&CancelToken) -> Result<T, ExecError>,
339 ) -> Effect<T> {
340 request.id = self.next_id.fetch_add(1, Ordering::SeqCst);
341 let mut stages = vec![EffectAdmissionStage::Validate];
342
343 let validated = if request.method.trim().is_empty() {
345 Err("effect method is empty".to_string())
346 } else {
347 guarded(|| self.hook.validate(&request))
348 .unwrap_or_else(|| Err("admission hook panicked during validate".into()))
349 };
350 if let Err(reason) = validated {
351 return self.deny(request, stages, EffectAdmissionStage::Validate, reason);
352 }
353
354 stages.push(EffectAdmissionStage::Approval);
356 let decision = guarded(|| self.hook.approve(&request))
357 .unwrap_or_else(|| AdmissionDecision::deny("admission hook panicked during approval"));
358 if let AdmissionDecision::Deny { reason } = decision {
359 return self.deny(request, stages, EffectAdmissionStage::Approval, reason);
360 }
361
362 if cancel.is_cancelled() {
364 return self.settle(
365 request,
366 stages,
367 SettleOutcome::Cancelled,
368 false,
369 false,
370 Some("cancelled before execute".into()),
371 None,
372 );
373 }
374
375 stages.push(EffectAdmissionStage::Execute);
377 emit(
378 &self.events,
379 ControlEvent::EffectStarted {
380 id: request.id,
381 kind: request.kind,
382 method: request.method.clone(),
383 },
384 );
385 match catch_unwind(AssertUnwindSafe(|| execute(cancel))) {
386 Ok(Ok(v)) => self.settle(
387 request,
388 stages,
389 SettleOutcome::Ok,
390 true,
391 false,
392 None,
393 Some(v),
394 ),
395 Ok(Err(ExecError::Cancelled)) => self.settle(
396 request,
397 stages,
398 SettleOutcome::Cancelled,
399 true,
400 false,
401 Some("cancelled during execute".into()),
402 None,
403 ),
404 Ok(Err(ExecError::Failed(reason))) => self.settle(
405 request,
406 stages,
407 SettleOutcome::Failed,
408 true,
409 false,
410 Some(reason),
411 None,
412 ),
413 Err(panic) => {
414 let reason = panic_text(&panic);
415 self.settle(
416 request,
417 stages,
418 SettleOutcome::Failed,
419 true,
420 true,
421 Some(reason),
422 None,
423 )
424 }
425 }
426 }
427
428 fn deny<T>(
429 &self,
430 request: EffectRequest,
431 stages: Vec<EffectAdmissionStage>,
432 stage: EffectAdmissionStage,
433 reason: String,
434 ) -> Effect<T> {
435 self.settle_with(
436 request,
437 stages,
438 SettleOutcome::Denied,
439 false,
440 false,
441 Some(reason),
442 None,
443 Some(stage),
444 )
445 }
446
447 #[allow(clippy::too_many_arguments)]
448 fn settle<T>(
449 &self,
450 request: EffectRequest,
451 stages: Vec<EffectAdmissionStage>,
452 outcome: SettleOutcome,
453 executed: bool,
454 crashed: bool,
455 reason: Option<String>,
456 value: Option<T>,
457 ) -> Effect<T> {
458 self.settle_with(
459 request, stages, outcome, executed, crashed, reason, value, None,
460 )
461 }
462
463 #[allow(clippy::too_many_arguments)]
464 fn settle_with<T>(
465 &self,
466 request: EffectRequest,
467 mut stages: Vec<EffectAdmissionStage>,
468 outcome: SettleOutcome,
469 executed: bool,
470 crashed: bool,
471 reason: Option<String>,
472 value: Option<T>,
473 denied_at: Option<EffectAdmissionStage>,
474 ) -> Effect<T> {
475 stages.push(EffectAdmissionStage::Settle);
476 let settlement = Settlement {
477 id: request.id,
478 kind: request.kind,
479 method: request.method,
480 target: request.target,
481 outcome,
482 stages,
483 executed,
484 crashed,
485 reason,
486 };
487 if self.capacity > 0 {
488 let mut j = lock(&self.journal);
489 while j.len() >= self.capacity {
490 j.pop_front();
491 }
492 j.push_back(settlement.clone());
493 }
494 let (id, kind, method) = (settlement.id, settlement.kind, settlement.method.clone());
495 let event = match (denied_at, crashed) {
496 (Some(stage), _) => ControlEvent::EffectDenied {
497 id,
498 kind,
499 method,
500 stage,
501 reason: settlement.reason.clone().unwrap_or_default(),
502 },
503 (None, true) => ControlEvent::EffectCrashed { id, kind, method },
504 (None, false) => ControlEvent::EffectStopped {
505 id,
506 kind,
507 method,
508 outcome,
509 executed,
510 },
511 };
512 emit(&self.events, event);
513 Effect { settlement, value }
514 }
515}
516
517fn guarded<R>(f: impl FnOnce() -> R) -> Option<R> {
518 catch_unwind(AssertUnwindSafe(f)).ok()
519}
520
521fn panic_text(p: &Box<dyn std::any::Any + Send>) -> String {
522 let msg = p
523 .downcast_ref::<&str>()
524 .map(|s| s.to_string())
525 .or_else(|| p.downcast_ref::<String>().cloned())
526 .unwrap_or_else(|| "non-string panic".into());
527 format!("executor panicked: {msg}")
528}
529
530pub(crate) fn lock<T>(m: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
531 m.lock().unwrap_or_else(|e| e.into_inner())
532}