Skip to main content

lash_core/triggers/
router.rs

1use super::*;
2
3pub fn deterministic_subscription_id(
4    owner_scope: &TriggerOwnerScope,
5    subscription_key: &str,
6) -> Result<String, PluginError> {
7    let digest = crate::stable_hash::stable_json_sha256_hex(&(
8        "lash.trigger-subscription",
9        1_u8,
10        owner_scope,
11        subscription_key,
12    ))
13    .map_err(|err| PluginError::Session(format!("failed to hash trigger identity: {err}")))?;
14    Ok(format!("trigger-subscription:v1:sha256:{digest}"))
15}
16
17pub fn trigger_subscription_definition_hash(
18    owner_scope: &TriggerOwnerScope,
19    draft: &TriggerSubscriptionDraft,
20) -> Result<String, PluginError> {
21    crate::stable_hash::stable_json_sha256_hex(&(
22        "lash.trigger-subscription-definition",
23        1_u8,
24        owner_scope,
25        draft,
26    ))
27    .map_err(|err| PluginError::Session(format!("failed to hash trigger definition: {err}")))
28}
29
30pub(super) fn reserve_in_memory_for_occurrence(
31    state: &mut InMemoryTriggerEventState,
32    occurrence: &TriggerOccurrenceRecord,
33    clock: &dyn crate::Clock,
34) -> Result<Vec<TriggerDeliveryReservation>, PluginError> {
35    let subscriptions = state
36        .subscriptions
37        .values()
38        .filter(|record| {
39            record.enabled
40                && !record.tombstoned
41                && record.source_type == occurrence.source_type
42                && record.source_key == occurrence.source_key
43                && occurrence
44                    .session_id
45                    .as_deref()
46                    .is_none_or(|session_id| record.registrant_session_id() == Some(session_id))
47        })
48        .cloned()
49        .collect::<Vec<_>>();
50    let mut reservations = Vec::new();
51    for subscription in subscriptions {
52        let process_id = deterministic_delivery_process_id(
53            &occurrence.occurrence_id,
54            &subscription.subscription_id,
55            &subscription.incarnation,
56            subscription.revision,
57        )?;
58        let key = (
59            occurrence.occurrence_id.clone(),
60            subscription.subscription_id.clone(),
61        );
62        let delivery = InMemoryTriggerDeliveryRecord {
63            occurrence_id: occurrence.occurrence_id.clone(),
64            subscription_id: subscription.subscription_id.clone(),
65            process_id,
66            created_at_ms: clock.timestamp_ms(),
67            subscription_snapshot: subscription.clone(),
68        };
69        state.deliveries.insert(key, delivery.clone());
70        reservations.push(TriggerDeliveryReservation {
71            occurrence: occurrence.clone(),
72            subscription,
73            process_id: delivery.process_id,
74            created_at_ms: delivery.created_at_ms,
75            reservation_status: TriggerDeliveryReservationStatus::Reserved,
76        });
77    }
78    Ok(reservations)
79}
80
81pub(super) fn default_enabled() -> bool {
82    true
83}
84
85pub fn default_trigger_source_key(
86    source_type: &str,
87    source: &serde_json::Value,
88) -> Result<String, PluginError> {
89    let digest = crate::stable_hash::stable_json_sha256_hex(&(source_type, source))
90        .map_err(|err| PluginError::Session(format!("failed to hash trigger source key: {err}")))?;
91    Ok(format!("source:{source_type}:sha256:{digest}"))
92}
93
94pub fn empty_trigger_source_key(source_type: &str) -> Result<String, PluginError> {
95    default_trigger_source_key(source_type, &serde_json::json!({}))
96}
97
98pub fn deterministic_occurrence_id(
99    request: &TriggerOccurrenceRequest,
100) -> Result<String, PluginError> {
101    let digest = crate::stable_hash::stable_json_sha256_hex(&(
102        request.source_type.as_str(),
103        request.source_key.as_str(),
104        request.idempotency_key.as_str(),
105    ))
106    .map_err(|err| PluginError::Session(format!("failed to hash trigger occurrence: {err}")))?;
107    Ok(format!("trigger:{digest}"))
108}
109
110pub fn deterministic_delivery_process_id(
111    occurrence_id: &str,
112    subscription_id: &str,
113    incarnation: &str,
114    revision: u64,
115) -> Result<String, PluginError> {
116    let digest = crate::stable_hash::stable_json_sha256_hex(&(
117        "lash.trigger-delivery",
118        1_u8,
119        occurrence_id,
120        subscription_id,
121        incarnation,
122        revision,
123    ))
124    .map_err(|err| PluginError::Session(format!("failed to hash trigger delivery: {err}")))?;
125    Ok(format!("process:trigger:{digest}"))
126}
127
128#[derive(Clone)]
129pub struct TriggerRouter {
130    store: Arc<dyn TriggerStore>,
131    process_registry: Option<Arc<dyn crate::ProcessRegistry>>,
132    process_work_driver: Option<crate::ProcessWorkDriver>,
133}
134
135impl TriggerRouter {
136    pub fn new(
137        store: Arc<dyn TriggerStore>,
138        process_registry: Option<Arc<dyn crate::ProcessRegistry>>,
139        process_work_driver: Option<crate::ProcessWorkDriver>,
140    ) -> Self {
141        Self {
142            store,
143            process_registry,
144            process_work_driver,
145        }
146    }
147
148    pub fn store(&self) -> Arc<dyn TriggerStore> {
149        Arc::clone(&self.store)
150    }
151
152    pub async fn emit(
153        &self,
154        request: TriggerOccurrenceRequest,
155        effect_controller: &dyn crate::RuntimeEffectController,
156    ) -> Result<TriggerEmitReport, PluginError> {
157        let TriggerIngressResult {
158            occurrence,
159            reservations,
160        } = self.store.ingest_occurrence(request).await?;
161        let Some(process_registry) = self.process_registry.as_ref() else {
162            let deliveries = reservations
163                .iter()
164                .map(|reservation| {
165                    let outcome = match reservation.reservation_status {
166                        TriggerDeliveryReservationStatus::Reserved => {
167                            TriggerDeliveryEmitOutcome::Failed {
168                                reason: "trigger delivery requires a process registry".to_string(),
169                            }
170                        }
171                        TriggerDeliveryReservationStatus::AlreadyReserved => {
172                            TriggerDeliveryEmitOutcome::AlreadyReserved
173                        }
174                    };
175                    reservation.emit_report(outcome)
176                })
177                .collect();
178            return Ok(TriggerEmitReport::new(occurrence.occurrence_id, deliveries));
179        };
180        let mut deliveries = Vec::new();
181        let mut started_any = false;
182        for reservation in reservations {
183            if reservation.reservation_status == TriggerDeliveryReservationStatus::AlreadyReserved {
184                deliveries
185                    .push(reservation.emit_report(TriggerDeliveryEmitOutcome::AlreadyReserved));
186                continue;
187            }
188            if let Err(err) = self
189                .start_delivery(
190                    &reservation,
191                    Arc::clone(process_registry),
192                    effect_controller,
193                )
194                .await
195            {
196                deliveries.push(reservation.emit_report(TriggerDeliveryEmitOutcome::Failed {
197                    reason: err.to_string(),
198                }));
199                continue;
200            }
201            started_any = true;
202            deliveries.push(reservation.emit_report(TriggerDeliveryEmitOutcome::Started));
203        }
204        if started_any && let Some(driver) = self.process_work_driver.as_ref() {
205            driver.claim_and_run_pending("trigger_delivery").await?;
206        }
207        Ok(TriggerEmitReport::new(occurrence.occurrence_id, deliveries))
208    }
209
210    pub(crate) async fn start_delivery(
211        &self,
212        reservation: &TriggerDeliveryReservation,
213        process_registry: Arc<dyn crate::ProcessRegistry>,
214        effect_controller: &dyn crate::RuntimeEffectController,
215    ) -> Result<(), PluginError> {
216        let subscription = &reservation.subscription;
217        let occurrence = &reservation.occurrence;
218        subscription
219            .payload_schema
220            .validate(&occurrence.payload)
221            .map_err(|err| {
222                PluginError::Session(format!(
223                    "invalid payload for trigger `{}`: {err}",
224                    subscription.subscription_key
225                ))
226            })?;
227        let args =
228            materialize_trigger_process_args(&subscription.input_template, &occurrence.payload)?;
229        let target = apply_trigger_inputs(subscription.target.clone(), args)?;
230        let originator_scope_id = subscription.registrant_scope_id();
231        let trigger_causal_ref = crate::CausalRef::TriggerOccurrence {
232            occurrence_id: occurrence.occurrence_id.clone(),
233            subscription_id: Some(subscription.subscription_id.clone()),
234            subscription_incarnation: Some(subscription.incarnation.clone()),
235            subscription_revision: Some(subscription.revision),
236        };
237        let trigger_occurrence_invocation = crate::runtime::causal::trigger_occurrence_invocation(
238            &originator_scope_id,
239            &occurrence.occurrence_id,
240        );
241        let registration = crate::ProcessRegistration::new(
242            reservation.process_id.clone(),
243            target.clone(),
244            // Trigger targets are journaled engine/tool rows, idempotent by
245            // process id, so recovery may re-execute them (ADR 0019).
246            crate::RecoveryDisposition::Rerunnable,
247            crate::ProcessProvenance::new(subscription.registrant.clone())
248                .with_caused_by(Some(trigger_causal_ref.clone())),
249        )
250        .with_identity(subscription.target_identity.clone())
251        .with_extra_event_types(subscription.event_types.clone())
252        .with_execution_env_ref(Some(subscription.env_ref.clone()))
253        .with_wake_target(subscription.wake_target.clone());
254        let descriptor_kind = subscription.target_identity.kind.clone();
255        let grant =
256            subscription
257                .wake_target
258                .clone()
259                .map(|session_scope| crate::ProcessStartGrant {
260                    session_scope,
261                    descriptor: crate::ProcessHandleDescriptor::new(
262                        Some(descriptor_kind.as_str()),
263                        subscription.target_label.as_deref(),
264                    ),
265                });
266        let execution_context = crate::ProcessExecutionContext::default()
267            .with_causal_invocation(Some(trigger_occurrence_invocation));
268        let command = crate::ProcessCommand::Start {
269            registration,
270            grant,
271            execution_context: Box::new(execution_context),
272        };
273        let effect_id = command.effect_id();
274        let invocation = crate::RuntimeInvocation::effect(
275            crate::RuntimeScope::new(originator_scope_id),
276            effect_id.clone(),
277            crate::RuntimeEffectKind::Process,
278            format!(
279                "trigger:{}:{}:{}:{}",
280                occurrence.occurrence_id,
281                subscription.subscription_id,
282                subscription.incarnation,
283                subscription.revision
284            ),
285        )
286        .with_caused_by(Some(trigger_causal_ref));
287        let outcome = effect_controller
288            .execute_effect(
289                crate::RuntimeEffectEnvelope::new(
290                    invocation,
291                    crate::RuntimeEffectCommand::process(command),
292                ),
293                crate::RuntimeEffectLocalExecutor::processes(
294                    process_registry,
295                    self.process_work_driver.clone(),
296                ),
297            )
298            .await?;
299        match outcome {
300            crate::RuntimeEffectOutcome::Process {
301                result: crate::ProcessEffectOutcome::Start { .. },
302            } => Ok(()),
303            other => Err(PluginError::Session(format!(
304                "trigger process start returned the wrong outcome: {}",
305                other.kind().as_str()
306            ))),
307        }
308    }
309}
310
311fn materialize_trigger_process_args(
312    input_template: &BTreeMap<String, TriggerInputBinding>,
313    event_payload: &serde_json::Value,
314) -> Result<serde_json::Map<String, serde_json::Value>, PluginError> {
315    let mut args = serde_json::Map::new();
316    for (input_name, input) in input_template {
317        let value = match input {
318            TriggerInputBinding::Event => event_payload.clone(),
319            TriggerInputBinding::Fixed { value } => value.clone(),
320        };
321        args.insert(input_name.to_string(), value);
322    }
323    Ok(args)
324}
325
326fn apply_trigger_inputs(
327    mut target: crate::ProcessInput,
328    args: serde_json::Map<String, serde_json::Value>,
329) -> Result<crate::ProcessInput, PluginError> {
330    match &mut target {
331        crate::ProcessInput::Engine { payload, .. } => {
332            let object = payload.as_object_mut().ok_or_else(|| {
333                PluginError::Session(
334                    "trigger engine target payload must be a JSON object".to_string(),
335                )
336            })?;
337            object.insert("args".to_string(), serde_json::Value::Object(args));
338            Ok(target)
339        }
340        other => Err(PluginError::Session(format!(
341            "trigger target must be an engine process, got {}",
342            other.engine_kind()
343        ))),
344    }
345}
346
347pub fn validate_trigger_occurrence_request(
348    request: &TriggerOccurrenceRequest,
349) -> Result<(), PluginError> {
350    if request.source_type.trim().is_empty() {
351        return Err(PluginError::Session(
352            "trigger occurrence requires source_type".to_string(),
353        ));
354    }
355    if request.source_key.trim().is_empty() {
356        return Err(PluginError::Session(
357            "trigger occurrence requires source_key".to_string(),
358        ));
359    }
360    if request.idempotency_key.trim().is_empty() {
361        return Err(PluginError::Session(
362            "trigger occurrence requires idempotency_key".to_string(),
363        ));
364    }
365    Ok(())
366}
367
368pub fn trigger_occurrence_request_hash(
369    request: &TriggerOccurrenceRequest,
370) -> Result<String, PluginError> {
371    crate::stable_hash::stable_json_sha256_hex(&(
372        request.source_type.as_str(),
373        request.source_key.as_str(),
374        &request.payload,
375        &request.source,
376    ))
377    .map_err(|err| PluginError::Session(format!("failed to hash trigger occurrence: {err}")))
378}
379
380#[cfg(test)]
381mod tests {
382    use super::*;
383
384    fn button_payload_schema() -> crate::LashSchema {
385        crate::LashSchema::any()
386    }
387
388    fn trigger_process_draft(source_key: &str, process_name: &str) -> TriggerSubscriptionDraft {
389        TriggerSubscriptionDraft::for_process(
390            format!("test/{process_name}"),
391            crate::ProcessExecutionEnvRef::new(format!("process-env:{process_name}")),
392            "ui.button.pressed",
393            source_key,
394            crate::ProcessInput::Engine {
395                kind: "test-engine".to_string(),
396                payload: serde_json::json!({ "process": process_name }),
397            },
398            crate::ProcessIdentity::new("test-engine").with_label(Some(process_name)),
399        )
400        .with_payload_schema(crate::LashSchema::any())
401    }
402
403    async fn register(
404        store: &InMemoryTriggerStore,
405        operation_id: &str,
406        draft: TriggerSubscriptionDraft,
407    ) -> TriggerSubscriptionRecord {
408        let outcome = store
409            .execute_command(
410                operation_id,
411                TriggerCommand::Register {
412                    owner_scope: TriggerOwnerScope::host("test").unwrap(),
413                    actor: crate::ProcessOriginator::host_scoped("test"),
414                    draft,
415                },
416            )
417            .await
418            .expect("execute registration")
419            .expect("register subscription");
420        let TriggerCommandOutcome::Mutation { receipt } = outcome else {
421            panic!("expected mutation receipt")
422        };
423        receipt.record_snapshot
424    }
425
426    fn button_occurrence(
427        source_key: impl Into<String>,
428        idempotency_key: impl Into<String>,
429    ) -> TriggerOccurrenceRequest {
430        TriggerOccurrenceRequest::new(
431            "ui.button.pressed",
432            source_key,
433            serde_json::json!({ "button": "Blue" }),
434            idempotency_key,
435        )
436    }
437
438    #[test]
439    fn trigger_catalog_rejects_duplicate_trigger_source_identity() {
440        let mut catalog = TriggerEventCatalog::new();
441        catalog
442            .declare(TriggerEvent::new(
443                "Button",
444                "ui.button",
445                "pressed",
446                button_payload_schema(),
447            ))
448            .expect("first trigger occurrence");
449
450        let err = catalog
451            .declare(TriggerEvent::new(
452                "AlternateButton",
453                "ui.button",
454                "pressed",
455                button_payload_schema(),
456            ))
457            .expect_err("duplicate public source identity should be rejected");
458
459        assert!(err.contains("duplicate trigger source `ui.button.pressed`"));
460    }
461
462    #[tokio::test]
463    async fn trigger_store_rejects_mismatched_target_label() {
464        let store = InMemoryTriggerStore::default();
465        let draft = TriggerSubscriptionDraft::for_process(
466            "mismatched-label",
467            crate::ProcessExecutionEnvRef::new("process-env:test"),
468            "ui.button.pressed",
469            "source-key",
470            crate::ProcessInput::External {
471                metadata: serde_json::json!({}),
472            },
473            crate::ProcessIdentity::new("external").with_label(Some("expected")),
474        )
475        .with_target_label("other");
476
477        let err = store
478            .execute_command(
479                "mismatched-label",
480                TriggerCommand::Register {
481                    owner_scope: TriggerOwnerScope::host("test").unwrap(),
482                    actor: crate::ProcessOriginator::host_scoped("test"),
483                    draft,
484                },
485            )
486            .await
487            .expect("store execution")
488            .expect_err("mismatched target labels should be rejected");
489        assert!(err.to_string().contains("target_label must match"));
490    }
491
492    #[tokio::test]
493    async fn trigger_emit_report_records_started_and_already_reserved_deliveries() {
494        let store = Arc::new(InMemoryTriggerStore::default());
495        let registry: Arc<dyn crate::ProcessRegistry> =
496            Arc::new(crate::TestLocalProcessRegistry::default());
497        let source_key = empty_trigger_source_key("ui.button.pressed").expect("source key");
498        let subscription = register(
499            store.as_ref(),
500            "started-register",
501            trigger_process_draft(&source_key, "started"),
502        )
503        .await;
504        let router = TriggerRouter::new(store, Some(Arc::clone(&registry)), None);
505        let controller = crate::InlineRuntimeEffectController::default();
506
507        let report = router
508            .emit(
509                button_occurrence(source_key.clone(), "button-blue-report"),
510                &controller,
511            )
512            .await
513            .expect("emit trigger");
514        assert_eq!(report.deliveries.len(), 1);
515        let delivery = &report.deliveries[0];
516        assert_eq!(delivery.occurrence_id, report.occurrence_id);
517        assert_eq!(delivery.subscription_id, subscription.subscription_id);
518        assert_eq!(delivery.outcome, TriggerDeliveryEmitOutcome::Started);
519        let record = registry
520            .get_process(&delivery.process_id)
521            .await
522            .expect("started process record");
523        assert!(matches!(
524            record.provenance.caused_by,
525            Some(crate::CausalRef::TriggerOccurrence {
526                occurrence_id,
527                subscription_id: Some(subscription_id),
528                ..
529            }) if occurrence_id == report.occurrence_id
530                && subscription_id == subscription.subscription_id
531        ));
532
533        let replay = router
534            .emit(
535                button_occurrence(source_key, "button-blue-report"),
536                &controller,
537            )
538            .await
539            .expect("replay trigger");
540        assert_eq!(replay.deliveries.len(), 1);
541        assert_eq!(
542            replay.deliveries[0].outcome,
543            TriggerDeliveryEmitOutcome::AlreadyReserved
544        );
545        assert_eq!(replay.deliveries[0].process_id, delivery.process_id);
546    }
547
548    #[tokio::test]
549    async fn trigger_emit_report_records_failed_delivery_outcome() {
550        let store = Arc::new(InMemoryTriggerStore::default());
551        let source_key = empty_trigger_source_key("ui.button.pressed").expect("source key");
552        let subscription = register(
553            store.as_ref(),
554            "failed-register",
555            trigger_process_draft(&source_key, "failed"),
556        )
557        .await;
558        let router = TriggerRouter::new(store, None, None);
559        let controller = crate::InlineRuntimeEffectController::default();
560
561        let report = router
562            .emit(
563                button_occurrence(source_key, "button-blue-failed"),
564                &controller,
565            )
566            .await
567            .expect("emit trigger");
568        assert_eq!(report.deliveries.len(), 1);
569        let delivery = &report.deliveries[0];
570        assert_eq!(delivery.subscription_id, subscription.subscription_id);
571        assert!(matches!(
572            &delivery.outcome,
573            TriggerDeliveryEmitOutcome::Failed { reason }
574                if reason.contains("process registry")
575        ));
576    }
577}