From 369dd0dac61e5430856ed9abf7129bbad3b75860 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Thu, 20 Apr 2023 20:05:12 +0200 Subject: [PATCH] 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 --- backend/src/main.rs | 35 ++- backend/windmill-worker/src/global_cache.rs | 260 ++++++++++++++++-- backend/windmill-worker/src/lib.rs | 4 +- .../windmill-worker/src/python_executor.rs | 18 +- backend/windmill-worker/src/worker.rs | 2 +- 5 files changed, 284 insertions(+), 35 deletions(-) diff --git a/backend/src/main.rs b/backend/src/main.rs index 51c9dd6ae0..411ccba489 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -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 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?; diff --git a/backend/windmill-worker/src/lib.rs b/backend/windmill-worker/src/lib.rs index 597e5f4f5e..ccaa3fd048 100644 --- a/backend/windmill-worker/src/lib.rs +++ b/backend/windmill-worker/src/lib.rs @@ -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::*; diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 3ad62be15e..03485ae81c 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -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) diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 73f93f8c36..a08d68140d 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -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");