diff --git a/backend/.sqlx/query-c4f1c14c3aae145b52ff39bdecd779cd3eba27a869fda68516dc68dc2abefd38.json b/backend/.sqlx/query-2c0c9312b8b326a3759566059d01c79efee28920a1a5afe2df92043526c1de82.json similarity index 55% rename from backend/.sqlx/query-c4f1c14c3aae145b52ff39bdecd779cd3eba27a869fda68516dc68dc2abefd38.json rename to backend/.sqlx/query-2c0c9312b8b326a3759566059d01c79efee28920a1a5afe2df92043526c1de82.json index de3bd165c6..b68a435417 100644 --- a/backend/.sqlx/query-c4f1c14c3aae145b52ff39bdecd779cd3eba27a869fda68516dc68dc2abefd38.json +++ b/backend/.sqlx/query-2c0c9312b8b326a3759566059d01c79efee28920a1a5afe2df92043526c1de82.json @@ -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" } diff --git a/backend/.sqlx/query-4fb3a4712d88afed40082d8d8bd63b5dedad61caa68e0e470252083d80df605f.json b/backend/.sqlx/query-4fb3a4712d88afed40082d8d8bd63b5dedad61caa68e0e470252083d80df605f.json new file mode 100644 index 0000000000..90bdd1c9b8 --- /dev/null +++ b/backend/.sqlx/query-4fb3a4712d88afed40082d8d8bd63b5dedad61caa68e0e470252083d80df605f.json @@ -0,0 +1,14 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE queue SET suspend = 0 WHERE id = $1", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "4fb3a4712d88afed40082d8d8bd63b5dedad61caa68e0e470252083d80df605f" +} diff --git a/backend/migrations/20240317084804_improve_cancelled_suspended_flow.down.sql b/backend/migrations/20240317084804_improve_cancelled_suspended_flow.down.sql new file mode 100644 index 0000000000..d2e8bb3d9f --- /dev/null +++ b/backend/migrations/20240317084804_improve_cancelled_suspended_flow.down.sql @@ -0,0 +1,2 @@ +-- Add down migration script here +ALTER TABLE resume_job DROP COLUMN approved; \ No newline at end of file diff --git a/backend/migrations/20240317084804_improve_cancelled_suspended_flow.up.sql b/backend/migrations/20240317084804_improve_cancelled_suspended_flow.up.sql new file mode 100644 index 0000000000..8af38d5fc4 --- /dev/null +++ b/backend/migrations/20240317084804_improve_cancelled_suspended_flow.up.sql @@ -0,0 +1,2 @@ +-- Add up migration script here +ALTER TABLE resume_job ADD COLUMN approved BOOLEAN NOT NULL DEFAULT true; \ No newline at end of file diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index 1995f8d9f1..e7f110f8fa 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -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, QueryOrBody(value): QueryOrBody, ) -> error::Result { + 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, + db: sqlx::Pool, + w_id: String, + job_id: Uuid, + resume_id: u32, + approver: QueryApprover, + secret: String, + authed: Option, + approved: bool, +) -> Result { 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, + 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, Extension(db): Extension, - Extension(rsmq): Extension>, - 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, -) -> error::Result { - 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, +) -> error::Result { + 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)] diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 9568a022bf..3bb864c2c1 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -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>, pub approver: Option, + pub approved: bool, pub resume_id: i32, } @@ -1359,7 +1360,7 @@ async fn push_next_flow_job .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 .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 } 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 } 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(()); } }