diff --git a/backend/.sqlx/query-03ae5b1c912b13a8a7aadf50cb4984a2ea952e782fd52eb3088454690bd13dd1.json b/backend/.sqlx/query-03ae5b1c912b13a8a7aadf50cb4984a2ea952e782fd52eb3088454690bd13dd1.json index 60404da57c..c26b190f79 100644 --- a/backend/.sqlx/query-03ae5b1c912b13a8a7aadf50cb4984a2ea952e782fd52eb3088454690bd13dd1.json +++ b/backend/.sqlx/query-03ae5b1c912b13a8a7aadf50cb4984a2ea952e782fd52eb3088454690bd13dd1.json @@ -39,7 +39,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-08f288d2781d823e109a9e5b8848234ca7d1efeee9661f3901f298da375e73f7.json b/backend/.sqlx/query-08f288d2781d823e109a9e5b8848234ca7d1efeee9661f3901f298da375e73f7.json index 2be39fce26..4bcf3c6ce3 100644 --- a/backend/.sqlx/query-08f288d2781d823e109a9e5b8848234ca7d1efeee9661f3901f298da375e73f7.json +++ b/backend/.sqlx/query-08f288d2781d823e109a9e5b8848234ca7d1efeee9661f3901f298da375e73f7.json @@ -130,28 +130,28 @@ }, { "ordinal": 25, - "name": "teams_command_script", - "type_info": "Text" - }, - { - "ordinal": 26, - "name": "teams_team_id", - "type_info": "Text" - }, - { - "ordinal": 27, - "name": "teams_team_name", - "type_info": "Text" - }, - { - "ordinal": 28, "name": "ai_models", "type_info": "VarcharArray" }, { - "ordinal": 29, + "ordinal": 26, "name": "code_completion_model", "type_info": "Varchar" + }, + { + "ordinal": 27, + "name": "teams_command_script", + "type_info": "Text" + }, + { + "ordinal": 28, + "name": "teams_team_id", + "type_info": "Text" + }, + { + "ordinal": 29, + "name": "teams_team_name", + "type_info": "Text" } ], "parameters": { @@ -185,10 +185,10 @@ true, true, true, - true, - true, - true, false, + true, + true, + true, true ] }, diff --git a/backend/.sqlx/query-29624682d687790dd199c4af759132d79fdb2982de111cb5fd43e3d9ecd0f15e.json b/backend/.sqlx/query-29624682d687790dd199c4af759132d79fdb2982de111cb5fd43e3d9ecd0f15e.json index 7623243f07..0967139406 100644 --- a/backend/.sqlx/query-29624682d687790dd199c4af759132d79fdb2982de111cb5fd43e3d9ecd0f15e.json +++ b/backend/.sqlx/query-29624682d687790dd199c4af759132d79fdb2982de111cb5fd43e3d9ecd0f15e.json @@ -68,7 +68,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-d2def87d7f7901eebc65082f7df5e0a33e5702b25c3db3affa06155e90480e42.json b/backend/.sqlx/query-2a510a8bec98055796f987d86c344ca116895d71de09db338ec09e425dcebe5e.json similarity index 63% rename from backend/.sqlx/query-d2def87d7f7901eebc65082f7df5e0a33e5702b25c3db3affa06155e90480e42.json rename to backend/.sqlx/query-2a510a8bec98055796f987d86c344ca116895d71de09db338ec09e425dcebe5e.json index 58f03dbec2..23337708f8 100644 --- a/backend/.sqlx/query-d2def87d7f7901eebc65082f7df5e0a33e5702b25c3db3affa06155e90480e42.json +++ b/backend/.sqlx/query-2a510a8bec98055796f987d86c344ca116895d71de09db338ec09e425dcebe5e.json @@ -1,50 +1,35 @@ { "db_name": "PostgreSQL", - "query": "SELECT * FROM job_perms WHERE job_id = $1 AND workspace_id = $2", + "query": "SELECT email, username, is_admin, is_operator, groups, folders FROM job_perms WHERE job_id = $1 AND workspace_id = $2", "describe": { "columns": [ { "ordinal": 0, - "name": "job_id", - "type_info": "Uuid" - }, - { - "ordinal": 1, "name": "email", "type_info": "Varchar" }, { - "ordinal": 2, + "ordinal": 1, "name": "username", "type_info": "Varchar" }, { - "ordinal": 3, + "ordinal": 2, "name": "is_admin", "type_info": "Bool" }, { - "ordinal": 4, + "ordinal": 3, "name": "is_operator", "type_info": "Bool" }, { - "ordinal": 5, - "name": "created_at", - "type_info": "Timestamp" - }, - { - "ordinal": 6, - "name": "workspace_id", - "type_info": "Varchar" - }, - { - "ordinal": 7, + "ordinal": 4, "name": "groups", "type_info": "TextArray" }, { - "ordinal": 8, + "ordinal": 5, "name": "folders", "type_info": "JsonbArray" } @@ -61,11 +46,8 @@ false, false, false, - false, - false, - false, false ] }, - "hash": "d2def87d7f7901eebc65082f7df5e0a33e5702b25c3db3affa06155e90480e42" + "hash": "2a510a8bec98055796f987d86c344ca116895d71de09db338ec09e425dcebe5e" } diff --git a/backend/.sqlx/query-488dd591096b2b47787afdc3a1d73917ed13269f2ee20b86df79fca2c8efe672.json b/backend/.sqlx/query-488dd591096b2b47787afdc3a1d73917ed13269f2ee20b86df79fca2c8efe672.json index 8622eb9125..e6d4991053 100644 --- a/backend/.sqlx/query-488dd591096b2b47787afdc3a1d73917ed13269f2ee20b86df79fca2c8efe672.json +++ b/backend/.sqlx/query-488dd591096b2b47787afdc3a1d73917ed13269f2ee20b86df79fca2c8efe672.json @@ -68,7 +68,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-4aaab98ebdaa90f1edf49ac96fba6c391c4d0054a618b861464ee37239f1f1e0.json b/backend/.sqlx/query-4aaab98ebdaa90f1edf49ac96fba6c391c4d0054a618b861464ee37239f1f1e0.json index 1ae47774dd..166d6b84a0 100644 --- a/backend/.sqlx/query-4aaab98ebdaa90f1edf49ac96fba6c391c4d0054a618b861464ee37239f1f1e0.json +++ b/backend/.sqlx/query-4aaab98ebdaa90f1edf49ac96fba6c391c4d0054a618b861464ee37239f1f1e0.json @@ -135,7 +135,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-55cb03040bc2a8c53dd7fbb42bbdcc40f463cbc52d94ed9315cf9a547d4c89f2.json b/backend/.sqlx/query-55cb03040bc2a8c53dd7fbb42bbdcc40f463cbc52d94ed9315cf9a547d4c89f2.json index 9243288c9d..14685a8bfa 100644 --- a/backend/.sqlx/query-55cb03040bc2a8c53dd7fbb42bbdcc40f463cbc52d94ed9315cf9a547d4c89f2.json +++ b/backend/.sqlx/query-55cb03040bc2a8c53dd7fbb42bbdcc40f463cbc52d94ed9315cf9a547d4c89f2.json @@ -130,28 +130,28 @@ }, { "ordinal": 25, - "name": "teams_command_script", - "type_info": "Text" - }, - { - "ordinal": 26, - "name": "teams_team_id", - "type_info": "Text" - }, - { - "ordinal": 27, - "name": "teams_team_name", - "type_info": "Text" - }, - { - "ordinal": 28, "name": "ai_models", "type_info": "VarcharArray" }, { - "ordinal": 29, + "ordinal": 26, "name": "code_completion_model", "type_info": "Varchar" + }, + { + "ordinal": 27, + "name": "teams_command_script", + "type_info": "Text" + }, + { + "ordinal": 28, + "name": "teams_team_id", + "type_info": "Text" + }, + { + "ordinal": 29, + "name": "teams_team_name", + "type_info": "Text" } ], "parameters": { @@ -185,10 +185,10 @@ true, true, true, - true, - true, - true, false, + true, + true, + true, true ] }, diff --git a/backend/.sqlx/query-61c29d684e8e683e839a6d7210b3b9b96854e5bfd752e45922c358c42ebea0c4.json b/backend/.sqlx/query-61c29d684e8e683e839a6d7210b3b9b96854e5bfd752e45922c358c42ebea0c4.json index 6354dae16b..1bc34cfbc0 100644 --- a/backend/.sqlx/query-61c29d684e8e683e839a6d7210b3b9b96854e5bfd752e45922c358c42ebea0c4.json +++ b/backend/.sqlx/query-61c29d684e8e683e839a6d7210b3b9b96854e5bfd752e45922c358c42ebea0c4.json @@ -59,7 +59,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-75451b6d48e4c26812ae64981d0d968b8fb0bf4374a2fccc167fa879bad7078f.json b/backend/.sqlx/query-75451b6d48e4c26812ae64981d0d968b8fb0bf4374a2fccc167fa879bad7078f.json index 287284b86b..98de576f13 100644 --- a/backend/.sqlx/query-75451b6d48e4c26812ae64981d0d968b8fb0bf4374a2fccc167fa879bad7078f.json +++ b/backend/.sqlx/query-75451b6d48e4c26812ae64981d0d968b8fb0bf4374a2fccc167fa879bad7078f.json @@ -59,7 +59,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-7e93c924e3fc51f8c26df26e5d09d60e3a3a40b90421aaf589c4c3bcc5a45ec8.json b/backend/.sqlx/query-7e93c924e3fc51f8c26df26e5d09d60e3a3a40b90421aaf589c4c3bcc5a45ec8.json deleted file mode 100644 index 044c752555..0000000000 --- a/backend/.sqlx/query-7e93c924e3fc51f8c26df26e5d09d60e3a3a40b90421aaf589c4c3bcc5a45ec8.json +++ /dev/null @@ -1,29 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "SELECT v2_job_completed.id AS \"id!\", flow_status AS \"flow_status!: Json\"\n FROM v2_job_completed\n INNER JOIN v2_job ON (v2_job_completed.id = v2_job.id)\n WHERE parent_job = $1 AND v2_job_completed.workspace_id = $2 AND flow_status IS NOT NULL", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id!", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "flow_status!: Json", - "type_info": "Jsonb" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Text" - ] - }, - "nullable": [ - false, - true - ] - }, - "hash": "7e93c924e3fc51f8c26df26e5d09d60e3a3a40b90421aaf589c4c3bcc5a45ec8" -} diff --git a/backend/.sqlx/query-804fc11e35f4afc0db194b6fe2594f91df7e588d4d2431bc85f4d8734920c8bf.json b/backend/.sqlx/query-804fc11e35f4afc0db194b6fe2594f91df7e588d4d2431bc85f4d8734920c8bf.json index 2179199223..967f53c987 100644 --- a/backend/.sqlx/query-804fc11e35f4afc0db194b6fe2594f91df7e588d4d2431bc85f4d8734920c8bf.json +++ b/backend/.sqlx/query-804fc11e35f4afc0db194b6fe2594f91df7e588d4d2431bc85f4d8734920c8bf.json @@ -32,7 +32,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-89de3ff8ab32e545efcbcda05f994cb1a32c4991cbd25046282d34272587d2de.json b/backend/.sqlx/query-89de3ff8ab32e545efcbcda05f994cb1a32c4991cbd25046282d34272587d2de.json index d90a9646de..e7f486d2a4 100644 --- a/backend/.sqlx/query-89de3ff8ab32e545efcbcda05f994cb1a32c4991cbd25046282d34272587d2de.json +++ b/backend/.sqlx/query-89de3ff8ab32e545efcbcda05f994cb1a32c4991cbd25046282d34272587d2de.json @@ -59,7 +59,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-ab04cda71f8e2be9acbecabe1ee5ef756b8e5c1955fbe111df9ee171dc262338.json b/backend/.sqlx/query-ab04cda71f8e2be9acbecabe1ee5ef756b8e5c1955fbe111df9ee171dc262338.json index 3284fb1846..aefb084bba 100644 --- a/backend/.sqlx/query-ab04cda71f8e2be9acbecabe1ee5ef756b8e5c1955fbe111df9ee171dc262338.json +++ b/backend/.sqlx/query-ab04cda71f8e2be9acbecabe1ee5ef756b8e5c1955fbe111df9ee171dc262338.json @@ -63,7 +63,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-d15f02f090b8d1a7e816fe11b2e0867540ab6bb02ac6bf82decc220dce0ab048.json b/backend/.sqlx/query-d15f02f090b8d1a7e816fe11b2e0867540ab6bb02ac6bf82decc220dce0ab048.json index ba3589bf24..7a54e8ea3f 100644 --- a/backend/.sqlx/query-d15f02f090b8d1a7e816fe11b2e0867540ab6bb02ac6bf82decc220dce0ab048.json +++ b/backend/.sqlx/query-d15f02f090b8d1a7e816fe11b2e0867540ab6bb02ac6bf82decc220dce0ab048.json @@ -40,7 +40,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-ec7836df5f9056ec70015800b7f4feaeb1b671120f5f8c98fca8c89c6587fc35.json b/backend/.sqlx/query-ec7836df5f9056ec70015800b7f4feaeb1b671120f5f8c98fca8c89c6587fc35.json index 81093c5f12..a99ef5bd04 100644 --- a/backend/.sqlx/query-ec7836df5f9056ec70015800b7f4feaeb1b671120f5f8c98fca8c89c6587fc35.json +++ b/backend/.sqlx/query-ec7836df5f9056ec70015800b7f4feaeb1b671120f5f8c98fca8c89c6587fc35.json @@ -54,7 +54,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/.sqlx/query-ff0403790674cdb07022af71c2377afbd8b3a660b3be27514b517c077c63c238.json b/backend/.sqlx/query-ff0403790674cdb07022af71c2377afbd8b3a660b3be27514b517c077c63c238.json index 6a4a3b3d75..95fca28ce0 100644 --- a/backend/.sqlx/query-ff0403790674cdb07022af71c2377afbd8b3a660b3be27514b517c077c63c238.json +++ b/backend/.sqlx/query-ff0403790674cdb07022af71c2377afbd8b3a660b3be27514b517c077c63c238.json @@ -63,7 +63,8 @@ "rust", "ansible", "csharp", - "oracledb" + "oracledb", + "nu" ] } } diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 0a87924669..a54011eef3 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -1803,6 +1803,7 @@ async fn handle_zombie_jobs(db: &Pool, base_internal_url: &str, worker *SCRIPT_TOKEN_EXPIRY, &job.email, &job.id, + None, ) .await .expect("could not create job token"); diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 695071b489..2c1e47f2af 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -315,7 +315,7 @@ mod suspend_resume { let second = completed.next().await.unwrap(); // print_job(second, &db).await; - let token = windmill_worker::create_token_for_owner(&db, "test-workspace", "u/test-user", "", 100, "", &Uuid::nil()).await.unwrap(); + let token = windmill_worker::create_token_for_owner(&db, "test-workspace", "u/test-user", "", 100, "", &Uuid::nil(), None).await.unwrap(); let secret = reqwest::get(format!( "http://localhost:{port}/api/w/test-workspace/jobs/job_signature/{second}/0?token={token}&approver=ruben" )) @@ -418,7 +418,7 @@ mod suspend_resume { /* ... and send a request resume it. */ let second = completed.next().await.unwrap(); - let token = windmill_worker::create_token_for_owner(&db, "test-workspace", "u/test-user", "", 100, "", &Uuid::nil()).await.unwrap(); + let token = windmill_worker::create_token_for_owner(&db, "test-workspace", "u/test-user", "", 100, "", &Uuid::nil(), None).await.unwrap(); let secret = reqwest::get(format!( "http://localhost:{port}/api/w/test-workspace/jobs/job_signature/{second}/0?token={token}" )) @@ -3806,6 +3806,7 @@ async fn test_result_format(db: Pool) { 100, "", &Uuid::nil(), + None, ) .await .unwrap(); diff --git a/backend/windmill-api/src/resources.rs b/backend/windmill-api/src/resources.rs index 2a3331ef9a..a3f7ee3d20 100644 --- a/backend/windmill-api/src/resources.rs +++ b/backend/windmill-api/src/resources.rs @@ -1208,7 +1208,11 @@ async fn update_resource_type( Ok(format!("resource_type {} updated", name)) } -#[cfg(any(feature = "postgres_trigger", feature = "mqtt_trigger", all(feature = "sqs_trigger", feature = "enterprise")))] +#[cfg(any( + feature = "postgres_trigger", + feature = "mqtt_trigger", + all(feature = "sqs_trigger", feature = "enterprise") +))] pub async fn try_get_resource_from_db_as( authed: ApiAuthed, user_db: Option, diff --git a/backend/windmill-common/src/auth.rs b/backend/windmill-common/src/auth.rs index 451157b115..d63fddc670 100644 --- a/backend/windmill-common/src/auth.rs +++ b/backend/windmill-common/src/auth.rs @@ -24,15 +24,12 @@ pub struct JWTAuthClaims { #[derive(Deserialize)] pub struct JobPerms { - pub workspace_id: String, - pub job_id: String, pub email: String, pub username: String, pub is_admin: bool, pub is_operator: bool, pub groups: Vec, pub folders: Vec, - pub created_at: chrono::NaiveDateTime, } impl From for Authed { diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index af5814aedc..a6979ad30e 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -142,16 +142,20 @@ fn format_pull_query(peek: String) -> String { raw_flow, script_entrypoint_override, preprocessed FROM v2_job WHERE id = (SELECT id FROM peek) - ) SELECT id, workspace_id, parent_job, created_by, started_at, scheduled_for, - runnable_id, runnable_path, args, canceled_by, - canceled_reason, kind, trigger, trigger_kind, permissioned_as, permissioned_as_email, - flow_status, script_lang, - same_worker, pre_run_error, visible_to_owner, - tag, concurrent_limit, concurrency_time_window_s, flow_innermost_root_job, - timeout, flow_step_id, cache_ttl, priority, raw_code, raw_lock, raw_flow, - script_entrypoint_override, preprocessed + ) SELECT j.id, j.workspace_id, j.parent_job, j.created_by, started_at, scheduled_for, + j.runnable_id, j.runnable_path, j.args, canceled_by, + canceled_reason, j.kind, j.trigger, j.trigger_kind, j.permissioned_as, j.permissioned_as_email, + flow_status, j.script_lang, + j.same_worker, j.pre_run_error, j.visible_to_owner, + j.tag, j.concurrent_limit, j.concurrency_time_window_s, j.flow_innermost_root_job, + j.timeout, j.flow_step_id, j.cache_ttl, j.priority, j.raw_code, j.raw_lock, j.raw_flow, + j.script_entrypoint_override, j.preprocessed, pj.runnable_path as parent_runnable_path, + p.email as permissioned_as_email, p.username as permissioned_as_username, p.is_admin as permissioned_as_is_admin, + p.is_operator as permissioned_as_is_operator, p.groups as permissioned_as_groups, p.folders as permissioned_as_folders FROM q, j - LEFT JOIN v2_job_status f USING (id)", + LEFT JOIN v2_job_status f USING (id) + LEFT JOIN job_perms p ON p.job_id = j.id + LEFT JOIN v2_job pj ON j.parent_job = pj.id", peek ); tracing::debug!("pull query: {}", r); diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 7f1473055a..394d04e366 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -2002,6 +2002,13 @@ pub struct PulledJob { pub raw_code: Option, pub raw_lock: Option, pub raw_flow: Option>>, + pub parent_runnable_path: Option, + pub permissioned_as_email: Option, + pub permissioned_as_username: Option, + pub permissioned_as_is_admin: Option, + pub permissioned_as_is_operator: Option, + pub permissioned_as_groups: Option>, + pub permissioned_as_folders: Option>, } impl std::ops::Deref for PulledJob { diff --git a/backend/windmill-worker/src/ansible_executor.rs b/backend/windmill-worker/src/ansible_executor.rs index a565ff6820..c98837ac5d 100644 --- a/backend/windmill-worker/src/ansible_executor.rs +++ b/backend/windmill-worker/src/ansible_executor.rs @@ -26,8 +26,8 @@ use crate::{ }, handle_child::handle_child, python_executor::{create_dependencies_dir, handle_python_reqs, uv_pip_compile, PyVersion}, - AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, - PROXY_ENVS, PY_INSTALL_DIR, TZ_ENV, + AuthedClient, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, + PY_INSTALL_DIR, TZ_ENV, }; lazy_static::lazy_static! { @@ -185,7 +185,8 @@ pub async fn handle_ansible_job( mem_peak: &mut i32, canceled_by: &mut Option, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, inner_content: &String, shared_mount: &str, base_internal_url: &str, @@ -263,7 +264,6 @@ pub async fn handle_ansible_job( }) .unwrap_or_else(|| vec![]); - let authed_client = client.get_authed().await; let mut nsjail_extra_mounts = vec![]; if let Some(r) = reqs { nsjail_extra_mounts = create_file_resources( @@ -272,7 +272,7 @@ pub async fn handle_ansible_job( job_dir, interpolated_args.as_ref(), &r, - &authed_client, + &client, db, ) .await?; @@ -311,7 +311,8 @@ remote_tmp={job_dir}/.ansible/tmp ); write_file(job_dir, "ansible.cfg", &ansible_cfg_content)?; - let mut reserved_variables = get_reserved_variables(job, &authed_client.token, db).await?; + let mut reserved_variables = + get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; let additional_python_paths_folders = additional_python_paths.join(":"); if !*DISABLE_NSJAIL { diff --git a/backend/windmill-worker/src/bash_executor.rs b/backend/windmill-worker/src/bash_executor.rs index 412014eb73..a496b327b5 100644 --- a/backend/windmill-worker/src/bash_executor.rs +++ b/backend/windmill-worker/src/bash_executor.rs @@ -46,7 +46,7 @@ use crate::{ OccupancyMetrics, }, handle_child::handle_child, - AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, + AuthedClient, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, POWERSHELL_CACHE_DIR, POWERSHELL_PATH, PROXY_ENVS, TZ_ENV, }; @@ -64,7 +64,8 @@ pub async fn handle_bash_job( canceled_by: &mut Option, job: &MiniPulledJob, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, content: &str, job_dir: &str, shared_mount: &str, @@ -135,8 +136,8 @@ exit $exit_status ); write_file(job_dir, "wrapper.sh", &script)?; - let token = client.get_token().await; - let mut reserved_variables = get_reserved_variables(job, &token, db).await?; + let mut reserved_variables = + get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; reserved_variables.insert("RUST_LOG".to_string(), "info".to_string()); let args = build_args_map(job, client, db).await?.map(Json); @@ -471,7 +472,8 @@ pub async fn handle_powershell_job( canceled_by: &mut Option, job: &MiniPulledJob, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, content: &str, job_dir: &str, shared_mount: &str, @@ -653,8 +655,8 @@ $env:PSModulePath = \"{};$PSModulePathBackup\"", ), )?; - let token = client.get_token().await; - let mut reserved_variables = get_reserved_variables(job, &token, db).await?; + let mut reserved_variables = + get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; reserved_variables.insert("RUST_LOG".to_string(), "info".to_string()); let _ = write_file(job_dir, "result.json", "")?; diff --git a/backend/windmill-worker/src/bigquery_executor.rs b/backend/windmill-worker/src/bigquery_executor.rs index b6367c8529..28aa1398d3 100644 --- a/backend/windmill-worker/src/bigquery_executor.rs +++ b/backend/windmill-worker/src/bigquery_executor.rs @@ -18,7 +18,7 @@ use crate::handle_child::run_future_with_polling_update_job_poller; use crate::sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args; use crate::{ common::{build_args_values, resolve_job_timeout}, - AuthedClientBackgroundTask, + AuthedClient, }; use gcp_auth::{AuthenticationManager, CustomServiceAccount}; @@ -207,7 +207,7 @@ use windmill_queue::MiniPulledJob; pub async fn do_bigquery( job: &MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, query: &str, db: &sqlx::Pool, mem_peak: &mut i32, @@ -223,8 +223,6 @@ pub async fn do_bigquery( let db_arg = if let Some(inline_db_res_path) = inline_db_res_path { Some( client - .get_authed() - .await .get_resource_value_interpolated::( &inline_db_res_path, Some(job.id.to_string()), @@ -267,13 +265,13 @@ pub async fn do_bigquery( .map_err(|x| Error::ExecutionErr(x.to_string()))? .args; - let (query, args_to_skip) = &sanitize_and_interpolate_unsafe_sql_args(query, &sig, &bigquery_args)?; + let (query, args_to_skip) = + &sanitize_and_interpolate_unsafe_sql_args(query, &sig, &bigquery_args)?; let queries = parse_sql_blocks(query); let mut statement_values: HashMap = HashMap::new(); - for arg in &sig { if args_to_skip.contains(&arg.name) { continue; diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 8e13e0d57f..ed6881b4fd 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -20,9 +20,9 @@ use crate::{ read_file_content, read_result, start_child_process, write_file_binary, OccupancyMetrics, }, handle_child::handle_child, - AuthedClientBackgroundTask, BUNFIG_INSTALL_SCOPES, BUN_BUNDLE_CACHE_DIR, BUN_CACHE_DIR, - BUN_PATH, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NODE_BIN_PATH, NODE_PATH, - NPM_CONFIG_REGISTRY, NPM_PATH, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, TZ_ENV, + AuthedClient, BUNFIG_INSTALL_SCOPES, BUN_BUNDLE_CACHE_DIR, BUN_CACHE_DIR, BUN_PATH, + DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NODE_BIN_PATH, NODE_PATH, NPM_CONFIG_REGISTRY, + NPM_PATH, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, TZ_ENV, }; #[cfg(windows)] @@ -825,7 +825,8 @@ pub async fn handle_bun_job( canceled_by: &mut Option, job: &MiniPulledJob, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, job_dir: &str, inner_content: &String, base_internal_url: &str, @@ -934,7 +935,7 @@ pub async fn handle_bun_job( &job.id, &job.workspace_id, Some(db), - &client.get_token().await, + &client.token, job.runnable_path(), job_dir, base_internal_url, @@ -1108,8 +1109,7 @@ try {{ Ok(()) as Result<()> }; let reserved_variables_f = async { - let client = client.get_authed().await; - let vars = get_reserved_variables(job, &client.token, db).await?; + let vars = get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; Ok(vars) as Result> }; let (_, reserved_variables) = tokio::try_join!(args_and_out_f, reserved_variables_f)?; @@ -1127,7 +1127,7 @@ try {{ build_loader( job_dir, base_internal_url, - &client.get_token().await, + &client.token, &job.workspace_id, job.runnable_path(), if annotation.nodejs { @@ -1145,7 +1145,7 @@ try {{ build_loader( job_dir, base_internal_url, - &client.get_token().await, + &client.token, &job.workspace_id, job.runnable_path(), if annotation.nodejs { diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 3b9eadae96..67df5e2f23 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -39,14 +39,11 @@ use windmill_common::{variables, DB}; use tokio::{io::AsyncWriteExt, process::Child, time::Instant}; -use crate::{ - AuthedClient, AuthedClientBackgroundTask, JOB_DEFAULT_TIMEOUT, MAX_RESULT_SIZE, - MAX_TIMEOUT_DURATION, PATH_ENV, -}; +use crate::{AuthedClient, JOB_DEFAULT_TIMEOUT, MAX_RESULT_SIZE, MAX_TIMEOUT_DURATION, PATH_ENV}; pub async fn build_args_map<'a>( job: &'a MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, db: &Pool, ) -> error::Result>>> { if let Some(args) = &job.args { @@ -74,7 +71,7 @@ pub fn check_executor_binary_exists( pub async fn build_args_values( job: &MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, db: &Pool, ) -> error::Result> { if let Some(args) = &job.args { @@ -86,7 +83,7 @@ pub async fn build_args_values( #[tracing::instrument(level = "trace", skip_all)] pub async fn create_args_and_out_file( - client: &AuthedClientBackgroundTask, + client: &AuthedClient, job: &MiniPulledJob, job_dir: &str, db: &Pool, @@ -126,7 +123,7 @@ lazy_static::lazy_static! { } pub async fn transform_json<'a>( - client: &AuthedClientBackgroundTask, + client: &AuthedClient, workspace: &str, vs: &'a HashMap>, job: &MiniPulledJob, @@ -150,9 +147,7 @@ pub async fn transform_json<'a>( let value = serde_json::from_str(inner_vs).map_err(|e| { error::Error::internal_err(format!("Error while parsing inner arg: {e:#}")) })?; - let transformed = - transform_json_value(&k, &client.get_authed().await, workspace, value, job, db) - .await?; + let transformed = transform_json_value(&k, &client, workspace, value, job, db).await?; let as_raw = serde_json::from_value(transformed).map_err(|e| { error::Error::internal_err(format!("Error while parsing inner arg: {e:#}")) })?; @@ -165,7 +160,7 @@ pub async fn transform_json<'a>( } pub async fn transform_json_as_values<'a>( - client: &AuthedClientBackgroundTask, + client: &AuthedClient, workspace: &str, vs: &'a HashMap>, job: &MiniPulledJob, @@ -178,9 +173,7 @@ pub async fn transform_json_as_values<'a>( let value = serde_json::from_str(inner_vs).map_err(|e| { error::Error::internal_err(format!("Error while parsing inner arg: {e:#}")) })?; - let transformed = - transform_json_value(&k, &client.get_authed().await, workspace, value, job, db) - .await?; + let transformed = transform_json_value(&k, &client, workspace, value, job, db).await?; let as_raw = serde_json::from_value(transformed).map_err(|e| { error::Error::internal_err(format!("Error while parsing inner arg: {e:#}")) })?; @@ -283,40 +276,14 @@ pub async fn transform_json_value( // let path = y.strip_prefix("$res:").unwrap(); } Value::String(y) if y.starts_with("$") => { - let flow_path = if let Some(uuid) = job.parent_job { - sqlx::query_scalar!("SELECT runnable_path FROM v2_job WHERE id = $1", uuid) - .fetch_optional(db) - .await? - .flatten() - } else { - None - }; - - let variables = variables::get_reserved_variables( - db, - &job.workspace_id, - &client.token, - &job.permissioned_as_email, - &job.created_by, - &job.id.to_string(), - &job.permissioned_as, - job.runnable_path.clone(), - job.parent_job.map(|x| x.to_string()), - flow_path, - job.schedule_path(), - job.flow_step_id.clone(), - job.flow_innermost_root_job.clone().map(|x| x.to_string()), - None, - Some(job.scheduled_for.clone()), - ) - .await; + let variables = get_reserved_variables(job, &client.token, &db, None).await?; let name = y.strip_prefix("$").unwrap(); let value = variables .iter() - .find(|x| x.name == name) - .map(|x| x.value.clone()) + .find(|x| x.0 == name) + .map(|x| x.1.clone()) .unwrap_or_else(|| y); Ok(json!(value)) } @@ -417,8 +384,11 @@ pub async fn get_reserved_variables( job: &MiniPulledJob, token: &str, db: &sqlx::Pool, + parent_runnable_path: Option, ) -> Result, Error> { - let flow_path = if let Some(uuid) = job.parent_job { + let flow_path = if parent_runnable_path.is_some() { + parent_runnable_path + } else if let Some(uuid) = job.parent_job { sqlx::query_scalar!("SELECT runnable_path FROM v2_job WHERE id = $1", uuid) .fetch_optional(db) .await? diff --git a/backend/windmill-worker/src/csharp_executor.rs b/backend/windmill-worker/src/csharp_executor.rs index 811443d3b8..9acb32b684 100644 --- a/backend/windmill-worker/src/csharp_executor.rs +++ b/backend/windmill-worker/src/csharp_executor.rs @@ -36,7 +36,7 @@ use crate::{ }; use crate::common::OccupancyMetrics; -use crate::AuthedClientBackgroundTask; +use crate::AuthedClient; #[cfg(windows)] use crate::SYSTEM_ROOT; @@ -433,7 +433,8 @@ pub async fn handle_csharp_job( _canceled_by: &mut Option, _job: &MiniPulledJob, _db: &sqlx::Pool, - _client: &AuthedClientBackgroundTask, + _client: &AuthedClient, + _parent_runnable_path: Option, _inner_content: &str, _job_dir: &str, _requirements_o: Option<&String>, @@ -452,7 +453,8 @@ pub async fn handle_csharp_job( canceled_by: &mut Option, job: &MiniPulledJob, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, inner_content: &str, job_dir: &str, requirements_o: Option<&String>, @@ -534,8 +536,8 @@ pub async fn handle_csharp_job( let logs2 = format!("{cache_logs}\n\n--- C# CODE EXECUTION ---\n"); append_logs(&job.id, &job.workspace_id, format!("{}\n", logs2), db).await; - let client = &client.get_authed().await; - let reserved_variables = get_reserved_variables(job, &client.token, db).await?; + let reserved_variables = + get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; let child = if !*DISABLE_NSJAIL { write_file( diff --git a/backend/windmill-worker/src/deno_executor.rs b/backend/windmill-worker/src/deno_executor.rs index 2bd979571c..2b93c86ab4 100644 --- a/backend/windmill-worker/src/deno_executor.rs +++ b/backend/windmill-worker/src/deno_executor.rs @@ -11,8 +11,8 @@ use crate::{ start_child_process, OccupancyMetrics, }, handle_child::handle_child, - AuthedClientBackgroundTask, DENO_CACHE_DIR, DENO_PATH, DISABLE_NSJAIL, HOME_ENV, - NPM_CONFIG_REGISTRY, PATH_ENV, TZ_ENV, + AuthedClient, DENO_CACHE_DIR, DENO_PATH, DISABLE_NSJAIL, HOME_ENV, NPM_CONFIG_REGISTRY, + PATH_ENV, TZ_ENV, }; use tokio::{fs::File, io::AsyncReadExt, process::Command}; use windmill_common::error::{self}; @@ -179,7 +179,8 @@ pub async fn handle_deno_job( canceled_by: &mut Option, job: &MiniPulledJob, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, job_dir: &str, inner_content: &String, base_internal_url: &str, @@ -318,21 +319,21 @@ try {{ Ok(()) as Result<()> }; let reserved_variables_f = async { - let client = client.get_authed().await; - let vars = get_reserved_variables(job, &client.token, db).await?; - Ok((vars, client.token)) as Result<(HashMap, String)> + let vars = get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; + Ok(vars) as Result> }; let (_, reserved_variables) = tokio::try_join!(args_and_out_f, reserved_variables_f)?; - Ok(reserved_variables) as error::Result<(HashMap, String)> + Ok(reserved_variables) as error::Result> }; - let ((reserved_variables, token), _, _) = tokio::try_join!( + let (reserved_variables, _, _) = tokio::try_join!( reserved_variables_args_out_f, write_wrapper_f, write_import_map_f )?; - let mut common_deno_proc_envs = get_common_deno_proc_envs(&token, base_internal_url).await; + let mut common_deno_proc_envs = + get_common_deno_proc_envs(&client.token, base_internal_url).await; if !*DISABLE_NSJAIL { common_deno_proc_envs.insert("HOME".to_string(), job_dir.to_string()); } diff --git a/backend/windmill-worker/src/go_executor.rs b/backend/windmill-worker/src/go_executor.rs index 23aee36fa4..43e88c59aa 100644 --- a/backend/windmill-worker/src/go_executor.rs +++ b/backend/windmill-worker/src/go_executor.rs @@ -19,8 +19,8 @@ use crate::{ start_child_process, OccupancyMetrics, }, handle_child::handle_child, - AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, GOPRIVATE, GOPROXY, - GO_BIN_CACHE_DIR, GO_CACHE_DIR, HOME_ENV, NSJAIL_PATH, PATH_ENV, TZ_ENV, + AuthedClient, DISABLE_NSJAIL, DISABLE_NUSER, GOPRIVATE, GOPROXY, GO_BIN_CACHE_DIR, + GO_CACHE_DIR, HOME_ENV, NSJAIL_PATH, PATH_ENV, TZ_ENV, }; const GO_REQ_SPLITTER: &str = "//go.sum\n"; @@ -37,7 +37,8 @@ pub async fn handle_go_job( canceled_by: &mut Option, job: &MiniPulledJob, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, inner_content: &str, job_dir: &str, requirements_o: Option<&String>, @@ -248,9 +249,8 @@ func Run(req Req) (interface{{}}, error){{ let logs2 = format!("{cache_logs}\n\n--- GO CODE EXECUTION ---\n"); append_logs(&job.id, &job.workspace_id, logs2, db).await; - let client = &client.get_authed().await; - - let reserved_variables = get_reserved_variables(job, &client.token, db).await?; + let reserved_variables = + get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; let child = if !*DISABLE_NSJAIL { let _ = write_file( diff --git a/backend/windmill-worker/src/graphql_executor.rs b/backend/windmill-worker/src/graphql_executor.rs index 98e5ef9a32..c755d33beb 100644 --- a/backend/windmill-worker/src/graphql_executor.rs +++ b/backend/windmill-worker/src/graphql_executor.rs @@ -13,7 +13,7 @@ use serde::Deserialize; use crate::common::{build_http_client, resolve_job_timeout, OccupancyMetrics}; use crate::handle_child::run_future_with_polling_update_job_poller; -use crate::{common::build_args_map, AuthedClientBackgroundTask}; +use crate::{common::build_args_map, AuthedClient}; #[derive(Deserialize)] struct GraphqlApi { @@ -35,7 +35,7 @@ struct GraphqlError { pub async fn do_graphql( job: &MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, query: &str, db: &sqlx::Pool, mem_peak: &mut i32, diff --git a/backend/windmill-worker/src/mssql_executor.rs b/backend/windmill-worker/src/mssql_executor.rs index ef06c8f0bf..659686268e 100644 --- a/backend/windmill-worker/src/mssql_executor.rs +++ b/backend/windmill-worker/src/mssql_executor.rs @@ -18,7 +18,7 @@ use windmill_queue::{append_logs, CanceledBy}; use crate::common::{build_args_values, OccupancyMetrics}; use crate::handle_child::run_future_with_polling_update_job_poller; use crate::sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args; -use crate::AuthedClientBackgroundTask; +use crate::AuthedClient; #[derive(Deserialize)] struct MssqlDatabase { @@ -36,7 +36,7 @@ lazy_static::lazy_static! { pub async fn do_mssql( job: &MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, query: &str, db: &sqlx::Pool, mem_peak: &mut i32, @@ -51,8 +51,6 @@ pub async fn do_mssql( let db_arg = if let Some(inline_db_res_path) = inline_db_res_path { Some( client - .get_authed() - .await .get_resource_value_interpolated::( &inline_db_res_path, Some(job.id.to_string()), @@ -133,7 +131,8 @@ pub async fn do_mssql( .map_err(|x| Error::ExecutionErr(x.to_string()))? .args; - let (query, args_to_skip) = &sanitize_and_interpolate_unsafe_sql_args(query, &sig, &mssql_args)?; + let (query, args_to_skip) = + &sanitize_and_interpolate_unsafe_sql_args(query, &sig, &mssql_args)?; let mut prepared_query = Query::new(query.to_owned()); for arg in &sig { diff --git a/backend/windmill-worker/src/mysql_executor.rs b/backend/windmill-worker/src/mysql_executor.rs index ca8fa1a5c9..2cab8ce8b2 100644 --- a/backend/windmill-worker/src/mysql_executor.rs +++ b/backend/windmill-worker/src/mysql_executor.rs @@ -21,7 +21,10 @@ use windmill_queue::CanceledBy; use windmill_queue::MiniPulledJob; use crate::{ - common::{build_args_values, OccupancyMetrics}, handle_child::run_future_with_polling_update_job_poller, sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args, AuthedClientBackgroundTask + common::{build_args_values, OccupancyMetrics}, + handle_child::run_future_with_polling_update_job_poller, + sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args, + AuthedClient, }; #[derive(Deserialize)] @@ -101,7 +104,7 @@ pub fn do_mysql_inner<'a>( pub async fn do_mysql( job: &MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, query: &str, db: &sqlx::Pool, mem_peak: &mut i32, @@ -115,14 +118,14 @@ pub async fn do_mysql( let inline_db_res_path = parse_db_resource(&query); let db_arg = if let Some(inline_db_res_path) = inline_db_res_path { - Some(client - .get_authed() - .await - .get_resource_value_interpolated::( - &inline_db_res_path, - Some(job.id.to_string()), - ) - .await?) + Some( + client + .get_resource_value_interpolated::( + &inline_db_res_path, + Some(job.id.to_string()), + ) + .await?, + ) } else { job_args.get("database").cloned() }; @@ -171,7 +174,8 @@ pub async fn do_mysql( } let arg_t = arg.otyp.clone().unwrap_or_else(|| "text".to_string()); let arg_n = arg.name.clone(); - let mysql_v = match job_args.get(arg.name.as_str()) + let mysql_v = match job_args + .get(arg.name.as_str()) .unwrap_or_else(|| &json!(null)) { Value::Null => mysql_async::Value::NULL, diff --git a/backend/windmill-worker/src/nu_executor.rs b/backend/windmill-worker/src/nu_executor.rs index 2859e016b8..e40e1029f4 100644 --- a/backend/windmill-worker/src/nu_executor.rs +++ b/backend/windmill-worker/src/nu_executor.rs @@ -13,8 +13,7 @@ use crate::{ create_args_and_out_file, get_reserved_variables, read_result, start_child_process, OccupancyMetrics, }, - handle_child, AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, PATH_ENV, - PROXY_ENVS, + handle_child, AuthedClient, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, }; const NSJAIL_CONFIG_RUN_NU_CONTENT: &str = include_str!("../nsjail/run.nu.config.proto"); @@ -29,7 +28,8 @@ lazy_static::lazy_static! { pub(crate) struct JobHandlerInput<'a> { pub base_internal_url: &'a str, pub canceled_by: &'a mut Option, - pub client: &'a AuthedClientBackgroundTask, + pub client: &'a AuthedClient, + pub parent_runnable_path: Option, pub db: &'a sqlx::Pool, pub envs: HashMap, pub inner_content: &'a str, @@ -221,14 +221,15 @@ async fn run<'a>( job_dir, shared_mount, client, + parent_runnable_path, envs, base_internal_url, .. }: &mut JobHandlerInput<'a>, // plugins: Vec<&'a str>, ) -> Result<(), Error> { - let client = &client.get_authed().await; - let reserved_variables = get_reserved_variables(job, &client.token, db).await?; + let reserved_variables = + get_reserved_variables(job, &client.token, db, parent_runnable_path.clone()).await?; let child = if !cfg!(windows) && !*DISABLE_NSJAIL { append_logs( &job.id, diff --git a/backend/windmill-worker/src/oracledb_executor.rs b/backend/windmill-worker/src/oracledb_executor.rs index 7f340b8525..2020e421df 100644 --- a/backend/windmill-worker/src/oracledb_executor.rs +++ b/backend/windmill-worker/src/oracledb_executor.rs @@ -23,7 +23,7 @@ use crate::{ common::{build_args_values, check_executor_binary_exists, OccupancyMetrics}, handle_child::run_future_with_polling_update_job_poller, sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args, - AuthedClientBackgroundTask, + AuthedClient, }; #[derive(Deserialize)] @@ -233,7 +233,8 @@ fn get_statement_values( } let arg_t = arg.otyp.clone().unwrap_or_else(|| "text".to_string()); let arg_n = arg.name.clone(); - let oracle_v: Box = match job_args.get(arg.name.as_str()) + let oracle_v: Box = match job_args + .get(arg.name.as_str()) .unwrap_or_else(|| &json!(null)) { // Value::Null => todo!(), @@ -293,7 +294,7 @@ fn get_statement_values( pub async fn do_oracledb( job: &MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, query: &str, db: &sqlx::Pool, mem_peak: &mut i32, @@ -315,8 +316,6 @@ pub async fn do_oracledb( let db_arg = if let Some(inline_db_res_path) = inline_db_res_path { Some( client - .get_authed() - .await .get_resource_value_interpolated::( &inline_db_res_path, Some(job.id.to_string()), diff --git a/backend/windmill-worker/src/pg_executor.rs b/backend/windmill-worker/src/pg_executor.rs index 2847389ca7..e9aa27eca9 100644 --- a/backend/windmill-worker/src/pg_executor.rs +++ b/backend/windmill-worker/src/pg_executor.rs @@ -37,7 +37,7 @@ use windmill_queue::{CanceledBy, MiniPulledJob}; use crate::common::{build_args_values, sizeof_val, OccupancyMetrics}; use crate::handle_child::run_future_with_polling_update_job_poller; use crate::sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args; -use crate::{AuthedClientBackgroundTask, MAX_RESULT_SIZE}; +use crate::{AuthedClient, MAX_RESULT_SIZE}; use bytes::Buf; use lazy_static::lazy_static; use urlencoding::encode; @@ -159,7 +159,7 @@ fn do_postgresql_inner<'a>( pub async fn do_postgresql( job: &MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, query: &str, db: &sqlx::Pool, mem_peak: &mut i32, @@ -175,8 +175,6 @@ pub async fn do_postgresql( let db_arg = if let Some(inline_db_res_path) = inline_db_res_path { Some( client - .get_authed() - .await .get_resource_value_interpolated::( &inline_db_res_path, Some(job.id.to_string()), @@ -293,8 +291,7 @@ pub async fn do_postgresql( let sig = parse_pgsql_sig(&query).map_err(|x| Error::ExecutionErr(x.to_string()))?; - let (query, _) = - &sanitize_and_interpolate_unsafe_sql_args(query, &sig.args, &pg_args)?; + let (query, _) = &sanitize_and_interpolate_unsafe_sql_args(query, &sig.args, &pg_args)?; let queries = parse_sql_blocks(query); diff --git a/backend/windmill-worker/src/php_executor.rs b/backend/windmill-worker/src/php_executor.rs index fa5eb58c16..19dce7f1b9 100644 --- a/backend/windmill-worker/src/php_executor.rs +++ b/backend/windmill-worker/src/php_executor.rs @@ -20,8 +20,8 @@ use crate::{ read_result, start_child_process, OccupancyMetrics, }, handle_child::handle_child, - AuthedClientBackgroundTask, COMPOSER_CACHE_DIR, COMPOSER_PATH, DISABLE_NSJAIL, DISABLE_NUSER, - NSJAIL_PATH, PHP_PATH, + AuthedClient, COMPOSER_CACHE_DIR, COMPOSER_PATH, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, + PHP_PATH, }; const NSJAIL_CONFIG_RUN_PHP_CONTENT: &str = include_str!("../nsjail/run.php.config.proto"); @@ -139,7 +139,8 @@ pub async fn handle_php_job( canceled_by: &mut Option, job: &MiniPulledJob, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, job_dir: &str, inner_content: &String, base_internal_url: &str, @@ -265,8 +266,8 @@ try {{ Ok(()) as Result<()> }; let reserved_variables_f = async { - let client = client.get_authed().await; - let vars = get_reserved_variables(job, &client.token, db).await?; + let vars = get_reserved_variables(job, &client.token, db, parent_runnable_path.clone()) + .await?; Ok(vars) as Result> }; let (_, reserved_variables) = tokio::try_join!(args_and_out_f, reserved_variables_f)?; diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index d39580d285..648c40a9a3 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -76,9 +76,8 @@ use crate::{ start_child_process, OccupancyMetrics, }, handle_child::handle_child, - AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, INSTANCE_PYTHON_VERSION, - NSJAIL_PATH, PATH_ENV, PIP_EXTRA_INDEX_URL, PIP_INDEX_URL, PROXY_ENVS, PY_INSTALL_DIR, TZ_ENV, - UV_CACHE_DIR, + AuthedClient, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, INSTANCE_PYTHON_VERSION, NSJAIL_PATH, + PATH_ENV, PIP_EXTRA_INDEX_URL, PIP_INDEX_URL, PROXY_ENVS, PY_INSTALL_DIR, TZ_ENV, UV_CACHE_DIR, }; // To change latest stable version: @@ -839,7 +838,8 @@ pub async fn handle_python_job( mem_peak: &mut i32, canceled_by: &mut Option, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, inner_content: &String, shared_mount: &str, base_internal_url: &str, @@ -1024,8 +1024,8 @@ except BaseException as e: tracing::debug!("Finished writing wrapper"); - let client = client.get_authed().await; - let mut reserved_variables = get_reserved_variables(job, &client.token, db).await?; + let mut reserved_variables = + get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; // Add /tmp/windmill/cache/python_xyz/global-site-packages to PYTHONPATH. // Usefull if certain wheels needs to be preinstalled before execution. diff --git a/backend/windmill-worker/src/result_processor.rs b/backend/windmill-worker/src/result_processor.rs index 966991d0ea..4c2e9fa7a0 100644 --- a/backend/windmill-worker/src/result_processor.rs +++ b/backend/windmill-worker/src/result_processor.rs @@ -235,7 +235,7 @@ async fn send_job_completed( canceled_by: Option, success: bool, cached_res_path: Option, - token: String, + token: &str, duration: Option, ) { let jc = JobCompleted { @@ -246,7 +246,7 @@ async fn send_job_completed( canceled_by, success, cached_res_path, - token, + token: token.to_string(), duration, }; job_completed_tx @@ -264,7 +264,7 @@ pub async fn process_result( mem_peak: i32, canceled_by: Option, cached_res_path: Option, - token: String, + token: &str, column_order: Option>, new_args: Option>>, db: &DB, diff --git a/backend/windmill-worker/src/rust_executor.rs b/backend/windmill-worker/src/rust_executor.rs index b3fbb08d35..d2a0c3217d 100644 --- a/backend/windmill-worker/src/rust_executor.rs +++ b/backend/windmill-worker/src/rust_executor.rs @@ -19,8 +19,8 @@ use crate::{ read_result, start_child_process, OccupancyMetrics, }, handle_child::handle_child, - AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, - PROXY_ENVS, RUST_CACHE_DIR, TZ_ENV, + AuthedClient, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, + RUST_CACHE_DIR, TZ_ENV, }; #[cfg(windows)] @@ -277,7 +277,8 @@ pub async fn handle_rust_job( canceled_by: &mut Option, job: &MiniPulledJob, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, inner_content: &str, job_dir: &str, requirements_o: Option<&String>, @@ -343,8 +344,8 @@ pub async fn handle_rust_job( let logs2 = format!("{cache_logs}\n\n--- RUST CODE EXECUTION ---\n"); append_logs(&job.id, &job.workspace_id, logs2, db).await; - let client = &client.get_authed().await; - let reserved_variables = get_reserved_variables(job, &client.token, db).await?; + let reserved_variables = + get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; let child = if !*DISABLE_NSJAIL { let _ = write_file( diff --git a/backend/windmill-worker/src/sanitized_sql_params.rs b/backend/windmill-worker/src/sanitized_sql_params.rs index c7eb3b6ef4..465b67658f 100644 --- a/backend/windmill-worker/src/sanitized_sql_params.rs +++ b/backend/windmill-worker/src/sanitized_sql_params.rs @@ -10,7 +10,10 @@ use windmill_parser_sql::{SANITIZED_ENUM_STR, SANITIZED_RAW_STRING_STR}; /// a number, that can contain underscores fn sanitize_identifier(arg: &Arg, input: &str) -> Result<(), error::Error> { if input.is_empty() { - return Err(error::Error::BadRequest(format!("Interpolated argument `{}` cannot be empty", arg.name))); + return Err(error::Error::BadRequest(format!( + "Interpolated argument `{}` cannot be empty", + arg.name + ))); } if input .chars() diff --git a/backend/windmill-worker/src/snowflake_executor.rs b/backend/windmill-worker/src/snowflake_executor.rs index 30f1188ce0..377c553d56 100644 --- a/backend/windmill-worker/src/snowflake_executor.rs +++ b/backend/windmill-worker/src/snowflake_executor.rs @@ -19,7 +19,7 @@ use serde::{Deserialize, Serialize}; use crate::common::{build_http_client, resolve_job_timeout, OccupancyMetrics}; use crate::handle_child::run_future_with_polling_update_job_poller; use crate::sanitized_sql_params::sanitize_and_interpolate_unsafe_sql_args; -use crate::{common::build_args_values, AuthedClientBackgroundTask}; +use crate::{common::build_args_values, AuthedClient}; #[derive(Serialize)] struct Claims { @@ -247,7 +247,7 @@ fn do_snowflake_inner<'a>( pub async fn do_snowflake( job: &MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, query: &str, db: &sqlx::Pool, mem_peak: &mut i32, @@ -263,8 +263,6 @@ pub async fn do_snowflake( let db_arg = if let Some(inline_db_res_path) = inline_db_res_path { Some( client - .get_authed() - .await .get_resource_value_interpolated::( &inline_db_res_path, Some(job.id.to_string()), diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 6f7bd1f16a..c696dbe6c6 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -160,20 +160,10 @@ use crate::bench::{benchmark_init, BenchmarkInfo, BenchmarkIter}; use windmill_common::add_time; -pub async fn create_token_for_owner_in_bg( - db: &Pool, - job: &MiniPulledJob, -) -> Arc> { - let rw_lock = Arc::new(RwLock::new(String::new())); +// struct Permission +pub async fn create_token(db: &DB, job: &MiniPulledJob, perms: Option) -> String { // skipping test runs if job.workspace_id != "" { - let mut locked = rw_lock.clone().write_owned().await; - let db = db.clone(); - let w_id = job.workspace_id.clone(); - let owner = job.permissioned_as.clone(); - let email = job.permissioned_as_email.clone(); - let job_id = job.id.clone(); - let label = if job.permissioned_as != format!("u/{}", job.created_by) && job.permissioned_as != job.created_by { @@ -181,23 +171,22 @@ pub async fn create_token_for_owner_in_bg( } else { "ephemeral-script".to_string() }; - tokio::spawn(async move { - let token = create_token_for_owner( - &db.clone(), - &w_id, - &owner, - &label, - *SCRIPT_TOKEN_EXPIRY, - &email, - &job_id, - ) - .warn_after_seconds(5) - .await - .expect("could not create job token"); - *locked = token; - }); - }; - return rw_lock; + create_token_for_owner( + &db, + &job.workspace_id, + &job.permissioned_as, + &label, + *SCRIPT_TOKEN_EXPIRY, + &job.permissioned_as_email, + &job.id, + perms, + ) + .warn_after_seconds(5) + .await + .expect("could not create job token") + } else { + return "".to_string(); + } } #[tracing::instrument(level = "trace", skip_all)] @@ -209,21 +198,27 @@ pub async fn create_token_for_owner( expires_in: u64, email: &str, job_id: &Uuid, + perms: Option, ) -> error::Result { // TODO: Bad implementation. We should not have access to this DB here. if let Some(token) = JOB_TOKEN.as_ref() { return Ok(token.clone()); } - let job_authed = match sqlx::query_as!( - JobPerms, - "SELECT * FROM job_perms WHERE job_id = $1 AND workspace_id = $2", - job_id, - w_id - ) + let job_perms = if perms.is_some() { + Ok(perms) + } else { + sqlx::query_as!( + JobPerms, + "SELECT email, username, is_admin, is_operator, groups, folders FROM job_perms WHERE job_id = $1 AND workspace_id = $2", + job_id, + w_id + ) .fetch_optional(db) .await - { + }; + + let job_authed = match job_perms { Ok(Some(jp)) => jp.into(), _ => { tracing::warn!("Could not get permissions for job {job_id} from job_perms table, getting permissions directly..."); @@ -460,25 +455,6 @@ pub const MAX_RESULT_SIZE: usize = 1024 * 1024 * 2; // 2MB pub const INIT_SCRIPT_TAG: &str = "init_script"; -pub struct AuthedClientBackgroundTask { - pub base_internal_url: String, - pub workspace: String, - pub token: Arc>, -} - -impl AuthedClientBackgroundTask { - pub async fn get_authed(&self) -> AuthedClient { - return AuthedClient { - base_internal_url: self.base_internal_url.clone(), - workspace: self.workspace.clone(), - token: self.get_token().await, - force_client: None, - }; - } - pub async fn get_token(&self) -> String { - return self.token.read().await.clone(); - } -} #[derive(Clone)] pub struct AuthedClient { pub base_internal_url: String, @@ -1322,35 +1298,43 @@ pub async fn run_worker( v2_job.created_by, v2_job_queue.started_at, scheduled_for, - runnable_path, - kind, - runnable_id, - canceled_reason, - canceled_by, - permissioned_as, - permissioned_as_email, - flow_status, + v2_job.runnable_path, + v2_job.kind, + v2_job.runnable_id, + v2_job_queue.canceled_reason, + v2_job_queue.canceled_by, + v2_job.permissioned_as, + v2_job.permissioned_as_email, + v2_job_status.flow_status, v2_job.tag, - script_lang, - same_worker, - pre_run_error, - concurrent_limit, - concurrency_time_window_s, - flow_innermost_root_job, - timeout, - flow_step_id, - cache_ttl, + v2_job.script_lang, + v2_job.same_worker, + v2_job.pre_run_error, + v2_job.concurrent_limit, + v2_job.concurrency_time_window_s, + v2_job.flow_innermost_root_job, + v2_job.timeout, + v2_job.flow_step_id, + v2_job.cache_ttl, v2_job_queue.priority, - preprocessed, - script_entrypoint_override, - trigger, - trigger_kind, - visible_to_owner, - raw_code, - raw_lock, - raw_flow - FROM v2_job_queue INNER JOIN v2_job ON v2_job.id = v2_job_queue.id LEFT JOIN v2_job_status ON v2_job_status.id = v2_job_queue.id WHERE v2_job_queue.id = $1 - ", + v2_job.preprocessed, + v2_job.script_entrypoint_override, + v2_job.trigger, + v2_job.trigger_kind, + v2_job.visible_to_owner, + v2_job.raw_code, + v2_job.raw_lock, + v2_job.raw_flow, + pj.runnable_path as parent_runnable_path, + p.email as permissioned_as_email, p.username as permissioned_as_username, p.is_admin as permissioned_as_is_admin, + p.is_operator as permissioned_as_is_operator, p.groups as permissioned_as_groups, p.folders as permissioned_as_folders + FROM v2_job_queue + INNER JOIN v2_job ON v2_job.id = v2_job_queue.id + LEFT JOIN v2_job_status ON v2_job_status.id = v2_job_queue.id + LEFT JOIN job_perms p ON p.job_id = v2_job.id + LEFT JOIN v2_job pj ON v2_job.parent_job = pj.id + WHERE v2_job_queue.id = $1 +", ) .bind(same_worker_job.job_id) .fetch_optional(db) @@ -1526,7 +1510,6 @@ pub async fn run_worker( .expect("send job completed END"); add_time!(bench, "sent job completed"); } else { - let token = create_token_for_owner_in_bg(&db, &job).await; add_outstanding_wait_time(&job, db, OUTSTANDING_WAIT_TIME_THRESHOLD_MS); #[cfg(feature = "prometheus")] @@ -1627,17 +1610,57 @@ pub async fn run_worker( .expect("could not create shared dir"); } - let authed_client = AuthedClientBackgroundTask { - base_internal_url: base_internal_url.to_string(), - token, - workspace: job.workspace_id.to_string(), - }; - #[cfg(feature = "prometheus")] let tag = job.tag.clone(); let is_init_script: bool = job.tag.as_str() == INIT_SCRIPT_TAG; - let PulledJob { job, raw_code, raw_lock, raw_flow } = job; + let PulledJob { + job, + raw_code, + raw_lock, + raw_flow, + parent_runnable_path, + permissioned_as_email, + permissioned_as_username, + permissioned_as_is_admin, + permissioned_as_is_operator, + permissioned_as_groups, + permissioned_as_folders, + } = job; + let job_perms = match ( + permissioned_as_email, + permissioned_as_username, + permissioned_as_is_admin, + permissioned_as_is_operator, + permissioned_as_groups, + permissioned_as_folders, + ) { + ( + Some(email), + Some(username), + Some(is_admin), + Some(is_operator), + Some(groups), + Some(folders), + ) => Some(JobPerms { + email, + username, + is_admin, + is_operator, + groups, + folders, + }), + _ => None, + }; + + let token = create_token(&db, &job, job_perms).await; + let authed_client = AuthedClient { + base_internal_url: base_internal_url.to_string(), + token, + workspace: job.workspace_id.to_string(), + force_client: None, + }; + let arc_job = Arc::new(job); add_time!(bench, "handle_queued_job START"); @@ -1678,6 +1701,7 @@ pub async fn run_worker( raw_code, raw_lock, raw_flow, + parent_runnable_path, db, &authed_client, &hostname, @@ -1698,7 +1722,7 @@ pub async fn run_worker( Err(err) => { handle_job_error( db, - &authed_client.get_authed().await, + &authed_client, arc_job.as_ref(), 0, None, @@ -1923,7 +1947,7 @@ pub struct JobCompleted { async fn do_nativets( job: &MiniPulledJob, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, env_code: String, code: String, db: &Pool, @@ -1968,8 +1992,9 @@ async fn handle_queued_job( raw_code: Option, raw_lock: Option, raw_flow: Option>>, + parent_runnable_path: Option, db: &DB, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, hostname: &str, worker_name: &str, worker_dir: &str, @@ -2081,17 +2106,15 @@ async fn handle_queued_job( }; let cached_res_path = if job.cache_ttl.is_some() { - Some(cached_result_path(db, &client.get_authed().await, &job, preview_data.as_ref()).await) + Some(cached_result_path(db, &client, &job, preview_data.as_ref()).await) } else { None }; if let Some(cached_res_path) = cached_res_path.as_ref() { - let authed_client = client.get_authed().await; - let cached_result_maybe = get_cached_resource_value_if_valid( db, - &authed_client, + &client, &job.id, &job.workspace_id, &cached_res_path, @@ -2113,7 +2136,7 @@ async fn handle_queued_job( canceled_by: None, success: true, cached_res_path: None, - token: authed_client.token, + token: client.token.clone(), duration: None, }) .await @@ -2132,7 +2155,7 @@ async fn handle_queued_job( job, &flow_data, db, - &client.get_authed().await, + &client, None, same_worker_tx, worker_dir, @@ -2193,7 +2216,7 @@ async fn handle_queued_job( worker_name, worker_dir, base_internal_url, - &client.get_token().await, + &client.token, occupancy_metrics, ) .await @@ -2209,7 +2232,7 @@ async fn handle_queued_job( worker_name, worker_dir, base_internal_url, - &client.get_token().await, + &client.token, occupancy_metrics, ) .await @@ -2223,7 +2246,7 @@ async fn handle_queued_job( worker_name, worker_dir, base_internal_url, - &client.get_token().await, + &client.token, occupancy_metrics, ) .await @@ -2246,6 +2269,7 @@ async fn handle_queued_job( preview_data, db, client, + parent_runnable_path, job_dir, worker_dir, &mut mem_peak, @@ -2284,7 +2308,7 @@ async fn handle_queued_job( mem_peak, canceled_by, cached_res_path, - client.get_token().await, + &client.token, column_order, new_args, db, @@ -2421,9 +2445,7 @@ async fn try_validate_schema( }; let sv = match job.runnable_id { - Some(hash) - if job.kind != JobKind::Preview && job.kind != JobKind::FlowPreview => - { + Some(hash) if job.kind != JobKind::Preview && job.kind != JobKind::FlowPreview => { sv_fut.cached(validators_cache, (sub_key, hash)).await? } _ => sv_fut.await?, @@ -2464,7 +2486,8 @@ async fn handle_code_execution_job( job: &MiniPulledJob, preview: Option>, db: &sqlx::Pool, - client: &AuthedClientBackgroundTask, + client: &AuthedClient, + parent_runnable_path: Option, job_dir: &str, #[allow(unused_variables)] worker_dir: &str, mem_peak: &mut i32, @@ -2754,7 +2777,8 @@ async fn handle_code_execution_job( ) .await; - let reserved_variables = get_reserved_variables(job, &client.get_token().await, db).await?; + let reserved_variables = + get_reserved_variables(job, &client.token, db, parent_runnable_path).await?; let env_code = format!( "const process = {{ env: {{}} }};\nconst BASE_URL = '{base_internal_url}';\nconst BASE_INTERNAL_URL = '{base_internal_url}';\nprocess.env['BASE_URL'] = BASE_URL;process.env['BASE_INTERNAL_URL'] = BASE_INTERNAL_URL;\n{}", @@ -2839,6 +2863,7 @@ mount {{ canceled_by, db, client, + parent_runnable_path, &code, &shared_mount, base_internal_url, @@ -2856,6 +2881,7 @@ mount {{ job, db, client, + parent_runnable_path, job_dir, &code, base_internal_url, @@ -2875,6 +2901,7 @@ mount {{ job, db, client, + parent_runnable_path, job_dir, &code, base_internal_url, @@ -2893,6 +2920,7 @@ mount {{ job, db, client, + parent_runnable_path, &code, job_dir, lock.as_ref(), @@ -2911,6 +2939,7 @@ mount {{ job, db, client, + parent_runnable_path, &code, job_dir, &shared_mount, @@ -2929,6 +2958,7 @@ mount {{ job, db, client, + parent_runnable_path, &code, job_dir, &shared_mount, @@ -2953,6 +2983,7 @@ mount {{ job, db, client, + parent_runnable_path, job_dir, &code, base_internal_url, @@ -2976,6 +3007,7 @@ mount {{ job, db, client, + parent_runnable_path, &code, job_dir, lock.as_ref(), @@ -3004,6 +3036,7 @@ mount {{ canceled_by, db, client, + parent_runnable_path, &code, &shared_mount, base_internal_url, @@ -3019,6 +3052,7 @@ mount {{ job, db, client, + parent_runnable_path, &code, job_dir, lock.as_ref(), @@ -3043,6 +3077,7 @@ mount {{ job, db, client, + parent_runnable_path, inner_content: &code, job_dir, requirements_o: lock.as_ref(), diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 0b1b049f4c..d8c9a0a410 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -2651,7 +2651,7 @@ async fn push_next_flow_job( { sqlx::query_as!( JobPerms, - "SELECT * FROM job_perms WHERE job_id = $1 AND workspace_id = $2", + "SELECT email, username, is_admin, is_operator, groups, folders FROM job_perms WHERE job_id = $1 AND workspace_id = $2", root_job, flow_job.workspace_id, ) @@ -2940,15 +2940,17 @@ async fn push_next_flow_job( .execute(&mut *tx) .await?; - tx.commit().warn_after_seconds(3).await?; - tracing::info!(id = %flow_job.id, root_id = %job_root, "all next flow jobs pushed: {uuids:?}"); - if continue_on_same_worker { if !is_one_uuid { return Err(Error::BadRequest( "Cannot continue on same worker with multiple jobs, parallel cannot be used in conjunction with same_worker".to_string(), )); } + } + tx.commit().warn_after_seconds(3).await?; + tracing::info!(id = %flow_job.id, root_id = %job_root, "all next flow jobs pushed: {uuids:?}"); + + if continue_on_same_worker { same_worker_tx .send(SameWorkerPayload { job_id: first_uuid, recoverable: true }) .await