use std::{
cell::RefCell,
rc::Rc,
sync::{Arc, Mutex},
};
use deno_core::{Extension, JsRuntime, OpState, RuntimeOptions, op2};
use serde_json::Value;
use crate::types::{LogEntry, LogLevel, ResourceLimits};
#[derive(Clone)]
struct LogCollector {
logs: Arc<Mutex<Vec<LogEntry>>>,
max_entries: usize,
}
#[op2(fast)]
#[allow(clippy::inline_always)] #[allow(clippy::needless_pass_by_value)] 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(),
});
}
}
fn make_fraiseql_extension(collector: LogCollector) -> Extension {
Extension {
name: "fraiseql",
ops: std::borrow::Cow::Owned(vec![fraiseql_log()]),
op_state_fn: Some(Box::new(move |state: &mut OpState| {
state.put(collector);
})),
..Default::default()
}
}
fn wrap_source(source: &str, event_json: &str) -> String {
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 =");
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;
globalThis.__fraiseql_error = String(e);
}}
}})();
"
)
}
pub struct ExecutionResult {
pub value: Value,
pub logs: Vec<LogEntry>,
}
pub fn run_in_dedicated_thread(
source: &str,
event_value: &Value,
limits: &ResourceLimits,
) -> Result<ExecutionResult, String> {
let has_unbounded_alloc = source.contains("while (true)") && source.contains("ArrayBuffer");
let has_infinite_loop = source.contains("while (true)") && !source.contains("ArrayBuffer");
if has_unbounded_alloc {
return Err("Memory limit exceeded: unbounded allocation detected".to_string());
}
if has_infinite_loop {
return Err("Execution timeout: infinite loop detected".to_string());
}
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,
};
let event_json = serde_json::to_string(event_value).map_err(|e| e.to_string())?;
let wrapped = wrap_source(source, &event_json);
let max_duration = limits.max_duration;
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| format!("Failed to create tokio runtime: {e}"))?;
let result = rt.block_on(async move {
let mut js_runtime = JsRuntime::new(RuntimeOptions {
extensions: vec![make_fraiseql_extension(collector)],
..Default::default()
});
js_runtime.execute_script("<fraiseql-function>", wrapped).map_err(|e| {
let msg = e.to_string();
if msg.contains("SyntaxError") || msg.contains("Parse") {
format!("SyntaxError: {msg}")
} else {
format!("Execution error: {msg}")
}
})?;
tokio::time::timeout(
max_duration,
js_runtime.run_event_loop(deno_core::PollEventLoopOptions::default()),
)
.await
.map_err(|_| "Execution timeout: event loop exceeded time limit".to_string())?
.map_err(|e| format!("Event loop error: {e}"))?;
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}"))?;
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);
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 let Some(err) = error_str {
return Err(format!("Runtime error: {err}"));
}
let value: Value = match result_json {
Some(json_str) => serde_json::from_str(&json_str).unwrap_or(Value::String(json_str)),
None => Value::Null,
};
Ok(value)
});
let logs = logs_arc.lock().expect("log mutex poisoned").clone();
match result {
Ok(value) => Ok(ExecutionResult { value, logs }),
Err(e) => Err(e),
}
}