mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-08 16:02:30 +00:00
fix: disapproval does not trigger flow error handler anymore
This commit is contained in:
+4
-3
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "\n INSERT INTO resume_job\n (id, resume_id, job, flow, value, approver)\n VALUES ($1, $2, $3, $4, $5, $6)\n ",
|
||||
"query": "\n INSERT INTO resume_job\n (id, resume_id, job, flow, value, approver, approved)\n VALUES ($1, $2, $3, $4, $5, $6, $7)\n ",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
@@ -10,10 +10,11 @@
|
||||
"Uuid",
|
||||
"Uuid",
|
||||
"Jsonb",
|
||||
"Varchar"
|
||||
"Varchar",
|
||||
"Bool"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "c4f1c14c3aae145b52ff39bdecd779cd3eba27a869fda68516dc68dc2abefd38"
|
||||
"hash": "2c0c9312b8b326a3759566059d01c79efee28920a1a5afe2df92043526c1de82"
|
||||
}
|
||||
+14
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"db_name": "PostgreSQL",
|
||||
"query": "UPDATE queue SET suspend = 0 WHERE id = $1",
|
||||
"describe": {
|
||||
"columns": [],
|
||||
"parameters": {
|
||||
"Left": [
|
||||
"Uuid"
|
||||
]
|
||||
},
|
||||
"nullable": []
|
||||
},
|
||||
"hash": "4fb3a4712d88afed40082d8d8bd63b5dedad61caa68e0e470252083d80df605f"
|
||||
}
|
||||
@@ -0,0 +1,2 @@
|
||||
-- Add down migration script here
|
||||
ALTER TABLE resume_job DROP COLUMN approved;
|
||||
@@ -0,0 +1,2 @@
|
||||
-- Add up migration script here
|
||||
ALTER TABLE resume_job ADD COLUMN approved BOOLEAN NOT NULL DEFAULT true;
|
||||
@@ -1142,7 +1142,16 @@ pub async fn resume_suspended_flow_as_owner(
|
||||
&flow.script_path.clone().unwrap_or_else(|| String::new()),
|
||||
)?;
|
||||
|
||||
insert_resume_job(0, job_id, &flow, value, Some(authed.username), &mut tx).await?;
|
||||
insert_resume_job(
|
||||
0,
|
||||
job_id,
|
||||
&flow,
|
||||
value,
|
||||
Some(authed.username),
|
||||
true,
|
||||
&mut tx,
|
||||
)
|
||||
.await?;
|
||||
|
||||
resume_immediately_if_relevant(flow, job_id, &mut tx).await?;
|
||||
|
||||
@@ -1157,6 +1166,23 @@ pub async fn resume_suspended_job(
|
||||
Query(approver): Query<QueryApprover>,
|
||||
QueryOrBody(value): QueryOrBody<serde_json::Value>,
|
||||
) -> error::Result<StatusCode> {
|
||||
resume_suspended_job_internal(
|
||||
value, db, w_id, job_id, resume_id, approver, secret, authed, true,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn resume_suspended_job_internal(
|
||||
value: Option<serde_json::Value>,
|
||||
db: sqlx::Pool<Postgres>,
|
||||
w_id: String,
|
||||
job_id: Uuid,
|
||||
resume_id: u32,
|
||||
approver: QueryApprover,
|
||||
secret: String,
|
||||
authed: Option<ApiAuthed>,
|
||||
approved: bool,
|
||||
) -> Result<StatusCode, Error> {
|
||||
let value = value.unwrap_or(serde_json::Value::Null);
|
||||
let mut tx = db.begin().await?;
|
||||
let key = get_workspace_key(&w_id, &mut tx).await?;
|
||||
@@ -1209,13 +1235,32 @@ pub async fn resume_suspended_job(
|
||||
job_id,
|
||||
&parent_flow_info,
|
||||
value,
|
||||
approver,
|
||||
approver.clone(),
|
||||
approved,
|
||||
&mut tx,
|
||||
)
|
||||
.await?;
|
||||
|
||||
resume_immediately_if_relevant(parent_flow_info, job_id, &mut tx).await?;
|
||||
|
||||
if !approved {
|
||||
sqlx::query!(
|
||||
"UPDATE queue SET suspend = 0 WHERE id = $1",
|
||||
parent_flow_info.id
|
||||
)
|
||||
.execute(&mut *tx)
|
||||
.await?;
|
||||
} else {
|
||||
resume_immediately_if_relevant(parent_flow_info, job_id, &mut tx).await?;
|
||||
}
|
||||
audit_log(
|
||||
&mut *tx,
|
||||
&approver.unwrap_or_else(|| "anonymous".to_string()),
|
||||
"jobs.approved",
|
||||
ActionKind::Update,
|
||||
&w_id,
|
||||
Some(&job_id.to_string()),
|
||||
None,
|
||||
)
|
||||
.await?;
|
||||
tx.commit().await?;
|
||||
Ok(StatusCode::CREATED)
|
||||
}
|
||||
@@ -1259,20 +1304,22 @@ async fn insert_resume_job<'c>(
|
||||
flow: &FlowInfo,
|
||||
value: serde_json::Value,
|
||||
approver: Option<String>,
|
||||
approved: bool,
|
||||
tx: &mut Transaction<'c, Postgres>,
|
||||
) -> error::Result<()> {
|
||||
sqlx::query!(
|
||||
r#"
|
||||
INSERT INTO resume_job
|
||||
(id, resume_id, job, flow, value, approver)
|
||||
VALUES ($1, $2, $3, $4, $5, $6)
|
||||
(id, resume_id, job, flow, value, approver, approved)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7)
|
||||
"#,
|
||||
Uuid::from_u128(job_id.as_u128() ^ resume_id as u128),
|
||||
resume_id as i32,
|
||||
job_id,
|
||||
flow.id,
|
||||
value,
|
||||
approver
|
||||
approver,
|
||||
approved
|
||||
)
|
||||
.execute(&mut **tx)
|
||||
.await?;
|
||||
@@ -1342,64 +1389,14 @@ async fn get_suspended_flow_info<'c>(
|
||||
pub async fn cancel_suspended_job(
|
||||
authed: Option<ApiAuthed>,
|
||||
Extension(db): Extension<DB>,
|
||||
Extension(rsmq): Extension<Option<rsmq_async::MultiplexedRsmq>>,
|
||||
Path((w_id, job, resume_id, secret)): Path<(String, Uuid, u32, String)>,
|
||||
Path((w_id, job_id, resume_id, secret)): Path<(String, Uuid, u32, String)>,
|
||||
Query(approver): Query<QueryApprover>,
|
||||
) -> error::Result<String> {
|
||||
let mut tx = db.begin().await?;
|
||||
let key = get_workspace_key(&w_id, &mut tx).await?;
|
||||
let mut mac = HmacSha256::new_from_slice(key.as_bytes()).map_err(to_anyhow)?;
|
||||
mac.update(job.as_bytes());
|
||||
mac.update(resume_id.to_be_bytes().as_ref());
|
||||
if let Some(approver) = approver.approver.clone() {
|
||||
mac.update(approver.as_bytes());
|
||||
}
|
||||
mac.verify_slice(hex::decode(secret)?.as_ref())
|
||||
.map_err(|_| anyhow::anyhow!("Invalid signature"))?;
|
||||
|
||||
let whom = approver.approver.unwrap_or_else(|| "unknown".to_string());
|
||||
let parent_flow_id = get_suspended_parent_flow_info(job, &mut tx).await?.id;
|
||||
|
||||
let parent_flow = get_job_internal(&db, w_id.as_str(), parent_flow_id).await?;
|
||||
let flow_status = parent_flow
|
||||
.flow_status()
|
||||
.ok_or_else(|| anyhow::anyhow!("unable to find the flow status in the flow job"))?;
|
||||
let trigger_email = match &parent_flow {
|
||||
Job::CompletedJob(job) => &job.email,
|
||||
Job::QueuedJob(job) => &job.email,
|
||||
};
|
||||
conditionally_require_authed_user(authed, flow_status, trigger_email)?;
|
||||
|
||||
let (mut tx, cjob) = windmill_queue::cancel_job(
|
||||
&whom,
|
||||
Some("approval request disapproved".to_string()),
|
||||
parent_flow_id,
|
||||
&w_id,
|
||||
tx,
|
||||
&db,
|
||||
rsmq,
|
||||
false,
|
||||
QueryOrBody(value): QueryOrBody<serde_json::Value>,
|
||||
) -> error::Result<StatusCode> {
|
||||
resume_suspended_job_internal(
|
||||
value, db, w_id, job_id, resume_id, approver, secret, authed, false,
|
||||
)
|
||||
.await?;
|
||||
if cjob.is_some() {
|
||||
audit_log(
|
||||
&mut *tx,
|
||||
&whom,
|
||||
"jobs.disapproval",
|
||||
ActionKind::Delete,
|
||||
&w_id,
|
||||
Some(&parent_flow_id.to_string()),
|
||||
None,
|
||||
)
|
||||
.await?;
|
||||
tx.commit().await?;
|
||||
|
||||
Ok(format!("Flow {parent_flow_id} of job {job} cancelled"))
|
||||
} else {
|
||||
Ok(format!(
|
||||
"Flow {parent_flow_id} of job {job} was not cancellable"
|
||||
))
|
||||
}
|
||||
.await
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
|
||||
@@ -572,7 +572,7 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
let module = get_module(&flow_job, module_index);
|
||||
|
||||
// tracing::error!(
|
||||
// "UPDATE FLOW STATUS 3: {module:#?} {skip_failure} {is_last_step} {success}"
|
||||
// "UPDATE FLOW STATUS 3: {module:#?} {unrecoverable} {} {is_last_step} {success} {skip_error_handler}", flow_job.canceled
|
||||
// );
|
||||
let should_continue_flow = match success {
|
||||
_ if stop_early => false,
|
||||
@@ -1145,6 +1145,7 @@ pub struct MergeArgs<'a> {
|
||||
pub struct ResumeRow {
|
||||
pub value: Json<Box<RawValue>>,
|
||||
pub approver: Option<String>,
|
||||
pub approved: bool,
|
||||
pub resume_id: i32,
|
||||
}
|
||||
|
||||
@@ -1359,7 +1360,7 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
.context("lock flow in queue")?;
|
||||
|
||||
let resumes = sqlx::query(
|
||||
"SELECT value, approver, resume_id FROM resume_job WHERE job = $1 ORDER BY created_at ASC",
|
||||
"SELECT value, approver, resume_id, approved FROM resume_job WHERE job = $1 ORDER BY created_at ASC",
|
||||
)
|
||||
.bind(last)
|
||||
.fetch_all(&mut *tx)
|
||||
@@ -1443,7 +1444,10 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
.await?;
|
||||
}
|
||||
|
||||
if resume_messages.len() >= required_events as usize {
|
||||
let is_disapproved = resumes
|
||||
.iter()
|
||||
.find(|x| x.as_ref().is_ok_and(|x| !x.approved));
|
||||
if is_disapproved.is_none() && resume_messages.len() >= required_events as usize {
|
||||
sqlx::query(
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'approvers'], $2)
|
||||
@@ -1482,7 +1486,8 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
} else if matches!(
|
||||
&status_module,
|
||||
FlowStatusModule::WaitingForPriorSteps { .. }
|
||||
) {
|
||||
) && is_disapproved.is_none()
|
||||
{
|
||||
sqlx::query(
|
||||
"UPDATE queue SET
|
||||
flow_status = JSONB_SET(flow_status, ARRAY['modules', flow_status->>'step'::text], $1),
|
||||
@@ -1514,53 +1519,42 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
} else {
|
||||
tx.commit().await?;
|
||||
|
||||
let success = false;
|
||||
let skipped = false;
|
||||
let (logs, error_name) = if let Some(disapprover) = is_disapproved {
|
||||
(
|
||||
format!(
|
||||
"Disapproved by {:?}",
|
||||
disapprover.as_ref().unwrap().approver
|
||||
),
|
||||
"SuspendedDisapproved",
|
||||
)
|
||||
} else {
|
||||
(
|
||||
"Timed out waiting to be resumed".to_string(),
|
||||
"SuspendedTimedOut",
|
||||
)
|
||||
};
|
||||
|
||||
let result: Value = json!({ "error": {"message": logs, "name": error_name}});
|
||||
|
||||
let logs = "Timed out waiting to be resumed".to_string();
|
||||
append_logs(flow_job.id, flow_job.workspace_id.clone(), logs.clone(), db).await;
|
||||
|
||||
let result = json!({ "error": {"message": logs, "name": "SuspendedTimeout"}});
|
||||
let canceled_by = if flow_job.canceled {
|
||||
Some(CanceledBy {
|
||||
username: flow_job.canceled_by.clone(),
|
||||
reason: flow_job.canceled_reason.clone(),
|
||||
job_completed_tx
|
||||
.send(SendResult::UpdateFlow {
|
||||
flow: flow_job.id,
|
||||
success: false,
|
||||
result: to_raw_value(&result),
|
||||
stop_early_override: None,
|
||||
w_id: flow_job.workspace_id.clone(),
|
||||
worker_dir: worker_dir.to_string(),
|
||||
token: client.token.clone(),
|
||||
})
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let _uuid = add_completed_job(
|
||||
db,
|
||||
&flow_job,
|
||||
success,
|
||||
skipped,
|
||||
Json(&result),
|
||||
0,
|
||||
canceled_by,
|
||||
rsmq.clone(),
|
||||
false,
|
||||
)
|
||||
.await?;
|
||||
if flow_job.is_flow_step {
|
||||
if let Some(parent_job) = flow_job.parent_job {
|
||||
job_completed_tx
|
||||
.send(SendResult::UpdateFlow {
|
||||
flow: parent_job,
|
||||
success: true,
|
||||
result: to_raw_value(&result),
|
||||
stop_early_override: Some(true),
|
||||
w_id: flow_job.workspace_id.clone(),
|
||||
worker_dir: worker_dir.to_string(),
|
||||
token: client.token.clone(),
|
||||
})
|
||||
.await
|
||||
.map_err(|e| {
|
||||
Error::InternalErr(format!(
|
||||
"error sending update flow message to job completed channel: {e}"
|
||||
))
|
||||
})?;
|
||||
}
|
||||
}
|
||||
.await
|
||||
.map_err(|e| {
|
||||
Error::InternalErr(format!(
|
||||
"error sending update flow message to job completed channel: {e}"
|
||||
))
|
||||
})?;
|
||||
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user