fraiseql-functions 2.16.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
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
//! Deno executor implementation — handles V8 isolation and module execution.
//!
//! # Architecture
//!
//! Each invocation spawns a fresh OS thread with a dedicated single-threaded Tokio
//! runtime.  This sidesteps the "cannot call `block_on` inside a Tokio runtime"
//! restriction while keeping V8 + its event loop on an uncontested thread.
//!
//! The result is returned to the outer async context via a `oneshot` channel.

use std::{
    cell::RefCell,
    rc::Rc,
    sync::{
        Arc, Mutex,
        atomic::{AtomicBool, Ordering},
    },
    time::{Duration, Instant},
};

use deno_core::{Extension, JsRuntime, JsRuntimeForSnapshot, OpState, RuntimeOptions, op2, v8};
use serde_json::Value;

use super::{ops::DenoHostContext, watchdog};
use crate::{
    host::dyn_context::DynHostContext,
    types::{LogEntry, LogLevel, ResourceLimits},
};

// ── V8 platform initialization ─────────────────────────────────────────────────

/// Initialize the process-wide V8 platform as the **unprotected** default
/// platform, once, before any isolate exists.
///
/// Without this, the first `JsRuntime::new` lazily initializes `deno_core`'s
/// *protected* default platform, whose JIT pages are hardened with memory
/// protection keys (PKU). Pkey access rights are per-thread and inherited only
/// by threads **descended from the initializing thread** — and every FraiseQL
/// invocation runs on its own short-lived thread, so the lazily-bound platform
/// died with invocation 1's thread and the second Deno invocation in any
/// process faulted on the first JIT page it touched (SIGSEGV, #969). A library
/// cannot guarantee that all future invocation threads descend from one common
/// initializer, which is exactly the case `rusty_v8` documents the unprotected
/// platform for. The cost is V8's pkey JIT-page hardening (defense-in-depth
/// against V8-level exploits); ordinary W^X page protection still applies.
fn ensure_v8_platform() {
    static V8_PLATFORM: std::sync::Once = std::sync::Once::new();
    V8_PLATFORM.call_once(|| {
        JsRuntime::init_platform(Some(
            v8::new_unprotected_default_platform(0, false).make_shared(),
        ));
    });
}

// ── Startup snapshot ───────────────────────────────────────────────────────────

/// A V8 startup snapshot of a fresh fraiseql isolate, built once per process (#1343).
///
/// Every invocation gets its own isolate. Without a snapshot, each one re-runs
/// `deno_core`'s built-in `JavaScript` to build the global environment — measured at
/// ~3.9 ms of a ~4.6 ms trivial invocation (`benches/deno_invocation_bench.rs`), against
/// ~0.4 ms for everything the guest itself does. A snapshot captures that environment
/// once and each later isolate deserialises it instead.
///
/// Built at runtime rather than in a build script so the crate does not compile V8 a
/// second time as a build dependency; the first invocation in the process pays for it.
/// The snapshot is made with the same extension the invocations register, which is what
/// `deno_core` requires of a runtime started from it. The extension carries ops only, no
/// `JavaScript`, so nothing guest-visible is frozen into it, and every invocation still
/// gets a fresh isolate and its own `OpState`.
fn startup_snapshot() -> &'static [u8] {
    static SNAPSHOT: std::sync::OnceLock<&'static [u8]> = std::sync::OnceLock::new();
    SNAPSHOT.get_or_init(|| {
        let collector = LogCollector {
            logs:        Arc::new(Mutex::new(Vec::new())),
            max_entries: 0,
        };
        let runtime = JsRuntimeForSnapshot::new(RuntimeOptions {
            extensions: vec![make_fraiseql_extension(collector, None)],
            ..Default::default()
        });
        Box::leak(runtime.snapshot())
    })
}

// ── Log collector state stored in OpState ─────────────────────────────────────

/// Shared log collector threaded into the `fraiseql_log` op via `OpState`.
#[derive(Clone)]
struct LogCollector {
    logs:        Arc<Mutex<Vec<LogEntry>>>,
    max_entries: usize,
}

// ── Op definitions ─────────────────────────────────────────────────────────────

/// Log a message from Deno guest code.
///
/// Called by guests as `Deno.core.ops.fraiseql_log(level, message)`.
/// Levels: 0 = debug, 1 = info, 2 = warn, 3 = error.
#[op2(fast)]
#[allow(clippy::inline_always)] // Reason: emitted by the `#[op2]` proc-macro, not our code
#[allow(clippy::needless_pass_by_value)] // Reason: `#[op2]` requires owned `String` for `#[string]`
fn fraiseql_log(state: Rc<RefCell<OpState>>, #[smi] level: u8, #[string] message: String) {
    let state = state.borrow();
    let collector = state.borrow::<LogCollector>();
    let mut logs = collector.logs.lock().expect("log mutex poisoned");
    if logs.len() < collector.max_entries {
        let log_level = match level {
            0 => LogLevel::Debug,
            2 => LogLevel::Warn,
            3 => LogLevel::Error,
            _ => LogLevel::Info,
        };
        logs.push(LogEntry {
            level: log_level,
            message,
            timestamp: chrono::Utc::now(),
        });
    }
}

// ── Extension builder ──────────────────────────────────────────────────────────

fn make_fraiseql_extension(
    collector: LogCollector,
    host: Option<Arc<dyn DynHostContext>>,
) -> Extension {
    use super::ops;
    Extension {
        name: "fraiseql",
        ops: std::borrow::Cow::Owned(vec![
            fraiseql_log(),
            ops::fraiseql_query(),
            ops::fraiseql_sql_query(),
            ops::fraiseql_http_request(),
            ops::fraiseql_storage_get(),
            ops::fraiseql_storage_put(),
            ops::fraiseql_send_email(),
            ops::fraiseql_auth_context(),
            ops::fraiseql_env_var(),
            ops::fraiseql_idempotency_token(),
            ops::fraiseql_cursor_get(),
            ops::fraiseql_cursor_advance(),
        ]),
        op_state_fn: Some(Box::new(move |state: &mut OpState| {
            state.put(collector);
            if let Some(host) = host {
                state.put(DenoHostContext(host));
            }
        })),
        ..Default::default()
    }
}

// ── Source preprocessing ────────────────────────────────────────────────────────

/// Wrap the user-provided source so the default-exported function is called with
/// the event data and the result stored in `globalThis.__fraiseql_result`.
///
/// Handles both:
///   `export default async (event) => { ... };`
///   `export default async function(event) { ... }`
///
/// The wrapped code also catches any thrown exception and stores it in
/// `globalThis.__fraiseql_error`.
fn wrap_source(source: &str, event_json: &str) -> String {
    // Convert "export default <expr>" to "const __fn = <expr>"
    let inner = source
        .replace("export default async function", "const __fn = async function")
        .replace("export default async", "const __fn = async")
        .replace("export default function", "const __fn = function")
        .replace("export default", "const __fn =");

    // Injected so a guest can mark a thrown error permanent structurally
    // (`throw Object.assign(new Error(msg), {{ fraiseqlPermanent: true }})`) and the
    // runtime dead-letters it immediately. Host-op errors already carry the marker
    // in their message; this folds the property form into the same signal.
    let marker = crate::types::PERMANENT_ERROR_MARKER;
    format!(
        r#"
{inner}
(async () => {{
    try {{
        const __event = {event_json};
        const __result = await __fn(__event);
        globalThis.__fraiseql_result = JSON.stringify(__result);
        globalThis.__fraiseql_error = null;
    }} catch (e) {{
        globalThis.__fraiseql_result = null;
        const __permanent = e && e.fraiseqlPermanent === true;
        globalThis.__fraiseql_error = (__permanent ? "{marker} " : "") + String(e);
    }}
}})();
"#
    )
}

// ── Execution result ────────────────────────────────────────────────────────────

/// Result returned from the blocking deno thread back to the async caller.
pub struct ExecutionResult {
    /// The value returned by the function (serialised → deserialised).
    pub value:  Value,
    /// Log entries captured during execution.
    pub logs:   Vec<LogEntry>,
    /// Where the invocation's time went, phase by phase (#1343).
    pub phases: PhaseTimings,
}

/// Wall-clock time of each phase of one invocation, in the order they run.
///
/// Reported separately because a single total hides which term moved: building the
/// isolate is a fixed cost paid on every invocation, while script and event-loop time
/// belong to the guest. `benches/deno_invocation_bench.rs` reports each one.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct PhaseTimings {
    /// Building the invocation's single-threaded tokio runtime.
    pub tokio_runtime: Duration,
    /// Constructing the `JsRuntime`: the V8 isolate, its context and the fraiseql ops.
    pub isolate:       Duration,
    /// Compiling and running the wrapped guest script up to its first `await`.
    pub script:        Duration,
    /// Driving the event loop until the guest's promise settles.
    pub event_loop:    Duration,
    /// Reading the result and error globals back out of the isolate.
    pub result:        Duration,
    /// Dropping the isolate and the tokio runtime once the result is out.
    pub teardown:      Duration,
}

impl PhaseTimings {
    /// The sum of the phases.
    #[must_use]
    pub fn total(&self) -> Duration {
        self.tokio_runtime
            + self.isolate
            + self.script
            + self.event_loop
            + self.result
            + self.teardown
    }
}

// ── Core execution (runs on a dedicated thread with its own Tokio runtime) ─────

/// Execute guest `JavaScript` inside a fresh `JsRuntime` + dedicated Tokio runtime.
///
/// This function is called from inside `std::thread::spawn`, so it is free to
/// create its own `current_thread` Tokio runtime and call `block_on` without
/// conflict.
///
/// # Errors
///
/// Returns an error string on syntax errors, runtime exceptions, timeout, or
/// resource-limit violations.
///
/// # Panics
///
/// Panics if the internal log mutex is poisoned (only possible if another thread
/// panicked while holding the lock, which cannot happen in normal operation).
pub fn run_in_dedicated_thread(
    source: &str,
    event_value: &Value,
    limits: &ResourceLimits,
    host: Option<Arc<dyn DynHostContext>>,
) -> Result<ExecutionResult, String> {
    // Resource limits are enforced for real against the V8 isolate (M-deno-limits):
    //
    //   * Memory: the isolate is created with a hard heap limit (`max_memory_bytes`) and a
    //     near-heap-limit callback that terminates execution when V8 approaches that limit.
    //   * CPU/time: a watchdog thread terminates the isolate after `max_duration`, catching tight
    //     synchronous loops (`while (true) {}`) that never yield to the async event loop, in
    //     addition to the event-loop `tokio::time::timeout` guard that catches async hangs
    //     (unresolved promises).
    //
    // The previous substring heuristics (matching `while (true)` / `ArrayBuffer` in
    // the source) were removed: they false-positived on those substrings appearing
    // in comments or string literals and missed every other form of DoS.

    // The platform must exist (and must be the unprotected one, #969) before
    // the first isolate in the process is created below.
    ensure_v8_platform();

    // Shared log storage
    let logs_arc: Arc<Mutex<Vec<LogEntry>>> = Arc::new(Mutex::new(Vec::new()));
    let collector = LogCollector {
        logs:        Arc::clone(&logs_arc),
        max_entries: limits.max_log_entries,
    };

    // Serialise the event data for injection into JS
    let event_json = serde_json::to_string(event_value).map_err(|e| e.to_string())?;

    // Build the wrapper script
    let wrapped = wrap_source(source, &event_json);

    // Extract before async move to avoid partial-move into the closure.
    let max_duration = limits.max_duration;
    // V8's heap limit is expressed in `usize`; saturate on 32-bit targets where the
    // configured `u64` limit could exceed the address space.
    let max_memory_bytes = usize::try_from(limits.max_memory_bytes).unwrap_or(usize::MAX);

    // Flags shared with the heap-limit callback and the watchdog thread so that,
    // once V8 execution is terminated, we can report *why* (memory vs. timeout)
    // rather than surfacing the opaque "execution terminated" V8 error.
    let mem_exceeded_run = Arc::new(AtomicBool::new(false));
    let timed_out_run = Arc::new(AtomicBool::new(false));

    // Create a single-threaded Tokio runtime for deno's event loop.
    // This is safe because we're in a fresh OS thread with no existing Tokio context.
    let mut phases = PhaseTimings::default();
    let phase_started = Instant::now();
    let rt = tokio::runtime::Builder::new_current_thread()
        .enable_all()
        .build()
        .map_err(|e| format!("Failed to create tokio runtime: {e}"))?;
    phases.tokio_runtime = phase_started.elapsed();

    let result = rt.block_on(async {
        // Hard heap limit enforced by V8: the isolate may not grow past
        // `max_memory_bytes`. `0` initial lets V8 pick a sane starting size.
        let create_params = v8::CreateParams::default().heap_limits(0, max_memory_bytes);

        let phase_started = Instant::now();
        let mut js_runtime = JsRuntime::new(RuntimeOptions {
            extensions: vec![make_fraiseql_extension(collector, host)],
            startup_snapshot: Some(startup_snapshot()),
            create_params: Some(create_params),
            ..Default::default()
        });
        phases.isolate = phase_started.elapsed();

        // Watchdog: terminate the isolate after `max_duration`. This catches tight
        // *synchronous* loops (`while (true) {}`) that never yield to the async
        // event loop, so the `tokio::time::timeout` below cannot fire. The handle
        // is thread-safe (Send + Sync) and may be invoked from any thread.
        //
        // ONE deadline covers the whole invocation — script evaluation AND the
        // event loop. The watchdog must stay armed across `run_event_loop`
        // (#804): a guest that `await`s a real host op resumes *inside* the
        // event loop, and a synchronous spin there never yields to tokio, so
        // the `tokio::time::timeout` future is never polled and only this
        // thread can terminate the isolate.
        let invocation_deadline = std::time::Instant::now() + max_duration;
        let isolate_handle = js_runtime.v8_isolate().thread_safe_handle();
        // Signalled, not polled (#1342). The watchdog blocks until the deadline or
        // until the invocation says it is done — whichever comes first — so the
        // `join` below costs nothing. It used to poll a flag on a 10 ms sleep, and
        // the join waited out the rest of that sleep on *every* invocation: ~9 ms,
        // about 40 % of the total, against a guest whose own work was under 1 ms.
        let watchdog_done = Arc::new(watchdog::WatchdogSignal::new());
        let watchdog_done_thread = Arc::clone(&watchdog_done);
        let timed_out_watchdog = Arc::clone(&timed_out_run);
        let watchdog = std::thread::spawn(move || {
            if watchdog_done_thread.wait_until(invocation_deadline)
                == watchdog::WatchdogOutcome::DeadlineReached
            {
                timed_out_watchdog.store(true, Ordering::Release);
                isolate_handle.terminate_execution();
            }
        });

        // Near-heap-limit callback: when V8 approaches the hard limit, flag the
        // condition and terminate execution. We bump the returned limit so V8 does
        // not `FatalProcessOutOfMemory` (which would abort the whole process)
        // before the termination exception propagates out of the running script.
        let mem_flag = Arc::clone(&mem_exceeded_run);
        let heap_handle = js_runtime.v8_isolate().thread_safe_handle();
        js_runtime.add_near_heap_limit_callback(move |current_limit, _initial| {
            mem_flag.store(true, Ordering::Release);
            heap_handle.terminate_execution();
            // Grant headroom so V8 can unwind via the termination exception instead
            // of crashing the process.
            current_limit.saturating_add(current_limit / 2).max(current_limit + 1)
        });

        // Helper to convert a raw V8/deno error into a meaningful message,
        // distinguishing limit-driven termination from ordinary failures.
        let classify = |raw: &str, mem: &Arc<AtomicBool>, time: &Arc<AtomicBool>| -> String {
            if mem.load(Ordering::Acquire) {
                "Memory limit exceeded: heap allocation exceeded the configured limit".to_string()
            } else if time.load(Ordering::Acquire) {
                "Execution timeout: script exceeded the configured time limit".to_string()
            } else if raw.contains("SyntaxError") || raw.contains("Parse") {
                format!("SyntaxError: {raw}")
            } else {
                format!("Execution error: {raw}")
            }
        };

        // Execute the wrapped script. The watchdog stays ARMED past this point
        // (#804): the guest's continuations after `await` run inside
        // `run_event_loop` below, and a synchronous spin there is only
        // stoppable by the watchdog's `terminate_execution`.
        let phase_started = Instant::now();
        let exec_outcome = js_runtime.execute_script("<fraiseql-function>", wrapped);
        phases.script = phase_started.elapsed();

        if let Err(e) = exec_outcome {
            // Stop and reap the watchdog before returning.
            watchdog_done.finish();
            let _ = watchdog.join();
            return Err(classify(&e.to_string(), &mem_exceeded_run, &timed_out_run));
        }

        // Drive the event loop to resolve Promises from async functions. The
        // tokio timeout covers *async idle* hangs (a pending promise that never
        // resolves — no JS running, so `terminate_execution` has nothing to
        // terminate); the still-armed watchdog covers *synchronous* spins that
        // never yield back to tokio. Both share the single invocation budget.
        let phase_started = Instant::now();
        let loop_outcome = tokio::time::timeout(
            invocation_deadline.saturating_duration_since(std::time::Instant::now()),
            js_runtime.run_event_loop(deno_core::PollEventLoopOptions::default()),
        )
        .await;
        phases.event_loop = phase_started.elapsed();

        // The event loop is done (or timed out): stop and reap the watchdog.
        watchdog_done.finish();
        let _ = watchdog.join();

        match loop_outcome {
            Err(_) => {
                return Err("Execution timeout: event loop exceeded time limit".to_string());
            },
            Ok(Err(e)) => {
                return Err(classify(&e.to_string(), &mem_exceeded_run, &timed_out_run));
            },
            Ok(Ok(())) => {},
        }

        // Retrieve the result stored in globalThis.__fraiseql_result
        let phase_started = Instant::now();
        let result_global = js_runtime
            .execute_script("<get-result>", "globalThis.__fraiseql_result")
            .map_err(|e| format!("Failed to read result: {e}"))?;

        let error_global = js_runtime
            .execute_script("<get-error>", "globalThis.__fraiseql_error")
            .map_err(|e| format!("Failed to read error: {e}"))?;

        // Inspect the values via a V8 handle scope
        let (result_json, error_str) = {
            let scope = &mut js_runtime.handle_scope();

            let result_local = deno_core::v8::Local::new(scope, result_global);
            let error_local = deno_core::v8::Local::new(scope, error_global);

            // Both undefined means the IIFE wrapper never completed (e.g. an unresolvable
            // Promise caused the event loop to drain without the try/catch finalising).
            if result_local.is_undefined() && error_local.is_undefined() {
                return Err("Execution incomplete: function did not produce a result \
                     (possible unresolved promise)"
                    .to_string());
            }

            let error_str = if error_local.is_null_or_undefined() {
                None
            } else {
                Some(error_local.to_rust_string_lossy(scope))
            };

            let result_json = if result_local.is_null_or_undefined() {
                None
            } else {
                Some(result_local.to_rust_string_lossy(scope))
            };

            (result_json, error_str)
        };

        // If the guest threw an error, propagate it
        if let Some(err) = error_str {
            return Err(format!("Runtime error: {err}"));
        }

        // Parse the JSON result
        let value: Value = match result_json {
            Some(json_str) => serde_json::from_str(&json_str).unwrap_or(Value::String(json_str)),
            None => Value::Null,
        };
        phases.result = phase_started.elapsed();

        // Timed here rather than left to the end of the block, so teardown is a phase
        // of its own instead of an unexplained gap in the total.
        let phase_started = Instant::now();
        drop(js_runtime);
        phases.teardown = phase_started.elapsed();

        Ok(value)
    });
    let phase_started = Instant::now();
    drop(rt);
    phases.teardown += phase_started.elapsed();

    // Collect logs
    let logs = logs_arc.lock().expect("log mutex poisoned").clone();

    match result {
        Ok(value) => Ok(ExecutionResult {
            value,
            logs,
            phases,
        }),
        Err(e) => Err(e),
    }
}

#[cfg(test)]
mod tests;