From 17a86e69077f847e708d2b15f5cd0a8e6db82ea1 Mon Sep 17 00:00:00 2001 From: Lucas Abel <22837557+uael@users.noreply.github.com> Date: Fri, 7 Feb 2025 12:27:43 +0100 Subject: [PATCH] backend: fix workflow as code (#5239) * backend: improve `/get_job_update` after v2 * backend: insert missing `workflow_as_code_status` on completion also insert `flow_status` from so we can remove the query when `_duration` is above 500 * backend: fix workflow_as_code after v2 * backend: add `workflow_as_code` worker test --- ...fff3e22de8332d0c6d8ca98d22d62137fe701.json | 67 +++++++ ...df664b4de9aabf1e0e219596b295d52438008.json | 31 ++++ ...7cfb058a35a81314db981bc953a9505725082.json | 16 ++ ...d0e53392e238144b039cfe042587d7fe1dea3.json | 15 ++ ...50d9d258c288883b2b5b0ab286f5cb50850b5.json | 16 ++ ...4844b4c1a38701190795a801fa43c74a0beeb.json | 22 +++ ...a665692a2192ee664216bdab04e6d2547422f.json | 12 ++ ...4cbcb90e5a8d94d628cc78d1f2271c43e56f1.json | 14 ++ ...1514_v2_queue_compatibility_view.down.sql} | 0 ...131514_v2_queue_compatibility_view.up.sql} | 2 +- ...completed_job_compatibility_view.down.sql} | 0 ...2_completed_job_compatibility_view.up.sql} | 2 +- backend/tests/fixtures/base.sql | 3 +- backend/tests/worker.rs | 65 +++++++ backend/windmill-api/src/db.rs | 10 ++ backend/windmill-api/src/jobs.rs | 170 ++++++++---------- backend/windmill-common/src/jobs.rs | 2 +- backend/windmill-common/src/scripts.rs | 3 +- backend/windmill-queue/src/jobs.rs | 68 +++---- backend/windmill-worker/src/worker.rs | 17 +- 20 files changed, 383 insertions(+), 152 deletions(-) create mode 100644 backend/.sqlx/query-0128194fb539809e15bee670864fff3e22de8332d0c6d8ca98d22d62137fe701.json create mode 100644 backend/.sqlx/query-3a9441fe8fef1605d02e92b65d1df664b4de9aabf1e0e219596b295d52438008.json create mode 100644 backend/.sqlx/query-3eb447ed317f3d8724b2309cfdf7cfb058a35a81314db981bc953a9505725082.json create mode 100644 backend/.sqlx/query-45d616c92ebcbe30a563e1fa7d2d0e53392e238144b039cfe042587d7fe1dea3.json create mode 100644 backend/.sqlx/query-56f7325e3b0316866714e76d94b50d9d258c288883b2b5b0ab286f5cb50850b5.json create mode 100644 backend/.sqlx/query-867d5c75ddc6c5d20136880c7294844b4c1a38701190795a801fa43c74a0beeb.json create mode 100644 backend/.sqlx/query-976295cc04007f3cf5d8ba8d3fea665692a2192ee664216bdab04e6d2547422f.json create mode 100644 backend/.sqlx/query-d067bf59ed65562f1efbfbc9c264cbcb90e5a8d94d628cc78d1f2271c43e56f1.json rename backend/migrations/{20250201145630_v2_queue_compatibility_view.down.sql => 20250205131514_v2_queue_compatibility_view.down.sql} (100%) rename backend/migrations/{20250201145630_v2_queue_compatibility_view.up.sql => 20250205131514_v2_queue_compatibility_view.up.sql} (95%) rename backend/migrations/{20250201145631_v2_completed_job_compatibility_view.down.sql => 20250205131515_v2_completed_job_compatibility_view.down.sql} (100%) rename backend/migrations/{20250201145631_v2_completed_job_compatibility_view.up.sql => 20250205131515_v2_completed_job_compatibility_view.up.sql} (94%) diff --git a/backend/.sqlx/query-0128194fb539809e15bee670864fff3e22de8332d0c6d8ca98d22d62137fe701.json b/backend/.sqlx/query-0128194fb539809e15bee670864fff3e22de8332d0c6d8ca98d22d62137fe701.json new file mode 100644 index 0000000000..045993ce2e --- /dev/null +++ b/backend/.sqlx/query-0128194fb539809e15bee670864fff3e22de8332d0c6d8ca98d22d62137fe701.json @@ -0,0 +1,67 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT\n c.id IS NOT NULL AS completed,\n q.id IS NOT NULL AND q.running AS running,\n SUBSTR(logs, GREATEST($1 - log_offset, 0)) AS logs,\n COALESCE(r.memory_peak, c.memory_peak) AS mem_peak,\n CASE\n -- flow step:\n WHEN flow_step_id IS NOT NULL THEN NULL\n -- completed:\n WHEN c.id IS NOT NULL THEN COALESCE(\n c.workflow_as_code_status || c.flow_status,\n c.workflow_as_code_status,\n c.flow_status\n )\n -- not completed:\n ELSE COALESCE(\n f.workflow_as_code_status || f.flow_status,\n f.workflow_as_code_status,\n f.flow_status\n )\n END AS \"flow_status: sqlx::types::Json>\",\n job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 AS log_offset,\n created_by AS \"created_by!\",\n CASE WHEN $4::BOOLEAN THEN (\n SELECT scalar_int FROM job_stats WHERE job_id = $3 AND metric_id = 'progress_perc'\n ) END AS progress\n FROM v2_job j\n LEFT JOIN v2_job_queue q USING (id)\n LEFT JOIN v2_job_runtime r USING (id)\n LEFT JOIN v2_job_status f USING (id)\n LEFT JOIN v2_job_completed c USING (id)\n LEFT JOIN job_logs ON job_logs.job_id = $3\n WHERE j.workspace_id = $2 AND j.id = $3", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "completed", + "type_info": "Bool" + }, + { + "ordinal": 1, + "name": "running", + "type_info": "Bool" + }, + { + "ordinal": 2, + "name": "logs", + "type_info": "Text" + }, + { + "ordinal": 3, + "name": "mem_peak", + "type_info": "Int4" + }, + { + "ordinal": 4, + "name": "flow_status: sqlx::types::Json>", + "type_info": "Jsonb" + }, + { + "ordinal": 5, + "name": "log_offset", + "type_info": "Int4" + }, + { + "ordinal": 6, + "name": "created_by!", + "type_info": "Varchar" + }, + { + "ordinal": 7, + "name": "progress", + "type_info": "Int4" + } + ], + "parameters": { + "Left": [ + "Int4", + "Text", + "Uuid", + "Bool" + ] + }, + "nullable": [ + null, + null, + null, + null, + null, + null, + false, + null + ] + }, + "hash": "0128194fb539809e15bee670864fff3e22de8332d0c6d8ca98d22d62137fe701" +} diff --git a/backend/.sqlx/query-3a9441fe8fef1605d02e92b65d1df664b4de9aabf1e0e219596b295d52438008.json b/backend/.sqlx/query-3a9441fe8fef1605d02e92b65d1df664b4de9aabf1e0e219596b295d52438008.json new file mode 100644 index 0000000000..4a789635ee --- /dev/null +++ b/backend/.sqlx/query-3a9441fe8fef1605d02e92b65d1df664b4de9aabf1e0e219596b295d52438008.json @@ -0,0 +1,31 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO v2_job_completed AS cj\n ( workspace_id\n , id\n , started_at\n , duration_ms\n , result\n , result_columns\n , canceled_by\n , canceled_reason\n , flow_status\n , workflow_as_code_status\n , memory_peak\n , status\n )\n SELECT q.workspace_id, q.id, started_at, COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(started_at, now()))))*1000), $3, $10, $5, $6,\n flow_status, workflow_as_code_status,\n $8, CASE WHEN $4::BOOL THEN 'canceled'::job_status\n WHEN $7::BOOL THEN 'skipped'::job_status\n WHEN $2::BOOL THEN 'success'::job_status\n ELSE 'failure'::job_status END AS status\n FROM v2_job_queue q LEFT JOIN v2_job_status USING (id) WHERE q.id = $1\n ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3 RETURNING duration_ms AS \"duration_ms!\"", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "duration_ms!", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Bool", + "Jsonb", + "Bool", + "Varchar", + "Text", + "Bool", + "Int4", + "Int8", + "TextArray" + ] + }, + "nullable": [ + false + ] + }, + "hash": "3a9441fe8fef1605d02e92b65d1df664b4de9aabf1e0e219596b295d52438008" +} diff --git a/backend/.sqlx/query-3eb447ed317f3d8724b2309cfdf7cfb058a35a81314db981bc953a9505725082.json b/backend/.sqlx/query-3eb447ed317f3d8724b2309cfdf7cfb058a35a81314db981bc953a9505725082.json new file mode 100644 index 0000000000..69266cc130 --- /dev/null +++ b/backend/.sqlx/query-3eb447ed317f3d8724b2309cfdf7cfb058a35a81314db981bc953a9505725082.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO v2_job_status (id, workflow_as_code_status)\n VALUES ($1, JSONB_SET('{}'::JSONB, array[$2], $3))\n ON CONFLICT (id) DO UPDATE SET workflow_as_code_status =\n COALESCE(EXCLUDED.workflow_as_code_status, '{}'::JSONB) || $3", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Uuid", + "Text", + "Jsonb" + ] + }, + "nullable": [] + }, + "hash": "3eb447ed317f3d8724b2309cfdf7cfb058a35a81314db981bc953a9505725082" +} diff --git a/backend/.sqlx/query-45d616c92ebcbe30a563e1fa7d2d0e53392e238144b039cfe042587d7fe1dea3.json b/backend/.sqlx/query-45d616c92ebcbe30a563e1fa7d2d0e53392e238144b039cfe042587d7fe1dea3.json new file mode 100644 index 0000000000..052238874a --- /dev/null +++ b/backend/.sqlx/query-45d616c92ebcbe30a563e1fa7d2d0e53392e238144b039cfe042587d7fe1dea3.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE v2_job_status SET\n workflow_as_code_status = jsonb_set(\n jsonb_set(\n COALESCE(workflow_as_code_status, '{}'::jsonb),\n array[$1],\n COALESCE(workflow_as_code_status->$1, '{}'::jsonb)\n ),\n array[$1, 'started_at'],\n to_jsonb(now()::text)\n )\n WHERE id = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "45d616c92ebcbe30a563e1fa7d2d0e53392e238144b039cfe042587d7fe1dea3" +} diff --git a/backend/.sqlx/query-56f7325e3b0316866714e76d94b50d9d258c288883b2b5b0ab286f5cb50850b5.json b/backend/.sqlx/query-56f7325e3b0316866714e76d94b50d9d258c288883b2b5b0ab286f5cb50850b5.json new file mode 100644 index 0000000000..6b47103c3a --- /dev/null +++ b/backend/.sqlx/query-56f7325e3b0316866714e76d94b50d9d258c288883b2b5b0ab286f5cb50850b5.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE v2_job_status SET\n workflow_as_code_status = jsonb_set(\n jsonb_set(\n COALESCE(workflow_as_code_status, '{}'::jsonb),\n array[$1],\n COALESCE(workflow_as_code_status->$1, '{}'::jsonb)\n ),\n array[$1, 'duration_ms'],\n to_jsonb($2::bigint)\n )\n WHERE id = $3", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Int8", + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "56f7325e3b0316866714e76d94b50d9d258c288883b2b5b0ab286f5cb50850b5" +} diff --git a/backend/.sqlx/query-867d5c75ddc6c5d20136880c7294844b4c1a38701190795a801fa43c74a0beeb.json b/backend/.sqlx/query-867d5c75ddc6c5d20136880c7294844b4c1a38701190795a801fa43c74a0beeb.json new file mode 100644 index 0000000000..9232baf39c --- /dev/null +++ b/backend/.sqlx/query-867d5c75ddc6c5d20136880c7294844b4c1a38701190795a801fa43c74a0beeb.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT workflow_as_code_status FROM v2_job_completed WHERE id = $1", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "workflow_as_code_status", + "type_info": "Jsonb" + } + ], + "parameters": { + "Left": [ + "Uuid" + ] + }, + "nullable": [ + true + ] + }, + "hash": "867d5c75ddc6c5d20136880c7294844b4c1a38701190795a801fa43c74a0beeb" +} diff --git a/backend/.sqlx/query-976295cc04007f3cf5d8ba8d3fea665692a2192ee664216bdab04e6d2547422f.json b/backend/.sqlx/query-976295cc04007f3cf5d8ba8d3fea665692a2192ee664216bdab04e6d2547422f.json new file mode 100644 index 0000000000..b32af76e07 --- /dev/null +++ b/backend/.sqlx/query-976295cc04007f3cf5d8ba8d3fea665692a2192ee664216bdab04e6d2547422f.json @@ -0,0 +1,12 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM _sqlx_migrations WHERE version=20250201145630 OR version=20250201145631", + "describe": { + "columns": [], + "parameters": { + "Left": [] + }, + "nullable": [] + }, + "hash": "976295cc04007f3cf5d8ba8d3fea665692a2192ee664216bdab04e6d2547422f" +} diff --git a/backend/.sqlx/query-d067bf59ed65562f1efbfbc9c264cbcb90e5a8d94d628cc78d1f2271c43e56f1.json b/backend/.sqlx/query-d067bf59ed65562f1efbfbc9c264cbcb90e5a8d94d628cc78d1f2271c43e56f1.json new file mode 100644 index 0000000000..53490d5b60 --- /dev/null +++ b/backend/.sqlx/query-d067bf59ed65562f1efbfbc9c264cbcb90e5a8d94d628cc78d1f2271c43e56f1.json @@ -0,0 +1,14 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO v2_job_status (id, workflow_as_code_status) VALUES ($1, '{}'::JSONB)\n ON CONFLICT (id) DO NOTHING", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "d067bf59ed65562f1efbfbc9c264cbcb90e5a8d94d628cc78d1f2271c43e56f1" +} diff --git a/backend/migrations/20250201145630_v2_queue_compatibility_view.down.sql b/backend/migrations/20250205131514_v2_queue_compatibility_view.down.sql similarity index 100% rename from backend/migrations/20250201145630_v2_queue_compatibility_view.down.sql rename to backend/migrations/20250205131514_v2_queue_compatibility_view.down.sql diff --git a/backend/migrations/20250201145630_v2_queue_compatibility_view.up.sql b/backend/migrations/20250205131514_v2_queue_compatibility_view.up.sql similarity index 95% rename from backend/migrations/20250201145630_v2_queue_compatibility_view.up.sql rename to backend/migrations/20250205131514_v2_queue_compatibility_view.up.sql index 76065477f9..c0b79bf9b6 100644 --- a/backend/migrations/20250201145630_v2_queue_compatibility_view.up.sql +++ b/backend/migrations/20250205131514_v2_queue_compatibility_view.up.sql @@ -21,7 +21,7 @@ SELECT CASE WHEN j.trigger_kind = 'schedule'::job_trigger_kind THEN j.trigger END AS schedule_path, j.permissioned_as, - s.flow_status, + COALESCE(s.flow_status, s.workflow_as_code_status) AS flow_status, j.raw_flow, j.flow_step_id IS NOT NULL AS is_flow_step, j.script_lang AS language, diff --git a/backend/migrations/20250201145631_v2_completed_job_compatibility_view.down.sql b/backend/migrations/20250205131515_v2_completed_job_compatibility_view.down.sql similarity index 100% rename from backend/migrations/20250201145631_v2_completed_job_compatibility_view.down.sql rename to backend/migrations/20250205131515_v2_completed_job_compatibility_view.down.sql diff --git a/backend/migrations/20250201145631_v2_completed_job_compatibility_view.up.sql b/backend/migrations/20250205131515_v2_completed_job_compatibility_view.up.sql similarity index 94% rename from backend/migrations/20250201145631_v2_completed_job_compatibility_view.up.sql rename to backend/migrations/20250205131515_v2_completed_job_compatibility_view.up.sql index dd5b11ed2c..51709c42c7 100644 --- a/backend/migrations/20250201145631_v2_completed_job_compatibility_view.up.sql +++ b/backend/migrations/20250205131515_v2_completed_job_compatibility_view.up.sql @@ -21,7 +21,7 @@ SELECT CASE WHEN j.trigger_kind = 'schedule'::job_trigger_kind THEN j.trigger END AS schedule_path, j.permissioned_as, - c.flow_status, + COALESCE(c.flow_status, c.workflow_as_code_status) AS flow_status, j.raw_flow, j.flow_step_id IS NOT NULL AS is_flow_step, j.script_lang AS language, diff --git a/backend/tests/fixtures/base.sql b/backend/tests/fixtures/base.sql index c5bec3454c..0e38b5c7bd 100644 --- a/backend/tests/fixtures/base.sql +++ b/backend/tests/fixtures/base.sql @@ -152,7 +152,8 @@ BEGIN -- v2_job_status: IF EXISTS(SELECT 1 FROM v2_job_status WHERE id = OLD.id) THEN SELECT * INTO job_status FROM v2_job_status WHERE id = OLD.id; - IF job_status.flow_status::TEXT IS DISTINCT FROM OLD.__flow_status::TEXT THEN + IF COALESCE(job_status.flow_status, job_status.workflow_as_code_status)::TEXT IS DISTINCT FROM OLD.__flow_status::TEXT + THEN RAISE EXCEPTION 'flow_status mismatch'; END IF; IF job_status.flow_leaf_jobs::TEXT IS DISTINCT FROM OLD.__leaf_jobs::TEXT THEN diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 2b4768c34e..12298795ad 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -3828,6 +3828,71 @@ async fn test_job_labels(db: Pool) { test(&["z", "a", "x"]).await; } +#[cfg(feature = "python")] +const WORKFLOW_AS_CODE: &str = r#" +from wmill import task + +import pandas as pd +import numpy as np + +@task() +def heavy_compute(n: int): + df = pd.DataFrame(np.random.randn(100, 4), columns=list('ABCD')) + return df.sum().sum() + +@task +def send_result(res: int, email: str): + print(f"Sending result {res} to {email}") + return "OK" + +def main(n: int): + l = [] + for i in range(n): + l.append(heavy_compute(i)) + print(l) + return [send_result(sum(l), "example@example.com"), n] +"#; + +#[cfg(feature = "python")] +#[sqlx::test(fixtures("base", "hello"))] +async fn test_workflow_as_code(db: Pool) { + initialize_tracing().await; + let server = ApiServer::start(db.clone()).await; + let port = server.addr.port(); + + // workflow as code require at least 2 workers: + let db = &db; + in_test_worker( + &db, + async move { + let job = RunJob::from(JobPayload::Code(RawCode { + language: ScriptLang::Python3, + content: WORKFLOW_AS_CODE.into(), + ..RawCode::default() + })) + .arg("n", json!(3)) + .run_until_complete(&db, port) + .await; + + assert_eq!(job.json_result().unwrap(), json!(["OK", 3])); + let workflow_as_code_status = sqlx::query_scalar!( + "SELECT workflow_as_code_status FROM v2_job_completed WHERE id = $1", + job.id + ) + .fetch_one(db) + .await + .unwrap() + .unwrap(); + assert_eq!( + workflow_as_code_status.get("name"), + Some(&json!("send_result")) + ); + }, + port, + ) + .await; +} + async fn test_for_versions>( version_flags: impl Iterator>>, test: impl Fn() -> F, diff --git a/backend/windmill-api/src/db.rs b/backend/windmill-api/src/db.rs index 5cf6b0073c..a2a641b802 100644 --- a/backend/windmill-api/src/db.rs +++ b/backend/windmill-api/src/db.rs @@ -188,6 +188,16 @@ pub async fn migrate(db: &DB) -> Result<(), Error> { tracing::info!("Could not remove sqlx migration with version=20250201145632: {err:#}"); } + // New version of `v2_as_queue` and `v2_as_completed_job` VIEWs. + if let Err(err) = sqlx::query!( + "DELETE FROM _sqlx_migrations WHERE version=20250201145630 OR version=20250201145631" + ) + .execute(db) + .await + { + tracing::info!("Could not remove sqlx migration with version=[20250201145630, 20250201145631] : {err:#}"); + } + match sqlx::migrate!("../migrations") .run_direct(&mut custom_migrator) .await diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index d1630a7017..84d6428518 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -3425,13 +3425,25 @@ pub async fn run_workflow_as_code( if !wkflow_query.skip_update.unwrap_or(false) { sqlx::query!( - "UPDATE v2_job_status SET flow_status = jsonb_set(COALESCE(flow_status, '{}'::jsonb), array[$1], jsonb_set(jsonb_set('{}'::jsonb, '{scheduled_for}', to_jsonb(now()::text)), '{name}', to_jsonb($3::text))) WHERE id = $2", - uuid.to_string(), + "INSERT INTO v2_job_status (id, workflow_as_code_status) + VALUES ($1, JSONB_SET('{}'::JSONB, array[$2], $3)) + ON CONFLICT (id) DO UPDATE SET workflow_as_code_status = + COALESCE(EXCLUDED.workflow_as_code_status, '{}'::JSONB) || $3", job_id, - entrypoint - ).execute(&mut *tx).await?; + uuid.to_string(), + serde_json::json!({ "scheduled_for": Utc::now(), "name": entrypoint }), + ) + .execute(&mut *tx) + .await?; } else { tracing::info!("Skipping update of flow status for job {job_id} in workspace {w_id}"); + sqlx::query!( + "INSERT INTO v2_job_status (id, workflow_as_code_status) VALUES ($1, '{}'::JSONB) + ON CONFLICT (id) DO NOTHING", + job_id, + ) + .execute(&mut *tx) + .await?; } if *CLOUD_HOSTED { @@ -5159,114 +5171,72 @@ async fn get_log_file(Path((_w_id, file_p)): Path<(String, String)>) -> error::R ))); } -#[derive(Deserialize, sqlx::FromRow)] -pub struct JobUpdateRow { - pub running: bool, - pub logs: Option, - pub mem_peak: Option, - pub flow_status: Option>>, - pub log_offset: Option, - pub created_by: String, -} async fn get_job_update( OptAuthed(opt_authed): OptAuthed, Extension(db): Extension, Path((w_id, job_id)): Path<(String, Uuid)>, - Query(JobUpdateQuery { running, log_offset, get_progress }): Query, -) -> error::JsonResult { + Query(JobUpdateQuery { log_offset, get_progress, .. }): Query, +) -> JsonResult { let record = sqlx::query!( "SELECT - running AS \"running!\", - substr(concat(coalesce(v2_as_queue.logs, ''), job_logs.logs), greatest($1 - job_logs.log_offset, 0)) AS logs, - mem_peak, - CASE WHEN is_flow_step is true then NULL else flow_status END AS \"flow_status: sqlx::types::Json>\", - job_logs.log_offset + char_length(job_logs.logs) + 1 AS log_offset, - created_by AS \"created_by!\" - FROM v2_as_queue - LEFT JOIN job_logs ON job_logs.job_id = v2_as_queue.id - WHERE v2_as_queue.workspace_id = $2 AND v2_as_queue.id = $3", + c.id IS NOT NULL AS completed, + q.id IS NOT NULL AND q.running AS running, + SUBSTR(logs, GREATEST($1 - log_offset, 0)) AS logs, + COALESCE(r.memory_peak, c.memory_peak) AS mem_peak, + CASE + -- flow step: + WHEN flow_step_id IS NOT NULL THEN NULL + -- completed: + WHEN c.id IS NOT NULL THEN COALESCE( + c.workflow_as_code_status || c.flow_status, + c.workflow_as_code_status, + c.flow_status + ) + -- not completed: + ELSE COALESCE( + f.workflow_as_code_status || f.flow_status, + f.workflow_as_code_status, + f.flow_status + ) + END AS \"flow_status: sqlx::types::Json>\", + job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 AS log_offset, + created_by AS \"created_by!\", + CASE WHEN $4::BOOLEAN THEN ( + SELECT scalar_int FROM job_stats WHERE job_id = $3 AND metric_id = 'progress_perc' + ) END AS progress + FROM v2_job j + LEFT JOIN v2_job_queue q USING (id) + LEFT JOIN v2_job_runtime r USING (id) + LEFT JOIN v2_job_status f USING (id) + LEFT JOIN v2_job_completed c USING (id) + LEFT JOIN job_logs ON job_logs.job_id = $3 + WHERE j.workspace_id = $2 AND j.id = $3", log_offset, &w_id, - job_id + job_id, + get_progress.unwrap_or(false) ) .fetch_optional(&db) - .await?; + .await? + .ok_or_else(|| Error::NotFound(format!("Job not found: {}", job_id)))?; - let progress: Option = if get_progress == Some(true) { - sqlx::query_scalar!( - "SELECT scalar_int FROM job_stats WHERE workspace_id = $1 AND job_id = $2 AND metric_id = $3", - &w_id, - job_id, - "progress_perc" - ) - .fetch_optional(&db) - .await?.and_then(|inner| inner) - } else { - None - }; - - if let Some(record) = record { - if opt_authed.is_none() && record.created_by != "anonymous" { - return Err(Error::BadRequest( - "As a non logged in user, you can only see jobs ran by anonymous users".to_string(), - )); - } - log_job_view(&db, opt_authed.as_ref(), &w_id, &job_id).await?; - Ok(Json(JobUpdate { - running: if !running && record.running { - Some(true) - } else { - None - }, - log_offset: record.log_offset, - completed: None, - new_logs: record.logs, - mem_peak: record.mem_peak, - progress, - flow_status: record - .flow_status - .map(|x: sqlx::types::Json>| x.0), - })) - } else { - let record = sqlx::query!( - "SELECT - substr(concat(coalesce(v2_as_completed_job.logs, ''), job_logs.logs), greatest($1 - job_logs.log_offset, 0)) AS logs, - mem_peak, - CASE WHEN is_flow_step is true then NULL else flow_status END AS \"flow_status: sqlx::types::Json>\", - job_logs.log_offset + char_length(job_logs.logs) + 1 AS log_offset, - created_by AS \"created_by!\" - FROM v2_as_completed_job - LEFT JOIN job_logs ON job_logs.job_id = v2_as_completed_job.id - WHERE v2_as_completed_job.workspace_id = $2 AND v2_as_completed_job.id = $3", - log_offset, - &w_id, - job_id - ) - .fetch_optional(&db) - .await?; - if let Some(record) = record { - if opt_authed.is_none() && record.created_by != "anonymous" { - return Err(Error::BadRequest( - "As a non logged in user, you can only see jobs ran by anonymous users" - .to_string(), - )); - } - log_job_view(&db, opt_authed.as_ref(), &w_id, &job_id).await?; - Ok(Json(JobUpdate { - running: Some(false), - completed: Some(true), - log_offset: record.log_offset, - new_logs: record.logs, - mem_peak: record.mem_peak, - progress, - flow_status: record - .flow_status - .map(|x: sqlx::types::Json>| x.0), - })) - } else { - Err(error::Error::NotFound(format!("Job not found: {}", job_id))) - } + if opt_authed.is_none() && record.created_by != "anonymous" { + return Err(Error::BadRequest( + "As a non logged in user, you can only see jobs ran by anonymous users".to_string(), + )); } + log_job_view(&db, opt_authed.as_ref(), &w_id, &job_id).await?; + Ok(Json(JobUpdate { + running: record.running, + completed: record.completed, + log_offset: record.log_offset, + new_logs: record.logs, + mem_peak: record.mem_peak, + progress: record.progress, + flow_status: record + .flow_status + .map(|x: sqlx::types::Json>| x.0), + })) } pub fn filter_list_completed_query( diff --git a/backend/windmill-common/src/jobs.rs b/backend/windmill-common/src/jobs.rs index 1c764b12df..49de1c2415 100644 --- a/backend/windmill-common/src/jobs.rs +++ b/backend/windmill-common/src/jobs.rs @@ -353,7 +353,7 @@ pub enum JobPayload { Noop, } -#[derive(Clone, Serialize, Deserialize, Debug)] +#[derive(Clone, Serialize, Deserialize, Debug, Default)] pub struct RawCode { pub content: String, pub path: Option, diff --git a/backend/windmill-common/src/scripts.rs b/backend/windmill-common/src/scripts.rs index 2356f5ee70..9d896e7f5b 100644 --- a/backend/windmill-common/src/scripts.rs +++ b/backend/windmill-common/src/scripts.rs @@ -24,11 +24,12 @@ use serde::{ser::SerializeSeq, Deserialize, Deserializer, Serialize}; use crate::utils::StripPath; -#[derive(Serialize, Deserialize, Debug, PartialEq, Copy, Clone, Hash, Eq, sqlx::Type)] +#[derive(Serialize, Deserialize, Debug, PartialEq, Copy, Clone, Hash, Eq, sqlx::Type, Default)] #[sqlx(type_name = "SCRIPT_LANG", rename_all = "lowercase")] #[serde(rename_all(serialize = "lowercase", deserialize = "lowercase"))] pub enum ScriptLang { Nativets, + #[default] Deno, Python3, Go, diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 48f66fcf5b..c76747af5a 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -591,28 +591,28 @@ pub async fn add_completed_job( , canceled_by , canceled_reason , flow_status + , workflow_as_code_status , memory_peak , status ) - VALUES ($1, $2, $3, COALESCE($12::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE($3, now()))))*1000), $5, $13, $7, $8, $9,\ - $11, CASE WHEN $6::BOOL THEN 'canceled'::job_status - WHEN $10::BOOL THEN 'skipped'::job_status - WHEN $4::BOOL THEN 'success'::job_status - ELSE 'failure'::job_status END) - ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $5 RETURNING duration_ms AS \"duration_ms!\"", - /* $1 */ queued_job.workspace_id, - /* $2 */ queued_job.id, - /* $3 */ queued_job.started_at, - /* $4 */ success, - /* $5 */ result as Json<&T>, - /* $6 */ canceled_by.is_some(), - /* $7 */ canceled_by.clone().map(|cb| cb.username).flatten(), - /* $8 */ canceled_by.clone().map(|cb| cb.reason).flatten(), - /* $9 */ &queued_job.flow_status as &Option>>, - /* $10 */ skipped, - /* $11 */ if mem_peak > 0 { Some(mem_peak) } else { None }, - /* $12 */ duration, - /* $13 */ result_columns as Option<&Vec>, + SELECT q.workspace_id, q.id, started_at, COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(started_at, now()))))*1000), $3, $10, $5, $6, + flow_status, workflow_as_code_status, + $8, CASE WHEN $4::BOOL THEN 'canceled'::job_status + WHEN $7::BOOL THEN 'skipped'::job_status + WHEN $2::BOOL THEN 'success'::job_status + ELSE 'failure'::job_status END AS status + FROM v2_job_queue q LEFT JOIN v2_job_status USING (id) WHERE q.id = $1 + ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3 RETURNING duration_ms AS \"duration_ms!\"", + /* $1 */ queued_job.id, + /* $2 */ success, + /* $3 */ result as Json<&T>, + /* $4 */ canceled_by.is_some(), + /* $5 */ canceled_by.clone().map(|cb| cb.username).flatten(), + /* $6 */ canceled_by.clone().map(|cb| cb.reason).flatten(), + /* $7 */ skipped, + /* $8 */ if mem_peak > 0 { Some(mem_peak) } else { None }, + /* $9 */ duration, + /* $10 */ result_columns as Option<&Vec>, ) .fetch_one(&mut *tx) .await @@ -633,28 +633,14 @@ pub async fn add_completed_job( } if !queued_job.is_flow_step { - if _duration > 500 - && (queued_job.job_kind == JobKind::Script - || queued_job.job_kind == JobKind::Preview) - { - if let Err(e) = sqlx::query!( - "UPDATE v2_job_completed SET flow_status = f.flow_status FROM v2_job_status f WHERE v2_job_completed.id = $1 AND f.id = $1 AND v2_job_completed.workspace_id = $2", - &queued_job.id, - &queued_job.workspace_id - ) - .execute(&mut *tx) - .await { - tracing::error!("Could not update job duration: {}", e); - } - } if let Some(parent_job) = queued_job.parent_job { - if let Err(e) = sqlx::query_scalar!( + let _ = sqlx::query_scalar!( "UPDATE v2_job_status SET - flow_status = jsonb_set( + workflow_as_code_status = jsonb_set( jsonb_set( - COALESCE(flow_status, '{}'::jsonb), + COALESCE(workflow_as_code_status, '{}'::jsonb), array[$1], - COALESCE(flow_status->$1, '{}'::jsonb) + COALESCE(workflow_as_code_status->$1, '{}'::jsonb) ), array[$1, 'duration_ms'], to_jsonb($2::bigint) @@ -665,9 +651,11 @@ pub async fn add_completed_job( parent_job ) .execute(&mut *tx) - .await { - tracing::error!("Could not update parent job flow_status: {}", e); - } + .await + .inspect_err(|e| tracing::error!( + "Could not update parent job `duration_ms` in workflow as code status: {}", + e, + )); } } // tracing::error!("Added completed job {:#?}", queued_job); diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index b1f2ad1f65..31c076de9d 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -2006,13 +2006,13 @@ async fn handle_queued_job( .warn_after_seconds(5) .await?; } else if let Some(parent_job) = job.parent_job { - if let Err(e) = sqlx::query_scalar!( + let _ = sqlx::query_scalar!( "UPDATE v2_job_status SET - flow_status = jsonb_set( + workflow_as_code_status = jsonb_set( jsonb_set( - COALESCE(flow_status, '{}'::jsonb), + COALESCE(workflow_as_code_status, '{}'::jsonb), array[$1], - COALESCE(flow_status->$1, '{}'::jsonb) + COALESCE(workflow_as_code_status->$1, '{}'::jsonb) ), array[$1, 'started_at'], to_jsonb(now()::text) @@ -2024,9 +2024,12 @@ async fn handle_queued_job( .execute(db) .warn_after_seconds(5) .await - { - tracing::error!("Could not update parent job started_at flow_status: {}", e); - } + .inspect_err(|e| { + tracing::error!( + "Could not update parent job `started_at` in workflow as code status: {}", + e + ) + }); } let started = Instant::now();