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 #[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 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 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 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(®istry)), 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}