mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-21 08:02:38 +00:00
* feat(nativets): bound fetch on a peer that never answers
deno_fetch applies no deadline of any kind. A peer that completes the TCP
handshake, accepts the request and then goes silent leaves `await fetch(...)`
pending indefinitely, holding its worker slot until the *job* timeout -- which
on self-hosted defaults to DEFAULT_SELFHOSTED_TIMEOUT, i.e. 7 days.
Nothing else catches this. Zombie-job detection keys off a stale
v2_job_runtime.ping, and a worker blocked inside a pending fetch keeps pinging
normally throughout: the worker is alive and healthy, only the work is dead.
What this bounds is the wait for a response to begin, and it stops there:
- a peer that never answers -> rejected after N seconds
- a peer slow to answer, but under N -> unaffected
- a body that then streams for an hour,
or is read slowly by the caller -> unaffected, always
That last line rules out the obvious implementation: AbortSignal.timeout(N)
around every fetch would bound the hang and break every streaming response and
long download. This is a hang detector, not a latency budget.
Default 300s via WINDMILL_FETCH_RESPONSE_TIMEOUT_SECS (0 disables), with a
per-script `//fetch_response_timeout <seconds>` annotation alongside the
existing //useragent and //proxy. Both nativets paths inherit it, since
eval_fetch_timeout and the dedicated-worker path in bun_executor both funnel
through create_nativets_runtime.
The ms value is clamped to i32::MAX: deno_web's setTimeout runs its delay
through webidl.converters.long, a 32-bit conversion that *wraps*, so a setting
past ~24.8 days would come out negative and fire immediately -- turning an
over-generous timeout into an instant one on every fetch.
The window covers connect, TLS and request upload as well as server think
time, so a very slow large upload is bounded by it too; the error message says
so rather than claiming the connection went silent.
Not covered: a body that stalls midway. Reaching that needs the response's
InnerBody, which deno_fetch keeps module-private, and every way to wrap it
from outside changes observable Response semantics (locking, bodyUsed,
double-consume errors). Left for a follow-up in deno_fetch itself.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JshNT5XVH78ZsTMvDWfFHb
* test(nativets): cover the instance-wide response-timeout env var
A typo in WINDMILL_FETCH_RESPONSE_TIMEOUT_SECS would compile, pass every
other test, and silently hand every operator the 300s default -- the same
class of silent-default failure the timeout itself exists to prevent. Its
own test binary, since a LazyLock resolves the value once per process.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JshNT5XVH78ZsTMvDWfFHb
* fix(nativets): inherit the caller's RequestInit, and clear the long-poll ceiling
Two problems with the first cut, both found in review.
`{ ...init, signal }` copied only own enumerable properties, but RequestInit is
a WebIDL dictionary whose members deno_fetch reads with plain property gets
that walk the prototype chain. Anything inherited or non-enumerable was
dropped: `fetch(url, Object.create({method: "POST"}))` silently became a GET.
Worse, a non-object init went from a loud TypeError to a silent GET, because
spreading "POST" yields {0:"P",1:"O",...} -- a valid dictionary with ignored
keys. Now the init is inherited from rather than copied, and a non-dictionary
is handed straight back to deno_fetch for its own TypeError.
The 300s default also sat at half of TIMEOUT_WAIT_RESULT (600s), which
run_wait_result long-polls against with no response headers. A script running
another job synchronously for 300-600s would have timed out client-side while
the server was still legitimately holding the request open -- the long-poll
risk class, instantiated inside the product and reachable without writing a
raw fetch. Default raised to 900s, with the constraint recorded where someone
would break it.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JshNT5XVH78ZsTMvDWfFHb
* fix(nativets): hand the caller's RequestInit to Request untouched
Carrying a WebIDL dictionary across by hand has no safe form, and both
previous attempts were wrong in opposite directions. Spreading a copy drops
inherited and non-enumerable members, and turns a non-object init from a loud
TypeError into a silent GET. Inheriting from it via Object.create fixes those
but makes the child the receiver, so an accessor on the original runs against
an object that lacks its private-field brand:
Cannot read private member #body from an object whose class did not
declare it
So don't carry it at all. fetch()'s own first act is `new Request(input,
init)`; doing that here hands the init to the same constructor, read exactly
as it would be without this wrapper, and our signal travels in an init we own.
`req.signal` is then deno's own resolution of init.signal over an input
Request's signal, which removes the hand-rolled version of that rule too.
The Request is built twice as a result, once here and once inside fetch. That
is cheap: cloneInnerRequest carries method, headers, redirect mode, clientRid
and blob entry, and a body is proxied rather than buffered -- a static body is
a shallow {body, consumed} copy sharing its bytes, a stream gets a one-chunk
pass-through.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JshNT5XVH78ZsTMvDWfFHb
* fix(nativets): keep an aborted fetch settling in the same tick
deno_fetch keeps its outer fetch non-async on purpose: "WPT has a test that
aborted fetch is settled in the same tick. This means we cannot wrap the
promise if it is already settled" (26_fetch.js). An `async` wrapper adopts
that promise through another one, so a rejection that used to land before any
microtask queued after the call now lands after it.
Made the wrapper non-async, with an early return that hands deno's settled
rejection straight back for an already-aborted signal, and no timer armed
there since there is no response to wait for. Construction still has to reject
rather than throw, so it is caught and returned as a rejection, which is what
the `async` was buying.
The comment claiming this matched deno_fetch's own `async function fetch` was
wrong on two counts -- that function is not async, and the wrapper was not
matching it. Replaced with the constraint that actually holds.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JshNT5XVH78ZsTMvDWfFHb
* fix(nativets): keep fetch's observable shape and reach intrinsics safely
Three ways the wrapper was distinguishable from the fetch it replaces, all
observable from a script sharing the isolate.
`.then` was an ordinary property lookup, so `Promise.prototype.then =
undefined` broke fetch after the request had already gone out. deno's own
modules reach intrinsics through primordials, and this file already captured
setTimeout, clearTimeout and Promise.reject for exactly that reason, so the
lookup was the odd one out. Now captured alongside them.
Declaring `init` without a default made `fetch.length` 2 where the standard
says 1. And the empty-call branch forwarded two explicit `undefined`s, so
deno's required-argument check saw two arguments and raised "Invalid URL:
'undefined'" instead of "1 argument required". Forwarding through
ReflectApply preserves the count.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JshNT5XVH78ZsTMvDWfFHb
* docs: clarify fetch timeout restart requirements
* fix: capture native fetch abort helpers
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
1105 lines
39 KiB
Rust
1105 lines
39 KiB
Rust
/*
|
|
* Author: Ruben Fiszel
|
|
* Copyright: Windmill Labs, Inc 2022
|
|
* This file and its contents are licensed under the AGPLv3 License.
|
|
* Please see the included NOTICE for copyright information and
|
|
* LICENSE-AGPL for a copy of the license.
|
|
*/
|
|
|
|
//! Isolated deno_core runtime for NativeTS script execution.
|
|
//!
|
|
//! This crate encapsulates all deno_core/V8 dependencies for executing
|
|
//! TypeScript scripts via the nativets runtime. By isolating this here,
|
|
//! deno_core compilation no longer blocks windmill-worker or windmill-api.
|
|
|
|
mod dedicated;
|
|
pub use dedicated::{ExecutingIsolate, PrewarmedIsolate, PrewarmedResult};
|
|
|
|
#[cfg(test)]
|
|
mod smoke_tests;
|
|
|
|
#[cfg(test)]
|
|
mod cert_tests;
|
|
|
|
use std::{
|
|
borrow::Cow,
|
|
cell::RefCell,
|
|
path::PathBuf,
|
|
rc::Rc,
|
|
sync::{Arc, LazyLock, Mutex},
|
|
};
|
|
|
|
// Re-export deno_telemetry for use by windmill-worker's otel proxy
|
|
pub use deno_telemetry;
|
|
|
|
use deno_ast::ParseParams;
|
|
use deno_core::{
|
|
op2, serde_v8, url,
|
|
v8::{self, IsolateHandle},
|
|
Extension, JsRuntime, OpState, PollEventLoopOptions, RuntimeOptions,
|
|
};
|
|
use deno_error::JsErrorBox;
|
|
use deno_fetch::FetchPermissions;
|
|
use deno_net::NetPermissions;
|
|
use deno_tls::{rustls::pki_types::CertificateDer, rustls::RootCertStore, RootCertStoreProvider};
|
|
use deno_web::{BlobStore, TimersPermission};
|
|
use itertools::Itertools;
|
|
use lazy_static::lazy_static;
|
|
use regex::Regex;
|
|
use serde_json::value::RawValue;
|
|
use sqlx::types::Json;
|
|
use tokio::sync::mpsc;
|
|
use uuid::Uuid;
|
|
|
|
use windmill_common::error::Error;
|
|
use windmill_common::result_stream::append_result_stream_db;
|
|
use windmill_common::worker::{write_file, Connection, WINDMILL_DIR};
|
|
|
|
// ── Snapshot-matched extensions ──────────────────────────────────────
|
|
//
|
|
// `deno_core` 0.352 validates that the snapshot's extension list is a
|
|
// *prefix* of the runtime's extension list (snapshot does not need an
|
|
// exact match — runtime is allowed to add extensions at the tail, but
|
|
// must not reorder or omit any that the snapshot baked in).
|
|
//
|
|
// Our snapshot (in build.rs) is the same eight deno_* extensions ending
|
|
// with this local `fetch` ext. The runtime adds one extra entry at the
|
|
// end — the windmill `ext` carrying our own ops — which is fine because
|
|
// it's after the snapshot prefix.
|
|
//
|
|
// This local `fetch` extension declaration must be present in both
|
|
// build.rs and lib.rs so the type passes through the `init()` macro.
|
|
// The ESM is already in the snapshot, so this `init()` call at runtime
|
|
// is a no-op for esm — the registration just records the ext.
|
|
deno_core::extension!(
|
|
fetch,
|
|
esm_entry_point = "ext:fetch/src/runtime.js",
|
|
esm = ["src/runtime.js"],
|
|
);
|
|
|
|
// ── Permission container ─────────────────────────────────────────────
|
|
|
|
pub struct PermissionsContainer;
|
|
|
|
impl FetchPermissions for PermissionsContainer {
|
|
#[inline(always)]
|
|
fn check_net_url(
|
|
&mut self,
|
|
_url: &deno_core::url::Url,
|
|
_api_name: &str,
|
|
) -> Result<(), deno_permissions::PermissionCheckError> {
|
|
Ok(())
|
|
}
|
|
|
|
#[inline(always)]
|
|
fn check_read<'a>(
|
|
&mut self,
|
|
path: Cow<'a, std::path::Path>,
|
|
_api_name: &str,
|
|
_get_path: &'a dyn deno_fs::GetPath,
|
|
) -> Result<deno_fs::CheckedPath<'a>, deno_io::fs::FsError> {
|
|
Ok(deno_fs::CheckedPath::Unresolved(path))
|
|
}
|
|
|
|
#[inline(always)]
|
|
fn check_write<'a>(
|
|
&mut self,
|
|
path: Cow<'a, std::path::Path>,
|
|
_api_name: &str,
|
|
_get_path: &'a dyn deno_fs::GetPath,
|
|
) -> Result<deno_fs::CheckedPath<'a>, deno_io::fs::FsError> {
|
|
Ok(deno_fs::CheckedPath::Unresolved(path))
|
|
}
|
|
|
|
#[inline(always)]
|
|
fn check_net_vsock(
|
|
&mut self,
|
|
_cid: u32,
|
|
_port: u32,
|
|
_api_name: &str,
|
|
) -> Result<(), deno_permissions::PermissionCheckError> {
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
impl TimersPermission for PermissionsContainer {
|
|
#[inline(always)]
|
|
fn allow_hrtime(&mut self) -> bool {
|
|
true
|
|
}
|
|
}
|
|
|
|
impl NetPermissions for PermissionsContainer {
|
|
fn check_read(
|
|
&mut self,
|
|
p: &str,
|
|
_api_name: &str,
|
|
) -> Result<PathBuf, deno_permissions::PermissionCheckError> {
|
|
Ok(PathBuf::from(p))
|
|
}
|
|
|
|
fn check_write(
|
|
&mut self,
|
|
p: &str,
|
|
_api_name: &str,
|
|
) -> Result<PathBuf, deno_permissions::PermissionCheckError> {
|
|
Ok(PathBuf::from(p))
|
|
}
|
|
|
|
fn check_net<T: AsRef<str>>(
|
|
&mut self,
|
|
_host: &(T, Option<u16>),
|
|
_api_name: &str,
|
|
) -> Result<(), deno_permissions::PermissionCheckError> {
|
|
Ok(())
|
|
}
|
|
|
|
fn check_write_path<'a>(
|
|
&mut self,
|
|
p: Cow<'a, std::path::Path>,
|
|
_api_name: &str,
|
|
) -> Result<Cow<'a, std::path::Path>, deno_permissions::PermissionCheckError> {
|
|
Ok(p)
|
|
}
|
|
|
|
fn check_vsock(
|
|
&mut self,
|
|
_cid: u32,
|
|
_port: u32,
|
|
_api_name: &str,
|
|
) -> Result<(), deno_permissions::PermissionCheckError> {
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
// ── Types ────────────────────────────────────────────────────────────
|
|
|
|
pub(crate) struct MainArgs {
|
|
pub(crate) args: Vec<Option<Box<RawValue>>>,
|
|
}
|
|
|
|
struct LogString {
|
|
pub s: mpsc::UnboundedSender<String>,
|
|
}
|
|
|
|
#[derive(Clone, Default)]
|
|
pub struct NativeAnnotation {
|
|
pub useragent: Option<String>,
|
|
pub proxy: Option<(String, Option<(String, String)>)>,
|
|
/// `//fetch_response_timeout <seconds>`: per-script override of
|
|
/// [`default_fetch_response_timeout_secs`]. `Some(0)` disables it for this
|
|
/// script; `None` leaves the default in force.
|
|
pub fetch_response_timeout_secs: Option<u64>,
|
|
}
|
|
|
|
/// How long `fetch()` waits for a response to begin, in seconds; `0` disables.
|
|
///
|
|
/// Covers everything up to the response headers and stops there, so a body may
|
|
/// then stream for any length of time. `src/runtime.js` holds the semantics.
|
|
///
|
|
/// Must exceed `TIMEOUT_WAIT_RESULT` (default 600), which holds synchronous job
|
|
/// calls open without headers. Raising that hot-reloaded instance setting may
|
|
/// also require raising `WINDMILL_FETCH_RESPONSE_TIMEOUT_SECS` in the deployment
|
|
/// and restarting workers: this environment value is cached for the process.
|
|
pub fn default_fetch_response_timeout_secs() -> u64 {
|
|
static SECS: LazyLock<u64> = LazyLock::new(|| {
|
|
std::env::var("WINDMILL_FETCH_RESPONSE_TIMEOUT_SECS")
|
|
.ok()
|
|
.and_then(|x| x.trim().parse::<u64>().ok())
|
|
.unwrap_or(900)
|
|
});
|
|
*SECS
|
|
}
|
|
|
|
/// Serializes V8 isolate creation as defense-in-depth against concurrent
|
|
/// creation races on x86_64 Linux. The primary fix is using the unprotected
|
|
/// V8 platform (see `setup_deno_runtime`).
|
|
static V8_ISOLATE_CREATE_LOCK: Mutex<()> = Mutex::new(());
|
|
|
|
/// Guard that terminates a running V8 isolate when dropped (e.g. on job cancellation).
|
|
struct IsolateDropGuard(Arc<Mutex<Option<IsolateHandle>>>);
|
|
|
|
impl Drop for IsolateDropGuard {
|
|
fn drop(&mut self) {
|
|
if let Some(handle) = self.0.lock().unwrap().take() {
|
|
handle.terminate_execution();
|
|
}
|
|
}
|
|
}
|
|
|
|
// ── Statics ──────────────────────────────────────────────────────────
|
|
|
|
static RUNTIME_SNAPSHOT: &[u8] = include_bytes!(concat!(env!("OUT_DIR"), "/FETCH_SNAPSHOT.bin"));
|
|
|
|
pub(crate) const WINDMILL_CLIENT: &str = include_str!("./windmill-client.js");
|
|
|
|
lazy_static::lazy_static! {
|
|
static ref ERROR_DIR: String = format!("{}/native_errors", *WINDMILL_DIR);
|
|
}
|
|
|
|
lazy_static! {
|
|
static ref RE_PROXY: Regex =
|
|
Regex::new(r"^(https?)://(([^:@\s]+):([^:@\s]+)@)?([^:@\s]+)(:(\d+))?$").unwrap();
|
|
}
|
|
|
|
lazy_static! {
|
|
/// Root cert store for the in-process nativets fetch runtime.
|
|
///
|
|
/// Unlike the Deno/Bun executors, nativets never spawns a child process, so
|
|
/// the CA env vars those executors forward (`DENO_CERT`, `DENO_TLS_CA_STORE`,
|
|
/// `SSL_CERT_FILE`/`NODE_EXTRA_CA_CERTS`) are never consumed by deno's CLI
|
|
/// layer. `deno_fetch` with `root_cert_store_provider: None` falls back to
|
|
/// the Mozilla webpki roots only, so corporate CAs fail with `UnknownIssuer`.
|
|
/// We read those env vars here and merge the certs into the default store.
|
|
///
|
|
/// Snapshotted once for the process lifetime, like the Deno executor's
|
|
/// `DENO_CERT`/`DENO_TLS_CA_STORE` lazy statics (`deno_executor.rs`): the
|
|
/// fetch root store is shared across all (potentially prewarmed) isolates, so
|
|
/// per-job CA reconfiguration is out of scope. A later `WORKER_CONFIG` reload
|
|
/// is not picked up until the process restarts.
|
|
static ref NATIVE_ROOT_CERT_STORE_PROVIDER: Option<Arc<dyn RootCertStoreProvider>> =
|
|
build_native_root_cert_store_provider();
|
|
}
|
|
|
|
struct NativeRootCertStoreProvider {
|
|
store: RootCertStore,
|
|
}
|
|
|
|
impl RootCertStoreProvider for NativeRootCertStoreProvider {
|
|
fn get_or_try_init(&self) -> Result<&RootCertStore, JsErrorBox> {
|
|
Ok(&self.store)
|
|
}
|
|
}
|
|
|
|
/// Resolve a CA-related env var the same way the child executors see it: the
|
|
/// worker's own process env, then the worker-group config (`env_vars_allowlist`
|
|
/// forwarded values + DB `env_vars_static` literals, resolved into
|
|
/// `WORKER_CONFIG.env_vars`). Child Deno/Bun jobs receive that config map via
|
|
/// `.envs(...)`, so nativets must consult it too or a CA set only through worker
|
|
/// config would silently not apply in-process.
|
|
fn resolve_ca_env_var(name: &str) -> Option<String> {
|
|
if let Ok(v) = std::env::var(name) {
|
|
if !v.is_empty() {
|
|
return Some(v);
|
|
}
|
|
}
|
|
windmill_common::worker::WORKER_CONFIG
|
|
.load()
|
|
.env_vars
|
|
.get(name)
|
|
.filter(|v| !v.is_empty())
|
|
.cloned()
|
|
}
|
|
|
|
/// Build a root cert store seeded with the Mozilla webpki roots plus any custom
|
|
/// CAs configured via env. Returns `None` when no custom CA is configured, which
|
|
/// preserves the previous default-only behaviour.
|
|
fn build_native_root_cert_store_provider() -> Option<Arc<dyn RootCertStoreProvider>> {
|
|
let mut store = deno_tls::create_default_root_cert_store();
|
|
let mut added = 0usize;
|
|
|
|
// File-path env vars, each pointing at a PEM bundle of one or more certs.
|
|
// `DENO_CERT` mirrors the Deno CLI; `SSL_CERT_FILE` is the OpenSSL standard
|
|
// also honoured by Bun/Node (via NODE_EXTRA_CA_CERTS). Dedupe by path because
|
|
// the tracing proxy points several of these at the same bundle, and rustls'
|
|
// RootCertStore::add does not dedupe — we'd otherwise trust the same root N times.
|
|
let mut seen_paths = std::collections::HashSet::new();
|
|
for var in ["DENO_CERT", "SSL_CERT_FILE", "NODE_EXTRA_CA_CERTS"] {
|
|
let Some(path) = resolve_ca_env_var(var).filter(|p| !p.is_empty()) else {
|
|
continue;
|
|
};
|
|
if !seen_paths.insert(path.clone()) {
|
|
continue;
|
|
}
|
|
match load_pem_certs_from_path(&path) {
|
|
Ok(certs) => {
|
|
for cert in certs {
|
|
if let Err(e) = store.add(cert) {
|
|
tracing::warn!("nativets: failed to add cert from {var}={path}: {e}");
|
|
} else {
|
|
added += 1;
|
|
}
|
|
}
|
|
}
|
|
Err(e) => tracing::warn!("nativets: failed to read CA file {var}={path}: {e}"),
|
|
}
|
|
}
|
|
|
|
// `DENO_TLS_CA_STORE=system` (comma-separated, may also contain `mozilla`)
|
|
// pulls in the OS trust store. Unlike the Deno CLI — where the list selects
|
|
// and orders the stores — this is purely additive: the Mozilla defaults are
|
|
// always seeded above, and `system` augments them. That is a strict superset
|
|
// of the public roots, which is what the corporate-CA use case needs.
|
|
if resolve_ca_env_var("DENO_TLS_CA_STORE")
|
|
.map(|v| v.split(',').any(|s| s.trim() == "system"))
|
|
.unwrap_or(false)
|
|
{
|
|
match deno_tls::deno_native_certs::load_native_certs() {
|
|
Ok(certs) => {
|
|
for cert in certs {
|
|
if store.add(CertificateDer::from(cert.0)).is_ok() {
|
|
added += 1;
|
|
}
|
|
}
|
|
}
|
|
Err(e) => tracing::warn!("nativets: failed to load system CA store: {e}"),
|
|
}
|
|
}
|
|
|
|
if added == 0 {
|
|
return None;
|
|
}
|
|
tracing::info!("nativets: loaded {added} custom CA cert(s) into fetch root store");
|
|
Some(Arc::new(NativeRootCertStoreProvider { store }))
|
|
}
|
|
|
|
fn load_pem_certs_from_path(path: &str) -> anyhow::Result<Vec<CertificateDer<'static>>> {
|
|
let file = std::fs::File::open(path)?;
|
|
let mut reader = std::io::BufReader::new(file);
|
|
deno_tls::load_certs(&mut reader).map_err(|e| anyhow::anyhow!(e))
|
|
}
|
|
|
|
// ── Public interface ─────────────────────────────────────────────────
|
|
|
|
/// Set up the deno_core/V8 runtime. Idempotent — safe to call multiple times.
|
|
/// Called automatically before JsRuntime creation, but can also be called
|
|
/// eagerly at startup for predictable initialization order.
|
|
pub fn setup_deno_runtime() -> anyhow::Result<()> {
|
|
use std::sync::Once;
|
|
static INIT: Once = Once::new();
|
|
|
|
let mut init_err: Option<String> = None;
|
|
INIT.call_once(|| {
|
|
// deno_fetch requires a TLS provider; install ring as default (idempotent).
|
|
let _ = rustls::crypto::ring::default_provider().install_default();
|
|
|
|
let unrecognized_v8_flags = deno_core::v8_set_flags(vec![
|
|
"--stack-size=1024".to_string(),
|
|
"--no-harmony-import-assertions".to_string(),
|
|
])
|
|
.into_iter()
|
|
.skip(1)
|
|
.collect::<Vec<_>>();
|
|
|
|
if !unrecognized_v8_flags.is_empty() {
|
|
init_err = Some(format!(
|
|
"Unrecognized V8 flags: {:?}",
|
|
unrecognized_v8_flags
|
|
));
|
|
}
|
|
|
|
// Use an unprotected platform that doesn't enforce thread-isolated allocations
|
|
// via Memory Protection Keys (pkeys). The default platform requires all V8-using
|
|
// threads to be descendants of the thread that called v8::Initialize, but tokio's
|
|
// spawn_blocking pool threads don't satisfy this. Without this, V8 crashes with
|
|
// SIGSEGV in WasmCodePointerTable::AllocateUninitializedEntry() on x86_64 Linux.
|
|
// See: https://github.com/denoland/deno_core/issues/952
|
|
let platform = deno_core::v8::new_unprotected_default_platform(0, false).make_shared();
|
|
deno_core::JsRuntime::init_platform(Some(platform), false);
|
|
});
|
|
|
|
if let Some(msg) = init_err {
|
|
println!("{msg}");
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
pub fn transpile_ts(expr: String) -> anyhow::Result<String> {
|
|
let parsed = deno_ast::parse_module(ParseParams {
|
|
specifier: url::Url::parse("file:///eval.ts")?,
|
|
capture_tokens: false,
|
|
scope_analysis: false,
|
|
media_type: deno_ast::MediaType::TypeScript,
|
|
maybe_syntax: None,
|
|
text: deno_core::ModuleCodeString::from(expr).into(),
|
|
})?;
|
|
Ok(parsed
|
|
.transpile(
|
|
&Default::default(),
|
|
&Default::default(),
|
|
&Default::default(),
|
|
)?
|
|
.into_source()
|
|
.text)
|
|
}
|
|
|
|
pub fn get_annotation(inner_content: &str) -> NativeAnnotation {
|
|
let mut res = NativeAnnotation::default();
|
|
|
|
let anns = inner_content
|
|
.lines()
|
|
.take_while(|x| x.starts_with("//"))
|
|
.map(|x| x.to_string().trim_start_matches("//").trim().to_string())
|
|
.collect_vec();
|
|
|
|
for ann in anns.iter() {
|
|
if ann.starts_with("useragent") {
|
|
res.useragent = Some(ann.trim_start_matches("useragent").trim().to_string());
|
|
} else if ann.starts_with("proxy") {
|
|
res.proxy = capture_proxy(ann.trim_start_matches("proxy").trim());
|
|
} else if ann.starts_with("fetch_response_timeout") {
|
|
// A typo falls back to the default, never to "no timeout".
|
|
res.fetch_response_timeout_secs = ann
|
|
.trim_start_matches("fetch_response_timeout")
|
|
.trim()
|
|
.parse::<u64>()
|
|
.ok();
|
|
}
|
|
}
|
|
res
|
|
}
|
|
|
|
fn capture_proxy(s: &str) -> Option<(String, Option<(String, String)>)> {
|
|
RE_PROXY.captures(s).map(|x| {
|
|
(
|
|
format!(
|
|
"{}://{}{}",
|
|
x.get(1).map(|x| x.as_str()).unwrap_or_default(),
|
|
x.get(5).map(|x| x.as_str()).unwrap_or_default(),
|
|
x.get(7)
|
|
.map(|x| format!(":{}", x.as_str()))
|
|
.unwrap_or_default(),
|
|
),
|
|
x.get(3).map(|y| {
|
|
(
|
|
y.as_str().to_string(),
|
|
x.get(4).map(|x| x.as_str().to_string()).unwrap_or_default(),
|
|
)
|
|
}),
|
|
)
|
|
})
|
|
}
|
|
|
|
fn write_error_expr(expr: &str, uuid: &Uuid) {
|
|
if let Err(e) = std::fs::create_dir_all(&*ERROR_DIR) {
|
|
tracing::error!("failed to create error dir {}: {e}", *ERROR_DIR);
|
|
return;
|
|
}
|
|
let dir_entries = match std::fs::read_dir(&*ERROR_DIR) {
|
|
Ok(entries) => entries.count(),
|
|
Err(_) => {
|
|
tracing::error!("failed to read error dir {}", *ERROR_DIR);
|
|
return;
|
|
}
|
|
};
|
|
|
|
if std::env::var("PRINT_NATIVE_ERRORS").is_ok() {
|
|
tracing::info!("native error for job {uuid}: {expr}");
|
|
}
|
|
if dir_entries >= 100 {
|
|
tracing::info!("Too many error files in {}, skipping write", *ERROR_DIR);
|
|
return;
|
|
}
|
|
|
|
let path = format!("/{uuid}.js");
|
|
tracing::info!(
|
|
"nativets job {uuid} failed, writing error expr to {}/{path} for debugging: {path}",
|
|
*ERROR_DIR
|
|
);
|
|
if let Err(e) = write_file(&ERROR_DIR, &path, expr) {
|
|
tracing::error!("failed to write error expr to file {path}: {e}");
|
|
}
|
|
}
|
|
|
|
use windmill_common::utils::unsafe_raw;
|
|
|
|
async fn append_result_stream(
|
|
conn: &Connection,
|
|
workspace_id: &str,
|
|
job_id: &Uuid,
|
|
nstream: &str,
|
|
offset: i32,
|
|
) -> windmill_common::error::Result<()> {
|
|
match conn {
|
|
Connection::Sql(db) => {
|
|
append_result_stream_db(db, workspace_id, job_id, nstream, offset).await?;
|
|
}
|
|
Connection::Http(client) => {
|
|
#[derive(serde::Serialize)]
|
|
struct ResultStreamBody<'a> {
|
|
result_stream: &'a str,
|
|
offset: i32,
|
|
}
|
|
let body = ResultStreamBody { result_stream: nstream, offset };
|
|
if let Err(e) = client
|
|
.post::<_, String>(
|
|
&format!(
|
|
"/api/w/{}/agent_workers/push_result_stream/{}",
|
|
workspace_id, job_id
|
|
),
|
|
None,
|
|
&body,
|
|
)
|
|
.await
|
|
{
|
|
tracing::error!(%job_id, %e, "error sending result stream for job {job_id}: {e}");
|
|
}
|
|
}
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
// ── ops ──────────────────────────────────────────────────────────────
|
|
|
|
#[op2]
|
|
#[serde]
|
|
fn op_get_static_args(op_state: Rc<RefCell<OpState>>) -> Vec<Option<String>> {
|
|
op_state
|
|
.borrow()
|
|
.borrow::<MainArgs>()
|
|
.args
|
|
.iter()
|
|
.map(|x| x.as_ref().map(|y| y.get().to_string()))
|
|
.collect_vec()
|
|
}
|
|
|
|
#[op2(fast)]
|
|
fn op_log(op_state: Rc<RefCell<OpState>>, #[string] log: &str) {
|
|
if let Err(e) = op_state
|
|
.borrow_mut()
|
|
.borrow_mut::<LogString>()
|
|
.s
|
|
.send(log.to_string())
|
|
{
|
|
tracing::error!("failed to send log: {e}");
|
|
}
|
|
}
|
|
|
|
// ── Shared V8 runtime creation ───────────────────────────────────────
|
|
|
|
pub(crate) struct CreatedRuntime {
|
|
pub(crate) js_runtime: JsRuntime,
|
|
pub(crate) log_receiver: mpsc::UnboundedReceiver<String>,
|
|
pub(crate) memory_limit_rx: mpsc::UnboundedReceiver<()>,
|
|
}
|
|
|
|
/// Create a JsRuntime with the standard nativets extensions, heap limit
|
|
/// callback, and log channel. Must be called on a blocking thread (not
|
|
/// on the async tokio runtime) because V8 isolate creation is
|
|
/// synchronous and potentially heavy.
|
|
pub(crate) fn create_nativets_runtime(
|
|
ann: NativeAnnotation,
|
|
initial_args: Vec<Option<Box<RawValue>>>,
|
|
) -> anyhow::Result<CreatedRuntime> {
|
|
let ops = vec![op_get_static_args(), op_log()];
|
|
let ext = Extension { name: "windmill", ops: ops.into(), ..Default::default() };
|
|
|
|
// deno_web's setTimeout puts its delay through `webidl.converters.long`,
|
|
// which wraps at 32 bits: past i32::MAX ms (~24.8 days) the delay comes out
|
|
// negative and fires immediately, so an over-generous setting would abort
|
|
// every fetch on the spot. Cap rather than wrap.
|
|
let fetch_response_timeout_ms = ann
|
|
.fetch_response_timeout_secs
|
|
.unwrap_or_else(default_fetch_response_timeout_secs)
|
|
.saturating_mul(1000)
|
|
.min(i32::MAX as u64);
|
|
|
|
let fetch_options = deno_fetch::Options {
|
|
root_cert_store_provider: NATIVE_ROOT_CERT_STORE_PROVIDER.clone(),
|
|
user_agent: ann.useragent.unwrap_or_else(|| "windmill/beta".to_string()),
|
|
proxy: ann.proxy.map(|x| deno_tls::Proxy::Http {
|
|
url: x.0,
|
|
basic_auth: x
|
|
.1
|
|
.map(|(username, password)| deno_tls::BasicAuth { username, password }),
|
|
}),
|
|
..Default::default()
|
|
};
|
|
|
|
let exts: Vec<Extension> = vec![
|
|
deno_telemetry::deno_telemetry::init(),
|
|
deno_webidl::deno_webidl::init(),
|
|
deno_url::deno_url::init(),
|
|
deno_console::deno_console::init(),
|
|
deno_web::deno_web::init::<PermissionsContainer>(Arc::new(BlobStore::default()), None),
|
|
// Registered after deno_web to keep the snapshot (build.rs) a prefix of
|
|
// the runtime extension list; deno_crypto declares deps = [deno_webidl, deno_web].
|
|
deno_crypto::deno_crypto::init(None),
|
|
deno_fetch::deno_fetch::init::<PermissionsContainer>(fetch_options),
|
|
deno_net::deno_net::init::<PermissionsContainer>(None, None),
|
|
fetch::init(),
|
|
ext,
|
|
];
|
|
|
|
let options = RuntimeOptions {
|
|
is_main: true,
|
|
extensions: exts,
|
|
create_params: Some(
|
|
deno_core::v8::CreateParams::default().heap_limits(0, 1024 * 1024 * 128),
|
|
),
|
|
startup_snapshot: Some(RUNTIME_SNAPSHOT),
|
|
module_loader: Some(Rc::new(deno_core::FsModuleLoader)),
|
|
extension_transpiler: None,
|
|
..Default::default()
|
|
};
|
|
|
|
let (memory_limit_tx, memory_limit_rx) = mpsc::unbounded_channel::<()>();
|
|
|
|
setup_deno_runtime().expect("V8 platform init failed");
|
|
|
|
let mut js_runtime = {
|
|
let _v8_lock = V8_ISOLATE_CREATE_LOCK
|
|
.lock()
|
|
.unwrap_or_else(|e| e.into_inner());
|
|
JsRuntime::new(options)
|
|
};
|
|
|
|
js_runtime.add_near_heap_limit_callback(move |x, y| {
|
|
tracing::error!("heap limit reached: {x} {y}");
|
|
if memory_limit_tx.send(()).is_err() {
|
|
tracing::warn!(
|
|
"memory limit notification channel closed - isolate may already be terminating"
|
|
);
|
|
}
|
|
y * 2
|
|
});
|
|
|
|
let (log_sender, log_receiver) = mpsc::unbounded_channel::<String>();
|
|
|
|
{
|
|
let op_state = js_runtime.op_state();
|
|
let mut op_state = op_state.borrow_mut();
|
|
op_state.put(PermissionsContainer {});
|
|
op_state.put(MainArgs { args: initial_args });
|
|
op_state.put(LogString { s: log_sender });
|
|
}
|
|
|
|
// Per-isolate JS init that can't run in the snapshot (runtime.js executes at
|
|
// snapshot-build time): the wall clock behind performance.timeOrigin and the
|
|
// fetch response timeout are both per-isolate values.
|
|
js_runtime
|
|
.execute_script(
|
|
"<wm_init>",
|
|
format!(
|
|
"globalThis.__wmInitPerIsolate({{ fetchResponseTimeoutMs: {fetch_response_timeout_ms} }})"
|
|
),
|
|
)
|
|
.map_err(windmill_common::error::to_anyhow)?;
|
|
|
|
Ok(CreatedRuntime { js_runtime, log_receiver, memory_limit_rx })
|
|
}
|
|
|
|
// ── Shared module-loading helpers ────────────────────────────────────
|
|
|
|
pub(crate) async fn load_client_module(
|
|
js_runtime: &mut JsRuntime,
|
|
env_code: &str,
|
|
) -> anyhow::Result<()> {
|
|
js_runtime
|
|
.load_side_es_module_from_code(
|
|
&deno_core::resolve_url("file:///windmill.ts")
|
|
.map_err(windmill_common::error::to_anyhow)?,
|
|
format!("{env_code}\n{WINDMILL_CLIENT}"),
|
|
)
|
|
.await
|
|
.map_err(windmill_common::error::to_anyhow)?;
|
|
Ok(())
|
|
}
|
|
|
|
pub(crate) async fn load_user_module(
|
|
js_runtime: &mut JsRuntime,
|
|
source: String,
|
|
) -> anyhow::Result<()> {
|
|
use anyhow::Context;
|
|
js_runtime
|
|
.load_side_es_module_from_code(
|
|
&deno_core::resolve_url("file:///eval.ts")
|
|
.map_err(windmill_common::error::to_anyhow)?,
|
|
source,
|
|
)
|
|
.await
|
|
.context("failed to load module")?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Extract a string result from a resolved V8 global and convert to `Box<RawValue>`.
|
|
pub(crate) fn extract_global_string(
|
|
js_runtime: &mut JsRuntime,
|
|
global: v8::Global<v8::Value>,
|
|
) -> Result<Box<RawValue>, String> {
|
|
let scope = &mut js_runtime.handle_scope();
|
|
let local = v8::Local::new(scope, global);
|
|
match serde_v8::from_v8::<Option<String>>(scope, local) {
|
|
Ok(s) => Ok(unsafe_raw(s.unwrap_or_else(|| "null".to_string()))),
|
|
Err(e) => Err(format!("failed to deserialize result: {e}")),
|
|
}
|
|
}
|
|
|
|
// ── eval_fetch_timeout ───────────────────────────────────────────────
|
|
|
|
/// Execute a NativeTS script using deno_core/V8.
|
|
///
|
|
/// Returns `(result, has_stream)` where `has_stream` indicates if the result
|
|
/// came from an async iterable stream.
|
|
///
|
|
/// `otel_initialized` should be `DENO_OTEL_INITIALIZED.load(SeqCst)`.
|
|
///
|
|
/// The caller (windmill-worker) is responsible for wrapping this in
|
|
/// `run_future_with_polling_update_job_poller` for job cancellation/polling.
|
|
#[allow(clippy::too_many_arguments)]
|
|
pub async fn eval_fetch_timeout(
|
|
env_code: String,
|
|
ts_expr: String,
|
|
js_expr: String,
|
|
args: Option<&Json<std::collections::HashMap<String, Box<RawValue>>>>,
|
|
script_entrypoint_override: Option<String>,
|
|
job_id: Uuid,
|
|
conn: &Connection,
|
|
w_id: &str,
|
|
load_client: bool,
|
|
otel_initialized: bool,
|
|
stream_notifier_update: Option<Arc<dyn Fn() + Send + Sync + 'static>>,
|
|
) -> windmill_common::error::Result<(Box<RawValue>, bool)> {
|
|
let isolate_handle: Arc<Mutex<Option<IsolateHandle>>> = Arc::new(Mutex::new(None));
|
|
let _isolate_guard = IsolateDropGuard(isolate_handle.clone());
|
|
let (append_logs_sender, mut append_logs_receiver) = mpsc::unbounded_channel::<String>();
|
|
let (result_stream_sender, mut result_stream_receiver) = mpsc::unbounded_channel::<String>();
|
|
|
|
let conn_ = conn.clone();
|
|
let w_id_ = w_id.to_string();
|
|
tokio::spawn(async move {
|
|
while let Some(log) = append_logs_receiver.recv().await {
|
|
windmill_queue::append_logs(&job_id, &w_id_, log, &conn_).await
|
|
}
|
|
});
|
|
|
|
let append_result_stream_fn = append_result_stream;
|
|
let conn_ = conn.clone();
|
|
let w_id_ = w_id.to_string();
|
|
tokio::spawn(async move {
|
|
let mut offset = -1;
|
|
while let Some(stream) = result_stream_receiver.recv().await {
|
|
offset += 1;
|
|
if let Err(e) = append_result_stream_fn(&conn_, &w_id_, &job_id, &stream, offset).await
|
|
{
|
|
tracing::error!("failed to append result stream: {e}");
|
|
}
|
|
}
|
|
});
|
|
|
|
let parsed_args = windmill_parser_ts::parse_deno_signature(
|
|
&ts_expr,
|
|
true,
|
|
false,
|
|
script_entrypoint_override.clone(),
|
|
)?
|
|
.args;
|
|
let spread = parsed_args
|
|
.into_iter()
|
|
.map(|x| {
|
|
args.as_ref()
|
|
.and_then(|args| args.0.get(&x.name).map(|x| x.clone()))
|
|
})
|
|
.collect::<Vec<_>>();
|
|
|
|
let ann = get_annotation(&ts_expr);
|
|
|
|
#[cfg(not(feature = "enterprise"))]
|
|
if ann.proxy.is_some() {
|
|
return Err(Error::ExecutionErr("Proxy is an EE feature".to_string()).into());
|
|
}
|
|
|
|
let mut extra_logs = String::new();
|
|
if ann.useragent.is_some() {
|
|
extra_logs.push_str(&format!("useragent: {}\n", ann.useragent.as_ref().unwrap()));
|
|
}
|
|
if ann.proxy.is_some() {
|
|
let (proxy, auth) = ann.proxy.as_ref().unwrap();
|
|
extra_logs.push_str(&format!(
|
|
"proxy: {proxy} (basic auth: {})\n",
|
|
auth.is_some()
|
|
));
|
|
}
|
|
|
|
let w_id_for_tracing = w_id.to_string();
|
|
let result_f = tokio::task::spawn_blocking(move || {
|
|
let CreatedRuntime { mut js_runtime, mut log_receiver, mut memory_limit_rx } =
|
|
create_nativets_runtime(ann, spread)?;
|
|
|
|
if otel_initialized {
|
|
if let Err(e) =
|
|
js_runtime.execute_script("<otel_bootstrap>", "globalThis.__bootstrapOtel()")
|
|
{
|
|
tracing::warn!("Failed to bootstrap OTEL telemetry: {}", e);
|
|
}
|
|
}
|
|
|
|
*isolate_handle.lock().unwrap_or_else(|e| e.into_inner()) =
|
|
Some(js_runtime.v8_isolate().thread_safe_handle());
|
|
|
|
let runtime = tokio::runtime::Builder::new_current_thread()
|
|
.enable_all()
|
|
.build()?;
|
|
|
|
let future = async {
|
|
if !extra_logs.is_empty() {
|
|
if let Err(e) = append_logs_sender.send(extra_logs) {
|
|
tracing::error!("failed to send extra logs: {e}");
|
|
}
|
|
}
|
|
let w_id_for_tracing = w_id_for_tracing;
|
|
let handle = tokio::spawn(async move {
|
|
let mut result_stream = String::new();
|
|
let mut is_stream = false;
|
|
while let Some(log) = log_receiver.recv().await {
|
|
use windmill_common::result_stream::extract_stream_from_logs;
|
|
use windmill_common::tracing_init::{OTEL_JOB_LOGS, OTEL_PREFIX};
|
|
|
|
// Mirror `process_streaming_log_lines` (EE) + the OTEL_JOB_LOGS
|
|
// hook from handle_child.rs, neither of which runs for nativets
|
|
// since nativets delivers logs in-process via the log channel.
|
|
for line in log.lines() {
|
|
tracing::info!(
|
|
target: "windmill:job_log",
|
|
job_id = ?job_id,
|
|
workspace_id = ?w_id_for_tracing,
|
|
"{line}"
|
|
);
|
|
if *OTEL_JOB_LOGS {
|
|
if let Some(otel_suffix) = line.strip_prefix(OTEL_PREFIX) {
|
|
tracing::event!(tracing::Level::INFO, otel_suffix);
|
|
}
|
|
}
|
|
}
|
|
|
|
if let Some(stream) = extract_stream_from_logs(&log.trim_end_matches("\n")) {
|
|
if !is_stream {
|
|
is_stream = true;
|
|
if let Some(ref f) = stream_notifier_update {
|
|
f();
|
|
}
|
|
}
|
|
|
|
result_stream.push_str(&stream);
|
|
if let Err(e) = result_stream_sender.send(stream) {
|
|
tracing::error!("failed to send result stream: {e}");
|
|
}
|
|
} else {
|
|
if let Err(e) = append_logs_sender.send(log) {
|
|
tracing::error!("failed to send log: {e}");
|
|
}
|
|
}
|
|
}
|
|
if !result_stream.is_empty() {
|
|
Some(result_stream)
|
|
} else {
|
|
None
|
|
}
|
|
});
|
|
|
|
let r = tokio::select! {
|
|
r = eval_fetch(&mut js_runtime, &js_expr, Some(env_code), script_entrypoint_override, load_client, &job_id, otel_initialized) => Ok(r),
|
|
_ = memory_limit_rx.recv() => Err(Error::ExecutionErr("Memory limit reached, killing isolate".to_string()))
|
|
};
|
|
*isolate_handle.lock().unwrap_or_else(|e| e.into_inner()) = None;
|
|
drop(js_runtime);
|
|
if let Ok(r) = r {
|
|
match handle.await {
|
|
Ok(Some(logs)) => {
|
|
// merge_result_stream: if main result is null but stream exists, use stream
|
|
match r {
|
|
Ok(raw) if raw.get() == "null" => Ok((unsafe_raw(logs), true)),
|
|
Ok(raw) => Ok((raw, true)),
|
|
Err(e) => Err(e),
|
|
}
|
|
}
|
|
Ok(None) => Ok(r.map(|r| (r, false))?),
|
|
Err(e) => Err(Error::ExecutionErr(e.to_string())),
|
|
}
|
|
} else {
|
|
r.map(|r| r.map(|r| (r, false)))?
|
|
}
|
|
};
|
|
let r = runtime.block_on(future)?;
|
|
|
|
Ok(r) as windmill_common::error::Result<(Box<RawValue>, bool)>
|
|
});
|
|
|
|
result_f.await.map_err(windmill_common::error::to_anyhow)?
|
|
}
|
|
|
|
async fn eval_fetch(
|
|
js_runtime: &mut JsRuntime,
|
|
expr: &str,
|
|
env_code: Option<String>,
|
|
script_entrypoint_override: Option<String>,
|
|
load_client: bool,
|
|
job_id: &Uuid,
|
|
otel_initialized: bool,
|
|
) -> windmill_common::error::Result<Box<RawValue>> {
|
|
if load_client {
|
|
if let Some(env_code) = env_code.as_ref() {
|
|
load_client_module(js_runtime, env_code).await?;
|
|
}
|
|
}
|
|
let source = format!("{}\n{expr}", env_code.unwrap_or_default());
|
|
if let Err(e) = load_user_module(js_runtime, source.clone()).await {
|
|
write_error_expr(expr, job_id);
|
|
return Err(e.into());
|
|
}
|
|
|
|
let result = execute_main(
|
|
js_runtime,
|
|
script_entrypoint_override.as_deref(),
|
|
otel_initialized,
|
|
Some(job_id),
|
|
)
|
|
.await;
|
|
|
|
match result {
|
|
Ok(raw) => Ok(raw),
|
|
Err(ExecuteError::Script(msg)) => {
|
|
write_error_expr(expr, job_id);
|
|
Err(Error::ExecutionErr(msg))
|
|
}
|
|
Err(ExecuteError::Js { message, stack, name, source: eval_source }) => {
|
|
write_error_expr(expr, job_id);
|
|
use windmill_common::worker::to_raw_value;
|
|
let stack_head = eval_source.and_then(|(file, line_no)| {
|
|
if file == "file:///eval.ts" {
|
|
source
|
|
.lines()
|
|
.nth(line_no.saturating_sub(1))
|
|
.map(|l| format!("{l}\n"))
|
|
} else {
|
|
None
|
|
}
|
|
});
|
|
let stack_s = format!(
|
|
"{}{}",
|
|
stack_head.unwrap_or_default(),
|
|
stack.as_deref().unwrap_or_default()
|
|
);
|
|
let stack = if stack_s.is_empty() {
|
|
None
|
|
} else {
|
|
Some(stack_s)
|
|
};
|
|
Err(Error::ExecutionRawError(to_raw_value(&serde_json::json!({
|
|
"message": message,
|
|
"stack": stack,
|
|
"name": name,
|
|
}))))
|
|
}
|
|
}
|
|
}
|
|
|
|
// ── Shared execution engine ──────────────────────────────────────────
|
|
|
|
pub(crate) enum ExecuteError {
|
|
/// Non-JS error (V8 internal, init failure, deserialization)
|
|
Script(String),
|
|
/// JS exception with structured error info
|
|
Js {
|
|
message: Option<String>,
|
|
stack: Option<String>,
|
|
name: Option<String>,
|
|
/// (file_name, line_number) from the first stack frame, if in user code
|
|
source: Option<(String, usize)>,
|
|
},
|
|
}
|
|
|
|
/// Execute the `main` function from the already-loaded `eval.ts` module.
|
|
///
|
|
/// Args must already be set in `MainArgs` in the runtime's OpState.
|
|
/// Modules (`windmill.ts` and `eval.ts`) must already be loaded.
|
|
pub(crate) async fn execute_main(
|
|
js_runtime: &mut JsRuntime,
|
|
entrypoint: Option<&str>,
|
|
_otel_initialized: bool,
|
|
_job_id: Option<&Uuid>,
|
|
) -> Result<Box<RawValue>, ExecuteError> {
|
|
let main_fn = entrypoint.unwrap_or("main");
|
|
|
|
#[cfg(all(feature = "private", feature = "enterprise"))]
|
|
let otel_context_inject = if _otel_initialized {
|
|
let trace_id = _job_id
|
|
.map(|id| id.as_simple().to_string())
|
|
.unwrap_or_default();
|
|
format!(
|
|
r#"globalThis.__enterSpan?.({{
|
|
isRecording: () => true,
|
|
spanContext: () => ({{ traceId: "{trace_id}", spanId: "ffffffffffffffff", traceFlags: 1 }})
|
|
}});"#
|
|
)
|
|
} else {
|
|
String::new()
|
|
};
|
|
|
|
#[cfg(not(all(feature = "private", feature = "enterprise")))]
|
|
let otel_context_inject = "";
|
|
|
|
let script = js_runtime
|
|
.execute_script(
|
|
"<anon>",
|
|
format!(
|
|
r#"
|
|
function isAsyncIterable(obj) {{
|
|
return obj != null && typeof obj[Symbol.asyncIterator] === 'function';
|
|
}}
|
|
|
|
function processStreamIterative(res) {{
|
|
const iterator = res[Symbol.asyncIterator]();
|
|
|
|
function processLoop() {{
|
|
return new Promise(function(resolve) {{
|
|
function step() {{
|
|
iterator.next().then(function(result) {{
|
|
if (!result.done) {{
|
|
const chunk = result.value;
|
|
console.log("WM_STREAM: " + chunk.replace(/\n/g, '\\n'));
|
|
step();
|
|
}} else {{
|
|
resolve("null");
|
|
}}
|
|
}}).catch(function(error) {{
|
|
resolve("null");
|
|
}});
|
|
}}
|
|
step();
|
|
}});
|
|
}}
|
|
|
|
return processLoop();
|
|
}}
|
|
|
|
{otel_context_inject}
|
|
|
|
// A slot is `null` only when the arg was not provided: pass `undefined` so the
|
|
// parameter default applies (JS defaults ignore `null`), matching the bun runner.
|
|
// A provided JSON `null` arrives as the string "null" and stays `null`.
|
|
let args = Deno.core.ops.op_get_static_args().map((arg) => arg === null ? undefined : JSON.parse(arg))
|
|
import("file:///eval.ts").then((module) => module.{main_fn}(...args))
|
|
.then(res => {{
|
|
if (isAsyncIterable(res)) {{
|
|
return processStreamIterative(res)
|
|
}} else {{
|
|
return JSON.stringify(res ?? null);
|
|
}}
|
|
}})
|
|
"#
|
|
),
|
|
)
|
|
.map_err(|e| ExecuteError::Script(format!("native script initialization: {e}")))?;
|
|
|
|
let fut = js_runtime.resolve(script);
|
|
let global = js_runtime
|
|
.with_event_loop_promise(fut, PollEventLoopOptions::default())
|
|
.await;
|
|
|
|
match global {
|
|
Ok(global) => {
|
|
extract_global_string(js_runtime, global).map_err(|e| ExecuteError::Script(e))
|
|
}
|
|
Err(deno_core::error::CoreError::Js(e)) => {
|
|
let source = e.frames.first().and_then(|f| {
|
|
f.file_name
|
|
.as_ref()
|
|
.map(|name| (name.clone(), f.line_number.unwrap_or(1) as usize))
|
|
});
|
|
Err(ExecuteError::Js { message: e.message, stack: e.stack, name: e.name, source })
|
|
}
|
|
Err(e) => Err(ExecuteError::Script(e.print_with_cause())),
|
|
}
|
|
}
|