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}