fix: fix flow progress monitor when using parallel branches and continuing long after

This commit is contained in:
Ruben Fiszel
2024-03-01 12:01:17 +01:00
parent 410ec2cd78
commit 07fb3754af
4 changed files with 94 additions and 20 deletions
@@ -0,0 +1,20 @@
{
"db_name": "PostgreSQL",
"query": "SELECT owner FROM workspace WHERE owner NOT LIKE '%@windmill.dev'",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "owner",
"type_info": "Varchar"
}
],
"parameters": {
"Left": []
},
"nullable": [
false
]
},
"hash": "2fa6e9c2cfc4f4ddbaa91793686feb9773b36993dd417019db51fa59fd4952c3"
}
@@ -0,0 +1,20 @@
{
"db_name": "PostgreSQL",
"query": "SELECT email FROM password WHERE super_admin IS true AND email NOT LIKE '%@windmill.dev'",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "email",
"type_info": "Varchar"
}
],
"parameters": {
"Left": []
},
"nullable": [
false
]
},
"hash": "3dd6ba886b214274d9453667e3732421251c3a0d84b0909206447a3927bad8b3"
}
@@ -0,0 +1,22 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM parallel_monitor_lock WHERE parent_flow_id = $1 RETURNING last_ping",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "last_ping",
"type_info": "Timestamptz"
}
],
"parameters": {
"Left": [
"Uuid"
]
},
"nullable": [
true
]
},
"hash": "fec3144df0d33bf8cfd4a629d85ba58d79f523994e75b7b4631f493d12933152"
}
+32 -20
View File
@@ -364,6 +364,17 @@ pub async fn update_flow_status_after_job_completion_internal<
branch_chosen: None,
}
};
let r = sqlx::query_scalar!(
"DELETE FROM parallel_monitor_lock WHERE parent_flow_id = $1 RETURNING last_ping",
flow,
).fetch_optional(db).await?;
if r.is_some() {
tracing::info!(
"parallel flow has removed lock on its parent, last ping was {:?}",
r.unwrap()
);
}
(true, Some(new_status))
} else {
tx.commit().await?;
@@ -377,13 +388,14 @@ pub async fn update_flow_status_after_job_completion_internal<
.await?;
}
sqlx::query!(
"UPDATE queue
SET last_ping = null
WHERE id = $1",
flow
).execute(db).await?;
)
.execute(db)
.await?;
let r = sqlx::query_scalar!(
"DELETE FROM parallel_monitor_lock WHERE parent_flow_id = $1 and job_id = $2 RETURNING last_ping",
@@ -391,7 +403,10 @@ pub async fn update_flow_status_after_job_completion_internal<
job_id_for_status
).fetch_optional(db).await?;
if r.is_some() {
tracing::info!("parallel flow has removed lock on its parent, last ping was {:?}", r.unwrap());
tracing::info!(
"parallel flow has removed lock on its parent, last ping was {:?}",
r.unwrap()
);
}
return Ok(None);
@@ -676,7 +691,7 @@ pub async fn update_flow_status_after_job_completion_internal<
0,
None,
rsmq.clone(),
true
true,
)
.await?;
} else {
@@ -694,7 +709,7 @@ pub async fn update_flow_status_after_job_completion_internal<
0,
None,
rsmq.clone(),
true
true,
)
.await?;
}
@@ -1064,7 +1079,6 @@ pub async fn handle_flow<R: rsmq_async::RsmqConnection + Send + Sync + Clone>(
.parse_flow_status()
.with_context(|| "Unable to parse flow status")?;
if !flow_job.is_flow_step
&& flow_job.schedule_path.is_some()
&& flow_job.script_path.is_some()
@@ -1144,7 +1158,7 @@ lazy_static::lazy_static! {
#[inline(always)]
fn potentially_crash_for_testing() {
#[cfg(feature = "flow_testing")]
if let Some(crash_at) = CRASH_FORCEFULLY_AT_STEP.as_ref() {
if let Some(crash_at) = CRASH_FORCEFULLY_AT_STEP.as_ref() {
let counter = CRASH_STEP_COUNTER.fetch_add(1, Ordering::SeqCst);
if &counter == crash_at {
panic!("CRASH#1 - expected crash for testing at step {}", crash_at);
@@ -1254,9 +1268,6 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
))
})?;
return Ok(());
}
}
@@ -1289,7 +1300,6 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
))
})?;
return Ok(());
}
}
@@ -1487,7 +1497,9 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
WHERE id = $1 AND last_ping = $2",
flow_job.id,
flow_job.last_ping
).execute(&mut *tx).await?;
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
return Ok(());
@@ -1518,7 +1530,7 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
0,
canceled_by,
rsmq.clone(),
false
false,
)
.await?;
if flow_job.is_flow_step {
@@ -2122,19 +2134,16 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
.await?;
};
potentially_crash_for_testing();
sqlx::query!(
"UPDATE queue
SET last_ping = null
WHERE id = $1",
flow_job.id
).execute(&mut tx).await?;
)
.execute(&mut tx)
.await?;
tx.commit().await?;
tracing::info!(id = %flow_job.id, root_id = %job_root, "all next flow jobs pushed: {uuids:?}");
@@ -2145,7 +2154,10 @@ async fn push_next_flow_job<R: rsmq_async::RsmqConnection + Send + Sync + Clone>
"Cannot continue on same worker with multiple jobs, parallel cannot be used in conjunction with same_worker".to_string(),
));
}
same_worker_tx.send(SameWorkerPayload { job_id: first_uuid, recoverable: true }).await.map_err(to_anyhow)?;
same_worker_tx
.send(SameWorkerPayload { job_id: first_uuid, recoverable: true })
.await
.map_err(to_anyhow)?;
}
return Ok(());
}