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