mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-07 00:01:49 +00:00
fix(backend): reschedule flow at first step end (#746)
* fix(backend): reschedule flow at first step end * fix * fix * fix * fix * sqlx prepare
This commit is contained in:
@@ -2306,6 +2306,28 @@
|
||||
},
|
||||
"query": "INSERT INTO workspace_settings\n (workspace_id, slack_team_id, slack_name)\n VALUES ($1, $2, $3) ON CONFLICT (workspace_id) DO UPDATE SET slack_team_id = $2, slack_name = $3"
|
||||
},
|
||||
"a241c56415759105ccbcbf7fff77287fa4ec2cc096c0060d14db421115d63e2d": {
|
||||
"describe": {
|
||||
"columns": [
|
||||
{
|
||||
"name": "exists",
|
||||
"ordinal": 0,
|
||||
"type_info": "Bool"
|
||||
}
|
||||
],
|
||||
"nullable": [
|
||||
null
|
||||
],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Text",
|
||||
"Text",
|
||||
"Timestamptz"
|
||||
]
|
||||
}
|
||||
},
|
||||
"query": "SELECT EXISTS (SELECT 1 FROM queue WHERE workspace_id = $1 AND schedule_path = $2 AND scheduled_for = $3)"
|
||||
},
|
||||
"a34066d4a1578a13b2e322e6936ae80a0239a79148f3edce65f51a93910a1a4b": {
|
||||
"describe": {
|
||||
"columns": [
|
||||
|
||||
+16
-1
@@ -25,7 +25,7 @@ use axum::{
|
||||
use chrono::{DateTime, Duration, FixedOffset};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use serde_json::{Map, Value};
|
||||
use sqlx::{FromRow, Postgres, Transaction};
|
||||
use sqlx::{query_scalar, FromRow, Postgres, Transaction};
|
||||
|
||||
pub fn workspaced_service() -> Router {
|
||||
Router::new()
|
||||
@@ -83,6 +83,20 @@ pub async fn push_scheduled_job<'c>(
|
||||
.expect("a schedule should have a next event")
|
||||
+ offset;
|
||||
|
||||
let already_exists: bool = query_scalar!(
|
||||
"SELECT EXISTS (SELECT 1 FROM queue WHERE workspace_id = $1 AND schedule_path = $2 AND scheduled_for = $3)",
|
||||
&schedule.workspace_id,
|
||||
&schedule.path,
|
||||
next
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?
|
||||
.unwrap_or(false);
|
||||
|
||||
if already_exists {
|
||||
return Ok(tx);
|
||||
}
|
||||
|
||||
let mut args: Option<Map<String, Value>> = None;
|
||||
|
||||
if let Some(args_v) = schedule.args {
|
||||
@@ -108,6 +122,7 @@ pub async fn push_scheduled_job<'c>(
|
||||
path: schedule.script_path,
|
||||
}
|
||||
};
|
||||
|
||||
let (_, tx) = push(
|
||||
tx,
|
||||
&schedule.workspace_id,
|
||||
|
||||
+15
-29
@@ -6,8 +6,8 @@ use crate::{
|
||||
error::{self, Error},
|
||||
flows::{FlowModule, FlowModuleValue, FlowValue, InputTransform, Retry, Suspend},
|
||||
jobs::{
|
||||
add_completed_job, add_completed_job_error, get_queued_job, postprocess_queued_job, push,
|
||||
script_path_to_payload, JobPayload, QueuedJob, RawCode,
|
||||
add_completed_job, add_completed_job_error, delete_job, get_queued_job, push,
|
||||
schedule_again_if_scheduled, script_path_to_payload, JobPayload, QueuedJob, RawCode,
|
||||
},
|
||||
js_eval::{eval_timeout, EvalCreds, IdContext},
|
||||
more_serde::is_default,
|
||||
@@ -374,6 +374,16 @@ pub async fn update_flow_status_after_job_completion(
|
||||
|
||||
tx.commit().await?;
|
||||
|
||||
if old_status.step == 0 && !flow_job.is_flow_step {
|
||||
schedule_again_if_scheduled(
|
||||
flow_job.schedule_path.clone(),
|
||||
flow_job.script_path.clone(),
|
||||
&w_id,
|
||||
db,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
|
||||
let done = if !should_continue_flow {
|
||||
let logs = if flow_job.canceled {
|
||||
"Flow job canceled".to_string()
|
||||
@@ -421,15 +431,7 @@ pub async fn update_flow_status_after_job_completion(
|
||||
};
|
||||
|
||||
if done {
|
||||
postprocess_queued_job(
|
||||
flow_job.is_flow_step,
|
||||
flow_job.schedule_path.clone(),
|
||||
flow_job.script_path.clone(),
|
||||
&w_id,
|
||||
flow,
|
||||
db,
|
||||
)
|
||||
.await?;
|
||||
let _ = delete_job(db, w_id, flow).await?;
|
||||
|
||||
if flow_job.same_worker && !keep_job_dir {
|
||||
let _ = tokio::fs::remove_dir_all(format!("{worker_dir}/{}", flow_job.id)).await;
|
||||
@@ -883,15 +885,7 @@ async fn push_next_flow_job(
|
||||
let result = is_cancelled.unwrap_or(json!({ "error": logs }));
|
||||
let _uuid =
|
||||
add_completed_job(db, &flow_job, success, skipped, result, logs).await?;
|
||||
postprocess_queued_job(
|
||||
false,
|
||||
flow_job.schedule_path.clone(),
|
||||
flow_job.script_path.clone(),
|
||||
&flow_job.workspace_id,
|
||||
flow_job.id,
|
||||
db,
|
||||
)
|
||||
.await?;
|
||||
let _ = delete_job(db, &flow_job.workspace_id, flow_job.id).await?;
|
||||
|
||||
return Ok(());
|
||||
}
|
||||
@@ -1236,15 +1230,7 @@ async fn jump_to_next_step(
|
||||
let skipped = false;
|
||||
let logs = "Forloop completed without iteration".to_string();
|
||||
let _uuid = add_completed_job(db, &new_job, success, skipped, json!([]), logs).await?;
|
||||
postprocess_queued_job(
|
||||
false,
|
||||
new_job.schedule_path,
|
||||
new_job.script_path,
|
||||
&new_job.workspace_id,
|
||||
new_job.id,
|
||||
db,
|
||||
)
|
||||
.await?;
|
||||
let _ = delete_job(db, &new_job.workspace_id, new_job.id).await?;
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user