add benchmark variants and fix db pool connection starvation (#7884)

Add new BENCHMARK_KIND variants (sequentialflow, scriptlogs, concurrencylimit,
concurrencykey, mixed, mixed_no_cc) for targeted performance testing. Fix shared
iteration counting across workers using a global atomic counter. Add job_perms
inserts and queue diagnostics for benchmark mode.

Move db connection setup to dedicated module and drop the initial connection pool
before creating the main one, preventing connection starvation when PostgreSQL
max_connections is low.

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-02-10 17:41:47 +01:00
committed by GitHub
parent f2359ee284
commit c2f8d5d686
65 changed files with 1032 additions and 370 deletions
@@ -12,8 +12,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -17,8 +17,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -46,11 +46,11 @@
]
},
"nullable": [
true,
true,
true,
true,
true,
false,
false,
false,
false,
false,
true,
true
]
@@ -30,8 +30,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -122,8 +122,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -40,8 +40,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "\n SELECT\n flow_version.id AS version,\n flow_version.value->>'early_return' as early_return, \n flow_version.value->>'preprocessor_module' IS NOT NULL as has_preprocessor, \n (flow_version.value->>'chat_input_enabled')::boolean as chat_input_enabled, \n flow.tag, \n flow.dedicated_worker, \n flow.on_behalf_of_email, \n flow.edited_by\n FROM \n flow_version\n INNER JOIN flow\n ON flow.path = flow_version.path AND\n flow.workspace_id = flow_version.workspace_id\n WHERE \n flow_version.workspace_id = $1 AND\n flow_version.path = $2 AND\n flow_version.id = $3\n ",
"query": "\n SELECT\n flow_version.id AS version,\n flow_version.value->>'early_return' as early_return,\n flow_version.value->>'preprocessor_module' IS NOT NULL as has_preprocessor,\n (flow_version.value->>'chat_input_enabled')::boolean as chat_input_enabled,\n flow.tag,\n flow.dedicated_worker,\n flow.on_behalf_of_email,\n flow.edited_by\n FROM\n flow_version\n INNER JOIN flow\n ON flow.path = flow_version.path AND\n flow.workspace_id = flow_version.workspace_id\n WHERE\n flow_version.workspace_id = $1 AND\n flow_version.path = $2 AND\n flow_version.id = $3\n ",
"describe": {
"columns": [
{
@@ -62,5 +62,5 @@
false
]
},
"hash": "a7468e9054beed88636786c5495ac3b9d9a6086ae6212ad9237b39a0346d7d26"
"hash": "209dc4c1b91eeab1c12ffcd9f9e16f315c689ca772c736b333dcdf07c8086087"
}
@@ -34,8 +34,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -68,8 +67,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -40,8 +40,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -16,8 +16,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -11,8 +11,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -11,8 +11,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -15,8 +15,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -12,8 +12,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -12,8 +12,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -16,8 +16,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -52,8 +51,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -30,8 +30,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -37,8 +37,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -32,8 +32,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -71,8 +70,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -16,8 +16,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -245,8 +245,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -35,8 +35,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -29,8 +29,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -40,8 +40,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -27,8 +27,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -35,8 +35,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -17,8 +17,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -17,8 +17,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -32,8 +32,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -30,8 +30,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -155,8 +155,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -185,8 +185,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -160,8 +160,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -21,8 +21,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -72,8 +71,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -105,8 +105,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -31,8 +31,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -11,8 +11,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -12,8 +12,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -105,8 +105,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -25,8 +25,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -185,8 +185,7 @@
"sqs",
"gcp",
"mqtt",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -31,8 +31,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -21,8 +21,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -72,8 +71,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -24,8 +24,7 @@
"mqtt",
"gcp",
"default_email",
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -21,8 +21,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
@@ -72,8 +71,7 @@
"name": "native_trigger_service",
"kind": {
"Enum": [
"nextcloud",
"google"
"nextcloud"
]
}
}
+159
View File
@@ -0,0 +1,159 @@
use windmill_common::{
error::{self, Error},
get_database_url, DatabaseUrl,
};
pub const DEFAULT_MAX_CONNECTIONS_SERVER: u32 = 50;
pub const DEFAULT_MAX_CONNECTIONS_WORKER: u32 = 5;
pub const DEFAULT_MAX_CONNECTIONS_INDEXER: u32 = 5;
pub async fn initial_connection() -> Result<sqlx::Pool<sqlx::Postgres>, error::Error> {
let connect_options = get_database_url().await?.connect_options().await?;
sqlx::postgres::PgPoolOptions::new()
.max_connections(2)
.connect_with(connect_options)
.await
.map_err(|err| Error::ConnectingToDatabase(err.to_string()))
}
pub async fn connect_db(
server_mode: bool,
indexer_mode: bool,
worker_mode: bool,
#[cfg(feature = "private")] mut killpill_rx: tokio::sync::broadcast::Receiver<()>,
) -> anyhow::Result<sqlx::Pool<sqlx::Postgres>> {
use anyhow::Context;
let database_url = get_database_url().await?;
let max_connections = match std::env::var("DATABASE_CONNECTIONS") {
Ok(n) => n.parse::<u32>().context("invalid DATABASE_CONNECTIONS")?,
Err(_) => {
if server_mode {
DEFAULT_MAX_CONNECTIONS_SERVER
} else if indexer_mode {
DEFAULT_MAX_CONNECTIONS_INDEXER
} else {
DEFAULT_MAX_CONNECTIONS_WORKER
+ std::env::var("NUM_WORKERS")
.ok()
.map(|x| x.parse().ok())
.flatten()
.unwrap_or(1)
- 1
}
}
};
let pool = connect(database_url.clone(), max_connections, worker_mode).await?;
#[cfg(all(feature = "enterprise", feature = "private"))]
let pool2 = pool.clone();
#[cfg(all(feature = "enterprise", feature = "private"))]
if let DatabaseUrl::IamRds(database_url) = database_url {
tokio::spawn(async move {
loop {
tokio::select! {
_ = killpill_rx.recv() => {
break;
}
_ = tokio::time::sleep(std::time::Duration::from_secs(10)) => {
let needs_refresh = {
let read_guard = database_url.read().await;
read_guard.needs_refresh()
};
if needs_refresh {
let new_url = tokio::time::timeout(std::time::Duration::from_secs(10), get_database_url()).await;
match new_url {
Ok(Ok(new_url)) => {
match new_url.connect_options().await {
Ok(connect_options) => {
pool2.set_connect_options(connect_options);
tracing::info!("Refreshed IAM RDS URL successfully");
}
Err(e) => {
tracing::error!("Error getting IAM RDS connect options, retrying in 10s: {}", e);
continue;
}
}
}
Ok(Err(e)) => {
tracing::error!("Error refreshing IAM RDS URL, trying again in 10s: {}", e);
continue;
}
Err(e) => {
tracing::error!("Timeout after 10s refreshing IAM RDS URL, trying again in 10 seconds: {}", e);
continue;
}
}
}
}
}
}
});
}
Ok(pool)
}
pub async fn connect(
database_url: DatabaseUrl,
max_connections: u32,
worker_mode: bool,
) -> Result<sqlx::Pool<sqlx::Postgres>, error::Error> {
use sqlx::Executor;
use std::time::Duration;
let mut pool_options = sqlx::postgres::PgPoolOptions::new()
.min_connections((max_connections / 5).clamp(1, max_connections))
.max_connections(max_connections)
.max_lifetime(Duration::from_secs(30 * 60)); // 30 mins
if worker_mode {
pool_options = pool_options.idle_timeout(Duration::from_secs(60));
}
pool_options
.after_connect(move |conn, _| {
if worker_mode {
Box::pin(async move {
if let Err(e) = conn
.execute(
r#"
SET enable_seqscan = OFF;
SET statement_timeout = '5min';
SET idle_in_transaction_session_timeout = '10min';
SET tcp_keepalives_idle = 300;
SET tcp_keepalives_interval = 60;
SET tcp_keepalives_count = 10;"#,
)
.await
{
tracing::error!("Error setting postgres settings: {}", e);
}
Ok(())
})
} else {
Box::pin(async move {
if let Err(e) = conn
.execute(
r#"
SET statement_timeout = '5min';
SET idle_in_transaction_session_timeout = '10min';
SET tcp_keepalives_idle = 300;
SET tcp_keepalives_interval = 60;
SET tcp_keepalives_count = 10;"#,
)
.await
{
tracing::error!("Error setting postgres settings: {}", e);
}
Ok(())
})
}
})
.connect_with(
database_url
.connect_options()
.await?
.statement_cache_capacity(400),
)
.await
.map_err(|err| Error::ConnectingToDatabase(err.to_string()))
}
+8 -3
View File
@@ -117,6 +117,7 @@ const BIND_ADDR_ENV: &str = "SERVER_BIND_ADDR";
#[cfg(target_os = "linux")]
mod cgroups;
mod db_connect;
#[cfg(feature = "private")]
pub mod ee;
mod ee_oss;
@@ -665,7 +666,7 @@ async fn windmill_main() -> anyhow::Result<()> {
} else {
println!("Connecting to database...");
let db = windmill_common::initial_connection().await?;
let db = crate::db_connect::initial_connection().await?;
let num_version = sqlx::query_scalar!("SELECT version()").fetch_one(&db).await;
@@ -770,8 +771,12 @@ async fn windmill_main() -> anyhow::Result<()> {
let conn = if mode == Mode::Agent {
conn
} else {
// This time we use a pool of connections
let db = windmill_common::connect_db(
// Drop the initial connection pool before creating the main one.
// With low PostgreSQL max_connections, both pools existing simultaneously
// can exhaust all available connection slots, causing connect_db to hang.
drop(conn);
let db = crate::db_connect::connect_db(
server_mode,
indexer_mode,
worker_mode,
+2 -13
View File
@@ -10,22 +10,11 @@ if [[ "$(uname)" == "Darwin" ]]; then
# Uncomment the git-based samael dependency
sed -i '' 's/^# \(samael = { git="https:\/\/github.com\/njaremko\/samael", rev="464d015e3ae393e4b5dd00b4d6baa1b617de0dd6", features = \["xmlsec"\] }\)/\1/' Cargo.toml
# Run cargo sqlx prepare with deno_core_mac
# Run cargo sqlx prepare with deno_core_mac
echo "Running cargo sqlx prepare with deno_core_mac..."
cargo sqlx prepare --workspace -- --all-targets --features all_sqlx_features,private,deno_core_mac
else
# Run cargo sqlx prepare
echo "Running cargo sqlx prepare..."
cargo sqlx prepare --workspace -- --all-targets --features all_sqlx_features,private
fi
# Undo the samael changes on macOS
if [[ "$(uname)" == "Darwin" ]]; then
echo "Reverting samael changes..."
# Uncomment the version-based samael dependency
sed -i '' 's/^#samael = { version="0.0.14", features = \["xmlsec"\] }/samael = { version="0.0.14", features = ["xmlsec"] }/' Cargo.toml
# Comment out the git-based samael dependency
sed -i '' 's/^\(samael = { git="https:\/\/github.com\/njaremko\/samael", rev="464d015e3ae393e4b5dd00b4d6baa1b617de0dd6", features = \["xmlsec"\] }\)/# \1/' Cargo.toml
cargo sqlx prepare --workspace -- --all-targets --features all_sqlx_features,ee
fi
BIN
View File
Binary file not shown.
+621 -3
View File
@@ -3,7 +3,20 @@ use crate::{
DB,
};
use serde::Serialize;
use std::collections::HashSet;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
use tokio::time::Instant;
use uuid::Uuid;
static BENCHMARK_INITIALIZED: AtomicBool = AtomicBool::new(false);
static SHARED_BENCH_ITERS: std::sync::LazyLock<Arc<AtomicU64>> =
std::sync::LazyLock::new(|| Arc::new(AtomicU64::new(0)));
pub fn shared_bench_iters() -> Arc<AtomicU64> {
SHARED_BENCH_ITERS.clone()
}
#[derive(Serialize)]
pub struct PoolStats {
@@ -36,6 +49,10 @@ pub struct BenchmarkInfo {
pub start: Instant,
#[serde(skip)]
pub iters: u64,
#[serde(skip)]
seen_top_level: HashSet<Uuid>,
#[serde(skip)]
pub shared_iters: Arc<AtomicU64>,
timings: Vec<BenchmarkIter>,
pub iter_durations: Vec<u64>,
pub total_duration: Option<u64>,
@@ -43,9 +60,11 @@ pub struct BenchmarkInfo {
}
impl BenchmarkInfo {
pub fn new() -> Self {
pub fn new(shared_iters: Arc<AtomicU64>) -> Self {
BenchmarkInfo {
iters: 0,
seen_top_level: HashSet::new(),
shared_iters,
timings: vec![],
start: Instant::now(),
iter_durations: vec![],
@@ -64,13 +83,23 @@ impl BenchmarkInfo {
}
}
pub fn add_iter(&mut self, bench: BenchmarkIter, inc_iters: bool) {
if inc_iters {
pub fn count_top_level(&mut self, job_id: Uuid) -> bool {
if self.seen_top_level.insert(job_id) {
self.iters += 1;
return true;
}
false
}
pub fn add_iter(&mut self, bench: BenchmarkIter, job_id: Uuid, is_top_level: bool) -> bool {
let newly_counted = is_top_level && self.seen_top_level.insert(job_id);
if newly_counted {
self.iters += 1;
}
let elapsed_total = bench.start.elapsed().as_nanos() as u64;
self.timings.push(bench);
self.iter_durations.push(elapsed_total);
newly_counted
}
pub fn write_to_file(&mut self, path: &str) -> anyhow::Result<()> {
@@ -118,12 +147,120 @@ impl BenchmarkIter {
}
}
pub async fn benchmark_verify(benchmark_jobs: i32, db: &DB) {
let benchmark_kind = std::env::var("BENCHMARK_KIND").unwrap_or("noop".to_string());
if benchmark_jobs <= 0 || benchmark_kind == "none" {
return;
}
// For flows, child jobs are created dynamically so only check top-level (parent_job IS NULL).
// "parallelflow" inserts only 1 top-level flow regardless of benchmark_jobs.
let expected_top_level = match benchmark_kind.as_str() {
"parallelflow" => 1i64,
_ => benchmark_jobs as i64,
};
let row = sqlx::query!(
"SELECT
COUNT(*) FILTER (WHERE status = 'success') AS succeeded,
COUNT(*) FILTER (WHERE status = 'failure') AS failed,
COUNT(*) FILTER (WHERE status = 'canceled') AS canceled
FROM v2_job_completed
JOIN v2_job USING (id)
WHERE v2_job.workspace_id = 'admins' AND v2_job.parent_job IS NULL",
)
.fetch_one(db)
.await
.expect("benchmark verify query failed");
let succeeded = row.succeeded.unwrap_or(0);
let failed = row.failed.unwrap_or(0);
let canceled = row.canceled.unwrap_or(0);
let total = succeeded + failed + canceled;
let remaining_in_queue = sqlx::query_scalar!(
"SELECT COUNT(*) FROM v2_job_queue WHERE workspace_id = 'admins'",
)
.fetch_one(db)
.await
.expect("benchmark verify queue query failed")
.unwrap_or(0);
println!("=== BENCHMARK VERIFICATION ===");
println!(" kind: {benchmark_kind}");
println!(" expected top-level: {expected_top_level}");
println!(" completed total: {total}");
println!(" succeeded: {succeeded}");
println!(" failed: {failed}");
println!(" canceled: {canceled}");
println!(" still in queue: {remaining_in_queue}");
if failed > 0 || canceled > 0 {
tracing::error!(
"BENCHMARK VERIFICATION FAILED: {failed} failed, {canceled} canceled out of {total} completed"
);
}
if remaining_in_queue > 0 {
tracing::warn!(
"BENCHMARK VERIFICATION: {remaining_in_queue} jobs still in queue after benchmark"
);
}
if succeeded != expected_top_level {
tracing::error!(
"BENCHMARK VERIFICATION FAILED: expected {expected_top_level} succeeded top-level jobs, got {succeeded}"
);
} else if failed == 0 && canceled == 0 && remaining_in_queue == 0 {
println!(" result: ALL PASSED");
}
println!("==============================");
}
pub async fn benchmark_init(benchmark_jobs: i32, db: &DB) {
use crate::{jobs::JobKind, scripts::ScriptLang};
// Only the first worker to reach this point runs init
if BENCHMARK_INITIALIZED.swap(true, Ordering::SeqCst) {
return;
}
let benchmark_kind = std::env::var("BENCHMARK_KIND").unwrap_or("noop".to_string());
if benchmark_jobs > 0 {
// Clean up data from previous benchmark runs
sqlx::query!("DELETE FROM v2_job_completed WHERE workspace_id = 'admins'")
.execute(db)
.await
.unwrap_or_else(|e| panic!("failed to clean up v2_job_completed: {e:#}"));
sqlx::query!("DELETE FROM v2_job_queue WHERE workspace_id = 'admins'")
.execute(db)
.await
.unwrap_or_else(|e| panic!("failed to clean up v2_job_queue: {e:#}"));
sqlx::query!("DELETE FROM v2_job_status WHERE id IN (SELECT id FROM v2_job WHERE workspace_id = 'admins')")
.execute(db)
.await
.unwrap_or_else(|e| panic!("failed to clean up v2_job_status: {e:#}"));
sqlx::query!("DELETE FROM v2_job_runtime WHERE id IN (SELECT id FROM v2_job WHERE workspace_id = 'admins')")
.execute(db)
.await
.unwrap_or_else(|e| panic!("failed to clean up v2_job_runtime: {e:#}"));
sqlx::query("DELETE FROM job_perms WHERE workspace_id = 'admins'")
.execute(db)
.await
.unwrap_or_else(|e| panic!("failed to clean up job_perms: {e:#}"));
sqlx::query!("DELETE FROM concurrency_key WHERE key LIKE 'bench_%' OR key LIKE 'u/admin/bench_%'")
.execute(db)
.await
.unwrap_or_else(|e| panic!("failed to clean up concurrency_key: {e:#}"));
sqlx::query!("DELETE FROM concurrency_counter WHERE concurrency_id LIKE 'bench_%' OR concurrency_id LIKE 'u/admin/bench_%'")
.execute(db)
.await
.unwrap_or_else(|e| panic!("failed to clean up concurrency_counter: {e:#}"));
sqlx::query!("DELETE FROM v2_job WHERE workspace_id = 'admins'")
.execute(db)
.await
.unwrap_or_else(|e| panic!("failed to clean up v2_job: {e:#}"));
let mut tx = db.begin().await.unwrap();
match benchmark_kind.as_str() {
"dedicated" => {
@@ -260,6 +397,477 @@ pub async fn benchmark_init(benchmark_jobs: i32, db: &DB) {
.await
.unwrap_or_else(|_e| panic!("failed to insert parallelflow jobs (4)"));
}
"sequentialflow" => {
let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_flow) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id",
None::<i64>,
None::<String>,
JobKind::FlowPreview as JobKind,
ScriptLang::Deno as ScriptLang,
"flow",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
serde_json::from_str::<serde_json::Value>(r#"
{
"modules": [
{
"id": "a",
"value": {
"type": "rawscript",
"content": "export async function main() { return 'a'; }",
"language": "deno",
"input_transforms": {}
}
},
{
"id": "b",
"value": {
"type": "rawscript",
"content": "export async function main() { return 'b'; }",
"language": "deno",
"input_transforms": {}
}
},
{
"id": "c",
"value": {
"type": "rawscript",
"content": "export async function main() { return 'c'; }",
"language": "deno",
"input_transforms": {}
}
},
{
"id": "d",
"value": {
"type": "rawscript",
"content": "export async function main() { return 'd'; }",
"language": "deno",
"input_transforms": {}
}
}
],
"preprocessor_module": null
}
"#).unwrap(),
benchmark_jobs,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert sequentialflow jobs (1)"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &uuids, "admins", "flow")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert sequentialflow jobs (2)"));
sqlx::query!(
"INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])",
&uuids
)
.execute(&mut *tx)
.await
.unwrap_or_else(|_e| panic!("failed to insert sequentialflow jobs (3)"));
sqlx::query!(
"INSERT INTO v2_job_status (id, flow_status) SELECT unnest($1::uuid[]), $2",
&uuids,
serde_json::from_str::<serde_json::Value>(
r#"
{
"step": 0,
"modules": [
{ "id": "a", "type": "WaitingForPriorSteps" },
{ "id": "b", "type": "WaitingForPriorSteps" },
{ "id": "c", "type": "WaitingForPriorSteps" },
{ "id": "d", "type": "WaitingForPriorSteps" }
],
"cleanup_module": {},
"failure_module": {
"id": "failure",
"type": "WaitingForPriorSteps"
},
"preprocessor_module": null
}
"#
)
.unwrap()
)
.execute(&mut *tx)
.await
.unwrap_or_else(|_e| panic!("failed to insert sequentialflow jobs (4)"));
}
"scriptlogs" => {
let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id",
None::<i64>,
None::<String>,
JobKind::Preview as JobKind,
ScriptLang::Deno as ScriptLang,
"deno",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
"export async function main() { for (let i = 0; i < 1000; i++) { console.log('benchmark log line ' + i); } return 'done'; }",
benchmark_jobs,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert scriptlogs jobs (1)"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &uuids, "admins", "deno")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert scriptlogs jobs (2)"));
sqlx::query!(
"INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])",
&uuids
)
.execute(&mut *tx)
.await
.unwrap_or_else(|_e| panic!("failed to insert scriptlogs jobs (3)"));
}
"concurrencylimit" => {
let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code, concurrent_limit, concurrency_time_window_s) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12 FROM generate_series(1, $13)) RETURNING id",
None::<i64>,
Some("u/admin/bench_conclimit"),
JobKind::Preview as JobKind,
ScriptLang::Deno as ScriptLang,
"deno",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
"export async function main() { return 'done'; }",
2i32,
0i32,
benchmark_jobs,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert concurrencylimit jobs (1)"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &uuids, "admins", "deno")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert concurrencylimit jobs (2)"));
sqlx::query!(
"INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])",
&uuids
)
.execute(&mut *tx)
.await
.unwrap_or_else(|_e| panic!("failed to insert concurrencylimit jobs (3)"));
let concurrency_id = "u/admin/bench_conclimit";
sqlx::query!(
"INSERT INTO concurrency_counter (concurrency_id, job_uuids) VALUES ($1, '{}'::jsonb) ON CONFLICT (concurrency_id) DO NOTHING",
concurrency_id,
)
.execute(&mut *tx)
.await
.unwrap_or_else(|_e| panic!("failed to insert concurrencylimit counter"));
for uuid in &uuids {
sqlx::query!(
"INSERT INTO concurrency_key (key, ended_at, job_id) VALUES ($1, now() - INTERVAL '1 hour', $2)",
concurrency_id,
uuid,
)
.execute(&mut *tx)
.await
.unwrap_or_else(|_e| panic!("failed to insert concurrencylimit key"));
}
}
"concurrencykey" => {
let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code, concurrent_limit, concurrency_time_window_s) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12 FROM generate_series(1, $13)) RETURNING id",
None::<i64>,
Some("u/admin/bench_conckey"),
JobKind::Preview as JobKind,
ScriptLang::Deno as ScriptLang,
"deno",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
"export async function main() { return 'done'; }",
1i32,
0i32,
benchmark_jobs,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert concurrencykey jobs (1)"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &uuids, "admins", "deno")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert concurrencykey jobs (2)"));
sqlx::query!(
"INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])",
&uuids
)
.execute(&mut *tx)
.await
.unwrap_or_else(|_e| panic!("failed to insert concurrencykey jobs (3)"));
let concurrency_id = "bench_shared_concurrency_key";
sqlx::query!(
"INSERT INTO concurrency_counter (concurrency_id, job_uuids) VALUES ($1, '{}'::jsonb) ON CONFLICT (concurrency_id) DO NOTHING",
concurrency_id,
)
.execute(&mut *tx)
.await
.unwrap_or_else(|_e| panic!("failed to insert concurrencykey counter"));
for uuid in &uuids {
sqlx::query!(
"INSERT INTO concurrency_key (key, ended_at, job_id) VALUES ($1, now() - INTERVAL '1 hour', $2)",
concurrency_id,
uuid,
)
.execute(&mut *tx)
.await
.unwrap_or_else(|_e| panic!("failed to insert concurrencykey key"));
}
}
"mixed" => {
let portion = benchmark_jobs / 5;
let remainder = benchmark_jobs % 5;
// 1) noop jobs
let noop_count = portion + remainder;
let noop_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9 FROM generate_series(1, $10)) RETURNING id",
None::<i64>,
None::<String>,
JobKind::Noop as JobKind,
ScriptLang::Deno as ScriptLang,
"deno",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
noop_count,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed noop jobs"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &noop_uuids, "admins", "deno")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed noop queue"));
sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &noop_uuids)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed noop runtime"));
// 2) sequentialflow jobs
if portion > 0 {
let sf_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_flow) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id",
None::<i64>,
None::<String>,
JobKind::FlowPreview as JobKind,
ScriptLang::Deno as ScriptLang,
"flow",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
serde_json::from_str::<serde_json::Value>(r#"{"modules":[{"id":"a","value":{"type":"rawscript","content":"export async function main() { return 'a'; }","language":"deno","input_transforms":{}}},{"id":"b","value":{"type":"rawscript","content":"export async function main() { return 'b'; }","language":"deno","input_transforms":{}}},{"id":"c","value":{"type":"rawscript","content":"export async function main() { return 'c'; }","language":"deno","input_transforms":{}}},{"id":"d","value":{"type":"rawscript","content":"export async function main() { return 'd'; }","language":"deno","input_transforms":{}}}],"preprocessor_module":null}"#).unwrap(),
portion,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed sequentialflow jobs"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &sf_uuids, "admins", "flow")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed sequentialflow queue"));
sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &sf_uuids)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed sequentialflow runtime"));
sqlx::query!(
"INSERT INTO v2_job_status (id, flow_status) SELECT unnest($1::uuid[]), $2",
&sf_uuids,
serde_json::from_str::<serde_json::Value>(r#"{"step":0,"modules":[{"id":"a","type":"WaitingForPriorSteps"},{"id":"b","type":"WaitingForPriorSteps"},{"id":"c","type":"WaitingForPriorSteps"},{"id":"d","type":"WaitingForPriorSteps"}],"cleanup_module":{},"failure_module":{"id":"failure","type":"WaitingForPriorSteps"},"preprocessor_module":null}"#).unwrap()
)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed sequentialflow status"));
}
// 3) scriptlogs jobs
if portion > 0 {
let sl_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id",
None::<i64>,
None::<String>,
JobKind::Preview as JobKind,
ScriptLang::Deno as ScriptLang,
"deno",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
"export async function main() { for (let i = 0; i < 1000; i++) { console.log('benchmark log line ' + i); } return 'done'; }",
portion,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed scriptlogs jobs"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &sl_uuids, "admins", "deno")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed scriptlogs queue"));
sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &sl_uuids)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed scriptlogs runtime"));
}
// 4) concurrencylimit jobs
if portion > 0 {
let cl_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code, concurrent_limit, concurrency_time_window_s) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12 FROM generate_series(1, $13)) RETURNING id",
None::<i64>,
Some("u/admin/bench_conclimit"),
JobKind::Preview as JobKind,
ScriptLang::Deno as ScriptLang,
"deno",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
"export async function main() { return 'done'; }",
2i32,
0i32,
portion,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit jobs"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &cl_uuids, "admins", "deno")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit queue"));
sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &cl_uuids)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit runtime"));
let cl_concurrency_id = "u/admin/bench_conclimit";
sqlx::query!(
"INSERT INTO concurrency_counter (concurrency_id, job_uuids) VALUES ($1, '{}'::jsonb) ON CONFLICT (concurrency_id) DO NOTHING",
cl_concurrency_id,
)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit counter"));
for uuid in &cl_uuids {
sqlx::query!(
"INSERT INTO concurrency_key (key, ended_at, job_id) VALUES ($1, now() - INTERVAL '1 hour', $2)",
cl_concurrency_id,
uuid,
)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencylimit key"));
}
}
// 5) concurrencykey jobs
if portion > 0 {
let ck_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code, concurrent_limit, concurrency_time_window_s) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12 FROM generate_series(1, $13)) RETURNING id",
None::<i64>,
Some("u/admin/bench_conckey"),
JobKind::Preview as JobKind,
ScriptLang::Deno as ScriptLang,
"deno",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
"export async function main() { return 'done'; }",
1i32,
0i32,
portion,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey jobs"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &ck_uuids, "admins", "deno")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey queue"));
sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &ck_uuids)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey runtime"));
let ck_concurrency_id = "bench_shared_concurrency_key";
sqlx::query!(
"INSERT INTO concurrency_counter (concurrency_id, job_uuids) VALUES ($1, '{}'::jsonb) ON CONFLICT (concurrency_id) DO NOTHING",
ck_concurrency_id,
)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey counter"));
for uuid in &ck_uuids {
sqlx::query!(
"INSERT INTO concurrency_key (key, ended_at, job_id) VALUES ($1, now() - INTERVAL '1 hour', $2)",
ck_concurrency_id,
uuid,
)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed concurrencykey key"));
}
}
}
"mixed_no_cc" => {
let portion = benchmark_jobs / 3;
let remainder = benchmark_jobs % 3;
// 1) noop jobs
let noop_count = portion + remainder;
let noop_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9 FROM generate_series(1, $10)) RETURNING id",
None::<i64>,
None::<String>,
JobKind::Noop as JobKind,
ScriptLang::Deno as ScriptLang,
"deno",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
noop_count,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc noop jobs"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &noop_uuids, "admins", "deno")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc noop queue"));
sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &noop_uuids)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc noop runtime"));
// 2) sequentialflow jobs
if portion > 0 {
let sf_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_flow) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id",
None::<i64>,
None::<String>,
JobKind::FlowPreview as JobKind,
ScriptLang::Deno as ScriptLang,
"flow",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
serde_json::from_str::<serde_json::Value>(r#"{"modules":[{"id":"a","value":{"type":"rawscript","content":"export async function main() { return 'a'; }","language":"deno","input_transforms":{}}},{"id":"b","value":{"type":"rawscript","content":"export async function main() { return 'b'; }","language":"deno","input_transforms":{}}},{"id":"c","value":{"type":"rawscript","content":"export async function main() { return 'c'; }","language":"deno","input_transforms":{}}},{"id":"d","value":{"type":"rawscript","content":"export async function main() { return 'd'; }","language":"deno","input_transforms":{}}}],"preprocessor_module":null}"#).unwrap(),
portion,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc sequentialflow jobs"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &sf_uuids, "admins", "flow")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc sequentialflow queue"));
sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &sf_uuids)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc sequentialflow runtime"));
sqlx::query!(
"INSERT INTO v2_job_status (id, flow_status) SELECT unnest($1::uuid[]), $2",
&sf_uuids,
serde_json::from_str::<serde_json::Value>(r#"{"step":0,"modules":[{"id":"a","type":"WaitingForPriorSteps"},{"id":"b","type":"WaitingForPriorSteps"},{"id":"c","type":"WaitingForPriorSteps"},{"id":"d","type":"WaitingForPriorSteps"}],"cleanup_module":{},"failure_module":{"id":"failure","type":"WaitingForPriorSteps"},"preprocessor_module":null}"#).unwrap()
)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc sequentialflow status"));
}
// 3) scriptlogs jobs
if portion > 0 {
let sl_uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id, raw_code) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9, $10 FROM generate_series(1, $11)) RETURNING id",
None::<i64>,
None::<String>,
JobKind::Preview as JobKind,
ScriptLang::Deno as ScriptLang,
"deno",
"admin",
"u/admin",
"admin@windmill.dev",
"admins",
"export async function main() { for (let i = 0; i < 1000; i++) { console.log('benchmark log line ' + i); } return 'done'; }",
portion,
)
.fetch_all(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc scriptlogs jobs"));
sqlx::query!("INSERT INTO v2_job_queue (id, workspace_id, scheduled_for, tag) SELECT unnest($1::uuid[]), $2, now(), $3", &sl_uuids, "admins", "deno")
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc scriptlogs queue"));
sqlx::query!("INSERT INTO v2_job_runtime (id) SELECT unnest($1::uuid[])", &sl_uuids)
.execute(&mut *tx)
.await.unwrap_or_else(|_e| panic!("failed to insert mixed_no_cc scriptlogs runtime"));
}
}
"none" => {}
_ => {
let uuids = sqlx::query_scalar!("INSERT INTO v2_job (id, runnable_id, runnable_path, kind, script_lang, tag, created_by, permissioned_as, permissioned_as_email, workspace_id) (SELECT gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8, $9 FROM generate_series(1, $10)) RETURNING id",
@@ -288,6 +896,16 @@ pub async fn benchmark_init(benchmark_jobs: i32, db: &DB) {
.unwrap_or_else(|_e| panic!("failed to insert noop jobs (3)"));
}
}
// Insert job_perms for all benchmark jobs so workers don't fall back to slow permission lookups
sqlx::query(
"INSERT INTO job_perms (job_id, email, username, is_admin, is_operator, groups, folders, workspace_id)
SELECT id, 'admin@windmill.dev', 'admin', true, false, ARRAY['all']::text[], ARRAY[]::jsonb[], 'admins'
FROM v2_job WHERE workspace_id = 'admins'"
)
.execute(&mut *tx)
.await
.unwrap_or_else(|e| panic!("failed to insert job_perms: {e:#}"));
tx.commit().await.unwrap();
}
}
+12 -161
View File
@@ -664,153 +664,6 @@ pub async fn get_database_url() -> Result<DatabaseUrl, Error> {
Ok(database_url.clone())
}
pub async fn initial_connection() -> Result<sqlx::Pool<sqlx::Postgres>, error::Error> {
let connect_options = get_database_url().await?.connect_options().await?;
sqlx::postgres::PgPoolOptions::new()
.max_connections(2)
.connect_with(connect_options)
.await
.map_err(|err| Error::ConnectingToDatabase(err.to_string()))
}
pub async fn connect_db(
server_mode: bool,
indexer_mode: bool,
worker_mode: bool,
#[cfg(feature = "private")] mut killpill_rx: tokio::sync::broadcast::Receiver<()>,
) -> anyhow::Result<sqlx::Pool<sqlx::Postgres>> {
use anyhow::Context;
let database_url = get_database_url().await?;
let max_connections = match std::env::var("DATABASE_CONNECTIONS") {
Ok(n) => n.parse::<u32>().context("invalid DATABASE_CONNECTIONS")?,
Err(_) => {
if server_mode {
DEFAULT_MAX_CONNECTIONS_SERVER
} else if indexer_mode {
DEFAULT_MAX_CONNECTIONS_INDEXER
} else {
DEFAULT_MAX_CONNECTIONS_WORKER
+ std::env::var("NUM_WORKERS")
.ok()
.map(|x| x.parse().ok())
.flatten()
.unwrap_or(1)
- 1
}
}
};
let pool = connect(database_url.clone(), max_connections, worker_mode).await?;
#[cfg(all(feature = "enterprise", feature = "private"))]
let pool2 = pool.clone();
#[cfg(all(feature = "enterprise", feature = "private"))]
if let DatabaseUrl::IamRds(database_url) = database_url {
tokio::spawn(async move {
loop {
tokio::select! {
_ = killpill_rx.recv() => {
break;
}
_ = tokio::time::sleep(std::time::Duration::from_secs(10)) => {
let needs_refresh = {
let read_guard = database_url.read().await;
read_guard.needs_refresh()
};
if needs_refresh {
let new_url = tokio::time::timeout(std::time::Duration::from_secs(10), get_database_url()).await;
match new_url {
Ok(Ok(new_url)) => {
match new_url.connect_options().await {
Ok(connect_options) => {
pool2.set_connect_options(connect_options);
tracing::info!("Refreshed IAM RDS URL successfully");
}
Err(e) => {
tracing::error!("Error getting IAM RDS connect options, retrying in 10s: {}", e);
continue;
}
}
}
Ok(Err(e)) => {
tracing::error!("Error refreshing IAM RDS URL, trying again in 10s: {}", e);
continue;
}
Err(e) => {
tracing::error!("Timeout after 10s refreshing IAM RDS URL, trying again in 10 seconds: {}", e);
continue;
}
}
}
}
}
}
});
}
Ok(pool)
}
pub async fn connect(
database_url: DatabaseUrl,
max_connections: u32,
worker_mode: bool,
) -> Result<sqlx::Pool<sqlx::Postgres>, error::Error> {
use sqlx::Executor;
use std::time::Duration;
sqlx::postgres::PgPoolOptions::new()
.min_connections((max_connections / 5).clamp(3, max_connections))
.max_connections(max_connections)
.max_lifetime(Duration::from_secs(30 * 60)) // 30 mins
.after_connect(move |conn, _| {
if worker_mode {
Box::pin(async move {
if let Err(e) = conn
.execute(
r#"
SET enable_seqscan = OFF;
SET statement_timeout = '5min';
SET idle_in_transaction_session_timeout = '10min';
SET tcp_keepalives_idle = 300;
SET tcp_keepalives_interval = 60;
SET tcp_keepalives_count = 10;"#,
)
.await
{
tracing::error!("Error setting postgres settings: {}", e);
}
Ok(())
})
} else {
Box::pin(async move {
if let Err(e) = conn
.execute(
r#"
SET statement_timeout = '5min';
SET idle_in_transaction_session_timeout = '10min';
SET tcp_keepalives_idle = 300;
SET tcp_keepalives_interval = 60;
SET tcp_keepalives_count = 10;"#,
)
.await
{
tracing::error!("Error setting postgres settings: {}", e);
}
Ok(())
})
}
})
.connect_with(
database_url
.connect_options()
.await?
.statement_cache_capacity(400),
)
.await
.map_err(|err| Error::ConnectingToDatabase(err.to_string()))
}
type Tag = String;
pub use db::DB;
@@ -852,11 +705,9 @@ impl ScriptHashInfo<ScriptRunnableSettingsHandle> {
self,
db: &DB,
) -> error::Result<ScriptHashInfo<ScriptRunnableSettingsInline>> {
let rs = runnable_settings::from_handle(
self.runnable_settings.runnable_settings_handle,
db,
)
.await?;
let rs =
runnable_settings::from_handle(self.runnable_settings.runnable_settings_handle, db)
.await?;
let (debouncing_settings, concurrency_settings) =
runnable_settings::prefetch_cached(&rs, db).await?;
@@ -1176,25 +1027,25 @@ pub fn get_flow_version_info_from_version<
_ => {
tracing::debug!("Fetching flow version info for {version} ({path})");
let mut conn = db.acquire().await?;
let flow_info =
let flow_info =
sqlx::query_as!(
FlowVersionInfo,
r#"
SELECT
flow_version.id AS version,
flow_version.value->>'early_return' as early_return,
flow_version.value->>'preprocessor_module' IS NOT NULL as has_preprocessor,
(flow_version.value->>'chat_input_enabled')::boolean as chat_input_enabled,
flow.tag,
flow.dedicated_worker,
flow.on_behalf_of_email,
flow_version.value->>'early_return' as early_return,
flow_version.value->>'preprocessor_module' IS NOT NULL as has_preprocessor,
(flow_version.value->>'chat_input_enabled')::boolean as chat_input_enabled,
flow.tag,
flow.dedicated_worker,
flow.on_behalf_of_email,
flow.edited_by
FROM
FROM
flow_version
INNER JOIN flow
ON flow.path = flow_version.path AND
flow.workspace_id = flow_version.workspace_id
WHERE
WHERE
flow_version.workspace_id = $1 AND
flow_version.path = $2 AND
flow_version.id = $3
+7
View File
@@ -145,6 +145,13 @@ lazy_static::lazy_static! {
tracing::info!("Mode not specified, defaulting to standalone");
Mode::Standalone
});
#[cfg(feature = "benchmark")]
let mode = {
if mode != Mode::Worker {
println!("Benchmark mode: forcing MODE=worker");
}
Mode::Worker
};
ModeAndAddons {
indexer: search_addon,
mode,
@@ -68,6 +68,7 @@ async fn process_jc(
stats_map: &JobStatsMap,
killpill_rx: &tokio::sync::broadcast::Receiver<()>,
#[cfg(feature = "benchmark")] bench: &mut BenchmarkIter,
#[cfg(feature = "benchmark")] bench_infos: &mut BenchmarkInfo,
) {
let success: bool = jc.success;
@@ -166,6 +167,13 @@ async fn process_jc(
if let Some(root_job) = root_job {
add_root_flow_job_to_otlp(&root_job, success);
#[cfg(feature = "benchmark")]
if bench_infos.count_top_level(root_job.id) {
bench_infos
.shared_iters
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
}
// Accumulate job stats if duration is available
@@ -209,7 +217,7 @@ pub fn start_background_processor(
let JobCompletedReceiver { bounded_rx, mut killpill_rx, unbounded_rx } = job_completed_rx;
#[cfg(feature = "benchmark")]
let mut infos = BenchmarkInfo::new();
let mut infos = BenchmarkInfo::new(windmill_common::bench::shared_bench_iters());
// Start periodic stats flush task
let db_clone = db.clone();
@@ -278,6 +286,10 @@ pub fn start_background_processor(
jc.job.kind,
JobKind::Dependencies | JobKind::FlowDependencies
);
#[cfg(feature = "benchmark")]
let bench_job_id = jc.job.id;
#[cfg(feature = "benchmark")]
let is_top_level_job = jc.job.parent_job.is_none();
process_jc(
jc,
@@ -291,6 +303,8 @@ pub fn start_background_processor(
&killpill_rx,
#[cfg(feature = "benchmark")]
&mut bench,
#[cfg(feature = "benchmark")]
&mut infos,
)
.warn_after_seconds(10)
.await;
@@ -315,7 +329,9 @@ pub fn start_background_processor(
#[cfg(feature = "benchmark")]
{
infos.add_iter(bench, true);
if infos.add_iter(bench, bench_job_id, is_top_level_job) {
infos.shared_iters.fetch_add(1, Ordering::Relaxed);
}
}
last_processing_duration
.store(time.elapsed().as_secs() as u16, Ordering::SeqCst);
@@ -365,6 +381,12 @@ pub fn start_background_processor(
{
tracing::error!("Error updating flow status after job completion for {flow} on {worker_name}: {e:#}");
}
#[cfg(feature = "benchmark")]
{
if infos.add_iter(bench, flow, true) {
infos.shared_iters.fetch_add(1, Ordering::Relaxed);
}
}
last_processing_duration
.store(time.elapsed().as_secs() as u16, Ordering::SeqCst);
}
+133 -62
View File
@@ -188,7 +188,7 @@ use crate::mssql_executor::do_mssql;
use crate::bigquery_executor::do_bigquery;
#[cfg(feature = "benchmark")]
use windmill_common::bench::{benchmark_init, BenchmarkInfo, BenchmarkIter};
use windmill_common::bench::{benchmark_init, benchmark_verify, BenchmarkInfo, BenchmarkIter};
use windmill_common::add_time;
@@ -1109,51 +1109,60 @@ fn start_interactive_worker_shell(
if let Ok(_) = killpill_rx.try_recv() {
tracing::info!("Received killpill, exiting worker shell");
break;
} else {
let pulled_job = match &conn {
Connection::Sql(db) => {
let common_worker_prefix = retrieve_common_worker_prefix(&worker_name);
let query = ("".to_string(), make_pull_query(&[common_worker_prefix]));
#[cfg(feature = "benchmark")]
let mut bench = windmill_common::bench::BenchmarkIter::new();
}
let job = pull(
&db,
false,
&worker_name,
Some(&query),
let pulled_job = tokio::select! {
_ = killpill_rx.recv() => {
tracing::info!("Received killpill during pull, exiting worker shell");
break;
}
result = async {
match &conn {
Connection::Sql(db) => {
let common_worker_prefix = retrieve_common_worker_prefix(&worker_name);
let query = ("".to_string(), make_pull_query(&[common_worker_prefix]));
#[cfg(feature = "benchmark")]
&mut bench,
)
.await;
let mut bench = windmill_common::bench::BenchmarkIter::new();
use PulledJobResultToJobErr::*;
match job {
Ok(j) => match j.to_pulled_job() {
Ok(j) => Ok(j
.clone()
.map(|job| NextJob::Sql { flow_runners: None, job })),
Err(MissingConcurrencyKey(jc))
| Err(ErrorWhilePreprocessing(jc)) => {
if let Err(err) = job_completed_tx.send_job(jc, true).await {
tracing::error!(
"An error occurred while sending job completed: {:#?}",
err
)
let job = pull(
&db,
false,
&worker_name,
Some(&query),
#[cfg(feature = "benchmark")]
&mut bench,
)
.await;
use PulledJobResultToJobErr::*;
match job {
Ok(j) => match j.to_pulled_job() {
Ok(j) => Ok(j
.clone()
.map(|job| NextJob::Sql { flow_runners: None, job })),
Err(MissingConcurrencyKey(jc))
| Err(ErrorWhilePreprocessing(jc)) => {
if let Err(err) = job_completed_tx.send_job(jc, true).await {
tracing::error!(
"An error occurred while sending job completed: {:#?}",
err
)
}
Ok(None)
}
Ok(None)
}
},
Err(err) => Err(err),
},
Err(err) => Err(err),
}
}
Connection::Http(client) => {
crate::agent_workers::pull_job(&client, None, Some(true))
.await
.map_err(|e| error::Error::InternalErr(e.to_string()))
.map(|x| x.map(|y| NextJob::Http(y)))
}
}
Connection::Http(client) => {
crate::agent_workers::pull_job(&client, None, Some(true))
.await
.map_err(|e| error::Error::InternalErr(e.to_string()))
.map(|x| x.map(|y| NextJob::Http(y)))
}
};
} => result,
};
match pulled_job {
Ok(Some(job)) => {
@@ -1233,7 +1242,6 @@ fn start_interactive_worker_shell(
tokio::time::sleep(Duration::from_millis(*SLEEP_QUEUE * 20)).await;
}
};
}
}
})
}
@@ -1571,6 +1579,7 @@ pub async fn run_worker(
// This is used to wake up the background processor when main loop is done and just waiting for new same workers jobs, and that bg processor is also not processing any jobs, bg processing can exit if no more same worker jobs
let wake_up_notify = Arc::new(tokio::sync::Notify::new());
let stats_map = JobStatsMap::default();
let send_result = match (conn, job_completed_rx) {
(Connection::Sql(db), Some(job_completed_receiver)) => Some(start_background_processor(
job_completed_receiver,
@@ -1615,7 +1624,10 @@ pub async fn run_worker(
let mut started = false;
#[cfg(feature = "benchmark")]
let mut infos = BenchmarkInfo::new();
let mut infos = BenchmarkInfo::new(windmill_common::bench::shared_bench_iters());
#[cfg(feature = "benchmark")]
let mut bench_empty_queue_count: u64 = 0;
#[cfg(feature = "benchmark")]
if let Some(db) = conn.as_sql() {
@@ -1807,15 +1819,60 @@ pub async fn run_worker(
// }
#[cfg(feature = "benchmark")]
if benchmark_jobs > 0 && infos.iters == benchmark_jobs as u64 {
tracing::info!("benchmark finished, exiting");
job_completed_tx
.kill()
.await
.expect("send kill to job completed tx");
break;
} else {
tracing::info!("benchmark not finished, still pulling jobs {}", infos.iters);
{
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);
job_completed_tx
.kill()
.await
.expect("send kill to job completed tx");
killpill_tx.send();
break;
} else if benchmark_jobs > 0 && bench_empty_queue_count > 2000 {
tracing::warn!(
"benchmark stalled: no jobs in queue for 2000 polls, exiting (total iters: {}, worker iters: {}/{})",
total_iters,
infos.iters,
benchmark_jobs
);
job_completed_tx
.kill()
.await
.expect("send kill to job completed tx");
killpill_tx.send();
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<String>, Option<uuid::Uuid>)>(
"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"
)
.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);
for (id, tag, running, kind, parent) in &rows {
tracing::info!(
" pending job: id={id}, tag={tag}, running={running}, kind={}, parent={:?}",
kind.as_deref().unwrap_or("?"), parent
);
}
tracing::info!(
"benchmark not finished (total: {}, worker: {}, queue: {})",
total_iters, infos.iters, total_remaining
);
}
Err(e) => {
tracing::info!("benchmark not finished (total: {}, worker: {}), queue query err: {e}", total_iters, infos.iters);
}
}
}
}
}
let next_job = {
@@ -2029,6 +2086,15 @@ pub async fn run_worker(
match next_job {
Ok(Some(job)) => {
#[cfg(feature = "benchmark")]
{
bench_empty_queue_count = 0;
}
#[cfg(feature = "benchmark")]
let is_top_level_job = job.parent_job.is_none() && !job.kind.is_flow();
#[cfg(feature = "benchmark")]
let bench_job_id = job.id;
#[cfg(feature = "prometheus")]
if let Some(wb) = worker_busy.as_ref() {
wb.set(1);
@@ -2074,7 +2140,7 @@ pub async fn run_worker(
if let Some(db) = conn.as_sql() {
infos.sample_pool(db.size(), db.num_idle() as u32);
}
infos.add_iter(bench, true);
infos.add_iter(bench, bench_job_id, is_top_level_job);
}
continue;
@@ -2147,7 +2213,7 @@ pub async fn run_worker(
if let Some(db) = conn.as_sql() {
infos.sample_pool(db.size(), db.num_idle() as u32);
}
infos.add_iter(bench, true);
infos.add_iter(bench, bench_job_id, is_top_level_job);
}
continue;
@@ -2450,7 +2516,7 @@ pub async fn run_worker(
if let Some(db) = conn.as_sql() {
infos.sample_pool(db.size(), db.num_idle() as u32);
}
infos.add_iter(bench, true);
infos.add_iter(bench, bench_job_id, is_top_level_job);
}
}
}
@@ -2477,11 +2543,12 @@ pub async fn run_worker(
#[cfg(feature = "benchmark")]
{
bench_empty_queue_count += 1;
add_time!(bench, "sleep because empty job queue");
if let Some(db) = conn.as_sql() {
infos.sample_pool(db.size(), db.num_idle() as u32);
}
infos.add_iter(bench, false);
infos.add_iter(bench, uuid::Uuid::nil(), false);
}
#[cfg(feature = "prometheus")]
_timer.map(|timer| {
@@ -2510,13 +2577,6 @@ pub async fn run_worker(
}
}
#[cfg(feature = "benchmark")]
{
infos
.write_to_file("profiling_main.json")
.expect("write to file profiling");
}
drop(dedicated_workers);
let has_dedicated_workers = !dedicated_handles.is_empty();
@@ -2537,6 +2597,17 @@ pub async fn run_worker(
tracing::error!("error in awaiting send_result process: {e:?}")
}
}
#[cfg(feature = "benchmark")]
{
infos
.write_to_file("profiling_main.json")
.expect("write to file profiling");
if let Some(db) = conn.as_sql() {
benchmark_verify(benchmark_jobs, db).await;
}
}
tracing::info!(worker = %worker_name, hostname = %hostname, "waiting for interactive_shell to finish");
if let Some(interactive_shell) = interactive_shell {
match tokio::time::timeout(Duration::from_secs(10), interactive_shell).await {
@@ -149,7 +149,6 @@
})
} else if (kind == 'app') {
const app = await AppService.getAppByPath({ workspace: $workspaceStore!, path })
console.log('app', app)
let result: { kind: Kind; path: string }[] = []
if (app.raw_app) {
const rawAppValue = app.value as { runnables?: Record<string, Runnable> }