diff --git a/backend/src/main.rs b/backend/src/main.rs index 2dbea11946..a56da225d9 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -22,7 +22,8 @@ use tokio::{ use windmill_api::{LICENSE_KEY, OAUTH_CLIENTS, SMTP_CLIENT}; use windmill_common::{utils::rd_string, METRICS_ADDR}; use windmill_worker::{ - BUN_CACHE_DIR, BUN_TMP_CACHE_DIR, DENO_CACHE_DIR, DENO_TMP_CACHE_DIR, GO_BIN_CACHE_DIR, + BUN_CACHE_DIR, BUN_TMP_CACHE_DIR, DENO_CACHE_DIR, DENO_CACHE_DIR_DEPS, DENO_CACHE_DIR_NPM, + DENO_TMP_CACHE_DIR, DENO_TMP_CACHE_DIR_DEPS, DENO_TMP_CACHE_DIR_NPM, GO_BIN_CACHE_DIR, GO_CACHE_DIR, GO_TMP_CACHE_DIR, HUB_CACHE_DIR, HUB_TMP_CACHE_DIR, LOCK_CACHE_DIR, PIP_CACHE_DIR, ROOT_TMP_CACHE_DIR, TAR_PIP_TMP_CACHE_DIR, }; @@ -343,12 +344,16 @@ pub async fn run_workers error::Result<()> { "tar {folder_name} does not exist in bucket" ))); } + // tracing::info!("B: {target} {folder}"); extract_pip_tar(&target, &folder).await?; tokio::fs::remove_file(&target).await?; @@ -282,8 +283,6 @@ pub async fn copy_cache_to_bucket_as_tar(bucket: &str) { #[cfg(feature = "enterprise")] pub async fn copy_denogo_cache_from_bucket_as_tar(bucket: &str) { - use tokio::fs::metadata; - tracing::info!("Copying deno,go,bun cache from bucket {bucket} as tar"); let start: Instant = Instant::now(); @@ -307,7 +306,7 @@ pub async fn copy_denogo_cache_from_bucket_as_tar(bucket: &str) { } if let Err(e) = execute_command( - ROOT_CACHE_DIR, + ROOT_TMP_CACHE_DIR, "tar", vec![ "-xpvf", @@ -320,12 +319,12 @@ pub async fn copy_denogo_cache_from_bucket_as_tar(bucket: &str) { return; } - if let Err(e) = - tokio::fs::remove_file(format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}")).await - { - tracing::info!("Failed to remove denogobuntar cache. Error: {:?}", e); - return; - }; + // if let Err(e) = + // tokio::fs::remove_file(format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}")).await + // { + // tracing::info!("Failed to remove denogobuntar cache. Error: {:?}", e); + // return; + // }; tracing::info!( "Finished copying denogobuntar from bucket {bucket} as tar, took: {:?}s", @@ -441,6 +440,7 @@ pub async fn untar_all_piptars() -> error::Result<()> { if metadata(&folder).await.is_ok() { continue; } + // tracing::info!("A: {path} {folder}"); extract_pip_tar(&path, &folder).await?; Ok(()) as error::Result<()> } { @@ -506,6 +506,7 @@ pub async fn copy_cache_to_tmp_cache() -> error::Result<()> { #[cfg(feature = "enterprise")] pub async fn execute_command(dir: &str, command: &str, args: Vec<&str>) -> error::Result<()> { + tracing::info!("Executing command: {command} {}", args.iter().join(" ")); match Command::new(command) .current_dir(dir) .args(args.clone()) diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 8fec2d31ce..27454b272d 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -135,17 +135,22 @@ pub const ROOT_TMP_CACHE_DIR: &str = "/tmp/windmill/tmpcache/"; pub const LOCK_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "lock"); pub const PIP_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "pip"); pub const DENO_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "deno"); +pub const DENO_CACHE_DIR_DEPS: &str = concatcp!(ROOT_CACHE_DIR, "deno/deps"); +pub const DENO_CACHE_DIR_NPM: &str = concatcp!(ROOT_CACHE_DIR, "deno/npm"); + pub const GO_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "go"); pub const BUN_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "bun"); pub const HUB_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "hub"); pub const GO_BIN_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "gobin"); + pub const TAR_PIP_TMP_CACHE_DIR: &str = concatcp!(ROOT_TMP_CACHE_DIR, "tar/pip"); pub const DENO_TMP_CACHE_DIR: &str = concatcp!(ROOT_TMP_CACHE_DIR, "deno"); pub const BUN_TMP_CACHE_DIR: &str = concatcp!(ROOT_TMP_CACHE_DIR, "bun"); pub const GO_TMP_CACHE_DIR: &str = concatcp!(ROOT_TMP_CACHE_DIR, "go"); pub const HUB_TMP_CACHE_DIR: &str = concatcp!(ROOT_TMP_CACHE_DIR, "hub"); - +pub const DENO_TMP_CACHE_DIR_DEPS: &str = concatcp!(ROOT_TMP_CACHE_DIR, "deno/deps"); +pub const DENO_TMP_CACHE_DIR_NPM: &str = concatcp!(ROOT_TMP_CACHE_DIR, "deno/npm"); const NUM_SECS_PING: u64 = 5; @@ -308,12 +313,12 @@ pub async fn run_worker, base_internal_url: &str, rsmq: Option, - sync_barrier: Arc>>, + _sync_barrier: Arc>>, ) { #[cfg(not(feature = "enterprise"))] if !*DISABLE_NSJAIL { @@ -529,9 +534,9 @@ pub async fn run_worker 1 { - create_barrier_for_all_workers(num_workers, sync_barrier.clone()).await; - } + // 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 { @@ -543,25 +548,26 @@ pub async fn run_worker 1 && S3_CACHE_BUCKET.is_some() { - let read_barrier = sync_barrier.read().await; - 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"); - }; - } + // // The barrier is to avoid the sync to bucket syncing partial folders + // #[cfg(feature = "enterprise")] + // if num_workers > 1 && S3_CACHE_BUCKET.is_some() { + // let read_barrier = sync_barrier.read().await; + // 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"); + // }; + // } let (do_break, next_job) = if first_run { (false, Ok(Some(QueuedJob::default()))) @@ -594,9 +600,9 @@ pub async fn run_worker { tracing::debug!("can_pull lock start"); let _lock = CAN_PULL.write().await; - if num_workers > 1 { - create_barrier_for_all_workers(num_workers, sync_barrier.clone()).await; - } + // if num_workers > 1 { + // 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 { @@ -787,21 +793,21 @@ 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 create_barrier_for_all_workers(num_workers: u32, sync_barrier: Arc>>) { +// 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, diff --git a/python-client/wmill/wmill/client.py b/python-client/wmill/wmill/client.py index df735f93f9..44b6cd961e 100644 --- a/python-client/wmill/wmill/client.py +++ b/python-client/wmill/wmill/client.py @@ -334,7 +334,7 @@ def set_variable(path: str, value: str) -> None: def get_state_path() -> str: - state_path = os.environ.get("WM_STATE_PATH_NEW") + state_path = os.environ.get("WM_STATE_PATH_NEW") or os.environ.get("WM_STATE_PATH") if state_path is None: raise Exception("State path not found") return state_path