feat: agent workers v2 using http (#5588)

This commit is contained in:
Ruben Fiszel
2025-04-10 02:56:11 +02:00
committed by GitHub
parent eabd3d1346
commit 63fa499015
83 changed files with 4852 additions and 2719 deletions

View File

@@ -10,7 +10,6 @@ use anyhow::anyhow;
use itertools::Itertools;
use regex::Regex;
use serde_json::value::RawValue;
use sqlx::{Pool, Postgres};
use tokio::{
fs::{metadata, DirBuilder, File},
io::AsyncReadExt,
@@ -27,15 +26,16 @@ use windmill_common::{
Error::{self},
},
utils::calculate_hash,
worker::{copy_dir_recursively, pad_string, write_file, PythonAnnotations, WORKER_CONFIG},
DB,
worker::{
copy_dir_recursively, pad_string, write_file, Connection, PythonAnnotations, WORKER_CONFIG,
},
};
#[cfg(feature = "enterprise")]
use windmill_common::variables::get_secret_value_as_admin;
use std::env::var;
use windmill_queue::{append_logs, CanceledBy};
use windmill_queue::{append_logs, CanceledBy, PrecomputedAgentInfo};
lazy_static::lazy_static! {
static ref PYTHON_PATH: Option<String> = var("PYTHON_PATH").ok().map(|v| {
@@ -77,6 +77,7 @@ use crate::{
start_child_process, OccupancyMetrics,
},
handle_child::handle_child,
worker_utils::ping_job_status,
AuthedClient, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, INSTANCE_PYTHON_VERSION, NSJAIL_PATH,
PATH_ENV, PIP_EXTRA_INDEX_URL, PIP_INDEX_URL, PROXY_ENVS, PY_INSTALL_DIR, TZ_ENV, UV_CACHE_DIR,
};
@@ -95,7 +96,7 @@ pub enum PyVersion {
}
impl PyVersion {
pub async fn from_instance_version(job_id: &Uuid, w_id: &str, db: &Pool<Postgres>) -> Self {
pub async fn from_instance_version(job_id: &Uuid, w_id: &str, conn: &Connection) -> Self {
let mut err = None;
let pyv = match INSTANCE_PYTHON_VERSION.read().await.clone() {
Some(v) => PyVersion::from_string_with_dots(&v).unwrap_or_else(|| {
@@ -108,7 +109,7 @@ impl PyVersion {
};
if let Some(msg) = err {
append_logs(job_id, w_id, &msg, db).await;
append_logs(job_id, w_id, &msg, conn).await;
tracing::error!(msg);
}
pyv
@@ -211,7 +212,7 @@ impl PyVersion {
job_id: &Uuid,
mem_peak: &mut i32,
// canceled_by: &mut Option<CanceledBy>,
db: &Pool<Postgres>,
conn: &Connection,
worker_name: &str,
w_id: &str,
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
@@ -221,7 +222,7 @@ impl PyVersion {
// }
let res = self
.get_python_inner(job_id, mem_peak, db, worker_name, w_id, occupancy_metrics)
.get_python_inner(job_id, mem_peak, conn, worker_name, w_id, occupancy_metrics)
.await;
if let Err(ref e) = res {
@@ -235,7 +236,7 @@ impl PyVersion {
format!(
"\nError while getting python from uv, falling back to system python: {e:?}"
),
db,
conn,
)
.await;
}
@@ -246,7 +247,7 @@ impl PyVersion {
job_id: &Uuid,
mem_peak: &mut i32,
// canceled_by: &mut Option<CanceledBy>,
db: &Pool<Postgres>,
conn: &Connection,
worker_name: &str,
w_id: &str,
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
@@ -257,7 +258,7 @@ impl PyVersion {
if py_path.is_err() {
// Install it
if let Err(err) = self
.install_python(job_id, mem_peak, db, worker_name, w_id, occupancy_metrics)
.install_python(job_id, mem_peak, conn, worker_name, w_id, occupancy_metrics)
.await
{
tracing::error!("Cannot install python: {err}");
@@ -283,13 +284,13 @@ impl PyVersion {
job_id: &Uuid,
mem_peak: &mut i32,
// canceled_by: &mut Option<CanceledBy>,
db: &Pool<Postgres>,
conn: &Connection,
worker_name: &str,
w_id: &str,
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
) -> error::Result<()> {
let v = self.to_string_with_dot();
append_logs(job_id, w_id, format!("\nINSTALLING PYTHON ({})", v), db).await;
append_logs(job_id, w_id, format!("\nINSTALLING PYTHON ({})", v), conn).await;
// Create dirs for newly installed python
// If we dont do this, NSJAIL will not be able to mount cache
// For the default version directory created during startup (main.rs)
@@ -337,10 +338,10 @@ impl PyVersion {
let child_process = start_child_process(child_cmd, "uv").await?;
append_logs(&job_id, &w_id, logs, db).await;
append_logs(&job_id, &w_id, logs, conn).await;
handle_child(
job_id,
db,
conn,
mem_peak,
&mut None,
child_process,
@@ -459,7 +460,7 @@ pub async fn uv_pip_compile(
mem_peak: &mut i32,
canceled_by: &mut Option<CanceledBy>,
job_dir: &str,
db: &Pool<Postgres>,
conn: &Connection,
worker_name: &str,
w_id: &str,
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
@@ -506,25 +507,27 @@ pub async fn uv_pip_compile(
let requirements = format!("# py{}\n{}", py_version.to_string_no_dot(), requirements);
#[cfg(feature = "enterprise")]
let requirements = replace_pip_secret(db, w_id, &requirements, worker_name, job_id).await?;
let requirements = replace_pip_secret(conn, w_id, &requirements, worker_name, job_id).await?;
let req_hash = format!("py-{}", calculate_hash(&requirements));
if !no_cache {
if let Some(cached) = sqlx::query_scalar!(
"SELECT lockfile FROM pip_resolution_cache WHERE hash = $1",
// Python version is included in hash,
// hash will be the different for every python version
req_hash
)
.fetch_optional(db)
.await?
{
logs.push_str(&format!(
"\nFound cached resolution: {req_hash}, on python version: {}",
py_version.to_string_with_dot()
));
return Ok(cached);
if let Some(db) = conn.as_sql() {
if let Some(cached) = sqlx::query_scalar!(
"SELECT lockfile FROM pip_resolution_cache WHERE hash = $1",
// Python version is included in hash,
// hash will be the different for every python version
req_hash
)
.fetch_optional(db)
.await?
{
logs.push_str(&format!(
"\nFound cached resolution: {req_hash}, on python version: {}",
py_version.to_string_with_dot()
));
return Ok(cached);
}
}
}
@@ -535,7 +538,7 @@ pub async fn uv_pip_compile(
{
// Make sure we have python runtime installed
py_version
.get_python(job_id, mem_peak, db, worker_name, w_id, occupancy_metrics)
.get_python(job_id, mem_peak, conn, worker_name, w_id, occupancy_metrics)
.await?;
let mut args = vec![
@@ -631,10 +634,10 @@ pub async fn uv_pip_compile(
}
let child_process = start_child_process(child_cmd, uv_cmd).await?;
append_logs(&job_id, &w_id, logs, db).await;
append_logs(&job_id, &w_id, logs, conn).await;
handle_child(
job_id,
db,
conn,
mem_peak,
canceled_by,
child_process,
@@ -671,11 +674,13 @@ pub async fn uv_pip_compile(
.collect::<Vec<String>>()
.join("\n")
);
sqlx::query!(
if let Some(db) = conn.as_sql() {
sqlx::query!(
"INSERT INTO pip_resolution_cache (hash, lockfile, expiration) VALUES ($1, $2, now() + ('3 days')::interval) ON CONFLICT (hash) DO UPDATE SET lockfile = $2",
req_hash,
lockfile
).fetch_optional(db).await?;
}
Ok(lockfile)
}
@@ -711,7 +716,7 @@ async fn postinstall(
additional_python_paths: &mut Vec<String>,
job_dir: &str,
job: &MiniPulledJob,
db: &sqlx::Pool<sqlx::Postgres>,
conn: &Connection,
) -> windmill_common::error::Result<()> {
// It is guranteed that additional_python_paths only contains paths within windmill/cache/
// All other paths you would usually expect in PYTHONPATH are NOT included. These are added in downstream
@@ -772,7 +777,7 @@ async fn postinstall(
&job.id,
&job.workspace_id,
"\n\nCopying some packages from cache to job_dir...\n".to_string(),
db,
conn,
)
.await;
// Remove PATHs we just moved
@@ -789,13 +794,20 @@ async fn get_python_path(
job_id: &Uuid,
w_id: &str,
mem_peak: &mut i32,
db: &sqlx::Pool<sqlx::Postgres>,
conn: &Connection,
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
) -> windmill_common::error::Result<String> {
let python_path = if let Some(python_path) = PYTHON_PATH.clone() {
python_path
} else if let Some(python_path) = py_version
.get_python(&job_id, mem_peak, db, worker_name, w_id, occupancy_metrics)
.get_python(
&job_id,
mem_peak,
conn,
worker_name,
w_id,
occupancy_metrics,
)
.await?
{
python_path
@@ -816,7 +828,7 @@ pub async fn handle_python_job(
job: &MiniPulledJob,
mem_peak: &mut i32,
canceled_by: &mut Option<CanceledBy>,
db: &sqlx::Pool<sqlx::Postgres>,
conn: &Connection,
client: &AuthedClient,
parent_runnable_path: Option<String>,
inner_content: &String,
@@ -825,6 +837,7 @@ pub async fn handle_python_job(
envs: HashMap<String, String>,
new_args: &mut Option<HashMap<String, Box<RawValue>>>,
occupancy_metrics: &mut OccupancyMetrics,
precomputed_agent_info: Option<PrecomputedAgentInfo>,
) -> windmill_common::error::Result<Box<RawValue>> {
let script_path = crate::common::use_flow_root_path(job.runnable_path());
@@ -835,12 +848,13 @@ pub async fn handle_python_job(
&job.workspace_id,
&script_path,
&job.id,
db,
conn,
worker_name,
worker_dir,
mem_peak,
canceled_by,
&mut Some(occupancy_metrics),
precomputed_agent_info,
)
.await?;
@@ -852,13 +866,13 @@ pub async fn handle_python_job(
&job.id,
&job.workspace_id,
mem_peak,
db,
conn,
&mut Some(occupancy_metrics),
)
.await?;
if !no_postinstall {
if let Err(e) = postinstall(&mut additional_python_paths, job_dir, job, db).await {
if let Err(e) = postinstall(&mut additional_python_paths, job_dir, job, conn).await {
tracing::error!("Postinstall stage has failed. Reason: {e}");
}
tracing::debug!("Finished deps postinstall stage");
@@ -872,7 +886,7 @@ pub async fn handle_python_job(
"\n\n--- PYTHON ({}) CODE EXECUTION ---\n",
py_version.to_string_with_dot()
),
db,
conn,
)
.await;
}
@@ -901,7 +915,7 @@ pub async fn handle_python_job(
let apply_preprocessor = pre_spread.is_some();
create_args_and_out_file(&client, job, job_dir, db).await?;
create_args_and_out_file(&client, job, job_dir, conn).await?;
tracing::debug!("Finished preparing wrapper");
let preprocessor = if let Some(pre_spread) = pre_spread {
@@ -1004,7 +1018,7 @@ except BaseException as e:
tracing::debug!("Finished writing wrapper");
let mut reserved_variables =
get_reserved_variables(job, &client.token, db, parent_runnable_path).await?;
get_reserved_variables(job, &client.token, conn, parent_runnable_path).await?;
// Add /tmp/windmill/cache/python_xyz/global-site-packages to PYTHONPATH.
// Usefull if certain wheels needs to be preinstalled before execution.
@@ -1129,7 +1143,7 @@ mount {{
handle_child(
&job.id,
db,
conn,
mem_peak,
canceled_by,
child,
@@ -1353,43 +1367,47 @@ async fn prepare_wrapper(
#[cfg(feature = "enterprise")]
async fn replace_pip_secret(
db: &DB,
conn: &Connection,
w_id: &str,
req: &str,
worker_name: &str,
job_id: &Uuid,
) -> error::Result<String> {
if PIP_SECRET_VARIABLE.is_match(req) {
let mut joined = "".to_string();
for req in req.lines() {
let nreq = if PIP_SECRET_VARIABLE.is_match(req) {
let capture = PIP_SECRET_VARIABLE.captures(req);
let variable = capture.unwrap().get(1).unwrap().as_str();
if !variable.contains("/PIP_SECRET_") {
return Err(error::Error::internal_err(format!(
if let Some(db) = conn.as_sql() {
if PIP_SECRET_VARIABLE.is_match(req) {
let mut joined = "".to_string();
for req in req.lines() {
let nreq = if PIP_SECRET_VARIABLE.is_match(req) {
let capture = PIP_SECRET_VARIABLE.captures(req);
let variable = capture.unwrap().get(1).unwrap().as_str();
if !variable.contains("/PIP_SECRET_") {
return Err(error::Error::internal_err(format!(
"invalid secret variable in pip requirements, (last part of path ma): {}",
req
)));
}
let secret = get_secret_value_as_admin(db, w_id, variable).await?;
tracing::info!(
worker = %worker_name,
job_id = %job_id,
workspace_id = %w_id,
"found secret variable in pip requirements: {}",
req
);
PIP_SECRET_VARIABLE
.replace(req, secret.as_str())
.to_string()
} else {
req.to_string()
};
joined.push_str(&nreq);
joined.push_str("\n");
}
}
let secret = get_secret_value_as_admin(db, w_id, variable).await?;
tracing::info!(
worker = %worker_name,
job_id = %job_id,
workspace_id = %w_id,
"found secret variable in pip requirements: {}",
req
);
PIP_SECRET_VARIABLE
.replace(req, secret.as_str())
.to_string()
} else {
req.to_string()
};
joined.push_str(&nreq);
joined.push_str("\n");
}
Ok(joined)
Ok(joined)
} else {
Ok(req.to_string())
}
} else {
Ok(req.to_string())
}
@@ -1402,12 +1420,13 @@ async fn handle_python_deps(
w_id: &str,
script_path: &str,
job_id: &Uuid,
db: &DB,
conn: &Connection,
worker_name: &str,
worker_dir: &str,
mem_peak: &mut i32,
canceled_by: &mut Option<CanceledBy>,
occupancy_metrics: &mut Option<&mut OccupancyMetrics>,
precomputed_agent_info: Option<PrecomputedAgentInfo>,
) -> error::Result<(PyVersion, Vec<String>)> {
create_dependencies_dir(job_dir).await;
@@ -1423,23 +1442,32 @@ async fn handle_python_deps(
let mut annotated_pyv = None;
let mut annotated_pyv_numeric = None;
let is_deployed = requirements_o.is_some();
let instance_pyv = PyVersion::from_instance_version(job_id, w_id, db).await;
let instance_pyv = PyVersion::from_instance_version(job_id, w_id, conn).await;
let annotations = windmill_common::worker::PythonAnnotations::parse(inner_content);
let requirements = match requirements_o {
Some(r) => r,
None => {
let mut already_visited = vec![];
requirements = windmill_parser_py_imports::parse_python_imports(
inner_content,
w_id,
script_path,
db,
&mut already_visited,
&mut annotated_pyv_numeric,
)
.await?
.join("\n");
requirements = match conn {
Connection::Sql(db) => windmill_parser_py_imports::parse_python_imports(
inner_content,
w_id,
script_path,
db,
&mut already_visited,
&mut annotated_pyv_numeric,
)
.await?
.join("\n"),
Connection::Http(_) => match precomputed_agent_info {
Some(PrecomputedAgentInfo::Python { py_version, requirements }) => {
annotated_pyv_numeric = py_version;
requirements.clone().unwrap_or_else(|| "".to_string())
}
_ => "".to_string(),
},
};
annotated_pyv = annotated_pyv_numeric.and_then(|v| PyVersion::from_numeric(v));
@@ -1450,7 +1478,7 @@ async fn handle_python_deps(
mem_peak,
canceled_by,
job_dir,
db,
conn,
worker_name,
w_id,
occupancy_metrics,
@@ -1491,7 +1519,7 @@ async fn handle_python_deps(
w_id,
mem_peak,
canceled_by,
db,
conn,
worker_name,
job_dir,
worker_dir,
@@ -1696,7 +1724,7 @@ pub async fn handle_python_reqs(
w_id: &str,
mem_peak: &mut i32,
_canceled_by: &mut Option<CanceledBy>,
db: &sqlx::Pool<sqlx::Postgres>,
conn: &Connection,
_worker_name: &str,
job_dir: &str,
worker_dir: &str,
@@ -1719,7 +1747,7 @@ pub async fn handle_python_reqs(
counter_arc: Arc<tokio::sync::Mutex<usize>>,
total_to_install: usize,
instant: std::time::Instant,
db: Pool<Postgres>,
conn: &Connection,
) {
#[cfg(not(all(feature = "enterprise", feature = "parquet", unix)))]
{
@@ -1748,7 +1776,7 @@ pub async fn handle_python_reqs(
if s3_push { " > (S3) " } else { "" },
instant.elapsed().as_millis(),
),
db,
conn,
)
.await;
// Drop lock, so next print success can fire
@@ -1810,7 +1838,7 @@ pub async fn handle_python_reqs(
&job_id,
w_id,
format!("\nenv deps from local cache: {}\n", in_cache.join(", ")),
db,
conn,
)
.await;
}
@@ -1824,7 +1852,7 @@ pub async fn handle_python_reqs(
let (_done_tx, mut done_rx) = tokio::sync::mpsc::channel::<()>(1);
let job_id_2 = job_id.clone();
let db_2 = db.clone();
let conn_2 = conn.clone();
let w_id_2 = w_id.to_string();
// Wheels to install
@@ -1874,9 +1902,12 @@ pub async fn handle_python_reqs(
*mem_peak_lock
};
// Notify server that we are still alive
// Detect if job has been canceled
let canceled = sqlx::query_scalar!(
let canceled = match conn_2 {
Connection::Sql(ref db) => {
sqlx::query_scalar!(
"UPDATE v2_job_runtime r SET
memory_peak = $1,
ping = now()
@@ -1885,17 +1916,25 @@ pub async fn handle_python_reqs(
RETURNING canceled_by IS NOT NULL AS \"canceled!\"",
mem_peak_actual,
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
});
)
.fetch_optional(db)
.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
})
}
Connection::Http(_) => {
if let Err(e) = ping_job_status(&conn_2, &job_id_2, Some(mem_peak_actual), None).await {
tracing::error!(%e, "error pinging job {job_id_2}: {e:#}");
}
false
}
};
if canceled {
@@ -1952,7 +1991,7 @@ pub async fn handle_python_reqs(
parallel_limit
));
}
append_logs(&job_id, w_id, logs, db).await;
append_logs(&job_id, w_id, logs, conn).await;
}
let semaphore = Arc::new(Semaphore::new(parallel_limit));
@@ -1964,7 +2003,14 @@ pub async fn handle_python_reqs(
let total_time = std::time::Instant::now();
let py_path = py_version
.get_python(job_id, mem_peak, db, _worker_name, w_id, _occupancy_metrics)
.get_python(
job_id,
mem_peak,
conn,
_worker_name,
w_id,
_occupancy_metrics,
)
.await?;
let has_work = req_with_penv.len() > 0;
@@ -1988,7 +2034,7 @@ pub async fn handle_python_reqs(
"started setup python dependencies"
);
let db = db.clone();
let conn = conn.clone();
let job_id = job_id.clone();
let job_dir = job_dir.to_owned();
let w_id = w_id.to_owned();
@@ -2037,7 +2083,7 @@ pub async fn handle_python_reqs(
counter_arc,
total_to_install,
start,
db
&conn
).await;
pids.lock().await.get_mut(i).and_then(|e| e.take());
@@ -2076,7 +2122,7 @@ pub async fn handle_python_reqs(
format!(
"\nError while spawning proccess:\n{e}",
),
db,
&conn,
)
.await;
pids.lock().await.get_mut(i).and_then(|e| e.take());
@@ -2127,7 +2173,7 @@ pub async fn handle_python_reqs(
"\nError while installing {}:\n{stderr_buf}",
&req
),
db,
&conn,
)
.await;
pids.lock().await.get_mut(i).and_then(|e| e.take());
@@ -2164,7 +2210,7 @@ pub async fn handle_python_reqs(
counter_arc,
total_to_install,
start,
db, //
&conn, //
)
.await;
@@ -2218,7 +2264,13 @@ pub async fn handle_python_reqs(
if has_work {
let total_time = total_time.elapsed().as_millis();
append_logs(&job_id, w_id, format!("\nenv set in {}ms", total_time), db).await;
append_logs(
&job_id,
w_id,
format!("\nenv set in {}ms", total_time),
conn,
)
.await;
}
*mem_peak = *mem_peak_thread_safe.lock().await;
@@ -2285,7 +2337,7 @@ pub async fn start_worker(
let mut mem_peak: i32 = 0;
let mut canceled_by: Option<CanceledBy> = None;
let context = variables::get_reserved_variables(
db,
&Connection::Sql(db.clone()),
w_id,
&token,
"dedicated_worker@windmill.dev",
@@ -2312,12 +2364,13 @@ pub async fn start_worker(
w_id,
script_path,
&Uuid::nil(),
db,
&Connection::Sql(db.clone()),
worker_name,
job_dir,
&mut mem_peak,
&mut canceled_by,
&mut None,
None,
)
.await?;
@@ -2400,7 +2453,7 @@ for line in sys.stdin:
}
let reserved_variables = windmill_common::variables::get_reserved_variables(
db,
&Connection::Sql(db.clone()),
w_id,
token,
"dedicated_worker",
@@ -2442,7 +2495,7 @@ for line in sys.stdin:
&Uuid::nil(),
w_id,
&mut mem_peak,
db,
&Connection::Sql(db.clone()),
&mut None,
)
.await?;