fix: make dedicated workers able to redeploy automatically

This commit is contained in:
Ruben Fiszel
2023-11-29 12:28:06 +01:00
parent 1a1d1db96f
commit e6d67f4e59
+76 -102
View File
@@ -1105,79 +1105,71 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
let is_flow_worker;
if let Some(flow_path) = _wp.path.strip_prefix("flow/") {
is_flow_worker = true;
loop {
let value = sqlx::query_scalar!(
"SELECT value FROM flow WHERE path = $1 AND workspace_id = $2",
flow_path,
_wp.workspace_id
)
.fetch_optional(db)
.await;
if let Ok(v) = value {
if let Some(v) = v {
let value = serde_json::from_value::<FlowValue>(v).map_err(|err| {
Error::InternalErr(format!(
"could not convert json to flow for {flow_path}: {err:?}"
))
let value = sqlx::query_scalar!(
"SELECT value FROM flow WHERE path = $1 AND workspace_id = $2",
flow_path,
_wp.workspace_id
)
.fetch_optional(db)
.await;
if let Ok(v) = value {
if let Some(v) = v {
let value = serde_json::from_value::<FlowValue>(v).map_err(|err| {
Error::InternalErr(format!(
"could not convert json to flow for {flow_path}: {err:?}"
))
});
if let Ok(flow) = value {
let workers = spawn_dedicated_workers_for_flow(
&flow.modules,
&_wp.path,
&_wp.workspace_id,
killpill_tx.clone(),
&killpill_rx,
db,
&worker_dir,
base_internal_url,
&worker_name,
&job_completed_tx,
)
.await;
workers.into_iter().for_each(|(path, sender, handle)| {
dedicated_handles.push(handle);
hm.insert(path, sender);
});
if let Ok(flow) = value {
let workers = spawn_dedicated_workers_for_flow(
&flow.modules,
&_wp.path,
&_wp.workspace_id,
killpill_tx.clone(),
&killpill_rx,
db,
&worker_dir,
base_internal_url,
&worker_name,
&job_completed_tx,
)
.await;
workers.into_iter().for_each(|(path, sender, handle)| {
dedicated_handles.push(handle);
hm.insert(path, sender);
});
break;
}
} else {
tracing::error!(
"flow present but value not found for {}, waiting 10s",
flow_path
);
}
} else {
tracing::error!("flow not found for {}, waiting 10s,", flow_path);
tracing::error!(
"flow present but value not found for dedicated worker. {}",
flow_path
);
}
tokio::time::sleep(Duration::from_millis(10000)).await;
} else {
tracing::error!("flow not found for dedicated worker: {}. Waiting for dependency job and expected to restart.", flow_path);
}
} else {
is_flow_worker = false;
loop {
if let Some((path, sender, handle)) = spawn_dedicated_worker(
SpawnWorker::Script { path: _wp.path.clone(), hash: None },
&_wp.workspace_id,
killpill_tx.clone(),
&killpill_rx,
db,
&worker_dir,
base_internal_url,
&worker_name,
&job_completed_tx,
None,
)
.await
{
dedicated_handles.push(handle);
hm.insert(path, sender);
break;
} else {
tracing::error!(
"failed to spawn dedicated worker for {}, script found but not in a compatible language. Retrying 10s",
if let Some((path, sender, handle)) = spawn_dedicated_worker(
SpawnWorker::Script { path: _wp.path.clone(), hash: None },
&_wp.workspace_id,
killpill_tx.clone(),
&killpill_rx,
db,
&worker_dir,
base_internal_url,
&worker_name,
&job_completed_tx,
None,
)
.await
{
dedicated_handles.push(handle);
hm.insert(path, sender);
} else {
tracing::error!(
"failed to spawn dedicated worker for {}, script not found",
_wp.path
);
tokio::time::sleep(Duration::from_millis(10000)).await;
}
}
}
(hm, is_flow_worker)
@@ -1851,7 +1843,7 @@ async fn spawn_dedicated_worker(
{
let (dedicated_worker_tx, dedicated_worker_rx) =
mpsc::channel::<Arc<QueuedJob>>(MAX_BUFFERED_DEDICATED_JOBS);
let mut killpill_rx = killpill_rx.resubscribe();
let killpill_rx = killpill_rx.resubscribe();
let db = db.clone();
let base_internal_url = base_internal_url.to_string();
let worker_name = worker_name.to_string();
@@ -1870,16 +1862,12 @@ async fn spawn_dedicated_worker(
let (content, lock, language, envs) = match sw {
SpawnWorker::Script { path, hash } => {
let r;
loop {
let q = if let Some(hash) = hash {
get_script_content_by_hash(&hash, &w_id, &db).await.map(
|r: ContentReqLangEnvs| {
Some((r.content, r.lockfile, r.language, r.envs))
},
)
} else {
sqlx::query_as::<_, (String, Option<String>, Option<ScriptLang>, Option<Vec<String>>)>(
let q = if let Some(hash) = hash {
get_script_content_by_hash(&hash, &w_id, &db).await.map(
|r: ContentReqLangEnvs| Some((r.content, r.lockfile, r.language, r.envs)),
)
} else {
sqlx::query_as::<_, (String, Option<String>, Option<ScriptLang>, Option<Vec<String>>)>(
"SELECT content, lock, language, envs FROM script WHERE path = $1 AND workspace_id = $2 AND
created_at = (SELECT max(created_at) FROM script WHERE path = $1 AND workspace_id = $2 AND
deleted = false AND lock IS not NULL AND lock_error_logs IS NULL)",
@@ -1889,37 +1877,23 @@ async fn spawn_dedicated_worker(
.fetch_optional(&db)
.await
.map_err(|e| Error::InternalErr(format!("expected content and lock: {e}")))
};
if let Ok(q) = q {
if let Some(wp) = q {
r = wp;
break;
} else {
tracing::error!(
"Failed to fetch script `{}` in workspace {} for dedicated worker. Retrying in 10s.",
path,
w_id
);
tokio::select! {
biased;
_ = killpill_rx.recv() => {
tracing::info!("Killing dedicated worker while it was attempting to fetch script");
return None;
}
_ = tokio::time::sleep(Duration::from_secs(10)) => {
continue;
}
}
}
};
if let Ok(q) = q {
if let Some(wp) = q {
wp
} else {
tracing::error!("Failed to fetch script for dedicated worker");
killpill_tx.send(()).expect("send");
tracing::error!(
"Failed to fetch script `{}` in workspace {} for dedicated worker.",
path,
w_id
);
return None;
}
} else {
tracing::error!("Failed to fetch script for dedicated worker");
killpill_tx.send(()).expect("send");
return None;
}
r
}
SpawnWorker::RawScript { content, lock, lang, .. } => (content, lock, Some(lang), None),
};