diff --git a/.github/DockerfileBackendTests b/.github/DockerfileBackendTests index 8acc762451..8783bb241c 100644 --- a/.github/DockerfileBackendTests +++ b/.github/DockerfileBackendTests @@ -42,7 +42,7 @@ RUN wget https://www.python.org/ftp/python/${PYTHON_VERSION}/Python-${PYTHON_VER RUN /usr/local/bin/python3 -m pip install pip-tools # Bun -COPY --from=oven/bun:1.2.23 /usr/local/bin/bun /usr/bin/bun +COPY --from=oven/bun:1.3.8 /usr/local/bin/bun /usr/bin/bun ARG TARGETPLATFORM diff --git a/.github/workflows/backend-test.yml b/.github/workflows/backend-test.yml index 0de231e9bb..1f51ee2fa4 100644 --- a/.github/workflows/backend-test.yml +++ b/.github/workflows/backend-test.yml @@ -44,7 +44,10 @@ jobs: go-version: 1.21.5 - uses: oven-sh/setup-bun@v2 with: - bun-version: 1.1.43 + bun-version: 1.3.8 + - uses: actions/setup-node@v4 + with: + node-version: '20' - uses: astral-sh/setup-uv@v6.2.1 with: version: "0.9.24" @@ -67,6 +70,111 @@ jobs: - name: Substitute EE code (EE logic is behind feature flag) run: | ./substitute_ee_code.sh --copy --dir ./windmill-ee-private + - name: Setup private npm registry with test package + working-directory: /tmp + run: | + set -e + + # Install Verdaccio globally + npm install -g verdaccio + + # Create Verdaccio config that requires authentication for @windmill-test packages + mkdir -p /tmp/verdaccio/storage + cat > /tmp/verdaccio/config.yaml << 'VERDACCIO_CONFIG' + storage: /tmp/verdaccio/storage + auth: + htpasswd: + file: /tmp/verdaccio/htpasswd + max_users: 100 + uplinks: + npmjs: + url: https://registry.npmjs.org/ + packages: + '@windmill-test/*': + access: $authenticated + publish: $authenticated + '@*/*': + access: $all + publish: $authenticated + proxy: npmjs + '**': + access: $all + publish: $authenticated + proxy: npmjs + server: + keepAliveTimeout: 60 + middlewares: + audit: + enabled: true + log: { type: stdout, format: pretty, level: warn } + VERDACCIO_CONFIG + + # Create empty htpasswd file (users will be created via API) + touch /tmp/verdaccio/htpasswd + + # Start Verdaccio in background + verdaccio --config /tmp/verdaccio/config.yaml & + VERDACCIO_PID=$! + + # Wait for Verdaccio to be ready + echo "Waiting for Verdaccio to start..." + for i in {1..30}; do + if curl -s http://localhost:4873/-/ping > /dev/null 2>&1; then + echo "Verdaccio is ready" + break + fi + sleep 1 + done + + # Login to get a token + echo "Getting auth token..." + RESPONSE=$(curl -s -X PUT \ + -H "Content-Type: application/json" \ + -d '{"name":"testuser","password":"testpass123"}' \ + http://localhost:4873/-/user/org.couchdb.user:testuser) + + echo "Auth response: $RESPONSE" + NPM_TOKEN=$(echo "$RESPONSE" | jq -r '.token') + + if [ -z "$NPM_TOKEN" ] || [ "$NPM_TOKEN" = "null" ]; then + echo "Failed to get NPM token from response" + exit 1 + fi + + echo "NPM_TOKEN=${NPM_TOKEN}" >> $GITHUB_ENV + echo "Got NPM token successfully: ${NPM_TOKEN:0:10}..." + + # Configure npm globally with the auth token + echo "//localhost:4873/:_authToken=${NPM_TOKEN}" > ~/.npmrc + echo "Configured ~/.npmrc with auth token" + + # Create a simple test package + mkdir -p /tmp/windmill-test-private-pkg + cat > /tmp/windmill-test-private-pkg/package.json << 'PKG_JSON' + { + "name": "@windmill-test/private-pkg", + "version": "1.0.0", + "main": "index.js" + } + PKG_JSON + cat > /tmp/windmill-test-private-pkg/index.js << 'PKG_JS' + module.exports.greet = (name) => `Hello from private package, ${name}!`; + PKG_JS + + # Publish to Verdaccio with auth + cd /tmp/windmill-test-private-pkg + echo "Publishing package..." + npm publish --registry http://localhost:4873 + echo "Package published successfully" + + # Verify the package requires auth by trying anonymous access (should fail) + rm -f ~/.npmrc + echo "Testing anonymous access (should fail)..." + if npm view @windmill-test/private-pkg --registry http://localhost:4873 2>/dev/null; then + echo "ERROR: Package should require authentication but anonymous access worked" + exit 1 + fi + echo "Verified: Package requires authentication for @windmill-test/private-pkg" - name: Cache DuckDB FFI module build uses: actions/cache@v3 with: @@ -84,9 +192,10 @@ jobs: RUST_LOG_STYLE: never CARGO_NET_GIT_FETCH_WITH_CLI: true WMDEBUG_FORCE_V0_WORKSPACE_DEPENDENCIES: 1 - WMDEBUG_FORCE_RUNNABLE_SETTINGS_V0: 1 + WMDEBUG_FORCE_RUNNABLE_SETTINGS_V0: 1 WMDEBUG_FORCE_NO_LEGACY_DEBOUNCING_COMPAT: 1 + TEST_NPM_REGISTRY: "http://localhost:4873/:_authToken=${{ env.NPM_TOKEN }}" run: | - deno --version && bun -v && go version && python3 --version + deno --version && bun -v && node --version && go version && python3 --version cd windmill-duckdb-ffi-internal && ./build_dev.sh && cd .. - DENO_PATH=$(which deno) BUN_PATH=$(which bun) GO_PATH=$(which go) UV_PATH=$(which uv) cargo test --features enterprise,deno_core,duckdb,license,python,rust,scoped_cache,parquet,private --all -- --nocapture + DENO_PATH=$(which deno) BUN_PATH=$(which bun) NODE_BIN_PATH=$(which node) GO_PATH=$(which go) UV_PATH=$(which uv) cargo test --features enterprise,deno_core,duckdb,license,python,rust,scoped_cache,parquet,private,private_registry_test --all -- --nocapture diff --git a/Dockerfile b/Dockerfile index e1e30876c4..8b15e5e79e 100644 --- a/Dockerfile +++ b/Dockerfile @@ -234,7 +234,7 @@ COPY --from=windmill_duckdb_ffi_internal_builder /windmill-duckdb-ffi-internal/t COPY --from=denoland/deno:2.2.1 --chmod=755 /usr/bin/deno /usr/bin/deno -COPY --from=oven/bun:1.2.23 /usr/local/bin/bun /usr/bin/bun +COPY --from=oven/bun:1.3.8 /usr/local/bin/bun /usr/bin/bun COPY --from=php:8.3.7-cli /usr/local/bin/php /usr/bin/php COPY --from=composer:2.7.6 /usr/bin/composer /usr/bin/composer diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 802e747f7b..8a3bc32d5e 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -15527,6 +15527,7 @@ dependencies = [ "sqlx", "strum 0.27.2", "systemstat", + "tempfile", "tikv-jemalloc-ctl", "tikv-jemalloc-sys", "tikv-jemallocator", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 4f24ea5c2a..c9f11920b0 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -90,6 +90,7 @@ zip = ["windmill-api/zip"] static_frontend = ["windmill-api/static_frontend"] scoped_cache = ["windmill-common/scoped_cache"] test_job_debouncing = [] +private_registry_test = [] # Languages python = ["windmill-worker/python", "windmill-api/python"] rust = ["windmill-worker/rust"] @@ -182,6 +183,7 @@ axum.workspace = true serde.workspace = true windmill-api-client.workspace = true deno_core = { workspace = true, features = ["include_js_files_for_snapshotting", "unsafe_use_unprotected_platform"] } +tempfile.workspace = true [workspace.dependencies] diff --git a/backend/tests/bun_jobs.rs b/backend/tests/bun_jobs.rs new file mode 100644 index 0000000000..cf3ce4d350 --- /dev/null +++ b/backend/tests/bun_jobs.rs @@ -0,0 +1,1358 @@ +mod common; +use crate::common::*; +use sqlx::postgres::Postgres; +use sqlx::Pool; +use windmill_common::jobs::{JobPayload, RawCode}; +use windmill_common::scripts::ScriptLang; + +// ============================================================================ +// Basic Execution Tests +// ============================================================================ + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_simple(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#" +export function main() { + return "hello world"; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!("hello world")); + Ok(()) +} + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_with_args(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#" +export function main(name: string, count: number) { + return `Hello ${name}, count: ${count}`; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = RunJob::from(job) + .arg("name", serde_json::json!("World")) + .arg("count", serde_json::json!(42)) + .run_until_complete(&db, false, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!("Hello World, count: 42")); + Ok(()) +} + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_return_types(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + // Test object return + { + let content = r#" +export function main() { + return { name: "test", value: 123 }; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!({"name": "test", "value": 123})); + } + + // Test array return + { + let content = r#" +export function main() { + return [1, 2, 3, "four"]; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!([1, 2, 3, "four"])); + } + + // Test BigInt serialization + { + let content = r#" +export function main() { + // Use BigInt literal notation to avoid JavaScript number precision loss + return 9007199254740993n; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + // BigInt should be serialized as string + assert_eq!(result, serde_json::json!("9007199254740993")); + } + + Ok(()) +} + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_async(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#" +export async function main() { + await new Promise(resolve => setTimeout(resolve, 100)); + return "async completed"; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!("async completed")); + Ok(()) +} + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_null_undefined_handling(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + // Test null return + { + let content = r#" +export function main() { + return null; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!(null)); + } + + // Test undefined return (should be serialized as null) + { + let content = r#" +export function main() { + return undefined; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!(null)); + } + + Ok(()) +} + +// ============================================================================ +// Error Handling Tests +// ============================================================================ + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_runtime_error(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#" +export function main() { + throw new Error("intentional error"); +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let completed = run_job_in_new_worker_until_complete(&db, false, job, port).await; + + assert!(!completed.success); + let result = completed.json_result().unwrap(); + // Error is wrapped: {"error": {"message": "...", "name": "...", "stack": "..."}} + let error = &result["error"]; + assert!(error["message"] + .as_str() + .unwrap() + .contains("intentional error")); + Ok(()) +} + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_missing_main_function(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#" +export function notMain() { + return "hello"; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let completed = run_job_in_new_worker_until_complete(&db, false, job, port).await; + + assert!(!completed.success); + let result = completed.json_result().unwrap(); + // Error is wrapped: {"error": {"message": "...", "name": "...", "stack": "..."}} + let error = &result["error"]; + assert!(error["message"] + .as_str() + .unwrap() + .contains("main function is missing")); + Ok(()) +} + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_syntax_error(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#" +export function main() { + return "unclosed string +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let completed = run_job_in_new_worker_until_complete(&db, false, job, port).await; + + assert!(!completed.success); + Ok(()) +} + +// ============================================================================ +// Annotation Mode Tests +// ============================================================================ + +#[sqlx::test(fixtures("base"))] +async fn test_bun_nodejs_mode(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#"//nodejs + +export function main() { + // Node.js specific API + return process.version.startsWith("v"); +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!(true)); + Ok(()) +} + +#[sqlx::test(fixtures("base"))] +async fn test_bun_nobundling_mode(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#"//nobundling + +export function main() { + return "nobundling works"; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!("nobundling works")); + Ok(()) +} + +// ============================================================================ +// Native Mode Tests (requires deno_core feature) +// ============================================================================ + +#[cfg(feature = "deno_core")] +#[sqlx::test(fixtures("base"))] +async fn test_bun_native_mode(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#"//native + +export function main() { + return "native execution"; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!("native execution")); + Ok(()) +} + +// ============================================================================ +// Relative Import Tests +// ============================================================================ + +#[sqlx::test(fixtures("base", "relative_bun"))] +async fn test_bun_relative_imports(db: Pool) -> anyhow::Result<()> { + let content = r#" +import { main as test1 } from "/f/system/same_folder_script.ts"; +import { main as test2 } from "./same_folder_script.ts"; +import { main as test3 } from "/f/system_relative/different_folder_script.ts"; +import { main as test4 } from "../system_relative/different_folder_script.ts"; + +export function main() { + return [test1(), test2(), test3(), test4()]; +} +"# + .to_string(); + + run_deployed_relative_imports(&db, content.clone(), ScriptLang::Bun).await?; + run_preview_relative_imports(&db, content, ScriptLang::Bun).await?; + Ok(()) +} + +#[sqlx::test(fixtures("base", "relative_bun"))] +async fn test_bun_nested_imports(db: Pool) -> anyhow::Result<()> { + // Test with absolute path (/f/...) + let content_absolute = r#" +import { main as test } from "/f/system_relative/nested_script.ts"; + +export function main() { + return test(); +} +"# + .to_string(); + + run_deployed_relative_imports(&db, content_absolute.clone(), ScriptLang::Bun).await?; + run_preview_relative_imports(&db, content_absolute, ScriptLang::Bun).await?; + + // Test with relative path (../...) + let content_relative = r#" +import { main as test } from "../system_relative/nested_script.ts"; + +export function main() { + return test(); +} +"# + .to_string(); + + run_preview_relative_imports(&db, content_relative, ScriptLang::Bun).await?; + Ok(()) +} + +// ============================================================================ +// Deeply Nested Import Tests +// ============================================================================ + +#[sqlx::test(fixtures("base", "bun_edge_cases"))] +async fn test_bun_deeply_nested_imports(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + // Test with absolute path import (/f/nested/level1.ts) + // Note: The fixture uses relative imports internally (level1 -> ./level2 -> ./level3) + { + let content = r#" +import { main as level1 } from "/f/nested/level1.ts"; + +export function main() { + return level1(); +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: Some("f/nested/test_deep".to_string()), + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + // level1 -> level2 -> level3, each adds to the chain + assert_eq!(result, serde_json::json!("level1 -> level2 -> level3")); + } + + // Test with relative path import (./level1.ts) + { + let content = r#" +import { main as level1 } from "./level1.ts"; + +export function main() { + return level1(); +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: Some("f/nested/test_deep_relative".to_string()), + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + // Same result: level1 -> level2 -> level3 + assert_eq!(result, serde_json::json!("level1 -> level2 -> level3")); + } + + Ok(()) +} + +#[sqlx::test(fixtures("base", "bun_edge_cases"))] +async fn test_bun_shared_imports_both_styles(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + // Test importing modules that use different import styles internally: + // - module_a uses relative import: ./shared.ts + // - module_b uses absolute import: /f/circular/shared.ts + // Both should work correctly + let content = r#" +import { getValue as getA } from "/f/circular/module_a.ts"; +import { getValue as getB } from "./module_b.ts"; + +export function main() { + return [getA(), getB()]; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: Some("f/circular/test_both".to_string()), + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + // Both modules should correctly import SHARED_VALUE + assert_eq!( + result, + serde_json::json!(["from_a_shared", "from_b_shared"]) + ); + Ok(()) +} + +// ============================================================================ +// Preprocessor Tests +// ============================================================================ + +#[sqlx::test(fixtures("base"))] +async fn test_bun_preprocessor_execution(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + // Test that the main function works correctly + // Note: Preprocessor execution requires specific job configuration + // (flow_step_id != "preprocessor" and preprocessed == Some(false)) + // which is not set by default in RawCode jobs + let content = r#" +export function main(x: number) { + return x + 10; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + // x=5, main adds 10 = 15 + let result = RunJob::from(job) + .arg("x", serde_json::json!(5)) + .run_until_complete(&db, false, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result, serde_json::json!(15)); + Ok(()) +} + +// ============================================================================ +// Wmill SDK Tests +// ============================================================================ + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_with_wmill_env_vars(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + // Test accessing Windmill environment variables (which the SDK uses internally) + // This validates that the execution environment is properly configured + let content = r#" +export function main() { + // WM_WORKSPACE and WM_TOKEN are injected by Windmill worker + return { + workspace: process.env.WM_WORKSPACE, + hasToken: process.env.WM_TOKEN !== undefined && process.env.WM_TOKEN.length > 0, + baseUrl: process.env.BASE_URL !== undefined + }; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result["workspace"], serde_json::json!("test-workspace")); + assert_eq!(result["hasToken"], serde_json::json!(true)); + assert_eq!(result["baseUrl"], serde_json::json!(true)); + Ok(()) +} + +// ============================================================================ +// Environment Variable Tests +// ============================================================================ + +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_env_vars(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + let content = r#" +export function main() { + return { + workspace: process.env.WM_WORKSPACE, + hasToken: !!process.env.WM_TOKEN, + }; +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = run_job_in_new_worker_until_complete(&db, false, job, port) + .await + .json_result() + .unwrap(); + + assert_eq!(result["workspace"], serde_json::json!("test-workspace")); + assert_eq!(result["hasToken"], serde_json::json!(true)); + Ok(()) +} + +// ============================================================================ +// Dedicated Worker Protocol Tests +// ============================================================================ + +mod dedicated_worker_protocol { + use crate::common::{parse_dedicated_worker_line, DedicatedWorkerResult}; + use std::io::{BufRead, BufReader, Write}; + use std::process::{Command, Stdio}; + use windmill_worker::{ + build_loader, generate_dedicated_worker_wrapper, BUN_DEDICATED_WORKER_ARGS, LoaderMode, + BUN_PATH, NODE_BIN_PATH, + }; + + /// Creates test worker files and optionally bundles for Node.js (like production) + /// Returns the path to the wrapper file to execute + fn create_test_worker_files( + dir: &std::path::Path, + script: &str, + arg_names: &[&str], + bundle_for_node: bool, + ) -> std::path::PathBuf { + let dir_str = dir.to_str().unwrap(); + std::fs::write(dir.join("main.ts"), script).unwrap(); + + if bundle_for_node { + // For Node.js: bundle to JavaScript first (like production's build_loader with LoaderMode::Node) + let wrapper = generate_dedicated_worker_wrapper(arg_names, "./main.js", None); + std::fs::write(dir.join("wrapper.mjs"), wrapper).unwrap(); + + // Use the exact same build_loader function as production + tokio::runtime::Runtime::new() + .unwrap() + .block_on(build_loader( + dir_str, + "http://localhost:8000", + "test_token", + "test-workspace", + "f/test/script", + LoaderMode::Node, + )) + .expect("build_loader failed"); + + // Run the bundler with bun (build_loader creates node_builder.ts) + let output = Command::new(BUN_PATH.as_str()) + .args(["run", dir.join("node_builder.ts").to_str().unwrap()]) + .current_dir(dir) + .output() + .expect("Failed to run bun build"); + + if !output.status.success() { + panic!( + "Bun build failed: {}", + String::from_utf8_lossy(&output.stderr) + ); + } + + // Bun outputs to wrapper.js, rename to wrapper.mjs for ES module + let bundled_path = dir.join("wrapper.js"); + let output_path = dir.join("wrapper_bundled.mjs"); + std::fs::rename(&bundled_path, &output_path).unwrap(); + output_path + } else { + // For Bun: use TypeScript directly (like production) + let wrapper = generate_dedicated_worker_wrapper(arg_names, "./main.ts", None); + let wrapper_path = dir.join("wrapper.mjs"); + std::fs::write(&wrapper_path, wrapper).unwrap(); + wrapper_path + } + } + + /// Helper to run a dedicated worker test with given runtime + fn run_worker_test( + runtime: &str, + script: &str, + arg_names: &[&str], + jobs: Vec, + ) -> Vec> { + let temp_dir = tempfile::tempdir().unwrap(); + + // Create files and get the wrapper path (bundled for node, raw for bun) + let wrapper_path = create_test_worker_files( + temp_dir.path(), + script, + arg_names, + runtime == "node", + ); + let wrapper_str = wrapper_path.to_str().unwrap(); + + // Build args matching production behavior + let (cmd, args): (&str, Vec<&str>) = match runtime { + "bun" => { + // Production: bun run -i --prefer-offline wrapper.mjs + let mut args: Vec<&str> = BUN_DEDICATED_WORKER_ARGS.to_vec(); + args.push(wrapper_str); + (BUN_PATH.as_str(), args) + } + "node" => { + // Production: node wrapper.mjs (after bundling to JS) + (NODE_BIN_PATH.as_str(), vec![wrapper_str]) + } + _ => panic!("Unknown runtime: {}", runtime), + }; + + let mut child = Command::new(cmd) + .args(args) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .current_dir(temp_dir.path()) + .spawn() + .expect("Failed to spawn worker process"); + + let mut stdin = child.stdin.take().unwrap(); + let stdout = child.stdout.take().unwrap(); + let mut reader = BufReader::new(stdout); + + // Wait for "start" signal + let mut start_line = String::new(); + reader.read_line(&mut start_line).unwrap(); + assert_eq!( + parse_dedicated_worker_line(start_line.trim()), + DedicatedWorkerResult::Start, + "Expected 'start', got: {}", + start_line.trim() + ); + + let mut results = Vec::new(); + + for job_args in jobs { + writeln!(stdin, "{}", job_args.to_string()).unwrap(); + stdin.flush().unwrap(); + + let mut response = String::new(); + reader.read_line(&mut response).unwrap(); + + match parse_dedicated_worker_line(response.trim()) { + DedicatedWorkerResult::Success(value) => results.push(Ok(value)), + DedicatedWorkerResult::Error(err) => { + let msg = err["message"].as_str().unwrap_or("Unknown error").to_string(); + results.push(Err(msg)); + } + other => panic!("Unexpected response: {:?}", other), + } + } + + writeln!(stdin, "end").unwrap(); + stdin.flush().unwrap(); + let _ = child.wait().expect("Worker process failed to exit"); + + results + } + + // ==================== Node.js Runtime Tests ==================== + + #[test] + fn test_dedicated_worker_nodejs_simple() { + let script = r#" +export function main(x: number, y: number): number { + return x + y; +} +"#; + let results = run_worker_test( + "node", + script, + &["x", "y"], + vec![serde_json::json!({"x": 5, "y": 3})], + ); + + assert_eq!(results.len(), 1); + assert_eq!(results[0], Ok(serde_json::json!(8))); + } + + #[test] + fn test_dedicated_worker_nodejs_multiple_jobs() { + let script = r#" +export function main(n: number): number { + return n * 2; +} +"#; + let jobs: Vec = (1..=5).map(|i| serde_json::json!({"n": i})).collect(); + let results = run_worker_test("node", script, &["n"], jobs); + + assert_eq!(results.len(), 5); + for (i, result) in results.iter().enumerate() { + let expected = ((i + 1) * 2) as i64; + assert_eq!(*result, Ok(serde_json::json!(expected))); + } + } + + #[test] + fn test_dedicated_worker_nodejs_error() { + let script = r#" +export function main(msg: string): never { + throw new Error(msg); +} +"#; + let results = run_worker_test( + "node", + script, + &["msg"], + vec![serde_json::json!({"msg": "test error"})], + ); + + assert_eq!(results.len(), 1); + assert!(results[0].is_err()); + assert_eq!(results[0], Err("test error".to_string())); + } + + // ==================== Bun Runtime Tests ==================== + + #[test] + fn test_dedicated_worker_bun_simple() { + let script = r#" +export function main(x: number, y: number): number { + return x + y; +} +"#; + let results = run_worker_test( + "bun", + script, + &["x", "y"], + vec![serde_json::json!({"x": 5, "y": 3})], + ); + + assert_eq!(results.len(), 1); + assert_eq!(results[0], Ok(serde_json::json!(8))); + } + + #[test] + fn test_dedicated_worker_bun_multiple_jobs() { + let script = r#" +export function main(n: number): number { + return n * 2; +} +"#; + let jobs: Vec = (1..=5).map(|i| serde_json::json!({"n": i})).collect(); + let results = run_worker_test("bun", script, &["n"], jobs); + + assert_eq!(results.len(), 5); + for (i, result) in results.iter().enumerate() { + let expected = ((i + 1) * 2) as i64; + assert_eq!(*result, Ok(serde_json::json!(expected))); + } + } + + #[test] + fn test_dedicated_worker_bun_error() { + let script = r#" +export function main(msg: string): never { + throw new Error(msg); +} +"#; + let results = run_worker_test( + "bun", + script, + &["msg"], + vec![serde_json::json!({"msg": "test error"})], + ); + + assert_eq!(results.len(), 1); + assert!(results[0].is_err()); + assert_eq!(results[0], Err("test error".to_string())); + } +} + +// ============================================================================ +// Private Registry Tests +// ============================================================================ + +/// Test that bun can install packages from a private npm registry with authentication. +/// The registry requires auth tokens for accessing @windmill-test/* packages. +/// Requires: +/// - `private_registry_test` feature enabled +/// - `TEST_NPM_REGISTRY` environment variable set to registry URL with auth token +/// Format: `http://registry-url/:_authToken=TOKEN` +#[cfg(feature = "private_registry_test")] +#[sqlx::test(fixtures("base"))] +async fn test_bun_job_private_npm_registry(db: Pool) -> anyhow::Result<()> { + use windmill_worker::NPM_CONFIG_REGISTRY; + + let registry_url = std::env::var("TEST_NPM_REGISTRY") + .expect("TEST_NPM_REGISTRY must be set when running private_registry_test"); + + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await?; + let port = server.addr.port(); + + // Set the private registry configuration + { + let mut registry = NPM_CONFIG_REGISTRY.write().await; + *registry = Some(registry_url.clone()); + } + + let content = r#" +import { greet } from "@windmill-test/private-pkg"; + +export function main(name: string) { + return greet(name); +} +"# + .to_owned(); + + let job = JobPayload::Code(RawCode { + hash: None, + content, + path: None, + language: ScriptLang::Bun, + lock: None, + concurrency_settings: + windmill_common::runnable_settings::ConcurrencySettings::default().into(), + debouncing_settings: windmill_common::runnable_settings::DebouncingSettings::default(), + cache_ttl: None, + cache_ignore_s3_path: None, + dedicated_worker: None, + }); + + let result = RunJob::from(job) + .arg("name", serde_json::json!("World")) + .run_until_complete(&db, false, port) + .await + .json_result() + .unwrap(); + + // Clean up + { + let mut registry = NPM_CONFIG_REGISTRY.write().await; + *registry = None; + } + + assert_eq!( + result, + serde_json::json!("Hello from private package, World!") + ); + Ok(()) +} + +/// Tests for RELATIVE_BUN_BUILDER (loader_builder.bun.js) +/// These tests verify Bun's behavior for import scanning and package.json generation. +/// Purpose: Catch regressions when upgrading Bun versions. +mod bun_builder_tests { + use std::process::{Command, Stdio}; + use windmill_worker::{BUN_PATH, RELATIVE_BUN_BUILDER, RELATIVE_BUN_LOADER}; + + /// Run the builder and return the generated package.json content + fn run_builder(main_ts_content: &str) -> serde_json::Value { + let temp_dir = tempfile::tempdir().unwrap(); + let dir = temp_dir.path(); + + // Write main.ts + std::fs::write(dir.join("main.ts"), main_ts_content).unwrap(); + + // Write build.js using the loader and builder constants directly + // Parameters are dummy values since tests don't use Windmill relative imports + let loader = RELATIVE_BUN_LOADER + .replace("W_ID", "test-workspace") + .replace("BASE_INTERNAL_URL", "http://localhost:8000") + .replace("TOKEN", "test-token") + .replace("CURRENT_PATH", "f/test/script") + .replace("RAW_GET_ENDPOINT", "raw"); + + let build_script = format!( + r#" +{loader} + +{RELATIVE_BUN_BUILDER} +"# + ); + std::fs::write(dir.join("build.js"), build_script).unwrap(); + + // Run bun build.js + let output = Command::new(BUN_PATH.as_str()) + .args(["run", "build.js"]) + .current_dir(dir) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .output() + .expect("Failed to run bun"); + + if !output.status.success() { + panic!( + "Builder failed:\nstdout: {}\nstderr: {}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + } + + // Read generated package.json + let package_json = std::fs::read_to_string(dir.join("package.json")) + .expect("package.json not generated"); + + serde_json::from_str(&package_json).expect("Invalid JSON in package.json") + } + + /// Test: scanImports() detects basic imports + #[test] + fn test_builder_simple_import() { + let main_ts = r#" +import lodash from "lodash"; +export function main() { return lodash; } +"#; + let pkg = run_builder(main_ts); + let deps = pkg["dependencies"].as_object().unwrap(); + + assert!(deps.contains_key("lodash"), "lodash should be in dependencies"); + assert_eq!(deps["lodash"], "latest"); + } + + /// Test: scanImports() preserves version info from versioned imports + #[test] + fn test_builder_versioned_import() { + let main_ts = r#" +import _ from "lodash@4.17.21"; +export function main() { return _; } +"#; + let pkg = run_builder(main_ts); + let deps = pkg["dependencies"].as_object().unwrap(); + + assert!(deps.contains_key("lodash"), "lodash should be in dependencies"); + assert_eq!(deps["lodash"], "4.17.21"); + } + + /// Test: scanImports() handles @scope/package correctly + #[test] + fn test_builder_scoped_package() { + let main_ts = r#" +import babel from "@babel/core"; +export function main() { return babel; } +"#; + let pkg = run_builder(main_ts); + let deps = pkg["dependencies"].as_object().unwrap(); + + assert!( + deps.contains_key("@babel/core"), + "@babel/core should be in dependencies" + ); + assert_eq!(deps["@babel/core"], "latest"); + } + + /// Test: scanImports() handles multiple imports + #[test] + fn test_builder_multiple_packages() { + let main_ts = r#" +import lodash from "lodash"; +import axios from "axios"; +import dayjs from "dayjs"; +export function main() { return { lodash, axios, dayjs }; } +"#; + let pkg = run_builder(main_ts); + let deps = pkg["dependencies"].as_object().unwrap(); + + assert!(deps.contains_key("lodash"), "lodash should be in dependencies"); + assert!(deps.contains_key("axios"), "axios should be in dependencies"); + assert!(deps.contains_key("dayjs"), "dayjs should be in dependencies"); + assert_eq!(deps.len(), 3, "Should have exactly 3 dependencies"); + } + + /// Test: isBuiltin() filters out Node.js builtins + #[test] + fn test_builder_builtin_skipped() { + let main_ts = r#" +import fs from "fs"; +import path from "path"; +import lodash from "lodash"; +export function main() { return { fs, path, lodash }; } +"#; + let pkg = run_builder(main_ts); + let deps = pkg["dependencies"].as_object().unwrap(); + + assert!( + !deps.contains_key("fs"), + "fs (builtin) should NOT be in dependencies" + ); + assert!( + !deps.contains_key("path"), + "path (builtin) should NOT be in dependencies" + ); + assert!(deps.contains_key("lodash"), "lodash should be in dependencies"); + assert_eq!(deps.len(), 1, "Should have exactly 1 dependency (lodash only)"); + } + + /// Test: semver.order() resolves version conflicts (picks lowest version) + #[test] + fn test_builder_version_conflict() { + // This test simulates what happens when the same package is imported with different versions + // The builder should use semver.order() to pick the lowest version + let main_ts = r#" +import a from "lodash@4.17.21"; +import b from "lodash@4.17.10"; +export function main() { return { a, b }; } +"#; + let pkg = run_builder(main_ts); + let deps = pkg["dependencies"].as_object().unwrap(); + + assert!(deps.contains_key("lodash"), "lodash should be in dependencies"); + // The builder sorts by semver and picks the first (lowest) version + assert_eq!( + deps["lodash"], "4.17.10", + "Should resolve to lower version 4.17.10" + ); + } +} diff --git a/backend/tests/common/mod.rs b/backend/tests/common/mod.rs index 2279b2c53c..ba89de3019 100644 --- a/backend/tests/common/mod.rs +++ b/backend/tests/common/mod.rs @@ -803,3 +803,47 @@ pub async fn rebuild_dmap(client: &windmill_api_client::Client) -> bool { .status() .is_success() } + +// ============================================================================ +// Dedicated Worker Protocol Helpers +// ============================================================================ + +/// Result from parsing a dedicated worker stdout line +#[derive(Debug, Clone, PartialEq)] +pub enum DedicatedWorkerResult { + /// Worker printed "start" indicating it's ready + Start, + /// Worker returned a successful result + Success(serde_json::Value), + /// Worker returned an error result + Error(serde_json::Value), + /// Line is not a protocol message (e.g., logs) + Other(String), +} + +/// Parse a line from dedicated worker stdout according to the protocol: +/// - "start" -> Ready signal +/// - "wm_res[success]:JSON" -> Success with result +/// - "wm_res[error]:JSON" -> Error with details +/// - anything else -> Other (logs) +pub fn parse_dedicated_worker_line(line: &str) -> DedicatedWorkerResult { + if line == "start" { + return DedicatedWorkerResult::Start; + } + + if let Some(json_str) = line.strip_prefix("wm_res[success]:") { + match serde_json::from_str(json_str) { + Ok(value) => return DedicatedWorkerResult::Success(value), + Err(_) => return DedicatedWorkerResult::Other(line.to_string()), + } + } + + if let Some(json_str) = line.strip_prefix("wm_res[error]:") { + match serde_json::from_str(json_str) { + Ok(value) => return DedicatedWorkerResult::Error(value), + Err(_) => return DedicatedWorkerResult::Other(line.to_string()), + } + } + + DedicatedWorkerResult::Other(line.to_string()) +} diff --git a/backend/tests/fixtures/bun_edge_cases.sql b/backend/tests/fixtures/bun_edge_cases.sql new file mode 100644 index 0000000000..46829ec9b3 --- /dev/null +++ b/backend/tests/fixtures/bun_edge_cases.sql @@ -0,0 +1,152 @@ +-- Fixture for Bun edge case tests +-- Tests deeply nested imports (level1 -> level2 -> level3) + +-- Level 3: Base script (deepest level) +INSERT INTO public.script(workspace_id, created_by, content, schema, summary, description, path, hash, language, lock) VALUES ( +'test-workspace', +'test-user', +' +export function main() { + return "level3"; +} +', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{},"required":[],"type":"object"}', +'', +'', +'f/nested/level3', 20001, 'bun', ''); + +-- Level 2: Imports level3 using RELATIVE path (./level3.ts) +INSERT INTO public.script(workspace_id, created_by, content, schema, summary, description, path, hash, language, lock) VALUES ( +'test-workspace', +'test-user', +' +import { main as level3 } from "./level3.ts"; + +export function main() { + return "level2 -> " + level3(); +} +', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{},"required":[],"type":"object"}', +'', +'', +'f/nested/level2', 20002, 'bun', ''); + +-- Level 1: Imports level2 using RELATIVE path (./level2.ts) +INSERT INTO public.script(workspace_id, created_by, content, schema, summary, description, path, hash, language, lock) VALUES ( +'test-workspace', +'test-user', +' +import { main as level2 } from "./level2.ts"; + +export function main() { + return "level1 -> " + level2(); +} +', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{},"required":[],"type":"object"}', +'', +'', +'f/nested/level1', 20003, 'bun', ''); + +-- Script with preprocessor function +INSERT INTO public.script(workspace_id, created_by, content, schema, summary, description, path, hash, language, lock, has_preprocessor) VALUES ( +'test-workspace', +'test-user', +' +export function preprocessor(value: number) { + return { value: value * 2 }; +} + +export function main(value: number) { + return value + 100; +} +', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{"value":{"type":"number"}},"required":["value"],"type":"object"}', +'Script with preprocessor', +'', +'f/edge_cases/with_preprocessor', 20004, 'bun', '', true); + +-- Script with nodejs annotation +INSERT INTO public.script(workspace_id, created_by, content, schema, summary, description, path, hash, language, lock) VALUES ( +'test-workspace', +'test-user', +'//nodejs + +export function main() { + return process.version; +} +', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{},"required":[],"type":"object"}', +'NodeJS mode script', +'', +'f/edge_cases/nodejs_mode', 20005, 'bun', ''); + +-- Script with nobundling annotation +INSERT INTO public.script(workspace_id, created_by, content, schema, summary, description, path, hash, language, lock) VALUES ( +'test-workspace', +'test-user', +'//nobundling + +export function main() { + return "no bundle"; +} +', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{},"required":[],"type":"object"}', +'No bundling mode script', +'', +'f/edge_cases/nobundling_mode', 20006, 'bun', ''); + +-- Script that uses circular-ish import pattern (A imports B, B imports C, test imports A and C) +INSERT INTO public.script(workspace_id, created_by, content, schema, summary, description, path, hash, language, lock) VALUES ( +'test-workspace', +'test-user', +' +export const SHARED_VALUE = "shared"; + +export function main() { + return SHARED_VALUE; +} +', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{},"required":[],"type":"object"}', +'', +'', +'f/circular/shared', 20007, 'bun', ''); + +-- module_a uses RELATIVE path import (./shared.ts) +INSERT INTO public.script(workspace_id, created_by, content, schema, summary, description, path, hash, language, lock) VALUES ( +'test-workspace', +'test-user', +' +import { SHARED_VALUE } from "./shared.ts"; + +export function getValue() { + return "from_a_" + SHARED_VALUE; +} + +export function main() { + return getValue(); +} +', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{},"required":[],"type":"object"}', +'', +'', +'f/circular/module_a', 20008, 'bun', ''); + +-- module_b uses ABSOLUTE path import (/f/circular/shared.ts) +INSERT INTO public.script(workspace_id, created_by, content, schema, summary, description, path, hash, language, lock) VALUES ( +'test-workspace', +'test-user', +' +import { SHARED_VALUE } from "/f/circular/shared.ts"; + +export function getValue() { + return "from_b_" + SHARED_VALUE; +} + +export function main() { + return getValue(); +} +', +'{"$schema":"https://json-schema.org/draft/2020-12/schema","properties":{},"required":[],"type":"object"}', +'', +'', +'f/circular/module_b', 20009, 'bun', ''); diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index aff0b6d650..e15aa09b46 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -16,13 +16,13 @@ use crate::{ common::{ build_command_with_isolation, create_args_and_out_file, get_reserved_variables, parse_npm_config, read_file, read_file_content, read_result, start_child_process, - write_file_binary, MaybeLock, OccupancyMetrics, StreamNotifier, - DEV_CONF_NSJAIL, + write_file_binary, MaybeLock, OccupancyMetrics, StreamNotifier, DEV_CONF_NSJAIL, }, + get_proxy_envs_for_lang, handle_child::handle_child, BUNFIG_INSTALL_SCOPES, BUN_BUNDLE_CACHE_DIR, BUN_CACHE_DIR, BUN_NO_CACHE, BUN_PATH, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NODE_BIN_PATH, NODE_PATH, NPM_CONFIG_REGISTRY, - NPM_PATH, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, TRACING_PROXY_CA_CERT_PATH, TZ_ENV, get_proxy_envs_for_lang, + NPM_PATH, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, TRACING_PROXY_CA_CERT_PATH, TZ_ENV, }; use windmill_common::{ client::AuthedClient, @@ -51,9 +51,9 @@ use windmill_common::s3_helpers::attempt_fetch_bytes; use windmill_parser::Typ; -const RELATIVE_BUN_LOADER: &str = include_str!("../loader.bun.js"); +pub const RELATIVE_BUN_LOADER: &str = include_str!("../loader.bun.js"); -const RELATIVE_BUN_BUILDER: &str = include_str!("../loader_builder.bun.js"); +pub const RELATIVE_BUN_BUILDER: &str = include_str!("../loader_builder.bun.js"); const NSJAIL_CONFIG_RUN_BUN_CONTENT: &str = include_str!("../nsjail/run.bun.config.proto"); @@ -64,6 +64,62 @@ pub const BUN_LOCKB_SPLIT_WINDOWS: &str = "\r\n//bun.lockb\r\n"; pub const EMPTY_FILE: &str = ""; +/// Bun args for dedicated worker (without the script path) +pub const BUN_DEDICATED_WORKER_ARGS: &[&str] = &["run", "-i", "--prefer-offline"]; + +/// Generate the dedicated worker wrapper content. +/// - `arg_names`: The argument names for the main function (e.g., ["x", "y"]) +/// - `main_import`: The import path for the main module (e.g., "./main.ts") +/// - `date_conversions`: Optional date conversion statements for Datetime args +pub fn generate_dedicated_worker_wrapper( + arg_names: &[&str], + main_import: &str, + date_conversions: Option<&str>, +) -> String { + let spread = arg_names.join(","); + let dates = date_conversions.unwrap_or(""); + let is_debug = std::env::var("RUST_LOG").is_ok_and(|x| x == "windmill=debug"); + let print_lines = if is_debug { + r#"console.log(line);"# + } else { + "" + }; + + format!( + r#" +import * as Main from "{main_import}"; +import * as Readline from "node:readline" + +BigInt.prototype.toJSON = function () {{ + return this.toString(); +}}; + +console.log('start'); + +function getArgs(line) {{ + let {{ {spread} }} = JSON.parse(line) + {dates} + return [ {spread} ]; +}} + +for await (const line of Readline.createInterface({{ input: process.stdin }})) {{ + {print_lines} + + if (line === "end") {{ + process.exit(0); + }} + try {{ + const args = getArgs(line); + const res = await Main.main(...args); + console.log("wm_res[success]:" + JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value)); + }} catch (e) {{ + console.log("wm_res[error]:" + JSON.stringify({{ message: e.message, name: e.name, stack: e.stack, line: line }})); + }} +}} +"# + ) +} + /// Returns (package.json, bun.lock(b), is_empty, is_binary) fn split_lockfile(lockfile: &str) -> (&str, Option<&str>, bool, bool) { if let Some(index) = lockfile.find(BUN_LOCK_SPLIT) { @@ -115,24 +171,25 @@ pub async fn gen_bun_lockfile( gen_bunfig(job_dir).await?; write_file(job_dir, "package.json", package_json_content.as_str())?; } else { + let loader = RELATIVE_BUN_LOADER + .replace("W_ID", w_id) + .replace("BASE_INTERNAL_URL", base_internal_url) + .replace("TOKEN", token) + .replace( + "CURRENT_PATH", + &crate::common::use_flow_root_path(script_path), + ) + .replace("RAW_GET_ENDPOINT", "raw"); + write_file( &job_dir, "build.js", &format!( r#" -{} +{loader} {RELATIVE_BUN_BUILDER} -"#, - RELATIVE_BUN_LOADER - .replace("W_ID", w_id) - .replace("BASE_INTERNAL_URL", base_internal_url) - .replace("TOKEN", token) - .replace( - "CURRENT_PATH", - &crate::common::use_flow_root_path(script_path) - ) - .replace("RAW_GET_ENDPOINT", "raw") +"# ), )?; @@ -382,14 +439,14 @@ pub async fn install_bun_lockfile( } #[derive(PartialEq)] -enum LoaderMode { +pub enum LoaderMode { Node, Bun, BunBundle, NodeBundle, BrowserBundle, } -async fn build_loader( +pub async fn build_loader( job_dir: &str, base_internal_url: &str, token: &str, @@ -406,13 +463,14 @@ async fn build_loader( &crate::common::use_flow_root_path(current_path), ) .replace("RAW_GET_ENDPOINT", "raw_unpinned"); + if mode == LoaderMode::Node { write_file( &job_dir, "node_builder.ts", &format!( r#" -{} +{loader} import {{ readdir }} from "node:fs/promises"; @@ -420,7 +478,6 @@ let fileNames = [] try {{ fileNames = await readdir("{job_dir}/node_modules") }} catch (e) {{ - }} try {{ @@ -437,8 +494,7 @@ try {{ console.log("Failed to build node bundle"); process.exit(1); }} -"#, - loader +"# ), )?; } else if mode == LoaderMode::Bun { @@ -449,11 +505,10 @@ try {{ r#" import {{ plugin }} from "bun"; -{} +{loader} plugin(p) -"#, - loader +"# ), )?; } else if mode == LoaderMode::BunBundle @@ -465,7 +520,7 @@ plugin(p) "node_builder.ts", &format!( r#" -{} +{loader} try {{ await Bun.build({{ @@ -486,7 +541,6 @@ try {{ process.exit(1); }} "#, - loader, if mode == LoaderMode::BunBundle { "bun" } else if mode == LoaderMode::NodeBundle { @@ -1726,55 +1780,21 @@ pub async fn start_worker( .map(|x| return format!("{x} = {x} ? new Date({x}) : undefined")) .join("\n"); - let spread = args.into_iter().map(|x| x.name).join(","); + let arg_names: Vec<&str> = args.iter().map(|x| x.name.as_str()).collect(); // logs.push_str(format!("infer args: {:?}\n", start.elapsed().as_micros()).as_str()); // we cannot use Bun.read and Bun.write because it results in an EBADF error on cloud - let is_debug = std::env::var("RUST_LOG").is_ok_and(|x| x == "windmill=debug"); - let print_lines = if is_debug { - r#"console.log(line);"# - } else { - "" - }; - let main_import = if codebase.is_some() { "./main.js" } else { "./main.ts" }; - let wrapper_content: String = format!( - r#" -import * as Main from "{main_import}"; -import * as Readline from "node:readline" - -BigInt.prototype.toJSON = function () {{ - return this.toString(); -}}; - -console.log('start'); - -function getArgs(line) {{ - let {{ {spread} }} = JSON.parse(line) - {dates} - return [ {spread} ]; -}} - -for await (const line of Readline.createInterface({{ input: process.stdin }})) {{ - {print_lines} - - if (line === "end") {{ - process.exit(0); - }} - try {{ - const args = getArgs(line); - const res = await Main.main(...args); - console.log("wm_res[success]:" + JSON.stringify(res ?? null, (key, value) => typeof value === 'undefined' ? null : value)); - }} catch (e) {{ - console.log("wm_res[error]:" + JSON.stringify({{ message: e.message, name: e.name, stack: e.stack, line: line }})); - }} -}} -"#, - ); + let dates_opt = if dates.is_empty() { + None + } else { + Some(dates.as_str()) + }; + let wrapper_content = generate_dedicated_worker_wrapper(&arg_names, main_import, dates_opt); write_file(job_dir, "wrapper.mjs", &wrapper_content)?; } @@ -1865,3 +1885,111 @@ for await (const line of Readline.createInterface({{ input: process.stdin }})) { .await } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_split_lockfile_text_unix() { + let lockfile = r#"{"dependencies":{"lodash":"^4.17.21"}} +//bun.lock +lockfile-content-here"#; + + let (pkg, lock, is_empty, is_binary) = split_lockfile(lockfile); + + assert_eq!(pkg, r#"{"dependencies":{"lodash":"^4.17.21"}}"#); + assert_eq!(lock, Some("lockfile-content-here")); + assert!(!is_empty); + assert!(!is_binary); + } + + #[test] + fn test_split_lockfile_text_windows() { + let lockfile = "{\"dependencies\":{}}\r\n//bun.lock\r\nlockfile-content"; + + let (pkg, lock, is_empty, is_binary) = split_lockfile(lockfile); + + assert_eq!(pkg, "{\"dependencies\":{}}"); + assert_eq!(lock, Some("lockfile-content")); + assert!(!is_empty); + assert!(!is_binary); + } + + #[test] + fn test_split_lockfile_binary_unix() { + let lockfile = r#"{"dependencies":{}} +//bun.lockb +YmluYXJ5LWNvbnRlbnQ="#; // base64 encoded "binary-content" + + let (pkg, lock, is_empty, is_binary) = split_lockfile(lockfile); + + assert_eq!(pkg, r#"{"dependencies":{}}"#); + assert_eq!(lock, Some("YmluYXJ5LWNvbnRlbnQ=")); + assert!(!is_empty); + assert!(is_binary); + } + + #[test] + fn test_split_lockfile_binary_windows() { + let lockfile = "{\"dependencies\":{}}\r\n//bun.lockb\r\nYmluYXJ5LWNvbnRlbnQ="; + + let (pkg, lock, is_empty, is_binary) = split_lockfile(lockfile); + + assert_eq!(pkg, "{\"dependencies\":{}}"); + assert_eq!(lock, Some("YmluYXJ5LWNvbnRlbnQ=")); + assert!(!is_empty); + assert!(is_binary); + } + + #[test] + fn test_split_lockfile_empty() { + let lockfile = r#"{"dependencies":{}} +//bun.lock +"#; + + let (pkg, lock, is_empty, is_binary) = split_lockfile(lockfile); + + assert_eq!(pkg, r#"{"dependencies":{}}"#); + assert_eq!(lock, Some(EMPTY_FILE)); + assert!(is_empty); + assert!(!is_binary); + } + + #[test] + fn test_split_lockfile_no_lock() { + let lockfile = r#"{"dependencies":{"lodash":"^4.17.21"}}"#; + + let (pkg, lock, is_empty, is_binary) = split_lockfile(lockfile); + + assert_eq!(pkg, r#"{"dependencies":{"lodash":"^4.17.21"}}"#); + assert!(lock.is_none()); + assert!(!is_empty); + assert!(!is_binary); + } + + #[test] + fn test_split_lockfile_multiline_package_json() { + let lockfile = r#"{ + "dependencies": { + "lodash": "^4.17.21" + } +} +//bun.lock +lockfile-content"#; + + let (pkg, lock, is_empty, is_binary) = split_lockfile(lockfile); + + assert_eq!( + pkg, + r#"{ + "dependencies": { + "lodash": "^4.17.21" + } +}"# + ); + assert_eq!(lock, Some("lockfile-content")); + assert!(!is_empty); + assert!(!is_binary); + } +} diff --git a/backend/windmill-worker/src/lib.rs b/backend/windmill-worker/src/lib.rs index fe6ca82b09..fed4e8d31b 100644 --- a/backend/windmill-worker/src/lib.rs +++ b/backend/windmill-worker/src/lib.rs @@ -95,8 +95,9 @@ pub use otel_tracing_proxy_ee::{load_internal_otel_exporter, DENO_OTEL_INITIALIZ pub use result_processor::handle_job_error; pub use bun_executor::{ - compute_bundle_local_and_remote_path, get_common_bun_proc_envs, install_bun_lockfile, - prebundle_bun_script, prepare_job_dir, + build_loader, compute_bundle_local_and_remote_path, generate_dedicated_worker_wrapper, + get_common_bun_proc_envs, install_bun_lockfile, prebundle_bun_script, prepare_job_dir, + BUN_DEDICATED_WORKER_ARGS, LoaderMode, RELATIVE_BUN_BUILDER, RELATIVE_BUN_LOADER, }; pub use deno_executor::generate_deno_lock; pub use prepare_deps::run_prepare_deps_cli; diff --git a/docker/DockerfileSlim b/docker/DockerfileSlim index bd07238b4c..84420b8be2 100644 --- a/docker/DockerfileSlim +++ b/docker/DockerfileSlim @@ -34,7 +34,7 @@ RUN mkdir -p /tmp/windmill/cache && \ rm -rf /tmp/build_cache && \ mkdir -p -m 777 /tmp/windmill/cache/uv -COPY --from=oven/bun:1.2.23 /usr/local/bin/bun /usr/bin/bun +COPY --from=oven/bun:1.3.8 /usr/local/bin/bun /usr/bin/bun # add the docker client to call docker from a worker if enabled COPY --from=docker:dind /usr/local/bin/docker /usr/local/bin/ diff --git a/docker/DockerfileSlimEe b/docker/DockerfileSlimEe index b127d1b782..80de5fd0c9 100644 --- a/docker/DockerfileSlimEe +++ b/docker/DockerfileSlimEe @@ -34,7 +34,7 @@ RUN mkdir -p /tmp/windmill/cache && \ rm -rf /tmp/build_cache && \ mkdir -p -m 777 /tmp/windmill/cache/uv -COPY --from=oven/bun:1.2.23 /usr/local/bin/bun /usr/bin/bun +COPY --from=oven/bun:1.3.8 /usr/local/bin/bun /usr/bin/bun # add the docker client to call docker from a worker if enabled COPY --from=docker:dind /usr/local/bin/docker /usr/local/bin/