From e6d67f4e5994360e7ae83c1303b31514b6062d1d Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Wed, 29 Nov 2023 12:28:00 +0100 Subject: [PATCH] fix: make dedicated workers able to redeploy automatically --- backend/windmill-worker/src/worker.rs | 178 +++++++++++--------------- 1 file changed, 76 insertions(+), 102 deletions(-) diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 9ba5c6cacb..0f1c5f13a5 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -1105,79 +1105,71 @@ pub async fn run_worker(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::(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::>(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, Option, Option>)>( + 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, Option, Option>)>( "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), };