fix: parallel flow with parallelism constraint could deadlock

This commit is contained in:
Ruben Fiszel
2024-04-06 13:25:28 +02:00
parent 08231c02d2
commit 9131d5cc40
15 changed files with 163 additions and 128 deletions
@@ -1,17 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO workspace_settings\n (workspace_id, slack_team_id, slack_name, slack_email)\n VALUES ($1, $2, $3, $4) ON CONFLICT (workspace_id) DO UPDATE SET slack_team_id = $2, slack_name = $3, slack_email = $4",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar",
"Varchar"
]
},
"nullable": []
},
"hash": "43b376a2eff086a32cd76e54361ce3631feee1565935d2a6ddbecc17950758d1"
}
@@ -1,18 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO resource\n (workspace_id, path, value, description, resource_type)\n VALUES ($1, $2, $3, $4, $5) ON CONFLICT (workspace_id, path) DO UPDATE SET value = $3",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Jsonb",
"Text",
"Varchar"
]
},
"nullable": []
},
"hash": "489a62b5943a7a21ce487aa7b72a63dfc6300dd93bc29f5ec4cb1bfc471ad0bf"
}
@@ -0,0 +1,20 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO variable\n (workspace_id, path, value, is_secret, description, account, is_oauth)\n VALUES ($1, $2, $3, $4, $5, $6, $7)\n ON CONFLICT (workspace_id, path) DO UPDATE SET value = $3",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar",
"Bool",
"Varchar",
"Int4",
"Bool"
]
},
"nullable": []
},
"hash": "52c8b4350235bdaab4df79e517d5e42a61a4e1e209d120b2c8bb31ebb7ce1e56"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "SELECT * FROM workspace_settings WHERE slack_team_id = $1 AND slack_command_script IS NOT NULL",
"query": "SELECT * FROM workspace_settings WHERE slack_team_id = $1",
"describe": {
"columns": [
{
@@ -144,5 +144,5 @@
true
]
},
"hash": "55cb03040bc2a8c53dd7fbb42bbdcc40f463cbc52d94ed9315cf9a547d4c89f2"
"hash": "5445083864b2b092b012e894bff7630a1d7b9deb8d33e9f909061f351f96844e"
}
@@ -1,20 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO variable\n (workspace_id, path, value, is_secret, description, account, is_oauth)\n VALUES ($1, $2, $3, $4, $5, $6, $7)\n ON CONFLICT (workspace_id, path) DO UPDATE SET value = $3",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar",
"Bool",
"Varchar",
"Int4",
"Bool"
]
},
"nullable": []
},
"hash": "55a2f170823f1d1abce76287d8817a6cf34de92b9b4079c00b75423a9ff835b9"
}
@@ -1,14 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE workspace_settings\n SET slack_team_id = null, slack_name = null WHERE workspace_id = $1",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text"
]
},
"nullable": []
},
"hash": "73ca6ea13b362b2569ea115d28e4e255cd5e9b990ffa89998ef24871d3a9717c"
}
@@ -0,0 +1,18 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO folder\n (workspace_id, name, display_name, owners, extra_perms)\n VALUES ($1, $2, $3, $4, $5) ON CONFLICT DO NOTHING",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar",
"VarcharArray",
"Jsonb"
]
},
"nullable": []
},
"hash": "95ebdfaf0510b9cad861568cd25d479759d4ea3c3ff4e136aad13a3521525372"
}
@@ -0,0 +1,14 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE workspace_settings\n SET slack_team_id = null, slack_name = null WHERE workspace_id = $1",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text"
]
},
"nullable": []
},
"hash": "bf1d8e043338867e1da1ed236ff6c85a566d5fd58d4b0d5c3a10454513811ba3"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE queue SET suspend = suspend - 1 WHERE parent_job = $1",
"query": "UPDATE queue SET suspend = suspend - 1 WHERE parent_job = $1 AND suspend > 0",
"describe": {
"columns": [],
"parameters": {
@@ -10,5 +10,5 @@
},
"nullable": []
},
"hash": "a4ae245dcf7e4b930cd45701db0b7c45f2a5797e8b6724bade5b964d4334c098"
"hash": "d42005ce0f8fbb2f65e8caa3cb06b6888b6c29e40ff575bb3893ea93e48450dc"
}
@@ -0,0 +1,18 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO resource\n (workspace_id, path, value, description, resource_type)\n VALUES ($1, $2, $3, $4, $5) ON CONFLICT (workspace_id, path) DO UPDATE SET value = $3",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Jsonb",
"Text",
"Varchar"
]
},
"nullable": []
},
"hash": "ea8ebb8d972fe99c960b5a69f794ee2b57bfb1914bf370c5b10313e45fa9b65f"
}
@@ -0,0 +1,17 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO workspace_settings\n (workspace_id, slack_team_id, slack_name, slack_email)\n VALUES ($1, $2, $3, $4) ON CONFLICT (workspace_id) DO UPDATE SET slack_team_id = $2, slack_name = $3, slack_email = $4",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar",
"Varchar"
]
},
"nullable": []
},
"hash": "eaa6e9dc4c0d3d6cc8152515019befa880bb3b69ff340d337edd14b65e74e2a3"
}
@@ -1,18 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO folder\n (workspace_id, name, display_name, owners, extra_perms)\n VALUES ($1, $2, $3, $4, $5) ON CONFLICT DO NOTHING",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Varchar",
"Varchar",
"Varchar",
"VarcharArray",
"Jsonb"
]
},
"nullable": []
},
"hash": "f2aee3ae39c90e40dd1835befc339e5381cc9104933cf8e90d840a9bf638ff52"
}
@@ -1,23 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "\n SELECT EXISTS (SELECT 1 \n FROM workspace_settings \n WHERE workspace_id <> $1 \n AND slack_command_script IS NOT NULL\n AND slack_team_id = $2\n AND (SELECT slack_command_script IS NOT NULL FROM workspace_settings WHERE workspace_id = $1))\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "exists",
"type_info": "Bool"
}
],
"parameters": {
"Left": [
"Text",
"Text"
]
},
"nullable": [
null
]
},
"hash": "fa59674af1d1a4ceb696fc883005ef114772f7d2ee0f60cb1358cdb7f0b5cd0c"
}
+5 -1
View File
@@ -1162,6 +1162,7 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
token,
} => {
// let r;
tracing::info!(parent_flow = %flow, "updating flow status");
if let Err(e) = update_flow_status_after_job_completion(
&db2,
&AuthedClient {
@@ -1185,7 +1186,7 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
)
.await
{
tracing::error!("Error updating flow status after job completion: {e}");
tracing::error!("Error updating flow status after job completion for {flow} on {worker_name2}: {e}");
}
}
SendResult::Kill => {
@@ -2259,6 +2260,7 @@ pub async fn process_completed_job<R: rsmq_async::RsmqConnection + Send + Sync +
let timer = _worker_flow_transition_duration
.as_ref()
.map(|x| x.start_timer());
tracing::info!(parent_flow = %parent_job, subflow = %job.id, "updating flow status (2)");
update_flow_status_after_job_completion(
db,
client,
@@ -2296,6 +2298,7 @@ pub async fn process_completed_job<R: rsmq_async::RsmqConnection + Send + Sync +
.await?;
if job.is_flow_step {
if let Some(parent_job) = job.parent_job {
tracing::error!(parent_flow = %parent_job, subflow = %job.id, "process completed job error, updating flow status");
update_flow_status_after_job_completion(
db,
client,
@@ -2409,6 +2412,7 @@ pub async fn handle_job_error<R: rsmq_async::RsmqConnection + Send + Sync + Clon
};
let wrapped_error = WrappedError { error: err.clone() };
tracing::error!(parent_flow = %flow, subflow = %job_status_to_update, "handle job error, updating flow status: {err:?}");
let updated_flow = update_flow_status_after_job_completion(
db,
client,
+67 -13
View File
@@ -116,6 +116,7 @@ pub async fn update_flow_status_after_job_completion<
{
Ok(j) => j,
Err(e) => {
tracing::error!("Error while updating flow status of {} after completion of {}, updating flow status again with error: {e}", nrec.flow,&nrec.job_id_for_status);
update_flow_status_after_job_completion_internal(
db,
client,
@@ -322,8 +323,17 @@ pub async fn update_flow_status_after_job_completion_internal<
flow
)
.fetch_one(&mut tx)
.await?
.await.map_err(|e| {
Error::InternalErr(format!(
"error while fetching iterator index: {e}"
))
})?
.ok_or_else(|| Error::InternalErr(format!("requiring an index in InProgress")))?;
tracing::info!(
"parallel iteration {job_id_for_status} of flow {flow} update nindex: {nindex} len: {len}",
nindex = nindex,
len = itered.len()
);
(nindex, itered.len() as i32)
}
(_, Some(BranchAllStatus { len, .. })) => {
@@ -336,7 +346,12 @@ pub async fn update_flow_status_after_job_completion_internal<
flow
)
.fetch_one(&mut tx)
.await?
.await
.map_err(|e| {
Error::InternalErr(format!(
"error while fetching branchall index: {e}"
))
})?
.ok_or_else(|| Error::InternalErr(format!("requiring an index in InProgress")))?;
(nindex, *len as i32)
}
@@ -351,7 +366,12 @@ pub async fn update_flow_status_after_job_completion_internal<
jobs.as_slice()
)
.fetch_all(&mut tx)
.await?
.await
.map_err(|e| {
Error::InternalErr(format!(
"error while fetching sucess from completed_jobs: {e}"
))
})?
.into_iter()
.all(|x| x)
{
@@ -375,13 +395,20 @@ pub async fn update_flow_status_after_job_completion_internal<
let r = sqlx::query_scalar!(
"DELETE FROM parallel_monitor_lock WHERE parent_flow_id = $1 RETURNING last_ping",
flow,
).fetch_optional(db).await?;
).fetch_optional(db).await.map_err(|e| {
Error::InternalErr(format!(
"error while deleting parallel_monitor_lock: {e}"
))
})?;
if r.is_some() {
tracing::info!(
"parallel flow has removed lock on its parent, last ping was {:?}",
r.unwrap()
);
}
tracing::info!(
"parallel iteration {job_id_for_status} of flow {flow} has finished",
);
(true, Some(new_status))
} else {
@@ -389,11 +416,14 @@ pub async fn update_flow_status_after_job_completion_internal<
if parallelism.is_some() {
sqlx::query!(
"UPDATE queue SET suspend = suspend - 1 WHERE parent_job = $1",
"UPDATE queue SET suspend = suspend - 1 WHERE parent_job = $1 AND suspend > 0",
flow
)
.execute(db)
.await?;
.await
.map_err(|e| {
Error::InternalErr(format!("error decreasing suspend: {e}"))
})?;
}
sqlx::query!(
@@ -476,7 +506,10 @@ pub async fn update_flow_status_after_job_completion_internal<
old_status.step
)
.fetch_optional(&mut tx)
.await?
.await
.map_err(|e| {
Error::InternalErr(format!("error while getting retry fromn step: {e}"))
})?
.flatten();
let retry = retry
@@ -510,7 +543,10 @@ pub async fn update_flow_status_after_job_completion_internal<
flow
)
.execute(&mut tx)
.await?;
.await
.map_err(|e| {
Error::InternalErr(format!("error while setting flow index for {flow}: {e}"))
})?;
old_status.step + 1
} else {
old_status.step
@@ -528,7 +564,11 @@ pub async fn update_flow_status_after_job_completion_internal<
flow
)
.fetch_one(&mut tx)
.await?;
.await.map_err(|e| {
Error::InternalErr(format!(
"error while fetching failure module: {e}"
))
})?;
sqlx::query!(
"UPDATE queue
@@ -541,7 +581,12 @@ pub async fn update_flow_status_after_job_completion_internal<
flow
)
.execute(&mut tx)
.await?;
.await
.map_err(|e| {
Error::InternalErr(format!(
"error while setting flow status in failure step: {e}"
))
})?;
} else {
sqlx::query!(
"UPDATE queue
@@ -552,7 +597,10 @@ pub async fn update_flow_status_after_job_completion_internal<
flow
)
.execute(&mut tx)
.await?;
.await
.map_err(|e| {
Error::InternalErr(format!("error while setting new flow status: {e}"))
})?;
if let Some(job_result) = new_status.job_result() {
sqlx::query!(
@@ -564,7 +612,11 @@ pub async fn update_flow_status_after_job_completion_internal<
flow
)
.execute(&mut tx)
.await?;
.await.map_err(|e| {
Error::InternalErr(format!(
"error while setting leaf jobs: {e}"
))
})?;
}
}
}
@@ -600,7 +652,7 @@ pub async fn update_flow_status_after_job_completion_internal<
.root_job
.map(|x| x.to_string())
.unwrap_or_else(|| "none".to_string());
tracing::info!(id = %flow_job.id, root_id = %job_root, "update flow status");
tracing::info!(id = %flow_job.id, root_id = %job_root, worker_name = %worker_name, "update flow status");
let module = get_module(&flow_job, module_index);
@@ -795,6 +847,8 @@ pub async fn update_flow_status_after_job_completion_internal<
}
if let Some(parent_job) = flow_job.parent_job {
tracing::info!(subflow_id = %flow_job.id, parent_id = %parent_job, worker_name = %worker_name, "subflow is finished, updating parent flow status");
return Ok(Some(RecUpdateFlowStatusAfterJobCompletion {
flow: parent_job,
job_id_for_status: flow,