Skip to main content

lash_core/
triggers.rs

1use std::collections::BTreeMap;
2use std::sync::{Arc, Mutex};
3
4use serde::{Deserialize, Serialize};
5
6use crate::plugin::PluginError;
7
8#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
9pub struct TriggerEvent {
10    pub resource_type: String,
11    pub alias: String,
12    pub event: String,
13    pub payload_schema: crate::LashSchema,
14}
15
16impl TriggerEvent {
17    pub fn new(
18        resource_type: impl Into<String>,
19        alias: impl Into<String>,
20        event: impl Into<String>,
21        payload_schema: crate::LashSchema,
22    ) -> Self {
23        Self {
24            resource_type: resource_type.into(),
25            alias: alias.into(),
26            event: event.into(),
27            payload_schema,
28        }
29    }
30
31    pub fn payload_schema(&self) -> &crate::LashSchema {
32        &self.payload_schema
33    }
34
35    pub fn key(&self) -> TriggerEventKey {
36        TriggerEventKey {
37            resource_type: self.resource_type.clone(),
38            alias: self.alias.clone(),
39            event: self.event.clone(),
40        }
41    }
42
43    pub fn source_type(&self) -> String {
44        trigger_event_type(&self.alias, &self.event)
45    }
46}
47
48#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
49pub struct TriggerEventKey {
50    pub resource_type: String,
51    pub alias: String,
52    pub event: String,
53}
54
55impl TriggerEventKey {
56    pub fn new(
57        resource_type: impl Into<String>,
58        alias: impl Into<String>,
59        event: impl Into<String>,
60    ) -> Self {
61        Self {
62            resource_type: resource_type.into(),
63            alias: alias.into(),
64            event: event.into(),
65        }
66    }
67
68    pub fn source_type(&self) -> String {
69        trigger_event_type(&self.alias, &self.event)
70    }
71}
72
73pub fn trigger_event_type(alias: &str, event: &str) -> String {
74    format!("{alias}.{event}")
75}
76
77#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
78pub struct TriggerEventCatalog {
79    events: BTreeMap<TriggerEventKey, TriggerEvent>,
80}
81
82impl TriggerEventCatalog {
83    pub fn new() -> Self {
84        Self::default()
85    }
86
87    pub fn declare(&mut self, event: TriggerEvent) -> Result<(), String> {
88        let key = event.key();
89        if self.events.contains_key(&key) {
90            return Err(format!(
91                "duplicate trigger occurrence `{}.{}.{}`",
92                key.resource_type, key.alias, key.event
93            ));
94        }
95        let source_type = event.source_type();
96        if let Some(existing) = self
97            .events
98            .values()
99            .find(|existing| existing.source_type() == source_type)
100        {
101            return Err(format!(
102                "duplicate trigger source `{source_type}` declared by `{}.{}.{}` and `{}.{}.{}`",
103                existing.resource_type,
104                existing.alias,
105                existing.event,
106                key.resource_type,
107                key.alias,
108                key.event
109            ));
110        }
111        self.events.insert(key, event);
112        Ok(())
113    }
114
115    pub fn from_events(events: impl IntoIterator<Item = TriggerEvent>) -> Result<Self, String> {
116        let mut catalog = Self::new();
117        for event in events {
118            catalog.declare(event)?;
119        }
120        Ok(catalog)
121    }
122
123    pub fn get(&self, resource_type: &str, alias: &str, event: &str) -> Option<&TriggerEvent> {
124        self.events
125            .get(&TriggerEventKey::new(resource_type, alias, event))
126    }
127
128    pub fn is_empty(&self) -> bool {
129        self.events.is_empty()
130    }
131
132    pub fn events(&self) -> impl Iterator<Item = &TriggerEvent> {
133        self.events.values()
134    }
135}
136
137#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
138#[serde(rename_all = "snake_case")]
139pub enum TriggerDeliveryEmitOutcome {
140    Started,
141    AlreadyReserved,
142    Failed { reason: String },
143}
144
145#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
146pub struct TriggerDeliveryEmitReport {
147    pub occurrence_id: String,
148    pub subscription_id: String,
149    pub process_id: String,
150    pub outcome: TriggerDeliveryEmitOutcome,
151}
152
153#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
154pub struct TriggerEmitReport {
155    #[serde(default, skip_serializing_if = "String::is_empty")]
156    pub occurrence_id: String,
157    #[serde(default, skip_serializing_if = "Vec::is_empty")]
158    pub deliveries: Vec<TriggerDeliveryEmitReport>,
159}
160
161impl TriggerEmitReport {
162    pub fn empty() -> Self {
163        Self::default()
164    }
165
166    fn new(occurrence_id: String, deliveries: Vec<TriggerDeliveryEmitReport>) -> Self {
167        Self {
168            occurrence_id,
169            deliveries,
170        }
171    }
172
173    pub fn started_process_ids(&self) -> Vec<String> {
174        self.deliveries
175            .iter()
176            .filter(|delivery| delivery.outcome == TriggerDeliveryEmitOutcome::Started)
177            .map(|delivery| delivery.process_id.clone())
178            .collect()
179    }
180}
181
182#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
183pub struct TriggerOccurrenceRequest {
184    pub source_type: String,
185    pub source_key: String,
186    #[serde(default)]
187    pub payload: serde_json::Value,
188    pub idempotency_key: String,
189    #[serde(default, skip_serializing_if = "Option::is_none")]
190    pub source: Option<serde_json::Value>,
191    /// Optional host routing scope. When present, only subscriptions
192    /// registered by this session can reserve deliveries for the occurrence.
193    #[serde(default, skip_serializing_if = "Option::is_none")]
194    pub session_id: Option<String>,
195}
196
197impl TriggerOccurrenceRequest {
198    pub fn new(
199        source_type: impl Into<String>,
200        source_key: impl Into<String>,
201        payload: serde_json::Value,
202        idempotency_key: impl Into<String>,
203    ) -> Self {
204        Self {
205            source_type: source_type.into(),
206            source_key: source_key.into(),
207            payload,
208            idempotency_key: idempotency_key.into(),
209            source: None,
210            session_id: None,
211        }
212    }
213
214    pub fn with_source(mut self, source: serde_json::Value) -> Self {
215        self.source = Some(source);
216        self
217    }
218
219    pub fn for_session(mut self, session_id: impl Into<String>) -> Self {
220        self.session_id = Some(session_id.into());
221        self
222    }
223}
224
225#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
226pub struct TriggerOccurrenceRecord {
227    pub occurrence_id: String,
228    pub source_type: String,
229    pub source_key: String,
230    #[serde(default)]
231    pub payload: serde_json::Value,
232    pub idempotency_key: String,
233    #[serde(default, skip_serializing_if = "Option::is_none")]
234    pub source: Option<serde_json::Value>,
235    #[serde(default, skip_serializing_if = "Option::is_none")]
236    pub session_id: Option<String>,
237    pub occurred_at_ms: u64,
238}
239
240#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
241pub struct TriggerOccurrenceFilter {
242    #[serde(default, skip_serializing_if = "Option::is_none")]
243    pub source_type: Option<String>,
244    #[serde(default, skip_serializing_if = "Option::is_none")]
245    pub source_key: Option<String>,
246    #[serde(default, skip_serializing_if = "Option::is_none")]
247    pub occurred_at_start_ms: Option<u64>,
248    #[serde(default, skip_serializing_if = "Option::is_none")]
249    pub occurred_at_end_ms: Option<u64>,
250}
251
252impl TriggerOccurrenceFilter {
253    pub fn for_source(source_type: impl Into<String>, source_key: impl Into<String>) -> Self {
254        Self {
255            source_type: Some(source_type.into()),
256            source_key: Some(source_key.into()),
257            ..Self::default()
258        }
259    }
260
261    pub fn matches(&self, record: &TriggerOccurrenceRecord) -> bool {
262        self.source_type
263            .as_deref()
264            .is_none_or(|source_type| record.source_type == source_type)
265            && self
266                .source_key
267                .as_deref()
268                .is_none_or(|source_key| record.source_key == source_key)
269            && self
270                .occurred_at_start_ms
271                .is_none_or(|start_ms| record.occurred_at_ms >= start_ms)
272            && self
273                .occurred_at_end_ms
274                .is_none_or(|end_ms| record.occurred_at_ms < end_ms)
275    }
276}
277
278#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
279#[serde(transparent)]
280pub struct TriggerEventType(String);
281
282impl TriggerEventType {
283    pub fn new(value: impl Into<String>) -> Self {
284        Self(value.into())
285    }
286
287    pub fn as_str(&self) -> &str {
288        &self.0
289    }
290}
291
292impl From<String> for TriggerEventType {
293    fn from(value: String) -> Self {
294        Self::new(value)
295    }
296}
297
298impl From<&str> for TriggerEventType {
299    fn from(value: &str) -> Self {
300        Self::new(value)
301    }
302}
303
304impl AsRef<str> for TriggerEventType {
305    fn as_ref(&self) -> &str {
306        self.as_str()
307    }
308}
309
310impl std::fmt::Display for TriggerEventType {
311    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
312        formatter.write_str(self.as_str())
313    }
314}
315
316#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
317pub struct TriggerRegistration {
318    pub handle: String,
319    pub source_key: String,
320    #[serde(default, skip_serializing_if = "Option::is_none")]
321    pub name: Option<String>,
322    pub source_type: TriggerEventType,
323    pub source: serde_json::Value,
324    pub target: TriggerTargetSummary,
325    #[serde(default = "default_enabled")]
326    pub enabled: bool,
327}
328
329#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
330pub struct TriggerTargetSummary {
331    pub label: Option<String>,
332    pub identity: crate::ProcessIdentity,
333    pub input: crate::ProcessInput,
334    #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
335    pub inputs: BTreeMap<String, TriggerInputBinding>,
336}
337
338#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
339#[serde(tag = "type", rename_all = "snake_case")]
340pub enum TriggerInputBinding {
341    Event,
342    Fixed { value: serde_json::Value },
343}
344
345#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
346pub struct TriggerSubscriptionDraft {
347    pub registrant: crate::ProcessOriginator,
348    pub env_ref: crate::ProcessExecutionEnvRef,
349    #[serde(default, skip_serializing_if = "Option::is_none")]
350    pub wake_target: Option<crate::SessionScope>,
351    #[serde(default, skip_serializing_if = "Option::is_none")]
352    pub name: Option<String>,
353    pub source_type: String,
354    pub source_key: String,
355    pub source: serde_json::Value,
356    pub payload_schema: crate::LashSchema,
357    pub target: crate::ProcessInput,
358    pub target_identity: crate::ProcessIdentity,
359    #[serde(default)]
360    pub event_types: Vec<crate::ProcessEventType>,
361    #[serde(default)]
362    pub input_template: BTreeMap<String, TriggerInputBinding>,
363    #[serde(default, skip_serializing_if = "Option::is_none")]
364    pub target_label: Option<String>,
365}
366
367impl TriggerSubscriptionDraft {
368    pub fn for_process(
369        registrant: crate::ProcessOriginator,
370        env_ref: crate::ProcessExecutionEnvRef,
371        source_type: impl Into<String>,
372        source_key: impl Into<String>,
373        target: crate::ProcessInput,
374        target_identity: crate::ProcessIdentity,
375    ) -> Self {
376        let target_label = target_identity.label.clone();
377        Self {
378            registrant,
379            env_ref,
380            wake_target: None,
381            name: None,
382            source_type: source_type.into(),
383            source_key: source_key.into(),
384            source: serde_json::Value::Object(serde_json::Map::new()),
385            payload_schema: crate::LashSchema::new(serde_json::Value::Object(
386                serde_json::Map::new(),
387            )),
388            target,
389            target_identity,
390            event_types: Vec::new(),
391            input_template: BTreeMap::new(),
392            target_label,
393        }
394    }
395
396    pub fn with_name(mut self, name: impl Into<String>) -> Self {
397        self.name = Some(name.into());
398        self
399    }
400
401    pub fn with_source(mut self, source: serde_json::Value) -> Self {
402        self.source = source;
403        self
404    }
405
406    pub fn with_payload_schema(mut self, payload_schema: crate::LashSchema) -> Self {
407        self.payload_schema = payload_schema;
408        self
409    }
410
411    pub fn with_wake_target(mut self, wake_target: crate::SessionScope) -> Self {
412        self.wake_target = Some(wake_target);
413        self
414    }
415
416    pub fn with_event_types(
417        mut self,
418        event_types: impl IntoIterator<Item = crate::ProcessEventType>,
419    ) -> Self {
420        self.event_types = event_types.into_iter().collect();
421        self
422    }
423
424    pub fn with_input_template(
425        mut self,
426        input_template: BTreeMap<String, TriggerInputBinding>,
427    ) -> Self {
428        self.input_template = input_template;
429        self
430    }
431
432    pub fn with_target_label(mut self, target_label: impl Into<String>) -> Self {
433        self.target_label = Some(target_label.into());
434        self
435    }
436
437    pub fn validate(&self) -> Result<(), PluginError> {
438        validate_trigger_subscription_target_label(
439            self.target_label.as_deref(),
440            self.target_identity.label.as_deref(),
441        )
442    }
443}
444
445#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
446pub struct TriggerSubscriptionRecord {
447    pub subscription_id: String,
448    pub registrant: crate::ProcessOriginator,
449    pub env_ref: crate::ProcessExecutionEnvRef,
450    #[serde(default, skip_serializing_if = "Option::is_none")]
451    pub wake_target: Option<crate::SessionScope>,
452    pub handle: String,
453    #[serde(default, skip_serializing_if = "Option::is_none")]
454    pub name: Option<String>,
455    pub source_type: String,
456    pub source_key: String,
457    pub source: serde_json::Value,
458    pub payload_schema: crate::LashSchema,
459    pub target: crate::ProcessInput,
460    pub target_identity: crate::ProcessIdentity,
461    #[serde(default)]
462    pub event_types: Vec<crate::ProcessEventType>,
463    #[serde(default)]
464    pub input_template: BTreeMap<String, TriggerInputBinding>,
465    #[serde(default, skip_serializing_if = "Option::is_none")]
466    pub target_label: Option<String>,
467    #[serde(default = "default_enabled")]
468    pub enabled: bool,
469    pub created_at_ms: u64,
470    pub updated_at_ms: u64,
471}
472
473impl TriggerSubscriptionRecord {
474    pub fn registrant_scope_id(&self) -> String {
475        self.registrant.scope_id()
476    }
477
478    pub fn registrant_session_id(&self) -> Option<&str> {
479        match &self.registrant {
480            crate::ProcessOriginator::Session { scope } => Some(scope.session_id.as_str()),
481            crate::ProcessOriginator::Host { .. } => None,
482        }
483    }
484}
485
486fn validate_trigger_subscription_target_label(
487    target_label: Option<&str>,
488    identity_label: Option<&str>,
489) -> Result<(), PluginError> {
490    match (target_label, identity_label) {
491        (Some(target_label), Some(identity_label)) if target_label != identity_label => {
492            Err(PluginError::Session(
493                "trigger target_label must match target_identity.label when both are present"
494                    .to_string(),
495            ))
496        }
497        _ => Ok(()),
498    }
499}
500
501impl From<&TriggerSubscriptionRecord> for TriggerRegistration {
502    fn from(route: &TriggerSubscriptionRecord) -> Self {
503        Self {
504            handle: route.handle.clone(),
505            source_key: route.source_key.clone(),
506            name: route.name.clone(),
507            source_type: TriggerEventType::new(route.source_type.clone()),
508            source: route.source.clone(),
509            target: TriggerTargetSummary {
510                label: route.target_label.clone(),
511                identity: route.target_identity.clone(),
512                input: route.target.clone(),
513                inputs: route.input_template.clone(),
514            },
515            enabled: route.enabled,
516        }
517    }
518}
519
520#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
521pub struct TriggerSubscriptionFilter {
522    #[serde(default, skip_serializing_if = "Option::is_none")]
523    pub registrant_scope_id: Option<String>,
524    #[serde(default, skip_serializing_if = "Option::is_none")]
525    pub session_id: Option<String>,
526    #[serde(default, skip_serializing_if = "Option::is_none")]
527    pub handle: Option<String>,
528    #[serde(default, skip_serializing_if = "Option::is_none")]
529    pub name: Option<String>,
530    #[serde(default, skip_serializing_if = "Option::is_none")]
531    pub source_type: Option<String>,
532    #[serde(default, skip_serializing_if = "Option::is_none")]
533    pub source_key: Option<String>,
534    #[serde(default, skip_serializing_if = "Option::is_none")]
535    pub target: Option<serde_json::Value>,
536    #[serde(default, skip_serializing_if = "Option::is_none")]
537    pub enabled: Option<bool>,
538}
539
540impl TriggerSubscriptionFilter {
541    pub fn for_session(session_id: impl Into<String>) -> Self {
542        Self {
543            session_id: Some(session_id.into()),
544            ..Self::default()
545        }
546    }
547
548    pub fn for_registrant_scope(scope_id: impl Into<String>) -> Self {
549        Self {
550            registrant_scope_id: Some(scope_id.into()),
551            ..Self::default()
552        }
553    }
554
555    pub fn for_source_type(source_type: impl Into<String>) -> Self {
556        Self {
557            source_type: Some(source_type.into()),
558            ..Self::default()
559        }
560    }
561
562    pub fn effective_registrant_scope_id(&self) -> Option<String> {
563        self.registrant_scope_id.clone()
564    }
565
566    pub fn matches(&self, record: &TriggerSubscriptionRecord) -> bool {
567        self.effective_registrant_scope_id()
568            .is_none_or(|scope_id| record.registrant_scope_id() == scope_id)
569            && self
570                .session_id
571                .as_deref()
572                .is_none_or(|session_id| record.registrant_session_id() == Some(session_id))
573            && self
574                .handle
575                .as_deref()
576                .is_none_or(|handle| record.handle == handle)
577            && self
578                .name
579                .as_deref()
580                .is_none_or(|name| record.name.as_deref() == Some(name))
581            && self
582                .source_type
583                .as_deref()
584                .is_none_or(|source_type| record.source_type == source_type)
585            && self
586                .source_key
587                .as_deref()
588                .is_none_or(|source_key| record.source_key == source_key)
589            && self.enabled.is_none_or(|enabled| record.enabled == enabled)
590            && self
591                .target
592                .as_ref()
593                .is_none_or(|target| record.target_identity.definition.as_ref() == Some(target))
594    }
595}
596
597#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
598#[serde(rename_all = "snake_case")]
599pub enum TriggerDeliveryReservationStatus {
600    Reserved,
601    AlreadyReserved,
602}
603
604#[derive(Clone, Debug, Serialize, Deserialize)]
605pub struct TriggerDeliveryReservation {
606    pub occurrence: TriggerOccurrenceRecord,
607    pub subscription: TriggerSubscriptionRecord,
608    pub process_id: String,
609    pub created_at_ms: u64,
610    pub reservation_status: TriggerDeliveryReservationStatus,
611}
612
613impl TriggerDeliveryReservation {
614    fn emit_report(&self, outcome: TriggerDeliveryEmitOutcome) -> TriggerDeliveryEmitReport {
615        TriggerDeliveryEmitReport {
616            occurrence_id: self.occurrence.occurrence_id.clone(),
617            subscription_id: self.subscription.subscription_id.clone(),
618            process_id: self.process_id.clone(),
619            outcome,
620        }
621    }
622}
623
624#[async_trait::async_trait]
625pub trait TriggerStore: Send + Sync {
626    fn durability_tier(&self) -> crate::DurabilityTier {
627        crate::DurabilityTier::Inline
628    }
629
630    async fn source_key_for_subscription(
631        &self,
632        source_type: &str,
633        source: &serde_json::Value,
634    ) -> Result<String, PluginError> {
635        default_trigger_source_key(source_type, source)
636    }
637
638    async fn register_subscription(
639        &self,
640        draft: TriggerSubscriptionDraft,
641    ) -> Result<TriggerSubscriptionRecord, PluginError>;
642
643    async fn list_subscriptions(
644        &self,
645        filter: TriggerSubscriptionFilter,
646    ) -> Result<Vec<TriggerSubscriptionRecord>, PluginError>;
647
648    async fn cancel_subscription(
649        &self,
650        registrant_scope_id: &str,
651        handle: &str,
652    ) -> Result<bool, PluginError>;
653
654    /// Set whether one scoped subscription participates in future delivery
655    /// reservations. Returns whether the stored value changed.
656    async fn set_subscription_enabled(
657        &self,
658        registrant_scope_id: &str,
659        handle: &str,
660        enabled: bool,
661    ) -> Result<bool, PluginError> {
662        if enabled {
663            return Err(PluginError::Session(
664                "re-enabling trigger subscriptions is unsupported by this store".to_string(),
665            ));
666        }
667        self.cancel_subscription(registrant_scope_id, handle).await
668    }
669
670    /// Delete one scoped subscription and its delivery records. Historical
671    /// occurrences remain available because they describe source events.
672    async fn delete_subscription(
673        &self,
674        _registrant_scope_id: &str,
675        _handle: &str,
676    ) -> Result<bool, PluginError> {
677        Err(PluginError::Session(
678            "deleting individual trigger subscriptions is unsupported by this store".to_string(),
679        ))
680    }
681
682    async fn delete_session_subscriptions(&self, session_id: &str) -> Result<usize, PluginError>;
683
684    async fn record_occurrence(
685        &self,
686        request: TriggerOccurrenceRequest,
687    ) -> Result<TriggerOccurrenceRecord, PluginError>;
688
689    async fn list_occurrences(
690        &self,
691        filter: TriggerOccurrenceFilter,
692    ) -> Result<Vec<TriggerOccurrenceRecord>, PluginError>;
693
694    async fn reserve_matching_deliveries(
695        &self,
696        occurrence_id: &str,
697    ) -> Result<Vec<TriggerDeliveryReservation>, PluginError>;
698
699    async fn list_deliveries_by_occurrence_id(
700        &self,
701        occurrence_id: &str,
702    ) -> Result<Vec<TriggerDeliveryReservation>, PluginError>;
703
704    async fn list_deliveries_by_subscription_id(
705        &self,
706        subscription_id: &str,
707    ) -> Result<Vec<TriggerDeliveryReservation>, PluginError>;
708
709    async fn list_deliveries_by_process_id(
710        &self,
711        process_id: &str,
712    ) -> Result<Vec<TriggerDeliveryReservation>, PluginError>;
713}
714
715pub struct InMemoryTriggerStore {
716    clock: Arc<dyn crate::Clock>,
717    state: Mutex<InMemoryTriggerEventState>,
718}
719
720impl InMemoryTriggerStore {
721    pub fn new() -> Self {
722        Self::with_clock(Arc::new(crate::SystemClock))
723    }
724
725    pub fn with_clock(clock: Arc<dyn crate::Clock>) -> Self {
726        Self {
727            clock,
728            state: Mutex::new(InMemoryTriggerEventState::default()),
729        }
730    }
731
732    fn list_deliveries(
733        &self,
734        matches: impl Fn(&InMemoryTriggerDeliveryRecord) -> bool,
735    ) -> Result<Vec<TriggerDeliveryReservation>, PluginError> {
736        let state = self
737            .state
738            .lock()
739            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
740        let mut deliveries = state
741            .deliveries
742            .values()
743            .filter(|delivery| matches(delivery))
744            .map(|delivery| {
745                in_memory_delivery_reservation(
746                    &state,
747                    delivery,
748                    TriggerDeliveryReservationStatus::AlreadyReserved,
749                )
750            })
751            .collect::<Result<Vec<_>, _>>()?;
752        deliveries.sort_by(|left, right| {
753            left.created_at_ms
754                .cmp(&right.created_at_ms)
755                .then_with(|| {
756                    left.occurrence
757                        .occurrence_id
758                        .cmp(&right.occurrence.occurrence_id)
759                })
760                .then_with(|| {
761                    left.subscription
762                        .subscription_id
763                        .cmp(&right.subscription.subscription_id)
764                })
765        });
766        Ok(deliveries)
767    }
768
769    #[cfg(any(test, feature = "testing"))]
770    pub(crate) fn delete_deliveries_by_process_ids(
771        &self,
772        process_ids: &std::collections::HashSet<String>,
773    ) -> Result<usize, PluginError> {
774        if process_ids.is_empty() {
775            return Ok(0);
776        }
777        let mut state = self
778            .state
779            .lock()
780            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
781        let before = state.deliveries.len();
782        state
783            .deliveries
784            .retain(|_, delivery| !process_ids.contains(&delivery.process_id));
785        Ok(before.saturating_sub(state.deliveries.len()))
786    }
787}
788
789impl Default for InMemoryTriggerStore {
790    fn default() -> Self {
791        Self::new()
792    }
793}
794
795#[derive(Default)]
796struct InMemoryTriggerEventState {
797    next_subscription_seq: u64,
798    subscriptions: BTreeMap<String, TriggerSubscriptionRecord>,
799    occurrences: BTreeMap<String, TriggerOccurrenceRecord>,
800    occurrence_id_by_idempotency_key: BTreeMap<String, String>,
801    occurrence_hashes: BTreeMap<String, String>,
802    deliveries: BTreeMap<(String, String), InMemoryTriggerDeliveryRecord>,
803}
804
805#[derive(Clone)]
806struct InMemoryTriggerDeliveryRecord {
807    occurrence_id: String,
808    subscription_id: String,
809    process_id: String,
810    created_at_ms: u64,
811}
812
813fn in_memory_delivery_reservation(
814    state: &InMemoryTriggerEventState,
815    delivery: &InMemoryTriggerDeliveryRecord,
816    reservation_status: TriggerDeliveryReservationStatus,
817) -> Result<TriggerDeliveryReservation, PluginError> {
818    let occurrence = state
819        .occurrences
820        .get(&delivery.occurrence_id)
821        .cloned()
822        .ok_or_else(|| {
823            PluginError::Session(format!(
824                "missing trigger occurrence `{}` for delivery",
825                delivery.occurrence_id
826            ))
827        })?;
828    let subscription = state
829        .subscriptions
830        .get(&delivery.subscription_id)
831        .cloned()
832        .ok_or_else(|| {
833            PluginError::Session(format!(
834                "missing trigger subscription `{}` for delivery",
835                delivery.subscription_id
836            ))
837        })?;
838    Ok(TriggerDeliveryReservation {
839        occurrence,
840        subscription,
841        process_id: delivery.process_id.clone(),
842        created_at_ms: delivery.created_at_ms,
843        reservation_status,
844    })
845}
846
847#[async_trait::async_trait]
848impl TriggerStore for InMemoryTriggerStore {
849    async fn register_subscription(
850        &self,
851        draft: TriggerSubscriptionDraft,
852    ) -> Result<TriggerSubscriptionRecord, PluginError> {
853        draft.validate()?;
854        let mut state = self
855            .state
856            .lock()
857            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
858        state.next_subscription_seq = state.next_subscription_seq.saturating_add(1);
859        let handle = format!("trigger:{}", state.next_subscription_seq);
860        let subscription_id = format!("subscription:{}", state.next_subscription_seq);
861        let now = self.clock.timestamp_ms();
862        let record = TriggerSubscriptionRecord {
863            subscription_id: subscription_id.clone(),
864            registrant: draft.registrant,
865            env_ref: draft.env_ref,
866            wake_target: draft.wake_target,
867            handle,
868            name: draft.name,
869            source_type: draft.source_type,
870            source_key: draft.source_key,
871            source: draft.source,
872            payload_schema: draft.payload_schema,
873            target: draft.target,
874            target_identity: draft.target_identity,
875            event_types: draft.event_types,
876            input_template: draft.input_template,
877            target_label: draft.target_label,
878            enabled: true,
879            created_at_ms: now,
880            updated_at_ms: now,
881        };
882        state.subscriptions.insert(subscription_id, record.clone());
883        Ok(record)
884    }
885
886    async fn list_subscriptions(
887        &self,
888        filter: TriggerSubscriptionFilter,
889    ) -> Result<Vec<TriggerSubscriptionRecord>, PluginError> {
890        let state = self
891            .state
892            .lock()
893            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
894        let mut records = state
895            .subscriptions
896            .values()
897            .filter(|record| filter.matches(record))
898            .cloned()
899            .collect::<Vec<_>>();
900        records.sort_by(|left, right| {
901            left.registrant_scope_id()
902                .cmp(&right.registrant_scope_id())
903                .then_with(|| left.handle.cmp(&right.handle))
904        });
905        Ok(records)
906    }
907
908    async fn cancel_subscription(
909        &self,
910        registrant_scope_id: &str,
911        handle: &str,
912    ) -> Result<bool, PluginError> {
913        self.set_subscription_enabled(registrant_scope_id, handle, false)
914            .await
915    }
916
917    async fn set_subscription_enabled(
918        &self,
919        registrant_scope_id: &str,
920        handle: &str,
921        enabled: bool,
922    ) -> Result<bool, PluginError> {
923        let mut state = self
924            .state
925            .lock()
926            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
927        let now = self.clock.timestamp_ms();
928        let Some(record) = state.subscriptions.values_mut().find(|record| {
929            record.registrant_scope_id() == registrant_scope_id && record.handle == handle
930        }) else {
931            return Ok(false);
932        };
933        let changed = record.enabled != enabled;
934        record.enabled = enabled;
935        record.updated_at_ms = now;
936        Ok(changed)
937    }
938
939    async fn delete_subscription(
940        &self,
941        registrant_scope_id: &str,
942        handle: &str,
943    ) -> Result<bool, PluginError> {
944        let mut state = self
945            .state
946            .lock()
947            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
948        let Some(subscription_id) = state
949            .subscriptions
950            .values()
951            .find(|record| {
952                record.registrant_scope_id() == registrant_scope_id && record.handle == handle
953            })
954            .map(|record| record.subscription_id.clone())
955        else {
956            return Ok(false);
957        };
958        state.subscriptions.remove(&subscription_id);
959        state.deliveries.retain(|(_, delivery_subscription_id), _| {
960            delivery_subscription_id != &subscription_id
961        });
962        Ok(true)
963    }
964
965    async fn delete_session_subscriptions(&self, session_id: &str) -> Result<usize, PluginError> {
966        let mut state = self
967            .state
968            .lock()
969            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
970        let before = state.subscriptions.len();
971        state
972            .subscriptions
973            .retain(|_, record| record.registrant_session_id() != Some(session_id));
974        Ok(before.saturating_sub(state.subscriptions.len()))
975    }
976
977    async fn record_occurrence(
978        &self,
979        request: TriggerOccurrenceRequest,
980    ) -> Result<TriggerOccurrenceRecord, PluginError> {
981        validate_trigger_occurrence_request(&request)?;
982        let mut state = self
983            .state
984            .lock()
985            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
986        let request_hash = trigger_occurrence_request_hash(&request)?;
987        if let Some(existing_id) = state
988            .occurrence_id_by_idempotency_key
989            .get(&request.idempotency_key)
990            .cloned()
991        {
992            let existing_hash = state
993                .occurrence_hashes
994                .get(&existing_id)
995                .cloned()
996                .unwrap_or_default();
997            if existing_hash != request_hash {
998                return Err(PluginError::Session(format!(
999                    "trigger occurrence idempotency conflict for `{}`",
1000                    request.idempotency_key
1001                )));
1002            }
1003            return state.occurrences.get(&existing_id).cloned().ok_or_else(|| {
1004                PluginError::Session(format!(
1005                    "missing trigger occurrence `{existing_id}` for idempotency key"
1006                ))
1007            });
1008        }
1009        let occurrence_id = deterministic_occurrence_id(&request)?;
1010        let record = TriggerOccurrenceRecord {
1011            occurrence_id: occurrence_id.clone(),
1012            source_type: request.source_type,
1013            source_key: request.source_key,
1014            payload: request.payload,
1015            idempotency_key: request.idempotency_key.clone(),
1016            source: request.source,
1017            session_id: request.session_id,
1018            occurred_at_ms: self.clock.timestamp_ms(),
1019        };
1020        state
1021            .occurrence_id_by_idempotency_key
1022            .insert(request.idempotency_key, occurrence_id.clone());
1023        state
1024            .occurrence_hashes
1025            .insert(occurrence_id.clone(), request_hash);
1026        state.occurrences.insert(occurrence_id, record.clone());
1027        Ok(record)
1028    }
1029
1030    async fn list_occurrences(
1031        &self,
1032        filter: TriggerOccurrenceFilter,
1033    ) -> Result<Vec<TriggerOccurrenceRecord>, PluginError> {
1034        let state = self
1035            .state
1036            .lock()
1037            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
1038        let mut records = state
1039            .occurrences
1040            .values()
1041            .filter(|record| filter.matches(record))
1042            .cloned()
1043            .collect::<Vec<_>>();
1044        records.sort_by(|left, right| {
1045            left.occurred_at_ms
1046                .cmp(&right.occurred_at_ms)
1047                .then_with(|| left.occurrence_id.cmp(&right.occurrence_id))
1048        });
1049        Ok(records)
1050    }
1051
1052    async fn reserve_matching_deliveries(
1053        &self,
1054        occurrence_id: &str,
1055    ) -> Result<Vec<TriggerDeliveryReservation>, PluginError> {
1056        let mut state = self
1057            .state
1058            .lock()
1059            .map_err(|_| PluginError::Session("trigger store lock poisoned".to_string()))?;
1060        let occurrence = state
1061            .occurrences
1062            .get(occurrence_id)
1063            .cloned()
1064            .ok_or_else(|| {
1065                PluginError::Session(format!("unknown trigger occurrence `{occurrence_id}`"))
1066            })?;
1067        let subscriptions = state
1068            .subscriptions
1069            .values()
1070            .filter(|record| {
1071                record.enabled
1072                    && record.source_type == occurrence.source_type
1073                    && record.source_key == occurrence.source_key
1074                    && occurrence
1075                        .session_id
1076                        .as_deref()
1077                        .is_none_or(|session_id| record.registrant_session_id() == Some(session_id))
1078            })
1079            .cloned()
1080            .collect::<Vec<_>>();
1081        let mut deliveries = Vec::new();
1082        for subscription in subscriptions {
1083            let process_id = deterministic_delivery_process_id(
1084                &occurrence.occurrence_id,
1085                &subscription.subscription_id,
1086            )?;
1087            let key = (
1088                occurrence.occurrence_id.clone(),
1089                subscription.subscription_id.clone(),
1090            );
1091            let (delivery, reservation_status) =
1092                if let Some(delivery) = state.deliveries.get(&key).cloned() {
1093                    (delivery, TriggerDeliveryReservationStatus::AlreadyReserved)
1094                } else {
1095                    let delivery = InMemoryTriggerDeliveryRecord {
1096                        occurrence_id: occurrence.occurrence_id.clone(),
1097                        subscription_id: subscription.subscription_id.clone(),
1098                        process_id,
1099                        created_at_ms: self.clock.timestamp_ms(),
1100                    };
1101                    state.deliveries.insert(key, delivery.clone());
1102                    (delivery, TriggerDeliveryReservationStatus::Reserved)
1103                };
1104            deliveries.push(TriggerDeliveryReservation {
1105                occurrence: occurrence.clone(),
1106                subscription,
1107                process_id: delivery.process_id,
1108                created_at_ms: delivery.created_at_ms,
1109                reservation_status,
1110            });
1111        }
1112        Ok(deliveries)
1113    }
1114
1115    async fn list_deliveries_by_occurrence_id(
1116        &self,
1117        occurrence_id: &str,
1118    ) -> Result<Vec<TriggerDeliveryReservation>, PluginError> {
1119        self.list_deliveries(|delivery| delivery.occurrence_id == occurrence_id)
1120    }
1121
1122    async fn list_deliveries_by_subscription_id(
1123        &self,
1124        subscription_id: &str,
1125    ) -> Result<Vec<TriggerDeliveryReservation>, PluginError> {
1126        self.list_deliveries(|delivery| delivery.subscription_id == subscription_id)
1127    }
1128
1129    async fn list_deliveries_by_process_id(
1130        &self,
1131        process_id: &str,
1132    ) -> Result<Vec<TriggerDeliveryReservation>, PluginError> {
1133        self.list_deliveries(|delivery| delivery.process_id == process_id)
1134    }
1135}
1136
1137fn default_enabled() -> bool {
1138    true
1139}
1140
1141pub fn default_trigger_source_key(
1142    source_type: &str,
1143    source: &serde_json::Value,
1144) -> Result<String, PluginError> {
1145    let digest = crate::stable_hash::stable_json_sha256_hex(&(source_type, source))
1146        .map_err(|err| PluginError::Session(format!("failed to hash trigger source key: {err}")))?;
1147    Ok(format!("source:{source_type}:sha256:{digest}"))
1148}
1149
1150pub fn empty_trigger_source_key(source_type: &str) -> Result<String, PluginError> {
1151    default_trigger_source_key(source_type, &serde_json::json!({}))
1152}
1153
1154pub fn deterministic_occurrence_id(
1155    request: &TriggerOccurrenceRequest,
1156) -> Result<String, PluginError> {
1157    let digest = crate::stable_hash::stable_json_sha256_hex(&(
1158        request.source_type.as_str(),
1159        request.source_key.as_str(),
1160        request.idempotency_key.as_str(),
1161    ))
1162    .map_err(|err| PluginError::Session(format!("failed to hash trigger occurrence: {err}")))?;
1163    Ok(format!("trigger:{digest}"))
1164}
1165
1166pub fn deterministic_delivery_process_id(
1167    occurrence_id: &str,
1168    subscription_id: &str,
1169) -> Result<String, PluginError> {
1170    let digest = crate::stable_hash::stable_json_sha256_hex(&(occurrence_id, subscription_id))
1171        .map_err(|err| PluginError::Session(format!("failed to hash trigger delivery: {err}")))?;
1172    Ok(format!("process:trigger:{digest}"))
1173}
1174
1175#[derive(Clone)]
1176pub struct TriggerRouter {
1177    store: Arc<dyn TriggerStore>,
1178    process_registry: Option<Arc<dyn crate::ProcessRegistry>>,
1179    process_work_driver: Option<crate::ProcessWorkDriver>,
1180}
1181
1182impl TriggerRouter {
1183    pub fn new(
1184        store: Arc<dyn TriggerStore>,
1185        process_registry: Option<Arc<dyn crate::ProcessRegistry>>,
1186        process_work_driver: Option<crate::ProcessWorkDriver>,
1187    ) -> Self {
1188        Self {
1189            store,
1190            process_registry,
1191            process_work_driver,
1192        }
1193    }
1194
1195    pub fn store(&self) -> Arc<dyn TriggerStore> {
1196        Arc::clone(&self.store)
1197    }
1198
1199    pub async fn emit(
1200        &self,
1201        request: TriggerOccurrenceRequest,
1202        effect_controller: &dyn crate::RuntimeEffectController,
1203    ) -> Result<TriggerEmitReport, PluginError> {
1204        let occurrence = self.store.record_occurrence(request).await?;
1205        let reservations = self
1206            .store
1207            .reserve_matching_deliveries(&occurrence.occurrence_id)
1208            .await?;
1209        let Some(process_registry) = self.process_registry.as_ref() else {
1210            let deliveries = reservations
1211                .iter()
1212                .map(|reservation| {
1213                    let outcome = match reservation.reservation_status {
1214                        TriggerDeliveryReservationStatus::Reserved => {
1215                            TriggerDeliveryEmitOutcome::Failed {
1216                                reason: "trigger delivery requires a process registry".to_string(),
1217                            }
1218                        }
1219                        TriggerDeliveryReservationStatus::AlreadyReserved => {
1220                            TriggerDeliveryEmitOutcome::AlreadyReserved
1221                        }
1222                    };
1223                    reservation.emit_report(outcome)
1224                })
1225                .collect();
1226            return Ok(TriggerEmitReport::new(occurrence.occurrence_id, deliveries));
1227        };
1228        let mut deliveries = Vec::new();
1229        let mut started_any = false;
1230        for reservation in reservations {
1231            if reservation.reservation_status == TriggerDeliveryReservationStatus::AlreadyReserved {
1232                deliveries
1233                    .push(reservation.emit_report(TriggerDeliveryEmitOutcome::AlreadyReserved));
1234                continue;
1235            }
1236            if let Err(err) = self
1237                .start_delivery(
1238                    &reservation,
1239                    Arc::clone(process_registry),
1240                    effect_controller,
1241                )
1242                .await
1243            {
1244                deliveries.push(reservation.emit_report(TriggerDeliveryEmitOutcome::Failed {
1245                    reason: err.to_string(),
1246                }));
1247                continue;
1248            }
1249            started_any = true;
1250            deliveries.push(reservation.emit_report(TriggerDeliveryEmitOutcome::Started));
1251        }
1252        if started_any && let Some(driver) = self.process_work_driver.as_ref() {
1253            driver.claim_and_run_pending("trigger_delivery").await?;
1254        }
1255        Ok(TriggerEmitReport::new(occurrence.occurrence_id, deliveries))
1256    }
1257
1258    pub(crate) async fn start_delivery(
1259        &self,
1260        reservation: &TriggerDeliveryReservation,
1261        process_registry: Arc<dyn crate::ProcessRegistry>,
1262        effect_controller: &dyn crate::RuntimeEffectController,
1263    ) -> Result<(), PluginError> {
1264        let subscription = &reservation.subscription;
1265        let occurrence = &reservation.occurrence;
1266        subscription
1267            .payload_schema
1268            .validate(&occurrence.payload)
1269            .map_err(|err| {
1270                PluginError::Session(format!(
1271                    "invalid payload for trigger `{}`: {err}",
1272                    subscription.handle
1273                ))
1274            })?;
1275        let args =
1276            materialize_trigger_process_args(&subscription.input_template, &occurrence.payload)?;
1277        let target = apply_trigger_inputs(subscription.target.clone(), args)?;
1278        let originator_scope_id = subscription.registrant_scope_id();
1279        let trigger_causal_ref = crate::CausalRef::TriggerOccurrence {
1280            occurrence_id: occurrence.occurrence_id.clone(),
1281            subscription_id: Some(subscription.subscription_id.clone()),
1282        };
1283        let trigger_occurrence_invocation = crate::runtime::causal::trigger_occurrence_invocation(
1284            &originator_scope_id,
1285            &occurrence.occurrence_id,
1286        );
1287        let registration = crate::ProcessRegistration::new(
1288            reservation.process_id.clone(),
1289            target.clone(),
1290            // Trigger targets are journaled engine/tool rows, idempotent by
1291            // process id, so recovery may re-execute them (ADR 0019).
1292            crate::RecoveryDisposition::Rerunnable,
1293            crate::ProcessProvenance::new(subscription.registrant.clone())
1294                .with_caused_by(Some(trigger_causal_ref.clone())),
1295        )
1296        .with_identity(subscription.target_identity.clone())
1297        .with_extra_event_types(subscription.event_types.clone())
1298        .with_execution_env_ref(Some(subscription.env_ref.clone()))
1299        .with_wake_target(subscription.wake_target.clone());
1300        let descriptor_kind = subscription.target_identity.kind.clone();
1301        let grant =
1302            subscription
1303                .wake_target
1304                .clone()
1305                .map(|session_scope| crate::ProcessStartGrant {
1306                    session_scope,
1307                    descriptor: crate::ProcessHandleDescriptor::new(
1308                        Some(descriptor_kind.as_str()),
1309                        subscription.target_label.as_deref(),
1310                    ),
1311                });
1312        let execution_context = crate::ProcessExecutionContext::default()
1313            .with_causal_invocation(Some(trigger_occurrence_invocation));
1314        let command = crate::ProcessCommand::Start {
1315            registration,
1316            grant,
1317            execution_context: Box::new(execution_context),
1318        };
1319        let effect_id = command.effect_id();
1320        let invocation = crate::RuntimeInvocation::effect(
1321            crate::RuntimeScope::new(originator_scope_id),
1322            effect_id.clone(),
1323            crate::RuntimeEffectKind::Process,
1324            format!(
1325                "trigger:{}:{}",
1326                occurrence.occurrence_id, subscription.subscription_id
1327            ),
1328        )
1329        .with_caused_by(Some(trigger_causal_ref));
1330        let outcome = effect_controller
1331            .execute_effect(
1332                crate::RuntimeEffectEnvelope::new(
1333                    invocation,
1334                    crate::RuntimeEffectCommand::process(command),
1335                ),
1336                crate::RuntimeEffectLocalExecutor::processes(
1337                    process_registry,
1338                    self.process_work_driver.clone(),
1339                ),
1340            )
1341            .await?;
1342        match outcome {
1343            crate::RuntimeEffectOutcome::Process {
1344                result: crate::ProcessEffectOutcome::Start { .. },
1345            } => Ok(()),
1346            other => Err(PluginError::Session(format!(
1347                "trigger process start returned the wrong outcome: {}",
1348                other.kind().as_str()
1349            ))),
1350        }
1351    }
1352}
1353
1354fn materialize_trigger_process_args(
1355    input_template: &BTreeMap<String, TriggerInputBinding>,
1356    event_payload: &serde_json::Value,
1357) -> Result<serde_json::Map<String, serde_json::Value>, PluginError> {
1358    let mut args = serde_json::Map::new();
1359    for (input_name, input) in input_template {
1360        let value = match input {
1361            TriggerInputBinding::Event => event_payload.clone(),
1362            TriggerInputBinding::Fixed { value } => value.clone(),
1363        };
1364        args.insert(input_name.to_string(), value);
1365    }
1366    Ok(args)
1367}
1368
1369fn apply_trigger_inputs(
1370    mut target: crate::ProcessInput,
1371    args: serde_json::Map<String, serde_json::Value>,
1372) -> Result<crate::ProcessInput, PluginError> {
1373    match &mut target {
1374        crate::ProcessInput::Engine { payload, .. } => {
1375            let object = payload.as_object_mut().ok_or_else(|| {
1376                PluginError::Session(
1377                    "trigger engine target payload must be a JSON object".to_string(),
1378                )
1379            })?;
1380            object.insert("args".to_string(), serde_json::Value::Object(args));
1381            Ok(target)
1382        }
1383        other => Err(PluginError::Session(format!(
1384            "trigger target must be an engine process, got {}",
1385            other.engine_kind()
1386        ))),
1387    }
1388}
1389
1390pub fn validate_trigger_occurrence_request(
1391    request: &TriggerOccurrenceRequest,
1392) -> Result<(), PluginError> {
1393    if request.source_type.trim().is_empty() {
1394        return Err(PluginError::Session(
1395            "trigger occurrence requires source_type".to_string(),
1396        ));
1397    }
1398    if request.source_key.trim().is_empty() {
1399        return Err(PluginError::Session(
1400            "trigger occurrence requires source_key".to_string(),
1401        ));
1402    }
1403    if request.idempotency_key.trim().is_empty() {
1404        return Err(PluginError::Session(
1405            "trigger occurrence requires idempotency_key".to_string(),
1406        ));
1407    }
1408    Ok(())
1409}
1410
1411pub fn trigger_occurrence_request_hash(
1412    request: &TriggerOccurrenceRequest,
1413) -> Result<String, PluginError> {
1414    crate::stable_hash::stable_json_sha256_hex(&(
1415        request.source_type.as_str(),
1416        request.source_key.as_str(),
1417        &request.payload,
1418        &request.source,
1419    ))
1420    .map_err(|err| PluginError::Session(format!("failed to hash trigger occurrence: {err}")))
1421}
1422
1423#[cfg(test)]
1424mod tests {
1425    use super::*;
1426
1427    fn button_payload_schema() -> crate::LashSchema {
1428        crate::LashSchema::any()
1429    }
1430
1431    fn trigger_process_draft(source_key: &str, process_name: &str) -> TriggerSubscriptionDraft {
1432        TriggerSubscriptionDraft::for_process(
1433            crate::ProcessOriginator::host(),
1434            crate::ProcessExecutionEnvRef::new(format!("process-env:{process_name}")),
1435            "ui.button.pressed",
1436            source_key,
1437            crate::ProcessInput::Engine {
1438                kind: "test-engine".to_string(),
1439                payload: serde_json::json!({ "process": process_name }),
1440            },
1441            crate::ProcessIdentity::new("test-engine").with_label(Some(process_name)),
1442        )
1443        .with_payload_schema(crate::LashSchema::any())
1444    }
1445
1446    fn button_occurrence(
1447        source_key: impl Into<String>,
1448        idempotency_key: impl Into<String>,
1449    ) -> TriggerOccurrenceRequest {
1450        TriggerOccurrenceRequest::new(
1451            "ui.button.pressed",
1452            source_key,
1453            serde_json::json!({ "button": "Blue" }),
1454            idempotency_key,
1455        )
1456    }
1457
1458    #[test]
1459    fn trigger_catalog_rejects_duplicate_trigger_source_identity() {
1460        let mut catalog = TriggerEventCatalog::new();
1461        catalog
1462            .declare(TriggerEvent::new(
1463                "Button",
1464                "ui.button",
1465                "pressed",
1466                button_payload_schema(),
1467            ))
1468            .expect("first trigger occurrence");
1469
1470        let err = catalog
1471            .declare(TriggerEvent::new(
1472                "AlternateButton",
1473                "ui.button",
1474                "pressed",
1475                button_payload_schema(),
1476            ))
1477            .expect_err("duplicate public source identity should be rejected");
1478
1479        assert!(err.contains("duplicate trigger source `ui.button.pressed`"));
1480    }
1481
1482    #[tokio::test]
1483    async fn trigger_store_rejects_mismatched_target_label() {
1484        let store = InMemoryTriggerStore::default();
1485        let draft = TriggerSubscriptionDraft::for_process(
1486            crate::ProcessOriginator::host(),
1487            crate::ProcessExecutionEnvRef::new("process-env:test"),
1488            "ui.button.pressed",
1489            "source-key",
1490            crate::ProcessInput::External {
1491                metadata: serde_json::json!({}),
1492            },
1493            crate::ProcessIdentity::new("external").with_label(Some("expected")),
1494        )
1495        .with_target_label("other");
1496
1497        let err = store
1498            .register_subscription(draft)
1499            .await
1500            .expect_err("mismatched target labels should be rejected");
1501        assert!(err.to_string().contains("target_label must match"));
1502    }
1503
1504    #[tokio::test]
1505    async fn trigger_emit_report_records_started_and_already_reserved_deliveries() {
1506        let store = Arc::new(InMemoryTriggerStore::default());
1507        let registry: Arc<dyn crate::ProcessRegistry> =
1508            Arc::new(crate::TestLocalProcessRegistry::default());
1509        let source_key = empty_trigger_source_key("ui.button.pressed").expect("source key");
1510        let subscription = store
1511            .register_subscription(trigger_process_draft(&source_key, "started"))
1512            .await
1513            .expect("register subscription");
1514        let router = TriggerRouter::new(store, Some(Arc::clone(&registry)), None);
1515        let controller = crate::InlineRuntimeEffectController;
1516
1517        let report = router
1518            .emit(
1519                button_occurrence(source_key.clone(), "button-blue-report"),
1520                &controller,
1521            )
1522            .await
1523            .expect("emit trigger");
1524        assert_eq!(report.deliveries.len(), 1);
1525        let delivery = &report.deliveries[0];
1526        assert_eq!(delivery.occurrence_id, report.occurrence_id);
1527        assert_eq!(delivery.subscription_id, subscription.subscription_id);
1528        assert_eq!(delivery.outcome, TriggerDeliveryEmitOutcome::Started);
1529        let record = registry
1530            .get_process(&delivery.process_id)
1531            .await
1532            .expect("started process record");
1533        assert!(matches!(
1534            record.provenance.caused_by,
1535            Some(crate::CausalRef::TriggerOccurrence {
1536                occurrence_id,
1537                subscription_id: Some(subscription_id),
1538            }) if occurrence_id == report.occurrence_id
1539                && subscription_id == subscription.subscription_id
1540        ));
1541
1542        let replay = router
1543            .emit(
1544                button_occurrence(source_key, "button-blue-report"),
1545                &controller,
1546            )
1547            .await
1548            .expect("replay trigger");
1549        assert_eq!(replay.deliveries.len(), 1);
1550        assert_eq!(
1551            replay.deliveries[0].outcome,
1552            TriggerDeliveryEmitOutcome::AlreadyReserved
1553        );
1554        assert_eq!(replay.deliveries[0].process_id, delivery.process_id);
1555    }
1556
1557    #[tokio::test]
1558    async fn trigger_emit_report_records_failed_delivery_outcome() {
1559        let store = Arc::new(InMemoryTriggerStore::default());
1560        let source_key = empty_trigger_source_key("ui.button.pressed").expect("source key");
1561        let subscription = store
1562            .register_subscription(trigger_process_draft(&source_key, "failed"))
1563            .await
1564            .expect("register subscription");
1565        let router = TriggerRouter::new(store, None, None);
1566        let controller = crate::InlineRuntimeEffectController;
1567
1568        let report = router
1569            .emit(
1570                button_occurrence(source_key, "button-blue-failed"),
1571                &controller,
1572            )
1573            .await
1574            .expect("emit trigger");
1575        assert_eq!(report.deliveries.len(), 1);
1576        let delivery = &report.deliveries[0];
1577        assert_eq!(delivery.subscription_id, subscription.subscription_id);
1578        assert!(matches!(
1579            &delivery.outcome,
1580            TriggerDeliveryEmitOutcome::Failed { reason }
1581                if reason.contains("process registry")
1582        ));
1583    }
1584}