remove denogobuncache

This commit is contained in:
Ruben Fiszel
2024-01-23 19:54:01 +01:00
parent 428c3161ad
commit 4b4effe767
3 changed files with 288 additions and 404 deletions

View File

@@ -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<()> {

View File

@@ -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,
};

View File

@@ -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::<u64>().ok());
pub static ref CAN_PULL: Arc<RwLock<()>> = Arc::new(RwLock::new(()));
pub static ref WORKER_EXECUTION_COUNT: Arc<RwLock<HashMap<String, IntCounter>>> = Arc::new(RwLock::new(HashMap::new()));
pub static ref WORKER_EXECUTION_DURATION_COUNTER: Arc<RwLock<HashMap<String, prometheus::Counter>>> = Arc::new(RwLock::new(HashMap::new()));
@@ -855,18 +851,6 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
ws.inc();
}
let (_copy_to_bucket_tx, mut copy_to_bucket_rx) = mpsc::channel::<()>(2);
#[cfg(feature = "enterprise")]
let mut copy_cache_from_bucket_handle: Option<tokio::task::JoinHandle<()>> = 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<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
tracing::warn!("S3 cache not available in the pro plan");
} else if crate::global_cache::worker_s3_bucket_sync_enabled(&db).await {
let bucket = s.to_string();
let worker_name2 = worker_name.clone();
//piptars can be fetched in background
handles.push(tokio::task::spawn(async move {
tracing::info!(worker = %worker_name2, "Started initial piptar sync in background");
copy_all_piptars_from_bucket(&bucket).await;
}));
//denogocache.tar need to be fetched in foreground, block workers until they fetched it
copy_denogo_cache_from_bucket_as_tar(s).await;
copy_all_piptars_from_bucket(&bucket).await;
if let Err(e) = untar_all_piptars().await {
tracing::error!("Failed to untar pip tarballs: {:?}", e);
}
}
}
}
@@ -1298,9 +1278,6 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
);
}
#[cfg(feature = "enterprise")]
let copy_tx = _copy_to_bucket_tx.clone();
if last_ping.elapsed().as_secs() > NUM_SECS_PING {
let tags = WORKER_CONFIG.read().await.worker_tags.clone();
@@ -1328,60 +1305,6 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
tracing::info!(worker = %worker_name, "vacuumed queue and completed_job");
}
#[cfg(feature = "enterprise")]
if i_worker == 1 && S3_CACHE_BUCKET.is_some() {
if matches!(get_license_plan().await, LicensePlan::Pro) {
tracing::warn!("S3 cache not available in the pro plan");
} else if last_sync.elapsed().as_secs() > *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<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
tokio::select! {
biased;
_ = killpill_rx.recv() => {
#[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)