feat(backend): global cache refactor for pip using tar for each dependency (#1443)

* cache refactor

* exclude tar from being synced to bucket

* run

* update

* update
This commit is contained in:
Ruben Fiszel
2023-04-20 20:05:12 +02:00
committed by GitHub
parent 3c98452f50
commit 369dd0dac6
5 changed files with 284 additions and 35 deletions
+30 -5
View File
@@ -14,9 +14,16 @@ use std::{
use git_version::git_version;
use monitor::handle_zombie_jobs_periodically;
use sqlx::{Pool, Postgres};
use tokio::{fs::DirBuilder, sync::RwLock};
use tokio::{
fs::{metadata, DirBuilder},
join,
sync::RwLock,
};
use windmill_common::{utils::rd_string, IS_READY, METRICS_ADDR};
use windmill_worker::{DENO_CACHE_DIR, GO_CACHE_DIR, PIP_CACHE_DIR, S3_CACHE_BUCKET};
use windmill_worker::{
DENO_CACHE_DIR, DENO_TMP_CACHE_DIR, GO_CACHE_DIR, GO_TMP_CACHE_DIR, PIP_CACHE_DIR,
ROOT_TMP_CACHE_DIR, S3_CACHE_BUCKET, TAR_PIP_TMP_CACHE_DIR,
};
const GIT_VERSION: &str = git_version!(args = ["--tag", "--always"], fallback = "unknown-version");
const DEFAULT_NUM_WORKERS: usize = 3;
@@ -271,7 +278,20 @@ pub async fn run_workers<R: rsmq_async::RsmqConnection + Send + Sync + Clone + '
let mut handles = Vec::with_capacity(num_workers as usize);
for x in [PIP_CACHE_DIR, DENO_CACHE_DIR, GO_CACHE_DIR] {
if metadata(&ROOT_TMP_CACHE_DIR).await.is_ok() {
if let Err(e) = tokio::fs::remove_dir_all(&ROOT_TMP_CACHE_DIR).await {
tracing::info!(error = %e, "Could not remove root tmp cache dir");
}
}
for x in [
PIP_CACHE_DIR,
DENO_CACHE_DIR,
GO_CACHE_DIR,
TAR_PIP_TMP_CACHE_DIR,
DENO_TMP_CACHE_DIR,
GO_TMP_CACHE_DIR,
] {
DirBuilder::new()
.recursive(true)
.create(x)
@@ -281,8 +301,13 @@ pub async fn run_workers<R: rsmq_async::RsmqConnection + Send + Sync + Clone + '
#[cfg(feature = "enterprise")]
if let Some(ref s) = S3_CACHE_BUCKET.clone() {
// We donwload the entire cache as tar
windmill_worker::copy_cache_from_bucket_as_tar(&s).await;
join!(
windmill_worker::copy_denogo_cache_from_bucket_as_tar(&s),
windmill_worker::copy_all_piptars_from_bucket(&s)
);
if let Err(e) = windmill_worker::untar_all_piptars().await {
tracing::info!("Failed to untar piptars. Error: {:?}", e);
}
}
IS_READY.store(true, Ordering::Relaxed);
+233 -27
View File
@@ -1,5 +1,5 @@
#[cfg(feature = "enterprise")]
use crate::{ROOT_CACHE_DIR, ROOT_TMP_CACHE_DIR, TAR_CACHE_RATE, TMP_DIR};
use crate::{ROOT_CACHE_DIR, ROOT_TMP_CACHE_DIR, TAR_CACHE_RATE, TAR_PIP_TMP_CACHE_DIR, TMP_DIR};
#[cfg(feature = "enterprise")]
use itertools::Itertools;
#[cfg(feature = "enterprise")]
@@ -14,7 +14,109 @@ use tokio::{process::Command, sync::mpsc::Sender, time::Instant};
use windmill_common::error;
#[cfg(feature = "enterprise")]
const TAR_CACHE_FILENAME: &str = "entirecache.tar";
const TAR_CACHE_FILENAME: &str = "denogocache.tar";
#[cfg(feature = "enterprise")]
pub async fn build_tar_and_push(bucket: &str, folder: String) -> error::Result<()> {
tracing::info!("Started building and pushing piptar {folder}");
let start = Instant::now();
let folder_name = folder.split("/").last().unwrap();
let tar_path = format!("{TAR_PIP_TMP_CACHE_DIR}/{folder_name}.tar",);
if let Err(e) = execute_command(
ROOT_TMP_CACHE_DIR,
"tar",
vec!["-c", "-f", &tar_path, &folder],
)
.await
{
tracing::info!("Failed to tar cache. Error: {:?}", e);
return Err(e);
}
let tar_metadata = tokio::fs::metadata(&tar_path).await;
if tar_metadata.is_err() || tar_metadata.as_ref().unwrap().len() == 0 {
tracing::info!("Failed to tar cache: {folder}");
return Err(error::Error::ExecutionErr(format!(
"Failed to tar cache: {folder}"
)));
}
if let Err(e) = execute_command(
ROOT_TMP_CACHE_DIR,
"rclone",
vec![
"copyto",
&tar_path,
&format!(":s3,env_auth=true:{bucket}/tar/pip/{folder_name}.tar"),
"-v",
"--size-only",
"--fast-list",
],
)
.await
{
tracing::info!("Failed to copy piptar {folder} to bucket. Error: {:?}", e);
return Err(e);
}
tracing::info!(
"Finished copying piptar {folder} to bucket {bucket} as tar, took: {:?}s. Size of tar: {}",
start.elapsed().as_secs(),
tar_metadata.unwrap().len()
);
Ok(())
}
#[cfg(feature = "enterprise")]
pub async fn pull_from_tar(bucket: &str, folder: String) -> error::Result<()> {
use tokio::fs::metadata;
let folder_name = folder.split("/").last().unwrap();
tracing::info!("Attempting to pull piptar {folder_name} from bucket");
let start = Instant::now();
let tar_path = format!("tar/pip/{folder_name}.tar");
let target = format!("{ROOT_TMP_CACHE_DIR}/{tar_path}");
if let Err(e) = execute_command(
ROOT_TMP_CACHE_DIR,
"rclone",
vec![
"copyto",
&format!(":s3,env_auth=true:{bucket}/{tar_path}"),
&target,
"-v",
"--size-only",
"--fast-list",
],
)
.await
{
tracing::info!(
"Failed to copy tar {folder_name} from bucket. Error: {:?}",
e
);
return Err(e);
}
if metadata(&target).await.is_err() {
tracing::info!(
"piptar {folder_name} not found in bucket. Took {:?}ms",
start.elapsed().as_millis()
);
return Err(error::Error::ExecutionErr(format!(
"tar {folder_name} does not exist in bucket"
)));
}
extract_pip_tar(&target, &folder).await?;
tracing::info!(
"Finished pulling and extracting {folder_name} from took {:?}ms",
start.elapsed().as_millis()
);
Ok(())
}
#[cfg(feature = "enterprise")]
pub async fn cache_global(bucket: &str, tx: Sender<()>) -> error::Result<()> {
@@ -47,6 +149,8 @@ pub async fn copy_cache_from_bucket(bucket: &str, tx: Sender<()>) -> error::Resu
"--exclude",
&format!("deno/gen/file/tmp/windmill/**"),
"--exclude",
&format!("pip/**"),
"--exclude",
&format!("{TAR_CACHE_FILENAME}"),
],
)
@@ -84,6 +188,10 @@ pub async fn copy_cache_to_bucket(bucket: &str) -> error::Result<()> {
&format!("deno/gen/file/tmp/windmill/**"),
"--exclude",
&format!("{TAR_CACHE_FILENAME}"),
"--exclude",
&format!("pip/**"),
"--exclude",
&format!("tar/**"),
],
)
.await
@@ -110,7 +218,6 @@ pub async fn copy_cache_to_bucket_as_tar(bucket: &str) {
"-c",
"-f",
&format!("{ROOT_TMP_CACHE_DIR}{TAR_CACHE_FILENAME}"),
"pip",
"go",
"deno",
],
@@ -160,22 +267,12 @@ pub async fn copy_cache_to_bucket_as_tar(bucket: &str) {
}
#[cfg(feature = "enterprise")]
pub async fn copy_cache_from_bucket_as_tar(bucket: &str) {
pub async fn copy_denogo_cache_from_bucket_as_tar(bucket: &str) {
use tokio::fs::metadata;
tracing::info!("Copying cache from bucket {bucket} as tar");
tracing::info!("Copying denogo cache from bucket {bucket} as tar");
if metadata(&ROOT_TMP_CACHE_DIR).await.is_ok() {
if let Err(e) = tokio::fs::remove_dir_all(&ROOT_TMP_CACHE_DIR).await {
tracing::info!(error = %e, "Could not remove root tmp cache dir");
}
}
tokio::fs::create_dir_all(&ROOT_TMP_CACHE_DIR)
.await
.expect("Could not create root tmp cache dir");
let elapsed = Instant::now();
let start: Instant = Instant::now();
if let Err(e) = execute_command(
ROOT_CACHE_DIR,
@@ -191,7 +288,7 @@ pub async fn copy_cache_from_bucket_as_tar(bucket: &str) {
)
.await
{
tracing::info!("Failed copy tar from cache. Error: {:?}", e);
tracing::info!("Failed copying denogo tar from cache. Error: {:?}", e);
return;
}
@@ -202,27 +299,29 @@ pub async fn copy_cache_from_bucket_as_tar(bucket: &str) {
)
.await
{
tracing::info!("Failed to untar cache. Error: {:?}", e);
tracing::info!("Failed to untar denogo. Error: {:?}", e);
return;
}
if let Err(e) =
tokio::fs::remove_dir_all(format!("{ROOT_CACHE_DIR}deno/gen/file/tmp/windmill")).await
{
tracing::info!("Failed to remove tmp gen windmill. Error: {:?}", e);
};
let denogen = format!("{ROOT_CACHE_DIR}deno/gen/file/tmp/windmill");
if metadata(&denogen).await.is_ok() {
let _ = tokio::fs::remove_dir_all(denogen).await;
}
if let Err(e) = tokio::fs::remove_file(format!("{ROOT_CACHE_DIR}{TAR_CACHE_FILENAME}")).await {
tracing::info!("Failed to remove tar cache. Error: {:?}", e);
tracing::info!("Failed to remove denotar cache. Error: {:?}", e);
return;
};
tracing::info!(
"Finished copying cache from bucket {bucket} as tar, took: {:?}s",
elapsed.elapsed().as_secs()
"Finished copying denogotar from bucket {bucket} as tar, took: {:?}s",
start.elapsed().as_secs()
);
for x in ["deno", "go", "pip"] {
tracing::info!("Copying denogo cache from bucket {bucket} as tar");
let start: Instant = Instant::now();
for x in ["deno", "go"] {
if let Err(e) = execute_command(
TMP_DIR,
"cp",
@@ -233,6 +332,40 @@ pub async fn copy_cache_from_bucket_as_tar(bucket: &str) {
tracing::info!(error = %e, "Could not copy root dir to tmp root dir");
}
}
tracing::info!(
"Finished copying untarred denogo to tmp cache, took: {:?}s",
start.elapsed().as_secs()
);
}
#[cfg(feature = "enterprise")]
pub async fn copy_all_piptars_from_bucket(bucket: &str) {
tracing::info!("Copying all piptars cache from bucket {bucket}");
let start = Instant::now();
if let Err(e) = execute_command(
ROOT_CACHE_DIR,
"rclone",
vec![
"copy",
&format!(":s3,env_auth=true:{bucket}/tar/pip/"),
&TAR_PIP_TMP_CACHE_DIR,
"-v",
"--size-only",
"--fast-list",
],
)
.await
{
tracing::info!("Failed transferring all piptars from cache. Error: {:?}", e);
return;
}
tracing::info!(
"Finished transferring piptars from bucket {bucket} as tar, took: {:?}s",
start.elapsed().as_secs()
);
}
// async fn check_if_bucket_syncable(bucket: &str) -> bool {
@@ -261,13 +394,84 @@ pub async fn copy_tmp_cache_to_cache() -> error::Result<()> {
ROOT_CACHE_DIR,
"--exclude",
TAR_CACHE_FILENAME,
"--exclude",
&format!("pip/**"),
"--exclude",
&format!("tar/**"),
],
)
.await?;
tracing::info!(
"Finished copying local tmp cache to local cache. Took {}ms",
start.elapsed().as_millis(),
);
let start = Instant::now();
if let Err(e) = untar_all_piptars().await {
tracing::info!("Failed to untar piptars. Error: {:?}", e);
}
tracing::info!(
"Finished untarring all piptars took: {:?}s",
start.elapsed().as_secs()
);
Ok(())
}
#[cfg(feature = "enterprise")]
pub async fn untar_all_piptars() -> error::Result<()> {
use tokio::fs::{self, metadata};
use crate::PIP_CACHE_DIR;
let start: Instant = Instant::now();
let mut entries = fs::read_dir(TAR_PIP_TMP_CACHE_DIR).await?;
while let Some(entry) = entries.next_entry().await? {
if let Err(e) = {
let path = entry.file_name().into_string().expect("Invalid path");
let folder = format!(
"{PIP_CACHE_DIR}/{}",
path.split('/')
.last()
.unwrap()
.strip_suffix(".tar")
.unwrap()
);
if metadata(&folder).await.is_ok() {
continue;
}
extract_pip_tar(&path, &folder).await?;
Ok(()) as error::Result<()>
} {
tracing::info!("Failed to extract pip tar. Error: {:?}", e);
}
}
tracing::info!(
"Finished copying local tmp cache to local cache. Took {}ms",
start.elapsed().as_millis(),
);
Ok(())
}
#[cfg(feature = "enterprise")]
pub async fn extract_pip_tar(tar: &str, folder: &str) -> error::Result<()> {
use tokio::fs;
let start: Instant = Instant::now();
fs::create_dir(&folder).await?;
if let Err(e) = execute_command(&folder, "tar", vec!["-xpvf", tar]).await {
tracing::info!("Failed to untar cache. Error: {:?}", e);
return Err(e);
}
tracing::info!(
"Finished extracting pip tar {folder}. Took {}ms",
start.elapsed().as_millis(),
);
Ok(())
}
@@ -283,6 +487,8 @@ pub async fn copy_cache_to_tmp_cache() -> error::Result<()> {
ROOT_TMP_CACHE_DIR,
"--exclude",
TAR_CACHE_FILENAME,
"--exclude",
&format!("pip/**"),
],
)
.await?;
+3 -1
View File
@@ -8,5 +8,7 @@ mod worker;
mod worker_flow;
#[cfg(feature = "enterprise")]
pub use global_cache::copy_cache_from_bucket_as_tar;
pub use global_cache::{
copy_all_piptars_from_bucket, copy_denogo_cache_from_bucket_as_tar, untar_all_piptars,
};
pub use worker::*;
+17 -1
View File
@@ -45,11 +45,14 @@ const NSJAIL_CONFIG_DOWNLOAD_PY_CONTENT: &str = include_str!("../nsjail/download
const NSJAIL_CONFIG_RUN_PYTHON3_CONTENT: &str = include_str!("../nsjail/run.python3.config.proto");
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::{
common::{read_result, set_logs},
create_args_and_out_file, get_reserved_variables, handle_child, write_file,
AuthedClientBackgroundTask, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, PATH_ENV,
PIP_CACHE_DIR,
PIP_CACHE_DIR, S3_CACHE_BUCKET,
};
pub async fn create_dependencies_dir(job_dir: &str) {
@@ -465,6 +468,14 @@ pub async fn handle_python_reqs(
continue;
}
#[cfg(feature = "enterprise")]
if let Some(ref bucket) = *S3_CACHE_BUCKET {
if pull_from_tar(bucket, venv_p.clone()).await.is_ok() {
req_paths.push(venv_p.clone());
continue;
}
}
logs.push_str("\n--- PIP INSTALL ---\n");
logs.push_str(&format!("\n{req} is being installed for the first time.\n It will be cached for all ulterior uses."));
@@ -548,6 +559,11 @@ pub async fn handle_python_reqs(
);
child?;
#[cfg(feature = "enterprise")]
if let Some(ref bucket) = *S3_CACHE_BUCKET {
let venv_p = venv_p.clone();
tokio::spawn(build_tar_and_push(bucket, venv_p));
}
req_paths.push(venv_p);
}
Ok(req_paths)
+1 -1
View File
@@ -127,7 +127,7 @@ pub const ROOT_TMP_CACHE_DIR: &str = "/tmp/windmill/tmpcache/";
pub const PIP_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "pip");
pub const DENO_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "deno");
pub const GO_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "go");
pub const PIP_TMP_CACHE_DIR: &str = concatcp!(ROOT_TMP_CACHE_DIR, "pip");
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 GO_TMP_CACHE_DIR: &str = concatcp!(ROOT_TMP_CACHE_DIR, "go");