From e2c28e42dbda0f7bf119efc8d587da0d30636a44 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Thu, 12 Feb 2026 14:23:52 +0000 Subject: [PATCH] fix: add private registries support for RUST + java home nit --- backend/src/monitor.rs | 21 +- .../windmill-common/src/global_settings.rs | 2 + backend/windmill-worker/src/java_executor.rs | 9 +- backend/windmill-worker/src/rust_executor.rs | 15 +- backend/windmill-worker/src/worker.rs | 190 ++++++++++-------- .../src/lib/components/instanceSettings.ts | 9 + 6 files changed, 154 insertions(+), 92 deletions(-) diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 09cc8ad66d..c936601624 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -57,10 +57,10 @@ use windmill_common::{ JOB_DEFAULT_TIMEOUT_SECS_SETTING, JWT_SECRET_SETTING, KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, MONITOR_LOGS_ON_OBJECT_STORE_SETTING, NPM_CONFIG_REGISTRY_SETTING, NUGET_CONFIG_SETTING, OTEL_SETTING, OTEL_TRACING_PROXY_SETTING, PIP_INDEX_URL_SETTING, - UV_INDEX_STRATEGY_SETTING, POWERSHELL_REPO_PAT_SETTING, POWERSHELL_REPO_URL_SETTING, REQUEST_SIZE_LIMIT_SETTING, REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING, RETENTION_PERIOD_SECS_SETTING, SAML_METADATA_SETTING, SCIM_TOKEN_SETTING, TIMEOUT_WAIT_RESULT_SETTING, + UV_INDEX_STRATEGY_SETTING, }, indexer::load_indexer_config, jwt::JWT_SECRET, @@ -86,10 +86,10 @@ use windmill_common::{client::AuthedClient, global_settings::APP_WORKSPACED_ROUT use windmill_queue::{cancel_job, get_queued_job_v2, SameWorkerPayload}; use windmill_worker::{ result_processor::handle_job_error, JobCompletedSender, OtelTracingProxySettings, - SameWorkerSender, BUNFIG_INSTALL_SCOPES, INSTANCE_PYTHON_VERSION, JOB_DEFAULT_TIMEOUT, - KEEP_JOB_DIR, MAVEN_REPOS, NO_DEFAULT_MAVEN, NPM_CONFIG_REGISTRY, NUGET_CONFIG, - OTEL_TRACING_PROXY_SETTINGS, PIP_EXTRA_INDEX_URL, PIP_INDEX_URL, POWERSHELL_REPO_PAT, - POWERSHELL_REPO_URL, UV_INDEX_STRATEGY, + SameWorkerSender, BUNFIG_INSTALL_SCOPES, CARGO_REGISTRIES, INSTANCE_PYTHON_VERSION, + JOB_DEFAULT_TIMEOUT, KEEP_JOB_DIR, MAVEN_REPOS, NO_DEFAULT_MAVEN, NPM_CONFIG_REGISTRY, + NUGET_CONFIG, OTEL_TRACING_PROXY_SETTINGS, PIP_EXTRA_INDEX_URL, PIP_INDEX_URL, + POWERSHELL_REPO_PAT, POWERSHELL_REPO_URL, UV_INDEX_STRATEGY, }; #[cfg(feature = "parquet")] @@ -328,6 +328,7 @@ pub async fn initial_load( reload_maven_repos_setting(&conn).await; reload_no_default_maven_setting(&conn).await; reload_ruby_repos_setting(&conn).await; + reload_cargo_registries_setting(&conn).await; } } @@ -1358,6 +1359,16 @@ pub async fn reload_ruby_repos_setting(conn: &Connection) { .await; } +pub async fn reload_cargo_registries_setting(conn: &Connection) { + reload_option_setting_with_tracing( + conn, + windmill_common::global_settings::CARGO_REGISTRIES_SETTING, + "CARGO_REGISTRIES", + CARGO_REGISTRIES.clone(), + ) + .await; +} + pub async fn reload_hub_api_secret_setting(conn: &Connection) { reload_option_setting_with_tracing( conn, diff --git a/backend/windmill-common/src/global_settings.rs b/backend/windmill-common/src/global_settings.rs index f1a1613ed2..3581b6de4d 100644 --- a/backend/windmill-common/src/global_settings.rs +++ b/backend/windmill-common/src/global_settings.rs @@ -16,6 +16,7 @@ pub const POWERSHELL_REPO_PAT_SETTING: &str = "powershell_repo_pat"; pub const MAVEN_REPOS_SETTING: &str = "maven_repos"; pub const NO_DEFAULT_MAVEN_SETTING: &str = "no_default_maven"; pub const RUBY_REPOS_SETTING: &str = "ruby_repos"; +pub const CARGO_REGISTRIES_SETTING: &str = "cargo_registries"; pub const EXTRA_PIP_INDEX_URL_SETTING: &str = "pip_extra_index_url"; pub const PIP_INDEX_URL_SETTING: &str = "pip_index_url"; @@ -82,6 +83,7 @@ pub const ENV_SETTINGS: &[&str] = &[ "GOPRIVATE", "GOPROXY", "NETRC", + "CARGO_REGISTRIES", "INSTANCE_PYTHON_VERSION", "PIP_INDEX_URL", "PIP_EXTRA_INDEX_URL", diff --git a/backend/windmill-worker/src/java_executor.rs b/backend/windmill-worker/src/java_executor.rs index b3c8bfec28..25ec1502ab 100644 --- a/backend/windmill-worker/src/java_executor.rs +++ b/backend/windmill-worker/src/java_executor.rs @@ -26,8 +26,8 @@ use crate::{ }, handle_child, universal_pkg_installer::{par_install_language_dependencies_all_at_once, RequiredDependency}, - COURSIER_CACHE_DIR, DISABLE_NSJAIL, DISABLE_NUSER, JAVA_CACHE_DIR, JAVA_REPOSITORY_DIR, - MAVEN_REPOS, NO_DEFAULT_MAVEN, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, + COURSIER_CACHE_DIR, DISABLE_NSJAIL, DISABLE_NUSER, JAVA_CACHE_DIR, JAVA_HOME_DIR, + JAVA_REPOSITORY_DIR, MAVEN_REPOS, NO_DEFAULT_MAVEN, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, }; use windmill_common::client::AuthedClient; @@ -184,6 +184,7 @@ pub async fn resolve<'a>( cmd.env_clear() .current_dir(job_dir.to_owned()) .env("PATH", PATH_ENV.as_str()) + .env("HOME", JAVA_HOME_DIR) .envs(PROXY_ENVS.clone()); // Configure proxies @@ -330,6 +331,7 @@ async fn install<'a>( cmd.env_clear() .current_dir(&job_dir) .env("PATH", PATH_ENV.as_str()) + .env("HOME", JAVA_HOME_DIR) .envs(PROXY_ENVS.clone()); // Configure proxies { @@ -494,6 +496,7 @@ async fn compile<'a>( cmd.env_clear() .current_dir(job_dir.to_owned()) .env("PATH", PATH_ENV.as_str()) + .env("HOME", JAVA_HOME_DIR) .env("BASE_INTERNAL_URL", base_internal_url) .envs(envs) .envs(reserved_variables) @@ -605,6 +608,7 @@ async fn run<'a>( cmd.env_clear() .current_dir(job_dir) .env("PATH", PATH_ENV.as_str()) + .env("HOME", JAVA_HOME_DIR) .env("BASE_INTERNAL_URL", base_internal_url) .envs(envs) .envs(reserved_variables) @@ -662,6 +666,7 @@ async fn run<'a>( cmd.env_clear() .current_dir(job_dir.to_owned()) .env("PATH", PATH_ENV.as_str()) + .env("HOME", JAVA_HOME_DIR) .env("BASE_INTERNAL_URL", base_internal_url) .envs(envs) .envs(reserved_variables); diff --git a/backend/windmill-worker/src/rust_executor.rs b/backend/windmill-worker/src/rust_executor.rs index 4c195ca24e..158e2d5c5c 100644 --- a/backend/windmill-worker/src/rust_executor.rs +++ b/backend/windmill-worker/src/rust_executor.rs @@ -25,8 +25,8 @@ use crate::{ }, get_proxy_envs_for_lang, handle_child::handle_child, - DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, RUST_CACHE_DIR, - TRACING_PROXY_CA_CERT_PATH, TZ_ENV, + CARGO_REGISTRIES, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, + RUST_CACHE_DIR, TRACING_PROXY_CA_CERT_PATH, TZ_ENV, }; use windmill_common::client::AuthedClient; use windmill_common::scripts::ScriptLang; @@ -135,6 +135,15 @@ pub fn __WINDMILL_RUN__(_args: __WINDMILL_ARGS__) -> Result anyhow::Result<()> { + if let Some(cargo_registries) = CARGO_REGISTRIES.read().await.clone() { + let cargo_dir = format!("{job_dir}/.cargo"); + create_dir_all(&cargo_dir).await?; + write_file(&cargo_dir, "config.toml", &cargo_registries)?; + } + Ok(()) +} + pub async fn generate_cargo_lockfile( job_id: &Uuid, code: &str, @@ -149,6 +158,7 @@ pub async fn generate_cargo_lockfile( check_executor_binary_exists("cargo", CARGO_PATH.as_str(), "rust")?; gen_cargo_crate(code, job_dir)?; + write_cargo_config(job_dir).await?; let mut gen_lockfile_cmd = Command::new(CARGO_PATH.as_str()); gen_lockfile_cmd @@ -503,6 +513,7 @@ pub async fn handle_rust_job( append_logs(&job.id, &job.workspace_id, logs1, conn).await; gen_cargo_crate(inner_content, job_dir)?; + write_cargo_config(job_dir).await?; if let Some(reqs) = requirements_o { if !reqs.is_empty() { diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index fad2878477..9221bfc590 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -215,6 +215,7 @@ pub const CSHARP_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "csharp"); pub const JAVA_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "java"); pub const COURSIER_CACHE_DIR: &str = concatcp!(JAVA_CACHE_DIR, "/coursier-cache"); pub const JAVA_REPOSITORY_DIR: &str = concatcp!(JAVA_CACHE_DIR, "/repository"); +pub const JAVA_HOME_DIR: &str = concatcp!(JAVA_CACHE_DIR, "/home"); // Ruby pub const RUBY_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "ruby"); @@ -581,6 +582,7 @@ lazy_static::lazy_static! { .and_then(|x| x.parse::().ok()) .unwrap_or(false)); pub static ref RUBY_REPOS: Arc>>> = Arc::new(RwLock::new(None)); + pub static ref CARGO_REGISTRIES: Arc>> = Arc::new(RwLock::new(None)); pub static ref PIP_EXTRA_INDEX_URL: Arc>> = Arc::new(RwLock::new(None)); pub static ref PIP_INDEX_URL: Arc>> = Arc::new(RwLock::new(None)); @@ -1165,84 +1167,84 @@ fn start_interactive_worker_shell( } => result, }; - match pulled_job { - Ok(Some(job)) => { - tracing::debug!(target: VERBOSE_TARGET, worker = %worker_name, hostname = %hostname, "started handling of job {}", job.id); - let job_dir = create_job_dir(&worker_dir, job.id).await; + match pulled_job { + Ok(Some(job)) => { + tracing::debug!(target: VERBOSE_TARGET, worker = %worker_name, hostname = %hostname, "started handling of job {}", job.id); + let job_dir = create_job_dir(&worker_dir, job.id).await; + #[cfg(feature = "benchmark")] + let mut bench = windmill_common::bench::BenchmarkIter::new(); + + let JobAndPerms { + job, + raw_code, + raw_lock, + raw_flow, + parent_runnable_path, + token, + precomputed_agent_info: precomputed_bundle, + flow_runners, + } = extract_job_and_perms(job, &conn).await; + + let authed_client = AuthedClient::new( + base_internal_url.to_owned(), + job.workspace_id.clone(), + token, + None, + ); + + let arc_job = Arc::new(job); + + let _ = handle_queued_job( + arc_job.clone(), + raw_code, + raw_lock, + raw_flow, + parent_runnable_path, + &conn, + &authed_client, + &hostname, + &worker_name, + &worker_dir, + &job_dir, + None, + &base_internal_url, + job_completed_tx.clone(), + &mut occupancy_metrics, + &mut killpill_rx, + precomputed_bundle, + flow_runners, #[cfg(feature = "benchmark")] - let mut bench = windmill_common::bench::BenchmarkIter::new(); + &mut bench, + ) + .await; - let JobAndPerms { - job, - raw_code, - raw_lock, - raw_flow, - parent_runnable_path, - token, - precomputed_agent_info: precomputed_bundle, - flow_runners, - } = extract_job_and_perms(job, &conn).await; - - let authed_client = AuthedClient::new( - base_internal_url.to_owned(), - job.workspace_id.clone(), - token, - None, - ); - - let arc_job = Arc::new(job); - - let _ = handle_queued_job( - arc_job.clone(), - raw_code, - raw_lock, - raw_flow, - parent_runnable_path, - &conn, - &authed_client, - &hostname, - &worker_name, - &worker_dir, - &job_dir, - None, - &base_internal_url, - job_completed_tx.clone(), - &mut occupancy_metrics, - &mut killpill_rx, - precomputed_bundle, - flow_runners, - #[cfg(feature = "benchmark")] - &mut bench, - ) - .await; - - last_executed_job = Some(Instant::now()); - } - Ok(None) => { - let now = Instant::now(); - let nap_time = match last_executed_job { - Some(last) - if now.duration_since(last).as_secs() - > TIMEOUT_TO_RESET_WORKER_SHELL_NAP_TIME_DURATION => - { - Duration::from_secs(WORKER_SHELL_NAP_TIME_DURATION) - } - _ => Duration::from_millis(*SLEEP_QUEUE * 10), - }; - tokio::select! { - _ = tokio::time::sleep(nap_time) => { - } - _ = killpill_rx.recv() => { - break; - } + last_executed_job = Some(Instant::now()); + } + Ok(None) => { + let now = Instant::now(); + let nap_time = match last_executed_job { + Some(last) + if now.duration_since(last).as_secs() + > TIMEOUT_TO_RESET_WORKER_SHELL_NAP_TIME_DURATION => + { + Duration::from_secs(WORKER_SHELL_NAP_TIME_DURATION) + } + _ => Duration::from_millis(*SLEEP_QUEUE * 10), + }; + tokio::select! { + _ = tokio::time::sleep(nap_time) => { + } + _ = killpill_rx.recv() => { + break; } } + } - Err(err) => { - tracing::error!(worker = %worker_name, hostname = %hostname, "Failed to pull jobs: {}", err); - tokio::time::sleep(Duration::from_millis(*SLEEP_QUEUE * 20)).await; - } - }; + Err(err) => { + tracing::error!(worker = %worker_name, hostname = %hostname, "Failed to pull jobs: {}", err); + tokio::time::sleep(Duration::from_millis(*SLEEP_QUEUE * 20)).await; + } + }; } }) } @@ -1821,9 +1823,15 @@ pub async fn run_worker( #[cfg(feature = "benchmark")] { - let total_iters = infos.shared_iters.load(std::sync::atomic::Ordering::Relaxed); + let total_iters = infos + .shared_iters + .load(std::sync::atomic::Ordering::Relaxed); if benchmark_jobs > 0 && total_iters >= benchmark_jobs as u64 { - tracing::info!("benchmark finished, exiting (total iters: {}, worker iters: {})", total_iters, infos.iters); + tracing::info!( + "benchmark finished, exiting (total iters: {}, worker iters: {})", + total_iters, + infos.iters + ); job_completed_tx .kill() .await @@ -1845,18 +1853,24 @@ pub async fn run_worker( break; } else if bench_empty_queue_count % 100 == 0 { if let Some(db) = conn.as_sql() { - let remaining = sqlx::query_as::<_, (uuid::Uuid, String, bool, Option, Option)>( + let remaining = sqlx::query_as::< + _, + (uuid::Uuid, String, bool, Option, Option), + >( "SELECT q.id, q.tag, q.running, j.kind::text, j.parent_job FROM v2_job_queue q JOIN v2_job j ON q.id = j.id - WHERE q.workspace_id = 'admins' LIMIT 10" + WHERE q.workspace_id = 'admins' LIMIT 10", ) .fetch_all(db) .await; match remaining { Ok(rows) => { let total_remaining = sqlx::query_scalar::<_, i64>( - "SELECT COUNT(*) FROM v2_job_queue WHERE workspace_id = 'admins'" - ).fetch_one(db).await.unwrap_or(0); + "SELECT COUNT(*) FROM v2_job_queue WHERE workspace_id = 'admins'", + ) + .fetch_one(db) + .await + .unwrap_or(0); for (id, tag, running, kind, parent) in &rows { tracing::info!( " pending job: id={id}, tag={tag}, running={running}, kind={}, parent={:?}", @@ -1865,7 +1879,9 @@ pub async fn run_worker( } tracing::info!( "benchmark not finished (total: {}, worker: {}, queue: {})", - total_iters, infos.iters, total_remaining + total_iters, + infos.iters, + total_remaining ); } Err(e) => { @@ -3112,9 +3128,17 @@ pub async fn handle_queued_job( #[cfg(not(feature = "enterprise"))] if let Connection::Sql(db) = conn { - if (job.concurrent_limit.is_some() || - windmill_common::runnable_settings::prefetch_cached_from_handle(job.runnable_settings_handle, db).await?.1.concurrent_limit.is_some()) - && !job.kind.is_dependency() { + if (job.concurrent_limit.is_some() + || windmill_common::runnable_settings::prefetch_cached_from_handle( + job.runnable_settings_handle, + db, + ) + .await? + .1 + .concurrent_limit + .is_some()) + && !job.kind.is_dependency() + { logs.push_str("---\n"); logs.push_str("WARNING: This job has concurrency limits enabled. Concurrency limits are an EE feature and the setting is ignored.\n"); logs.push_str("---\n"); diff --git a/frontend/src/lib/components/instanceSettings.ts b/frontend/src/lib/components/instanceSettings.ts index e28dc6ec44..82a39762d9 100644 --- a/frontend/src/lib/components/instanceSettings.ts +++ b/frontend/src/lib/components/instanceSettings.ts @@ -450,6 +450,15 @@ export const settings: Record = { storage: 'setting', ee_only: '' }, + { + label: 'Cargo registries', + description: 'Write a .cargo/config.toml to set custom Cargo registries and credentials', + key: 'cargo_registries', + fieldType: 'codearea', + codeAreaLang: 'toml', + storage: 'setting', + ee_only: '' + }, { label: 'PowerShell Repository URL', description: 'Add private PowerShell repository URL',