diff --git a/.env b/.env index 9d065ec296..ef5ce85cc8 100644 --- a/.env +++ b/.env @@ -5,7 +5,7 @@ WM_IMAGE=ghcr.io/windmill-labs/windmill:main WM_LICENSE_KEY="" # For Enterprise Edition, comment the 2 lines above and uncomment below -# WINDMILL_IMAGE=ghcr.io/windmill-labs/windmill-ee:main +# WM_IMAGE=ghcr.io/windmill-labs/windmill-ee:main # WM_LICENSE_KEY=".." diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index d70a1c8788..8cc608c80a 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -1,7 +1,10 @@ use std::sync::Arc; +#[cfg(feature = "enterprise")] use chrono::Timelike; +#[cfg(feature = "enterprise")] use futures::StreamExt; + use futures::{stream, Stream}; use serde::Deserialize; use serde_json::json; @@ -13,7 +16,7 @@ use tokio::{ use windmill_api_client::types::{ CreateFlowBody, EditSchedule, NewSchedule, RawScript, ScriptArgs, }; -use windmill_common::worker::{WORKER_CONFIG, load_worker_config}; +use windmill_common::worker::WORKER_CONFIG; use windmill_common::{ flow_status::{FlowStatus, FlowStatusModule}, flows::{FlowModule, FlowModuleValue, FlowValue, InputTransform}, @@ -126,6 +129,7 @@ impl ApiServer { } async fn close(self) -> anyhow::Result<()> { + println!("closing api server"); let Self { tx, task, .. } = self; drop(tx); task.await.unwrap() @@ -892,7 +896,8 @@ impl RunJob { let uuid = self.push(db).await; let listener = listen_for_completed_jobs(db).await; in_test_worker(db, listener.find(&uuid), port).await; - completed_job(uuid, db).await + let r = completed_job(uuid, db).await; + r } } @@ -933,7 +938,6 @@ async fn in_test_worker( .await .expect("worker timed out") .expect("worker panicked"); - res } @@ -960,9 +964,10 @@ fn spawn_test_worker( let tx2 = tx.clone(); let future = async move { let base_internal_url = format!("http://localhost:{}", port); + { let mut wc = WORKER_CONFIG.write().await; - *wc = load_worker_config(&db).await.unwrap(); - drop(wc); + (*wc).worker_tags = windmill_common::worker::DEFAULT_TAGS.clone(); + } windmill_worker::run_worker::( &db, worker_instance, @@ -1670,7 +1675,6 @@ echo "hello $msg" .arg("msg", json!("world")) .run_until_complete(&db, port) .await; - assert_eq!(job.json_result(), Some(json!("hello world"))); } @@ -2603,6 +2607,7 @@ async fn test_rust_client(db: Pool) { } +#[cfg(feature = "enterprise")] #[sqlx::test(fixtures("base"))] async fn test_script_schedule_handlers(db: Pool) { initialize_tracing().await; @@ -2736,6 +2741,7 @@ async fn test_script_schedule_handlers(db: Pool) { } +#[cfg(feature = "enterprise")] #[sqlx::test(fixtures("base"))] async fn test_flow_schedule_handlers(db: Pool) { initialize_tracing().await; diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index 2451930ab1..59cec4b8f0 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -33,7 +33,7 @@ pub mod worker; pub mod tracing_init; pub const DEFAULT_MAX_CONNECTIONS_SERVER: u32 = 50; -pub const DEFAULT_MAX_CONNECTIONS_WORKER: u32 = 5; +pub const DEFAULT_MAX_CONNECTIONS_WORKER: u32 = 10; lazy_static::lazy_static! { pub static ref METRICS_ADDR: Option = std::env::var("METRICS_ADDR") diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index f29e7ca314..ceca209d90 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -119,16 +119,23 @@ fn process_custom_tags(tags: Vec) -> (Vec, HashMap, pub dedicated_worker: Option, diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 5f32f973d1..dca78b227f 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -1164,9 +1164,9 @@ async fn pull_single_job_and_mark_as_running_no_concurrency_limit< * suspend_until is non-null * and suspend = 0 when the resume messages are received * or suspend_until <= now() if it has timed out */ - let config = WORKER_CONFIG.read().await; - let tags = config.worker_tags.as_slice(); - + let config = WORKER_CONFIG.read().await.clone(); + let tags = config.worker_tags.clone(); + drop(config); let r = if suspend_first { sqlx::query_as::<_, QueuedJob>("UPDATE queue SET running = true @@ -1188,14 +1188,12 @@ async fn pull_single_job_and_mark_as_running_no_concurrency_limit< } else { None }; - drop(config); if r.is_none() { // #[cfg(feature = "benchmark")] // let instant = Instant::now(); - let config = WORKER_CONFIG.read().await; - let tags = config.worker_tags.as_slice(); + let tags = WORKER_CONFIG.read().await.worker_tags.clone(); let r = sqlx::query_as::<_, QueuedJob>( "UPDATE queue diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index f066eccd96..ba88e79349 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -742,13 +742,12 @@ pub async fn run_worker NUM_SECS_PING { - let wc = WORKER_CONFIG.read().await; - let tags = wc.worker_tags.as_slice(); + let tags = WORKER_CONFIG.read().await.worker_tags.clone(); sqlx::query!( "UPDATE worker_ping SET ping_at = now(), jobs_executed = $1, custom_tags = $2 WHERE worker = $3", jobs_executed, - tags, + tags.as_slice(), &worker_name ) .execute(db)