From 7e35d9989aab74cd91f676c679b36e98033f1176 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jakub=20Ko=C5=82odziejczak?= <31549762+mrl5@users.noreply.github.com> Date: Wed, 12 Oct 2022 23:02:27 +0200 Subject: [PATCH] feat(backend): supercache for python heavy dependencies in alpha --- backend/src/worker.rs | 265 ++++++++++++++++++++++++-------- nsjail/run.python3.config.proto | 4 +- 2 files changed, 207 insertions(+), 62 deletions(-) diff --git a/backend/src/worker.rs b/backend/src/worker.rs index 6fa5323af6..12bf34da12 100644 --- a/backend/src/worker.rs +++ b/backend/src/worker.rs @@ -49,6 +49,7 @@ use futures::{ use async_recursion::async_recursion; const TMP_DIR: &str = "/tmp/windmill"; +const PIP_SUPERCACHE_DIR: &str = "/tmp/windmill/cache/pip_permanent"; const PIP_CACHE_DIR: &str = "/tmp/windmill/cache/pip"; const DENO_CACHE_DIR: &str = "/tmp/windmill/cache/deno"; const GO_CACHE_DIR: &str = "/tmp/windmill/cache/go"; @@ -95,7 +96,13 @@ pub async fn run_worker( let worker_dir = format!("{TMP_DIR}/{worker_name}"); tracing::debug!(worker_dir = %worker_dir, worker_name = %worker_name, "Creating worker dir"); - for x in [&worker_dir, PIP_CACHE_DIR, DENO_CACHE_DIR, GO_CACHE_DIR] { + for x in [ + &worker_dir, + PIP_SUPERCACHE_DIR, + PIP_CACHE_DIR, + DENO_CACHE_DIR, + GO_CACHE_DIR, + ] { DirBuilder::new() .recursive(true) .create(x) @@ -151,6 +158,10 @@ pub async fn run_worker( let go_path = std::env::var("GO_PATH").unwrap_or_else(|_| "/usr/bin/go".to_string()); let python_path = std::env::var("PYTHON_PATH").unwrap_or_else(|_| "/usr/local/bin/python3".to_string()); + let python_heavy_deps = std::env::var("PYTHON_HEAVY_DEPS") + .map(|x| x.split(',').map(|x| x.to_string()).collect::>()) + .unwrap_or_else(|_| vec![]); + let python_version = std::env::var("PYTHON_VERSION").unwrap_or_else(|_| "3.10".to_string()); let nsjail_path = std::env::var("NSJAIL_PATH").unwrap_or_else(|_| "nsjail".to_string()); let path_env = std::env::var("PATH").unwrap_or_else(|_| String::new()); let gopath_env = std::env::var("GOPATH").unwrap_or_else(|_| String::new()); @@ -162,6 +173,8 @@ pub async fn run_worker( deno_path, go_path, python_path, + python_version, + python_heavy_deps, nsjail_path, path_env, gopath_env, @@ -382,6 +395,8 @@ struct Envs { deno_path: String, go_path: String, python_path: String, + python_version: String, + python_heavy_deps: Vec, nsjail_path: String, path_env: String, gopath_env: String, @@ -1031,6 +1046,8 @@ async fn handle_python_job( Envs { nsjail_path, python_path, + python_heavy_deps, + python_version, path_env, pip_extra_index_url, pip_index_url, @@ -1053,6 +1070,8 @@ async fn handle_python_job( create_dependencies_dir(job_dir).await; + let mut additional_python_paths: Vec = vec![]; + if requirements.len() > 0 { if !disable_nsjail { let _ = write_file( @@ -1066,14 +1085,15 @@ async fn handle_python_job( ) .await?; } - let _ = write_file(job_dir, "requirements.txt", &requirements).await?; - tracing::info!( - worker_name = %worker_name, - job_id = %job.id, - workspace_id = %job.workspace_id, - "started setup python dependencies" - ); + let mut heavy_deps = vec!["numpy".to_string(), "pandas".to_string()]; + heavy_deps.extend(python_heavy_deps.into_iter().map(|s| s.to_string())); + + let (heavy, regular): (Vec<&str>, Vec<&str>) = requirements + .split("\n") + .partition(|d| heavy_deps.iter().any(|hd| d.starts_with(hd))); + + let _ = write_file(job_dir, "requirements.txt", ®ular.join("\n")).await?; let mut vars = vec![]; if let Some(url) = pip_extra_index_url { @@ -1085,58 +1105,79 @@ async fn handle_python_job( if let Some(host) = pip_trusted_host { vars.push(("TRUSTED_HOST", host)); } - let child = if !disable_nsjail { - Command::new(nsjail_path) - .current_dir(job_dir) - .env_clear() - .envs(vars) - .args(vec!["--config", "download.config.proto"]) - .stdout(Stdio::piped()) - .stderr(Stdio::piped()) - .spawn()? - } else { - let mut args = vec![ - "-m", - "pip", - "install", - "--no-color", - "--isolated", - "--no-warn-conflicts", - "--disable-pip-version-check", - "-t", - "./dependencies", - "-r", - "./requirements.txt", - ]; - if let Some(url) = pip_extra_index_url { - args.extend(["--extra-index-url", url]); - } - if let Some(url) = pip_index_url { - args.extend(["--index-url", url]); - } - if let Some(host) = pip_trusted_host { - args.extend(["--trusted-host", host]); - } - Command::new(python_path) - .current_dir(job_dir) - .env_clear() - .args(args) - .stdout(Stdio::piped()) - .stderr(Stdio::piped()) - .spawn()? - }; - logs.push_str("\n--- PIP DEPENDENCIES INSTALL ---\n"); - let child = handle_child(&job.id, db, logs, timeout, child).await; - tracing::info!( - worker_name = %worker_name, - job_id = %job.id, - workspace_id = %job.workspace_id, - is_ok = child.is_ok(), - "finished setting up python dependencies {}", - job.id - ); - child?; + if heavy.len() > 0 { + logs.push_str(&format!( + "\nheavy deps detected, using supercache for: {heavy:?}" + )); + additional_python_paths = + handle_python_heavy_reqs(python_path, heavy, vars.clone(), job, logs, db, timeout) + .await?; + } + + if regular.len() > 0 { + tracing::info!( + worker_name = %worker_name, + job_id = %job.id, + workspace_id = %job.workspace_id, + "started setup python dependencies" + ); + + let child = if !disable_nsjail { + Command::new(nsjail_path) + .current_dir(job_dir) + .env_clear() + .envs(vars) + .args(vec!["--config", "download.config.proto"]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn()? + } else { + let mut args = vec![ + "-m", + "pip", + "install", + "--no-color", + "--isolated", + "--no-warn-conflicts", + "--disable-pip-version-check", + "-t", + "./dependencies", + "-r", + "./requirements.txt", + ]; + if let Some(url) = pip_extra_index_url { + args.extend(["--extra-index-url", url]); + } + if let Some(url) = pip_index_url { + args.extend(["--index-url", url]); + } + if let Some(host) = pip_trusted_host { + args.extend(["--trusted-host", host]); + } + Command::new(python_path) + .current_dir(job_dir) + .env_clear() + .args(args) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn()? + }; + + logs.push_str("\n--- PIP DEPENDENCIES INSTALL ---\n"); + let child = handle_child(&job.id, db, logs, timeout, child).await; + tracing::info!( + worker_name = %worker_name, + job_id = %job.id, + workspace_id = %job.workspace_id, + is_ok = child.is_ok(), + "finished setting up python dependencies {}", + job.id + ); + child?; + } else { + logs.push_str("\nskipping pip install since not needed"); + }; } logs.push_str("\n\n--- PYTHON CODE EXECUTION ---\n"); @@ -1204,18 +1245,45 @@ with open("result.json", 'w') as f: write_file(job_dir, "main.py", &wrapper_content).await?; let mut reserved_variables = get_reserved_variables(job, &token, &base_url, db).await?; + let additional_python_paths_folders = additional_python_paths + .iter() + .map(|x| format!(":{PIP_SUPERCACHE_DIR}/{x}/lib/python{python_version}/site-packages")) + .join(""); if !disable_nsjail { + let shared_deps = additional_python_paths + .into_iter() + .map(|pp| { + format!( + r#" +mount {{ + src: "{pp}" + dst: "{pp}" + is_bind: true + rw: false +}} + "# + ) + }) + .join("\n"); let _ = write_file( job_dir, "run.config.proto", &NSJAIL_CONFIG_RUN_PYTHON3_CONTENT .replace("{JOB_DIR}", job_dir) .replace("{CLONE_NEWUSER}", &(!disable_nuser).to_string()) - .replace("{SHARED_MOUNT}", shared_mount), + .replace("{SHARED_MOUNT}", shared_mount) + .replace("{SHARED_DEPENDENCIES}", shared_deps.as_str()) + .replace( + "{ADDITIONAL_PYTHON_PATHS}", + additional_python_paths_folders.as_str(), + ), ) .await?; } else { - reserved_variables.insert("PYTHONPATH".to_string(), format!("{job_dir}/dependencies")); + reserved_variables.insert( + "PYTHONPATH".to_string(), + format!("{job_dir}/dependencies{additional_python_paths_folders}"), + ); } tracing::info!( @@ -1229,6 +1297,7 @@ with open("result.json", 'w') as f: Command::new(nsjail_path) .current_dir(job_dir) .env_clear() + // inject PYTHONPATH here - for some reason I had to do it in nsjail conf .envs(reserved_variables) .env("PATH", path_env) .env("BASE_INTERNAL_URL", base_internal_url) @@ -1858,6 +1927,60 @@ pub async fn handle_zombie_jobs_periodically( } } +async fn handle_python_heavy_reqs( + python_path: &String, + heavy_requirements: Vec<&str>, + vars: Vec<(&str, &String)>, + job: &QueuedJob, + logs: &mut String, + db: &sqlx::Pool, + timeout: i32, +) -> error::Result> { + let mut req_paths: Vec = vec![]; + for req in heavy_requirements { + // todo: handle many reqs + let venv_p = format!("{PIP_SUPERCACHE_DIR}/{req}"); + if metadata(&venv_p).await.is_ok() { + tracing::info!("already exists: {:?}", &venv_p); + req_paths.push(venv_p); + continue; + } + + logs.push_str("\n--- PIP SUPERCACHE INSTALL ---\n"); + logs.push_str(&format!("\nthe heavy dependency {req} is being installed for the first time.\nIt will take a bit longer but further execution will be much faster!")); + + logs.push_str("venv creation\n"); + + let child = Command::new(python_path) + .current_dir(PIP_SUPERCACHE_DIR) + .env_clear() + .args(vec![ + "-c", + &format!("import venv; venv.create(\"./{req}\");"), + ]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn()?; + + handle_child(&job.id, db, logs, timeout, child).await?; + + logs.push_str("pip install\n"); + + let child = Command::new(python_path) + .current_dir(format!("{PIP_SUPERCACHE_DIR}/{}", &req)) + .env_clear() + .envs(vars.clone()) + .args(vec!["-m", "pip", "install", &req]) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn()?; + handle_child(&job.id, db, logs, timeout, child).await?; + + req_paths.push(venv_p); + } + Ok(req_paths) +} + #[cfg(test)] mod tests { use futures::Stream; @@ -2292,6 +2415,26 @@ def main(): assert_eq!(result, serde_json::json!("hello world")); } + #[sqlx::test(fixtures("base"))] + async fn test_python_job_heavy_dep(db: DB) { + initialize_tracing().await; + + let content = r#" +import numpy as np + +def main(): + a = np.arange(15).reshape(3, 5) + return len(a) + "# + .to_owned(); + + let job = JobPayload::Code(RawCode { content, path: None, language: ScriptLang::Python3 }); + + let result = run_job_in_new_worker_until_complete(&db, job).await; + + assert_eq!(result, serde_json::json!(3)); + } + #[sqlx::test(fixtures("base"))] async fn test_python_job_with_imports(db: DB) { initialize_tracing().await; diff --git a/nsjail/run.python3.config.proto b/nsjail/run.python3.config.proto index fcba25894d..3ac95645a0 100644 --- a/nsjail/run.python3.config.proto +++ b/nsjail/run.python3.config.proto @@ -125,9 +125,11 @@ mount { {SHARED_MOUNT} +{SHARED_DEPENDENCIES} + iface_no_lo: true -envar: "PYTHONPATH=/tmp/dependencies" +envar: "PYTHONPATH=/tmp/dependencies{ADDITIONAL_PYTHON_PATHS}" envar: "HOME=/tmp"