fraiseql-functions 2.15.0

Serverless functions runtime for FraiseQL — WASM and Deno backends
Documentation
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
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
//! Core types for function execution.

use std::{path::PathBuf, time::Duration};

use serde::{Deserialize, Serialize};
use sha2::Digest;

/// Marker a function (or a host op) puts in a thrown error's message to signal the
/// failure is **permanent** — do not retry, dead-letter immediately.
///
/// Both runtimes surface a guest failure as a string, so permanence travels as this
/// sentinel substring: when the error message contains it, the runtime classifies
/// the failure as a client error (4xx) — which the durable dispatcher dead-letters
/// on the first attempt — rather than the default transient `Unsupported` (501).
///
/// A guest can also throw `Object.assign(new Error(msg), { fraiseqlPermanent: true })`
/// (Deno); the wrapper folds that into this marker. Host ops that already know a
/// failure is permanent (e.g. `send_email` on a denied identity or a rejected
/// recipient) prepend it automatically.
pub const PERMANENT_ERROR_MARKER: &str = "[fraiseql:permanent]";

/// Supported runtime types for serverless functions.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[non_exhaustive]
pub enum RuntimeType {
    /// `WebAssembly` Component Model runtime.
    Wasm,
    /// Deno (`JavaScript`/`TypeScript` via V8) runtime.
    Deno,
}

impl RuntimeType {
    /// Get supported file extensions for this runtime.
    #[must_use]
    pub const fn supported_extensions(&self) -> &[&str] {
        match self {
            RuntimeType::Wasm => &[".wasm"],
            RuntimeType::Deno => &[".js", ".ts", ".mjs", ".mts"],
        }
    }
}

/// A compiled function module ready for execution.
#[derive(Debug, Clone)]
pub struct FunctionModule {
    /// Unique name for this function.
    pub name:        String,
    /// Hash of the module source (for caching).
    pub source_hash: String,
    /// Compiled bytecode or source text.
    pub bytecode:    bytes::Bytes,
    /// Which runtime executes this module.
    pub runtime:     RuntimeType,
}

impl FunctionModule {
    /// Create a new WASM module from compiled bytecode.
    pub fn from_bytecode(name: String, bytecode: bytes::Bytes) -> Self {
        let source_hash = hex::encode(sha2::Sha256::digest(&bytecode));
        Self {
            name,
            source_hash,
            bytecode,
            runtime: RuntimeType::Wasm,
        }
    }

    /// Create a new source-based module (JavaScript/TypeScript).
    #[must_use]
    pub fn from_source(name: String, source: String, runtime: RuntimeType) -> Self {
        let bytecode = bytes::Bytes::from(source);
        let source_hash = hex::encode(sha2::Sha256::digest(&bytecode));
        Self {
            name,
            source_hash,
            bytecode,
            runtime,
        }
    }
}

/// Trigger event payload for a function invocation.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct EventPayload {
    /// Type of trigger: "mutation", "subscription", "cron", "webhook", etc.
    pub trigger_type: String,
    /// Entity name (e.g., "User", "Post").
    pub entity:       String,
    /// Event kind (e.g., "created", "updated", "deleted").
    pub event_kind:   String,
    /// Event data (JSON).
    pub data:         serde_json::Value,
    /// Timestamp when the event occurred.
    pub timestamp:    chrono::DateTime<chrono::Utc>,
}

/// The least-privilege authority ceiling a function's `fraiseql_query` bridge
/// writes run under (#594) — the function's `run_as`.
///
/// This is the same authority model scheduled sources use
/// (`fraiseql_core::schema::RunAs`, `docs/architecture/sources.md:88-109`), applied
/// to event-dispatched functions: a *ceiling* the function can never exceed. A
/// [`FunctionDefinition`] with no `run_as` runs **fail-closed** — its host's
/// `fraiseql_query` executes under an anonymous `SecurityContext::system_job` identity with no
/// roles/scopes/tenant, so RLS and field-authorization deny every write until an
/// operator grants a ceiling. Granting authority is a deliberate act, never a
/// default.
///
/// It is a distinct type from the core `RunAs` so the base `fraiseql-functions`
/// crate (used by the CLI, codegen, and authoring) need not depend on
/// `fraiseql-core`; the two share an identical JSON shape and the wiring layer
/// (`host-live`) converts to a `SecurityContext` via `FunctionDefinition::identity`.
///
/// Those three items live behind the `host-live` feature, which is what pulls in
/// `fraiseql-core`; they are named rather than linked so this page resolves in the
/// default build, which is the one a bare `cargo doc` produces (#1199).
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct RunAs {
    /// Roles granted to the function's background write identity (the RBAC ceiling).
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub roles: Vec<String>,

    /// Scopes granted to the function's background write identity.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub scopes: Vec<String>,

    /// The single tenant this function's bridge writes are scoped to, if any. Unset
    /// ⇒ global/system (NULL tenant).
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub tenant: Option<String>,
}

#[cfg(feature = "host-live")]
impl RunAs {
    /// Build the background [`SecurityContext`](fraiseql_core::security::SecurityContext)
    /// this ceiling grants, identified as `system_job:<job_id>` and correlated by
    /// `request_id`. Mirrors
    /// [`SourceDefinition::identity`](fraiseql_core::schema::SourceDefinition::identity):
    /// the roles/scopes/tenant are the ceiling, the [`ActorType::SystemJob`] is
    /// recorded for audit, never an authorization input.
    ///
    /// [`ActorType::SystemJob`]: fraiseql_core::security::ActorType::SystemJob
    #[must_use]
    pub fn identity(
        &self,
        job_id: impl Into<String>,
        request_id: impl Into<String>,
    ) -> fraiseql_core::security::SecurityContext {
        fraiseql_core::security::SecurityContext::system_job(
            job_id,
            request_id,
            self.roles.clone(),
            self.scopes.clone(),
            self.tenant.clone().map(fraiseql_core::types::TenantId::from),
        )
    }
}

/// Definition of a serverless function for deployment and execution.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct FunctionDefinition {
    /// Unique name for this function.
    pub name:       String,
    /// Trigger type and configuration (e.g., "after:mutation:createUser", "cron:0 * * * *",
    /// "<http:GET:/users/:id>").
    pub trigger:    String,
    /// Which runtime executes this function.
    pub runtime:    RuntimeType,
    /// Optional timeout in milliseconds (overrides defaults).
    /// - For `before:mutation` triggers: defaults to 500ms
    /// - For other triggers: defaults to 5s
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub timeout_ms: Option<u64>,

    /// The authority ceiling this function's `fraiseql_query` bridge writes run
    /// under (#594). Absent ⇒ **fail-closed**: the bridge runs under an anonymous
    /// identity and RLS/field-authz deny writes. See [`RunAs`] and
    /// `identity()` (feature `host-live`).
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub run_as: Option<RunAs>,

    /// Declarative `when` predicates (#597) — a conjunction the dispatcher evaluates
    /// on the row images before firing an `after:mutation`/`after:capture` function.
    /// Empty ⇒ always fire (back-compat). See
    /// [`TriggerPredicate`](crate::triggers::mutation::TriggerPredicate).
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub when: Vec<crate::triggers::mutation::TriggerPredicate>,

    /// Fire-and-forget opt-out for durable dispatch.
    ///
    /// After-mutation function dispatch is durable by default: a transient
    /// failure is retried with backoff and, once retries are exhausted, the
    /// invocation is dead-lettered so money- and send-path work is never
    /// silently lost. Set `re_runnable = true` for work that is safe to simply
    /// re-run later (e.g. LLM scoring) — such dispatch stays fire-and-forget with
    /// no retry or dead-letter overhead. See ADR 0015 for the rationale.
    #[serde(default)]
    pub re_runnable: bool,

    /// Per-function retry policy for durable dispatch.
    ///
    /// `None` uses the server default (overridable via `FRAISEQL_FUNCTIONS_RETRY_*`
    /// environment variables). Ignored when [`re_runnable`](Self::re_runnable) is
    /// `true`. Reuses the observer subsystem's [`RetryConfig`] so retry semantics
    /// are identical across both subsystems.
    ///
    /// [`RetryConfig`]: fraiseql_observers::RetryConfig
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub retry: Option<fraiseql_observers::RetryConfig>,
}

impl FunctionDefinition {
    /// Create a new function definition.
    #[must_use]
    pub fn new(name: &str, trigger: &str, runtime: RuntimeType) -> Self {
        Self {
            name: name.to_string(),
            trigger: trigger.to_string(),
            runtime,
            timeout_ms: None,
            run_as: None,
            when: Vec::new(),
            re_runnable: false,
            retry: None,
        }
    }

    /// Set the [`run_as`](Self::run_as) authority ceiling (#594).
    #[must_use]
    pub fn with_run_as(mut self, run_as: RunAs) -> Self {
        self.run_as = Some(run_as);
        self
    }

    /// The background [`SecurityContext`] this function's `fraiseql_query` bridge
    /// writes run under (#594).
    ///
    /// Built from [`run_as`](Self::run_as) via [`SecurityContext::system_job`], so a
    /// function write is audited as `system_job:<function-name>` under
    /// [`ActorType::SystemJob`] — the same envelope a source write carries. **Absent
    /// `run_as` yields a fail-closed identity** (no roles, no scopes, no tenant →
    /// every authz/RLS decision denies). `request_id` correlates one dispatch
    /// (typically its per-dispatch idempotency token).
    ///
    /// [`SecurityContext`]: fraiseql_core::security::SecurityContext
    /// [`SecurityContext::system_job`]: fraiseql_core::security::SecurityContext::system_job
    /// [`ActorType::SystemJob`]: fraiseql_core::security::ActorType::SystemJob
    #[cfg(feature = "host-live")]
    #[must_use]
    pub fn identity(
        &self,
        request_id: impl Into<String>,
    ) -> fraiseql_core::security::SecurityContext {
        self.run_as.clone().unwrap_or_default().identity(&self.name, request_id)
    }

    /// Mark this function as re-runnable (fire-and-forget) dispatch.
    ///
    /// See [`re_runnable`](Self::re_runnable) and ADR 0015.
    #[must_use]
    pub const fn re_runnable(mut self) -> Self {
        self.re_runnable = true;
        self
    }

    /// Set a custom timeout for this function.
    #[must_use]
    pub const fn with_timeout(mut self, timeout_ms: u64) -> Self {
        self.timeout_ms = Some(timeout_ms);
        self
    }

    /// Get the effective timeout for this function.
    #[must_use]
    pub fn effective_timeout(&self) -> Duration {
        match self.timeout_ms {
            Some(ms) => Duration::from_millis(ms),
            None => {
                // before:mutation defaults to 500ms; others default to 5s
                if self.trigger.starts_with("before:mutation") {
                    Duration::from_millis(500)
                } else {
                    Duration::from_secs(5)
                }
            },
        }
    }

    /// The module file this definition resolves to under `module_dir`, if one is
    /// there: `<module_dir>/<name>.<ext>` for each extension the declared runtime
    /// supports, first existing file wins.
    ///
    /// The **one** definition of where a function's code lives. The server's
    /// `build_functions_subsystem` loads through it, `fraiseql functions invoke`
    /// resolves through it, and the compiler checks through it (#1325) — three sites
    /// that each carried their own copy of this loop, which is three chances for the
    /// compiler to approve a layout the server cannot load.
    ///
    /// Returns `None` when no module is present, which the caller reports in its own
    /// terms: a compile error, a boot failure, or a harness error.
    #[must_use]
    pub fn resolve_module_path(&self, module_dir: &std::path::Path) -> Option<PathBuf> {
        self.runtime
            .supported_extensions()
            .iter()
            .map(|extension| module_dir.join(format!("{}{extension}", self.name)))
            .find(|path| path.exists())
    }

    /// The module-file candidates this definition names, rendered for a diagnostic:
    /// `<module_dir>/<name>.{ext,ext}`.
    #[must_use]
    pub fn module_path_pattern(&self, module_dir: &std::path::Path) -> String {
        format!(
            "{}/{}.{{{}}}",
            module_dir.display(),
            self.name,
            self.runtime
                .supported_extensions()
                .iter()
                .map(|ext| ext.trim_start_matches('.'))
                .collect::<Vec<_>>()
                .join(","),
        )
    }

    /// Check if this function is a before:mutation trigger.
    #[must_use]
    pub fn is_before_mutation(&self) -> bool {
        self.trigger.starts_with("before:mutation:")
    }

    /// Check if this function is an after:mutation trigger.
    #[must_use]
    pub fn is_after_mutation(&self) -> bool {
        self.trigger.starts_with("after:mutation:")
    }

    /// Check if this function is an after:storage trigger.
    #[must_use]
    pub fn is_after_storage(&self) -> bool {
        self.trigger.starts_with("after:storage:")
    }

    /// Check if this function is an after:ingest trigger.
    #[must_use]
    pub fn is_after_ingest(&self) -> bool {
        self.trigger == "after:ingest" || self.trigger.starts_with("after:ingest:")
    }

    /// Check if this function is a cron trigger.
    #[must_use]
    pub fn is_cron(&self) -> bool {
        self.trigger.starts_with("cron:")
    }

    /// Check if this function is an HTTP trigger.
    #[must_use]
    pub fn is_http(&self) -> bool {
        self.trigger.starts_with("http:")
    }
}

/// The compiled schema's `"functions"` section — the whole of what an author
/// declares about functions, and the single definition of that shape.
///
/// It is a **sibling** of the core `CompiledSchema` in `schema.compiled.json`, not a
/// field of it: `fraiseql-core` knows nothing about functions, and the dependency
/// runs the other way (`fraiseql-functions` optionally uses core, never the
/// reverse). The server reads it through `ExtendedCompiledSchema`; the CLI writes it
/// there and reads it back in `fraiseql functions invoke`.
///
/// ```json
/// {
///   "functions": {
///     "module_dir": "functions",
///     "definitions": [
///       { "name": "on_create_user", "trigger": "after:mutation:createUser", "runtime": "Wasm" }
///     ]
///   }
/// }
/// ```
///
/// # One definition, deliberately
///
/// This shape used to exist three times: `fraiseql-server`'s `FunctionsConfig`, the
/// CLI harness's private `FunctionsSection` ("a local mirror … so the harness does
/// not pull the whole server crate in"), and — had #1325 followed the same instinct —
/// an authoring mirror in `IntermediateSchema`. Three copies of a rule lose the
/// fourth. It lives here, in the crate that owns [`FunctionDefinition`], so the
/// producer and both consumers deserialize the same struct.
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct FunctionsConfig {
    /// Directory containing the function modules (`.wasm`, `.js`, `.ts`), resolved
    /// relative to the server's working directory.
    ///
    /// Required on the wire. The *authoring* default (`functions/`) lives in the
    /// compiler, which always emits an explicit value — so a hand-written compiled
    /// schema that omits it still fails loudly rather than silently resolving
    /// somewhere the author did not choose.
    pub module_dir: PathBuf,

    /// The declared functions.
    #[serde(default)]
    pub definitions: Vec<FunctionDefinition>,

    /// Which dead-letter store backs function dispatch (#598): `"memory"` (the
    /// default — dead-letters vanish on restart) or `"postgres"` (durable, survives
    /// a restart; requires a database pool). Overridable by the
    /// `FRAISEQL_FUNCTIONS_DLQ_STORE` env var. Absent ⇒ memory.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub dlq_store: Option<String>,
}

/// Log level for structured logging.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub enum LogLevel {
    /// Debug level.
    Debug,
    /// Info level.
    Info,
    /// Warning level.
    Warn,
    /// Error level.
    Error,
}

/// A single log entry from function execution.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct LogEntry {
    /// Log level.
    pub level:     LogLevel,
    /// Log message.
    pub message:   String,
    /// When the log was written.
    pub timestamp: chrono::DateTime<chrono::Utc>,
}

/// Result of a function invocation.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct FunctionResult {
    /// Return value from the function (may be None if function returns void).
    pub value:             Option<serde_json::Value>,
    /// All logs captured during execution.
    pub logs:              Vec<LogEntry>,
    /// Total execution duration.
    pub duration:          Duration,
    /// Peak memory usage in bytes.
    pub memory_peak_bytes: u64,
}

/// Resource limits for function execution.
#[derive(Debug, Clone)]
pub struct ResourceLimits {
    /// Maximum memory allocation in bytes.
    pub max_memory_bytes: u64,
    /// Maximum execution duration.
    pub max_duration:     Duration,
    /// Maximum number of log entries to capture.
    pub max_log_entries:  usize,
}

impl Default for ResourceLimits {
    fn default() -> Self {
        Self {
            max_memory_bytes: 128 * 1024 * 1024,      // 128 MB
            max_duration:     Duration::from_secs(5), // 5 seconds
            max_log_entries:  10_000,
        }
    }
}

#[cfg(test)]
mod tests;