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