feat: Parallelize uv install (#4774)
* Implement MVP of Parallel uv installation * Implement PY_CONCURRENT_DOWNLOADS * Remove Flock for uv installs * Make S3 pull/push parallel * Refactor and allow to Cancel installation * Dont print S3 in output if disabled * Implement better error handling * Polishing * More polishing * Implement error-handler for kill_tx_2.send() * Fix and Format prev merge * Presubscribe to all kill_tx's We do it now before first event could fire Meaning no events can be lost anymore * Early print errors and safer error handling * Return Err if installation failed * Final changes * Return error instead of just printing it * Safer the way to acquire permit * Fix compilation error * Remove double error logs
This commit is contained in:
@@ -3,8 +3,10 @@ use std::{
|
||||
fs,
|
||||
path::Path,
|
||||
process::Stdio,
|
||||
sync::Arc
|
||||
};
|
||||
|
||||
use anyhow::anyhow;
|
||||
use itertools::Itertools;
|
||||
use regex::Regex;
|
||||
use serde_json::value::RawValue;
|
||||
@@ -13,6 +15,8 @@ use tokio::{
|
||||
fs::{metadata, DirBuilder, File},
|
||||
io::AsyncReadExt,
|
||||
process::Command,
|
||||
sync::Semaphore,
|
||||
task,
|
||||
};
|
||||
use uuid::Uuid;
|
||||
#[cfg(all(feature = "enterprise", feature = "parquet"))]
|
||||
@@ -37,6 +41,9 @@ lazy_static::lazy_static! {
|
||||
static ref UV_PATH: String =
|
||||
std::env::var("UV_PATH").unwrap_or_else(|_| "/usr/local/bin/uv".to_string());
|
||||
|
||||
static ref PY_CONCURRENT_DOWNLOADS: usize =
|
||||
std::env::var("PY_CONCURRENT_DOWNLOADS").ok().map(|flag| flag.parse().unwrap_or(20)).unwrap_or(20);
|
||||
|
||||
static ref FLOCK_PATH: String =
|
||||
std::env::var("FLOCK_PATH").unwrap_or_else(|_| "/usr/bin/flock".to_string());
|
||||
static ref NON_ALPHANUM_CHAR: Regex = regex::Regex::new(r"[^0-9A-Za-z=.-]").unwrap();
|
||||
@@ -1108,51 +1115,27 @@ lazy_static::lazy_static! {
|
||||
static ref PIP_SECRET_VARIABLE: Regex = Regex::new(r"\$\{PIP_SECRET:([^\s\}]+)\}").unwrap();
|
||||
}
|
||||
|
||||
/// pip install, include cached or pull from S3
|
||||
pub async fn handle_python_reqs(
|
||||
requirements: Vec<&str>,
|
||||
job_id: &Uuid,
|
||||
/// Spawn process of uv install
|
||||
/// Can be wrapped by nsjail depending on configuration
|
||||
#[inline]
|
||||
async fn spawn_uv_install(
|
||||
w_id: &str,
|
||||
mem_peak: &mut i32,
|
||||
canceled_by: &mut Option<CanceledBy>,
|
||||
db: &sqlx::Pool<sqlx::Postgres>,
|
||||
worker_name: &str,
|
||||
req: &str,
|
||||
venv_p: &str,
|
||||
job_dir: &str,
|
||||
worker_dir: &str,
|
||||
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
|
||||
// TODO: Remove (Deprecated)
|
||||
mut no_uv_install: bool,
|
||||
is_ansible: bool,
|
||||
) -> error::Result<Vec<String>> {
|
||||
let mut req_paths: Vec<String> = vec![];
|
||||
let mut vars = vec![("PATH", PATH_ENV.as_str())];
|
||||
let pip_extra_index_url;
|
||||
let pip_index_url;
|
||||
|
||||
no_uv_install |= *USE_PIP_INSTALL;
|
||||
|
||||
if no_uv_install && !is_ansible {
|
||||
append_logs(&job_id, w_id, "\nFallback to pip (Deprecated!)\n", db).await;
|
||||
tracing::warn!("Fallback to pip");
|
||||
}
|
||||
|
||||
(pip_extra_index_url, pip_index_url): (Option<String>, Option<String>),
|
||||
no_uv_install: bool,
|
||||
) -> Result<tokio::process::Child, Error> {
|
||||
if !*DISABLE_NSJAIL {
|
||||
pip_extra_index_url = PIP_EXTRA_INDEX_URL
|
||||
.read()
|
||||
.await
|
||||
.clone()
|
||||
.map(handle_ephemeral_token);
|
||||
tracing::info!(
|
||||
workspace_id = %w_id,
|
||||
"starting nsjail"
|
||||
);
|
||||
|
||||
let mut vars = vec![("PATH", PATH_ENV.as_str())];
|
||||
if let Some(url) = pip_extra_index_url.as_ref() {
|
||||
vars.push(("EXTRA_INDEX_URL", url));
|
||||
}
|
||||
|
||||
pip_index_url = PIP_INDEX_URL
|
||||
.read()
|
||||
.await
|
||||
.clone()
|
||||
.map(handle_ephemeral_token);
|
||||
|
||||
if let Some(url) = pip_index_url.as_ref() {
|
||||
vars.push(("INDEX_URL", url));
|
||||
}
|
||||
@@ -1163,6 +1146,265 @@ pub async fn handle_python_reqs(
|
||||
vars.push(("TRUSTED_HOST", host));
|
||||
}
|
||||
|
||||
vars.push(("REQ", &req));
|
||||
vars.push(("TARGET", venv_p));
|
||||
|
||||
let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str());
|
||||
nsjail_cmd
|
||||
.current_dir(job_dir)
|
||||
.env_clear()
|
||||
.envs(vars)
|
||||
.envs(PROXY_ENVS.clone())
|
||||
.args(vec!["--config", "download.config.proto"])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
start_child_process(nsjail_cmd, NSJAIL_PATH.as_str()).await
|
||||
} else {
|
||||
let fssafe_req = NON_ALPHANUM_CHAR.replace_all(&req, "_").to_string();
|
||||
#[cfg(unix)]
|
||||
let req = if no_uv_install {
|
||||
format!("'{}'", req)
|
||||
} else {
|
||||
req.to_owned()
|
||||
};
|
||||
|
||||
#[cfg(windows)]
|
||||
let req = format!("{}", req);
|
||||
|
||||
let mut command_args = if no_uv_install {
|
||||
vec![
|
||||
PYTHON_PATH.as_str(),
|
||||
"-m",
|
||||
"pip",
|
||||
"install",
|
||||
&req,
|
||||
"-I",
|
||||
"--no-deps",
|
||||
"--no-color",
|
||||
"--isolated",
|
||||
"--no-warn-conflicts",
|
||||
"--disable-pip-version-check",
|
||||
"-t",
|
||||
venv_p,
|
||||
]
|
||||
} else {
|
||||
vec![
|
||||
UV_PATH.as_str(),
|
||||
"pip",
|
||||
"install",
|
||||
&req,
|
||||
"--no-deps",
|
||||
"--no-color",
|
||||
// "-p",
|
||||
// "3.11",
|
||||
// Prevent uv from discovering configuration files.
|
||||
"--no-config",
|
||||
"--link-mode=copy",
|
||||
"--system",
|
||||
// Prefer main index over extra
|
||||
// https://docs.astral.sh/uv/pip/compatibility/#packages-that-exist-on-multiple-indexes
|
||||
// TODO: Use env variable that can be toggled from UI
|
||||
"--index-strategy",
|
||||
"unsafe-best-match",
|
||||
"--target",
|
||||
venv_p,
|
||||
"--no-cache",
|
||||
"-q",
|
||||
]
|
||||
};
|
||||
|
||||
if let Some(url) = pip_extra_index_url.as_ref() {
|
||||
url.split(",").for_each(|url| {
|
||||
command_args.extend(["--extra-index-url", url]);
|
||||
});
|
||||
}
|
||||
|
||||
if let Some(url) = pip_index_url.as_ref() {
|
||||
command_args.extend(["--index-url", url]);
|
||||
}
|
||||
if let Some(cert_path) = PIP_INDEX_CERT.as_ref() {
|
||||
command_args.extend(["--cert", cert_path]);
|
||||
}
|
||||
if let Some(host) = PIP_TRUSTED_HOST.as_ref() {
|
||||
command_args.extend(["--trusted-host", &host]);
|
||||
}
|
||||
|
||||
let mut envs = vec![("PATH", PATH_ENV.as_str())];
|
||||
envs.push(("HOME", HOME_ENV.as_str()));
|
||||
|
||||
tracing::debug!("uv pip install command: {:?}", command_args);
|
||||
|
||||
#[cfg(unix)]
|
||||
{
|
||||
if no_uv_install {
|
||||
let mut flock_cmd = Command::new(FLOCK_PATH.as_str());
|
||||
flock_cmd
|
||||
.env_clear()
|
||||
.envs(PROXY_ENVS.clone())
|
||||
.envs(envs)
|
||||
.args([
|
||||
"-x",
|
||||
&format!(
|
||||
"{}/{}-{}.lock",
|
||||
LOCK_CACHE_DIR,
|
||||
if no_uv_install { "pip" } else { "py311" },
|
||||
fssafe_req
|
||||
),
|
||||
"--command",
|
||||
&command_args.join(" "),
|
||||
])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
start_child_process(flock_cmd, FLOCK_PATH.as_str()).await
|
||||
} else {
|
||||
let mut cmd = Command::new(command_args[0]);
|
||||
cmd
|
||||
.env_clear()
|
||||
.envs(PROXY_ENVS.clone())
|
||||
.envs(envs)
|
||||
.args(&command_args[1..])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
start_child_process(cmd, UV_PATH.as_str()).await
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
{
|
||||
let installer_path = if no_uv_install { command_args[0] } else { "uv" };
|
||||
let mut cmd: Command = Command::new(&installer_path);
|
||||
cmd.env_clear()
|
||||
.envs(envs)
|
||||
.envs(PROXY_ENVS.clone())
|
||||
.env("SystemRoot", SYSTEM_ROOT.as_str())
|
||||
.env(
|
||||
"TMP",
|
||||
std::env::var("TMP").unwrap_or_else(|_| String::from("/tmp")),
|
||||
)
|
||||
.args(&command_args[1..])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
start_child_process(cmd, installer_path).await
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// length = 5
|
||||
/// value = "foo"
|
||||
/// output = "foo "
|
||||
/// 12345
|
||||
fn pad_string(value: &str, total_length: usize) -> String {
|
||||
if value.len() >= total_length {
|
||||
value.to_string() // Return the original string if it's already long enough
|
||||
} else {
|
||||
let padding_needed = total_length - value.len();
|
||||
format!("{value}{}", " ".repeat(padding_needed)) // Pad with spaces
|
||||
}
|
||||
}
|
||||
|
||||
/// pip install, include cached or pull from S3
|
||||
pub async fn handle_python_reqs(
|
||||
requirements: Vec<&str>,
|
||||
job_id: &Uuid,
|
||||
w_id: &str,
|
||||
_mem_peak: &mut i32,
|
||||
_canceled_by: &mut Option<CanceledBy>,
|
||||
db: &sqlx::Pool<sqlx::Postgres>,
|
||||
_worker_name: &str,
|
||||
job_dir: &str,
|
||||
worker_dir: &str,
|
||||
_occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
|
||||
// TODO: Remove (Deprecated)
|
||||
mut no_uv_install: bool,
|
||||
is_ansible: bool,
|
||||
) -> error::Result<Vec<String>> {
|
||||
|
||||
let counter_arc = Arc::new(tokio::sync::Mutex::new(0));
|
||||
// Append logs with line like this:
|
||||
// [9/21] + requests==2.32.3 << (S3) | in 57ms
|
||||
#[allow(unused_assignments)]
|
||||
async fn print_success(
|
||||
mut s3_pull: bool,
|
||||
mut s3_push: bool,
|
||||
job_id: &Uuid,
|
||||
w_id: &str,
|
||||
req: &str,
|
||||
req_tl: usize,
|
||||
counter_arc: Arc<tokio::sync::Mutex<usize>>,
|
||||
total_to_install: usize,
|
||||
instant: std::time::Instant,
|
||||
db: Pool<Postgres>,
|
||||
) {
|
||||
#[cfg(not(all(feature = "enterprise", feature = "parquet")))]
|
||||
{
|
||||
(s3_pull, s3_push) = (false, false);
|
||||
}
|
||||
|
||||
#[cfg(all(feature = "enterprise", feature = "parquet"))]
|
||||
if OBJECT_STORE_CACHE_SETTINGS.read().await.is_none() {
|
||||
(s3_pull, s3_push) = (false, false);
|
||||
}
|
||||
|
||||
let mut counter = counter_arc.lock().await;
|
||||
*counter += 1;
|
||||
|
||||
append_logs(
|
||||
job_id,
|
||||
w_id,
|
||||
format!(
|
||||
"\n{}+ {}{}{}| in {}ms",
|
||||
pad_string(&format!("[{}/{total_to_install}]", counter), 9),
|
||||
// Because we want to align to max len [999/999] we take ^
|
||||
// 123456789
|
||||
pad_string(&req, req_tl + 1),
|
||||
// Margin to the right ^
|
||||
if s3_pull { "<< (S3) " } else { "" },
|
||||
if s3_push { " > (S3) " } else { "" },
|
||||
instant.elapsed().as_millis(),
|
||||
),
|
||||
db,
|
||||
)
|
||||
.await;
|
||||
// Drop lock, so next print success can fire
|
||||
}
|
||||
no_uv_install |= *USE_PIP_INSTALL;
|
||||
|
||||
if no_uv_install && !is_ansible {
|
||||
append_logs(&job_id, w_id, "\nFallback to pip (Deprecated!)\n", db).await;
|
||||
tracing::warn!("Fallback to pip");
|
||||
}
|
||||
// Parallelism level (N)
|
||||
let parallel_limit = if no_uv_install {
|
||||
1
|
||||
} else {
|
||||
// Semaphore will panic if value less then 1
|
||||
PY_CONCURRENT_DOWNLOADS.clamp(1, 30)
|
||||
};
|
||||
|
||||
tracing::info!(
|
||||
workspace_id = %w_id,
|
||||
// is_ok = out,
|
||||
"Parallel limit: {}, job: {}",
|
||||
parallel_limit,
|
||||
job_id
|
||||
);
|
||||
|
||||
let pip_indexes = (
|
||||
PIP_EXTRA_INDEX_URL
|
||||
.read()
|
||||
.await
|
||||
.clone()
|
||||
.map(handle_ephemeral_token),
|
||||
PIP_INDEX_URL
|
||||
.read()
|
||||
.await
|
||||
.clone()
|
||||
.map(handle_ephemeral_token),
|
||||
);
|
||||
|
||||
|
||||
// Prepare NSJAIL
|
||||
if !*DISABLE_NSJAIL {
|
||||
let _ = write_file(
|
||||
job_dir,
|
||||
"download.config.proto",
|
||||
@@ -1184,12 +1426,18 @@ pub async fn handle_python_reqs(
|
||||
)?;
|
||||
};
|
||||
|
||||
// Cached paths
|
||||
let mut req_with_penv: Vec<(String, String)> = vec![];
|
||||
|
||||
// Requirements to pull (not cached)
|
||||
let mut req_paths: Vec<String> = vec![];
|
||||
// Find out if there is already cached dependencies
|
||||
// If so, skip them
|
||||
for req in requirements {
|
||||
// Ignore python version annotation backed into lockfile
|
||||
if req.starts_with('#') {
|
||||
continue;
|
||||
}
|
||||
// TODO: Remove
|
||||
let py_prefix = if no_uv_install {
|
||||
PIP_CACHE_DIR
|
||||
} else {
|
||||
@@ -1201,308 +1449,324 @@ pub async fn handle_python_reqs(
|
||||
req.replace(' ', "").replace('/', "").replace(':', "")
|
||||
);
|
||||
if metadata(&venv_p).await.is_ok() {
|
||||
// If dir exists skip installation and push path to output
|
||||
req_paths.push(venv_p);
|
||||
} else {
|
||||
req_with_penv.push((req.to_string(), venv_p));
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(all(feature = "enterprise", feature = "parquet"))]
|
||||
enum PullFromTar {
|
||||
Pulled(String),
|
||||
NotPulled(String, String),
|
||||
}
|
||||
let (kill_tx, ..) = tokio::sync::broadcast::channel::<()>(1);
|
||||
let kill_rxs: Vec<tokio::sync::broadcast::Receiver<()>> =
|
||||
(0..req_with_penv.len()).map(|_| kill_tx.subscribe()).collect();
|
||||
|
||||
#[cfg(all(feature = "enterprise", feature = "parquet"))]
|
||||
if req_with_penv.len() > 0 {
|
||||
if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() {
|
||||
let (done_tx, mut done_rx) = tokio::sync::mpsc::channel(1);
|
||||
let job_id_2 = job_id.clone();
|
||||
let db_2 = db.clone();
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = tokio::time::sleep(tokio::time::Duration::from_secs(5)) => {
|
||||
if let Err(e) = sqlx::query_scalar!("UPDATE queue SET last_ping = now() WHERE id = $1", &job_id_2)
|
||||
.execute(&db_2)
|
||||
.await {
|
||||
tracing::error!("failed to update last_ping: {}", e);
|
||||
}
|
||||
}
|
||||
_ = done_rx.recv() => {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
// ________ Read comments at the end of the function to get more context
|
||||
let (_done_tx, mut done_rx) = tokio::sync::mpsc::channel::<()>(1);
|
||||
|
||||
let start = std::time::Instant::now();
|
||||
let futures = req_with_penv
|
||||
.clone()
|
||||
.into_iter()
|
||||
.map(|(req, venv_p)| {
|
||||
let os = os.clone();
|
||||
async move {
|
||||
if pull_from_tar(os, venv_p.clone(), no_uv_install)
|
||||
.await
|
||||
.is_ok()
|
||||
{
|
||||
PullFromTar::Pulled(venv_p.to_string())
|
||||
} else {
|
||||
PullFromTar::NotPulled(req.to_string(), venv_p.to_string())
|
||||
let job_id_2 = job_id.clone();
|
||||
let db_2 = db.clone();
|
||||
let w_id_2 = w_id.to_string();
|
||||
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = tokio::time::sleep(tokio::time::Duration::from_secs(1)) => {
|
||||
// Notify server that we are still alive
|
||||
// Detect if job has been canceled
|
||||
let canceled =
|
||||
sqlx::query_scalar::<_, bool>
|
||||
(r#"
|
||||
|
||||
UPDATE queue
|
||||
SET last_ping = now()
|
||||
WHERE id = $1
|
||||
RETURNING canceled
|
||||
|
||||
"#)
|
||||
.bind(job_id_2)
|
||||
.fetch_optional(&db_2)
|
||||
.await
|
||||
.unwrap_or_else(|e| {
|
||||
tracing::error!(%e, "error updating job {job_id_2}: {e:#}");
|
||||
Some(false)
|
||||
})
|
||||
.unwrap_or_else(|| {
|
||||
// if the job is not in queue, it can only be in the completed_job so it is already complete
|
||||
false
|
||||
});
|
||||
|
||||
if canceled {
|
||||
|
||||
tracing::info!(
|
||||
// If there is listener on other side,
|
||||
workspace_id = %w_id_2,
|
||||
"cancelling installations",
|
||||
);
|
||||
|
||||
if let Err(ref e) = kill_tx.send(()){
|
||||
tracing::error!(
|
||||
// If there is listener on other side,
|
||||
workspace_id = %w_id_2,
|
||||
"failed to send done: Probably receiving end closed too early or have not opened yet\n{}",
|
||||
// If there is no listener, it will be dropped safely
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
let results = futures::future::join_all(futures).await;
|
||||
req_with_penv.clear();
|
||||
done_tx.send(()).await.expect("failed to send done");
|
||||
let mut pulled = vec![];
|
||||
for result in results {
|
||||
match result {
|
||||
PullFromTar::Pulled(venv_p) => {
|
||||
pulled.push(venv_p.split("/").last().unwrap_or_default().to_string());
|
||||
req_paths.push(venv_p);
|
||||
}
|
||||
PullFromTar::NotPulled(req, venv_p) => {
|
||||
req_with_penv.push((req, venv_p));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if pulled.len() > 0 {
|
||||
append_logs(
|
||||
&job_id,
|
||||
&w_id,
|
||||
format!(
|
||||
"pulled {} from distributed cache in {}ms",
|
||||
pulled.join(", "),
|
||||
start.elapsed().as_millis()
|
||||
),
|
||||
db,
|
||||
)
|
||||
.await;
|
||||
// Once done_tx is dropped, this will be fired
|
||||
_ = done_rx.recv() => break
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
for (req, venv_p) in req_with_penv {
|
||||
let mut logs1 = String::new();
|
||||
// tl = total_length
|
||||
// "small".len == 5
|
||||
// "middle".len == 6
|
||||
// "largest".len == 7
|
||||
// ==> req_tl = 7
|
||||
let mut req_tl = 0;
|
||||
// Wheels to install
|
||||
let total_to_install = req_with_penv.len();
|
||||
if total_to_install > 0 {
|
||||
let mut logs = String::new();
|
||||
// Do we use UV?
|
||||
if no_uv_install {
|
||||
logs1.push_str("\n\n--- PIP INSTALL ---\n");
|
||||
logs.push_str("\n\n--- PIP INSTALL ---\n");
|
||||
} else {
|
||||
logs1.push_str("\n\n--- UV PIP INSTALL ---\n");
|
||||
logs.push_str("\n\n--- UV PIP INSTALL ---\n");
|
||||
}
|
||||
logs1.push_str(&format!("\n{req} is being installed for the first time.\n It will be cached for all ulterior uses."));
|
||||
append_logs(&job_id, w_id, logs1, db).await;
|
||||
|
||||
logs.push_str("\nTo be installed: \n\n");
|
||||
for (req, _) in &req_with_penv {
|
||||
if req.len() > req_tl {
|
||||
req_tl = req.len();
|
||||
}
|
||||
logs.push_str(&format!("{} \n", &req));
|
||||
}
|
||||
|
||||
// Do we use Nsjail?
|
||||
if !*DISABLE_NSJAIL {
|
||||
logs.push_str(&format!("\nStarting isolated installation... ({} tasks in parallel ) \n", parallel_limit));
|
||||
} else {
|
||||
logs.push_str(&format!("\nStarting installation... ({} tasks in parallel) \n", parallel_limit));
|
||||
}
|
||||
append_logs(&job_id, w_id, logs, db).await;
|
||||
}
|
||||
|
||||
let semaphore = Arc::new(Semaphore::new(parallel_limit));
|
||||
let mut handles = Vec::with_capacity(total_to_install);
|
||||
|
||||
for ((req, venv_p), mut kill_rx) in req_with_penv.iter().zip(kill_rxs.into_iter()) {
|
||||
let permit = semaphore.clone().acquire_owned().await; // Acquire a permit
|
||||
|
||||
if let Err(_) = permit {
|
||||
tracing::error!(
|
||||
workspace_id = %w_id,
|
||||
"Cannot acquire permit on semaphore, that can only mean that semaphore has been closed."
|
||||
);
|
||||
break;
|
||||
}
|
||||
|
||||
let permit = permit.unwrap();
|
||||
|
||||
tracing::info!(
|
||||
workspace_id = %w_id,
|
||||
"started setup python dependencies"
|
||||
);
|
||||
|
||||
let child = if !*DISABLE_NSJAIL {
|
||||
let db = db.clone();
|
||||
let job_id = job_id.clone();
|
||||
let job_dir = job_dir.to_owned();
|
||||
let w_id = w_id.to_owned();
|
||||
let req = req.clone();
|
||||
let venv_p = venv_p.clone();
|
||||
let counter_arc = counter_arc.clone();
|
||||
let pip_indexes = pip_indexes.clone();
|
||||
|
||||
handles.push(task::spawn(async move {
|
||||
// permit will be dropped anyway if this thread exits at any point
|
||||
// so we dont have to drop it manually
|
||||
// but we need to move permit into scope to take ownership
|
||||
let _permit = permit;
|
||||
|
||||
tracing::info!(
|
||||
workspace_id = %w_id,
|
||||
"starting nsjail"
|
||||
// is_ok = out,
|
||||
"started thread to install wheel {}",
|
||||
job_id
|
||||
);
|
||||
let mut vars = vars.clone();
|
||||
let req = req.to_string();
|
||||
vars.push(("REQ", &req));
|
||||
vars.push(("TARGET", &venv_p));
|
||||
let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str());
|
||||
nsjail_cmd
|
||||
.current_dir(job_dir)
|
||||
.env_clear()
|
||||
.envs(vars)
|
||||
.envs(PROXY_ENVS.clone())
|
||||
.args(vec!["--config", "download.config.proto"])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
start_child_process(nsjail_cmd, NSJAIL_PATH.as_str()).await?
|
||||
} else {
|
||||
let fssafe_req = NON_ALPHANUM_CHAR.replace_all(&req, "_").to_string();
|
||||
#[cfg(unix)]
|
||||
let req = if no_uv_install {
|
||||
format!("'{}'", req)
|
||||
} else {
|
||||
req.clone()
|
||||
};
|
||||
|
||||
#[cfg(windows)]
|
||||
let req = format!("{}", req);
|
||||
|
||||
let mut command_args = if no_uv_install {
|
||||
vec![
|
||||
PYTHON_PATH.as_str(),
|
||||
"-m",
|
||||
"pip",
|
||||
"install",
|
||||
&req,
|
||||
"-I",
|
||||
"--no-deps",
|
||||
"--no-color",
|
||||
"--isolated",
|
||||
"--no-warn-conflicts",
|
||||
"--disable-pip-version-check",
|
||||
"-t",
|
||||
venv_p.as_str(),
|
||||
]
|
||||
} else {
|
||||
vec![
|
||||
UV_PATH.as_str(),
|
||||
"pip",
|
||||
"install",
|
||||
&req,
|
||||
"--no-deps",
|
||||
"--no-color",
|
||||
// "-p",
|
||||
// "3.11",
|
||||
// Prevent uv from discovering configuration files.
|
||||
"--no-config",
|
||||
"--link-mode=copy",
|
||||
// TODO: Doublecheck it
|
||||
"--system",
|
||||
// Prefer main index over extra
|
||||
// https://docs.astral.sh/uv/pip/compatibility/#packages-that-exist-on-multiple-indexes
|
||||
// TODO: Use env variable that can be toggled from UI
|
||||
"--index-strategy",
|
||||
"unsafe-best-match",
|
||||
"--target",
|
||||
venv_p.as_str(),
|
||||
"--no-cache",
|
||||
]
|
||||
};
|
||||
let pip_extra_index_url = PIP_EXTRA_INDEX_URL
|
||||
.read()
|
||||
.await
|
||||
.clone()
|
||||
.map(handle_ephemeral_token);
|
||||
|
||||
if let Some(url) = pip_extra_index_url.as_ref() {
|
||||
url.split(",").for_each(|url| {
|
||||
command_args.extend(["--extra-index-url", url]);
|
||||
});
|
||||
}
|
||||
let pip_index_url = PIP_INDEX_URL
|
||||
.read()
|
||||
.await
|
||||
.clone()
|
||||
.map(handle_ephemeral_token);
|
||||
|
||||
if let Some(url) = pip_index_url.as_ref() {
|
||||
command_args.extend(["--index-url", url]);
|
||||
}
|
||||
if let Some(cert_path) = PIP_INDEX_CERT.as_ref() {
|
||||
command_args.extend(["--cert", cert_path]);
|
||||
}
|
||||
if let Some(host) = PIP_TRUSTED_HOST.as_ref() {
|
||||
command_args.extend(["--trusted-host", &host]);
|
||||
}
|
||||
|
||||
let mut envs = vec![("PATH", PATH_ENV.as_str())];
|
||||
|
||||
envs.push(("HOME", HOME_ENV.as_str()));
|
||||
|
||||
tracing::debug!("uv pip install command: {:?}", command_args);
|
||||
|
||||
#[cfg(unix)]
|
||||
let start = std::time::Instant::now();
|
||||
#[cfg(all(feature = "enterprise", feature = "parquet"))]
|
||||
{
|
||||
if no_uv_install {
|
||||
let mut flock_cmd = Command::new(FLOCK_PATH.as_str());
|
||||
flock_cmd
|
||||
.env_clear()
|
||||
.envs(PROXY_ENVS.clone())
|
||||
.envs(envs)
|
||||
.args([
|
||||
"-x",
|
||||
&format!(
|
||||
"{}/{}-{}.lock",
|
||||
LOCK_CACHE_DIR,
|
||||
if no_uv_install { "pip" } else { "py311" },
|
||||
fssafe_req
|
||||
),
|
||||
"--command",
|
||||
&command_args.join(" "),
|
||||
])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
start_child_process(flock_cmd, FLOCK_PATH.as_str()).await?
|
||||
} else {
|
||||
let mut cmd = Command::new(command_args[0]);
|
||||
cmd.env_clear()
|
||||
.envs(PROXY_ENVS.clone())
|
||||
.envs(envs)
|
||||
.args(&command_args[1..])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
start_child_process(cmd, UV_PATH.as_str()).await?
|
||||
if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() {
|
||||
tokio::select! {
|
||||
// Cancel was called on the job
|
||||
_ = kill_rx.recv() => return Err(anyhow::anyhow!("S3 pull was canceled")),
|
||||
|
||||
pull = pull_from_tar(os, venv_p.clone(), no_uv_install) => {
|
||||
if let Err(e) = pull {
|
||||
tracing::info!(
|
||||
workspace_id = %w_id,
|
||||
"No tarball was found on S3 or different problem occured {job_id}:\n{e}",
|
||||
);
|
||||
} else {
|
||||
print_success(
|
||||
true,
|
||||
false,
|
||||
&job_id,
|
||||
&w_id,
|
||||
&req,
|
||||
req_tl,
|
||||
counter_arc,
|
||||
total_to_install,
|
||||
start,
|
||||
db
|
||||
).await;
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
{
|
||||
let installer_path = if no_uv_install { command_args[0] } else { "uv" };
|
||||
let mut cmd: Command = Command::new(&installer_path);
|
||||
cmd.env_clear()
|
||||
.envs(envs)
|
||||
.envs(PROXY_ENVS.clone())
|
||||
.env("SystemRoot", SYSTEM_ROOT.as_str())
|
||||
.env(
|
||||
"TMP",
|
||||
std::env::var("TMP").unwrap_or_else(|_| String::from("/tmp")),
|
||||
let mut uv_install_proccess = match spawn_uv_install(
|
||||
&w_id,
|
||||
&req,
|
||||
&venv_p,
|
||||
&job_dir,
|
||||
pip_indexes,
|
||||
no_uv_install,
|
||||
).await {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
append_logs(
|
||||
&job_id,
|
||||
w_id,
|
||||
format!(
|
||||
"\nError while spawning proccess:\n{e}",
|
||||
),
|
||||
db,
|
||||
)
|
||||
.args(&command_args[1..])
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
start_child_process(cmd, installer_path).await?
|
||||
}
|
||||
};
|
||||
.await;
|
||||
return Err(e.into());
|
||||
}
|
||||
};
|
||||
|
||||
let child = handle_child(
|
||||
&job_id,
|
||||
db,
|
||||
mem_peak,
|
||||
canceled_by,
|
||||
child,
|
||||
false,
|
||||
worker_name,
|
||||
&w_id,
|
||||
&format!("uv pip install {req}"),
|
||||
None,
|
||||
false,
|
||||
occupancy_metrics,
|
||||
)
|
||||
.await;
|
||||
tracing::info!(
|
||||
workspace_id = %w_id,
|
||||
is_ok = child.is_ok(),
|
||||
"finished setting up python dependencies {}",
|
||||
job_id
|
||||
);
|
||||
if child.is_err() {
|
||||
|
||||
let mut stderr = uv_install_proccess
|
||||
.stderr
|
||||
.take()
|
||||
.ok_or(anyhow!("Cannot take stderr from uv_install_proccess"))?;
|
||||
|
||||
tokio::select! {
|
||||
// Canceled
|
||||
_ = kill_rx.recv() => {
|
||||
uv_install_proccess.kill().await?;
|
||||
return Err(anyhow::anyhow!("uv pip install was canceled"));
|
||||
}
|
||||
// Finished
|
||||
exitstatus = uv_install_proccess.wait() => match exitstatus {
|
||||
Ok(status) => if !status.success() {
|
||||
tracing::warn!(
|
||||
workspace_id = %w_id,
|
||||
"uv install {} did not succeed, exit status: {:?}",
|
||||
&req,
|
||||
status.code()
|
||||
);
|
||||
|
||||
let mut buf = String::new();
|
||||
stderr.read_to_string(&mut buf).await.unwrap_or_else(|_|{
|
||||
buf = "Cannot read stderr to string".to_owned();
|
||||
0
|
||||
});
|
||||
|
||||
append_logs(
|
||||
&job_id,
|
||||
w_id,
|
||||
format!(
|
||||
"\nError while installing {}:\n{buf}",
|
||||
&req
|
||||
),
|
||||
db,
|
||||
)
|
||||
.await;
|
||||
return Err(anyhow!(buf));
|
||||
},
|
||||
Err(e) => {
|
||||
tracing::error!(
|
||||
workspace_id = %w_id,
|
||||
"Cannot wait for uv_install_proccess, ExitStatus is Err: {e:?}",
|
||||
);
|
||||
return Err(e.into());
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
print_success(
|
||||
false,
|
||||
true,
|
||||
&job_id,
|
||||
&w_id,
|
||||
&req,
|
||||
req_tl,
|
||||
counter_arc,
|
||||
total_to_install,
|
||||
start,
|
||||
db, //
|
||||
)
|
||||
.await;
|
||||
|
||||
#[cfg(all(feature = "enterprise", feature = "parquet"))]
|
||||
if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() {
|
||||
if matches!(get_license_plan().await, LicensePlan::Pro) {
|
||||
tracing::warn!("S3 cache not available in the pro plan");
|
||||
} else {
|
||||
tokio::spawn(build_tar_and_push(os, venv_p.clone(), no_uv_install));
|
||||
}
|
||||
}
|
||||
|
||||
tracing::info!(
|
||||
workspace_id = %w_id,
|
||||
// is_ok = out,
|
||||
"finished setting up python dependency {}",
|
||||
job_id
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}));
|
||||
}
|
||||
|
||||
let mut failed = false;
|
||||
for (handle, (_, venv_p)) in handles.into_iter().zip(req_with_penv.into_iter()) {
|
||||
if let Err(e) = handle.await.unwrap_or(Err(anyhow!("Problem by joining handle"))) {
|
||||
failed = true;
|
||||
tracing::warn!(
|
||||
workspace_id = %w_id,
|
||||
"Installation failed: {:?}",
|
||||
e
|
||||
);
|
||||
if let Err(e) = fs::remove_dir_all(&venv_p) {
|
||||
tracing::warn!(
|
||||
workspace_id = %w_id,
|
||||
"failed to remove cache dir: {:?}",
|
||||
"Failed to remove cache dir: {:?}",
|
||||
e
|
||||
);
|
||||
}
|
||||
} else {
|
||||
req_paths.push(venv_p);
|
||||
}
|
||||
child?;
|
||||
|
||||
#[cfg(all(feature = "enterprise", feature = "parquet"))]
|
||||
if let Some(os) = OBJECT_STORE_CACHE_SETTINGS.read().await.clone() {
|
||||
if matches!(get_license_plan().await, LicensePlan::Pro) {
|
||||
tracing::warn!("S3 cache not available in the pro plan");
|
||||
} else {
|
||||
let venv_p = venv_p.clone();
|
||||
tokio::spawn(build_tar_and_push(os, venv_p, no_uv_install));
|
||||
}
|
||||
}
|
||||
req_paths.push(venv_p);
|
||||
}
|
||||
Ok(req_paths)
|
||||
|
||||
// Usually done_tx will drop after this return
|
||||
// If there is listener on other side,
|
||||
// it will be triggered
|
||||
// If there is no listener, it will be dropped safely
|
||||
return if failed {
|
||||
Err(anyhow!("Installation did not succeed, check logs").into())
|
||||
} else {
|
||||
Ok(req_paths)
|
||||
};
|
||||
}
|
||||
|
||||
#[cfg(feature = "enterprise")]
|
||||
@@ -1510,8 +1774,6 @@ use crate::JobCompletedSender;
|
||||
#[cfg(feature = "enterprise")]
|
||||
use crate::{common::build_envs_map, dedicated_worker::handle_dedicated_process};
|
||||
#[cfg(feature = "enterprise")]
|
||||
use tokio::sync::mpsc::Receiver;
|
||||
#[cfg(feature = "enterprise")]
|
||||
use windmill_common::variables;
|
||||
|
||||
#[cfg(feature = "enterprise")]
|
||||
@@ -1527,7 +1789,7 @@ pub async fn start_worker(
|
||||
script_path: &str,
|
||||
token: &str,
|
||||
job_completed_tx: JobCompletedSender,
|
||||
jobs_rx: Receiver<std::sync::Arc<QueuedJob>>,
|
||||
jobs_rx: tokio::sync::mpsc::Receiver<std::sync::Arc<QueuedJob>>,
|
||||
killpill_rx: tokio::sync::broadcast::Receiver<()>,
|
||||
) -> error::Result<()> {
|
||||
let mut mem_peak: i32 = 0;
|
||||
|
||||
Reference in New Issue
Block a user