mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-08 16:03:27 +00:00
feat: Add possibility to delete flow step results when the flow is complete (#2806)
* feat: Add possibility to delete flow step results when the flow is complete * Add third layer of tabs and gate feature to EE
This commit is contained in:
@@ -27,7 +27,7 @@ use tokio::sync::mpsc::Sender;
|
||||
use tracing::instrument;
|
||||
use uuid::Uuid;
|
||||
use windmill_common::flow_status::{
|
||||
ApprovalConditions, FlowStatusModuleWParent, Iterator, JobResult,
|
||||
ApprovalConditions, FlowCleanupModule, FlowStatusModuleWParent, Iterator, JobResult,
|
||||
};
|
||||
use windmill_common::flows::add_virtual_items_if_necessary;
|
||||
use windmill_common::jobs::{
|
||||
@@ -233,20 +233,18 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
(false, false)
|
||||
} else {
|
||||
let row = sqlx::query(
|
||||
"
|
||||
SELECT raw_flow->'modules'->$1::int->'stop_after_if'->>'expr' as stop_early_expr,
|
||||
(raw_flow->'modules'->$1::int->'stop_after_if'->>'skip_if_stopped')::bool as skip_if_stopped,
|
||||
args
|
||||
FROM queue
|
||||
WHERE id = $2
|
||||
").bind(
|
||||
old_status.step)
|
||||
.bind(
|
||||
flow
|
||||
)
|
||||
.fetch_one(db)
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("retrieval of stop_early_expr from state: {e}")))?;
|
||||
"SELECT
|
||||
raw_flow->'modules'->$1::int->'stop_after_if'->>'expr' as stop_early_expr,
|
||||
(raw_flow->'modules'->$1::int->'stop_after_if'->>'skip_if_stopped')::bool as skip_if_stopped,
|
||||
args
|
||||
FROM queue
|
||||
WHERE id = $2"
|
||||
)
|
||||
.bind(old_status.step)
|
||||
.bind(flow)
|
||||
.fetch_one(db)
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("retrieval of stop_early_expr from state: {e}")))?;
|
||||
let r = SkipIfStopped::from_row(&row)?;
|
||||
|
||||
let stop_early = success
|
||||
@@ -294,34 +292,30 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
let (nindex, len) = match (iterator, branchall) {
|
||||
(Some(Iterator { itered, .. }), _) => {
|
||||
let nindex = sqlx::query_scalar!(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'iterator', 'index'], ((flow_status->'modules'->$1::int->'iterator'->>'index')::int + 1)::text::jsonb)
|
||||
WHERE id = $2
|
||||
RETURNING (flow_status->'modules'->$1::int->'iterator'->>'index')::int
|
||||
",
|
||||
old_status.step,
|
||||
flow
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?
|
||||
.ok_or_else(|| Error::InternalErr(format!("requiring an index in InProgress")))?;
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'iterator', 'index'], ((flow_status->'modules'->$1::int->'iterator'->>'index')::int + 1)::text::jsonb)
|
||||
WHERE id = $2
|
||||
RETURNING (flow_status->'modules'->$1::int->'iterator'->>'index')::int",
|
||||
old_status.step,
|
||||
flow
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?
|
||||
.ok_or_else(|| Error::InternalErr(format!("requiring an index in InProgress")))?;
|
||||
(nindex, itered.len() as i32)
|
||||
}
|
||||
(_, Some(BranchAllStatus { len, .. })) => {
|
||||
let nindex = sqlx::query_scalar!(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'branchall', 'branch'], ((flow_status->'modules'->$1::int->'branchall'->>'branch')::int + 1)::text::jsonb)
|
||||
WHERE id = $2
|
||||
RETURNING (flow_status->'modules'->$1::int->'branchall'->>'branch')::int
|
||||
",
|
||||
old_status.step,
|
||||
flow
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?
|
||||
.ok_or_else(|| Error::InternalErr(format!("requiring an index in InProgress")))?;
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'branchall', 'branch'], ((flow_status->'modules'->$1::int->'branchall'->>'branch')::int + 1)::text::jsonb)
|
||||
WHERE id = $2
|
||||
RETURNING (flow_status->'modules'->$1::int->'branchall'->>'branch')::int",
|
||||
old_status.step,
|
||||
flow
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?
|
||||
.ok_or_else(|| Error::InternalErr(format!("requiring an index in InProgress")))?;
|
||||
(nindex, *len as i32)
|
||||
}
|
||||
_ => Err(Error::InternalErr(format!(
|
||||
@@ -331,12 +325,8 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
if nindex == len {
|
||||
let new_status = if skip_loop_failures
|
||||
|| sqlx::query_scalar!(
|
||||
"
|
||||
SELECT success
|
||||
FROM completed_job
|
||||
WHERE id = ANY($1)
|
||||
",
|
||||
jobs.as_slice(),
|
||||
"SELECT success FROM completed_job WHERE id = ANY($1)",
|
||||
jobs.as_slice()
|
||||
)
|
||||
.fetch_all(&mut tx)
|
||||
.await?
|
||||
@@ -353,7 +343,6 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
}
|
||||
} else {
|
||||
success = false;
|
||||
|
||||
FlowStatusModule::Failure {
|
||||
id: module_status.id(),
|
||||
job: job_id_for_status.clone(),
|
||||
@@ -435,11 +424,9 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
|
||||
let step_counter = if inc_step_counter {
|
||||
sqlx::query!(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['step'], $1)
|
||||
WHERE id = $2
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['step'], $1)
|
||||
WHERE id = $2",
|
||||
json!(old_status.step + 1),
|
||||
flow
|
||||
)
|
||||
@@ -458,18 +445,16 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
if let Some(new_status) = new_status.as_ref() {
|
||||
if is_failure_step {
|
||||
let parent_module = sqlx::query_scalar!(
|
||||
"SELECT flow_status->'failure_module'->>'parent_module' FROM queue WHERE id = $1",
|
||||
flow
|
||||
)
|
||||
"SELECT flow_status->'failure_module'->>'parent_module' FROM queue WHERE id = $1",
|
||||
flow
|
||||
)
|
||||
.fetch_one(&mut tx)
|
||||
.await?;
|
||||
|
||||
sqlx::query!(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['failure_module'], $1)
|
||||
WHERE id = $2
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['failure_module'], $1)
|
||||
WHERE id = $2",
|
||||
json!(FlowStatusModuleWParent {
|
||||
parent_module,
|
||||
module_status: new_status.clone()
|
||||
@@ -480,11 +465,9 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
.await?;
|
||||
} else {
|
||||
sqlx::query!(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2)
|
||||
WHERE id = $3
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2)
|
||||
WHERE id = $3",
|
||||
old_status.step.to_string(),
|
||||
json!(new_status),
|
||||
flow
|
||||
@@ -494,11 +477,9 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
|
||||
if let Some(job_result) = new_status.job_result() {
|
||||
sqlx::query!(
|
||||
"
|
||||
UPDATE queue
|
||||
SET leaf_jobs = JSONB_SET(coalesce(leaf_jobs, '{}'::jsonb), ARRAY[$1::TEXT], $2)
|
||||
WHERE COALESCE((SELECT root_job FROM queue WHERE id = $3), $3) = id
|
||||
",
|
||||
"UPDATE queue
|
||||
SET leaf_jobs = JSONB_SET(coalesce(leaf_jobs, '{}'::jsonb), ARRAY[$1::TEXT], $2)
|
||||
WHERE COALESCE((SELECT root_job FROM queue WHERE id = $3), $3) = id",
|
||||
new_status.id(),
|
||||
json!(job_result),
|
||||
flow
|
||||
@@ -519,12 +500,10 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
|
||||
if matches!(&new_status, Some(FlowStatusModule::Success { .. })) {
|
||||
sqlx::query(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = flow_status - 'retry'
|
||||
WHERE id = $1
|
||||
RETURNING flow_status
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = flow_status - 'retry'
|
||||
WHERE id = $1
|
||||
RETURNING flow_status",
|
||||
)
|
||||
.bind(flow)
|
||||
.execute(&mut tx)
|
||||
@@ -608,6 +587,26 @@ pub async fn update_flow_status_after_job_completion_internal<
|
||||
} else {
|
||||
"Flow job completed with error".to_string()
|
||||
};
|
||||
|
||||
#[cfg(feature = "enterprise")]
|
||||
if flow_job.parent_job.is_none() {
|
||||
// run the cleanup step only when the root job is complete
|
||||
let cleanup_module = retrieve_cleanup_module(flow, db).await?;
|
||||
if cleanup_module.flow_jobs_to_clean.len() > 0 {
|
||||
tracing::debug!(
|
||||
"Cleaning up jobs arguments, result and logs as they were marked as delete_after_use {:?}",
|
||||
cleanup_module.flow_jobs_to_clean
|
||||
);
|
||||
sqlx::query!(
|
||||
"UPDATE completed_job
|
||||
SET logs = '##DELETED##', args = '{}'::jsonb, result = '{}'::jsonb
|
||||
WHERE id = ANY($1)",
|
||||
&cleanup_module.flow_jobs_to_clean,
|
||||
)
|
||||
.execute(db)
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
if flow_job.canceled {
|
||||
add_completed_job_error(
|
||||
db,
|
||||
@@ -740,12 +739,9 @@ async fn retrieve_flow_jobs_results(
|
||||
job_uuids: &Vec<Uuid>,
|
||||
) -> error::Result<Box<RawValue>> {
|
||||
let results = sqlx::query(
|
||||
"
|
||||
SELECT result, id
|
||||
"SELECT result, id
|
||||
FROM completed_job
|
||||
WHERE id = ANY($1)
|
||||
AND workspace_id = $2
|
||||
",
|
||||
WHERE id = ANY($1) AND workspace_id = $2",
|
||||
)
|
||||
.bind(job_uuids.as_slice())
|
||||
.bind(w_id)
|
||||
@@ -793,11 +789,9 @@ async fn compute_skip_loop_failures_and_parallelism(
|
||||
db: &DB,
|
||||
) -> Result<(Option<bool>, Option<i32>), Error> {
|
||||
sqlx::query_as(
|
||||
"
|
||||
SELECT (raw_flow->'modules'->$1->'value'->>'skip_failures')::bool, (raw_flow->'modules'->$1->'value'->>'parallelism')::int
|
||||
FROM queue
|
||||
WHERE id = $2
|
||||
",
|
||||
"SELECT (raw_flow->'modules'->$1->'value'->>'skip_failures')::bool, (raw_flow->'modules'->$1->'value'->>'parallelism')::int
|
||||
FROM queue
|
||||
WHERE id = $2",
|
||||
)
|
||||
.bind(step)
|
||||
.bind(flow)
|
||||
@@ -814,11 +808,9 @@ async fn compute_skip_branchall_failure<'c>(
|
||||
db: &DB,
|
||||
) -> Result<Option<bool>, Error> {
|
||||
sqlx::query_as(
|
||||
"
|
||||
SELECT (raw_flow->'modules'->$1->'value'->'branches'->$2->>'skip_failure')::bool
|
||||
FROM queue
|
||||
WHERE id = $3
|
||||
",
|
||||
"SELECT (raw_flow->'modules'->$1->'value'->'branches'->$2->>'skip_failure')::bool
|
||||
FROM queue
|
||||
WHERE id = $3",
|
||||
)
|
||||
.bind(step)
|
||||
.bind(branch as i32)
|
||||
@@ -834,11 +826,9 @@ async fn has_failure_module<'c>(
|
||||
tx: &mut sqlx::Transaction<'c, sqlx::Postgres>,
|
||||
) -> Result<bool, Error> {
|
||||
sqlx::query_scalar::<_, Option<bool>>(
|
||||
"
|
||||
SELECT raw_flow->'failure_module' != 'null'::jsonb
|
||||
FROM queue
|
||||
WHERE id = $1
|
||||
",
|
||||
"SELECT raw_flow->'failure_module' != 'null'::jsonb
|
||||
FROM queue
|
||||
WHERE id = $1",
|
||||
)
|
||||
.bind(flow)
|
||||
.fetch_one(&mut **tx)
|
||||
@@ -847,6 +837,27 @@ async fn has_failure_module<'c>(
|
||||
.map(|v| v.unwrap_or(false))
|
||||
}
|
||||
|
||||
async fn retrieve_cleanup_module<'c>(flow_uuid: Uuid, db: &DB) -> Result<FlowCleanupModule, Error> {
|
||||
tracing::warn!("Retrieving cleanup module of flow {}", flow_uuid);
|
||||
let raw_value = sqlx::query_scalar!(
|
||||
"SELECT flow_status->'cleanup_module' as cleanup_module
|
||||
FROM queue
|
||||
WHERE id = $1",
|
||||
flow_uuid,
|
||||
)
|
||||
.fetch_one(db)
|
||||
.await
|
||||
.map_err(|e| Error::InternalErr(format!("error during retrieval of cleanup module: {e}")))?;
|
||||
|
||||
raw_value
|
||||
.clone()
|
||||
.and_then(|rv| serde_json::from_value::<FlowCleanupModule>(rv).ok())
|
||||
.ok_or(Error::InternalErr(format!(
|
||||
"Unable to parse flow cleanup module {:?}",
|
||||
raw_value
|
||||
)))
|
||||
}
|
||||
|
||||
fn next_retry(retry: &Retry, status: &RetryStatus) -> Option<(u16, Duration)> {
|
||||
(status.fail_count <= MAX_RETRY_ATTEMPTS)
|
||||
.then(|| &retry)
|
||||
@@ -1100,7 +1111,7 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
|
||||
let flow_job_args = flow_job.get_args();
|
||||
|
||||
// if this is an empty module of if the module has aleady been completed, successfully, update the parent flow
|
||||
// if this is an empty module of if the module has already been completed, successfully, update the parent flow
|
||||
if flow.modules.is_empty() || matches!(status_module, FlowStatusModule::Success { .. }) {
|
||||
let r;
|
||||
return update_flow_status_after_job_completion(
|
||||
@@ -1281,12 +1292,9 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
user_groups_required: user_groups_required,
|
||||
};
|
||||
sqlx::query(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status =
|
||||
JSONB_SET(flow_status, ARRAY['approval_conditions'], $1)
|
||||
WHERE id = $2
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['approval_conditions'], $1)
|
||||
WHERE id = $2",
|
||||
)
|
||||
.bind(json!(approval_conditions))
|
||||
.bind(flow_job.id)
|
||||
@@ -1296,12 +1304,9 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
|
||||
if resume_messages.len() >= required_events as usize {
|
||||
sqlx::query(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status =
|
||||
JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'approvers'], $2)
|
||||
WHERE id = $3
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT, 'approvers'], $2)
|
||||
WHERE id = $3",
|
||||
)
|
||||
.bind(status.step - 1)
|
||||
.bind(json!(resumes
|
||||
@@ -1321,11 +1326,9 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
|
||||
// Remove the approval conditions from the flow status
|
||||
sqlx::query(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = flow_status - 'approval_conditions'
|
||||
WHERE id = $1
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = flow_status - 'approval_conditions'
|
||||
WHERE id = $1",
|
||||
)
|
||||
.bind(flow_job.id)
|
||||
.execute(&mut *tx)
|
||||
@@ -1340,13 +1343,11 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
FlowStatusModule::WaitingForPriorSteps { .. }
|
||||
) {
|
||||
sqlx::query(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', flow_status->>'step'::text], $1)
|
||||
, suspend = $2
|
||||
, suspend_until = now() + $3
|
||||
WHERE id = $4
|
||||
",
|
||||
"UPDATE queue SET
|
||||
flow_status = JSONB_SET(flow_status, ARRAY['modules', flow_status->>'step'::text], $1),
|
||||
suspend = $2,
|
||||
suspend_until = now() + $3
|
||||
WHERE id = $4",
|
||||
)
|
||||
.bind(json!(FlowStatusModule::WaitingForEvents { id: status_module.id(), count: required_events, job: last }))
|
||||
.bind((required_events - resume_messages.len() as u16) as i32)
|
||||
@@ -1498,11 +1499,9 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
scheduled_for_o = Some(from_now(retry_in));
|
||||
status.retry.failed_jobs.push(job.clone());
|
||||
sqlx::query(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['retry'], $1)
|
||||
WHERE id = $2
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['retry'], $1)
|
||||
WHERE id = $2",
|
||||
)
|
||||
.bind(json!(RetryStatus { fail_count, ..status.retry.clone() }))
|
||||
.bind(flow_job.id)
|
||||
@@ -1536,11 +1535,9 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
let retry = &module.retry.clone().unwrap_or_default();
|
||||
if retry.has_attempts() {
|
||||
sqlx::query(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['retry'], $1)
|
||||
WHERE id = $2
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['retry'], $1)
|
||||
WHERE id = $2",
|
||||
)
|
||||
.bind(json!(RetryStatus { fail_count: 0, failed_jobs: vec![] }))
|
||||
.bind(flow_job.id)
|
||||
@@ -1562,11 +1559,9 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
&& status.retry.fail_count == 0 =>
|
||||
{
|
||||
sqlx::query(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['retry'], $1)
|
||||
WHERE id = $2
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['retry'], $1)
|
||||
WHERE id = $2",
|
||||
)
|
||||
.bind(json!(RetryStatus { fail_count: 0, failed_jobs: vec![] }))
|
||||
.bind(flow_job.id)
|
||||
@@ -1671,18 +1666,22 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
NextFlowTransform::Continue(job_payload, next_state) => (job_payload, next_state),
|
||||
NextFlowTransform::EmptyInnerFlows => {
|
||||
sqlx::query(
|
||||
r#"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2)
|
||||
WHERE id = $3
|
||||
"#,
|
||||
)
|
||||
.bind(status.step)
|
||||
.bind(json!(FlowStatusModule::Success { id: status_module.id(), job: Uuid::nil(), flow_jobs: Some(vec![]), branch_chosen: None, approvers: vec![] }))
|
||||
.bind(flow_job.id)
|
||||
.execute(db)
|
||||
.await?;
|
||||
// flow is reprocessed by the worker in a state where the module has completed succesfully.
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2)
|
||||
WHERE id = $3",
|
||||
)
|
||||
.bind(status.step)
|
||||
.bind(json!(FlowStatusModule::Success {
|
||||
id: status_module.id(),
|
||||
job: Uuid::nil(),
|
||||
flow_jobs: Some(vec![]),
|
||||
branch_chosen: None,
|
||||
approvers: vec![]
|
||||
}))
|
||||
.bind(flow_job.id)
|
||||
.execute(db)
|
||||
.await?;
|
||||
// flow is reprocessed by the worker in a state where the module has completed successfully.
|
||||
// The next steps are pull -> handle flow -> push next flow job -> update flow status since module status is success
|
||||
same_worker_tx
|
||||
.send(flow_job.id)
|
||||
@@ -1846,11 +1845,9 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
if let FlowModuleValue::ForloopFlow { parallelism: Some(p), .. } = &module.value {
|
||||
if i as u16 >= *p {
|
||||
sqlx::query!(
|
||||
"
|
||||
UPDATE queue
|
||||
SET suspend = $1, suspend_until = now() + interval '14 day', running = true
|
||||
WHERE id = $2
|
||||
",
|
||||
"UPDATE queue
|
||||
SET suspend = $1, suspend_until = now() + interval '14 day', running = true
|
||||
WHERE id = $2",
|
||||
(i as u16 - p + 1) as i32,
|
||||
uuid,
|
||||
)
|
||||
@@ -1858,6 +1855,22 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
|
||||
if payload_tag.delete_after_use {
|
||||
let uuid_singleton_json = serde_json::to_value(&[uuid])
|
||||
.map_err(|e| error::Error::InternalErr(format!("Unable to serialize uuid: {e}")))?;
|
||||
|
||||
sqlx::query(
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(flow_status, ARRAY['cleanup_module', 'flow_jobs_to_clean'], COALESCE(flow_status->'cleanup_module'->'flow_jobs_to_clean', '[]'::jsonb) || $1)
|
||||
WHERE id = $2",
|
||||
)
|
||||
.bind(uuid_singleton_json)
|
||||
.bind(root_job.unwrap_or(flow_job.id))
|
||||
.execute(&mut inner_tx)
|
||||
.await?;
|
||||
}
|
||||
|
||||
tx = inner_tx;
|
||||
uuids.push(uuid);
|
||||
}
|
||||
@@ -1931,13 +1944,10 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
|
||||
if i >= flow.modules.len() {
|
||||
sqlx::query!(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(
|
||||
JSONB_SET(flow_status, ARRAY['failure_module'], $1),
|
||||
ARRAY['step'], $2)
|
||||
WHERE id = $3
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(
|
||||
JSONB_SET(flow_status, ARRAY['failure_module'], $1), ARRAY['step'], $2)
|
||||
WHERE id = $3",
|
||||
json!(FlowStatusModuleWParent {
|
||||
parent_module: Some(current_id.clone()),
|
||||
module_status: new_status.clone()
|
||||
@@ -1945,23 +1955,20 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
|
||||
json!(i),
|
||||
flow_job.id
|
||||
)
|
||||
.execute(db)
|
||||
.execute(&mut tx)
|
||||
.await?;
|
||||
} else {
|
||||
sqlx::query!(
|
||||
"
|
||||
UPDATE queue
|
||||
SET flow_status = JSONB_SET(
|
||||
JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2),
|
||||
ARRAY['step'], $3)
|
||||
WHERE id = $4
|
||||
",
|
||||
"UPDATE queue
|
||||
SET flow_status = JSONB_SET(
|
||||
JSONB_SET(flow_status, ARRAY['modules', $1::TEXT], $2), ARRAY['step'], $3)
|
||||
WHERE id = $4",
|
||||
i as i32,
|
||||
json!(new_status),
|
||||
json!(i),
|
||||
flow_job.id
|
||||
)
|
||||
.execute(db)
|
||||
.execute(&mut tx)
|
||||
.await?;
|
||||
};
|
||||
|
||||
@@ -2085,6 +2092,7 @@ enum NextStatus {
|
||||
struct JobPayloadWithTag {
|
||||
payload: JobPayload,
|
||||
tag: Option<String>,
|
||||
delete_after_use: bool,
|
||||
}
|
||||
enum ContinuePayload {
|
||||
SingleJob(JobPayloadWithTag),
|
||||
@@ -2137,20 +2145,26 @@ async fn compute_next_flow_transform(
|
||||
ContinuePayload::SingleJob(JobPayloadWithTag {
|
||||
payload: JobPayload::Identity,
|
||||
tag: None,
|
||||
delete_after_use: false,
|
||||
}),
|
||||
NextStatus::NextStep,
|
||||
));
|
||||
}
|
||||
let trivial_next_job = |payload| {
|
||||
Ok(NextFlowTransform::Continue(
|
||||
ContinuePayload::SingleJob(JobPayloadWithTag { payload, tag: None }),
|
||||
ContinuePayload::SingleJob(JobPayloadWithTag {
|
||||
payload,
|
||||
tag: None,
|
||||
delete_after_use: false,
|
||||
}),
|
||||
NextStatus::NextStep,
|
||||
))
|
||||
};
|
||||
let delete_after_use = module.delete_after_use.unwrap_or(false);
|
||||
match &module.value {
|
||||
FlowModuleValue::Identity => trivial_next_job(JobPayload::Identity),
|
||||
FlowModuleValue::Flow { path, .. } => {
|
||||
let payload = flow_to_payload(path);
|
||||
let payload = flow_to_payload(path, &delete_after_use);
|
||||
Ok(NextFlowTransform::Continue(
|
||||
ContinuePayload::SingleJob(payload),
|
||||
NextStatus::NextStep,
|
||||
@@ -2185,6 +2199,7 @@ async fn compute_next_flow_transform(
|
||||
concurrency_time_window_s,
|
||||
module,
|
||||
tag,
|
||||
&delete_after_use,
|
||||
);
|
||||
Ok(NextFlowTransform::Continue(
|
||||
ContinuePayload::SingleJob(payload),
|
||||
@@ -2374,6 +2389,7 @@ async fn compute_next_flow_transform(
|
||||
restarted_from: None,
|
||||
},
|
||||
tag: None,
|
||||
delete_after_use: delete_after_use,
|
||||
}),
|
||||
NextStatus::NextLoopIteration {
|
||||
next: ns,
|
||||
@@ -2413,6 +2429,7 @@ async fn compute_next_flow_transform(
|
||||
restarted_from: None,
|
||||
},
|
||||
tag: None,
|
||||
delete_after_use: delete_after_use,
|
||||
}
|
||||
};
|
||||
ContinuePayload::ForloopJobs { n: itered.len(), payload }
|
||||
@@ -2513,6 +2530,7 @@ async fn compute_next_flow_transform(
|
||||
restarted_from: None,
|
||||
},
|
||||
tag: None,
|
||||
delete_after_use: delete_after_use,
|
||||
}),
|
||||
NextStatus::BranchChosen(branch),
|
||||
))
|
||||
@@ -2562,6 +2580,7 @@ async fn compute_next_flow_transform(
|
||||
restarted_from: None,
|
||||
},
|
||||
tag: None,
|
||||
delete_after_use: delete_after_use,
|
||||
}
|
||||
})
|
||||
.collect(),
|
||||
@@ -2627,6 +2646,7 @@ async fn compute_next_flow_transform(
|
||||
restarted_from: None,
|
||||
},
|
||||
tag: None,
|
||||
delete_after_use: delete_after_use,
|
||||
}),
|
||||
NextStatus::NextBranchStep(NextBranch { status: branch_status, flow_jobs }),
|
||||
))
|
||||
@@ -2641,8 +2661,9 @@ async fn payload_from_simple_module(
|
||||
module: &FlowModule,
|
||||
inner_path: Option<String>,
|
||||
) -> Result<JobPayloadWithTag, Error> {
|
||||
let delete_after_use = module.delete_after_use.unwrap_or(false);
|
||||
Ok(match value {
|
||||
FlowModuleValue::Flow { path, .. } => flow_to_payload(path),
|
||||
FlowModuleValue::Flow { path, .. } => flow_to_payload(path, &delete_after_use),
|
||||
FlowModuleValue::Script { path: script_path, hash: script_hash, .. } => {
|
||||
script_to_payload(script_hash, script_path, db, flow_job, module).await?
|
||||
}
|
||||
@@ -2664,6 +2685,7 @@ async fn payload_from_simple_module(
|
||||
concurrency_time_window_s,
|
||||
module,
|
||||
tag,
|
||||
&delete_after_use,
|
||||
),
|
||||
_ => unreachable!("is simple flow"),
|
||||
})
|
||||
@@ -2678,6 +2700,7 @@ fn raw_script_to_payload(
|
||||
concurrency_time_window_s: &Option<i32>,
|
||||
module: &FlowModule,
|
||||
tag: &Option<String>,
|
||||
delete_after_use: &bool,
|
||||
) -> JobPayloadWithTag {
|
||||
JobPayloadWithTag {
|
||||
payload: JobPayload::Code(RawCode {
|
||||
@@ -2691,12 +2714,13 @@ fn raw_script_to_payload(
|
||||
dedicated_worker: None,
|
||||
}),
|
||||
tag: tag.clone(),
|
||||
delete_after_use: *delete_after_use,
|
||||
}
|
||||
}
|
||||
|
||||
fn flow_to_payload(path: &str) -> JobPayloadWithTag {
|
||||
fn flow_to_payload(path: &str, delete_after_use: &bool) -> JobPayloadWithTag {
|
||||
let payload = JobPayload::Flow { path: path.to_string(), dedicated_worker: None };
|
||||
JobPayloadWithTag { payload, tag: None }
|
||||
JobPayloadWithTag { payload, tag: None, delete_after_use: *delete_after_use }
|
||||
}
|
||||
|
||||
async fn script_to_payload(
|
||||
@@ -2706,7 +2730,7 @@ async fn script_to_payload(
|
||||
flow_job: &QueuedJob,
|
||||
module: &FlowModule,
|
||||
) -> Result<JobPayloadWithTag, Error> {
|
||||
let (payload, tag) = if script_hash.is_none() {
|
||||
let (payload, tag, delete_after_use) = if script_hash.is_none() {
|
||||
script_path_to_payload(script_path, &db, &flow_job.workspace_id).await?
|
||||
} else {
|
||||
let hash = script_hash.clone().unwrap();
|
||||
@@ -2719,6 +2743,7 @@ async fn script_to_payload(
|
||||
language,
|
||||
dedicated_worker,
|
||||
priority,
|
||||
delete_after_use,
|
||||
) = script_hash_to_tag_and_limits(&hash, &mut tx, &flow_job.workspace_id).await?;
|
||||
(
|
||||
JobPayload::ScriptHash {
|
||||
@@ -2732,9 +2757,13 @@ async fn script_to_payload(
|
||||
priority,
|
||||
},
|
||||
tag,
|
||||
delete_after_use,
|
||||
)
|
||||
};
|
||||
Ok(JobPayloadWithTag { payload, tag })
|
||||
// the module value overrides the value set at the script level. Defaults to false if both are unset.
|
||||
let final_delete_after_user =
|
||||
module.delete_after_use.unwrap_or(false) || delete_after_use.unwrap_or(false);
|
||||
Ok(JobPayloadWithTag { payload, tag, delete_after_use: final_delete_after_user })
|
||||
}
|
||||
|
||||
async fn get_transform_context(
|
||||
|
||||
Reference in New Issue
Block a user