perf: complete a job in one statement on the common path (#11355)

* perf: complete a job in one statement on the common path

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: take completion locks in one order on every path

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: keep a losing zombie completion from touching its wac parent

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: leave a flow's ping alone when a step completes during its cancel

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* test: probe only this test's completion for the lock wait

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix: stamp a wac child's kept duration when its completed row exists

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TJMJSJ2bDhYh7Shoh78Yyb

* perf: leave the parent ping out of completions with no flow to ping

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TJMJSJ2bDhYh7Shoh78Yyb

* docs: note that the two completion statements must stay in step

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TJMJSJ2bDhYh7Shoh78Yyb

* chore: update ee-repo-ref to 7a256cf353db7cf64a60a09fa0de7f3a8b27f626

This commit updates the EE repository reference after PR #830 was merged in windmill-ee-private.

Previous ee-repo-ref: 497137acb65e521568d46f3cbe1d66359f7f87ec

New ee-repo-ref: 7a256cf353db7cf64a60a09fa0de7f3a8b27f626

Automated by sync-ee-ref workflow.

---------

Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
This commit is contained in:
Ruben Fiszel
2026-09-26 00:18:38 +02:00
committed by GitHub
co-authored by Claude Opus 5.5 windmill-internal-app[bot]
parent e2be584ca5
commit bebd762194
15 changed files with 714 additions and 335 deletions
@@ -1,31 +0,0 @@
{
"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 , worker\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::text::jsonb, $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 q.worker\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::text::jsonb RETURNING duration_ms AS \"duration_ms!\"",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "duration_ms!",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Uuid",
"Bool",
"Text",
"Bool",
"Varchar",
"Text",
"Bool",
"Int4",
"Int8",
"TextArray"
]
},
"nullable": [
false
]
},
"hash": "042ff3003bf82d11a78a1074a63379c19880e4d66dbf4dd5391f053b6dfa6b01"
}
@@ -1,23 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO v2_job_completed\n (workspace_id, id, started_at, duration_ms, result, memory_peak, status, worker)\n SELECT q.workspace_id, q.id, q.started_at,\n COALESCE((EXTRACT('epoch' FROM now()) - EXTRACT('epoch' FROM COALESCE(q.started_at, now()))) * 1000, 0)::bigint,\n $2::jsonb, r.memory_peak, 'failure'::job_status, q.worker\n FROM v2_job_queue q\n LEFT JOIN v2_job_runtime r ON r.id = q.id\n WHERE q.id = $1\n ON CONFLICT (id) DO UPDATE SET status = 'failure', result = $2::jsonb\n RETURNING duration_ms AS \"duration_ms!\"",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "duration_ms!",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Uuid",
"Jsonb"
]
},
"nullable": [
false
]
},
"hash": "2abc2a5830130b2b4b32983407abeea41923ba4122fa666c6e8d8b06dcb5a71f"
}
@@ -1,27 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "\n WITH step_index AS (\n SELECT idx::text AS idx\n FROM v2_job_status,\n jsonb_array_elements(flow_status->'modules') WITH ORDINALITY arr(elem, idx)\n WHERE id = $1\n AND elem->>'id' = $5\n LIMIT 1\n ), completed AS (\n INSERT INTO v2_job_completed\n (workspace_id, id, started_at, duration_ms, result,\n flow_status, workflow_as_code_status, status, worker)\n SELECT\n q.workspace_id, q.id, q.started_at,\n (EXTRACT('epoch' FROM now()) - EXTRACT('epoch' FROM COALESCE(q.started_at, now()))) * 1000,\n $3::text::jsonb,\n CASE WHEN si.idx IS NOT NULL\n THEN jsonb_set(\n s.flow_status,\n ARRAY['modules', (si.idx::int - 1)::text],\n $6::jsonb\n )\n ELSE s.flow_status\n END,\n s.workflow_as_code_status,\n 'skipped'::job_status,\n q.worker\n FROM v2_job_queue q\n LEFT JOIN v2_job_status s ON s.id = q.id\n LEFT JOIN step_index si ON true\n WHERE q.id = $1\n ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = EXCLUDED.result\n RETURNING 1 AS x\n ), _deleted AS (\n DELETE FROM v2_job_queue WHERE id = $1\n ), _logged AS (\n INSERT INTO job_logs (logs, job_id, workspace_id)\n VALUES ($4, $1, $2)\n ON CONFLICT (job_id) DO UPDATE SET logs = concat(job_logs.logs, EXCLUDED.logs)\n )\n SELECT x FROM completed\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "x",
"type_info": "Int4"
}
],
"parameters": {
"Left": [
"Uuid",
"Varchar",
"Text",
"Text",
"Text",
"Jsonb"
]
},
"nullable": [
null
]
},
"hash": "4461fe84370e7f07b4423ea5e71b913dd6ee9c655f9ec7087cc0fec9b9c3099a"
}
@@ -0,0 +1,22 @@
{
"db_name": "PostgreSQL",
"query": "SELECT COALESCE(c.duration_ms, COALESCE((EXTRACT('epoch' FROM now()) - EXTRACT('epoch' FROM COALESCE(q.started_at, now()))) * 1000, 0)::bigint) AS \"duration_ms!\"\n FROM v2_job_queue q LEFT JOIN v2_job_completed c ON c.id = q.id WHERE q.id = $1",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "duration_ms!",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Uuid"
]
},
"nullable": [
null
]
},
"hash": "538f12535bd1b93a3a9baac91dc0d4ee0f3c949a00f62d38686e75d22e9d7884"
}
@@ -0,0 +1,23 @@
{
"db_name": "PostgreSQL",
"query": "SELECT COALESCE(c.duration_ms, COALESCE($2::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(q.started_at, now()))))*1000)::bigint) AS \"duration_ms!\"\n FROM v2_job_queue q LEFT JOIN v2_job_completed c ON c.id = q.id WHERE q.id = $1",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "duration_ms!",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Uuid",
"Int8"
]
},
"nullable": [
null
]
},
"hash": "67bde191639648d1f225a046eac57f4055900fdf0b1f65979deb1c1c64092a8e"
}
@@ -0,0 +1,39 @@
{
"db_name": "PostgreSQL",
"query": "WITH deleted AS (\n DELETE FROM v2_job_queue WHERE id = $1\n RETURNING id, workspace_id, started_at, worker, canceled_by, canceled_reason\n ), completed AS (\n 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 , worker\n )\n SELECT d.workspace_id, d.id, d.started_at,\n COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(d.started_at, now()))))*1000),\n $3::text::jsonb, $10,\n CASE WHEN $4::BOOL THEN $5 ELSE d.canceled_by END,\n CASE WHEN $4::BOOL THEN $6 WHEN d.canceled_by IS NOT NULL THEN d.canceled_reason END,\n s.flow_status, s.workflow_as_code_status, $8,\n CASE WHEN $4::BOOL OR d.canceled_by IS NOT NULL 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,\n d.worker\n FROM deleted d LEFT JOIN v2_job_status s ON s.id = d.id\n ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3::text::jsonb,\n canceled_by = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL\n THEN EXCLUDED.canceled_by ELSE cj.canceled_by END,\n canceled_reason = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL\n THEN EXCLUDED.canceled_reason ELSE cj.canceled_reason END\n RETURNING duration_ms\n ), parent_ping AS (\n UPDATE v2_job_runtime r SET ping = now()\n FROM v2_job_queue q\n WHERE r.id = $11 AND q.id = r.id AND q.workspace_id = $12 AND q.canceled_by IS NULL\n AND EXISTS (SELECT 1 FROM completed)\n AND NOT EXISTS (SELECT 1 FROM deleted WHERE canceled_by IS NOT NULL)\n )\n SELECT c.duration_ms AS \"duration_ms!\",\n EXISTS (SELECT 1 FROM deleted WHERE canceled_by IS NOT NULL) AS \"carried_cancel!\"\n FROM completed c",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "duration_ms!",
"type_info": "Int8"
},
{
"ordinal": 1,
"name": "carried_cancel!",
"type_info": "Bool"
}
],
"parameters": {
"Left": [
"Uuid",
"Bool",
"Text",
"Bool",
"Varchar",
"Text",
"Bool",
"Int4",
"Int8",
"TextArray",
"Uuid",
"Text"
]
},
"nullable": [
false,
null
]
},
"hash": "700cb5b65f46da7b3e34566a4fd4ea27005eb98f5d0d41999161082152df0f40"
}
@@ -0,0 +1,31 @@
{
"db_name": "PostgreSQL",
"query": "WITH deleted AS (\n DELETE FROM v2_job_queue WHERE id = $1\n RETURNING id, workspace_id, started_at, worker, canceled_by, canceled_reason\n ), completed AS (\n 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 , worker\n )\n SELECT d.workspace_id, d.id, d.started_at,\n COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(d.started_at, now()))))*1000),\n $3::text::jsonb, $10,\n CASE WHEN $4::BOOL THEN $5 ELSE d.canceled_by END,\n CASE WHEN $4::BOOL THEN $6 WHEN d.canceled_by IS NOT NULL THEN d.canceled_reason END,\n s.flow_status, s.workflow_as_code_status, $8,\n CASE WHEN $4::BOOL OR d.canceled_by IS NOT NULL 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,\n d.worker\n FROM deleted d LEFT JOIN v2_job_status s ON s.id = d.id\n ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3::text::jsonb,\n canceled_by = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL\n THEN EXCLUDED.canceled_by ELSE cj.canceled_by END,\n canceled_reason = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL\n THEN EXCLUDED.canceled_reason ELSE cj.canceled_reason END\n RETURNING duration_ms\n )\n SELECT duration_ms AS \"duration_ms!\" FROM completed",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "duration_ms!",
"type_info": "Int8"
}
],
"parameters": {
"Left": [
"Uuid",
"Bool",
"Text",
"Bool",
"Varchar",
"Text",
"Bool",
"Int4",
"Int8",
"TextArray"
]
},
"nullable": [
false
]
},
"hash": "9212902d4eb2834ab3b429968a7a85c53c49e4ef0f42c146324d5c71fbd1e21d"
}
@@ -1,16 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE v2_job_completed SET status = 'canceled'::job_status, canceled_by = $2, canceled_reason = $3 WHERE id = $1",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Uuid",
"Varchar",
"Text"
]
},
"nullable": []
},
"hash": "9a2345e35c9f44578aa7dde994fdde0be9e0d192e0aee224bff78d8bd2ccbd8b"
}
@@ -0,0 +1,27 @@
{
"db_name": "PostgreSQL",
"query": "\n WITH step_index AS (\n SELECT idx::text AS idx\n FROM v2_job_status,\n jsonb_array_elements(flow_status->'modules') WITH ORDINALITY arr(elem, idx)\n WHERE id = $1\n AND elem->>'id' = $5\n LIMIT 1\n ), target AS (\n -- The queue row is locked before the completed row's key, as every completion does.\n SELECT id FROM v2_job_queue WHERE id = $1 FOR UPDATE\n ), completed AS (\n INSERT INTO v2_job_completed\n (workspace_id, id, started_at, duration_ms, result,\n flow_status, workflow_as_code_status, status, worker)\n SELECT\n q.workspace_id, q.id, q.started_at,\n (EXTRACT('epoch' FROM now()) - EXTRACT('epoch' FROM COALESCE(q.started_at, now()))) * 1000,\n $3::text::jsonb,\n CASE WHEN si.idx IS NOT NULL\n THEN jsonb_set(\n s.flow_status,\n ARRAY['modules', (si.idx::int - 1)::text],\n $6::jsonb\n )\n ELSE s.flow_status\n END,\n s.workflow_as_code_status,\n 'skipped'::job_status,\n q.worker\n FROM v2_job_queue q\n LEFT JOIN v2_job_status s ON s.id = q.id\n LEFT JOIN step_index si ON true\n WHERE q.id IN (SELECT id FROM target)\n ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = EXCLUDED.result\n RETURNING 1 AS x\n ), _deleted AS (\n DELETE FROM v2_job_queue WHERE id IN (SELECT id FROM target)\n ), _logged AS (\n INSERT INTO job_logs (logs, job_id, workspace_id)\n VALUES ($4, $1, $2)\n ON CONFLICT (job_id) DO UPDATE SET logs = concat(job_logs.logs, EXCLUDED.logs)\n )\n SELECT x FROM completed\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "x",
"type_info": "Int4"
}
],
"parameters": {
"Left": [
"Uuid",
"Varchar",
"Text",
"Text",
"Text",
"Jsonb"
]
},
"nullable": [
null
]
},
"hash": "c4254f50fc6f74aae9e23b461a6c69c181eba0210cc068b64e67b0a518fe1cba"
}
@@ -0,0 +1,24 @@
{
"db_name": "PostgreSQL",
"query": "WITH deleted AS (\n DELETE FROM v2_job_queue WHERE id = $1 RETURNING id, workspace_id, started_at, worker\n )\n INSERT INTO v2_job_completed\n (workspace_id, id, started_at, duration_ms, result, memory_peak, status, worker)\n SELECT d.workspace_id, d.id, d.started_at, $3, $2::jsonb, r.memory_peak,\n 'failure'::job_status, d.worker\n FROM deleted d\n LEFT JOIN v2_job_runtime r ON r.id = d.id\n ON CONFLICT (id) DO UPDATE SET status = 'failure', result = $2::jsonb\n RETURNING id",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "id",
"type_info": "Uuid"
}
],
"parameters": {
"Left": [
"Uuid",
"Jsonb",
"Int8"
]
},
"nullable": [
false
]
},
"hash": "d9b7e45213aebced9cebe8ce1278d38136bc73c3517d0e0fbad369666504930c"
}
@@ -1,28 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM v2_job_queue WHERE id = $1 RETURNING canceled_by, canceled_reason",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "canceled_by",
"type_info": "Varchar"
},
{
"ordinal": 1,
"name": "canceled_reason",
"type_info": "Text"
}
],
"parameters": {
"Left": [
"Uuid"
]
},
"nullable": [
true,
true
]
},
"hash": "ed2eb32d445164e7ab5fb735e01f86505bc499320e08d1aad15a9a82c3823c0b"
}
+1 -1
View File
@@ -1 +1 @@
de73db2bacfdc3eaa2e63b1827178bc198d54e5c
7a256cf353db7cf64a60a09fa0de7f3a8b27f626
+50 -35
View File
@@ -6210,8 +6210,8 @@ async fn handle_zombie_jobs(db: &Pool<Postgres>, base_internal_url: &str, node_n
}
/// Force-complete a zombie job that handle_job_error failed to complete.
/// This is a minimal fallback: it inserts a failed completed job and deletes
/// from the queue in a single transaction, without schedule pushing or
/// This is a minimal fallback: it moves the job from the queue to a failed
/// completed job in a single transaction, without schedule pushing or
/// error handler logic. The one thing it keeps is the WAC parent notification,
/// deliberately inside the transaction: if that fails, the whole completion
/// rolls back and the job waits for the next sweep, which is cheaper than a
@@ -6248,52 +6248,67 @@ async fn force_complete_zombie_job(
let mut tx = db.begin().await?;
// Locks in the order of every other completion: a WAC parent's rows, then this job's queue
// row, then its completed row (see `record_child_completion`). A worker completing the same
// job concurrently would otherwise deadlock against this.
let duration_ms = sqlx::query_scalar!(
"INSERT INTO v2_job_completed
(workspace_id, id, started_at, duration_ms, result, memory_peak, status, worker)
SELECT q.workspace_id, q.id, q.started_at,
COALESCE((EXTRACT('epoch' FROM now()) - EXTRACT('epoch' FROM COALESCE(q.started_at, now()))) * 1000, 0)::bigint,
$2::jsonb, r.memory_peak, 'failure'::job_status, q.worker
FROM v2_job_queue q
LEFT JOIN v2_job_runtime r ON r.id = q.id
WHERE q.id = $1
ON CONFLICT (id) DO UPDATE SET status = 'failure', result = $2::jsonb
RETURNING duration_ms AS \"duration_ms!\"",
"SELECT COALESCE(c.duration_ms, COALESCE((EXTRACT('epoch' FROM now()) - EXTRACT('epoch' FROM COALESCE(q.started_at, now()))) * 1000, 0)::bigint) AS \"duration_ms!\"
FROM v2_job_queue q LEFT JOIN v2_job_completed c ON c.id = q.id WHERE q.id = $1",
job_id,
error_value,
)
.fetch_optional(&mut *tx)
.await?;
let Some(duration_ms) = duration_ms else {
return Ok(());
};
// A WAC parent parked on this job must learn of the failure here too, or it
// waits out its whole suspend window and runs the task again.
let mut wac_parent_ready = false;
if let Some(duration_ms) = duration_ms {
let parent = sqlx::query!(
"SELECT parent_job, flow_step_id FROM v2_job WHERE id = $1",
job_id
let parent = sqlx::query!(
"SELECT parent_job, flow_step_id FROM v2_job WHERE id = $1",
job_id
)
.fetch_optional(&mut *tx)
.await?;
if let Some(parent_job) = parent
.filter(|j| j.flow_step_id.is_none())
.and_then(|j| j.parent_job)
{
wac_parent_ready = windmill_common::wac::record_child_completion(
&mut tx,
&parent_job,
job_id,
false,
duration_ms,
&error_value.to_string(),
)
.fetch_optional(&mut *tx)
.await?;
if let Some(parent_job) = parent
.filter(|j| j.flow_step_id.is_none())
.and_then(|j| j.parent_job)
{
wac_parent_ready = windmill_common::wac::record_child_completion(
&mut tx,
&parent_job,
job_id,
false,
duration_ms,
&error_value.to_string(),
)
.await?;
}
}
sqlx::query!("DELETE FROM v2_job_queue WHERE id = $1", job_id)
.execute(&mut *tx)
.await?;
let completed = sqlx::query_scalar!(
"WITH deleted AS (
DELETE FROM v2_job_queue WHERE id = $1 RETURNING id, workspace_id, started_at, worker
)
INSERT INTO v2_job_completed
(workspace_id, id, started_at, duration_ms, result, memory_peak, status, worker)
SELECT d.workspace_id, d.id, d.started_at, $3, $2::jsonb, r.memory_peak,
'failure'::job_status, d.worker
FROM deleted d
LEFT JOIN v2_job_runtime r ON r.id = d.id
ON CONFLICT (id) DO UPDATE SET status = 'failure', result = $2::jsonb
RETURNING id",
job_id,
error_value,
duration_ms,
)
.fetch_optional(&mut *tx)
.await?;
// Completed by someone else while this waited on the WAC parent: roll back, so the parent
// keeps what the winning completion recorded.
if completed.is_none() {
return Ok(());
}
tx.commit().await?;
+195
View File
@@ -0,0 +1,195 @@
//! A flow step's completion is the progress that keeps its flow from being reaped as a zombie: it
//! must refresh the parent's ping, unless the parent is being canceled.
use serde_json::json;
use sqlx::{types::Json, Pool, Postgres};
use uuid::Uuid;
use windmill_queue::{add_completed_job, get_mini_completed_job};
const W_ID: &str = "test-workspace";
async fn insert_running(
db: &Pool<Postgres>,
kind: &str,
parent: Option<Uuid>,
step_id: Option<&str>,
) -> anyhow::Result<Uuid> {
let id = Uuid::new_v4();
sqlx::query(
"INSERT INTO v2_job (id, workspace_id, created_by, created_at, permissioned_as, \
permissioned_as_email, kind, runnable_path, script_lang, tag, visible_to_owner, \
parent_job, flow_step_id) \
VALUES ($1, $2, 'test-user', now(), 'u/test-user', 'test@windmill.dev', \
$3::job_kind, 'u/test-user/f', 'bash', 'bash', true, $4, $5)",
)
.bind(id)
.bind(W_ID)
.bind(kind)
.bind(parent)
.bind(step_id)
.execute(db)
.await?;
sqlx::query(
"INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, running, started_at, tag) \
VALUES ($1, $2, now(), true, now(), 'bash')",
)
.bind(id)
.bind(W_ID)
.execute(db)
.await?;
sqlx::query("INSERT INTO v2_job_runtime (id, ping) VALUES ($1, now() - interval '1 hour')")
.bind(id)
.execute(db)
.await?;
Ok(id)
}
async fn complete_step_and_read_parent_ping_age(
db: &Pool<Postgres>,
parent_canceled: bool,
) -> anyhow::Result<f64> {
let parent = insert_running(db, "flow", None, None).await?;
if parent_canceled {
sqlx::query("UPDATE v2_job_queue SET canceled_by = 'test-user' WHERE id = $1")
.bind(parent)
.execute(db)
.await?;
}
let step = insert_running(db, "script", Some(parent), Some("a")).await?;
let job = get_mini_completed_job(&step, W_ID, db).await?.unwrap();
add_completed_job(
db,
&job,
true,
false,
Json(&json!("done")),
None,
0,
None,
false,
None,
false,
)
.await?;
let (status, queued): (String, bool) = sqlx::query_as(
"SELECT status::text, EXISTS (SELECT 1 FROM v2_job_queue WHERE id = $1) \
FROM v2_job_completed WHERE id = $1",
)
.bind(step)
.fetch_one(db)
.await?;
assert_eq!(status, "success");
assert!(!queued, "the step left the queue");
Ok(sqlx::query_scalar(
"SELECT EXTRACT(EPOCH FROM now() - ping)::float8 FROM v2_job_runtime WHERE id = $1",
)
.bind(parent)
.fetch_one(db)
.await?)
}
#[sqlx::test(fixtures("base"))]
async fn a_step_completion_refreshes_its_running_flow_ping(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
let age = complete_step_and_read_parent_ping_age(&db, false).await?;
assert!(age < 60.0, "the parent's ping was refreshed: {age}s old");
let age = complete_step_and_read_parent_ping_age(&db, true).await?;
assert!(
age > 3000.0,
"a canceled parent's ping is left alone: {age}s old"
);
Ok(())
}
async fn ping_age(db: &Pool<Postgres>, id: Uuid) -> anyhow::Result<f64> {
Ok(sqlx::query_scalar(
"SELECT EXTRACT(EPOCH FROM now() - ping)::float8 FROM v2_job_runtime WHERE id = $1",
)
.bind(id)
.fetch_one(db)
.await?)
}
/// Waits for the completion to be stuck on the queue row of `id`.
async fn wait_until_blocked_on(db: &Pool<Postgres>, id: Uuid) -> anyhow::Result<()> {
for _ in 0..100 {
let blocked: bool = sqlx::query_scalar(
"SELECT EXISTS (SELECT 1 FROM pg_locks l JOIN v2_job_queue q \
ON q.ctid = ('(' || l.page || ',' || l.tuple || ')')::tid \
WHERE NOT l.granted AND l.locktype = 'tuple' AND q.id = $1 \
AND l.database = (SELECT oid FROM pg_database WHERE datname = current_database())) \
OR EXISTS (SELECT 1 FROM pg_stat_activity \
WHERE datname = current_database() AND wait_event_type = 'Lock' \
AND query LIKE '%DELETE FROM v2_job_queue WHERE id = $1%')",
)
.bind(id)
.fetch_one(db)
.await?;
if blocked {
return Ok(());
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
anyhow::bail!("the completion never blocked on the step's queue row")
}
/// A cancel of the flow marks the flow and then its steps in one transaction. A step completing
/// while that is uncommitted waits on its own row, and must then leave the canceled flow's ping
/// alone although it began before the cancel was visible.
#[sqlx::test(fixtures("base"))]
async fn a_flow_cancel_landing_during_a_step_completion_keeps_the_flow_ping(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
let parent = insert_running(&db, "flow", None, None).await?;
let step = insert_running(&db, "script", Some(parent), Some("a")).await?;
let job = get_mini_completed_job(&step, W_ID, &db).await?.unwrap();
let mut cancel = db.begin().await?;
sqlx::query(
"UPDATE v2_job_queue SET canceled_by = 'test-user', canceled_reason = 'stop' \
WHERE id = ANY($1)",
)
.bind(vec![parent, step])
.execute(&mut *cancel)
.await?;
let completing = tokio::spawn({
let db = db.clone();
async move {
add_completed_job(
&db,
&job,
true,
false,
Json(&json!("done")),
None,
0,
None,
false,
None,
false,
)
.await
}
});
wait_until_blocked_on(&db, step).await?;
cancel.commit().await?;
completing.await??;
let status: String =
sqlx::query_scalar("SELECT status::text FROM v2_job_completed WHERE id = $1")
.bind(step)
.fetch_one(&db)
.await?;
assert_eq!(status, "canceled");
let age = ping_age(&db, parent).await?;
assert!(
age > 3000.0,
"the canceled flow's ping is left alone: {age}s old"
);
Ok(())
}
+302 -174
View File
@@ -1177,74 +1177,76 @@ async fn commit_completed_job<T: Serialize + Send + Sync + ValidableJson>(
let serialized_result = result.serialized_json();
let sanitized_result = strip_json_nul(serialized_result.as_ref());
let labels = result.wm_labels();
let is_scheduled =
completed_job.schedule_path().is_some() && completed_job.runnable_path.is_some();
let wac_parent = (!completed_job.is_flow_step())
.then_some(completed_job.parent_job)
.flatten();
let monitor_parent = (completed_job.is_flow_step() && flow_is_done)
.then_some(completed_job.parent_job)
.flatten();
let completion = Completion {
completed_job,
success,
skipped,
result: sanitized_result.as_ref(),
result_columns,
mem_peak,
canceled_by,
duration,
};
if labels.is_none()
&& !has_concurrent_limit
&& wac_parent.is_none()
&& monitor_parent.is_none()
&& !is_scheduled
{
let Some(duration) = completion.execute(&mut *db.acquire().await?).await? else {
return Err(not_in_queue_error(db, job_id).await);
};
log_completed_job(completed_job, duration, success);
return Ok((None, duration, false, false));
}
let mut tx = db.begin().warn_after_seconds(10).await?;
let duration = sqlx::query_scalar!(
"INSERT INTO v2_job_completed AS cj
( workspace_id
, id
, started_at
, duration_ms
, result
, result_columns
, canceled_by
, canceled_reason
, flow_status
, workflow_as_code_status
, memory_peak
, status
, worker
)
SELECT q.workspace_id, q.id, started_at, COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(started_at, now()))))*1000), $3::text::jsonb, $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,
q.worker
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::text::jsonb RETURNING duration_ms AS \"duration_ms!\"",
/* $1 */ completed_job.id,
/* $2 */ success,
/* $3 */ sanitized_result.as_ref(),
/* $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<String>>,
// The parent's rows are locked ahead of the child's own queue row (see
// `record_child_completion` for the order this must keep), so the duration it stamps is read
// before the completion: the one the completed row will hold. `now()` is fixed for the
// transaction, and a completed row already there keeps its own duration.
let mut wac_parent_ready = false;
if let Some(parent_job) = wac_parent {
let Some(duration) = sqlx::query_scalar!(
"SELECT COALESCE(c.duration_ms, COALESCE($2::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(q.started_at, now()))))*1000)::bigint) AS \"duration_ms!\"
FROM v2_job_queue q LEFT JOIN v2_job_completed c ON c.id = q.id WHERE q.id = $1",
job_id,
duration,
)
.fetch_optional(&mut *tx)
.warn_after_seconds(10)
.await
.map_err(|e| Error::internal_err(format!("Could not add completed job {job_id}: {e:#}")))?;
let duration = if let Some(duration) = duration {
duration
} else {
let already_inserted = sqlx::query_scalar!(
"SELECT EXISTS(SELECT 1 FROM v2_job_completed WHERE id = $1)",
job_id
.await?
else {
return Err(not_in_queue_error(&mut *tx, job_id).await);
};
wac_parent_ready = windmill_common::wac::record_child_completion(
&mut tx,
&parent_job,
&completed_job.id,
success,
duration,
sanitized_result.as_ref(),
)
.fetch_one(&mut *tx)
.warn_after_seconds(10)
.await
.map_err(|e| Error::internal_err(format!("Could not add completed job {job_id}: {e:#}")))?
.unwrap_or(false);
.await?;
}
if already_inserted {
return Err(Error::AlreadyCompleted(format!(
"The queued job {job_id} is already completed."
)));
} else {
return Err(Error::AlreadyCompleted(format!(
"There is no queued job anymore for {job_id} but there is no completed job either."
)));
}
let Some(duration) = completion.execute(&mut *tx).await? else {
return Err(not_in_queue_error(&mut *tx, job_id).await);
};
if let Some(mut labels) = result.wm_labels() {
if let Some(mut labels) = labels {
// A `\u0000` inside a wm_labels entry decodes to a real NUL that the
// `text[]` column rejects, which would abort this same transaction (and
// roll back the sanitized result insert) exactly like an unsanitized
@@ -1266,78 +1268,19 @@ async fn commit_completed_job<T: Serialize + Send + Sync + ValidableJson>(
.map_err(|e| Error::InternalErr(format!("Could not update job labels: {e:#}")))?;
}
// Before `delete_job`: the parent's rows are locked ahead of the child's own
// queue row (see `record_child_completion` for the order this must keep).
let mut wac_parent_ready = false;
if !completed_job.is_flow_step() {
if let Some(parent_job) = completed_job.parent_job {
wac_parent_ready = windmill_common::wac::record_child_completion(
&mut tx,
&parent_job,
&completed_job.id,
success,
duration,
sanitized_result.as_ref(),
)
.warn_after_seconds(10)
.await?;
}
}
let mut _skip_downstream_error_handlers = false;
let (ntx, canceled_at_delete) = delete_job(tx, &job_id).warn_after_seconds(10).await?;
tx = ntx;
// `canceled_by` is only what the worker last read from the queue row. Deleting the row waits
// for a cancel still being written, so the row it removed is the final word on whether this
// job was canceled.
if canceled_by.is_none() {
if let Some(canceled) = canceled_at_delete {
sqlx::query!(
"UPDATE v2_job_completed SET status = 'canceled'::job_status, canceled_by = $2, \
canceled_reason = $3 WHERE id = $1",
job_id,
canceled.username,
canceled.reason,
)
.execute(&mut *tx)
.warn_after_seconds(10)
.await?;
}
}
// tracing::error!("3 {:?}", start.elapsed());
if completed_job.is_flow_step() {
if let Some(parent_job) = completed_job.parent_job {
// persist the flow last progress timestamp to avoid zombie flow jobs
tracing::debug!(
"Persisting flow last progress timestamp to flow job: {:?}",
parent_job
);
sqlx::query!(
"UPDATE v2_job_runtime r SET
ping = now()
FROM v2_job_queue q
WHERE r.id = $1 AND q.id = r.id
AND q.workspace_id = $2
AND canceled_by IS NULL",
if let Some(parent_job) = monitor_parent {
let r = sqlx::query_scalar!(
"UPDATE parallel_monitor_lock SET last_ping = now() WHERE parent_flow_id = $1 and job_id = $2 RETURNING 1",
parent_job,
&completed_job.workspace_id
)
.execute(&mut *tx)
.warn_after_seconds(10)
.await?;
if flow_is_done {
let r = sqlx::query_scalar!(
"UPDATE parallel_monitor_lock SET last_ping = now() WHERE parent_flow_id = $1 and job_id = $2 RETURNING 1",
parent_job,
&completed_job.id
).fetch_optional(&mut *tx).warn_after_seconds(10).await?;
if r.is_some() {
tracing::info!(
"parallel flow iteration is done, setting parallel monitor last ping lock for job {}",
&completed_job.id
);
}
&completed_job.id
).fetch_optional(&mut *tx).warn_after_seconds(10).await?;
if r.is_some() {
tracing::info!(
"parallel flow iteration is done, setting parallel monitor last ping lock for job {}",
&completed_job.id
);
}
}
} else {
@@ -1509,6 +1452,238 @@ async fn commit_completed_job<T: Serialize + Send + Sync + ValidableJson>(
tx.commit().warn_after_seconds(10).await?;
log_completed_job(completed_job, duration, success);
// tracing::info!("completed job: {:?}", start.elapsed().as_micros());
Ok((
None,
duration,
_skip_downstream_error_handlers,
wac_parent_ready,
))
}
/// What a job's completion writes, for `Completion::execute`.
struct Completion<'a> {
completed_job: &'a MiniCompletedJob,
success: bool,
skipped: bool,
result: &'a str,
result_columns: Option<&'a Vec<String>>,
mem_peak: i32,
canceled_by: &'a Option<CanceledBy>,
duration: Option<i64>,
}
impl Completion<'_> {
/// Moves the job from the queue to the completed jobs and refreshes a flow step's parent
/// ping, as one statement. Returns `None` when the job was no longer in the queue.
///
/// The completion takes its cancellation from the queue row it deletes, not only from
/// `canceled_by`: that is what the worker last read, and the delete waits for a cancel still
/// being written, so the deleted row is the final word on whether the job was canceled.
///
/// It locks the queue row before the completed row's key. Any other writer completing a job
/// (the monitor's zombie fallback, debounce) must take them in the same order, or the two
/// deadlock.
async fn execute(&self, conn: &mut sqlx::PgConnection) -> error::Result<Option<i64>> {
let Completion {
completed_job,
success,
skipped,
result,
result_columns,
mem_peak,
canceled_by,
duration,
} = *self;
#[cfg(feature = "prometheus")]
if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) {
QUEUE_DELETE_COUNT.inc();
}
otel_incr_queue_delete_count();
let err = |e: sqlx::Error| {
Error::internal_err(format!(
"Could not add completed job {}: {e:#}",
completed_job.id
))
};
// A step's completion is progress of its flow, and keeps the flow from being reaped as a
// zombie. Any other completion runs the statement without the ping: Postgres sets up every
// write of a plan, so an unused ping would cost about a tenth of the completion. The two
// statements differ only by the ping; a change to the delete or the insert goes in both.
let Some(parent_to_ping) = completed_job
.is_flow_step()
.then_some(completed_job.parent_job)
.flatten()
else {
return sqlx::query_scalar!(
"WITH deleted AS (
DELETE FROM v2_job_queue WHERE id = $1
RETURNING id, workspace_id, started_at, worker, canceled_by, canceled_reason
), completed AS (
INSERT INTO v2_job_completed AS cj
( workspace_id
, id
, started_at
, duration_ms
, result
, result_columns
, canceled_by
, canceled_reason
, flow_status
, workflow_as_code_status
, memory_peak
, status
, worker
)
SELECT d.workspace_id, d.id, d.started_at,
COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(d.started_at, now()))))*1000),
$3::text::jsonb, $10,
CASE WHEN $4::BOOL THEN $5 ELSE d.canceled_by END,
CASE WHEN $4::BOOL THEN $6 WHEN d.canceled_by IS NOT NULL THEN d.canceled_reason END,
s.flow_status, s.workflow_as_code_status, $8,
CASE WHEN $4::BOOL OR d.canceled_by IS NOT NULL THEN 'canceled'::job_status
WHEN $7::BOOL THEN 'skipped'::job_status
WHEN $2::BOOL THEN 'success'::job_status
ELSE 'failure'::job_status END,
d.worker
FROM deleted d LEFT JOIN v2_job_status s ON s.id = d.id
ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3::text::jsonb,
canceled_by = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL
THEN EXCLUDED.canceled_by ELSE cj.canceled_by END,
canceled_reason = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL
THEN EXCLUDED.canceled_reason ELSE cj.canceled_reason END
RETURNING duration_ms
)
SELECT duration_ms AS \"duration_ms!\" FROM completed",
/* $1 */ completed_job.id,
/* $2 */ success,
/* $3 */ result,
/* $4 */ canceled_by.is_some(),
/* $5 */ canceled_by.as_ref().and_then(|cb| cb.username.as_deref()),
/* $6 */ canceled_by.as_ref().and_then(|cb| cb.reason.as_deref()),
/* $7 */ skipped,
/* $8 */ if mem_peak > 0 { Some(mem_peak) } else { None },
/* $9 */ duration,
/* $10 */ result_columns as Option<&Vec<String>>,
)
.fetch_optional(&mut *conn)
.warn_after_seconds(10)
.await
.map_err(err);
};
// A cancel of the flow marks the flow and then its steps. When the delete waited on it,
// the parent row read by this statement's snapshot predates it, so the ping moves to a
// statement of its own that sees the cancel.
let completed = sqlx::query!(
"WITH deleted AS (
DELETE FROM v2_job_queue WHERE id = $1
RETURNING id, workspace_id, started_at, worker, canceled_by, canceled_reason
), completed AS (
INSERT INTO v2_job_completed AS cj
( workspace_id
, id
, started_at
, duration_ms
, result
, result_columns
, canceled_by
, canceled_reason
, flow_status
, workflow_as_code_status
, memory_peak
, status
, worker
)
SELECT d.workspace_id, d.id, d.started_at,
COALESCE($9::bigint, (EXTRACT('epoch' FROM (now())) - EXTRACT('epoch' FROM (COALESCE(d.started_at, now()))))*1000),
$3::text::jsonb, $10,
CASE WHEN $4::BOOL THEN $5 ELSE d.canceled_by END,
CASE WHEN $4::BOOL THEN $6 WHEN d.canceled_by IS NOT NULL THEN d.canceled_reason END,
s.flow_status, s.workflow_as_code_status, $8,
CASE WHEN $4::BOOL OR d.canceled_by IS NOT NULL THEN 'canceled'::job_status
WHEN $7::BOOL THEN 'skipped'::job_status
WHEN $2::BOOL THEN 'success'::job_status
ELSE 'failure'::job_status END,
d.worker
FROM deleted d LEFT JOIN v2_job_status s ON s.id = d.id
ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, result = $3::text::jsonb,
canceled_by = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL
THEN EXCLUDED.canceled_by ELSE cj.canceled_by END,
canceled_reason = CASE WHEN NOT $4::BOOL AND EXCLUDED.canceled_by IS NOT NULL
THEN EXCLUDED.canceled_reason ELSE cj.canceled_reason END
RETURNING duration_ms
), parent_ping AS (
UPDATE v2_job_runtime r SET ping = now()
FROM v2_job_queue q
WHERE r.id = $11 AND q.id = r.id AND q.workspace_id = $12 AND q.canceled_by IS NULL
AND EXISTS (SELECT 1 FROM completed)
AND NOT EXISTS (SELECT 1 FROM deleted WHERE canceled_by IS NOT NULL)
)
SELECT c.duration_ms AS \"duration_ms!\",
EXISTS (SELECT 1 FROM deleted WHERE canceled_by IS NOT NULL) AS \"carried_cancel!\"
FROM completed c",
/* $1 */ completed_job.id,
/* $2 */ success,
/* $3 */ result,
/* $4 */ canceled_by.is_some(),
/* $5 */ canceled_by.as_ref().and_then(|cb| cb.username.as_deref()),
/* $6 */ canceled_by.as_ref().and_then(|cb| cb.reason.as_deref()),
/* $7 */ skipped,
/* $8 */ if mem_peak > 0 { Some(mem_peak) } else { None },
/* $9 */ duration,
/* $10 */ result_columns as Option<&Vec<String>>,
/* $11 */ parent_to_ping,
/* $12 */ &completed_job.workspace_id,
)
.fetch_optional(&mut *conn)
.warn_after_seconds(10)
.await
.map_err(err)?;
let Some(completed) = completed else {
return Ok(None);
};
if completed.carried_cancel {
sqlx::query!(
"UPDATE v2_job_runtime r SET
ping = now()
FROM v2_job_queue q
WHERE r.id = $1 AND q.id = r.id
AND q.workspace_id = $2
AND canceled_by IS NULL",
parent_to_ping,
&completed_job.workspace_id
)
.execute(&mut *conn)
.warn_after_seconds(10)
.await?;
}
Ok(Some(completed.duration_ms))
}
}
/// The error for a completion that found no queue row to complete.
async fn not_in_queue_error<'e>(conn: impl PgExecutor<'e>, job_id: Uuid) -> Error {
let already_inserted = sqlx::query_scalar!(
"SELECT EXISTS(SELECT 1 FROM v2_job_completed WHERE id = $1)",
job_id
)
.fetch_one(conn)
.warn_after_seconds(10)
.await;
match already_inserted {
Err(e) => Error::internal_err(format!("Could not add completed job {job_id}: {e:#}")),
Ok(Some(true)) => {
Error::AlreadyCompleted(format!("The queued job {job_id} is already completed."))
}
Ok(_) => Error::AlreadyCompleted(format!(
"There is no queued job anymore for {job_id} but there is no completed job either."
)),
}
}
fn log_completed_job(completed_job: &MiniCompletedJob, duration: i64, success: bool) {
let job_id = completed_job.id;
tracing::info!(
%job_id,
root_job = ?completed_job.flow_innermost_root_job.map(|x| x.to_string()).unwrap_or_else(|| String::new()),
@@ -1527,13 +1702,6 @@ async fn commit_completed_job<T: Serialize + Send + Sync + ValidableJson>(
"inserted completed job: {} (success: {success})",
completed_job.id
);
// tracing::info!("completed job: {:?}", start.elapsed().as_micros());
Ok((
None,
duration,
_skip_downstream_error_handlers,
wac_parent_ready,
))
}
async fn check_result_size<T: ValidableJson>(
@@ -5190,46 +5358,6 @@ async fn extract_result_from_job_result(
}
}
/// Also reports the cancellation the deleted row carried, if any. Unlike a plain read of the queue
/// row, this waits for a cancel that is still being written, so it is the last word on one.
pub async fn delete_job<'c>(
mut tx: Transaction<'c, Postgres>,
job_id: &Uuid,
) -> windmill_common::error::Result<(Transaction<'c, Postgres>, Option<CanceledBy>)> {
#[cfg(feature = "prometheus")]
if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) {
QUEUE_DELETE_COUNT.inc();
}
otel_incr_queue_delete_count();
let job_removed = sqlx::query!(
"DELETE FROM v2_job_queue WHERE id = $1 RETURNING canceled_by, canceled_reason",
job_id,
)
.fetch_optional(&mut *tx)
.await;
let canceled = match &job_removed {
Err(job_removed) => {
tracing::error!(
"Job {job_id} could not be deleted: {job_removed}. This is not necessarily an error, as the job might have been deleted by another process such as in the case of cancelling"
);
None
}
Ok(None) => {
tracing::error!("Job {job_id} could not be deleted, no row was removed. This is not necessarily an error, as the job might have been deleted by another process such as in the case of cancelling");
None
}
Ok(Some(row)) => row.canceled_by.as_ref().map(|username| CanceledBy {
username: Some(username.clone()),
reason: row.canceled_reason.clone(),
}),
};
tracing::debug!("Job {job_id} deleted");
Ok((tx, canceled))
}
pub async fn job_is_complete(db: &DB, id: Uuid, w_id: &str) -> error::Result<bool> {
Ok(sqlx::query_scalar!(
"SELECT EXISTS(SELECT 1 FROM v2_job_completed WHERE id = $1 AND workspace_id = $2)",