diff --git a/.github/workflows/benchmark.yml b/.github/workflows/benchmark.yml index c794e3e735..f4cb3ce1ba 100644 --- a/.github/workflows/benchmark.yml +++ b/.github/workflows/benchmark.yml @@ -19,15 +19,25 @@ jobs: --health-cmd pg_isready --health-interval 10s --health-timeout 5s --health-retries 5 windmill: - image: ghcr.io/windmill-labs/windmill:main + image: ghcr.io/windmill-labs/windmill-ee:main env: DATABASE_URL: postgres://postgres:changeme@postgres:5432/windmill + LICENSE_KEY: ${{ secrets.WM_LICENSE_KEY_CI }} options: >- --pull always --health-interval 10s --health-timeout 5s --health-retries 5 --health-cmd "curl http://localhost:8000/api/version" ports: - 8000:8000 + windmill-worker: + image: ghcr.io/windmill-labs/windmill-ee:main + env: + DATABASE_URL: postgres://postgres:changeme@postgres:5432/windmill + DISABLE_SERVER: true + DEDICATED_WORKER: "admins:f/benchmarks/dedicated" + LICENSE_KEY: ${{ secrets.WM_LICENSE_KEY_CI }} + options: >- + --pull always steps: - uses: denoland/setup-deno@v1 with: @@ -37,14 +47,9 @@ jobs: ref: benchmarks - name: benchmark timeout-minutes: 10 - run: deno run --unstable -A + run: deno run --unstable -A -r https://raw.githubusercontent.com/windmill-labs/windmill/${GITHUB_REF##ref/head/}/benchmarks/benchmark_suite.ts - --host http://localhost:8000 - -r - -e admin@windmill.dev - -p changeme -c https://raw.githubusercontent.com/windmill-labs/windmill/${GITHUB_REF##ref/head/}/benchmarks/suite_config.json - --branch ${GITHUB_REF##ref/head/} - name: Push changes run: | pwd diff --git a/backend/windmill-api/src/workers.rs b/backend/windmill-api/src/workers.rs index 51f540c91a..d444312fc6 100644 --- a/backend/windmill-api/src/workers.rs +++ b/backend/windmill-api/src/workers.rs @@ -22,11 +22,6 @@ use windmill_common::{ DB, }; -#[cfg(feature = "benchmark")] -use std::sync::atomic::Ordering; -#[cfg(feature = "benchmark")] -use windmill_queue::IDLE_WORKERS; - use crate::{db::ApiAuthed, utils::require_super_admin}; pub fn global_service() -> Router { @@ -40,10 +35,7 @@ pub fn global_service() -> Router { "/worker_group/:name", post(update_worker_group).delete(delete_worker_group), ); - #[cfg(feature = "benchmark")] - return router.route("/toggle", get(toggle)); - #[cfg(not(feature = "benchmark"))] return router; } @@ -85,12 +77,6 @@ async fn list_worker_pings( Ok(Json(rows)) } -#[cfg(feature = "benchmark")] -async fn toggle(Query(query): Query) -> JsonResult { - IDLE_WORKERS.store(query.disable, Ordering::Relaxed); - Ok(Json(IDLE_WORKERS.load(Ordering::Relaxed))) -} - async fn get_custom_tags() -> Json> { Json(ALL_TAGS.read().await.clone().into()) } diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 067ae0c229..ca9cabb890 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -6,7 +6,7 @@ * LICENSE-AGPL for a copy of the license. */ -use std::{collections::HashMap, sync::atomic::AtomicBool, vec}; +use std::{collections::HashMap, vec}; use anyhow::Context; use async_recursion::async_recursion; @@ -70,10 +70,6 @@ lazy_static::lazy_static! { ) .unwrap(); - // When compiled in 'benchmark' mode, this flags is exposed via the /workers/toggle endpoint - // and make it possible to disable to current active workers (such that they don't pull any) - // jobs from the queue - pub static ref IDLE_WORKERS: AtomicBool = AtomicBool::new(false); } #[cfg(feature = "enterprise")] diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index d4466622f1..b3fa1fe62e 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -62,9 +62,6 @@ use crate::bun_executor::start_worker; use windmill_queue::{add_completed_job, add_completed_job_error}; -#[cfg(feature = "benchmark")] -use windmill_queue::IDLE_WORKERS; - use crate::{ bash_executor::{handle_bash_job, handle_powershell_job, ANSI_ESCAPE_RE}, bun_executor::{gen_lockfile, handle_bun_job}, @@ -818,78 +815,67 @@ pub async fn run_worker { - #[cfg(feature = "enterprise")] - if let Some(copy_cache_from_bucket_handle) = copy_cache_from_bucket_handle.as_ref() { - if !copy_cache_from_bucket_handle.is_finished() { - copy_cache_from_bucket_handle.abort(); - } + tokio::select! { + biased; + _ = killpill_rx.recv() => { + #[cfg(feature = "enterprise")] + if let Some(copy_cache_from_bucket_handle) = copy_cache_from_bucket_handle.as_ref() { + if !copy_cache_from_bucket_handle.is_finished() { + copy_cache_from_bucket_handle.abort(); } - #[cfg(feature = "enterprise")] - for handle in &handles { - if !handle.is_finished() { - handle.abort(); - } + } + #[cfg(feature = "enterprise")] + for handle in &handles { + if !handle.is_finished() { + handle.abort(); } - println!("received killpill for worker {}", i_worker); - break - }, - _ = copy_to_bucket_rx.recv() => { - tracing::debug!("can_pull lock start"); - let _lock = CAN_PULL.write().await; - // if num_workers > 1 { - // create_barrier_for_all_workers(num_workers, sync_barrier.clone()).await; - // } - //Arc::new(tokio::sync::Barrier::new(num_workers as usize + 1)); - #[cfg(feature = "enterprise")] - if let Err(e) = copy_tmp_cache_to_cache().await { - tracing::error!(worker = %worker_name, "failed to sync tmp cache to cache: {}", e); - } - tracing::debug!("can_pull lock end"); - Ok(None) - }, - Some(job_id) = same_worker_rx.recv() => { - sqlx::query_as::<_, QueuedJob>("SELECT * FROM queue WHERE id = $1") - .bind(job_id) - .fetch_optional(db) - .await - .map_err(|_| Error::InternalErr("Impossible to fetch same_worker job".to_string())) - }, - (job, timer) = { - let timer = if *METRICS_ENABLED { Some(worker_pull_duration.start_timer()) } else { None }; - let suspend_first = if last_checked_suspended.elapsed().as_secs() > 3 { - last_checked_suspended = Instant::now(); - true - } else { false }; - pull(&db, rsmq.clone(), suspend_first).map(|x| (x, timer)) - } => { - add_time!(timing, loop_start, "post pull"); + } + println!("received killpill for worker {}", i_worker); + break + }, + _ = copy_to_bucket_rx.recv() => { + tracing::debug!("can_pull lock start"); + let _lock = CAN_PULL.write().await; + // if num_workers > 1 { + // create_barrier_for_all_workers(num_workers, sync_barrier.clone()).await; + // } + //Arc::new(tokio::sync::Barrier::new(num_workers as usize + 1)); + #[cfg(feature = "enterprise")] + if let Err(e) = copy_tmp_cache_to_cache().await { + tracing::error!(worker = %worker_name, "failed to sync tmp cache to cache: {}", e); + } + tracing::debug!("can_pull lock end"); + Ok(None) + }, + Some(job_id) = same_worker_rx.recv() => { + sqlx::query_as::<_, QueuedJob>("SELECT * FROM queue WHERE id = $1") + .bind(job_id) + .fetch_optional(db) + .await + .map_err(|_| Error::InternalErr("Impossible to fetch same_worker job".to_string())) + }, + (job, timer) = { + let timer = if *METRICS_ENABLED { Some(worker_pull_duration.start_timer()) } else { None }; + let suspend_first = if last_checked_suspended.elapsed().as_secs() > 3 { + last_checked_suspended = Instant::now(); + true + } else { false }; + pull(&db, rsmq.clone(), suspend_first).map(|x| (x, timer)) + } => { + add_time!(timing, loop_start, "post pull"); - timer.map(|timer| { - let duration_pull_s = timer.stop_and_record(); - worker_pull_duration_counter.inc_by(duration_pull_s); - }); - job + timer.map(|timer| { + let duration_pull_s = timer.stop_and_record(); + worker_pull_duration_counter.inc_by(duration_pull_s); + }); + job - }, - } + }, } }; diff --git a/benchmarks/benchmark_noop.ts b/benchmarks/benchmark_oneoff.ts similarity index 65% rename from benchmarks/benchmark_noop.ts rename to benchmarks/benchmark_oneoff.ts index 8204941dd2..6d51e0ec5b 100644 --- a/benchmarks/benchmark_noop.ts +++ b/benchmarks/benchmark_oneoff.ts @@ -2,25 +2,14 @@ /// import { Command } from "https://deno.land/x/cliffy@v0.25.7/command/mod.ts"; -import { sleep } from "https://deno.land/x/sleep@v1.2.1/mod.ts"; -import * as windmill from "https://deno.land/x/windmill@v1.38.5/mod.ts"; import { UpgradeCommand } from "https://deno.land/x/cliffy@v0.25.7/command/upgrade/upgrade_command.ts"; import { DenoLandProvider } from "https://deno.land/x/cliffy@v0.25.7/command/upgrade/mod.ts"; -export { - DenoLandProvider, - UpgradeCommand, -} from "https://deno.land/x/cliffy@v0.25.7/command/upgrade/mod.ts"; -async function login(email: string, password: string): Promise { - return await windmill.UserService.login({ - requestBody: { - email: email, - password: password, - }, - }); -} +import { sleep } from "https://deno.land/x/sleep@v1.2.1/mod.ts"; -export const VERSION = "v1.167.0"; +import * as windmill from "https://deno.land/x/windmill@v1.174.0/mod.ts"; + +import { VERSION, createBenchScript, getFlowPayload, login } from "./lib.ts"; export async function main({ host, @@ -28,21 +17,21 @@ export async function main({ password, token, workspace, + kind, jobs, - batches, }: { host: string; email?: string; password?: string; token?: string; workspace: string; + kind: string; jobs: number; - batches: number; }) { windmill.setClient("", host); console.log( - "Started benchmark with NOOP jobs with options", + "Started benchmark with options", JSON.stringify( { host, @@ -76,48 +65,62 @@ export async function main({ windmill.setClient(final_token, host); const enc = (s: string) => new TextEncoder().encode(s); - console.log("Disabling workers before loading jobs"); - const disable_workers = await fetch( - config.server + "/api/workers/toggle?disable=true", - { - method: "GET", - headers: { ["Authorization"]: "Bearer " + config.token }, - } - ); - if (!disable_workers.ok) { - console.error( - "Unable to disable workers. Is the Windmill server running in benchmark mode?" - ); + if (["deno", "python", "go", "bash", "dedicated", "bun"].includes(kind)) { + await createBenchScript(kind, workspace); } - const jobsSent = jobs; - const batch_num = batches; - console.log(`Bulk creating ${jobsSent} jobs in ${batch_num} batches`); + let jobsSent = jobs; + console.log(`Bulk creating ${jobsSent} jobs`); const start_create = Date.now(); - const all_create_operations = []; - for (let i = 0; i < batch_num; i++) { - all_create_operations.push( - fetch( - config.server + - "/api/w/" + - config.workspace_id + - `/jobs/add_noop_jobs/${jobsSent / batch_num}`, - { - method: "POST", - headers: { ["Authorization"]: "Bearer " + config.token }, - } - ) - ); + let body: string; + if (kind === "noop") { + body = JSON.stringify({ + kind: "noop", + }); + } else if ( + ["deno", "python", "go", "bash", "dedicated", "bun"].includes(kind) + ) { + body = JSON.stringify({ + kind: "script", + path: "f/benchmarks/" + kind, + dedicated_worker: kind === "dedicated", + }); + } else if (["2steps", "onebranch", "branchallparrallel"].includes(kind)) { + const payload = getFlowPayload(kind); + body = JSON.stringify({ + kind: "flow", + flow_value: payload.value, + }); + } else { + throw new Error("Unknown script pattern " + kind); } - await Promise.all(all_create_operations); + const response = await fetch( + config.server + + "/api/w/" + + config.workspace_id + + `/jobs/add_batch_jobs/${jobsSent}`, + { + method: "POST", + headers: { + ["Authorization"]: "Bearer " + config.token, + "Content-Type": "application/json", + }, + body, + } + ); + if (!response.ok) { + throw new Error("Failed to create jobs: " + response.statusText); + } const end_create = Date.now(); const create_duration = end_create - start_create; console.log( - `Jobs successfully added to the queue in ${create_duration}s. Windmill will start pulling them\n` + `Jobs successfully added to the queue in ${ + create_duration / 1000 + }s. Windmill will start pulling them\n` ); - const start = Date.now(); + let start = Date.now(); let queue_length = jobsSent; let lastElapsed = 0; @@ -151,30 +154,23 @@ export async function main({ }/${jobsSent} (thr: inst ${instThr} - avg ${avgThr}) | queue: ${queue_length} \r` ) ); - }, 100); - - console.log("Enabling workers to start processing jobs"); - const enable_workers = await fetch( - config.server + "/api/workers/toggle?disable=false", - { - method: "GET", - headers: { ["Authorization"]: "Bearer " + config.token }, - } - ); - if (!enable_workers.ok) { - console.error( - "Unable to disable workers. Is the Windmill server running in benchmark mode?" - ); - } + }, 10); while (queue_length > 0) { - await sleep(0.1); + if (queue_length < jobsSent && jobsSent === jobs) { + // reset start time to when the first job was picked up + start = Date.now(); + jobsSent = queue_length; + } + await sleep(0.01); } clearInterval(updateState); const total_duration_sec = (Date.now() - start) / 1000.0; - console.log(`jobs: ${jobsSent}`); + + await sleep(0.1); + console.log(`\njobs: ${jobsSent}`); console.log(`duration: ${total_duration_sec}s`); console.log(`avg. throughput (jobs/time): ${jobsSent / total_duration_sec}`); @@ -204,8 +200,16 @@ if (import.meta.main) { .option("--host ", "The windmill host to benchmark.", { default: "http://127.0.0.1:8000", }) - .option("-e --email ", "The email to use to login.") - .option("-p --password ", "The password to use to login.") + .option("-e --email ", "The email to use to login.", { + default: "admin@windmill.dev", + }) + .option( + "-p --password ", + "The password to use to login.", + { + default: "changeme", + } + ) .env( "WM_TOKEN=", "The token to use when talking to the API server. Preferred over manual login." @@ -223,14 +227,16 @@ if (import.meta.main) { "The workspace to spawn scripts from.", { default: "admins" } ) - .option("-j --jobs ", "Number of NOOP jobs to create.", { + .option( + "--kind ", + "Specifiy the benchmark kind among: deno, identity, python, go, bash, dedicated, bun, noop, 2steps, onebranch, branchallparrallel", + { + required: true, + } + ) + .option("-j --jobs ", "Number of jobs to create.", { default: 10000, }) - .option( - "-b --batches ", - "Number of batches to create all the jobs.", - { default: 1 } - ) .action(main) .command( "upgrade", diff --git a/benchmarks/benchmark_suite.ts b/benchmarks/benchmark_suite.ts index b091f1e9c9..1572f510bf 100644 --- a/benchmarks/benchmark_suite.ts +++ b/benchmarks/benchmark_suite.ts @@ -2,15 +2,23 @@ import { Command } from "https://deno.land/x/cliffy@v0.25.7/command/mod.ts"; import { UpgradeCommand } from "https://deno.land/x/cliffy@v0.25.7/command/upgrade/upgrade_command.ts"; import { DenoLandProvider } from "https://deno.land/x/cliffy@v0.25.7/command/upgrade/mod.ts"; -const VERSION = "1.167.0"; +import { main as runBenchmark } from "./benchmark_oneoff.ts"; + +import { drawGraph, drawGraphMulti } from "./graph.ts"; +import { VERSION } from "./lib.ts"; type Config = { benchmarks: [ { graph_title: string; - name: string; - jobs: number | undefined; - type: "noop" | "flow" | "deno" | "python" | "go" | "bash"; + kind: string; + jobs: number; + } + ]; + extra_graphs?: [ + { + graph_title: string; + kinds: string[]; } ]; }; @@ -22,7 +30,6 @@ async function main({ token, workspace, configPath, - branch, }: { host: string; email?: string; @@ -30,26 +37,7 @@ async function main({ token?: string; workspace: string; configPath: string; - branch?: string; }) { - const { main: runNoopBenchmark } = await import( - branch !== undefined - ? `https://raw.githubusercontent.com/windmill-labs/windmill/${branch}/benchmarks/benchmark_noop.ts` - : "./benchmark_noop.ts" - ); - - const { main: runBenchmark } = await import( - branch !== undefined - ? `https://raw.githubusercontent.com/windmill-labs/windmill/${branch}/benchmarks/main.ts` - : "./main.ts" - ); - - const { drawGraph } = await import( - branch !== undefined - ? `https://raw.githubusercontent.com/windmill-labs/windmill/${branch}/benchmarks/graph.ts` - : "./graph.ts" - ); - async function getConfig(configPath: string): Promise { if (configPath.startsWith("http")) { const response = await fetch(configPath); @@ -64,51 +52,19 @@ async function main({ for (const benchmark of config.benchmarks) { try { console.log( - "%cRunning benchmark " + benchmark.name, + "%cRunning benchmark " + benchmark.kind, "font-weight: bold;" ); - let result: - | { - throughput: number; - } - | undefined; - if (benchmark.type === "noop") { - result = await runNoopBenchmark({ - host, - email, - password, - token, - workspace, - jobs: 1000, - batches: 1, - }); - } else { - result = await runBenchmark({ - host, - email, - password, - token, - workspace, - workers: 1, - seconds: benchmark.type === "flow" ? 2 : 5, - metrics: "http://localhost:8001/metrics", - maximumThroughput: Infinity, - zombieTimeout: 90000, - histogramBuckets: [], - scriptPattern: [ - "deno", - "python", - "go", - "bash", - "dedicated", - ].includes(benchmark.type) - ? benchmark.type - : "deno", - useFlows: benchmark.type === "flow", - hideProgress: true, - }); - } + const result = await runBenchmark({ + host, + email, + password, + token, + workspace, + kind: benchmark.kind, + jobs: benchmark.jobs, + }); if (!result) { throw new Error("No result returned"); @@ -118,7 +74,7 @@ async function main({ ts: Date.now(), }; let data: (typeof stat)[] = []; - const jsonFilePath = `${benchmark.name}.json`; + const jsonFilePath = `${benchmark.kind}_benchmark.json`; try { const existing = await Deno.readTextFile(jsonFilePath); data = JSON.parse(existing); @@ -131,15 +87,42 @@ async function main({ data.slice(-10).map((d) => ({ ...d, date: new Date(d.ts) })), benchmark.graph_title ); - await Deno.writeTextFile(`${benchmark.name}.svg`, svg); + await Deno.writeTextFile(`${benchmark.kind}_benchmark.svg`, svg); } catch (err) { - console.error("Failed to run benchmark", benchmark.name, err); + console.error("Failed to run benchmark", benchmark.kind, err); } } + for (const extraGraph of config.extra_graphs || []) { + const data: { + value: number; + ts: number; + date: Date; + kind: string; + }[] = []; + for (const kind of extraGraph.kinds) { + try { + const existing = await Deno.readTextFile(`${kind}_benchmark.json`); + const existingData = JSON.parse(existing) + .map((d: { value: number; ts: number }) => ({ + ...d, + date: new Date(d.ts), + kind, + })) + .slice(-10); + data.push(...existingData); + } catch (err) { + console.log("Error while loading", kind, "benchmark data", err); + } + } + const svg = drawGraphMulti(data, extraGraph.graph_title); + await Deno.writeTextFile(`${extraGraph.kinds.join("_vs_")}.svg`, svg); + } + Deno.exit(0); // JSDOM from drawGraph doesn't exit cleanly } catch (err) { - return console.error(`Failed to read config file ${configPath}: ${err}`); + console.error(`Failed to read config file ${configPath}: ${err}`); + Deno.exit(0); // JSDOM from drawGraph doesn't exit cleanly } } @@ -150,8 +133,12 @@ await new Command() .option("--host ", "The windmill host to benchmark.", { default: "http://127.0.0.1:8000", }) - .option("-e --email ", "The email to use to login.") - .option("-p --password ", "The password to use to login.") + .option("-e --email ", "The email to use to login.", { + default: "admin@windmill.dev", + }) + .option("-p --password ", "The password to use to login.", { + default: "changeme", + }) .env( "WM_TOKEN=", "The token to use when talking to the API server. Preferred over manual login." @@ -172,10 +159,6 @@ await new Command() .option("-c --config-path ", "The path of the config file", { required: true, }) - .option( - "--branch ", - "The branch to use when running remotely." - ) .action(main) .command( "upgrade", diff --git a/benchmarks/graph.ts b/benchmarks/graph.ts index 474533581c..f6131383c5 100644 --- a/benchmarks/graph.ts +++ b/benchmarks/graph.ts @@ -107,6 +107,167 @@ export function drawGraph(data: DataPoint[], title: string) { return body.node().innerHTML; } +interface DataPointMulti extends DataPoint { + kind: string; +} + +export function drawGraphMulti(data: DataPointMulti[], title: string) { + const context = { + jsdom: new JSDOM(""), + }; + const { window } = context.jsdom; + const { document } = window; + + const body = d3.select(document).select("body"); + + const width = 400; + const height = 200; + + const marginTop = 20; + const marginRight = 100; + const marginBottom = 30; + const marginLeft = 60; + + let svg = body + .append("svg") + .attr("xmlns", "http://www.w3.org/2000/svg") + .attr("width", width + marginLeft + marginRight) + .attr("height", height + marginTop + marginBottom); + + svg + .append("rect") + .attr("width", "100%") + .attr("height", "100%") + .attr("fill", "white"); + + svg = svg + .append("g") + .attr("transform", "translate(" + marginLeft + "," + marginTop + ")"); + + const x = d3 + .scaleTime() + .domain( + d3.extent(data, function (d: DataPoint) { + return d.date; + }) + ) + .nice() + .range([0, width]); + + const xAxis = d3.axisBottom(x).ticks(5); + + svg + .append("g") + .attr("transform", "translate(0," + height + ")") + .call(xAxis); + + // Add Y axis + const y = d3 + .scaleLinear() + .domain([ + 0, + d3.max(data, function (d: DataPoint) { + return +d.value; + }) * 1.5, + ]) + .range([height, 0]) + .nice(); + svg.append("g").call(d3.axisLeft(y)); + + svg + .append("text") + .attr("text-anchor", "middle") + .attr("style", "font-size: 12px") + .attr("transform", "rotate(-90)") + .attr("y", -marginLeft + 20) + .attr("x", -height / 2) + .text("[jobs/s]"); + + svg + .append("text") + .attr("text-anchor", "middle") + .attr("style", "font-size: 16px") + .attr("y", 0) + .attr("x", width / 2) + .text(title); + + const sumstat = d3.group(data, function (d: DataPointMulti) { + return d.kind; + }); + + const keys = Array.from(sumstat.keys()); + + const color = d3 + .scaleOrdinal() + .domain(keys) + .range([ + "#e41a1c", + "#377eb8", + "#4daf4a", + "#984ea3", + "#ff7f00", + "#ffff33", + "#a65628", + "#f781bf", + "#999999", + ]); + + // Add the line + svg + .selectAll("path.line") + .data(sumstat) + .join("path") + .attr("class", "line") + .attr("fill", "none") + .attr("stroke", function (d) { + return color(d[0]); + }) + .attr("stroke-width", 1.5) + .attr("d", (d) => { + return d3 + .line() + .x((d) => x(d.date)) + .y((d) => y(d.value))(d[1]); + }); + + const size = 15; + svg + .selectAll(".dot") + .data(keys) + .enter() + .append("rect") + .attr("class", "dot") + .attr("x", 400) + .attr("y", function (d, i) { + return 5 + i * (size + 5); + }) + .attr("width", size) + .attr("height", size) + .style("fill", function (d) { + return color(d); + }); + svg + .selectAll(".label") + .data(keys) + .enter() + .append("text") + .attr("class", "label") + .attr("x", 400 + size * 1.2) + .attr("y", function (d, i) { + return 5 + i * (size + 5) + size / 2; + }) + .style("fill", function (d) { + return color(d); + }) + .text(function (d) { + return d; + }) + .attr("text-anchor", "left") + .style("alignment-baseline", "middle"); + + return body.node().innerHTML; +} + if (import.meta.main) { const svg = drawGraph( [ @@ -122,6 +283,33 @@ if (import.meta.main) { "test" ); + const svg2 = drawGraphMulti( + [ + { + value: 10, + date: new Date(86400000), + kind: "test", + }, + { + value: 12, + date: new Date(86400000 * 2), + kind: "test", + }, + { + value: 8, + date: new Date(86400000), + kind: "test2", + }, + { + value: 9, + date: new Date(86400000 * 2), + kind: "test2", + }, + ], + "test" + ); + console.log(svg); + console.log(svg2); Deno.exit(0); } diff --git a/benchmarks/lib.ts b/benchmarks/lib.ts new file mode 100644 index 0000000000..c0cdd029ba --- /dev/null +++ b/benchmarks/lib.ts @@ -0,0 +1,244 @@ +import { sleep } from "https://deno.land/x/sleep@v1.2.1/mod.ts"; +import * as windmill from "https://deno.land/x/windmill@v1.174.0/mod.ts"; +import * as api from "https://deno.land/x/windmill@v1.174.0/windmill-api/index.ts"; + +export const VERSION = "v1.174.0"; + +export async function login(email: string, password: string): Promise { + return await windmill.UserService.login({ + requestBody: { + email: email, + password: password, + }, + }); +} + +async function waitForDeployment(workspace: string, hash: string) { + const maxTries = 20; + for (let i = 0; i < maxTries; i++) { + const resp = await windmill.ScriptService.getScriptDeploymentStatus({ + workspace, + hash, + }); + if (resp.lock !== null) { + return; + } + await sleep(0.5); + } + throw new Error("Script did not deploy in time"); +} + +async function waitForDedicatedWorker(workspace: string, path: string) { + const query = windmill.JobService.runWaitResultScriptByPath({ + workspace, + path, + requestBody: { + args: {}, + }, + }); + const timeout = new Promise((_, reject) => { + setTimeout(() => { + reject("Timeout"); + }, 15000); + }); + await Promise.race([query, timeout]); +} + +export async function createBenchScript( + scriptPattern: string, + workspace: string +) { + const path = `f/benchmarks/${scriptPattern}`; + const exists = await windmill.ScriptService.existsScriptByPath({ + workspace, + path, + }); + + if (exists) { + await windmill.ScriptService.deleteScriptByPath({ + workspace, + path, + }); + } + + let scriptContent: string; + let language: string; + if (scriptPattern === "python") { + scriptContent = + 'import os\n\ndef main():\n return os.environ.get("WM_JOB_ID")'; + language = "python3"; + } else if (scriptPattern === "go") { + scriptContent = + 'package inner\nimport "os"\nfunc main() (string, error) { return os.Getenv("WM_JOB_ID"), nil }'; + language = "go"; + } else if (scriptPattern === "bash") { + scriptContent = "echo $WM_JOB_ID"; + language = "bash"; + } else if (scriptPattern === "dedicated" || scriptPattern === "bun") { + scriptContent = 'export function main(){ return Bun.env["WM_JOB_ID"]; }'; + language = "bun"; + } else if (scriptPattern === "deno") { + scriptContent = + 'export function main(){ return Deno.env.get("WM_JOB_ID"); }'; + language = "deno"; + } else { + throw new Error( + "Could not create script for script pattern " + scriptPattern + ); + } + + const hash = await windmill.ScriptService.createScript({ + workspace, + requestBody: { + path, + content: scriptContent, + summary: scriptPattern + " benchmark", + description: "", + language: language as api.NewScript.language, + dedicated_worker: scriptPattern === "dedicated", + }, + }); + + await waitForDeployment(workspace, hash); + + console.log("Created benchmark script at path", path); + + if (scriptPattern === "dedicated") { + await waitForDedicatedWorker(workspace, path); + } +} + +export const getFlowPayload = (flowPattern: string): api.FlowPreview => { + if (flowPattern == "branchone") { + return { + path: "branchone", + args: {}, + value: { + modules: [ + { + id: "a", + value: { + input_transforms: {}, + language: api.RawScript.language.DENO, + type: "rawscript", + content: + 'export function main(){ return Deno.env.get("WM_FLOW_JOB_ID"); }', + }, + }, + { + id: "b", + value: { + type: "branchone", + branches: [], + default: [ + { + id: "c", + value: { + input_transforms: { + x: { + type: "javascript", + expr: "results.a", + }, + }, + language: api.RawScript.language.DENO, + type: "rawscript", + content: "export function main(x: string){ return x; }", + }, + }, + ], + }, + }, + ], + }, + }; + } else if (flowPattern == "branchallparrallel") { + return { + path: "branchall", + args: {}, + value: { + modules: [ + { + id: "a", + value: { + input_transforms: {}, + language: api.RawScript.language.DENO, + type: "rawscript", + content: + 'export function main(){ return Deno.env.get("WM_FLOW_JOB_ID"); }', + }, + }, + { + id: "b", + value: { + type: "branchall", + parallel: true, + branches: [ + { + modules: [ + { + id: "c", + value: { + input_transforms: { + x: { + type: "javascript", + expr: "results.a", + }, + }, + language: api.RawScript.language.DENO, + type: "rawscript", + content: "export function main(x: string){ return x; }", + }, + }, + ], + }, + { + modules: [ + { + id: "d", + value: { + input_transforms: { + x: { + type: "javascript", + expr: "results.a", + }, + }, + language: api.RawScript.language.DENO, + type: "rawscript", + content: "export function main(x: string){ return x; }", + }, + }, + ], + }, + ], + }, + }, + ], + }, + }; + } else { + return { + path: "2steps", + args: {}, + value: { + modules: [ + { + id: "a", + value: { + input_transforms: {}, + language: api.RawScript.language.DENO, + type: "rawscript", + content: + 'export function main(){ return Deno.env.get("WM_FLOW_JOB_ID"); }', + }, + }, + { + id: "b", + value: { + type: "identity", + }, + }, + ], + }, + }; + } +}; diff --git a/benchmarks/main.ts b/benchmarks/main.ts index d3b7d292ac..700be94c27 100644 --- a/benchmarks/main.ts +++ b/benchmarks/main.ts @@ -3,11 +3,11 @@ import { Command } from "https://deno.land/x/cliffy@v0.25.7/command/mod.ts"; import { sleep } from "https://deno.land/x/sleep@v1.2.1/mod.ts"; -import * as windmill from "https://deno.land/x/windmill@v1.167.0/mod.ts"; -import * as api from "https://deno.land/x/windmill@v1.167.0/windmill-api/index.ts"; +import * as windmill from "https://deno.land/x/windmill@v1.174.0/mod.ts"; import { Action } from "./action.ts"; import { UpgradeCommand } from "https://deno.land/x/cliffy@v0.25.7/command/upgrade/upgrade_command.ts"; import { DenoLandProvider } from "https://deno.land/x/cliffy@v0.25.7/command/upgrade/mod.ts"; +import { VERSION, createBenchScript } from "./lib.ts"; export { DenoLandProvider, UpgradeCommand, @@ -22,8 +22,6 @@ async function login(email: string, password: string): Promise { }); } -export const VERSION = "v1.174.0"; - export async function main({ host, workers: num_workers, @@ -46,7 +44,6 @@ export async function main({ continous, max, custom, - hideProgress, }: { host: string; workers: number; @@ -69,7 +66,6 @@ export async function main({ continous?: boolean; max?: number; custom?: string; - hideProgress?: boolean; }) { windmill.setClient("", host); const versionResp = await fetch(`${host}/api/version`); @@ -130,7 +126,6 @@ export async function main({ scriptPattern, zombieTimeout, continous, - hideProgress, }, null, 4 @@ -175,59 +170,16 @@ export async function main({ scriptPattern, continous, custom: custom_content, - hideProgress, }; if ( !useFlows && (scriptPattern === undefined || - ["deno", "python", "go", "bash"].includes(scriptPattern)) + ["deno", "python", "go", "bash", "bun", "dedicated"].includes( + scriptPattern + )) ) { - console.log("Creating benchmark script..."); - const path = `f/benchmarks/${scriptPattern || "deno"}`; - const exists = await windmill.ScriptService.existsScriptByPath({ - workspace, - path, - }); - - if (exists) { - await windmill.ScriptService.deleteScriptByPath({ - workspace, - path, - }); - } - - let scriptContent: string; - let language: string; - if (scriptPattern === "python") { - scriptContent = - 'import os\n\ndef main():\n return os.environ.get("WM_JOB_ID")'; - language = "python3"; - } else if (scriptPattern === "go") { - scriptContent = - 'package inner\nimport "os"\nfunc main() (string, error) { return os.Getenv("WM_JOB_ID"), nil }'; - language = "go"; - } else if (scriptPattern === "bash") { - scriptContent = "echo $WM_JOB_ID"; - language = "bash"; - } else { - scriptContent = - 'export function main(){ return Deno.env.get("WM_JOB_ID"); }'; - language = "deno"; - } - - await windmill.ScriptService.createScript({ - workspace, - requestBody: { - path, - content: scriptContent, - summary: (scriptPattern || "deno") + " benchmark", - description: "", - language: language as api.NewScript.language, - }, - }); - - await sleep(5); // make sure script is created + await createBenchScript(scriptPattern || "deno", workspace); } let workers: Worker[] = new Array(num_workers); @@ -312,13 +264,27 @@ export async function main({ ); const shutdown_start = Date.now(); - let zombie_jobs = 0; - let incorrect_results = 0; + // let zombie_jobs = 0; + // let incorrect_results = 0; + // workers.forEach((worker, i) => { + // const l = (evt: MessageEvent) => { + // if (evt.data.type === "zombie_jobs") { + // zombie_jobs += evt.data.zombie_jobs; + // incorrect_results += evt.data.incorrect_results; + // worker.removeEventListener("message", l); + // workers = workers.filter((w) => w != worker); + // jobsSent[i] = evt.data.jobs_sent; + // worker.terminate(); + // } + // }; + // worker.addEventListener("message", l); + // worker.postMessage( + // Number.isSafeInteger(zombieTimeout) ? zombieTimeout : 90000 + // ); + // }); workers.forEach((worker, i) => { const l = (evt: MessageEvent) => { - if (evt.data.type === "zombie_jobs") { - zombie_jobs += evt.data.zombie_jobs; - incorrect_results += evt.data.incorrect_results; + if (evt.data.type === "done") { worker.removeEventListener("message", l); workers = workers.filter((w) => w != worker); jobsSent[i] = evt.data.jobs_sent; @@ -326,15 +292,32 @@ export async function main({ } }; worker.addEventListener("message", l); - worker.postMessage( - Number.isSafeInteger(zombieTimeout) ? zombieTimeout : 90000 - ); + worker.postMessage("done"); }); console.log("waiting for shutdown\n"); while (workers.length > 0) { await sleep(0.1); } + + let queue_length = await getQueueCount(); + const updateQueue = setInterval(async () => { + queue_length = ( + await ( + await fetch( + host + "/api/w/" + config.workspace_id + "/jobs/queue/count", + { headers: { ["Authorization"]: "Bearer " + config.token } } + ) + ).json() + ).database_length; + await Deno.stdout.write(enc(`queue length: ${queue_length}\r`)); + }, 100); + while (queue_length > 0) { + await sleep(0.1); + } + + clearInterval(updateQueue); + sum = jobsSent.reduce((a, b) => a + b, 0); const tts = (Date.now() - shutdown_start) / 1000; @@ -344,8 +327,8 @@ export async function main({ console.log("time (s + tts):", time); console.log("throughput /s (jobs/time):", sum / time); - console.log("zombie jobs: ", zombie_jobs); - console.log("incorrect results: ", incorrect_results); + // console.log("zombie jobs: ", zombie_jobs); + // console.log("incorrect results: ", incorrect_results); console.log( "queue length:", ( @@ -495,7 +478,7 @@ if (import.meta.main) { ) .option( "--script-pattern ", - "Use a different script pattern among: deno, identity, python, go, bash (Default deno)" + "Use a different script pattern among: deno, identity, python, go, bash, dedicated, bun (Default deno)" ) .option("--custom ", "Use custom actions during bench") .option( diff --git a/benchmarks/suite_config.json b/benchmarks/suite_config.json index 05cbe35d6e..8ccc26fb0f 100644 --- a/benchmarks/suite_config.json +++ b/benchmarks/suite_config.json @@ -1,34 +1,64 @@ { "benchmarks": [ { - "name": "noop_benchmark", "graph_title": "noop throughput benchmark (single worker)", - "type": "noop" + "kind": "noop", + "jobs": 5000 }, { - "name": "flow_benchmark", "graph_title": "flow throughput benchmark (single worker)", - "type": "flow" + "kind": "2steps", + "jobs": 250 + }, + { + "graph_title": "dedicated throughput benchmark (single worker)", + "kind": "dedicated", + "jobs": 2000 }, { - "name": "deno_benchmark", "graph_title": "deno throughput benchmark (single worker)", - "type": "deno" + "kind": "deno", + "jobs": 500 + }, + { + "graph_title": "bun throughput benchmark (single worker)", + "kind": "bun", + "jobs": 500 }, { - "name": "python_benchmark", "graph_title": "python throughput benchmark (single worker)", - "type": "python" + "kind": "python", + "jobs": 500 }, { - "name": "go_benchmark", "graph_title": "go throughput benchmark (single worker)", - "type": "go" + "kind": "go", + "jobs": 500 }, { - "name": "bash_benchmark", "graph_title": "bash throughput benchmark (single worker)", - "type": "bash" + "kind": "bash", + "jobs": 500 + } + ], + "extra_graphs": [ + { + "graph_title": "go vs python vs deno vs bun vs bash", + "kinds": [ + "go", + "python", + "deno", + "bun", + "bash" + ] + }, + { + "graph_title": "bun vs dedicated vs noop", + "kinds": [ + "bun", + "dedicated", + "noop" + ] } ] } \ No newline at end of file diff --git a/benchmarks/worker.ts b/benchmarks/worker.ts index 19db5a1b5c..cd56d0fb45 100644 --- a/benchmarks/worker.ts +++ b/benchmarks/worker.ts @@ -1,10 +1,21 @@ /// /// import { sleep } from "https://deno.land/x/sleep@v1.2.1/sleep.ts"; -import * as windmill from "https://deno.land/x/windmill@v1.167.0/mod.ts"; -import * as api from "https://deno.land/x/windmill@v1.167.0/windmill-api/index.ts"; -import { Job } from "https://deno.land/x/windmill@v1.167.0/windmill-api/index.ts"; +import * as windmill from "https://deno.land/x/windmill@v1.174.0/mod.ts"; +import * as api from "https://deno.land/x/windmill@v1.174.0/windmill-api/index.ts"; import { Action, evaluate } from "./action.ts"; +import { getFlowPayload } from "./lib.ts"; + +async function getQueueCount() { + return ( + await ( + await fetch( + config.server + "/api/w/" + config.workspace_id + "/jobs/queue/count", + { headers: { ["Authorization"]: "Bearer " + config.token } } + ) + ).json() + ).database_length; +} const promise = new Promise<{ workspace_id: string; @@ -46,11 +57,11 @@ let cont = true; let total_spawned = 0; const start_time: number = Date.now(); -let complete_timeout = Infinity; +// let complete_timeout = Infinity; self.onmessage = (evt) => { cont = false; - complete_timeout = evt.data; + // complete_timeout = evt.data; }; const updateStatusInterval = setInterval(() => { @@ -85,141 +96,8 @@ while (cont) { await evaluate(config.custom); continue; } else if (config.useFlows) { - let payload: api.FlowPreview; - if (config.flowPattern == "branchone") { - payload = { - path: "branchone", - args: {}, - value: { - modules: [ - { - id: "a", - value: { - input_transforms: {}, - language: api.RawScript.language.DENO, - type: "rawscript", - content: - 'export function main(){ return Deno.env.get("WM_FLOW_JOB_ID"); }', - }, - }, - { - id: "b", - value: { - type: "branchone", - branches: [], - default: [ - { - id: "c", - value: { - input_transforms: { - x: { - type: "javascript", - expr: "results.a", - }, - }, - language: api.RawScript.language.DENO, - type: "rawscript", - content: "export function main(x: string){ return x; }", - }, - }, - ], - }, - }, - ], - }, - }; - } else if (config.flowPattern == "branchallparrallel") { - payload = { - path: "branchall", - args: {}, - value: { - modules: [ - { - id: "a", - value: { - input_transforms: {}, - language: api.RawScript.language.DENO, - type: "rawscript", - content: - 'export function main(){ return Deno.env.get("WM_FLOW_JOB_ID"); }', - }, - }, - { - id: "b", - value: { - type: "branchall", - parallel: true, - branches: [ - { - modules: [ - { - id: "c", - value: { - input_transforms: { - x: { - type: "javascript", - expr: "results.a", - }, - }, - language: api.RawScript.language.DENO, - type: "rawscript", - content: - "export function main(x: string){ return x; }", - }, - }, - ], - }, - { - modules: [ - { - id: "d", - value: { - input_transforms: { - x: { - type: "javascript", - expr: "results.a", - }, - }, - language: api.RawScript.language.DENO, - type: "rawscript", - content: - "export function main(x: string){ return x; }", - }, - }, - ], - }, - ], - }, - }, - ], - }, - }; - } else { - payload = { - path: "2steps", - args: {}, - value: { - modules: [ - { - id: "a", - value: { - input_transforms: {}, - language: api.RawScript.language.DENO, - type: "rawscript", - content: - 'export function main(){ return Deno.env.get("WM_FLOW_JOB_ID"); }', - }, - }, - { - id: "b", - value: { - type: "identity", - }, - }, - ], - }, - }; - } + const payload = getFlowPayload(config.flowPattern); + uuid = await windmill.JobService.runFlowPreview({ workspace: config.workspace_id, requestBody: payload, @@ -261,99 +139,97 @@ while (cont) { clearInterval(updateStatusInterval); -const end_time = Date.now() + complete_timeout; +// const end_time = Date.now() + complete_timeout; -let incorrect_results = 0; -const enc = (s: string) => new TextEncoder().encode(s); +// let incorrect_results = 0; +// const enc = (s: string) => new TextEncoder().encode(s); -async function getQueueCount() { - return ( - await ( - await fetch( - config.server + "/api/w/" + config.workspace_id + "/jobs/queue/count", - { headers: { ["Authorization"]: "Bearer " + config.token } } - ) - ).json() - ).database_length; -} +// let last_queue_length = await getQueueCount(); +// console.log(`waiting for ${last_queue_length} jobs to complete...`); -let last_queue_length = await getQueueCount(); -console.log(`waiting for ${last_queue_length} jobs to complete...`); +// while ( +// outstanding.length > 0 && +// last_queue_length > 0 && +// Date.now() < end_time +// ) { +// try { +// if (!config.hideProgress) { +// await Deno.stdout.write( +// enc( +// "\rwaiting for jobs to complete: outstanding " + +// outstanding.length + +// " - queue" + +// last_queue_length + +// "\n" +// ) +// ); +// } +// last_queue_length = await getQueueCount(); -while ( - outstanding.length > 0 && - last_queue_length > 0 && - Date.now() < end_time -) { - try { - if (!config.hideProgress) { - await Deno.stdout.write( - enc( - "\rwaiting for jobs to complete: outstanding " + - outstanding.length + - " - queue" + - last_queue_length + - "\n" - ) - ); - } - last_queue_length = await getQueueCount(); +// const uuid = outstanding.shift()!; - const uuid = outstanding.shift()!; +// let r: Job; +// try { +// r = await windmill.JobService.getJob({ +// workspace: config.workspace_id, +// id: uuid, +// }); +// } catch (e) { +// console.log("job not found: " + uuid + " " + e.message); +// continue; +// } +// if (r.type == "QueuedJob") { +// outstanding.push(uuid); - let r: Job; - try { - r = await windmill.JobService.getJob({ - workspace: config.workspace_id, - id: uuid, - }); - } catch (e) { - console.log("job not found: " + uuid + " " + e.message); - continue; - } - if (r.type == "QueuedJob") { - outstanding.push(uuid); +// if (!config.hideProgress) { +// await Deno.stdout.write( +// enc(`uuid: ${uuid}, queue length: ${last_queue_length}\r`) +// ); +// } +// } else { +// r = r as api.CompletedJob; +// try { +// if ( +// ![ +// "httpversion", +// "identity", +// "httpslow", +// "noop", +// "dedicated", +// ].includes(config.scriptPattern) && +// r.result != uuid +// ) { +// console.log( +// "job did not return correct UUID: " + +// r.result + +// " != " + +// uuid + +// "job: \n" + +// JSON.stringify(r, null, 2) +// ); +// incorrect_results++; +// } else { +// // console.log(r.result); +// } +// } catch (e) { +// console.log("error during wait: ", e); +// outstanding.push(uuid); +// } +// } +// } catch (e) { +// console.log("error while waiting for outstanding jobs, sleeing: ", e); +// await sleep(0.5); +// } +// } - if (!config.hideProgress) { - await Deno.stdout.write( - enc(`uuid: ${uuid}, queue length: ${last_queue_length}\r`) - ); - } - } else { - r = r as api.CompletedJob; - try { - if ( - !["httpversion", "identity", "httpslow", "noop"].includes( - config.scriptPattern - ) && - r.result != uuid - ) { - console.log( - "job did not return correct UUID: " + - r.result + - " != " + - uuid + - "job: \n" + - JSON.stringify(r, null, 2) - ); - incorrect_results++; - } else { - // console.log(r.result); - } - } catch (e) { - console.log("error during wait: ", e); - outstanding.push(uuid); - } - } - } catch (e) { - console.log("error while waiting for outstanding jobs, sleeing: ", e); - await sleep(0.5); - } -} +// self.postMessage({ +// type: "zombie_jobs", +// zombie_jobs: outstanding.length, +// incorrect_results, +// jobs_sent: total_spawned, +// }); self.postMessage({ - type: "zombie_jobs", - zombie_jobs: outstanding.length, - incorrect_results, + type: "done", jobs_sent: total_spawned, });