diff --git a/backend/.sqlx/query-123c0608e229c29187009b7961355ddd99c4ad1f46b876dd86e372b84d806ecd.json b/backend/.sqlx/query-123c0608e229c29187009b7961355ddd99c4ad1f46b876dd86e372b84d806ecd.json index 1b8084742c..7718e05ccf 100644 --- a/backend/.sqlx/query-123c0608e229c29187009b7961355ddd99c4ad1f46b876dd86e372b84d806ecd.json +++ b/backend/.sqlx/query-123c0608e229c29187009b7961355ddd99c4ad1f46b876dd86e372b84d806ecd.json @@ -37,7 +37,6 @@ "bash", "postgresql", "nativets", - "Nativets", "bun", "mysql", "bigquery", diff --git a/backend/.sqlx/query-25bef6a248f3ee0ea2cbcc376c217cbcf1013ae311c36b42d423bf6a02fa016c.json b/backend/.sqlx/query-25bef6a248f3ee0ea2cbcc376c217cbcf1013ae311c36b42d423bf6a02fa016c.json index bf591ef11c..9d082a6772 100644 --- a/backend/.sqlx/query-25bef6a248f3ee0ea2cbcc376c217cbcf1013ae311c36b42d423bf6a02fa016c.json +++ b/backend/.sqlx/query-25bef6a248f3ee0ea2cbcc376c217cbcf1013ae311c36b42d423bf6a02fa016c.json @@ -67,7 +67,6 @@ "bash", "postgresql", "nativets", - "Nativets", "bun", "mysql", "bigquery", diff --git a/backend/.sqlx/query-2de52e1f3226ca9281b6e25f74d4d05f4509cb87c875234bfc7b310a012e4d40.json b/backend/.sqlx/query-2de52e1f3226ca9281b6e25f74d4d05f4509cb87c875234bfc7b310a012e4d40.json new file mode 100644 index 0000000000..0174706d37 --- /dev/null +++ b/backend/.sqlx/query-2de52e1f3226ca9281b6e25f74d4d05f4509cb87c875234bfc7b310a012e4d40.json @@ -0,0 +1,73 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO queue (id, script_hash, script_path, job_kind, language, tag, created_by, permissioned_as, email, scheduled_for, workspace_id) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Int8", + "Varchar", + { + "Custom": { + "name": "job_kind", + "kind": { + "Enum": [ + "script", + "preview", + "flow", + "dependencies", + "flowpreview", + "script_hub", + "identity", + "flowdependencies", + "http", + "graphql", + "postgresql", + "noop", + "appdependencies" + ] + } + } + }, + { + "Custom": { + "name": "script_lang", + "kind": { + "Enum": [ + "python3", + "deno", + "go", + "bash", + "postgresql", + "nativets", + "bun", + "mysql", + "bigquery", + "snowflake", + "graphql", + "powershell" + ] + } + } + }, + "Varchar", + "Varchar", + "Varchar", + "Varchar", + "Timestamptz", + "Varchar", + "Int4" + ] + }, + "nullable": [ + false + ] + }, + "hash": "2de52e1f3226ca9281b6e25f74d4d05f4509cb87c875234bfc7b310a012e4d40" +} diff --git a/backend/.sqlx/query-438b5b5d29b05846c2e074cad2404e797527841cd97dba80c271cbefafae65cc.json b/backend/.sqlx/query-438b5b5d29b05846c2e074cad2404e797527841cd97dba80c271cbefafae65cc.json index a5dee163e5..1166260449 100644 --- a/backend/.sqlx/query-438b5b5d29b05846c2e074cad2404e797527841cd97dba80c271cbefafae65cc.json +++ b/backend/.sqlx/query-438b5b5d29b05846c2e074cad2404e797527841cd97dba80c271cbefafae65cc.json @@ -28,7 +28,6 @@ "bash", "postgresql", "nativets", - "Nativets", "bun", "mysql", "bigquery", diff --git a/backend/.sqlx/query-5cd89ab614d3cac80fb81627267ee85b191263989b0c78b2bfce77e796e96825.json b/backend/.sqlx/query-5cd89ab614d3cac80fb81627267ee85b191263989b0c78b2bfce77e796e96825.json index 1517c8d1d4..eabf671894 100644 --- a/backend/.sqlx/query-5cd89ab614d3cac80fb81627267ee85b191263989b0c78b2bfce77e796e96825.json +++ b/backend/.sqlx/query-5cd89ab614d3cac80fb81627267ee85b191263989b0c78b2bfce77e796e96825.json @@ -42,7 +42,6 @@ "bash", "postgresql", "nativets", - "Nativets", "bun", "mysql", "bigquery", diff --git a/backend/.sqlx/query-65835f2e5ad38f7cc6b147dadfef6f580f15bca96d9746c9359e98ca793f8f1f.json b/backend/.sqlx/query-65835f2e5ad38f7cc6b147dadfef6f580f15bca96d9746c9359e98ca793f8f1f.json index bfe7c41f64..c52efca4c0 100644 --- a/backend/.sqlx/query-65835f2e5ad38f7cc6b147dadfef6f580f15bca96d9746c9359e98ca793f8f1f.json +++ b/backend/.sqlx/query-65835f2e5ad38f7cc6b147dadfef6f580f15bca96d9746c9359e98ca793f8f1f.json @@ -42,7 +42,6 @@ "bash", "postgresql", "nativets", - "Nativets", "bun", "mysql", "bigquery", diff --git a/backend/.sqlx/query-9e8c3ff3d6b31e366e15beda1e96e03e870ccc3b353401439bc0ed8ff219249b.json b/backend/.sqlx/query-9e8c3ff3d6b31e366e15beda1e96e03e870ccc3b353401439bc0ed8ff219249b.json index 6308bf3bb2..fead4ba250 100644 --- a/backend/.sqlx/query-9e8c3ff3d6b31e366e15beda1e96e03e870ccc3b353401439bc0ed8ff219249b.json +++ b/backend/.sqlx/query-9e8c3ff3d6b31e366e15beda1e96e03e870ccc3b353401439bc0ed8ff219249b.json @@ -60,7 +60,6 @@ "bash", "postgresql", "nativets", - "Nativets", "bun", "mysql", "bigquery", diff --git a/backend/.sqlx/query-b224cdd1221fc9e7227ef8e8c025eedc09bceb82bca349f3e31c8513ebbf0192.json b/backend/.sqlx/query-b224cdd1221fc9e7227ef8e8c025eedc09bceb82bca349f3e31c8513ebbf0192.json index 38f81da395..c90719118a 100644 --- a/backend/.sqlx/query-b224cdd1221fc9e7227ef8e8c025eedc09bceb82bca349f3e31c8513ebbf0192.json +++ b/backend/.sqlx/query-b224cdd1221fc9e7227ef8e8c025eedc09bceb82bca349f3e31c8513ebbf0192.json @@ -42,7 +42,6 @@ "bash", "postgresql", "nativets", - "Nativets", "bun", "mysql", "bigquery", diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 3cb66a5ffb..e2f017e3d0 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -24,6 +24,7 @@ use axum::{ }; use base64::Engine; use bytes::Bytes; +use chrono::Utc; use hmac::Mac; use hyper::{header::CONTENT_TYPE, http, HeaderMap, Request, StatusCode}; use serde::{de::DeserializeOwned, Deserialize, Serialize}; @@ -41,7 +42,7 @@ use windmill_common::{ flows::FlowValue, jobs::{script_path_to_payload, JobKind, JobPayload, QueuedJob, RawCode}, oauth2::HmacSha256, - scripts::{ScriptHash, ScriptLang}, + scripts::{Script, ScriptHash, ScriptLang}, users::username_to_permissioned_as, utils::{not_found_if_none, now_from_db, paginate, require_admin, Pagination, StripPath}, }; @@ -105,7 +106,7 @@ pub fn workspaced_service() -> Router { .layer(cors.clone()), ) .route("/run/preview", post(run_preview_job)) - .route("/add_noop_jobs/:n", post(add_noop_jobs)) + .route("/add_batch_jobs/:n", post(add_batch_jobs)) .route("/run/preview_flow", post(run_preview_flow_job)) .route("/list", get(list_jobs)) .route("/queue/list", get(list_queue_jobs)) @@ -2357,53 +2358,131 @@ async fn run_preview_job( Ok((StatusCode::CREATED, uuid.to_string())) } +#[derive(Deserialize)] +struct BatchInfo { + kind: String, + flow_value: Option, + path: Option, + dedicated_worker: Option, +} + #[tracing::instrument(level = "trace", skip_all)] -async fn add_noop_jobs( +async fn add_batch_jobs( authed: ApiAuthed, Extension(db): Extension, Extension(rsmq): Extension>, Path((w_id, n)): Path<(String, i32)>, -) -> error::JsonResult> { + Json(batch_info): Json, +) -> error::JsonResult> { require_super_admin(&db, &authed.email).await?; - let mut tx = PushIsolationLevel::IsolatedRoot(db.clone(), rsmq); - let mut uuids: Vec = Vec::new(); - for _ in 0..n { - let (uuid, ntx) = push( - &db, - tx, - &w_id, - JobPayload::Noop, - serde_json::Map::new(), - &authed.username, - &authed.email, - username_to_permissioned_as(&authed.username), - None, - None, - None, - None, - None, - false, - false, - None, - true, - None, - None, - None, - ) - .await?; - tx = PushIsolationLevel::Transaction(ntx); - uuids.push(uuid.to_string()); - } - match tx { - PushIsolationLevel::Transaction(tx) => { - tx.commit().await?; + let (hash, path, job_kind, language, dedicated_worker) = match batch_info.kind.as_str() { + "script" => { + let script = sqlx::query_as::<_, Script>( + "select * from script where path = $1 and workspace_id = $2", + ) + .bind(&batch_info.path) + .bind(&w_id) + .fetch_optional(&db) + .await? + .ok_or_else(|| { + error::Error::BadRequest(format!("Script not found: {:?}", batch_info.path)) + })?; + ( + Some(script.hash), + batch_info.path, + JobKind::Script, + Some(script.language), + batch_info.dedicated_worker, + ) } - _ => (), - } + "flow" => { + let mut tx = PushIsolationLevel::IsolatedRoot(db.clone(), rsmq); + + let mut uuids: Vec = Vec::new(); + if batch_info.flow_value.is_none() { + return Err(error::Error::BadRequest( + "Flow value is required for batch flow".to_string(), + )); + } + for _ in 0..n { + let (uuid, ntx) = push( + &db, + tx, + &w_id, + JobPayload::RawFlow { + value: batch_info.flow_value.clone().unwrap(), + path: None, + }, + serde_json::Map::new(), + &authed.username, + &authed.email, + username_to_permissioned_as(&authed.username), + None, + None, + None, + None, + None, + false, + false, + None, + true, + None, + None, + None, + ) + .await?; + tx = PushIsolationLevel::Transaction(ntx); + uuids.push(uuid); + } + match tx { + PushIsolationLevel::Transaction(tx) => { + tx.commit().await?; + } + _ => (), + } + return Ok(Json(uuids)); + } + "noop" => (None, None, JobKind::Noop, None, None), + _ => { + return Err(error::Error::BadRequest(format!( + "Invalid batch kind: {}", + batch_info.kind + ))) + } + }; + + let language = language.unwrap_or(ScriptLang::Deno); + + let tag = if let Some(dedicated_worker) = dedicated_worker { + if dedicated_worker && path.is_some() { + format!("{}:{}", w_id, path.clone().unwrap()) + } else { + format!("{}", language.as_str()) + } + } else { + format!("{}", language.as_str()) + }; + + let uuids = sqlx::query_scalar!("INSERT INTO queue (id, script_hash, script_path, job_kind, language, tag, created_by, permissioned_as, email, scheduled_for, workspace_id) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id", + hash.map(|h| h.0), + path, + job_kind.clone() as JobKind, + language as ScriptLang, + tag, + authed.username, + authed.email, + username_to_permissioned_as(&authed.username), + Utc::now(), + w_id, + n + ) + .fetch_all(&db) + .await?; Ok(Json(uuids)) } + async fn run_preview_flow_job( authed: ApiAuthed, Extension(db): Extension,