diff --git a/backend/.sqlx/query-6dafcc89668fb0e5740f23264b515b1724c36031da39d53ae6c329a479bdf8aa.json b/backend/.sqlx/query-0987f164fe3d64bf0a6a4e9699c2d1f339670da3ed1c08203d54d45798d022dd.json similarity index 83% rename from backend/.sqlx/query-6dafcc89668fb0e5740f23264b515b1724c36031da39d53ae6c329a479bdf8aa.json rename to backend/.sqlx/query-0987f164fe3d64bf0a6a4e9699c2d1f339670da3ed1c08203d54d45798d022dd.json index 1c0035f7cb..8aecbed333 100644 --- a/backend/.sqlx/query-6dafcc89668fb0e5740f23264b515b1724c36031da39d53ae6c329a479bdf8aa.json +++ b/backend/.sqlx/query-0987f164fe3d64bf0a6a4e9699c2d1f339670da3ed1c08203d54d45798d022dd.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT\n external_id,\n workspace_id,\n service_name AS \"service_name!: ServiceName\",\n script_path,\n is_flow,\n webhook_token_hash,\n service_config,\n error,\n created_at,\n updated_at,\n summary\n FROM\n native_trigger\n WHERE\n workspace_id = $1\n AND service_name = $2\n AND script_path = $3\n AND is_flow = $4\n LIMIT 1\n ", + "query": "\n SELECT\n external_id,\n workspace_id,\n service_name AS \"service_name!: ServiceName\",\n script_path,\n is_flow,\n webhook_token_hash,\n service_config,\n error,\n created_at,\n updated_at,\n summary,\n enabled\n FROM\n native_trigger\n WHERE\n workspace_id = $1\n AND service_name = $2\n AND script_path = $3\n AND is_flow = $4\n LIMIT 1\n ", "describe": { "columns": [ { @@ -68,6 +68,11 @@ "ordinal": 10, "name": "summary", "type_info": "Varchar" + }, + { + "ordinal": 11, + "name": "enabled", + "type_info": "Bool" } ], "parameters": { @@ -100,8 +105,9 @@ true, false, false, - true + true, + false ] }, - "hash": "6dafcc89668fb0e5740f23264b515b1724c36031da39d53ae6c329a479bdf8aa" + "hash": "0987f164fe3d64bf0a6a4e9699c2d1f339670da3ed1c08203d54d45798d022dd" } diff --git a/backend/.sqlx/query-1048d1c95270ce1f36c02bce31a2bc8a88935c613bd213b7156299811377db8e.json b/backend/.sqlx/query-64dce306efc0d9989542dba5bf003dc94b9337a8d3b69805308258d1fa145e63.json similarity index 63% rename from backend/.sqlx/query-1048d1c95270ce1f36c02bce31a2bc8a88935c613bd213b7156299811377db8e.json rename to backend/.sqlx/query-64dce306efc0d9989542dba5bf003dc94b9337a8d3b69805308258d1fa145e63.json index 89a02933e0..d234750140 100644 --- a/backend/.sqlx/query-1048d1c95270ce1f36c02bce31a2bc8a88935c613bd213b7156299811377db8e.json +++ b/backend/.sqlx/query-64dce306efc0d9989542dba5bf003dc94b9337a8d3b69805308258d1fa145e63.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n INSERT INTO native_trigger (\n external_id,\n workspace_id,\n service_name,\n script_path,\n is_flow,\n webhook_token_hash,\n service_config,\n summary\n ) VALUES (\n $1, $2, $3, $4, $5, $6, $7, $8\n )\n ON CONFLICT (external_id, workspace_id, service_name)\n DO UPDATE SET script_path = $4, is_flow = $5, webhook_token_hash = $6, service_config = $7, summary = $8, error = NULL, updated_at = NOW()\n ", + "query": "\n INSERT INTO native_trigger (\n external_id,\n workspace_id,\n service_name,\n script_path,\n is_flow,\n webhook_token_hash,\n service_config,\n summary,\n enabled\n ) VALUES (\n $1, $2, $3, $4, $5, $6, $7, $8, $9\n )\n ON CONFLICT (external_id, workspace_id, service_name)\n DO UPDATE SET script_path = $4, is_flow = $5, webhook_token_hash = $6, service_config = $7, summary = $8, error = NULL, updated_at = NOW()\n ", "describe": { "columns": [], "parameters": { @@ -23,10 +23,11 @@ "Bool", "Varchar", "Jsonb", - "Varchar" + "Varchar", + "Bool" ] }, "nullable": [] }, - "hash": "1048d1c95270ce1f36c02bce31a2bc8a88935c613bd213b7156299811377db8e" + "hash": "64dce306efc0d9989542dba5bf003dc94b9337a8d3b69805308258d1fa145e63" } diff --git a/backend/.sqlx/query-b0775af41a9b54cce040bf37cae770e63dab12a1383a0855d105c061e4e4ca48.json b/backend/.sqlx/query-a2122e8520268919e2ddc85ef46b5f9322229bb2a1e104b6eabaf2a697c2776a.json similarity index 68% rename from backend/.sqlx/query-b0775af41a9b54cce040bf37cae770e63dab12a1383a0855d105c061e4e4ca48.json rename to backend/.sqlx/query-a2122e8520268919e2ddc85ef46b5f9322229bb2a1e104b6eabaf2a697c2776a.json index 5889ab4ee7..0f43afed39 100644 --- a/backend/.sqlx/query-b0775af41a9b54cce040bf37cae770e63dab12a1383a0855d105c061e4e4ca48.json +++ b/backend/.sqlx/query-a2122e8520268919e2ddc85ef46b5f9322229bb2a1e104b6eabaf2a697c2776a.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT\n nt.external_id,\n nt.workspace_id,\n nt.service_name AS \"service_name!: ServiceName\",\n nt.script_path,\n nt.is_flow,\n nt.webhook_token_hash,\n nt.service_config,\n nt.error,\n nt.created_at,\n nt.updated_at,\n nt.summary\n FROM\n native_trigger nt\n WHERE\n nt.workspace_id = $1 AND\n nt.service_name = $2 AND\n ($5::text IS NULL OR nt.script_path = $5) AND\n ($6::bool IS NULL OR nt.is_flow = $6) AND\n (\n (nt.is_flow = false AND EXISTS (\n SELECT 1 FROM script s\n WHERE s.workspace_id = nt.workspace_id\n AND s.path = nt.script_path\n ))\n OR\n (nt.is_flow = true AND EXISTS (\n SELECT 1 FROM flow f\n WHERE f.workspace_id = nt.workspace_id\n AND f.path = nt.script_path\n ))\n )\n LIMIT $3\n OFFSET $4\n ", + "query": "\n SELECT\n nt.external_id,\n nt.workspace_id,\n nt.service_name AS \"service_name!: ServiceName\",\n nt.script_path,\n nt.is_flow,\n nt.webhook_token_hash,\n nt.service_config,\n nt.error,\n nt.created_at,\n nt.updated_at,\n nt.summary,\n nt.enabled\n FROM\n native_trigger nt\n WHERE\n nt.workspace_id = $1 AND\n nt.service_name = $2 AND\n ($5::text IS NULL OR nt.script_path = $5) AND\n ($6::bool IS NULL OR nt.is_flow = $6) AND\n (\n (nt.is_flow = false AND EXISTS (\n SELECT 1 FROM script s\n WHERE s.workspace_id = nt.workspace_id\n AND s.path = nt.script_path\n ))\n OR\n (nt.is_flow = true AND EXISTS (\n SELECT 1 FROM flow f\n WHERE f.workspace_id = nt.workspace_id\n AND f.path = nt.script_path\n ))\n )\n LIMIT $3\n OFFSET $4\n ", "describe": { "columns": [ { @@ -68,6 +68,11 @@ "ordinal": 10, "name": "summary", "type_info": "Varchar" + }, + { + "ordinal": 11, + "name": "enabled", + "type_info": "Bool" } ], "parameters": { @@ -102,8 +107,9 @@ true, false, false, - true + true, + false ] }, - "hash": "b0775af41a9b54cce040bf37cae770e63dab12a1383a0855d105c061e4e4ca48" + "hash": "a2122e8520268919e2ddc85ef46b5f9322229bb2a1e104b6eabaf2a697c2776a" } diff --git a/backend/.sqlx/query-a3871035319012d72679f132696146d5275cdd4313961c76223ed0f3fec7dca3.json b/backend/.sqlx/query-a3871035319012d72679f132696146d5275cdd4313961c76223ed0f3fec7dca3.json new file mode 100644 index 0000000000..5f2bb07e24 --- /dev/null +++ b/backend/.sqlx/query-a3871035319012d72679f132696146d5275cdd4313961c76223ed0f3fec7dca3.json @@ -0,0 +1,35 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT enabled\n FROM native_trigger\n WHERE\n workspace_id = $1\n AND service_name = $2\n AND external_id = $3\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "enabled", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Text", + { + "Custom": { + "name": "native_trigger_service", + "kind": { + "Enum": [ + "nextcloud", + "google", + "github" + ] + } + } + }, + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "a3871035319012d72679f132696146d5275cdd4313961c76223ed0f3fec7dca3" +} diff --git a/backend/.sqlx/query-15014ce696cf2af4f719a537a4e3ca5b322cc130a35a91f8b8854f5ebdf25ad2.json b/backend/.sqlx/query-c7766afaea3e187824698cae2b090af9f49ded98cd1824f1ac91cc5dbe709bc1.json similarity index 84% rename from backend/.sqlx/query-15014ce696cf2af4f719a537a4e3ca5b322cc130a35a91f8b8854f5ebdf25ad2.json rename to backend/.sqlx/query-c7766afaea3e187824698cae2b090af9f49ded98cd1824f1ac91cc5dbe709bc1.json index 0af901b2b8..6a9de7e35b 100644 --- a/backend/.sqlx/query-15014ce696cf2af4f719a537a4e3ca5b322cc130a35a91f8b8854f5ebdf25ad2.json +++ b/backend/.sqlx/query-c7766afaea3e187824698cae2b090af9f49ded98cd1824f1ac91cc5dbe709bc1.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT\n external_id,\n workspace_id,\n service_name AS \"service_name!: ServiceName\",\n script_path,\n is_flow,\n webhook_token_hash,\n service_config,\n error,\n created_at,\n updated_at,\n summary\n FROM\n native_trigger\n WHERE\n workspace_id = $1\n AND service_name = $2\n AND external_id = $3\n ", + "query": "\n SELECT\n external_id,\n workspace_id,\n service_name AS \"service_name!: ServiceName\",\n script_path,\n is_flow,\n webhook_token_hash,\n service_config,\n error,\n created_at,\n updated_at,\n summary,\n enabled\n FROM\n native_trigger\n WHERE\n workspace_id = $1\n AND service_name = $2\n AND external_id = $3\n ", "describe": { "columns": [ { @@ -68,6 +68,11 @@ "ordinal": 10, "name": "summary", "type_info": "Varchar" + }, + { + "ordinal": 11, + "name": "enabled", + "type_info": "Bool" } ], "parameters": { @@ -99,8 +104,9 @@ true, false, false, - true + true, + false ] }, - "hash": "15014ce696cf2af4f719a537a4e3ca5b322cc130a35a91f8b8854f5ebdf25ad2" + "hash": "c7766afaea3e187824698cae2b090af9f49ded98cd1824f1ac91cc5dbe709bc1" } diff --git a/backend/.sqlx/query-fa42f89d9494056e1e7ad904b844203f1132aa5ae5eff9194c45e0259f4bee76.json b/backend/.sqlx/query-fa42f89d9494056e1e7ad904b844203f1132aa5ae5eff9194c45e0259f4bee76.json new file mode 100644 index 0000000000..bf8fa238a0 --- /dev/null +++ b/backend/.sqlx/query-fa42f89d9494056e1e7ad904b844203f1132aa5ae5eff9194c45e0259f4bee76.json @@ -0,0 +1,28 @@ +{ + "db_name": "PostgreSQL", + "query": "\n UPDATE native_trigger\n SET enabled = $1\n WHERE\n workspace_id = $2\n AND service_name = $3\n AND external_id = $4\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Bool", + "Text", + { + "Custom": { + "name": "native_trigger_service", + "kind": { + "Enum": [ + "nextcloud", + "google", + "github" + ] + } + } + }, + "Text" + ] + }, + "nullable": [] + }, + "hash": "fa42f89d9494056e1e7ad904b844203f1132aa5ae5eff9194c45e0259f4bee76" +} diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 3fe73c9b62..285c031f79 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -15846,6 +15846,7 @@ dependencies = [ "reqwest 0.13.4", "serde", "serde_json", + "serial_test", "sqlx", "tempfile", "tokio", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 9ef1d099ee..9e8d62316f 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -108,6 +108,10 @@ debug = "line-tables-only" [profile.dev.package."*"] debug = false +# The expression parser's fallback call chain exhausts worker stacks without optimization. +[profile.dev.package.php-parser-rs] +opt-level = 1 + [profile.release] lto = "thin" debug = "line-tables-only" diff --git a/backend/migrations/20260908113909_native_trigger_enabled.down.sql b/backend/migrations/20260908113909_native_trigger_enabled.down.sql new file mode 100644 index 0000000000..75f26b22c4 --- /dev/null +++ b/backend/migrations/20260908113909_native_trigger_enabled.down.sql @@ -0,0 +1 @@ +ALTER TABLE native_trigger DROP COLUMN IF EXISTS enabled; diff --git a/backend/migrations/20260908113909_native_trigger_enabled.up.sql b/backend/migrations/20260908113909_native_trigger_enabled.up.sql new file mode 100644 index 0000000000..15139498d7 --- /dev/null +++ b/backend/migrations/20260908113909_native_trigger_enabled.up.sql @@ -0,0 +1 @@ +ALTER TABLE native_trigger ADD COLUMN enabled BOOLEAN NOT NULL DEFAULT true; diff --git a/backend/src/main.rs b/backend/src/main.rs index f84812151d..118795bd17 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -1944,6 +1944,13 @@ async fn process_notify_event( if let Err(e) = reload_license_key(&db.into()).await { tracing::error!("Failed to reload license key: {e:#}"); } + // The worker-group cache override is Enterprise-only, and nothing else + // re-reads the plan for it: the periodic settings pass runs ahead of + // reload_license_key, so it would see the plan this event just replaced. + #[cfg(feature = "parquet")] + if worker_mode { + crate::monitor::reload_cache_object_store_override_with_retry(db).await; + } } DEFAULT_TAGS_PER_WORKSPACE_SETTING => { if let Err(e) = load_tag_per_workspace_enabled(db).await { diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 9aae2cca06..867d2c8520 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -398,6 +398,8 @@ pub async fn initial_load( additional_python_paths: None, pip_local_dependencies: None, native_mode, + // an agent worker never reads its group's config, only its token + object_store_cache_config: None, })); } } @@ -5294,6 +5296,7 @@ pub async fn reload_worker_config(db: &DB, tx: KillpillSender, kill_if_change: b .dedicated_workers .as_ref() .is_some_and(|dws| !dws.is_empty()); + if **wc != config || has_dedicated { if kill_if_change { if has_dedicated @@ -5341,6 +5344,37 @@ pub async fn reload_worker_config(db: &DB, tx: KillpillSender, kill_if_change: b store_pull_query(&config).await; WORKER_CONFIG.store(std::sync::Arc::new(config)); } + + // After the store, so a retry that wakes mid-build reads the config being applied + // rather than the one it replaced. Unconditional rather than gated on the value + // changing, so that a pass triggered by anything else — a license-plan change, most + // of all — still re-evaluates the entitlement. + #[cfg(feature = "parquet")] + reload_cache_object_store_override_with_retry(db).await; + } +} + +/// Apply this worker group's dependency-cache object store, retrying once shortly after a build +/// that failed for a reason that may pass — the periodic settings reload behind it is 12h apart, +/// which is a long time for a whole group to cache nothing but locally. +#[cfg(feature = "parquet")] +pub async fn reload_cache_object_store_override_with_retry(db: &DB) { + let settings = WORKER_CONFIG.load().object_store_cache_config.clone(); + if matches!( + windmill_object_store::reload_cache_object_store_override(db, settings).await, + ObjectStoreReload::Later + ) { + let db = db.clone(); + tokio::spawn(async move { + tokio::time::sleep(Duration::from_secs(10)).await; + if windmill_object_store::cache_object_store_override_failed().await { + // Re-read rather than reuse: the group config may have changed while we slept, + // and installing the settings this retry was born with would pin the worker to a + // store the group no longer asks for. + let settings = WORKER_CONFIG.load().object_store_cache_config.clone(); + windmill_object_store::reload_cache_object_store_override(&db, settings).await; + } + }); } } diff --git a/backend/summarized_schema.txt b/backend/summarized_schema.txt index 67758cef32..a1e0b38518 100644 --- a/backend/summarized_schema.txt +++ b/backend/summarized_schema.txt @@ -149,7 +149,7 @@ mcp_oauth_server_code: code(char), client_id(char), user_email(char), workspace_ FK: (client_id) -> mcp_oauth_server_client(client_id) metrics: id(char), value(jsonb), created_at(ts) mqtt_trigger: mqtt_resource_path(char), subscribe_topics(jsonb[]), client_version(mqtt_client_version), v5_config(jsonb), v3_config(jsonb), client_id(char), path(char), script_path(char), is_flow(bool), workspace_id(char), edited_by(char), email(char), edited_at(ts), extra_perms(jsonb), server_id(char), last_server_ping(ts), error(text), error_handler_path(char), error_handler_args(jsonb), retry(jsonb), mode(trigger_mode), labels(text[]) -native_trigger: external_id(char), workspace_id(char), service_name(native_trigger_service), script_path(char), is_flow(bool), webhook_token_hash(char), service_config(jsonb), error(text), created_at(ts), updated_at(ts) +native_trigger: external_id(char), workspace_id(char), service_name(native_trigger_service), script_path(char), is_flow(bool), webhook_token_hash(char), service_config(jsonb), error(text), created_at(ts), updated_at(ts), enabled(bool) FK: (workspace_id) -> workspace(id) nats_trigger: path(char), nats_resource_path(char), subjects(char), stream_name(char), consumer_name(char), use_jetstream(bool), script_path(char), is_flow(bool), workspace_id(char), edited_by(char), email(char), edited_at(ts), extra_perms(jsonb), server_id(char), last_server_ping(ts), error(text), error_handler_path(char), error_handler_args(jsonb), retry(jsonb), mode(trigger_mode), labels(text[]) FK: (workspace_id) -> workspace(id) diff --git a/backend/tests/agent_workers.rs b/backend/tests/agent_workers.rs index 791556b035..4d6fffa346 100644 --- a/backend/tests/agent_workers.rs +++ b/backend/tests/agent_workers.rs @@ -293,6 +293,7 @@ async fn test_agent_worker_volume_e2e(db: Pool) -> anyhow::Result<()> let lfs_config = json!({ "type": "FilesystemStorage", "root_path": storage_root, + "volume_storage": "primary", "public_resource": null, "advanced_permissions": null }); @@ -306,7 +307,11 @@ async fn test_agent_worker_volume_e2e(db: Pool) -> anyhow::Result<()> .await?; // 2. Pre-populate the volume with a file - let vol_dir = storage_dir.path().join("volumes").join("test-vol"); + let vol_dir = storage_dir + .path() + .join("volumes") + .join("test-workspace") + .join("test-vol"); std::fs::create_dir_all(&vol_dir)?; std::fs::write(vol_dir.join("hello.txt"), b"hello from volume")?; @@ -343,6 +348,7 @@ async fn test_agent_worker_volume_e2e(db: Pool) -> anyhow::Result<()> // 4. GET /file/* — download the existing file let resp = http .get(format!("{vol_base}/file/hello.txt")) + .query(&[("worker_name", "test-worker-1")]) .send() .await?; assert!( @@ -360,6 +366,7 @@ async fn test_agent_worker_volume_e2e(db: Pool) -> anyhow::Result<()> // 5. PUT /file/* — upload a new file let resp = http .put(format!("{vol_base}/file/output.txt")) + .query(&[("worker_name", "test-worker-1")]) .body(b"written by agent worker".to_vec()) .send() .await?; @@ -432,7 +439,8 @@ async fn test_agent_worker_volume_http_worker_e2e(db: Pool) -> anyhow: "type": "FilesystemStorage", "root_path": storage_root, "public_resource": null, - "advanced_permissions": null + "advanced_permissions": null, + "volume_storage": "primary" }); sqlx::query!( @@ -444,20 +452,24 @@ async fn test_agent_worker_volume_http_worker_e2e(db: Pool) -> anyhow: .await?; // 2. Pre-populate the volume with a file - let vol_dir = storage_dir.path().join("volumes").join("test-vol"); + let vol_dir = storage_dir + .path() + .join("volumes") + .join("test-workspace") + .join("test-vol"); std::fs::create_dir_all(&vol_dir)?; std::fs::write(vol_dir.join("hello.txt"), b"hello from volume")?; // 3. Push the job, then run worker with HTTP connection (bun tag) - let code = r#"// volume: test-vol /tmp/data + let code = r#"// volume: test-vol data import { readFileSync, writeFileSync, existsSync } from "fs"; export function main() { - const content = readFileSync("/tmp/data/hello.txt", "utf-8"); - writeFileSync("/tmp/data/output.txt", "written by agent worker"); + const content = readFileSync("data/hello.txt", "utf-8"); + writeFileSync("data/output.txt", "written by agent worker"); return { read_content: content, - output_exists: existsSync("/tmp/data/output.txt"), + output_exists: existsSync("data/output.txt"), }; }"#; @@ -528,6 +540,7 @@ async fn test_agent_worker_volume_release(db: Pool) -> anyhow::Result< let lfs_config = json!({ "type": "FilesystemStorage", "root_path": storage_root, + "volume_storage": "primary", "public_resource": null, "advanced_permissions": null }); diff --git a/backend/tests/nativets_dedicated.rs b/backend/tests/nativets_dedicated.rs index 1c6987a480..1e9aef3228 100644 --- a/backend/tests/nativets_dedicated.rs +++ b/backend/tests/nativets_dedicated.rs @@ -12,7 +12,7 @@ mod prewarmed_isolate_tests { use windmill_worker::{build_loader, LoaderMode, BUN_PATH}; fn default_annotation() -> NativeAnnotation { - NativeAnnotation { useragent: None, proxy: None } + NativeAnnotation::default() } /// Bundle a TypeScript script into JS suitable for `PrewarmedIsolate`. diff --git a/backend/tests/volume_tests.rs b/backend/tests/volume_tests.rs index 8ecf6e98b6..e93979e9c7 100644 --- a/backend/tests/volume_tests.rs +++ b/backend/tests/volume_tests.rs @@ -537,6 +537,64 @@ fn test_asset_kind_volume_variant() { #[cfg(feature = "parquet")] #[sqlx::test(fixtures("base"))] async fn test_volume_sql_worker_e2e(db: Pool) -> anyhow::Result<()> { + let code = r#"// volume: test-vol data + +import { readFileSync, writeFileSync, existsSync } from "fs"; + +export function main() { + const content = readFileSync("data/hello.txt", "utf-8"); + writeFileSync("data/output.txt", "written by sql worker"); + return { + read_content: content, + output_exists: existsSync("data/output.txt"), + }; +}"#; + run_volume_with_default_stack(db, ScriptLang::Bun, code).await +} + +#[cfg(all(feature = "parquet", feature = "private", feature = "php"))] +#[sqlx::test(fixtures("base"))] +async fn test_php_volume_with_default_stack(db: Pool) -> anyhow::Result<()> { + // Nested calls exercise the parser's fallback chain; flattening them weakens this guard. + let code = r#" $content, + "output_exists" => in_array("output.txt", array_values(array_diff(scandir("data"), ['.', '..']))), + ]; +}"#; + + run_volume_with_default_stack(db, ScriptLang::Php, code).await +} + +#[cfg(feature = "parquet")] +async fn run_volume_with_default_stack( + db: Pool, + language: ScriptLang, + code: &'static str, +) -> anyhow::Result<()> { + // CI raises RUST_MIN_STACK; keep the worker at Tokio's default to catch regressions. + tokio::task::spawn_blocking(move || { + tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .thread_stack_size(2 * 1024 * 1024) + .enable_all() + .build()? + .block_on(run_volume_sql_worker_e2e(db, language, code)) + }) + .await? +} + +#[cfg(feature = "parquet")] +async fn run_volume_sql_worker_e2e( + db: Pool, + language: ScriptLang, + code: &str, +) -> anyhow::Result<()> { initialize_tracing().await; let server = ApiServer::start(db.clone()).await?; let port = server.addr.port(); @@ -571,24 +629,11 @@ async fn test_volume_sql_worker_e2e(db: Pool) -> anyhow::Result<()> { std::fs::write(vol_dir.join("hello.txt"), b"hello from volume")?; // 3. Push the job and run with SQL-connected worker - let code = r#"// volume: test-vol data - -import { readFileSync, writeFileSync, existsSync } from "fs"; - -export function main() { - const content = readFileSync("data/hello.txt", "utf-8"); - writeFileSync("data/output.txt", "written by sql worker"); - return { - read_content: content, - output_exists: existsSync("data/output.txt"), - }; -}"#; - let job = JobPayload::Code(RawCode { hash: None, content: code.to_string(), path: None, - language: ScriptLang::Bun, + language, lock: None, cache_ttl: None, cache_ignore_s3_path: None, diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 419be6fdf3..22f95aefb8 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -3091,16 +3091,16 @@ async fn test_php_job(db: Pool) -> anyhow::Result<()> { let server = ApiServer::start(db.clone()).await?; let port = server.addr.port(); - let content = r#" + let content = r#"// schema_validation error::Result<()> { + if !has_masked_object_store_secret(config) { + return Ok(()); + } + let stored = sqlx::query_as!( + Config, + "SELECT name, config FROM config WHERE name = $1", + name + ) + .fetch_optional(db) + .await? + .map(|c| c.config); + restore_object_store_secrets(config, stored.as_ref()); + Ok(()) +} + +fn has_masked_object_store_secret(config: &serde_json::Value) -> bool { + let Some(store) = config.get("object_store_cache_config") else { + return false; + }; + OBJECT_STORE_SECRET_KEYS + .iter() + .any(|k| store.get(k).and_then(|v| v.as_str()) == Some(OBJECT_STORE_SECRET_MASK)) +} + +/// The half of [`restore_masked_object_store_secrets`] after the read. +fn restore_object_store_secrets( + config: &mut serde_json::Value, + stored: Option<&serde_json::Value>, +) { + let stored = stored + .and_then(|c| c.get("object_store_cache_config")) + .and_then(|v| v.as_object()) + .cloned(); + let Some(store) = config + .get_mut("object_store_cache_config") + .and_then(|v| v.as_object_mut()) + else { + return; + }; + for key in OBJECT_STORE_SECRET_KEYS { + if store.get(*key).and_then(|v| v.as_str()) != Some(OBJECT_STORE_SECRET_MASK) { + continue; + } + match stored.as_ref().and_then(|s| s.get(*key)) { + Some(secret) => store.insert(key.to_string(), secret.clone()), + // Nothing to put back: drop the mask rather than store it. + None => store.remove(*key), + }; + } +} + async fn list_worker_groups( authed: ApiAuthed, Extension(db): Extension, ) -> error::JsonResult> { - let mut configs_raw = - sqlx::query_as!(Config, "SELECT name, config FROM config WHERE name LIKE 'worker__%'") - .fetch_all(&db) - .await?; + let mut configs_raw = sqlx::query_as!( + Config, + "SELECT name, config FROM config WHERE name LIKE 'worker__%'" + ) + .fetch_all(&db) + .await?; // Remove the 'worker__' prefix from all config names for config in configs_raw.iter_mut() { if let Some(name) = &config.name { @@ -75,44 +177,12 @@ async fn list_worker_groups( } } } - // Worker-group configs are instance-global and expose env_vars_static (may hold - // secrets); a job token (capped at workspace admin) gets the obfuscated view even - // when its identity is a superadmin. See is_instance_admin (GHSA-hfh4-cx4h-3fcr). - let configs = if !is_instance_admin(&authed) { - let mut obfuscated_configs: Vec = vec![]; - for config in configs_raw { - let config_value_opt = config.config.as_object().map(|obj| obj.to_owned()); - if let Some(mut config_value) = config_value_opt { - if let Some(env_var_map) = config_value - .get("env_vars_static") - .map(|obj| obj.as_object()) - .flatten() - { - let mut new_env_var_map: serde_json::Map = - serde_json::Map::new(); - for (key, value) in env_var_map { - new_env_var_map.insert( - key.to_owned(), - // we know the value is a string here, so we to_string() it and take -2 to remove the quotes - serde_json::json!("*".repeat(value.to_string().len() - 2)), - ); - } - config_value.insert( - "env_vars_static".to_string(), - serde_json::Value::Object(new_env_var_map), - ); - } - obfuscated_configs.push(Config { - name: config.name, - config: serde_json::Value::Object(config_value), - }) - } + if !is_instance_admin(&authed) { + for config in configs_raw.iter_mut() { + obfuscate_worker_config(&mut config.config); } - obfuscated_configs - } else { - configs_raw - }; - Ok(Json(configs)) + } + Ok(Json(configs_raw)) } async fn get_config( @@ -122,10 +192,20 @@ async fn get_config( ) -> error::JsonResult> { require_devops_role(&db, &authed).await?; - let config = sqlx::query_as!(Config, "SELECT name, config FROM config WHERE name = $1", name) - .fetch_optional(&db) - .await? - .map(|c| c.config); + let mut config = sqlx::query_as!( + Config, + "SELECT name, config FROM config WHERE name = $1", + name + ) + .fetch_optional(&db) + .await? + .map(|c| c.config); + + if !is_instance_admin(&authed) { + if let Some(config) = config.as_mut() { + obfuscate_worker_config(config); + } + } Ok(Json(config)) } @@ -134,10 +214,14 @@ async fn update_config( Path(name): Path, Extension(db): Extension, authed: ApiAuthed, - Json(config): Json, + Json(mut config): Json, ) -> error::Result { require_devops_role(&db, &authed).await?; + if name.starts_with("worker__") { + restore_masked_object_store_secrets(&db, &name, &mut config).await?; + } + #[cfg(not(feature = "enterprise"))] let config = if name.starts_with("worker__") { // In CE, only allow setting worker_tags, cache_clear, init_bash, and native_mode @@ -321,9 +405,14 @@ async fn list_configs( Extension(db): Extension, ) -> error::JsonResult> { require_devops_role(&db, &authed).await?; - let configs = sqlx::query_as!(Config, "SELECT name, config FROM config") + let mut configs = sqlx::query_as!(Config, "SELECT name, config FROM config") .fetch_all(&db) .await?; + if !is_instance_admin(&authed) { + for config in configs.iter_mut() { + obfuscate_worker_config(&mut config.config); + } + } Ok(Json(configs)) } @@ -414,3 +503,49 @@ async fn list_all_dedicated_with_deps( Ok(Json(result)) } + +#[cfg(test)] +mod tests { + use super::*; + + /// The mask an obfuscated read hands out must never be storable as the credential itself: + /// a devops user who is not an instance admin edits the group from that view, and a worker + /// that cannot build its override degrades to a local-only cache without failing a job, so + /// the breakage would go unnoticed. + #[test] + fn masked_secrets_survive_a_save_from_the_obfuscated_view() { + let stored = serde_json::json!({ + "object_store_cache_config": { + "type": "S3", "bucket": "cache", "access_key": "AKIA", "secret_key": "s3cr3t" + }, + "env_vars_static": { "TOKEN": "hunter2" }, + }); + + let mut shown = stored.clone(); + obfuscate_worker_config(&mut shown); + let store = &shown["object_store_cache_config"]; + assert_eq!(store["secret_key"], OBJECT_STORE_SECRET_MASK); + assert_eq!(store["access_key"], OBJECT_STORE_SECRET_MASK); + assert_eq!(store["bucket"], "cache"); + assert_ne!(shown["env_vars_static"]["TOKEN"], "hunter2"); + + let mut saved = shown.clone(); + saved["object_store_cache_config"]["bucket"] = serde_json::json!("other"); + restore_object_store_secrets(&mut saved, Some(&stored)); + let store = &saved["object_store_cache_config"]; + assert_eq!(store["secret_key"], "s3cr3t"); + assert_eq!(store["access_key"], "AKIA"); + assert_eq!(store["bucket"], "other"); + } + + #[test] + fn a_mask_with_nothing_behind_it_is_dropped_rather_than_stored() { + let mut saved = serde_json::json!({ + "object_store_cache_config": { "type": "S3", "secret_key": OBJECT_STORE_SECRET_MASK } + }); + restore_object_store_secrets(&mut saved, None); + assert!(saved["object_store_cache_config"] + .get("secret_key") + .is_none()); + } +} diff --git a/backend/windmill-api-integration-tests/tests/native_triggers.rs b/backend/windmill-api-integration-tests/tests/native_triggers.rs index 06c47cbb14..059b3aa2d2 100644 --- a/backend/windmill-api-integration-tests/tests/native_triggers.rs +++ b/backend/windmill-api-integration-tests/tests/native_triggers.rs @@ -18,15 +18,16 @@ use windmill_common::{ variables::{build_crypt, encrypt}, }; use windmill_native_triggers::{ - classify_read_failure, decrypt_oauth_data, delete_native_trigger, - delete_workspace_integration, get_workspace_integration, + classify_read_failure, decrypt_oauth_data, delete_native_trigger, delete_workspace_integration, + get_workspace_integration, github::GitHub, google::{parse_stop_channel_params, should_renew_channel}, - http_error_status, list_native_triggers, map_external_error, + grant_refused, http_error_status, list_native_triggers, map_external_error, + native_trigger_is_enabled, nextcloud::NextCloud, - grant_refused, require_native_integration_use, store_native_trigger, - store_workspace_integration, External, ExternalReadFailure, HttpRequestError, - NativeTriggerConfig, OAuthConfig, ServiceName, + require_native_integration_use, set_native_trigger_enabled, store_native_trigger, + store_workspace_integration, update_native_trigger, External, ExternalReadFailure, + HttpRequestError, NativeTriggerConfig, OAuthConfig, ServiceName, }; // ============================================================================ @@ -456,6 +457,7 @@ async fn test_delete_integration_full_cascade(db: Pool) -> anyhow::Res &trigger_config, json!({"triggerType": "drive"}), None, + true, ) .await?; @@ -546,6 +548,7 @@ async fn test_cleanup_preserves_triggers(db: Pool) -> anyhow::Result<( &trigger_config, json!({"triggerType": "drive"}), None, + true, ) .await?; @@ -603,6 +606,7 @@ async fn test_rename_moves_native_trigger(db: Pool) -> anyhow::Result< }, json!({"event": "OCP\\Files\\Events\\Node\\NodeCreatedEvent"}), None, + true, ) .await?; // An unrelated trigger already sitting on the target path must not be reported as moved. @@ -619,6 +623,7 @@ async fn test_rename_moves_native_trigger(db: Pool) -> anyhow::Result< }, json!({"event": "OCP\\Files\\Events\\Node\\NodeCreatedEvent"}), None, + true, ) .await?; @@ -784,7 +789,10 @@ fn test_refresh_failures_blame_only_the_grant_they_refuse() { let ok = Some(StatusCode::OK); assert!(grant_refused(ok, r#"{"error":"bad_refresh_token"}"#)); assert!(grant_refused(ok, r#"{"error":"invalid_grant"}"#)); - assert!(!grant_refused(ok, r#"{"access_token":"t","token_type":"bearer"}"#)); + assert!(!grant_refused( + ok, + r#"{"access_token":"t","token_type":"bearer"}"# + )); } /// A service that is busy or broken has not refused anything, and callers react differently to @@ -826,7 +834,9 @@ fn test_transient_service_failures_are_not_refusals() { body: r#"{"message":"Must have admin rights to Repository."}"#.to_string(), }); assert!( - map_external_error(refused).to_string().contains("admin rights"), + map_external_error(refused) + .to_string() + .contains("admin rights"), "a real 403 keeps its guidance" ); } @@ -850,3 +860,115 @@ fn test_only_service_failures_degrade_the_read() { "a non-provider error must pass through unmapped" ); } + +/// The pause switch a webhook delivery is gated on. A trigger arrives enabled, survives an +/// unrelated edit, and an unknown one reads as enabled so a delivery Windmill cannot place is +/// never silently dropped. +#[sqlx::test(migrations = "../migrations", fixtures("base"))] +async fn test_native_trigger_enabled_toggle(db: Pool) -> anyhow::Result<()> { + insert_test_script(&db, "f/test/handler").await?; + let config = NativeTriggerConfig { + script_path: "f/test/handler".to_string(), + is_flow: false, + webhook_token: "abcdefghij1234567890".to_string(), + }; + store_native_trigger( + &db, + "test-workspace", + ServiceName::Nextcloud, + "ext-1", + &config, + json!({"event": "OCA\\Files\\Event\\LoadAdditionalScriptsEvent"}), + None, + true, + ) + .await?; + + assert!( + native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-1").await?, + "a new trigger fires" + ); + + assert!( + set_native_trigger_enabled( + &db, + "test-workspace", + ServiceName::Nextcloud, + "ext-1", + false + ) + .await? + ); + assert!( + !native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-1").await? + ); + + // Saving a configuration must not resume a trigger someone paused. + update_native_trigger( + &db, + "test-workspace", + ServiceName::Nextcloud, + "ext-1", + &config, + None, + Some("edited"), + ) + .await?; + assert!( + !native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-1").await?, + "an edit leaves the pause in place" + ); + + // A recreate registers a fresh trigger and must be able to come up already paused, in one + // write, rather than being enabled for as long as it takes a second call to arrive. + store_native_trigger( + &db, + "test-workspace", + ServiceName::Nextcloud, + "ext-2", + &config, + json!({}), + None, + false, + ) + .await?; + assert!( + !native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-2").await? + ); + + // The conflict branch is a re-registration of a trigger that already exists, so it carries no + // opinion about the pause. + store_native_trigger( + &db, + "test-workspace", + ServiceName::Nextcloud, + "ext-2", + &config, + json!({}), + None, + true, + ) + .await?; + assert!( + !native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "ext-2").await?, + "re-registering leaves the pause in place" + ); + + assert!( + !set_native_trigger_enabled( + &db, + "test-workspace", + ServiceName::Nextcloud, + "unknown", + true + ) + .await?, + "nothing to toggle" + ); + assert!( + native_trigger_is_enabled(&db, "test-workspace", ServiceName::Nextcloud, "unknown").await?, + "a trigger Windmill has no row for is not treated as paused" + ); + + Ok(()) +} diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 5eddd2c885..bb71d9ea29 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -19126,6 +19126,49 @@ paths: schema: type: string + /w/{workspace}/native_triggers/{service_name}/setenabled/{external_id}: + post: + summary: set enabled state of native trigger + description: | + Enables or disables a native trigger. A disabled trigger stays registered on the + external service but starts no job when it fires. + Requires write access to the script or flow that the trigger is associated with. + operationId: setNativeTriggerEnabled + tags: + - native_trigger + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - name: service_name + in: path + required: true + schema: + $ref: "#/components/schemas/NativeServiceName" + - name: external_id + in: path + required: true + description: The external ID of the trigger from the external service + schema: + type: string + requestBody: + description: updated enabled state + required: true + content: + application/json: + schema: + type: object + properties: + enabled: + type: boolean + required: + - enabled + responses: + "200": + description: native trigger enabled state updated + content: + text/plain: + schema: + type: string + /w/{workspace}/native_triggers/{service_name}/list: get: summary: list native triggers @@ -36095,6 +36138,9 @@ components: type: string nullable: true description: Short summary to be displayed when listed + enabled: + type: boolean + description: Whether the trigger starts a job when it fires required: - external_id - workspace_id @@ -36102,6 +36148,7 @@ components: - script_path - is_flow - service_config + - enabled NativeTriggerWithExternal: type: object @@ -36133,6 +36180,9 @@ components: type: string nullable: true description: Short summary to be displayed when listed + enabled: + type: boolean + description: Whether the trigger starts a job when it fires external_data: type: object nullable: true @@ -36153,6 +36203,7 @@ components: - script_path - is_flow - service_config + - enabled - external_data WorkspaceIntegrations: @@ -36238,6 +36289,12 @@ components: type: string nullable: true description: Short summary to be displayed when listed + enabled: + type: boolean + description: >- + Whether the trigger starts a job when it fires. Honoured on create only, so a + trigger can be registered already paused; an update ignores it and setenabled is + the only way to change an existing trigger's state. Defaults to true. required: - script_path - is_flow diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 8e361dcb2e..ccac957521 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -6511,7 +6511,7 @@ pub async fn run_flow_by_path( Query(run_query): Query, args: RawWebhookArgs, ) -> error::Result<(StatusCode, String)> { - let (args, trigger_metadata) = get_args_and_trigger_metadata( + let (args, trigger_metadata) = match get_args_and_trigger_metadata( &db, &authed, RunnableId::from_flow_path(flow_path.to_path()), @@ -6519,7 +6519,15 @@ pub async fn run_flow_by_path( &w_id, args, ) - .await?; + .await? + { + WebhookRun::Run(args, trigger_metadata) => (args, trigger_metadata), + // 200 rather than an error: services drop a webhook that keeps failing, and disabling a + // trigger in Windmill must not cost it its registration. + WebhookRun::TriggerDisabled => { + return Ok((StatusCode::OK, NATIVE_TRIGGER_DISABLED_MSG.to_string())) + } + }; let (uuid, _, _, _) = push_flow_job_by_path_into_queue( authed, @@ -6944,7 +6952,7 @@ pub async fn run_script_by_path( Query(run_query): Query, args: RawWebhookArgs, ) -> error::Result<(StatusCode, String)> { - let (args, trigger_metadata) = get_args_and_trigger_metadata( + let (args, trigger_metadata) = match get_args_and_trigger_metadata( &db, &authed, RunnableId::from_script_path(script_path.to_path()), @@ -6952,7 +6960,13 @@ pub async fn run_script_by_path( &w_id, args, ) - .await?; + .await? + { + WebhookRun::Run(args, trigger_metadata) => (args, trigger_metadata), + WebhookRun::TriggerDisabled => { + return Ok((StatusCode::OK, NATIVE_TRIGGER_DISABLED_MSG.to_string())) + } + }; let (uuid, _, _) = push_script_job_by_path_into_queue( authed, @@ -6970,6 +6984,16 @@ pub async fn run_script_by_path( Ok((StatusCode::CREATED, uuid.to_string())) } +/// What a webhook delivery resolved to: the arguments to run with, or nothing to run. +pub enum WebhookRun { + Run(PushArgsOwned, Option), + /// The native trigger this delivery belongs to is disabled. + TriggerDisabled, +} + +const NATIVE_TRIGGER_DISABLED_MSG: &str = + "This trigger is disabled in Windmill, so no job was created"; + #[allow(unused)] pub async fn get_args_and_trigger_metadata( db: &DB, @@ -6978,14 +7002,21 @@ pub async fn get_args_and_trigger_metadata( run_query: &RunJobQuery, w_id: &str, args: RawWebhookArgs, -) -> error::Result<(PushArgsOwned, Option)> { +) -> error::Result { use windmill_common::triggers::TriggerMetadata; // Build trigger metadata if this is a native trigger request #[cfg(feature = "native_trigger")] let (trigger_metadata, native_args) = if let Some(service_name_str) = &run_query.service_name { - use crate::native_triggers::{prepare_native_trigger_args, ServiceName}; + use crate::native_triggers::{ + native_trigger_is_enabled, prepare_native_trigger_args, ServiceName, + }; let service_name = ServiceName::try_from(service_name_str.to_owned())?; + if let Some(external_id) = run_query.trigger_external_id.as_deref() { + if !native_trigger_is_enabled(db, w_id, service_name, external_id).await? { + return Ok(WebhookRun::TriggerDisabled); + } + } let metadata = Some(TriggerMetadata::new( run_query.trigger_external_id.clone(), service_name.as_job_trigger_kind(), @@ -7021,7 +7052,7 @@ pub async fn get_args_and_trigger_metadata( .await? }; - Ok((args, trigger_metadata)) + Ok(WebhookRun::Run(args, trigger_metadata)) } #[derive(Deserialize)] diff --git a/backend/windmill-api/src/workspaces_export.rs b/backend/windmill-api/src/workspaces_export.rs index 869c4be68e..ff457ff35e 100644 --- a/backend/windmill-api/src/workspaces_export.rs +++ b/backend/windmill-api/src/workspaces_export.rs @@ -1434,9 +1434,11 @@ pub(crate) async fn tarball_workspace( // Native triggers (Nextcloud, Google Drive, GitHub) are never // cloned into a fork — a fork only has one if its owner created // it there, so it's always "fork-only" and keeps its own mode. - // No parent-value substitution applies; we only strip the - // webhook token hash. - let native_ignore_keys = vec!["webhook_token_hash"]; + // No parent-value substitution applies; we strip the webhook + // token hash, and `enabled`, which is operational state a sync + // deliberately does not carry — whether a trigger is paused + // belongs to the workspace it runs in, not to the code. + let native_ignore_keys = vec!["webhook_token_hash", "enabled"]; for trigger in native_triggers { let trigger_str = &to_string_without_metadata( diff --git a/backend/windmill-common/src/instance_config.rs b/backend/windmill-common/src/instance_config.rs index 7ea6dbe5ae..c3060e9a89 100644 --- a/backend/windmill-common/src/instance_config.rs +++ b/backend/windmill-common/src/instance_config.rs @@ -916,6 +916,14 @@ pub struct WorkerGroupConfig { pub autoscaling: Option, #[serde(skip_serializing_if = "Option::is_none")] pub native_mode: Option, + /// Object store this group's dependency cache uses instead of the instance one. Same shape + /// as the instance `object_store_cache_config`. + #[serde(skip_serializing_if = "Option::is_none")] + #[cfg_attr( + feature = "instance_config_schema", + schemars(schema_with = "opaque_json_schema") + )] + pub object_store_cache_config: Option, /// Catch-all for fields not yet covered by typed fields. #[serde(flatten)] diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index bd6284f64b..0f3a043c34 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -392,6 +392,7 @@ lazy_static::lazy_static! { pip_local_dependencies: Default::default(), env_vars: Default::default(), native_mode: false, + object_store_cache_config: Default::default(), }); pub static ref WORKER_PULL_QUERIES: arc_swap::ArcSwap> = arc_swap::ArcSwap::from_pointee(vec![]); @@ -2376,6 +2377,7 @@ pub async fn load_worker_config( .or_else(|| load_additional_python_paths_from_env()), env_vars: resolved_env_vars, native_mode, + object_store_cache_config: config.object_store_cache_config, }) } @@ -2465,6 +2467,7 @@ pub struct WorkerConfigOpt { pub env_vars_static: Option>, pub env_vars_allowlist: Option>, pub native_mode: Option, + pub object_store_cache_config: Option, } impl Default for WorkerConfigOpt { @@ -2483,6 +2486,7 @@ impl Default for WorkerConfigOpt { env_vars_static: Default::default(), env_vars_allowlist: Default::default(), native_mode: Default::default(), + object_store_cache_config: Default::default(), } } } @@ -2501,12 +2505,18 @@ pub struct WorkerConfig { pub pip_local_dependencies: Option>, pub env_vars: HashMap, pub native_mode: bool, + /// Object store this group's dependency cache uses instead of the instance one, as stored + /// in the group config. Raw JSON: `windmill-common` cannot depend on the object store crate + /// that parses it, and comparing the raw value is what tells a reload the store changed. + pub object_store_cache_config: Option, } impl std::fmt::Debug for WorkerConfig { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "WorkerConfig {{ worker_tags: {:?}, priority_tags_sorted: {:?}, dedicated_worker: {:?}, dedicated_workers: {:?}, init_bash: {:?}, periodic_script_bash: {:?}, periodic_script_interval_seconds: {:?}, cache_clear: {:?}, additional_python_paths: {:?}, pip_local_dependencies: {:?}, env_vars: {:?}, native_mode: {:?} }}", - self.worker_tags, self.priority_tags_sorted, self.dedicated_worker, self.dedicated_workers, self.init_bash, self.periodic_script_bash, self.periodic_script_interval_seconds, self.cache_clear, self.additional_python_paths, self.pip_local_dependencies, self.env_vars.iter().map(|(k, v)| format!("{}: {}{} ({} chars)", k, &v[..3.min(v.len())], "***", v.len())).collect::>().join(", "), self.native_mode) + write!(f, "WorkerConfig {{ worker_tags: {:?}, priority_tags_sorted: {:?}, dedicated_worker: {:?}, dedicated_workers: {:?}, init_bash: {:?}, periodic_script_bash: {:?}, periodic_script_interval_seconds: {:?}, cache_clear: {:?}, additional_python_paths: {:?}, pip_local_dependencies: {:?}, env_vars: {:?}, native_mode: {:?}, object_store_cache_config: {} }}", + self.worker_tags, self.priority_tags_sorted, self.dedicated_worker, self.dedicated_workers, self.init_bash, self.periodic_script_bash, self.periodic_script_interval_seconds, self.cache_clear, self.additional_python_paths, self.pip_local_dependencies, self.env_vars.iter().map(|(k, v)| format!("{}: {}{} ({} chars)", k, &v[..3.min(v.len())], "***", v.len())).collect::>().join(", "), self.native_mode, + // holds bucket credentials + self.object_store_cache_config.as_ref().map(|_| "***").unwrap_or("None")) } } diff --git a/backend/windmill-native-triggers/src/handler.rs b/backend/windmill-native-triggers/src/handler.rs index cc6c61be81..27a50ba07a 100644 --- a/backend/windmill-native-triggers/src/handler.rs +++ b/backend/windmill-native-triggers/src/handler.rs @@ -1,8 +1,8 @@ use crate::{ classify_read_failure, decrypt_oauth_data, delete_native_trigger, delete_token_by_hash, get_native_trigger, list_native_triggers, lock::TriggerLock, map_external_error, - map_external_error_with, rotate_webhook_token, store_native_trigger, - sync::EXTERNAL_TRIGGER_MISSING_ERROR, update_native_trigger_error, + map_external_error_with, rotate_webhook_token, set_native_trigger_enabled, + store_native_trigger, sync::EXTERNAL_TRIGGER_MISSING_ERROR, update_native_trigger_error, update_native_trigger_if_runnable_unchanged, webhook_token_label, webhook_token_scopes, External, ExternalReadFailure, NativeTrigger, NativeTriggerConfig, NativeTriggerData, ServiceName, @@ -239,6 +239,7 @@ async fn create_native_trigger( &config, service_config, data.summary.as_deref(), + data.enabled, ) .await?; @@ -603,6 +604,87 @@ async fn delete_native_trigger_handler( Ok(format!("Native trigger deleted")) } +#[derive(Debug, Deserialize)] +pub struct SetEnabledPayload { + pub enabled: bool, +} + +/// Pause or resume a trigger, without touching its registration on the external service. +/// +/// Leaving the webhook registered is what makes this reversible: services drop or deactivate a +/// subscription that keeps failing, so a paused trigger keeps answering deliveries normally and +/// simply starts no job. +async fn set_native_trigger_enabled_handler( + Extension(service_name): Extension, + authed: ApiAuthed, + Extension(db): Extension, + Extension(user_db): Extension, + Path((workspace_id, external_id)): Path<(String, String)>, + Json(payload): Json, +) -> Result { + let mut tx = user_db.begin(&authed).await?; + + let existing = get_native_trigger(&mut *tx, &workspace_id, service_name, &external_id) + .await? + .ok_or_else(|| Error::NotFound(format!("Native trigger not found: {}", external_id)))?; + + check_scopes(&authed, || { + format!("native_triggers:write:{}", &existing.script_path) + })?; + require_is_writer_on_runnable( + &authed, + &existing.script_path, + existing.is_flow, + &workspace_id, + db.clone(), + ) + .await?; + + let updated = set_native_trigger_enabled( + &mut *tx, + &workspace_id, + service_name, + &external_id, + payload.enabled, + ) + .await?; + + // The read above takes no row lock, so a concurrent delete can land in between; reporting + // success then would tell the caller a trigger that is gone had been paused. + if !updated { + return Err(Error::NotFound(format!( + "Native trigger not found: {}", + external_id + ))); + } + + audit_log( + &mut *tx, + &authed, + &format!( + "native_triggers.{}.{}", + service_name, + if payload.enabled { "enable" } else { "disable" } + ), + ActionKind::Update, + &workspace_id, + Some(&external_id), + None, + ) + .await?; + + tx.commit().await?; + + Ok(format!( + "Native trigger {}", + if payload.enabled { + "enabled" + } else { + "disabled" + } + )) +} + async fn list_native_triggers_handler( Extension(service_name): Extension, authed: ApiAuthed, @@ -642,6 +724,10 @@ pub fn service_routes(handler: T) -> Router { .route( "/delete/{external_id}", delete(delete_native_trigger_handler::), + ) + .route( + "/setenabled/{external_id}", + post(set_native_trigger_enabled_handler::), ); standard_routes diff --git a/backend/windmill-native-triggers/src/lib.rs b/backend/windmill-native-triggers/src/lib.rs index 4400091ed6..570319ee0e 100644 --- a/backend/windmill-native-triggers/src/lib.rs +++ b/backend/windmill-native-triggers/src/lib.rs @@ -226,6 +226,10 @@ pub struct NativeTrigger { pub created_at: DateTime, pub updated_at: DateTime, pub summary: Option, + /// Whether incoming webhooks for this trigger start a job. Operational state: a create sets + /// its initial value and `setenabled` is its only mutator afterwards, so saving a + /// configuration can never silently re-enable a trigger someone paused. + pub enabled: bool, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -235,12 +239,20 @@ pub struct NativeTriggerConfig { pub webhook_token: String, } +fn default_true() -> bool { + true +} + #[derive(Debug, Serialize, Deserialize)] pub struct NativeTriggerData { pub script_path: String, pub is_flow: bool, pub service_config: C, pub summary: Option, + /// Honoured on create only, so a trigger can be registered already paused in one request. + /// An update ignores it: `setenabled` is the only way to change an existing trigger's state. + #[serde(default = "default_true")] + pub enabled: bool, } #[derive(Debug, Clone, FromRow, Serialize, Deserialize)] @@ -1179,11 +1191,15 @@ pub async fn store_native_trigger<'c, E: sqlx::Executor<'c, Database = Postgres> config: &NativeTriggerConfig, service_config: C, summary: Option<&str>, + enabled: bool, ) -> Result<()> { use windmill_common::auth::hash_token; let webhook_token_hash = hash_token(&config.webhook_token); + // `enabled` is set by the INSERT alone: writing it here rather than in a follow-up statement + // is what keeps a trigger created paused from ever being visible, and therefore runnable, in + // any other state. The conflict branch leaves it untouched for the mirror-image reason. sqlx::query!( r#" INSERT INTO native_trigger ( @@ -1194,9 +1210,10 @@ pub async fn store_native_trigger<'c, E: sqlx::Executor<'c, Database = Postgres> is_flow, webhook_token_hash, service_config, - summary + summary, + enabled ) VALUES ( - $1, $2, $3, $4, $5, $6, $7, $8 + $1, $2, $3, $4, $5, $6, $7, $8, $9 ) ON CONFLICT (external_id, workspace_id, service_name) DO UPDATE SET script_path = $4, is_flow = $5, webhook_token_hash = $6, service_config = $7, summary = $8, error = NULL, updated_at = NOW() @@ -1209,6 +1226,7 @@ pub async fn store_native_trigger<'c, E: sqlx::Executor<'c, Database = Postgres> webhook_token_hash, sqlx::types::Json(service_config) as _, summary, + enabled, ) .execute(db) .await?; @@ -1405,7 +1423,8 @@ pub async fn get_native_trigger<'c, E: sqlx::Executor<'c, Database = Postgres>>( error, created_at, updated_at, - summary + summary, + enabled FROM native_trigger WHERE @@ -1444,7 +1463,8 @@ pub async fn get_native_trigger_by_script<'c, E: sqlx::Executor<'c, Database = P error, created_at, updated_at, - summary + summary, + enabled FROM native_trigger WHERE @@ -1491,7 +1511,8 @@ pub async fn list_native_triggers<'c, E: sqlx::Executor<'c, Database = Postgres> nt.error, nt.created_at, nt.updated_at, - nt.summary + nt.summary, + nt.enabled FROM native_trigger nt WHERE @@ -1555,6 +1576,71 @@ pub async fn update_native_trigger_error<'c, E: sqlx::Executor<'c, Database = Po Ok(()) } +/// Pause or resume a trigger. Returns `false` when there is no such trigger. +/// +/// Callers MUST have verified write access to the trigger's runnable: this writes operational +/// state and performs no authorization of its own. +pub async fn set_native_trigger_enabled<'c, E: sqlx::Executor<'c, Database = Postgres>>( + db: E, + workspace_id: &str, + service_name: ServiceName, + external_id: &str, + enabled: bool, +) -> Result { + // `updated_at` is the row version `record_reregistration` conditions on, so leave it alone: + // pausing a trigger must not make a registration that is mid-flight discard its result. + let updated = sqlx::query!( + r#" + UPDATE native_trigger + SET enabled = $1 + WHERE + workspace_id = $2 + AND service_name = $3 + AND external_id = $4 + "#, + enabled, + workspace_id, + service_name as ServiceName, + external_id, + ) + .execute(db) + .await? + .rows_affected(); + + Ok(updated > 0) +} + +/// Whether a webhook arriving for this trigger should start a job. +/// +/// A trigger Windmill no longer knows about counts as enabled: the token in the URL is what +/// authorizes the run, and this is a pause switch, not a second authorization check. It reads +/// nothing a caller could not already learn from the trigger it is delivering for, so it needs no +/// authorization of its own — but it also grants none, and must not be used as one. +pub async fn native_trigger_is_enabled<'c, E: sqlx::Executor<'c, Database = Postgres>>( + db: E, + workspace_id: &str, + service_name: ServiceName, + external_id: &str, +) -> Result { + let enabled = sqlx::query_scalar!( + r#" + SELECT enabled + FROM native_trigger + WHERE + workspace_id = $1 + AND service_name = $2 + AND external_id = $3 + "#, + workspace_id, + service_name as ServiceName, + external_id, + ) + .fetch_optional(db) + .await?; + + Ok(enabled.unwrap_or(true)) +} + pub async fn update_native_trigger_service_config< 'c, E: sqlx::Executor<'c, Database = Postgres>, @@ -1701,6 +1787,11 @@ pub async fn delete_workspace_integration( /// /// `external_id` is optional because during CREATE we don't have it yet /// (it's returned by the external service). During UPDATE, we have it. +/// +/// Every registered URL MUST end up carrying it: it is the only thing a delivery identifies its +/// trigger by, so a service that leaves it out ships a disable switch that silently does nothing. +/// A handler that returns `None` from `service_config_from_create_response` gets this for free — +/// `create_native_trigger` then runs the `update` cycle that re-registers with the assigned id. pub fn generate_webhook_service_url( base_url: &str, w_id: &str, diff --git a/backend/windmill-native-triggers/src/rename.rs b/backend/windmill-native-triggers/src/rename.rs index a7168f5c7e..caae99028c 100644 --- a/backend/windmill-native-triggers/src/rename.rs +++ b/backend/windmill-native-triggers/src/rename.rs @@ -135,6 +135,7 @@ async fn reregister_one( is_flow: trigger.is_flow, service_config, summary: trigger.summary.clone(), + enabled: trigger.enabled, }; // The token is scoped to the runnable path and only its hash is kept, so pointing the webhook diff --git a/backend/windmill-object-store/Cargo.toml b/backend/windmill-object-store/Cargo.toml index 3602007d8c..90af4f7e7c 100644 --- a/backend/windmill-object-store/Cargo.toml +++ b/backend/windmill-object-store/Cargo.toml @@ -58,3 +58,4 @@ aws-credential-types = { workspace = true, optional = true } tempfile.workspace = true tokio = { workspace = true, features = ["rt", "macros"] } object_store.workspace = true +serial_test = "3" diff --git a/backend/windmill-object-store/src/lib.rs b/backend/windmill-object-store/src/lib.rs index cc588d485f..c403625e00 100644 --- a/backend/windmill-object-store/src/lib.rs +++ b/backend/windmill-object-store/src/lib.rs @@ -162,37 +162,211 @@ impl From> for ExpirableObjectStore { #[cfg(feature = "parquet")] lazy_static::lazy_static! { pub static ref OBJECT_STORE_SETTINGS: Arc>> = Arc::new(RwLock::new(None)); + + /// Worker-group override of the store backing the *dependency cache* only: venvs, language + /// bundles and compiled binaries, which a worker both writes and reads back itself. + /// Everything the server also reads — job results, logs, codebases, app assets — stays on + /// [`OBJECT_STORE_SETTINGS`], which a worker-local redirect would make unreachable. + static ref CACHE_OBJECT_STORE_OVERRIDE: Arc>> = Arc::new(RwLock::new(None)); + + /// The config [`CACHE_OBJECT_STORE_OVERRIDE`] was built from, so a rebuild that fails for a + /// config already being served can keep serving it. Locked after the store, never before. + static ref CACHE_OVERRIDE_APPLIED: Arc>> = Arc::new(RwLock::new(None)); + + /// Held across a whole [`reload_cache_object_store_override`], build included, so that the + /// override's flag, store and applied config only ever move together. + static ref CACHE_OVERRIDE_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::new(()); +} + +/// Whether a worker-group cache override is configured, held apart from the store it built so +/// that a configured-but-broken override reads as "no cache store" instead of silently falling +/// back to the instance bucket the operator redirected away from. +#[cfg(feature = "parquet")] +static CACHE_OBJECT_STORE_OVERRIDDEN: std::sync::atomic::AtomicBool = + std::sync::atomic::AtomicBool::new(false); + +/// Bumped by every [`reload_cache_object_store_override`] at entry. Builds are slow and several +/// callers race — a config change, the retry behind a failed one, a license-plan change — and the +/// lock alone would only order them by arrival, so a reload that lost its claim while waiting +/// drops out rather than installing a store the group has already moved off. +#[cfg(feature = "parquet")] +static CACHE_OVERRIDE_GENERATION: std::sync::atomic::AtomicU64 = + std::sync::atomic::AtomicU64::new(0); + +#[cfg(feature = "parquet")] +async fn resolve_object_store( + settings_lock: &RwLock>, +) -> Option> { + let settings = settings_lock.read().await; + let Some(s) = settings.as_ref() else { + return None; + }; + match &s.refresh { + Some(refresh) if refresh.refresh_needed() => { + let refresh = refresh.clone(); + let refreshed_from = s.store.clone(); + drop(settings); + let new_store = refresh.refresh().await?; + let mut settings = settings_lock.write().await; + match settings.as_ref() { + // A reload may have installed a different store while the credentials were + // being minted; that one reflects newer config, so the refresh is stale. + Some(current) if !Arc::ptr_eq(¤t.store, &refreshed_from) => { + Some(current.store.clone()) + } + Some(_) => { + let arc = new_store.store.clone(); + *settings = Some(new_store); + Some(arc) + } + // Cleared while refreshing. + None => None, + } + } + _ => Some(s.store.clone()), + } } #[cfg(feature = "parquet")] pub async fn get_object_store() -> Option> { - let settings = OBJECT_STORE_SETTINGS.read().await; - if let Some(s) = settings.as_ref() { - match &s.refresh { - Some(refresh) => { - if refresh.refresh_needed() { - let refresh = refresh.clone(); - drop(settings); - let new_store = refresh.refresh().await; - if let Some(new_store) = new_store { - let mut s3_cache_settings = OBJECT_STORE_SETTINGS.write().await; - let arc = new_store.store.clone(); - *s3_cache_settings = Some(new_store); - return Some(arc); - } else { - return None; - } - } else { - return Some(s.store.clone()); - } - } - None => { - return Some(s.store.clone()); - } - } - } else { - return None; + resolve_object_store(&OBJECT_STORE_SETTINGS).await +} + +/// The store the dependency cache reads and writes: the worker group's override when it has one, +/// the instance object store otherwise. Anything the server must also reach goes through +/// [`get_object_store`] instead. +#[cfg(feature = "parquet")] +pub async fn get_cache_object_store() -> Option> { + if CACHE_OBJECT_STORE_OVERRIDDEN.load(std::sync::atomic::Ordering::Relaxed) { + return resolve_object_store(&CACHE_OBJECT_STORE_OVERRIDE).await; } + resolve_object_store(&OBJECT_STORE_SETTINGS).await +} + +/// True when an override is configured but has no usable store. The caller's short retry is the +/// fast path back; this is the backstop, and the full settings reload it rides on is 12h apart by +/// default (`SETTINGS_RELOAD_PERIOD_SECS`), so an outage outlasting the retry keeps the group's +/// dependency cache local until then or until someone edits the group config. +#[cfg(feature = "parquet")] +pub async fn cache_object_store_override_failed() -> bool { + CACHE_OBJECT_STORE_OVERRIDDEN.load(std::sync::atomic::Ordering::Relaxed) + && CACHE_OBJECT_STORE_OVERRIDE.read().await.is_none() +} + +/// Apply the `object_store_cache_config` of this worker's group. `None` (or JSON null) drops the +/// override and returns the worker to the instance object store. +/// +/// Returns [`ObjectStoreReload::Later`] when the store did not build for a reason that may pass — +/// the caller is expected to retry shortly, as `initial_load` does for the instance store. +#[cfg(feature = "parquet")] +pub async fn reload_cache_object_store_override( + db: &windmill_common::DB, + settings: Option, +) -> ObjectStoreReload { + use std::sync::atomic::Ordering; + use windmill_common::ee_oss::{get_license_plan, LicensePlan}; + + // Claim a generation, then take the lock: every state transition below happens inside one + // critical section, and a caller that lost its claim while waiting drops out rather than + // installing what the group has already moved off. + let generation = CACHE_OVERRIDE_GENERATION.fetch_add(1, Ordering::SeqCst) + 1; + let _transition = CACHE_OVERRIDE_LOCK.lock().await; + if CACHE_OVERRIDE_GENERATION.load(Ordering::SeqCst) != generation { + return ObjectStoreReload::Never; + } + + // DISABLE_S3_STORE turns off the instance object store for this process; a group override + // must not be a way back in. + let store_disabled = std::env::var("DISABLE_S3_STORE") + .ok() + .is_some_and(|x| x == "1" || x == "true"); + + let Some(settings) = settings.filter(|v| !v.is_null() && !store_disabled) else { + if CACHE_OBJECT_STORE_OVERRIDDEN.swap(false, Ordering::Relaxed) { + clear_cache_object_store_override().await; + tracing::info!( + "Worker group object store cache override removed, falling back to the instance object store" + ); + } + return ObjectStoreReload::Never; + }; + + // Enterprise-only, so anything else — Community, including a CE build reaching this through + // the config-as-code API, and Pro — must not get a store, and a plan that stops being + // Enterprise must drop one loaded while it still was. + if !matches!(get_license_plan().await, LicensePlan::Enterprise) { + tracing::error!( + "Object store cache override requires an enterprise license, ignoring it for this worker group" + ); + if CACHE_OBJECT_STORE_OVERRIDDEN.swap(false, Ordering::Relaxed) { + clear_cache_object_store_override().await; + } + return ObjectStoreReload::Never; + } + + apply_cache_object_store_override(db, settings).await +} + +/// The half of [`reload_cache_object_store_override`] past the entitlement gate: build the store +/// and commit it. Split out so the commit rules are testable without a license plan. +#[cfg(feature = "parquet")] +async fn apply_cache_object_store_override( + db: &windmill_common::DB, + settings: serde_json::Value, +) -> ObjectStoreReload { + use std::sync::atomic::Ordering; + + // Claim the override before building it: until a store is in place the dependency cache + // must stay local-only rather than reach for the instance bucket. + CACHE_OBJECT_STORE_OVERRIDDEN.store(true, Ordering::Relaxed); + + let (store, reload) = match serde_json::from_value::(settings.clone()) { + Ok(setting) => match build_object_store_from_settings(setting, Some(db)).await { + Ok(store) => (Some(store), ObjectStoreReload::Never), + Err(e) => { + tracing::error!( + "Error building the worker group object store cache override, the dependency cache stays local to this worker until it builds: {e:?}" + ); + (None, ObjectStoreReload::Later) + } + }, + // A malformed config will read the same on every retry. + Err(e) => { + tracing::error!( + "Error parsing the worker group object store cache override, the dependency cache stays local to this worker: {e:?}" + ); + (None, ObjectStoreReload::Never) + } + }; + + let mut current = CACHE_OBJECT_STORE_OVERRIDE.write().await; + match store { + Some(store) => { + *current = Some(store); + *CACHE_OVERRIDE_APPLIED.write().await = Some(settings); + tracing::info!( + "Dependency cache of this worker group now uses its own object store, not the instance one" + ); + } + // A rebuild that failed for the config already being served leaves that store in place: + // the group is entitled to it, and dropping it would take the whole group's cache local + // over a transient error. A *different* config failing must still clear, or the worker + // would keep writing to the bucket the operator redirected it away from. + None if current.is_some() + && CACHE_OVERRIDE_APPLIED.read().await.as_ref() == Some(&settings) => {} + None => { + *current = None; + *CACHE_OVERRIDE_APPLIED.write().await = None; + } + } + reload +} + +/// Drop the override store and the config it was built from, in that lock order. +#[cfg(feature = "parquet")] +async fn clear_cache_object_store_override() { + *CACHE_OBJECT_STORE_OVERRIDE.write().await = None; + *CACHE_OVERRIDE_APPLIED.write().await = None; } #[cfg(feature = "parquet")] @@ -2361,10 +2535,100 @@ mod tests { .contains("Error building filesystem object store")); } + /// A worker group override that is configured but has no usable store must leave the + /// dependency cache with no object store at all. Falling back to the instance one would + /// write the group's cache into the bucket the operator redirected it away from. + // Serialized with the other test that swaps OBJECT_STORE_SETTINGS: the store is + // process-global and CI runs this binary with --test-threads=10. + #[cfg(feature = "parquet")] + #[tokio::test] + #[serial_test::serial(object_store_settings)] + async fn test_get_cache_object_store_override() { + use object_store::{path::Path, ObjectStore, PutPayload}; + use std::sync::atomic::Ordering; + + async fn marker_of(store: &Arc) -> String { + let bytes = store.get(&Path::from("marker")).await.unwrap(); + String::from_utf8(bytes.bytes().await.unwrap().to_vec()).unwrap() + } + + let instance_dir = tempfile::tempdir().unwrap(); + let instance = build_filesystem_client(instance_dir.path().to_str().unwrap()).unwrap(); + instance + .put(&Path::from("marker"), PutPayload::from("instance")) + .await + .unwrap(); + + let group_dir = tempfile::tempdir().unwrap(); + let group = build_filesystem_client(group_dir.path().to_str().unwrap()).unwrap(); + group + .put(&Path::from("marker"), PutPayload::from("group")) + .await + .unwrap(); + + *OBJECT_STORE_SETTINGS.write().await = Some(ExpirableObjectStore::from(instance)); + + let store = get_cache_object_store().await.unwrap(); + assert_eq!(marker_of(&store).await, "instance"); + assert!(!cache_object_store_override_failed().await); + + CACHE_OBJECT_STORE_OVERRIDDEN.store(true, Ordering::Relaxed); + *CACHE_OBJECT_STORE_OVERRIDE.write().await = Some(ExpirableObjectStore::from(group)); + let store = get_cache_object_store().await.unwrap(); + assert_eq!(marker_of(&store).await, "group"); + + // Configured but unbuilt, as a failed reload leaves it: no store at all, rather than the + // instance bucket the operator redirected the group away from. + *CACHE_OBJECT_STORE_OVERRIDE.write().await = None; + assert!(get_cache_object_store().await.is_none()); + assert!(cache_object_store_override_failed().await); + + // The teardown branch returns before the pool is used, so a lazy one is enough. + let db = sqlx::postgres::PgPool::connect_lazy("postgres://localhost/unused").unwrap(); + reload_cache_object_store_override(&db, None).await; + let store = get_cache_object_store().await.unwrap(); + assert_eq!(marker_of(&store).await, "instance"); + assert!(!cache_object_store_override_failed().await); + + *OBJECT_STORE_SETTINGS.write().await = None; + } + + /// A rebuild is triggered by any edit to the group config, not only by editing the store, so + /// a build that fails for the config already installed must leave it alone — otherwise a + /// renamed worker tag plus one flaky token mint takes the whole group's cache local. A + /// *different* config failing still has to clear it. + #[cfg(feature = "parquet")] + #[tokio::test] + #[serial_test::serial(object_store_settings)] + async fn test_failed_rebuild_keeps_the_store_serving_the_same_config() { + let db = sqlx::postgres::PgPool::connect_lazy("postgres://localhost/unused").unwrap(); + let dir = tempfile::tempdir().unwrap(); + let settings = serde_json::json!({ + "type": "Filesystem", "root_path": dir.path().to_str().unwrap() + }); + + apply_cache_object_store_override(&db, settings.clone()).await; + assert!(get_cache_object_store().await.is_some()); + + // Same config, now unbuildable: the store it already produced stays. + dir.close().unwrap(); + apply_cache_object_store_override(&db, settings).await; + assert!(get_cache_object_store().await.is_some()); + + // A different config that will not build must not leave the old bucket in place. + let moved = serde_json::json!({ "type": "Filesystem", "root_path": "/proc/nonexistent" }); + apply_cache_object_store_override(&db, moved).await; + assert!(get_cache_object_store().await.is_none()); + assert!(cache_object_store_override_failed().await); + + reload_cache_object_store_override(&db, None).await; + } + // --- get_logs_from_store test --- #[cfg(feature = "parquet")] #[tokio::test] + #[serial_test::serial(object_store_settings)] async fn test_get_logs_from_store_with_filesystem() { use futures::StreamExt; use object_store::{path::Path, ObjectStore, PutPayload}; diff --git a/backend/windmill-runtime-nativets/src/lib.rs b/backend/windmill-runtime-nativets/src/lib.rs index 6e80b1405e..d44213d3bc 100644 --- a/backend/windmill-runtime-nativets/src/lib.rs +++ b/backend/windmill-runtime-nativets/src/lib.rs @@ -26,7 +26,7 @@ use std::{ cell::RefCell, path::PathBuf, rc::Rc, - sync::{Arc, Mutex}, + sync::{Arc, LazyLock, Mutex}, }; // Re-export deno_telemetry for use by windmill-worker's otel proxy @@ -182,10 +182,33 @@ struct LogString { pub s: mpsc::UnboundedSender, } -#[derive(Clone)] +#[derive(Clone, Default)] pub struct NativeAnnotation { pub useragent: Option, pub proxy: Option<(String, Option<(String, String)>)>, + /// `//fetch_response_timeout `: 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, +} + +/// 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 = LazyLock::new(|| { + std::env::var("WINDMILL_FETCH_RESPONSE_TIMEOUT_SECS") + .ok() + .and_then(|x| x.trim().parse::().ok()) + .unwrap_or(900) + }); + *SECS } /// Serializes V8 isolate creation as defense-in-depth against concurrent @@ -401,7 +424,7 @@ pub fn transpile_ts(expr: String) -> anyhow::Result { } pub fn get_annotation(inner_content: &str) -> NativeAnnotation { - let mut res = NativeAnnotation { useragent: None, proxy: None }; + let mut res = NativeAnnotation::default(); let anns = inner_content .lines() @@ -414,6 +437,13 @@ pub fn get_annotation(inner_content: &str) -> NativeAnnotation { 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::() + .ok(); } } res @@ -554,6 +584,16 @@ pub(crate) fn create_nativets_runtime( 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()), @@ -625,10 +665,15 @@ pub(crate) fn create_nativets_runtime( } // Per-isolate JS init that can't run in the snapshot (runtime.js executes at - // snapshot-build time): currently seeds performance.timeOrigin via - // setTimeOrigin(), which must read this isolate's wall clock. + // snapshot-build time): the wall clock behind performance.timeOrigin and the + // fetch response timeout are both per-isolate values. js_runtime - .execute_script("", "globalThis.__wmInitPerIsolate()") + .execute_script( + "", + format!( + "globalThis.__wmInitPerIsolate({{ fetchResponseTimeoutMs: {fetch_response_timeout_ms} }})" + ), + ) .map_err(windmill_common::error::to_anyhow)?; Ok(CreatedRuntime { js_runtime, log_receiver, memory_limit_rx }) diff --git a/backend/windmill-runtime-nativets/src/runtime.js b/backend/windmill-runtime-nativets/src/runtime.js index 01bec09983..88e10f765a 100644 --- a/backend/windmill-runtime-nativets/src/runtime.js +++ b/backend/windmill-runtime-nativets/src/runtime.js @@ -30,9 +30,115 @@ import * as performance from "ext:deno_web/15_performance.js"; import "ext:deno_web/16_image_data.js"; import "ext:deno_fetch/27_eventsource.js"; +// deno_fetch applies no deadline, so a peer that accepts a request and then +// never answers leaves `await fetch(...)` pending until the job timeout, which +// self-hosted defaults to 7 days. +const ORIGINAL_FETCH = fetch.fetch; + +// deno_web's timers reject any `this` other than undefined/globalThis, so +// `timers.setTimeout(...)` passes the module namespace and throws "Illegal +// invocation". +const setTimeoutUnbound = timers.setTimeout; +const clearTimeoutUnbound = timers.clearTimeout; +// Captured before user code shares the isolate and could redefine them, the +// way deno's own modules reach intrinsics through primordials. +const PromiseReject = Promise.reject.bind(Promise); +const promiseThen = Function.prototype.call.bind(Promise.prototype.then); +const abortSignalAny = abortSignal.AbortSignal.any.bind(abortSignal.AbortSignal); +const abortControllerAbort = Function.prototype.call.bind( + abortSignal.AbortController.prototype.abort, +); +const ReflectApply = Reflect.apply; + +// Installed per isolate by __wmInitPerIsolate; 0 disables. Only a backstop +// for the impossible case of fetch running before that init. +let fetchResponseTimeoutMs = 900_000; + +function fetchResponseTimeoutError(requestUrl, timeoutMs) { + let target; + try { + const parsed = new url.URL(requestUrl); + // Query and fragment routinely carry tokens, and this reaches a job log. + target = parsed.origin + parsed.pathname; + } catch { + target = "the request target"; + } + return new domException.DOMException( + `fetch to ${target} timed out: no response headers arrived within ` + + `${Math.round(timeoutMs / 1000)}s (this covers connect, request upload ` + + `and the wait for the server to start replying; once a response begins ` + + `it is never interrupted). Change it per script with ` + + `"//fetch_response_timeout " (0 disables), or instance-wide ` + + `with WINDMILL_FETCH_RESPONSE_TIMEOUT_SECS.`, + "TimeoutError", + ); +} + globalThis.atob = base64.atob; globalThis.btoa = base64.btoa; -globalThis.fetch = fetch.fetch; +// Not `async`, for the same reason deno_fetch's own outer fetch isn't: WPT +// pins that an aborted fetch settles in the same tick, which adopting its +// promise through another one would break. Construction still has to reject +// rather than throw, so it is caught and handed back as a rejection. +globalThis.fetch = function fetch(input, init = undefined) { + const timeoutMs = fetchResponseTimeoutMs; + // Forwarded with the original argument count, so deno still sees an empty + // call as empty and raises its own "1 argument required". The default on + // `init` is what keeps `fetch.length` at 1, as the standard has it. + if (!(timeoutMs > 0) || arguments.length < 1) { + return ReflectApply(ORIGINAL_FETCH, undefined, arguments); + } + + let req; + let controller; + let signal; + try { + // RequestInit is a WebIDL dictionary: copying it drops inherited and + // non-enumerable members, and inheriting from it runs accessors against the + // wrong receiver. Hand it to the same Request constructor fetch() would, + // and carry our own signal in an init of our own. + req = new request.Request(input, init); + + // `req.signal` is deno's own resolution of init.signal over an input + // Request's signal, so combining with it preserves the caller's abort and + // reason while ours only adds a ceiling. + controller = new abortSignal.AbortController(); + signal = abortSignalAny([req.signal, controller.signal]); + } catch (e) { + return PromiseReject(e); + } + + // Already aborted: hand back deno's own settled rejection untouched, and arm + // nothing -- there is no response to wait for. + if (signal.aborted) { + return ORIGINAL_FETCH(req, { signal }); + } + + let timer = setTimeoutUnbound(() => { + timer = undefined; + abortControllerAbort(controller, fetchResponseTimeoutError(req.url, timeoutMs)); + }, timeoutMs); + // Disarmed on headers, never on body completion: a response that has begun + // arriving must be free to stream for as long as it needs. + const disarm = () => { + if (timer !== undefined) { + clearTimeoutUnbound(timer); + timer = undefined; + } + }; + + return promiseThen( + ORIGINAL_FETCH(req, { signal }), + (res) => { + disarm(); + return res; + }, + (e) => { + disarm(); + throw e; + }, + ); +}; globalThis.Request = request.Request; globalThis.Response = response.Response; globalThis.Blob = file.Blob; @@ -123,7 +229,11 @@ Object.assign(globalThis, { // Per-isolate init, invoked from Rust after the snapshot is restored (this // module body runs at snapshot-build time, not per isolate). -globalThis.__wmInitPerIsolate = () => { +globalThis.__wmInitPerIsolate = (config) => { + if (config != null && typeof config.fetchResponseTimeoutMs === "number") { + fetchResponseTimeoutMs = config.fetchResponseTimeoutMs; + } + // setTimeOrigin() seeds performance.timeOrigin from the isolate's wall clock; // without it timeOrigin is undefined and `timeOrigin + performance.now()` is NaN. performance.setTimeOrigin(); diff --git a/backend/windmill-runtime-nativets/src/smoke_tests.rs b/backend/windmill-runtime-nativets/src/smoke_tests.rs index b581a9f1fe..f3da460227 100644 --- a/backend/windmill-runtime-nativets/src/smoke_tests.rs +++ b/backend/windmill-runtime-nativets/src/smoke_tests.rs @@ -24,7 +24,7 @@ use crate::{transpile_ts, NativeAnnotation, PrewarmedIsolate, PrewarmedResult}; /// positional args, and return the isolate's result + captured logs. async fn run_ts(ts: &str, arg_names: &[&str], args: serde_json::Value) -> PrewarmedResult { let js = transpile_ts(ts.to_string()).expect("transpile_ts failed"); - let ann = NativeAnnotation { useragent: None, proxy: None }; + let ann = NativeAnnotation::default(); let arg_names: Vec = arg_names.iter().map(|s| s.to_string()).collect(); let mut iso = PrewarmedIsolate::spawn(String::new(), js, ann, arg_names, None); iso.wait_ready().await.expect("isolate failed to pre-warm"); @@ -232,7 +232,7 @@ export async function main(i: number): Promise { for i in 0..N { let js = js.clone(); let h = tokio::spawn(async move { - let ann = NativeAnnotation { useragent: None, proxy: None }; + let ann = NativeAnnotation::default(); let mut iso = PrewarmedIsolate::spawn(String::new(), js, ann, vec!["i".to_string()], None); iso.wait_ready().await.expect("pre-warm failed"); diff --git a/backend/windmill-runtime-nativets/tests/annotation.rs b/backend/windmill-runtime-nativets/tests/annotation.rs new file mode 100644 index 0000000000..d8319a87a5 --- /dev/null +++ b/backend/windmill-runtime-nativets/tests/annotation.rs @@ -0,0 +1,37 @@ +//! `//fetch_response_timeout` parsing. Cheap (no V8), so it stays out of the +//! e2e file. + +use windmill_runtime_nativets::get_annotation; + +#[test] +fn a_value_is_parsed_and_zero_stays_distinct_from_absent() { + // Collapsing `Some(0)` to `None` would silently reinstate the default on a + // script that explicitly asked for no limit. + assert_eq!( + get_annotation("//native\n//fetch_response_timeout 30\n").fetch_response_timeout_secs, + Some(30) + ); + assert_eq!( + get_annotation("//fetch_response_timeout 0\n").fetch_response_timeout_secs, + Some(0) + ); + assert_eq!( + get_annotation("//native\n").fetch_response_timeout_secs, + None + ); +} + +#[test] +fn a_malformed_value_falls_back_to_the_default_not_to_no_timeout() { + for src in [ + "//fetch_response_timeout abc\n", + "//fetch_response_timeout\n", + "//fetch_response_timeout -5\n", + ] { + assert_eq!( + get_annotation(src).fetch_response_timeout_secs, + None, + "{src:?} should leave the default in force" + ); + } +} diff --git a/backend/windmill-runtime-nativets/tests/fetch_response_timeout.rs b/backend/windmill-runtime-nativets/tests/fetch_response_timeout.rs new file mode 100644 index 0000000000..5118b5f93f --- /dev/null +++ b/backend/windmill-runtime-nativets/tests/fetch_response_timeout.rs @@ -0,0 +1,593 @@ +//! The nativets `fetch()` response timeout. +//! +//! Two halves of one contract, and a fix that satisfies only the first is worse +//! than no fix: +//! +//! 1. a peer that accepts a request and never answers is given up on +//! 2. a response that has begun arriving is never cut off, however long it +//! takes in total +//! +//! (2) rules out the obvious implementation — `AbortSignal.timeout(N)` around +//! every fetch would satisfy (1) and break every streaming response and long +//! download. +//! +//! Hermetic: loopback listeners, no egress. + +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpListener; +use tokio::sync::Mutex; + +use windmill_runtime_nativets::{transpile_ts, NativeAnnotation, PrewarmedIsolate}; + +/// The annotation is in whole seconds, so the tests scale around this. +const TIMEOUT_SECS: u64 = 2; + +async fn run_with_timeout_secs(ts: &str, secs: u64) -> Result { + let js = transpile_ts(ts.to_string()).expect("transpile_ts failed"); + let ann = + NativeAnnotation { fetch_response_timeout_secs: Some(secs), ..NativeAnnotation::default() }; + let mut iso = PrewarmedIsolate::spawn(String::new(), js, ann, vec![], None); + iso.wait_ready().await.expect("isolate failed to pre-warm"); + let res = iso + .start_execution("{}".to_string()) + .wait() + .await + .expect("isolate panicked"); + res.result.map(|raw| raw.get().to_string()) +} + +/// A peer that reads the request and then answers nothing, closing only after +/// `close_after` so no test leaves an isolate wedged on a pending fetch. +/// +/// The socket is held rather than dropped on accept: dropping it sends a FIN, +/// which surfaces as a connection error — the easy failure, not this one. +async fn spawn_silent_peer(seen: Arc>>, close_after: Duration) -> u16 { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + tokio::spawn(async move { + while let Ok((mut sock, _)) = listener.accept().await { + let seen = seen.clone(); + tokio::spawn(async move { + let mut buf = [0u8; 8192]; + if let Ok(n) = sock.read(&mut buf).await { + seen.lock().await.extend_from_slice(&buf[..n]); + } + tokio::time::sleep(close_after).await; + drop(sock); + }); + } + }); + port +} + +/// A peer that records the request it received and answers 200 immediately. +async fn spawn_echo_peer(seen: Arc>>) -> u16 { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + tokio::spawn(async move { + while let Ok((mut sock, _)) = listener.accept().await { + let seen = seen.clone(); + tokio::spawn(async move { + let mut buf = [0u8; 8192]; + if let Ok(n) = sock.read(&mut buf).await { + seen.lock().await.extend_from_slice(&buf[..n]); + } + let _ = sock + .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 0\r\n\r\n") + .await; + }); + } + }); + port +} + +/// A peer that responds after `headers_after`, then dribbles a chunked body out +/// over `chunks * chunk_every`. +async fn spawn_streaming_peer( + headers_after: Duration, + chunks: usize, + chunk_every: Duration, +) -> u16 { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + tokio::spawn(async move { + while let Ok((mut sock, _)) = listener.accept().await { + tokio::spawn(async move { + let mut buf = [0u8; 8192]; + let _ = sock.read(&mut buf).await; + tokio::time::sleep(headers_after).await; + if sock + .write_all( + b"HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\nTransfer-Encoding: chunked\r\n\r\n", + ) + .await + .is_err() + { + return; + } + for _ in 0..chunks { + tokio::time::sleep(chunk_every).await; + if sock.write_all(b"1\r\nx\r\n").await.is_err() { + return; + } + } + let _ = sock.write_all(b"0\r\n\r\n").await; + }); + } + }); + port +} + +const POST_TO_SILENT_PEER: &str = r#" +export async function main(): Promise { + const res = await fetch("http://127.0.0.1:{port}/orders", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ id: 1 }), + }); + return res.status; +} +"#; + +fn post_script(port: u16) -> String { + POST_TO_SILENT_PEER.replace("{port}", &port.to_string()) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_peer_that_never_answers_is_given_up_on() { + let seen = Arc::new(Mutex::new(Vec::new())); + let port = spawn_silent_peer(seen.clone(), Duration::from_secs(60)).await; + + let started = Instant::now(); + let err = run_with_timeout_secs(&post_script(port), TIMEOUT_SECS) + .await + .expect_err("fetch against a peer that never answers must not resolve"); + let elapsed = started.elapsed(); + + assert!( + elapsed >= Duration::from_secs(TIMEOUT_SECS), + "gave up after {elapsed:?}, before the configured {TIMEOUT_SECS}s -- \ + the timeout is firing on something other than the wait for a response", + ); + assert!( + elapsed < Duration::from_secs(TIMEOUT_SECS + 15), + "took {elapsed:?} to give up", + ); + + // The message has to stand on its own in a job log: a hung request is + // otherwise indistinguishable from a slow one. + assert!( + err.contains("no response headers arrived within"), + "error should explain what timed out, got: {err}", + ); + assert!( + err.contains("/orders") && err.contains("fetch_response_timeout"), + "error should name the target and how to change the limit, got: {err}", + ); + + // Without this the test would also pass if the request never went out. + let body = String::from_utf8_lossy(&seen.lock().await.clone()).to_string(); + assert!( + body.contains("POST /orders"), + "peer should have received the request, got: {body}", + ); +} + +/// The counterfactual for the test above: with the timeout disabled, the same +/// script against the same peer is still running well past the point the +/// timeout would have fired. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn without_a_timeout_the_same_request_keeps_running() { + let seen = Arc::new(Mutex::new(Vec::new())); + // Closes eventually, so the isolate unwinds instead of pinning a blocking + // task for the rest of the test binary's life. + let port = spawn_silent_peer(seen, Duration::from_secs(TIMEOUT_SECS * 3)).await; + + let still_running = tokio::time::timeout( + Duration::from_secs(TIMEOUT_SECS * 2), + run_with_timeout_secs(&post_script(port), 0), + ) + .await + .is_err(); + + assert!( + still_running, + "with the timeout disabled the request should still have been pending \ + at {}s -- if it ends on its own, the test above proves nothing", + TIMEOUT_SECS * 2, + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_streaming_body_outliving_the_timeout_is_not_cut_off() { + // Headers land fast, then the body trickles well past the timeout. A + // total-duration timeout fails here; that is the point of the test. + let port = + spawn_streaming_peer(Duration::from_millis(200), 10, Duration::from_millis(500)).await; + + let ts = format!( + r#" +export async function main(): Promise {{ + const res = await fetch("http://127.0.0.1:{port}/stream"); + return `${{res.status}}:${{(await res.text()).length}}`; +}} +"# + ); + + let started = Instant::now(); + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("a streaming response must not be interrupted"); + let elapsed = started.elapsed(); + + assert_eq!(out, "\"200:10\"", "full body should arrive intact"); + assert!( + elapsed > Duration::from_secs(TIMEOUT_SECS), + "the transfer ({elapsed:?}) has to outlast the {TIMEOUT_SECS}s timeout \ + for this to be exercising anything", + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_slow_but_answering_peer_is_not_cut_off() { + // A quarter of the window rather than half: CI runs this at + // --test-threads=10 alongside other V8 isolates, and this margin is what + // absorbs executor starvation. + let port = spawn_streaming_peer( + Duration::from_millis(1_000 * TIMEOUT_SECS / 4), + 1, + Duration::from_millis(10), + ) + .await; + + let ts = format!( + r#" +export async function main(): Promise {{ + const res = await fetch("http://127.0.0.1:{port}/slow"); + await res.text(); + return res.status; +}} +"# + ); + + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("a slow but answering peer must not be cut off"); + assert_eq!(out, "200"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn an_over_large_timeout_does_not_wrap_into_an_instant_one() { + // deno_web's setTimeout puts its delay through `webidl.converters.long`, + // which wraps at 32 bits. Unclamped, this ~46-day setting wraps negative and + // aborts immediately -- asking for a longer leash would kill every fetch. + let port = spawn_streaming_peer(Duration::from_millis(50), 1, Duration::from_millis(10)).await; + + let ts = format!( + r#" +export async function main(): Promise {{ + const res = await fetch("http://127.0.0.1:{port}/ok"); + await res.text(); + return res.status; +}} +"# + ); + + let out = run_with_timeout_secs(&ts, 4_000_000) + .await + .expect("an over-large timeout must not abort the request"); + assert_eq!(out, "200"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_caller_abort_still_wins_with_its_own_reason() { + let seen = Arc::new(Mutex::new(Vec::new())); + let port = spawn_silent_peer(seen, Duration::from_secs(60)).await; + + let ts = format!( + r#" +export async function main(): Promise {{ + const ac = new AbortController(); + setTimeout(() => ac.abort(new Error("caller_abort_marker")), 200); + try {{ + await fetch("http://127.0.0.1:{port}/probe", {{ signal: ac.signal }}); + return "unexpectedly resolved"; + }} catch (e) {{ + return String((e as Error).message); + }} +}} +"# + ); + + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("script should catch its own abort"); + assert!( + out.contains("caller_abort_marker"), + "caller's abort reason should survive being combined with ours, got: {out}", + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn the_timeout_is_catchable_as_a_timeout_error() { + // Scripts that retry on transient failures need to recognise this one; + // `TimeoutError` matches AbortSignal.timeout()'s reason. + let seen = Arc::new(Mutex::new(Vec::new())); + let port = spawn_silent_peer(seen, Duration::from_secs(60)).await; + + let ts = format!( + r#" +export async function main(): Promise {{ + try {{ + await fetch("http://127.0.0.1:{port}/probe"); + return "unexpectedly resolved"; + }} catch (e) {{ + return (e as Error).name; + }} +}} +"# + ); + + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("script should catch the timeout"); + assert_eq!(out, "\"TimeoutError\""); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn an_init_whose_members_are_inherited_is_not_flattened() { + // RequestInit is a WebIDL dictionary and deno_fetch reads its members with + // plain property gets, so they may sit on the prototype chain or be + // non-enumerable. Object spread copies neither, which would silently + // downgrade this POST to a GET and drop the header. + let seen = Arc::new(Mutex::new(Vec::new())); + let port = spawn_echo_peer(seen.clone()).await; + + let ts = format!( + r#" +declare const Object: any; +export async function main(): Promise {{ + const base = {{ method: "POST", headers: {{ "x-probe": "yes" }} }}; + const res = await fetch("http://127.0.0.1:{port}/inherited", Object.create(base)); + return res.status; +}} +"# + ); + + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("request should succeed"); + assert_eq!(out, "200"); + + let body = String::from_utf8_lossy(&seen.lock().await.clone()).to_string(); + assert!( + body.starts_with("POST /inherited"), + "inherited `method` should survive, got: {body}", + ); + assert!( + body.to_lowercase().contains("x-probe: yes"), + "inherited `headers` should survive, got: {body}", + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_non_dictionary_init_still_fails_loudly() { + // deno's dictionary converter throws on a non-object init. Copying members + // into a fresh object instead of inheriting would turn `"POST"` into + // {0:"P",1:"O",...} -- a valid dictionary with ignored keys, i.e. a silent + // GET where the caller used to get a TypeError. + let seen = Arc::new(Mutex::new(Vec::new())); + let port = spawn_echo_peer(seen).await; + + let ts = format!( + r#" +export async function main(): Promise {{ + try {{ + await fetch("http://127.0.0.1:{port}/x", "POST" as any); + return "unexpectedly resolved"; + }} catch (e) {{ + return (e as Error).name; + }} +}} +"# + ); + + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("script should catch the error"); + assert_eq!(out, "\"TypeError\""); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn an_accessor_backed_init_reads_against_its_own_receiver() { + // A getter on the init must run with the object it was defined on as `this`, + // or a private field is unreachable and it throws. Carrying the init across + // by inheritance rather than by handing it to Request would break this. + let seen = Arc::new(Mutex::new(Vec::new())); + let port = spawn_echo_peer(seen.clone()).await; + + let ts = format!( + r#" +class Init {{ + #method = "POST"; + #body = "from-private-field"; + get method(): string {{ return this.#method; }} + get body(): string {{ return this.#body; }} +}} +export async function main(): Promise {{ + const res = await fetch("http://127.0.0.1:{port}/accessor", new Init() as any); + return res.status; +}} +"# + ); + + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("an accessor-backed init must not throw"); + assert_eq!(out, "200"); + + let body = String::from_utf8_lossy(&seen.lock().await.clone()).to_string(); + assert!( + body.starts_with("POST /accessor") && body.contains("from-private-field"), + "getter-provided method and body should both reach the wire, got: {body}", + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_request_input_keeps_its_body_and_headers() { + // The wrapper builds a Request and hands that to fetch, so the body survives + // one more construction than it used to -- deno proxies it rather than + // consuming it, and this pins that. + let seen = Arc::new(Mutex::new(Vec::new())); + let port = spawn_echo_peer(seen.clone()).await; + + let ts = format!( + r#" +export async function main(): Promise {{ + const req = new Request("http://127.0.0.1:{port}/from-request", {{ + method: "PUT", + headers: {{ "x-probe": "yes" }}, + body: "payload-body", + }}); + const res = await fetch(req); + return res.status; +}} +"# + ); + + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("a Request input must work"); + assert_eq!(out, "200"); + + let body = String::from_utf8_lossy(&seen.lock().await.clone()).to_string(); + assert!( + body.starts_with("PUT /from-request") + && body.to_lowercase().contains("x-probe: yes") + && body.contains("payload-body"), + "method, headers and body should all survive, got: {body}", + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn an_already_aborted_fetch_settles_in_the_same_tick() { + // deno_fetch keeps its outer fetch non-async and returns an already-settled + // rejection untouched, because WPT pins that an aborted fetch settles in the + // same tick. Adopting it through another promise pushes the rejection behind + // any microtask queued after the call. + let ts = r#" +export async function main(): Promise { + const order: string[] = []; + const ac = new AbortController(); + ac.abort(); + const f = fetch("http://127.0.0.1:1/x", { signal: ac.signal }) + .catch(() => { order.push("fetch"); }); + Promise.resolve().then(() => { order.push("queued-after"); }); + await f; + await new Promise((r) => setTimeout(r, 0)); + return order.join(","); +} +"#; + + let out = run_with_timeout_secs(ts, TIMEOUT_SECS) + .await + .expect("script should run"); + assert_eq!( + out, "\"fetch,queued-after\"", + "the rejection must land before a microtask queued after the call", + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn the_wrapper_keeps_fetch_s_own_shape() { + // The wrapper is indistinguishable from deno's fetch on three counts a + // script can observe: its arity, the error for an empty call, and not + // depending on a mutable `Promise.prototype.then` the way an ordinary + // property lookup would. + let seen = Arc::new(Mutex::new(Vec::new())); + let port = spawn_echo_peer(seen).await; + + let ts = format!( + r#" +declare const Promise: any; +declare const AbortSignal: any; +declare const AbortController: any; +export async function main(): Promise {{ + const arity = (fetch as any).length; + + // The message, not just the type: forwarding two explicit `undefined`s + // would still throw a TypeError, just deno's invalid-URL one instead of + // its required-argument one. + let emptyCall = "resolved"; + try {{ + await (fetch as any)(); + }} catch (e) {{ + emptyCall = (e as Error).message.includes("1 argument required") + ? "required-argument" + : `other(${{(e as Error).message}})`; + }} + + // Patched only across the call: the wrapper reaches for these while + // building its return value, and awaiting under a patched Promise + // prototype would instead measure V8 treating it as a plain thenable. + const originalThen = Promise.prototype.then; + const originalAny = AbortSignal.any; + const originalAbort = AbortController.prototype.abort; + let pending: any; + try {{ + Promise.prototype.then = undefined; + (AbortSignal as any).any = undefined; + (AbortController.prototype as any).abort = undefined; + pending = fetch("http://127.0.0.1:{port}/shape"); + }} finally {{ + Promise.prototype.then = originalThen; + (AbortSignal as any).any = originalAny; + (AbortController.prototype as any).abort = originalAbort; + }} + const status = (await pending).status; + + return `${{arity}}:${{emptyCall}}:${{status}}`; +}} +"# + ); + + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("script should run"); + assert_eq!(out, "\"1:required-argument:200\""); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn the_timer_aborts_through_a_captured_intrinsic() { + // The timeout has to fire for this one: `AbortController.prototype.abort` + // is patched out for the whole wait, so an ordinary lookup would throw + // inside the timer callback and leave the request pending forever. + let seen = Arc::new(Mutex::new(Vec::new())); + let port = spawn_silent_peer(seen, Duration::from_secs(TIMEOUT_SECS * 5)).await; + + let ts = format!( + r#" +declare const AbortController: any; +export async function main(): Promise {{ + const originalAbort = AbortController.prototype.abort; + AbortController.prototype.abort = undefined; + try {{ + await fetch("http://127.0.0.1:{port}/patched-abort"); + return "unexpectedly resolved"; + }} catch (e) {{ + return (e as Error).name; + }} finally {{ + AbortController.prototype.abort = originalAbort; + }} +}} +"# + ); + + let out = run_with_timeout_secs(&ts, TIMEOUT_SECS) + .await + .expect("the timeout must still fire"); + assert_eq!(out, "\"TimeoutError\""); +} diff --git a/backend/windmill-runtime-nativets/tests/otel_e2e.rs b/backend/windmill-runtime-nativets/tests/otel_e2e.rs index cadddb2fcf..22fd18de90 100644 --- a/backend/windmill-runtime-nativets/tests/otel_e2e.rs +++ b/backend/windmill-runtime-nativets/tests/otel_e2e.rs @@ -136,7 +136,7 @@ export async function main(): Promise {{ "# ); let js = transpile_ts(ts).expect("transpile failed"); - let ann = NativeAnnotation { useragent: None, proxy: None }; + let ann = NativeAnnotation::default(); let mut iso = PrewarmedIsolate::spawn(String::new(), js, ann, vec![], None); iso.wait_ready().await.expect("isolate failed to pre-warm"); diff --git a/backend/windmill-runtime-nativets/tests/response_timeout_env.rs b/backend/windmill-runtime-nativets/tests/response_timeout_env.rs new file mode 100644 index 0000000000..a206c63134 --- /dev/null +++ b/backend/windmill-runtime-nativets/tests/response_timeout_env.rs @@ -0,0 +1,13 @@ +//! Its own test binary: the setting is read once per process through a +//! `LazyLock`, so nothing else may resolve it first. Adding a second test to +//! this file breaks that isolation. + +use windmill_runtime_nativets::default_fetch_response_timeout_secs; + +#[test] +fn the_env_var_is_what_operators_actually_set() { + // A typo in the variable's name would compile, pass every other test, and + // silently hand every operator the built-in default instead. + std::env::set_var("WINDMILL_FETCH_RESPONSE_TIMEOUT_SECS", "17"); + assert_eq!(default_fetch_response_timeout_secs(), 17); +} diff --git a/backend/windmill-worker/src/ai/utils.rs b/backend/windmill-worker/src/ai/utils.rs index 6e80551650..51a00e65a7 100644 --- a/backend/windmill-worker/src/ai/utils.rs +++ b/backend/windmill-worker/src/ai/utils.rs @@ -26,11 +26,11 @@ use windmill_queue::{flow_status::get_step_of_flow_status, MiniPulledJob}; use crate::parse_sig_of_lang; -pub fn parse_raw_script_schema( +pub async fn parse_raw_script_schema( content: &str, language: &ScriptLang, ) -> Result, Error> { - let main_arg_signature = parse_sig_of_lang(content, Some(&language), None)? + let main_arg_signature = parse_sig_of_lang(content, Some(&language), None).await? .ok_or_else(|| Error::BadConfig(format!( "Cannot parse signature for language {:?}. The language parser may not be enabled in this build.", language diff --git a/backend/windmill-worker/src/ai_executor.rs b/backend/windmill-worker/src/ai_executor.rs index 2c29721ef6..cd17c57f59 100644 --- a/backend/windmill-worker/src/ai_executor.rs +++ b/backend/windmill-worker/src/ai_executor.rs @@ -621,7 +621,7 @@ pub async fn handle_ai_agent_job( (schema, input_transforms, derived_description) } FlowModuleValue::RawScript { content, language, input_transforms, .. } => { - let schema = Some(parse_raw_script_schema(&content, &language)?); + let schema = Some(parse_raw_script_schema(&content, &language).await?); (schema, input_transforms, None) } FlowModuleValue::AIAgent { input_transforms, .. } => { diff --git a/backend/windmill-worker/src/global_cache.rs b/backend/windmill-worker/src/global_cache.rs index 5e9666722a..3f6205f100 100644 --- a/backend/windmill-worker/src/global_cache.rs +++ b/backend/windmill-worker/src/global_cache.rs @@ -147,14 +147,17 @@ pub fn extract_tar(tar: bytes::Bytes, folder: &str) -> error::Result<()> { Ok(()) } -/// Two-tier cache load: check local disk first, then fall back to instance object store. +/// Two-tier cache load: check local disk first, then fall back to the shared object store. +/// +/// "Shared" is the worker group's own store when its config overrides one, the instance store +/// otherwise — see [`windmill_object_store::get_cache_object_store`]. /// Returns `(hit, log_message)`. pub async fn load_cache(bin_path: &str, _remote_path: &str, is_dir: bool) -> (bool, String) { if tokio::fs::metadata(&bin_path).await.is_ok() { (true, format!("loaded from local cache: {}\n", bin_path)) } else { #[cfg(all(feature = "enterprise", feature = "parquet"))] - if let Some(os) = windmill_object_store::get_object_store().await { + if let Some(os) = windmill_object_store::get_cache_object_store().await { let started = std::time::Instant::now(); if let Ok(mut x) = windmill_object_store::attempt_fetch_bytes(os, _remote_path).await { @@ -200,13 +203,15 @@ pub async fn load_cache(bin_path: &str, _remote_path: &str, is_dir: bool) -> (bo } } -/// Whether this worker can push to the instance object store at all — the features are +/// Whether this worker can push to the shared object store at all — the features are /// compiled in and a store is loaded. False on builds without them, where `save_cache` /// only ever writes to the worker's own disk. pub async fn object_store_available() -> bool { #[cfg(all(feature = "enterprise", feature = "parquet"))] { - windmill_object_store::get_object_store().await.is_some() + windmill_object_store::get_cache_object_store() + .await + .is_some() } #[cfg(not(all(feature = "enterprise", feature = "parquet")))] { @@ -214,14 +219,14 @@ pub async fn object_store_available() -> bool { } } -/// Whether a binary/bundle is in the instance object store, ignoring the local cache. +/// Whether a binary/bundle is in the shared object store, ignoring the local cache. /// /// The deploy-time prebuild asks this rather than [`exists_in_cache`]: a copy on the /// building worker's own disk is exactly the state the prebuild exists to fix, so /// answering from it would latch a failed upload into a permanent skip. pub async fn exists_in_object_store(_remote_path: &str) -> bool { #[cfg(all(feature = "enterprise", feature = "parquet"))] - if let Some(os) = windmill_object_store::get_object_store().await { + if let Some(os) = windmill_object_store::get_cache_object_store().await { return os .head(&windmill_object_store::object_store_reexports::Path::from( _remote_path, @@ -241,18 +246,18 @@ pub async fn ensure_pushed_to_object_store(remote_path: &str) -> error::Result<( return Ok(()); } Err(error::Error::ExecutionErr(format!( - "the binary was built but did not reach the instance object store at {remote_path}, \ + "the binary was built but did not reach the object store at {remote_path}, \ so no other worker can load it" ))) } -/// Check whether a binary/bundle exists in local cache or instance object store. +/// Check whether a binary/bundle exists in local cache or the shared object store. pub async fn exists_in_cache(bin_path: &str, _remote_path: &str) -> bool { if tokio::fs::metadata(&bin_path).await.is_ok() { return true; } else { #[cfg(all(feature = "enterprise", feature = "parquet"))] - if let Some(os) = windmill_object_store::get_object_store().await { + if let Some(os) = windmill_object_store::get_cache_object_store().await { return os .get(&windmill_object_store::object_store_reexports::Path::from( _remote_path, @@ -264,7 +269,7 @@ pub async fn exists_in_cache(bin_path: &str, _remote_path: &str) -> bool { } } -/// Two-tier cache write: upload to instance object store, then copy to local disk. +/// Two-tier cache write: upload to the shared object store, then copy to local disk. pub async fn save_cache( local_cache_path: &str, _remote_cache_path: &str, @@ -275,7 +280,7 @@ pub async fn save_cache( let mut _cached_to_s3 = false; #[cfg(all(feature = "enterprise", feature = "parquet"))] - if let Some(os) = windmill_object_store::get_object_store().await { + if let Some(os) = windmill_object_store::get_cache_object_store().await { use windmill_object_store::object_store_reexports::Path; let file_to_cache = if is_dir { let tar_path = format!( diff --git a/backend/windmill-worker/src/php_executor.rs b/backend/windmill-worker/src/php_executor.rs index 1ede0933cb..03328b82cb 100644 --- a/backend/windmill-worker/src/php_executor.rs +++ b/backend/windmill-worker/src/php_executor.rs @@ -14,7 +14,7 @@ use windmill_common::{ }; use windmill_queue::MiniPulledJob; -use windmill_parser::Typ; +use windmill_parser::{MainArgSignature, Typ}; use windmill_queue::{append_logs, CanceledBy}; use crate::{ @@ -40,6 +40,43 @@ lazy_static::lazy_static! { const COMPOSER_LOCK_SPLIT: &str = "\nLOCK\n"; +static PHP_PARSER_SLOT: tokio::sync::Semaphore = tokio::sync::Semaphore::const_new(1); + +pub(crate) async fn parse_php_signature( + code: &str, + main_override: Option, +) -> Result { + parse_php_signature_with_slot(code, main_override, &PHP_PARSER_SLOT).await +} + +async fn parse_php_signature_with_slot( + code: &str, + main_override: Option, + slot: &'static tokio::sync::Semaphore, +) -> Result { + let acquire = slot.acquire(); + tokio::pin!(acquire); + // Retain the acquisition across the warning to preserve its FIFO queue position. + let permit = tokio::select! { + permit = &mut acquire => permit, + _ = tokio::time::sleep(std::time::Duration::from_secs(1)) => { + tracing::warn!("Waiting over a second for PHP signature parser capacity"); + acquire.await + } + } + .map_err(to_anyhow)?; + let code = code.to_owned(); + tokio::task::spawn_blocking(move || { + // Parsing walks the entire AST. Keep its stack off async workers and retain + // the process-wide CPU limit even if the awaiting job is cancelled. + let _permit = permit; + windmill_parser_php::parse_php_signature(&code, main_override) + }) + .await + .map_err(|e| error::Error::internal_err(format!("PHP signature parsing task failed: {e}")))? + .map_err(Into::into) +} + pub fn parse_php_imports(code: &str) -> anyhow::Result> { let find_requirements = code .lines() @@ -333,11 +370,9 @@ pub async fn handle_php_job( let main_override = job.script_entrypoint_override.as_deref(); let write_wrapper_f = async { - let args = windmill_parser_php::parse_php_signature( - inner_content, - main_override.map(ToString::to_string), - )? - .args; + let args = parse_php_signature(inner_content, main_override.map(ToString::to_string)) + .await? + .args; let args_to_include = args .iter() @@ -491,3 +526,51 @@ try {{ .await?; read_result(job_dir, None).await } + +#[cfg(test)] +mod tests { + use super::parse_php_signature_with_slot; + + #[test] + fn cancelled_parse_retains_slot_until_blocking_work_finishes() { + static SLOT: tokio::sync::Semaphore = tokio::sync::Semaphore::const_new(1); + let runtime = tokio::runtime::Builder::new_current_thread() + .max_blocking_threads(1) + .enable_all() + .build() + .unwrap(); + let (release_tx, release_rx) = std::sync::mpsc::channel(); + let (started_tx, started_rx) = tokio::sync::oneshot::channel(); + let blocker = runtime.spawn_blocking(move || { + started_tx.send(()).unwrap(); + let _ = release_rx.recv(); + }); + + runtime.block_on(async { + started_rx.await.unwrap(); + let parse = tokio::spawn(parse_php_signature_with_slot( + ") { use crate::global_cache::build_tar_and_push; - use windmill_object_store::get_object_store; + use windmill_object_store::get_cache_object_store; while let Some(task) = rx.recv().await { - if let Some(os) = get_object_store().await { + if let Some(os) = get_cache_object_store().await { match build_tar_and_push(os, task.venv_path.clone(), task.cache_dir, None, false).await { Ok(()) => { @@ -144,7 +144,7 @@ pub fn has_relative_imports(content: &str) -> bool { use crate::global_cache::pull_from_tar; #[cfg(all(feature = "enterprise", feature = "parquet"))] -use windmill_object_store::OBJECT_STORE_SETTINGS; +use windmill_object_store::get_cache_object_store; use crate::{ common::{ @@ -2449,7 +2449,7 @@ pub async fn handle_python_reqs( } #[cfg(all(feature = "enterprise", feature = "parquet"))] - if OBJECT_STORE_SETTINGS.read().await.is_none() { + if get_cache_object_store().await.is_none() { (s3_pull, s3_push) = (false, false); } @@ -2907,7 +2907,7 @@ pub async fn handle_python_reqs( #[cfg(all(feature = "enterprise", feature = "parquet"))] if is_not_pro { - if let Some(os) = windmill_object_store::get_object_store().await { + if let Some(os) = windmill_object_store::get_cache_object_store().await { tokio::select! { // Cancel was called on the job _ = kill_rx.recv() => return Err(Error::from(anyhow::anyhow!("S3 pull was canceled"))), diff --git a/backend/windmill-worker/src/universal_pkg_installer.rs b/backend/windmill-worker/src/universal_pkg_installer.rs index f36a5e6fc9..1126c94d17 100644 --- a/backend/windmill-worker/src/universal_pkg_installer.rs +++ b/backend/windmill-worker/src/universal_pkg_installer.rs @@ -333,7 +333,7 @@ pub async fn par_install_language_dependencies_all_at_once< mark_success(path.clone(), job_id, w_id).await; #[cfg(all(feature = "enterprise", feature = "parquet"))] { - if let Some(os) = windmill_object_store::get_object_store().await { + if let Some(os) = windmill_object_store::get_cache_object_store().await { let language_name = _language_name.to_owned(); tokio::spawn(async move { if let Err(e) = crate::global_cache::build_tar_and_push( @@ -790,7 +790,7 @@ async fn try_install_one_detached<'a, T: Clone + std::marker::Send + Sync + 'a + #[cfg(all(feature = "enterprise", feature = "parquet"))] let s3_pull_future = if is_not_pro { - if let Some(os) = windmill_object_store::get_object_store().await { + if let Some(os) = windmill_object_store::get_cache_object_store().await { Some(crate::global_cache::pull_from_tar( os, dep.path.clone(), @@ -893,7 +893,7 @@ async fn try_install_one_detached<'a, T: Clone + std::marker::Send + Sync + 'a + #[cfg(all(feature = "enterprise", feature = "parquet"))] { - if let Some(os) = windmill_object_store::get_object_store().await { + if let Some(os) = windmill_object_store::get_cache_object_store().await { let language_name = _language_name.to_string(); let platform_agnostic = _platform_agnostic; let path = dep.path.clone(); @@ -954,8 +954,7 @@ async fn print_success( } #[cfg(all(feature = "enterprise", feature = "parquet"))] - if windmill_object_store::OBJECT_STORE_SETTINGS - .read() + if windmill_object_store::get_cache_object_store() .await .is_none() { diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 8f11502f62..d893ec071b 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -3867,7 +3867,8 @@ pub async fn run_worker( let job_result = windmill_common::log_context::with_log_context( log_ctx, async { - let result = handle_queued_job( + // Keep large job-phase futures boxed to limit debug polling frames. + let result = Box::pin(handle_queued_job( arc_job.clone(), raw_code, raw_lock, @@ -3888,7 +3889,7 @@ pub async fn run_worker( flow_runners, #[cfg(feature = "benchmark")] &mut bench, - ) + )) .await; record_job_span_status(&result); result @@ -5159,7 +5160,7 @@ async fn try_validate_schema( code, language, job.script_entrypoint_override.clone(), - )? { + ).await? { Ok(Some(schema_validator_from_main_arg_sig(&sig))) } else { Err(anyhow!("Job was expected to validate the arguments schema, but no schema was provided and couldn't be inferred from the script for language `{language:?}`. Try removing schema validation for this job").into()) @@ -5527,7 +5528,7 @@ async fn handle_code_execution_job( .await?; let language = language.clone(); - let result = run_language_executor( + let result = Box::pin(run_language_executor( job, conn, client, @@ -5552,7 +5553,7 @@ async fn handle_code_execution_job( &modules, false, in_pipeline, - ) + )) .await; record_declared_warehouse_write(job, conn, code, &result).await; result @@ -6466,7 +6467,7 @@ mount {{ .await; if let Connection::Sql(db) = conn { - volume_setup = crate::volume_oss::setup_volumes_sql_worker( + volume_setup = Box::pin(crate::volume_oss::setup_volumes_sql_worker( &volume_mounts, db, &job.workspace_id, @@ -6479,10 +6480,10 @@ mount {{ language, &mut envs, &mut shared_mount, - ) + )) .await?; } else if let Connection::Http(http) = conn { - volume_setup = crate::volume_oss::setup_volumes_http_worker( + volume_setup = Box::pin(crate::volume_oss::setup_volumes_http_worker( &volume_mounts, http, &job.workspace_id, @@ -6495,7 +6496,7 @@ mount {{ language, &mut envs, &mut shared_mount, - ) + )) .await?; } } @@ -6991,7 +6992,7 @@ mount {{ if let Some(ref vol_client) = volume_setup.client { if let Connection::Sql(db) = conn { - crate::volume_oss::sync_volumes_sql_worker( + Box::pin(crate::volume_oss::sync_volumes_sql_worker( &volume_setup.states, &volume_setup.writable, vol_client, @@ -7001,13 +7002,13 @@ mount {{ worker_name, conn, result.is_ok(), - ) + )) .await; } } if let Connection::Http(http) = conn { - crate::volume_oss::sync_volumes_http_worker( + Box::pin(crate::volume_oss::sync_volumes_http_worker( &volume_setup.states, &volume_setup.writable, http, @@ -7016,7 +7017,7 @@ mount {{ worker_name, conn, result.is_ok(), - ) + )) .await; } @@ -7051,7 +7052,7 @@ mount {{ result } -pub fn parse_sig_of_lang( +pub async fn parse_sig_of_lang( code: &str, language: Option<&ScriptLang>, main_override: Option, @@ -7086,10 +7087,9 @@ pub fn parse_sig_of_lang( ScriptLang::DuckDb => Some(windmill_parser_sql::parse_duckdb_sig(code)?), ScriptLang::OracleDB => Some(windmill_parser_sql::parse_oracledb_sig(code)?), #[cfg(feature = "php")] - ScriptLang::Php => Some(windmill_parser_php::parse_php_signature( - code, - main_override, - )?), + ScriptLang::Php => { + Some(crate::php_executor::parse_php_signature(code, main_override).await?) + } #[cfg(not(feature = "php"))] ScriptLang::Php => None, #[cfg(feature = "rust")] diff --git a/backend/windmill-worker/src/worker_lockfiles.rs b/backend/windmill-worker/src/worker_lockfiles.rs index 9bba9a99e6..1f3b9f1d56 100644 --- a/backend/windmill-worker/src/worker_lockfiles.rs +++ b/backend/windmill-worker/src/worker_lockfiles.rs @@ -174,7 +174,7 @@ async fn handle_build_binary_job( if !crate::global_cache::object_store_available().await { return Ok(to_raw_value_owned(json!({ "status": "skipped", - "reason": "this worker cannot reach the instance object store, so the binary \ + "reason": "this worker cannot reach an object store, so the binary \ would not be shared with other workers", }))); } diff --git a/cli/src/commands/trigger/trigger.ts b/cli/src/commands/trigger/trigger.ts index b011a0f0b7..4e6b0eeb16 100644 --- a/cli/src/commands/trigger/trigger.ts +++ b/cli/src/commands/trigger/trigger.ts @@ -224,9 +224,12 @@ export async function pushTrigger( } } +// `enabled` is operational state a sync deliberately does not carry: the server strips it from +// the workspace export and the push below never sends it, so a created trigger comes up enabled +// and pausing one stays a local decision. type NativeTriggerFile = Omit< NativeTrigger, - "external_id" | "workspace_id" | "error" + "external_id" | "workspace_id" | "error" | "enabled" >; export async function pushNativeTrigger( @@ -264,6 +267,7 @@ export async function pushNativeTrigger( service_config: result.service_config, error: result.error, summary: result.summary, + enabled: result.enabled, }; log.debug(`Native trigger ${serviceName}/${externalId} exists on remote`); } catch { diff --git a/frontend/src/lib/components/ObjectStoreConfigSettings.svelte b/frontend/src/lib/components/ObjectStoreConfigSettings.svelte index 14658a4228..9ddd7f7fdc 100644 --- a/frontend/src/lib/components/ObjectStoreConfigSettings.svelte +++ b/frontend/src/lib/components/ObjectStoreConfigSettings.svelte @@ -1,3 +1,42 @@ + +