Skip to main content

fraiseql_functions/observer/
mod.rs

1//! Function observer for executing functions in response to events.
2
3#[cfg(test)]
4mod tests;
5
6use std::{collections::HashMap, sync::Arc};
7
8use fraiseql_error::Result;
9
10use crate::{
11    HostContext,
12    runtime::FunctionRuntime,
13    types::{EventPayload, FunctionModule, FunctionResult, ResourceLimits, RuntimeType},
14};
15
16/// Executes functions in response to trigger events.
17///
18/// This observer integrates with the fraiseql-observers action execution pipeline.
19/// It receives trigger events, looks up the corresponding function module,
20/// selects the appropriate runtime, and executes the function.
21pub struct FunctionObserver {
22    runtimes: HashMap<RuntimeType, Arc<dyn std::any::Any + Send + Sync>>,
23}
24
25impl FunctionObserver {
26    /// Create a new function observer.
27    #[must_use]
28    pub fn new() -> Self {
29        Self {
30            runtimes: HashMap::new(),
31        }
32    }
33
34    /// Register a runtime for a specific runtime type.
35    ///
36    /// # Errors
37    ///
38    /// Returns `Err` if runtime registration fails.
39    pub fn register_runtime<R: FunctionRuntime + 'static>(
40        &mut self,
41        runtime_type: RuntimeType,
42        runtime: R,
43    ) {
44        self.runtimes.insert(runtime_type, Arc::new(runtime));
45    }
46
47    /// Execute a function module in response to an event.
48    ///
49    /// Dispatches to the appropriate runtime based on the module's `runtime` field.
50    ///
51    /// # Errors
52    ///
53    /// Returns `Err` if:
54    /// - No runtime is registered for the module's runtime type
55    /// - The runtime fails to execute the module
56    pub async fn invoke<H>(
57        &self,
58        module: &FunctionModule,
59        #[allow(unused_variables)] // Reason: used only when runtime features are enabled
60        event: EventPayload,
61        #[allow(unused_variables)] // Reason: used only when runtime features are enabled
62        host: &H,
63        #[allow(unused_variables)] // Reason: used only when runtime features are enabled
64        limits: ResourceLimits,
65    ) -> Result<FunctionResult>
66    where
67        H: HostContext + ?Sized,
68    {
69        #[allow(unused_variables)] // Reason: used only when runtime features are enabled
70        let runtime_box = self.runtimes.get(&module.runtime).ok_or_else(|| {
71            fraiseql_error::FraiseQLError::Unsupported {
72                message: format!("No runtime registered for {:?}", module.runtime),
73            }
74        })?;
75
76        // Dispatch based on runtime type
77        #[allow(unreachable_patterns)] // Reason: pattern reachability depends on features
78        match module.runtime {
79            RuntimeType::Wasm => {
80                #[cfg(feature = "runtime-wasm")]
81                {
82                    let runtime = runtime_box
83                        .downcast_ref::<crate::runtime::wasm::WasmRuntime>()
84                        .ok_or_else(|| fraiseql_error::FraiseQLError::Unsupported {
85                        message: "Invalid WASM runtime".to_string(),
86                    })?;
87                    runtime.invoke(module, event, host, limits).await
88                }
89                #[cfg(not(feature = "runtime-wasm"))]
90                {
91                    Err(fraiseql_error::FraiseQLError::Unsupported {
92                        message: "WASM runtime not enabled".to_string(),
93                    })
94                }
95            },
96            RuntimeType::Deno => {
97                #[cfg(feature = "runtime-deno")]
98                {
99                    let runtime = runtime_box
100                        .downcast_ref::<crate::runtime::deno::DenoRuntime>()
101                        .ok_or_else(|| fraiseql_error::FraiseQLError::Unsupported {
102                        message: "Invalid Deno runtime".to_string(),
103                    })?;
104                    runtime.invoke(module, event, host, limits).await
105                }
106                #[cfg(not(feature = "runtime-deno"))]
107                {
108                    Err(fraiseql_error::FraiseQLError::Unsupported {
109                        message: "Deno runtime not enabled".to_string(),
110                    })
111                }
112            },
113        }
114    }
115
116    /// Find all `after:mutation` triggers that match a given entity event.
117    ///
118    /// Returns the list of matching [`AfterMutationTrigger`]s from `registry`.
119    /// Used by after-mutation dispatchers to know which functions to invoke.
120    /// This method is allocation-free when no triggers match.
121    ///
122    /// [`AfterMutationTrigger`]: crate::triggers::mutation::AfterMutationTrigger
123    #[must_use]
124    pub fn find_after_mutation_triggers(
125        &self,
126        registry: &crate::triggers::registry::TriggerRegistry,
127        event: &crate::triggers::mutation::EntityEvent,
128    ) -> Vec<crate::triggers::mutation::AfterMutationTrigger> {
129        registry.after_mutation_triggers.find(&event.entity, event.event_kind)
130    }
131
132    /// Find `after:capture` triggers (#366) matching an externally-captured change.
133    ///
134    /// The change-log reader calls this for captured rows; the match is on entity +
135    /// event kind, same as after:mutation, but against the separate capture matcher.
136    #[must_use]
137    pub fn find_after_capture_triggers(
138        &self,
139        registry: &crate::triggers::registry::TriggerRegistry,
140        event: &crate::triggers::mutation::EntityEvent,
141    ) -> Vec<crate::triggers::mutation::AfterMutationTrigger> {
142        registry.after_capture_triggers.find(&event.entity, event.event_kind)
143    }
144
145    /// Execute a function module with a full I/O-capable host context.
146    ///
147    /// Unlike [`invoke`](Self::invoke) — which snapshots the host into a
148    /// sync-only view and so returns `Unsupported` for async I/O — this routes to
149    /// the runtime's `invoke_with_context`, giving the guest live HTTP / query /
150    /// storage access through `host`. The after:mutation dispatcher uses this so
151    /// side-effecting functions (webhooks, external provisioning, LLM scoring) can
152    /// reach the network.
153    ///
154    /// Dispatches by the module's [`runtime`](FunctionModule::runtime) field:
155    /// WASM modules route to
156    /// [`WasmRuntime::invoke_with_context`](crate::runtime::wasm::WasmRuntime::invoke_with_context),
157    /// Deno (`JavaScript`/`TypeScript`) modules to
158    /// [`DenoRuntime::invoke_with_context`](crate::runtime::deno::DenoRuntime::invoke_with_context).
159    /// An event whose module targets a runtime that is not compiled in returns
160    /// `Unsupported`.
161    ///
162    /// # Errors
163    ///
164    /// Returns `Err` if no runtime is registered for the module's runtime type,
165    /// the registered runtime is of the wrong concrete type, or guest execution
166    /// fails.
167    #[cfg(any(feature = "runtime-wasm", feature = "runtime-deno"))]
168    pub async fn invoke_with_context(
169        &self,
170        module: &FunctionModule,
171        event: EventPayload,
172        #[allow(unused_variables)]
173        // Reason: used only when the matching runtime feature is enabled
174        host: Arc<dyn crate::host::dyn_context::DynHostContext>,
175        #[allow(unused_variables)]
176        // Reason: used only when the matching runtime feature is enabled
177        limits: ResourceLimits,
178    ) -> Result<FunctionResult> {
179        #[allow(unused_variables)] // Reason: used only when the matching runtime feature is enabled
180        let runtime_box = self.runtimes.get(&module.runtime).ok_or_else(|| {
181            fraiseql_error::FraiseQLError::Unsupported {
182                message: format!(
183                    "No runtime registered for {:?} on the I/O host-context dispatch path",
184                    module.runtime
185                ),
186            }
187        })?;
188
189        #[allow(unreachable_patterns)] // Reason: pattern reachability depends on features
190        match module.runtime {
191            RuntimeType::Wasm => {
192                #[cfg(feature = "runtime-wasm")]
193                {
194                    let runtime = runtime_box
195                        .downcast_ref::<crate::runtime::wasm::WasmRuntime>()
196                        .ok_or_else(|| fraiseql_error::FraiseQLError::Unsupported {
197                        message: "Invalid WASM runtime".to_string(),
198                    })?;
199                    runtime.invoke_with_context(module, event, host, limits).await
200                }
201                #[cfg(not(feature = "runtime-wasm"))]
202                {
203                    Err(fraiseql_error::FraiseQLError::Unsupported {
204                        message: "WASM runtime not enabled".to_string(),
205                    })
206                }
207            },
208            RuntimeType::Deno => {
209                #[cfg(feature = "runtime-deno")]
210                {
211                    let runtime = runtime_box
212                        .downcast_ref::<crate::runtime::deno::DenoRuntime>()
213                        .ok_or_else(|| fraiseql_error::FraiseQLError::Unsupported {
214                        message: "Invalid Deno runtime".to_string(),
215                    })?;
216                    runtime.invoke_with_context(module, event, host, limits).await
217                }
218                #[cfg(not(feature = "runtime-deno"))]
219                {
220                    Err(fraiseql_error::FraiseQLError::Unsupported {
221                        message: "Deno runtime not enabled".to_string(),
222                    })
223                }
224            },
225        }
226    }
227}
228
229impl Default for FunctionObserver {
230    fn default() -> Self {
231        Self::new()
232    }
233}