From 4b4effe767b2dbdcaabc3b70648062b3c35aed2d Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Tue, 23 Jan 2024 19:54:01 +0100 Subject: [PATCH] remove denogobuncache --- backend/windmill-worker/src/global_cache.rs | 576 +++++++++--------- .../windmill-worker/src/python_executor.rs | 3 +- backend/windmill-worker/src/worker.rs | 113 +--- 3 files changed, 288 insertions(+), 404 deletions(-) diff --git a/backend/windmill-worker/src/global_cache.rs b/backend/windmill-worker/src/global_cache.rs index 07c52c670b..4845e3187f 100644 --- a/backend/windmill-worker/src/global_cache.rs +++ b/backend/windmill-worker/src/global_cache.rs @@ -1,9 +1,9 @@ #[cfg(feature = "enterprise")] -use crate::{ROOT_CACHE_DIR, ROOT_TMP_CACHE_DIR, TAR_CACHE_RATE, TAR_PIP_TMP_CACHE_DIR, TMP_DIR}; +use crate::{ROOT_CACHE_DIR, ROOT_TMP_CACHE_DIR, TAR_PIP_TMP_CACHE_DIR}; #[cfg(feature = "enterprise")] use itertools::Itertools; -#[cfg(feature = "enterprise")] -use rand::Rng; +// #[cfg(feature = "enterprise")] +// use rand::Rng; #[cfg(feature = "enterprise")] use std::process::Stdio; #[cfg(feature = "enterprise")] @@ -13,13 +13,13 @@ use windmill_common::DB; use windmill_common::global_settings::WORKER_S3_BUCKET_SYNC; #[cfg(feature = "enterprise")] -use tokio::{process::Command, sync::mpsc::Sender, time::Instant}; +use tokio::{process::Command, time::Instant}; #[cfg(feature = "enterprise")] use windmill_common::error; -#[cfg(feature = "enterprise")] -const TAR_CACHE_FILENAME: &str = "denogocache.tar"; +// #[cfg(feature = "enterprise")] +// const TAR_CACHE_FILENAME: &str = "denogocache.tar"; #[cfg(feature = "enterprise")] pub async fn build_tar_and_push(bucket: &str, folder: String) -> error::Result<()> { @@ -126,103 +126,103 @@ pub async fn pull_from_tar(bucket: &str, folder: String) -> error::Result<()> { Ok(()) } -#[cfg(feature = "enterprise")] -pub async fn cache_global(bucket: &str, tx: Sender<()>) -> error::Result<()> { - copy_cache_from_bucket(bucket, tx).await?; - copy_cache_to_bucket(bucket).await?; +// #[cfg(feature = "enterprise")] +// pub async fn cache_global(bucket: &str, _tx: Sender<()>) -> error::Result<()> { +// // copy_cache_from_bucket(bucket, tx).await?; +// // copy_cache_to_bucket(bucket).await?; - // this is to prevent excessive tar upload. 1/100*15min = each worker sync its tar once per day on average - if rand::thread_rng().gen_range(0..*TAR_CACHE_RATE) == 0 { - copy_cache_to_bucket_as_tar(bucket).await; - } - Ok(()) -} +// // this is to prevent excessive tar upload. 1/100*15min = each worker sync its tar once per day on average +// if rand::thread_rng().gen_range(0..*TAR_CACHE_RATE) == 0 { +// copy_cache_to_bucket_as_tar(bucket).await; +// } +// Ok(()) +// } -#[cfg(feature = "enterprise")] -pub async fn copy_cache_from_bucket(bucket: &str, tx: Sender<()>) -> error::Result<()> { - tracing::info!("Copying cache from bucket in the background {bucket}"); - let bucket = bucket.to_string(); +// #[cfg(feature = "enterprise")] +// pub async fn copy_cache_from_bucket(bucket: &str, tx: Sender<()>) -> error::Result<()> { +// tracing::info!("Copying cache from bucket in the background {bucket}"); +// let bucket = bucket.to_string(); - let start = Instant::now(); +// let start = Instant::now(); - if let Err(e) = execute_command( - ROOT_TMP_CACHE_DIR, - "rclone", - vec![ - "copy", - &format!(":s3,env_auth=true:{bucket}"), - &ROOT_TMP_CACHE_DIR, - // "-l", - "--size-only", - "--fast-list", - "--filter", - "+ deno/npm/**", - "--filter", - "+ deno/deps/**", - // "--filter", - // "+ bun/**", - "--filter", - "+ go/**", - "--filter", - "+ tar/**", - "--filter", - "- *", - ], - ) - .await - { - tracing::info!("Failed to copy cache from bucket. Error: {:?}", e); - return Err(e); - } +// if let Err(e) = execute_command( +// ROOT_TMP_CACHE_DIR, +// "rclone", +// vec![ +// "copy", +// &format!(":s3,env_auth=true:{bucket}"), +// &ROOT_TMP_CACHE_DIR, +// // "-l", +// "--size-only", +// "--fast-list", +// "--filter", +// "+ deno/npm/**", +// "--filter", +// "+ deno/deps/**", +// // "--filter", +// // "+ bun/**", +// "--filter", +// "+ go/**", +// "--filter", +// "+ tar/**", +// "--filter", +// "- *", +// ], +// ) +// .await +// { +// tracing::info!("Failed to copy cache from bucket. Error: {:?}", e); +// return Err(e); +// } - tracing::info!( - "Finished copying cache from bucket {bucket}, took {:?}s", - start.elapsed().as_secs() - ); +// tracing::info!( +// "Finished copying cache from bucket {bucket}, took {:?}s", +// start.elapsed().as_secs() +// ); - tx.send(()).await.expect("can send copy cache signal"); +// tx.send(()).await.expect("can send copy cache signal"); - Ok(()) -} +// Ok(()) +// } -#[cfg(feature = "enterprise")] -pub async fn copy_cache_to_bucket(bucket: &str) -> error::Result<()> { - tracing::info!("Copying cache to bucket {bucket}"); - let start = Instant::now(); +// #[cfg(feature = "enterprise")] +// pub async fn copy_cache_to_bucket(bucket: &str) -> error::Result<()> { +// tracing::info!("Copying cache to bucket {bucket}"); +// let start = Instant::now(); - if let Err(e) = execute_command( - ROOT_TMP_CACHE_DIR, - "rclone", - vec![ - "copy", - &ROOT_TMP_CACHE_DIR, - &format!(":s3,env_auth=true:{bucket}"), - // "-l", - "--size-only", - "--fast-list", - "--filter", - "+ deno/npm/**", - "--filter", - "+ deno/deps/**", - // "--filter", - // "+ bun/**", - "--filter", - "+ go/**", - "--filter", - "- *", - ], - ) - .await - { - tracing::info!("Failed to copy cache to bucket. Error: {:?}", e); - return Err(e); - } - tracing::info!( - "Finished copying cache to bucket {bucket}, took: {:?}s", - start.elapsed().as_secs() - ); - Ok(()) -} +// if let Err(e) = execute_command( +// ROOT_TMP_CACHE_DIR, +// "rclone", +// vec![ +// "copy", +// &ROOT_TMP_CACHE_DIR, +// &format!(":s3,env_auth=true:{bucket}"), +// // "-l", +// "--size-only", +// "--fast-list", +// "--filter", +// "+ deno/npm/**", +// "--filter", +// "+ deno/deps/**", +// // "--filter", +// // "+ bun/**", +// "--filter", +// "+ go/**", +// "--filter", +// "- *", +// ], +// ) +// .await +// { +// tracing::info!("Failed to copy cache to bucket. Error: {:?}", e); +// return Err(e); +// } +// tracing::info!( +// "Finished copying cache to bucket {bucket}, took: {:?}s", +// start.elapsed().as_secs() +// ); +// Ok(()) +// } #[cfg(feature = "enterprise")] pub async fn worker_s3_bucket_sync_enabled(db: &DB) -> bool { @@ -243,145 +243,145 @@ pub async fn worker_s3_bucket_sync_enabled(db: &DB) -> bool { } } -#[cfg(feature = "enterprise")] -pub async fn copy_cache_to_bucket_as_tar(bucket: &str) { - tracing::info!("Copying cache to bucket {bucket} as tar"); - let start = Instant::now(); +// #[cfg(feature = "enterprise")] +// pub async fn copy_cache_to_bucket_as_tar(bucket: &str) { +// tracing::info!("Copying cache to bucket {bucket} as tar"); +// let start = Instant::now(); - if let Err(e) = execute_command( - ROOT_TMP_CACHE_DIR, - "tar", - vec![ - "-c", - "-f", - &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), - "go", - "deno/npm", - "deno/deps", // "bun", - ], - ) - .await - { - tracing::info!("Failed to tar cache. Error: {:?}", e); - return; - } +// if let Err(e) = execute_command( +// ROOT_TMP_CACHE_DIR, +// "tar", +// vec![ +// "-c", +// "-f", +// &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), +// "go", +// "deno/npm", +// "deno/deps", // "bun", +// ], +// ) +// .await +// { +// tracing::info!("Failed to tar cache. Error: {:?}", e); +// return; +// } - let tar_metadata = - tokio::fs::metadata(format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}")).await; - if tar_metadata.is_err() || tar_metadata.as_ref().unwrap().len() == 0 { - tracing::info!("Failed to tar cache"); - return; - } +// let tar_metadata = +// tokio::fs::metadata(format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}")).await; +// if tar_metadata.is_err() || tar_metadata.as_ref().unwrap().len() == 0 { +// tracing::info!("Failed to tar cache"); +// return; +// } - if let Err(e) = execute_command( - ROOT_TMP_CACHE_DIR, - "rclone", - vec![ - "copyto", - &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), - &format!(":s3,env_auth=true:{bucket}/{TAR_CACHE_FILENAME}"), - "-v", - "--size-only", - "--fast-list", - "--s3-no-check-bucket", - ], - ) - .await - { - tracing::info!("Failed to copy tar to bucket. Error: {:?}", e); - return; - } +// if let Err(e) = execute_command( +// ROOT_TMP_CACHE_DIR, +// "rclone", +// vec![ +// "copyto", +// &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), +// &format!(":s3,env_auth=true:{bucket}/{TAR_CACHE_FILENAME}"), +// "-v", +// "--size-only", +// "--fast-list", +// "--s3-no-check-bucket", +// ], +// ) +// .await +// { +// tracing::info!("Failed to copy tar to bucket. 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 tar cache. Error: {:?}", e); - }; +// if let Err(e) = +// tokio::fs::remove_file(format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}")).await +// { +// tracing::info!("Failed to remove tar cache. Error: {:?}", e); +// }; - tracing::info!( - "Finished copying cache to bucket {bucket} as tar, took: {:?}s. Size of new tar: {}", - start.elapsed().as_secs(), - tar_metadata.unwrap().len() - ); -} +// tracing::info!( +// "Finished copying cache to bucket {bucket} as tar, took: {:?}s. Size of new tar: {}", +// start.elapsed().as_secs(), +// tar_metadata.unwrap().len() +// ); +// } -#[cfg(feature = "enterprise")] -pub async fn copy_denogo_cache_from_bucket_as_tar(bucket: &str) { - tracing::info!("Copying deno,go,bun cache from bucket {bucket} as tar"); +// #[cfg(feature = "enterprise")] +// pub async fn copy_denogo_cache_from_bucket_as_tar(bucket: &str) { +// tracing::info!("Copying deno,go,bun cache from bucket {bucket} as tar"); - let mut start: Instant = Instant::now(); +// let mut start: Instant = Instant::now(); - if let Err(e) = execute_command( - ROOT_TMP_CACHE_DIR, - "rclone", - vec![ - "copyto", - &format!(":s3,env_auth=true:{bucket}/{TAR_CACHE_FILENAME}"), - &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), - "-v", - "--size-only", - "--fast-list", - ], - ) - .await - { - tracing::info!("Failed copying deno,go,bun tar from cache. Error: {:?}", e); - return; - } +// if let Err(e) = execute_command( +// ROOT_TMP_CACHE_DIR, +// "rclone", +// vec![ +// "copyto", +// &format!(":s3,env_auth=true:{bucket}/{TAR_CACHE_FILENAME}"), +// &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), +// "-v", +// "--size-only", +// "--fast-list", +// ], +// ) +// .await +// { +// tracing::info!("Failed copying deno,go,bun tar from cache. Error: {:?}", e); +// return; +// } - tracing::info!( - "Finished copying denogobun tar for from bucket as tar. took {}s", - start.elapsed().as_secs() - ); +// tracing::info!( +// "Finished copying denogobun tar for from bucket as tar. took {}s", +// start.elapsed().as_secs() +// ); - start = Instant::now(); +// start = Instant::now(); - if let Err(e) = execute_command( - ROOT_CACHE_DIR, - "tar", - vec![ - "-xpvf", - &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), - ], - ) - .await - { - tracing::info!("Failed to untar denogobun tar to cache. Error: {:?}", e); - return; - } +// if let Err(e) = execute_command( +// ROOT_CACHE_DIR, +// "tar", +// vec![ +// "-xpvf", +// &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), +// ], +// ) +// .await +// { +// tracing::info!("Failed to untar denogobun tar to cache. Error: {:?}", e); +// return; +// } - tracing::info!( - "Finished untaring denogobun tar to cache. took: {}s", - start.elapsed().as_secs() - ); +// tracing::info!( +// "Finished untaring denogobun tar to cache. took: {}s", +// start.elapsed().as_secs() +// ); - start = Instant::now(); +// start = Instant::now(); - if let Err(e) = execute_command( - ROOT_TMP_CACHE_DIR, - "tar", - vec![ - "-xpvf", - &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), - ], - ) - .await - { - tracing::info!("Failed to untar denogobun tar to tmpcache. Error: {:?}", e); - return; - } +// if let Err(e) = execute_command( +// ROOT_TMP_CACHE_DIR, +// "tar", +// vec![ +// "-xpvf", +// &format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"), +// ], +// ) +// .await +// { +// tracing::info!("Failed to untar denogobun tar to tmpcache. Error: {:?}", e); +// return; +// } - tracing::info!( - "Finished untaring denogobun tar to /tmpcache. took: {}s", - start.elapsed().as_secs() - ); +// tracing::info!( +// "Finished untaring denogobun tar to /tmpcache. took: {}s", +// start.elapsed().as_secs() +// ); - 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); - }; -} +// 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); +// }; +// } #[cfg(feature = "enterprise")] pub async fn copy_all_piptars_from_bucket(bucket: &str) { @@ -413,57 +413,43 @@ pub async fn copy_all_piptars_from_bucket(bucket: &str) { ); } -// async fn check_if_bucket_syncable(bucket: &str) -> bool { -// match Command::new("rclone") -// .arg("lsf") -// .arg(format!(":s3,env_auth=true:{bucket}/NOSYNC")) +// #[cfg(feature = "enterprise")] +// pub async fn copy_tmp_cache_to_cache() -> error::Result<()> { +// let start: Instant = Instant::now(); +// execute_command( +// TMP_DIR, +// "rclone", +// vec![ +// "sync", +// ROOT_TMP_CACHE_DIR, +// ROOT_CACHE_DIR, +// // "-l", +// "--filter", +// "+ deno/npm/**", +// "--filter", +// "+ deno/deps/**", +// // "--filter", +// // "+ bun/**", +// "--filter", +// "+ go/**", +// "--filter", +// "- *", +// ], +// ) +// .await?; -// .arg("-vv") -// .arg("--fast-list") -// .stdin(Stdio::null()) -// .stdout(Stdio::null()) -// .output() -// .await; -// return true; +// tracing::info!( +// "Finished copying local tmp cache to local cache. Took {}ms", +// start.elapsed().as_millis(), +// ); + +// if let Err(e) = untar_all_piptars().await { +// tracing::info!("Failed to untar piptars. Error: {:?}", e); +// } + +// Ok(()) // } -#[cfg(feature = "enterprise")] -pub async fn copy_tmp_cache_to_cache() -> error::Result<()> { - let start: Instant = Instant::now(); - execute_command( - TMP_DIR, - "rclone", - vec![ - "sync", - ROOT_TMP_CACHE_DIR, - ROOT_CACHE_DIR, - // "-l", - "--filter", - "+ deno/npm/**", - "--filter", - "+ deno/deps/**", - // "--filter", - // "+ bun/**", - "--filter", - "+ go/**", - "--filter", - "- *", - ], - ) - .await?; - - tracing::info!( - "Finished copying local tmp cache to local cache. Took {}ms", - start.elapsed().as_millis(), - ); - - if let Err(e) = untar_all_piptars().await { - tracing::info!("Failed to untar piptars. Error: {:?}", e); - } - - Ok(()) -} - #[cfg(feature = "enterprise")] pub async fn untar_all_piptars() -> error::Result<()> { use tokio::fs::{self, metadata}; @@ -524,36 +510,36 @@ pub async fn extract_pip_tar(tar: &str, folder: &str) -> error::Result<()> { Ok(()) } -#[cfg(feature = "enterprise")] -pub async fn copy_cache_to_tmp_cache() -> error::Result<()> { - let start: Instant = Instant::now(); - execute_command( - TMP_DIR, - "rclone", - vec![ - "sync", - ROOT_CACHE_DIR, - ROOT_TMP_CACHE_DIR, - // "-l", - "--filter", - "+ deno/npm/**", - "--filter", - "+ deno/deps/**", - // "--filter", - // "+ bun/**", - "--filter", - "+ go/**", - "--filter", - "- *", - ], - ) - .await?; - tracing::info!( - "Finished copying local cache to local tmp cache. Took {}ms", - start.elapsed().as_millis() - ); - Ok(()) -} +// #[cfg(feature = "enterprise")] +// pub async fn copy_cache_to_tmp_cache() -> error::Result<()> { +// let start: Instant = Instant::now(); +// execute_command( +// TMP_DIR, +// "rclone", +// vec![ +// "sync", +// ROOT_CACHE_DIR, +// ROOT_TMP_CACHE_DIR, +// // "-l", +// "--filter", +// "+ deno/npm/**", +// "--filter", +// "+ deno/deps/**", +// // "--filter", +// // "+ bun/**", +// "--filter", +// "+ go/**", +// "--filter", +// "- *", +// ], +// ) +// .await?; +// tracing::info!( +// "Finished copying local cache to local tmp cache. Took {}ms", +// start.elapsed().as_millis() +// ); +// Ok(()) +// } #[cfg(feature = "enterprise")] pub async fn execute_command(dir: &str, command: &str, args: Vec<&str>) -> error::Result<()> { diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 249e4d1a76..3e60949241 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -44,7 +44,7 @@ const NSJAIL_CONFIG_RUN_PYTHON3_CONTENT: &str = include_str!("../nsjail/run.pyth const RELATIVE_PYTHON_LOADER: &str = include_str!("../loader.py"); #[cfg(feature = "enterprise")] -use crate::global_cache::{build_tar_and_push, pull_from_tar}; +use crate::global_cache::pull_from_tar; #[cfg(feature = "enterprise")] use crate::S3_CACHE_BUCKET; @@ -54,6 +54,7 @@ use crate::{ create_args_and_out_file, get_reserved_variables, handle_child, read_result, set_logs, start_child_process, write_file, }, + global_cache::build_tar_and_push, AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, HTTPS_PROXY, HTTP_PROXY, LOCK_CACHE_DIR, NO_PROXY, NSJAIL_PATH, PATH_ENV, PIP_CACHE_DIR, PIP_EXTRA_INDEX_URL, TZ_ENV, }; diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index f95eeb229b..55d8de9d1b 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -68,10 +68,7 @@ use async_recursion::async_recursion; use rand::Rng; #[cfg(feature = "enterprise")] -use crate::global_cache::{ - cache_global, copy_all_piptars_from_bucket, copy_cache_to_tmp_cache, - copy_denogo_cache_from_bucket_as_tar, copy_tmp_cache_to_cache, -}; +use crate::global_cache::{copy_all_piptars_from_bucket, untar_all_piptars}; use windmill_queue::{add_completed_job, add_completed_job_error}; @@ -306,7 +303,6 @@ lazy_static::lazy_static! { .ok() .and_then(|x| x.parse::().ok()); - pub static ref CAN_PULL: Arc> = Arc::new(RwLock::new(())); pub static ref WORKER_EXECUTION_COUNT: Arc>> = Arc::new(RwLock::new(HashMap::new())); pub static ref WORKER_EXECUTION_DURATION_COUNTER: Arc>> = Arc::new(RwLock::new(HashMap::new())); @@ -855,18 +851,6 @@ pub async fn run_worker(2); - - #[cfg(feature = "enterprise")] - let mut copy_cache_from_bucket_handle: Option> = None; - - #[cfg(feature = "enterprise")] - let mut last_sync = Instant::now() - + Duration::from_secs(rand::thread_rng().gen_range(0..*GLOBAL_CACHE_INTERVAL)); - - #[cfg(feature = "enterprise")] - let mut handles = Vec::with_capacity(2); - #[cfg(feature = "enterprise")] if i_worker == 1 { if let Some(ref s) = S3_CACHE_BUCKET.clone() { @@ -874,15 +858,11 @@ pub async fn run_worker NUM_SECS_PING { let tags = WORKER_CONFIG.read().await.worker_tags.clone(); @@ -1328,60 +1305,6 @@ 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()) - { - last_sync = Instant::now(); - - if crate::global_cache::worker_s3_bucket_sync_enabled(&db).await { - tracing::debug!("CAN PULL LOCK START"); - let _lock = CAN_PULL.write().await; - - tracing::info!("Started syncing cache"); - // 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 { - copy_cache_from_bucket_handle = Some(tokio::task::spawn(async move { - if let Some(ref s) = S3_CACHE_BUCKET.clone() { - if let Err(e) = cache_global(s, copy_tx).await { - tracing::error!("failed to sync cache: {}", e); - } - } - })); - } - tracing::info!("Ended syncing cache sync part"); - } - } - } - - // // 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 next_job = { // println!("2: {:?}", instant.elapsed()); #[cfg(feature = "benchmark")] @@ -1392,36 +1315,10 @@ pub async fn run_worker { - #[cfg(feature = "enterprise")] - if let Some(copy_cache_from_bucket_handle) = copy_cache_from_bucket_handle.as_ref() { - if !copy_cache_from_bucket_handle.is_finished() { - copy_cache_from_bucket_handle.abort(); - } - } - #[cfg(feature = "enterprise")] - for handle in &handles { - if !handle.is_finished() { - handle.abort(); - } - } println!("received killpill for worker {}", i_worker); job_completed_tx.0.send(SendResult::Kill).await.unwrap(); break }, - _ = copy_to_bucket_rx.recv() => { - 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; - // } - //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"); - Ok(None) - }, Some(job_id) = same_worker_rx.recv() => { sqlx::query_as::<_, QueuedJob>("SELECT * FROM queue WHERE id = $1") .bind(job_id)