diff --git a/backend/.sqlx/query-c0977384d5400e41f1d2aab6dea42c202bad922a6829c4768e5cf1fb6d90f810.json b/backend/.sqlx/query-c0977384d5400e41f1d2aab6dea42c202bad922a6829c4768e5cf1fb6d90f810.json new file mode 100644 index 0000000000..a75bb2f8fc --- /dev/null +++ b/backend/.sqlx/query-c0977384d5400e41f1d2aab6dea42c202bad922a6829c4768e5cf1fb6d90f810.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE queue\n SET running = false\n , started_at = null\n , scheduled_for = $1\n WHERE id = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Timestamptz", + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "c0977384d5400e41f1d2aab6dea42c202bad922a6829c4768e5cf1fb6d90f810" +} diff --git a/backend/.sqlx/query-a0ce703def7e976947513029874fb571893c75730e0a8feeb609852423404cf0.json b/backend/.sqlx/query-effb8e396c5811d3c45f880a6bf33ca564d3109901c93d4f312397e21656ce27.json similarity index 58% rename from backend/.sqlx/query-a0ce703def7e976947513029874fb571893c75730e0a8feeb609852423404cf0.json rename to backend/.sqlx/query-effb8e396c5811d3c45f880a6bf33ca564d3109901c93d4f312397e21656ce27.json index e249554fbc..f5aa2759fd 100644 --- a/backend/.sqlx/query-a0ce703def7e976947513029874fb571893c75730e0a8feeb609852423404cf0.json +++ b/backend/.sqlx/query-effb8e396c5811d3c45f880a6bf33ca564d3109901c93d4f312397e21656ce27.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT CAST(ROUND(AVG(duration_ms) / 1000, 0) AS BIGINT) AS avg_duration_s FROM\n (SELECT duration_ms FROM completed_job WHERE script_path = $1\n ORDER BY started_at\n DESC LIMIT 10) AS t", + "query": "SELECT CAST(ROUND(AVG(duration_ms) / 1000, 0) AS BIGINT) AS avg_duration_s FROM\n (SELECT duration_ms FROM concurrency_key LEFT JOIN completed_job ON completed_job.id = concurrency_key.job_id WHERE key = $1 AND ended_at IS NOT NULL\n ORDER BY ended_at\n DESC LIMIT 10) AS t", "describe": { "columns": [ { @@ -18,5 +18,5 @@ null ] }, - "hash": "a0ce703def7e976947513029874fb571893c75730e0a8feeb609852423404cf0" + "hash": "effb8e396c5811d3c45f880a6bf33ca564d3109901c93d4f312397e21656ce27" } diff --git a/backend/windmill-api/src/flows.rs b/backend/windmill-api/src/flows.rs index 74b199a485..8950dea80e 100644 --- a/backend/windmill-api/src/flows.rs +++ b/backend/windmill-api/src/flows.rs @@ -504,7 +504,7 @@ async fn update_flow( nf.visible_to_runner_only.unwrap_or(false), ) .execute(&mut tx) - .await?; + .await.map_err(|e| error::Error::InternalErr(format!("Error updating flow due to flow update: {e}")))?; if nf.path != flow_path { check_schedule_conflict(tx.transaction_mut(), &w_id, &nf.path).await?; @@ -520,7 +520,7 @@ async fn update_flow( .bind(&flow_path) .bind(&w_id) .fetch_all(&mut tx) - .await?; + .await.map_err(|e| error::Error::InternalErr(format!("Error updating flow due to related schedules update: {e}")))?; let schedule = sqlx::query_as::<_, Schedule>( "UPDATE schedule SET path = $1, script_path = $1 WHERE path = $2 AND workspace_id = $3 AND is_flow IS true RETURNING *") @@ -528,7 +528,7 @@ async fn update_flow( .bind(&flow_path) .bind(&w_id) .fetch_optional(&mut tx) - .await?; + .await.map_err(|e| error::Error::InternalErr(format!("Error updating flow due to related schedule update: {e}")))?; if let Some(schedule) = schedule { clear_schedule(tx.transaction_mut(), &flow_path, &w_id).await?; @@ -618,14 +618,24 @@ async fn update_flow( w_id ) .execute(&mut new_tx) - .await?; + .await + .map_err(|e| { + error::Error::InternalErr(format!( + "Error updating flow due to updating dependency job field: {e}" + )) + })?; if let Some(old_dep_job) = old_dep_job { sqlx::query!( "UPDATE queue SET canceled = true WHERE id = $1", old_dep_job ) .execute(&mut new_tx) - .await?; + .await + .map_err(|e| { + error::Error::InternalErr(format!( + "Error updating flow due to cancelling dependency job: {e}" + )) + })?; } new_tx.commit().await?; diff --git a/backend/windmill-api/src/schedule.rs b/backend/windmill-api/src/schedule.rs index b869cc352b..a410fb9d61 100644 --- a/backend/windmill-api/src/schedule.rs +++ b/backend/windmill-api/src/schedule.rs @@ -793,6 +793,7 @@ pub async fn clear_schedule<'c>( path: &str, w_id: &str, ) -> Result<()> { + tracing::info!("Clearing schedule {}", path); sqlx::query!( "DELETE FROM queue WHERE schedule_path = $1 AND running = false AND workspace_id = $2 AND is_flow_step = false", path, diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 8be2763840..04c1c79655 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -1773,13 +1773,14 @@ pub async fn pull( let job_uuid: Uuid = pulled_job.id; let avg_script_duration: Option = sqlx::query_scalar!( "SELECT CAST(ROUND(AVG(duration_ms) / 1000, 0) AS BIGINT) AS avg_duration_s FROM - (SELECT duration_ms FROM completed_job WHERE script_path = $1 - ORDER BY started_at + (SELECT duration_ms FROM concurrency_key LEFT JOIN completed_job ON completed_job.id = concurrency_key.job_id WHERE key = $1 AND ended_at IS NOT NULL + ORDER BY ended_at DESC LIMIT 10) AS t", - job_script_path + job_concurrency_key ) .fetch_one(&mut tx) .await?; + tracing::info!("avg script duration computed: {:?}", avg_script_duration); // optimal scheduling is: 'older_job_in_concurrency_time_window_started_timestamp + script_avg_duration + concurrency_time_window_s' let estimated_next_schedule_timestamp = ((min_started_at @@ -1827,13 +1828,15 @@ pub async fn pull( tx.commit().await?; } else { // if using posgtres, then we're able to re-queue the entire batch of scheduled job for this script_path, so we do it - sqlx::query(&format!( + sqlx::query!( "UPDATE queue SET running = false , started_at = null - , scheduled_for = '{estimated_next_schedule_timestamp}' - WHERE (id = '{job_uuid}') OR (script_path = '{job_script_path}' AND running = false AND scheduled_for <= now())" - )) + , scheduled_for = $1 + WHERE id = $2", + estimated_next_schedule_timestamp, + job_uuid, + ) .fetch_all(&mut tx) .await .map_err(|e| Error::InternalErr(format!("Could not update and re-queue job {job_uuid}. The job will be marked as running but it is not running: {e}")))?;