From 62d196ecece158f47d52cdbe11089f57992df589 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Wed, 19 Apr 2023 15:16:44 +0200 Subject: [PATCH] add barrier when num workers > 1 --- backend/src/main.rs | 17 +++-- backend/tests/worker.rs | 2 +- backend/windmill-common/src/lib.rs | 4 +- backend/windmill-worker/src/global_cache.rs | 24 +++++-- backend/windmill-worker/src/lib.rs | 1 + backend/windmill-worker/src/worker.rs | 69 +++++++++++++++------ 6 files changed, 88 insertions(+), 29 deletions(-) diff --git a/backend/src/main.rs b/backend/src/main.rs index fba806f435..51c9dd6ae0 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -14,9 +14,9 @@ use std::{ use git_version::git_version; use monitor::handle_zombie_jobs_periodically; use sqlx::{Pool, Postgres}; -use tokio::sync::RwLock; +use tokio::{fs::DirBuilder, sync::RwLock}; use windmill_common::{utils::rd_string, IS_READY, METRICS_ADDR}; -use windmill_worker::S3_CACHE_BUCKET; +use windmill_worker::{DENO_CACHE_DIR, GO_CACHE_DIR, PIP_CACHE_DIR, S3_CACHE_BUCKET}; const GIT_VERSION: &str = git_version!(args = ["--tag", "--always"], fallback = "unknown-version"); const DEFAULT_NUM_WORKERS: usize = 3; @@ -271,6 +271,14 @@ pub async fn run_workers = Arc::new(std::sync::atomic::AtomicBool::new(false)); + pub static ref IS_READY: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false); } #[cfg(feature = "tokio")] diff --git a/backend/windmill-worker/src/global_cache.rs b/backend/windmill-worker/src/global_cache.rs index d906cdf70e..2929e4b7ee 100644 --- a/backend/windmill-worker/src/global_cache.rs +++ b/backend/windmill-worker/src/global_cache.rs @@ -45,7 +45,9 @@ pub async fn copy_cache_from_bucket(bucket: &str, tx: Sender<()>) -> error::Resu "--size-only", "--fast-list", "--exclude", - &format!("\"/{TAR_CACHE_FILENAME},/deno/gen/file/**\""), + &format!("deno/gen/file/tmp/windmill/**"), + "--exclude", + &format!("{TAR_CACHE_FILENAME}"), ], ) .await @@ -79,7 +81,9 @@ pub async fn copy_cache_to_bucket(bucket: &str) -> error::Result<()> { "--size-only", "--fast-list", "--exclude", - &format!("\"/{TAR_CACHE_FILENAME},/deno/gen/file/**\""), + &format!("deno/gen/file/tmp/windmill/**"), + "--exclude", + &format!("{TAR_CACHE_FILENAME}"), ], ) .await @@ -251,7 +255,13 @@ pub async fn copy_tmp_cache_to_cache() -> error::Result<()> { execute_command( TMP_DIR, "rclone", - vec!["sync", ROOT_TMP_CACHE_DIR, ROOT_CACHE_DIR], + vec![ + "sync", + ROOT_TMP_CACHE_DIR, + ROOT_CACHE_DIR, + "--exclude", + TAR_CACHE_FILENAME, + ], ) .await?; tracing::info!( @@ -267,7 +277,13 @@ pub async fn copy_cache_to_tmp_cache() -> error::Result<()> { execute_command( TMP_DIR, "rclone", - vec!["sync", ROOT_CACHE_DIR, ROOT_TMP_CACHE_DIR], + vec![ + "sync", + ROOT_CACHE_DIR, + ROOT_TMP_CACHE_DIR, + "--exclude", + TAR_CACHE_FILENAME, + ], ) .await?; tracing::info!( diff --git a/backend/windmill-worker/src/lib.rs b/backend/windmill-worker/src/lib.rs index ac0fdea068..597e5f4f5e 100644 --- a/backend/windmill-worker/src/lib.rs +++ b/backend/windmill-worker/src/lib.rs @@ -7,5 +7,6 @@ mod python_executor; mod worker; mod worker_flow; +#[cfg(feature = "enterprise")] pub use global_cache::copy_cache_from_bucket_as_tar; pub use worker::*; diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index f26692c71c..733312b8f5 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -232,6 +232,8 @@ lazy_static::lazy_static! { .map(|e| Some(e)) .unwrap_or(None); + pub static ref CAN_PULL: Arc> = Arc::new(RwLock::new(())); + } //only matter if CLOUD_HOSTED @@ -284,7 +286,7 @@ pub async fn run_worker, base_internal_url: &str, rsmq: Option, - sync_barrier: RwLock>>, + sync_barrier: Arc>>, ) { #[cfg(not(feature = "enterprise"))] if !*DISABLE_NSJAIL { @@ -298,13 +300,11 @@ pub async fn run_worker *GLOBAL_CACHE_INTERVAL && (copy_cache_from_bucket_handle.is_none() || copy_cache_from_bucket_handle.as_ref().unwrap().is_finished()) { + + tracing::debug!("CAN PULL LOCK START"); + let _lock = CAN_PULL.write().await; + tracing::info!("Started syncing cache"); last_sync = Instant::now(); + if num_workers > 1 { + create_barrier_for_all_workers(num_workers, sync_barrier.clone()).await; + } if let Err(e) = copy_cache_to_tmp_cache().await { tracing::error!("failed to copy cache to tmp cache: {}", e); } else { @@ -484,12 +491,20 @@ pub async fn run_worker 1 { + if num_workers > 1 && S3_CACHE_BUCKET.is_some() { let read_barrier = sync_barrier.read().await; - let barrier = read_barrier.clone(); - if let Some(b) = barrier.as_ref() { + if let Some(b) = read_barrier.as_ref() { + tracing::debug!("worker #{i_worker} waiting for barrier"); b.wait().await; + tracing::debug!("worker #{i_worker} done waiting for barrier"); + drop(read_barrier); + // wait for barrier to be reset + let _ = CAN_PULL.read().await; + tracing::debug!("worker #{i_worker} done waiting for lock"); + } else { + tracing::debug!("worker #{i_worker} no barrier"); }; } @@ -509,19 +524,17 @@ pub async fn run_worker { + tracing::debug!("CAN PULL LOCK START"); + let _lock = CAN_PULL.write().await; if num_workers > 1 { - let mut barrier = sync_barrier.write().await; - let arc_barrier = Arc::new(Some(tokio::sync::Barrier::new(num_workers as usize))); - *barrier = arc_barrier.clone(); - if let Some(b) = arc_barrier.as_ref() { - b.wait().await; - }; + create_barrier_for_all_workers(num_workers, sync_barrier.clone()).await; } //Arc::new(tokio::sync::Barrier::new(num_workers as usize + 1)); #[cfg(feature = "enterprise")] if let Err(e) = copy_tmp_cache_to_cache().await { tracing::error!(worker = %worker_name, "failed to sync tmp cache to cache: {}", e); } + tracing::debug!("CAN PULL LOCK END"); (false, Ok(None)) }, Some(job_id) = same_worker_rx.recv() => { @@ -580,7 +593,12 @@ pub async fn run_worker>>) { + tracing::debug!("acquiring write lock"); + let mut barrier = sync_barrier.write().await; + *barrier = Some(tokio::sync::Barrier::new(num_workers as usize)); + drop(barrier); + tracing::debug!("dropped write lock"); + if let Some(b) = sync_barrier.read().await.as_ref() { + tracing::debug!("leader worker waiting for barrier"); + b.wait().await; + tracing::debug!("leader worker done waiting for barrier"); + }; + let mut barrier = sync_barrier.write().await; + *barrier = None; + tracing::debug!("leader worker done waiting for"); +} pub async fn handle_job_error( db: &Pool, client: &AuthedClient,