fraiseql_functions/triggers/mutation/mod.rs
1//! Mutation triggers: `after:mutation` and `before:mutation`.
2//!
3//! ## `after:mutation` Triggers
4//!
5//! Fire asynchronously after a mutation completes (insert, update, or delete).
6//! The function receives the old and new row data. Failures do not block the mutation.
7//!
8//! ## `before:mutation` Triggers
9//!
10//! Fire synchronously before a mutation executes. The function can:
11//! - Return `Proceed(modified_input)` to allow the mutation with possibly modified input
12//! - Return `Abort(error_message)` to cancel the mutation
13//!
14//! Multiple before-hooks execute in declaration order. The first abort short-circuits remaining
15//! hooks.
16//!
17//! **Timeout**: Defaults to 500ms (shorter than general function timeout of 5s)
18//! because before-hooks are on the critical mutation path.
19
20use std::collections::HashMap;
21
22use serde::{Deserialize, Serialize};
23
24use crate::types::EventPayload;
25
26/// Types of mutations that can trigger events.
27#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
28#[non_exhaustive]
29pub enum EventKind {
30 /// Insert operation.
31 Insert,
32 /// Update operation.
33 Update,
34 /// Delete operation.
35 Delete,
36}
37
38impl EventKind {
39 /// Convert to string representation.
40 #[must_use]
41 pub const fn as_str(&self) -> &'static str {
42 match self {
43 EventKind::Insert => "insert",
44 EventKind::Update => "update",
45 EventKind::Delete => "delete",
46 }
47 }
48}
49
50impl std::fmt::Display for EventKind {
51 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
52 f.write_str(self.as_str())
53 }
54}
55
56/// Entity event with old and new row data.
57///
58/// Represents a mutation event from the database. Used by the observer pipeline
59/// to dispatch to `after:mutation` triggers asynchronously.
60///
61/// # Dispatch Semantics
62///
63/// - Fire after mutation completes (mutation response already sent)
64/// - Async dispatch: doesn't block mutation response
65/// - Failure doesn't affect mutation (error logged only)
66/// - Execution order: in declaration order from schema
67#[derive(Debug, Clone)]
68pub struct EntityEvent {
69 /// Entity type (e.g., "User", "Post").
70 pub entity: String,
71 /// Kind of mutation.
72 pub event_kind: EventKind,
73 /// Old row data (None for Insert).
74 pub old: Option<serde_json::Value>,
75 /// New row data (None for Delete).
76 pub new: Option<serde_json::Value>,
77 /// Timestamp of the event.
78 pub timestamp: chrono::DateTime<chrono::Utc>,
79}
80
81/// A single declarative field/transition predicate on an after:mutation trigger
82/// (#597).
83///
84/// The `when` condition that decides whether a function fires, evaluated by the
85/// dispatcher against the built payload **before** any runtime spins.
86///
87/// Deliberately small: a `field` plus exactly one operator — `eq` (state) or
88/// `changed_to` (transition). Anything richer stays guest code; this is a dispatch
89/// filter, not a rules engine. A list of predicates is a **conjunction** (all must
90/// hold); an empty list always fires (back-compat).
91///
92/// ```jsonc
93/// { "field": "status", "changed_to": "approved" } // UPDATE-only transition test
94/// { "field": "kind", "eq": "standard" } // state test (INSERT + UPDATE)
95/// ```
96#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
97#[serde(deny_unknown_fields)]
98pub struct TriggerPredicate {
99 /// The field (a JSON key in the row image) to test.
100 pub field: String,
101
102 /// **State test**: the field currently equals this value. Evaluated on the
103 /// after-image (INSERT/UPDATE) or, for a DELETE, the pre-image. An absent field
104 /// never equals a value (missing ⇒ `false`).
105 #[serde(default, skip_serializing_if = "Option::is_none")]
106 pub eq: Option<serde_json::Value>,
107
108 /// **Transition test** (UPDATE-only): the field *changed to* this value —
109 /// `old.field != v && new.field == v`. A DELETE (no after-image) never matches;
110 /// `changed_to` on a non-`update` trigger is a load error.
111 #[serde(default, skip_serializing_if = "Option::is_none")]
112 pub changed_to: Option<serde_json::Value>,
113}
114
115impl TriggerPredicate {
116 /// Evaluate this predicate against a row's `old`/`new` images.
117 ///
118 /// - `eq`: the current image (`new`, or `old` on a DELETE) has `field == v`. Missing field ⇒
119 /// `false`.
120 /// - `changed_to`: `old.field != v && new.field == v`.
121 ///
122 /// A predicate with neither operator set (rejected at load) never matches.
123 #[must_use]
124 pub fn matches(
125 &self,
126 old: Option<&serde_json::Value>,
127 new: Option<&serde_json::Value>,
128 ) -> bool {
129 if let Some(value) = &self.eq {
130 // Prefer the after-image; fall back to the pre-image for a DELETE.
131 let image = new.or(old);
132 return image.and_then(|row| row.get(&self.field)) == Some(value);
133 }
134 if let Some(value) = &self.changed_to {
135 let before = old.and_then(|row| row.get(&self.field));
136 let after = new.and_then(|row| row.get(&self.field));
137 return before != Some(value) && after == Some(value);
138 }
139 false
140 }
141
142 /// Validate the predicate at load time against the trigger's `operation`
143 /// (`insert`/`update`/`delete`, or `None` for all).
144 ///
145 /// # Errors
146 ///
147 /// - Neither `eq` nor `changed_to` set, or both set (exactly one operator).
148 /// - `changed_to` on a non-`update` trigger (a transition needs a before + after).
149 pub fn validate(&self, operation: Option<&str>) -> Result<(), String> {
150 match (&self.eq, &self.changed_to) {
151 (Some(_), Some(_)) => Err(format!(
152 "predicate on field `{}` sets both `eq` and `changed_to` — use exactly one",
153 self.field
154 )),
155 (None, None) => Err(format!(
156 "predicate on field `{}` sets neither `eq` nor `changed_to`",
157 self.field
158 )),
159 (None, Some(_)) if operation != Some("update") => Err(format!(
160 "predicate on field `{}` uses `changed_to`, which is UPDATE-only, but the \
161 trigger operation is `{}`",
162 self.field,
163 operation.unwrap_or("all")
164 )),
165 _ => Ok(()),
166 }
167 }
168}
169
170/// Whether *all* predicates in a conjunction hold for a row's `old`/`new` images.
171/// An empty conjunction always holds (back-compat: no `when` ⇒ always fire).
172#[must_use]
173pub fn predicates_match(
174 predicates: &[TriggerPredicate],
175 old: Option<&serde_json::Value>,
176 new: Option<&serde_json::Value>,
177) -> bool {
178 predicates.iter().all(|predicate| predicate.matches(old, new))
179}
180
181/// Trigger that fires after a mutation completes.
182///
183/// When a mutation completes, the observer pipeline emits an `EntityEvent`.
184/// If an `AfterMutationTrigger` matches the entity type and event kind,
185/// the corresponding function is invoked asynchronously without blocking
186/// the mutation response.
187///
188/// # Matching
189///
190/// - Must match `entity_type` exactly
191/// - If `event_filter` is `None`, matches all event kinds (Insert/Update/Delete)
192/// - If `event_filter` is `Some`, matches only that specific event kind
193/// - `predicates` (the `when` clause) must all hold on the row images (#597)
194///
195/// # Dispatch
196///
197/// - Invoked in declaration order from `schema.compiled.json`
198/// - Spawned as an async task (mutation response returns immediately)
199/// - Function execution timeout: 5s default (can be overridden per function)
200/// - Failure doesn't affect mutation (error logged to tracing subscriber)
201#[derive(Debug, Clone)]
202pub struct AfterMutationTrigger {
203 /// Name of the function to invoke.
204 pub function_name: String,
205 /// Entity type to trigger on (e.g., "User").
206 pub entity_type: String,
207 /// Optional filter on event kind (None = all).
208 pub event_filter: Option<EventKind>,
209 /// The `when` conjunction (#597); empty ⇒ always fire (back-compat).
210 pub predicates: Vec<TriggerPredicate>,
211}
212
213impl AfterMutationTrigger {
214 /// Check if this trigger matches the given entity and event kind. Does **not**
215 /// evaluate the `when` predicates — the dispatcher applies
216 /// [`predicates_hold`](Self::predicates_hold) against the payload afterwards.
217 #[must_use]
218 pub fn matches(&self, entity: &str, event_kind: EventKind) -> bool {
219 self.entity_type == entity && self.event_filter.is_none_or(|filter| filter == event_kind)
220 }
221
222 /// Whether this trigger's `when` predicates all hold for the event's row images
223 /// (#597). Evaluated by the dispatcher before spawning the runtime — a `false`
224 /// result means the function does not fire (no dispatch record at all).
225 #[must_use]
226 pub fn predicates_hold(&self, event: &EntityEvent) -> bool {
227 predicates_match(&self.predicates, event.old.as_ref(), event.new.as_ref())
228 }
229
230 /// Build an `EventPayload` from an entity event.
231 #[must_use]
232 pub fn build_payload(&self, event: &EntityEvent) -> EventPayload {
233 EventPayload {
234 trigger_type: format!("after:mutation:{}", self.function_name),
235 entity: event.entity.clone(),
236 event_kind: event.event_kind.to_string(),
237 data: serde_json::json!({
238 "event_kind": event.event_kind.as_str(),
239 "old": event.old,
240 "new": event.new,
241 }),
242 timestamp: event.timestamp,
243 }
244 }
245}
246
247/// Result of a before-mutation trigger execution.
248///
249/// # Semantics
250///
251/// - `Proceed`: Allows mutation to continue with the provided input
252/// - Input may be modified from original
253/// - Passed to next trigger in chain (if any)
254/// - `Abort`: Prevents mutation from executing
255/// - Returns error to client immediately
256/// - Short-circuits remaining triggers in chain
257/// - Side-effects from aborted triggers are NOT rolled back
258///
259/// # Important: Side-Effects Not Rolled Back
260///
261/// If a `before:mutation` trigger abort is triggered, any side-effects
262/// (HTTP calls, storage writes, logs) from earlier triggers in the chain
263/// are NOT rolled back. Only the mutation itself is prevented.
264///
265/// This is by design: function side-effects are intended to be independent
266/// of mutation success. For example, if a function logs an audit entry and
267/// then a later trigger aborts, the audit entry remains.
268#[derive(Debug, Clone, Serialize, Deserialize)]
269#[non_exhaustive]
270pub enum BeforeMutationResult {
271 /// Proceed with the mutation using the provided (possibly modified) input.
272 Proceed(serde_json::Value),
273 /// Abort the mutation with an error message.
274 Abort(String),
275}
276
277/// Trigger that fires before a mutation executes.
278#[derive(Debug, Clone)]
279pub struct BeforeMutationTrigger {
280 /// Name of the function to invoke.
281 pub function_name: String,
282 /// Name of the mutation to trigger on (e.g., "createUser").
283 pub mutation_name: String,
284}
285
286impl BeforeMutationTrigger {
287 /// Check if this trigger matches the given mutation.
288 #[must_use]
289 pub fn matches(&self, mutation: &str) -> bool {
290 self.mutation_name == mutation
291 }
292}
293
294/// Chain of before-mutation triggers for a single mutation.
295///
296/// Executes multiple `before:mutation` triggers in declaration order.
297/// Each trigger can modify the input and pass it to the next trigger,
298/// or abort the mutation by returning an error.
299///
300/// # Execution Semantics
301///
302/// - Synchronous: blocks the mutation (execution is on the hot path)
303/// - Sequential: triggers execute in declaration order
304/// - Propagating: each trigger receives the modified input from previous trigger
305/// - Short-circuit: first abort stops the chain immediately
306/// - Default timeout: 500ms per trigger (shorter than general 5s default)
307/// - Side-effects: any side-effects from aborted triggers are NOT rolled back
308///
309/// # Example
310///
311/// ```ignore
312/// let chain = BeforeMutationChain {
313/// triggers: vec![
314/// validateInput, // checks required fields
315/// checkDuplicates, // checks uniqueness
316/// auditLog, // logs the attempt
317/// ]
318/// };
319///
320/// let result = chain.execute(input, &observer).await?;
321/// match result {
322/// Proceed(modified) => { /* mutation continues */ }
323/// Abort(error) => { /* mutation cancelled */ }
324/// }
325/// ```
326#[derive(Debug, Clone)]
327pub struct BeforeMutationChain {
328 /// Triggers in declaration order.
329 pub triggers: Vec<BeforeMutationTrigger>,
330}
331
332impl BeforeMutationChain {
333 /// Execute the before-mutation chain with the given input.
334 ///
335 /// Runs all triggers in declaration order. Each trigger receives the
336 /// (possibly modified) output of the previous trigger as its input.
337 /// The first `Abort` short-circuits the chain.
338 ///
339 /// # Convention for function return values
340 ///
341 /// Functions signal their intent via the returned JSON object:
342 /// - `{"abort": "message"}` → abort the mutation with `message`
343 /// - `{"input": {...}}` → proceed with modified input
344 /// - Any other value (or `null`) → proceed with the input unchanged
345 ///
346 /// # Errors
347 ///
348 /// Returns `Err` if a trigger's function name is not found in `modules`, or if
349 /// function execution itself returns an error.
350 pub async fn execute<H>(
351 &self,
352 input: serde_json::Value,
353 modules: &std::collections::HashMap<String, crate::types::FunctionModule>,
354 observer: &crate::observer::FunctionObserver,
355 host: &H,
356 limits: crate::types::ResourceLimits,
357 ) -> fraiseql_error::Result<BeforeMutationResult>
358 where
359 H: crate::HostContext + ?Sized,
360 {
361 let mut current = input;
362 for trigger in &self.triggers {
363 let module = modules.get(&trigger.function_name).ok_or_else(|| {
364 fraiseql_error::FraiseQLError::Validation {
365 message: format!(
366 "before:mutation function '{}' not found in module registry",
367 trigger.function_name,
368 ),
369 path: None,
370 }
371 })?;
372
373 let payload = crate::types::EventPayload {
374 trigger_type: format!("before:mutation:{}", trigger.mutation_name),
375 entity: trigger.mutation_name.clone(),
376 event_kind: "before".to_string(),
377 data: current.clone(),
378 timestamp: chrono::Utc::now(),
379 };
380
381 let result = observer.invoke(module, payload, host, limits.clone()).await?;
382
383 match result.value {
384 Some(ref v) if v.get("abort").is_some() => {
385 let msg = v["abort"]
386 .as_str()
387 .unwrap_or("Aborted by before:mutation trigger")
388 .to_string();
389 return Ok(BeforeMutationResult::Abort(msg));
390 },
391 Some(ref v) if v.get("input").is_some() => {
392 current = v["input"].clone();
393 },
394 _ => {},
395 }
396 }
397 Ok(BeforeMutationResult::Proceed(current))
398 }
399}
400
401/// Matcher for efficiently finding triggers by (`entity_type`, `event_kind`).
402///
403/// Uses a nested `HashMap` for O(1) lookup:
404/// - `entity_type` → `event_kind` → `Vec<AfterMutationTrigger>`
405/// - When `event_kind` is None (matches all), stored separately for fallback
406///
407/// # Integration with `FunctionObserver`
408///
409/// When the `FunctionObserver` receives an `EntityEvent` from the mutation pipeline,
410/// it calls `find()` to get all matching `AfterMutationTrigger`s. For each matching
411/// trigger, the observer spawns an async task to invoke the function without blocking
412/// the mutation response. Task completion is tracked to prevent leaks on shutdown.
413///
414/// # Example
415///
416/// ```ignore
417/// let mut matcher = TriggerMatcher::new();
418/// matcher.add(AfterMutationTrigger {
419/// function_name: "onUserCreated".to_string(),
420/// entity_type: "User".to_string(),
421/// event_filter: Some(EventKind::Insert),
422/// });
423///
424/// // Later, when a User insert occurs:
425/// let triggers = matcher.find("User", EventKind::Insert);
426/// for trigger in triggers {
427/// // Spawn async task to invoke function
428/// }
429/// ```
430#[derive(Debug, Clone)]
431pub struct TriggerMatcher {
432 /// Map of `entity_type` → `event_kind` → triggers
433 specific: HashMap<String, HashMap<String, Vec<AfterMutationTrigger>>>,
434 /// Map of `entity_type` → triggers that match all event kinds
435 all_kinds: HashMap<String, Vec<AfterMutationTrigger>>,
436}
437
438impl TriggerMatcher {
439 /// Create a new empty trigger matcher.
440 #[must_use]
441 pub fn new() -> Self {
442 Self {
443 specific: HashMap::new(),
444 all_kinds: HashMap::new(),
445 }
446 }
447
448 /// Add a trigger to the matcher.
449 pub fn add(&mut self, trigger: AfterMutationTrigger) {
450 match trigger.event_filter {
451 Some(event_kind) => {
452 self.specific
453 .entry(trigger.entity_type.clone())
454 .or_default()
455 .entry(event_kind.as_str().to_string())
456 .or_default()
457 .push(trigger);
458 },
459 None => {
460 self.all_kinds.entry(trigger.entity_type.clone()).or_default().push(trigger);
461 },
462 }
463 }
464
465 /// Find all triggers matching the given entity and event kind.
466 #[must_use]
467 pub fn find(&self, entity: &str, event_kind: EventKind) -> Vec<AfterMutationTrigger> {
468 let event_str = event_kind.as_str();
469 let mut result = Vec::new();
470
471 // Get specific triggers for this event kind
472 if let Some(entity_map) = self.specific.get(entity) {
473 if let Some(triggers) = entity_map.get(event_str) {
474 result.extend(triggers.clone());
475 }
476 }
477
478 // Get all-kinds triggers for this entity
479 if let Some(triggers) = self.all_kinds.get(entity) {
480 result.extend(triggers.clone());
481 }
482
483 result
484 }
485}
486
487impl Default for TriggerMatcher {
488 fn default() -> Self {
489 Self::new()
490 }
491}
492
493#[cfg(test)]
494mod tests;