445 lines
12 KiB
TypeScript
445 lines
12 KiB
TypeScript
/// <reference no-default-lib="true" />
|
|
/// <reference lib="deno.window" />
|
|
|
|
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";
|
|
|
|
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";
|
|
|
|
import { VERSION, createBenchScript, createWacBenchScript, getFlowPayload, login, WAC_KINDS, STEPS_PER_WORKFLOW } from "./lib.ts";
|
|
|
|
async function verifyOutputs(uuids: string[], workspace: string) {
|
|
console.log("Verifying outputs");
|
|
let incorrectResults = 0;
|
|
for (const uuid of uuids) {
|
|
try {
|
|
const job = await windmill.JobService.getCompletedJob({
|
|
workspace,
|
|
id: uuid,
|
|
});
|
|
if (!job.success) {
|
|
console.log(`Job ${uuid} did not complete`);
|
|
incorrectResults++;
|
|
}
|
|
if (job.result !== uuid) {
|
|
console.log(`Job ${uuid} did not output the correct value: ${JSON.stringify(job)}`);
|
|
incorrectResults++;
|
|
}
|
|
} catch (_) {
|
|
console.log(`Job ${uuid} did not complete`);
|
|
incorrectResults++;
|
|
}
|
|
}
|
|
console.log(`Incorrect results: ${incorrectResults}`);
|
|
}
|
|
|
|
export const NON_TEST_TAGS = ["deno", "python", "go", "bash", "dedicated", "bun", "nativets", "dedicated_nativets", "flow"]
|
|
|
|
const FLOW_COMPARISON_KINDS = ["flow_seq_2_bun", "flow_par_2_bun", "flow_seq_3_bun"];
|
|
export async function main({
|
|
host,
|
|
email,
|
|
password,
|
|
token,
|
|
workspace,
|
|
kind,
|
|
jobs,
|
|
noVerify,
|
|
}: {
|
|
host: string;
|
|
email?: string;
|
|
password?: string;
|
|
token?: string;
|
|
workspace: string;
|
|
kind: string;
|
|
jobs: number;
|
|
noVerify?: boolean;
|
|
}) {
|
|
windmill.setClient("", host);
|
|
|
|
console.log(
|
|
"Started benchmark with options",
|
|
JSON.stringify(
|
|
{
|
|
host,
|
|
email,
|
|
workspace,
|
|
kind,
|
|
jobs,
|
|
noVerify,
|
|
},
|
|
null,
|
|
4
|
|
)
|
|
);
|
|
|
|
const config = {
|
|
token: "",
|
|
server: host,
|
|
workspace_id: workspace,
|
|
};
|
|
|
|
let final_token: string;
|
|
if (!token) {
|
|
if (email && password) {
|
|
final_token = await login(email, password);
|
|
} else {
|
|
console.error("Token or email with password are required.");
|
|
return;
|
|
}
|
|
} else {
|
|
final_token = token;
|
|
}
|
|
|
|
config.token = final_token;
|
|
windmill.setClient(final_token, host);
|
|
const enc = (s: string) => new TextEncoder().encode(s);
|
|
|
|
async function getQueueCount(tags?: string[]) {
|
|
return (
|
|
await (
|
|
await fetch(
|
|
config.server + "/api/w/" + config.workspace_id + "/jobs/queue/count" + (tags && tags.length > 0 ? "?tags=" + tags.join(",") : ""),
|
|
{ headers: { ["Authorization"]: "Bearer " + config.token } }
|
|
)
|
|
).json()
|
|
).database_length;
|
|
}
|
|
|
|
async function getFlowStepCount(
|
|
workspace: string,
|
|
path: string
|
|
): Promise<number> {
|
|
const response = await fetch(
|
|
`${config.server}/api/w/${workspace}/flows/get/${path}`,
|
|
{ headers: { ["Authorization"]: "Bearer " + config.token } }
|
|
);
|
|
|
|
const data = await response.json();
|
|
let stepCount = 0;
|
|
|
|
for (const mod of data.value.modules) {
|
|
if (mod.value.type === "flow" && mod.value.path) {
|
|
const subFlowCount = await getFlowStepCount(workspace, mod.value.path);
|
|
stepCount += subFlowCount;
|
|
} else {
|
|
stepCount += 1;
|
|
}
|
|
}
|
|
|
|
return stepCount;
|
|
}
|
|
|
|
let pastJobs = 0;
|
|
async function getCompletedJobsCount(tags?: string[]): Promise<number> {
|
|
const completedJobs = (
|
|
await (
|
|
await fetch(
|
|
host + "/api/w/" + config.workspace_id + "/jobs/completed/count" + (tags && tags.length > 0 ? "?tags=" + tags.join(",") : ""),
|
|
{ headers: { ["Authorization"]: "Bearer " + config.token } }
|
|
)
|
|
).json()
|
|
).database_length;
|
|
return completedJobs - pastJobs;
|
|
}
|
|
|
|
if (
|
|
["deno", "python", "go", "bash", "dedicated", "bun", "nativets", "dedicated_nativets"].includes(
|
|
kind
|
|
)
|
|
) {
|
|
await createBenchScript(kind, workspace);
|
|
} else if (WAC_KINDS.includes(kind)) {
|
|
await createWacBenchScript(kind, workspace);
|
|
}
|
|
|
|
|
|
const jobsSent = jobs;
|
|
console.log(`Bulk creating ${jobsSent} jobs`);
|
|
|
|
const start_create = Date.now();
|
|
let nStepsFlow = 0;
|
|
let body: string;
|
|
if (kind === "noop") {
|
|
body = JSON.stringify({
|
|
kind: "noop",
|
|
});
|
|
} else if (
|
|
["deno", "python", "go", "bash", "dedicated", "bun", "nativets", "dedicated_nativets"].includes(
|
|
kind
|
|
)
|
|
) {
|
|
body = JSON.stringify({
|
|
kind: "script",
|
|
path: "f/benchmarks/" + kind,
|
|
});
|
|
} else if (WAC_KINDS.includes(kind)) {
|
|
// WAC v2 scripts are deployed as bun scripts, run via script path
|
|
nStepsFlow = STEPS_PER_WORKFLOW[kind] ?? 0;
|
|
body = JSON.stringify({
|
|
kind: "script",
|
|
path: "f/benchmarks/" + kind,
|
|
});
|
|
} else if (FLOW_COMPARISON_KINDS.includes(kind)) {
|
|
nStepsFlow = STEPS_PER_WORKFLOW[kind] ?? 0;
|
|
const payload = getFlowPayload(kind);
|
|
body = JSON.stringify({
|
|
kind: "flow",
|
|
flow_value: payload.value,
|
|
});
|
|
} else if (["2steps", "bigscriptinflow"].includes(kind)) {
|
|
nStepsFlow = kind == "2steps" ? 2 : 1;
|
|
const payload = getFlowPayload(kind);
|
|
body = JSON.stringify({
|
|
kind: "flow",
|
|
flow_value: payload.value,
|
|
});
|
|
} else if (kind.startsWith("flow:")) {
|
|
console.log("Detected custom flow ");
|
|
let flow_path = kind.substring(5);
|
|
nStepsFlow = await getFlowStepCount(config.workspace_id, flow_path);
|
|
console.log(`Total steps of flow including sub-flows: ${nStepsFlow}`);
|
|
body = JSON.stringify({
|
|
kind: "flow",
|
|
path: flow_path,
|
|
});
|
|
} else if (kind.startsWith("script:")) {
|
|
console.log("Detected custom script");
|
|
body = JSON.stringify({
|
|
kind: "script",
|
|
path: kind.substring(7),
|
|
});
|
|
} else if (kind == "bigrawscript") {
|
|
noVerify = true;
|
|
body = JSON.stringify({
|
|
kind: "rawscript",
|
|
rawscript: {
|
|
language: api.RawScript.language.BASH,
|
|
content: "# let's bloat that bash script, 3.. 2.. 1.. BOOM\n".repeat(100) + "echo \"$WM_FLOW_JOB_ID\"\n",
|
|
},
|
|
});
|
|
} else {
|
|
throw new Error("Unknown script pattern " + kind);
|
|
}
|
|
|
|
let testOtherTag = false;
|
|
if (testOtherTag) {
|
|
const otherTagTodo = 2000000;
|
|
|
|
let parsed = JSON.parse(body);
|
|
parsed.tag = "test";
|
|
let nbody = JSON.stringify(parsed);
|
|
let response2 = await fetch(
|
|
config.server +
|
|
"/api/w/" +
|
|
config.workspace_id +
|
|
`/jobs/add_batch_jobs/${otherTagTodo}`,
|
|
{
|
|
method: "POST",
|
|
headers: {
|
|
["Authorization"]: "Bearer " + config.token,
|
|
"Content-Type": "application/json",
|
|
},
|
|
body: nbody,
|
|
}
|
|
);
|
|
if (!response2.ok) {
|
|
throw new Error(
|
|
"Failed to create jobs: " +
|
|
response2.statusText +
|
|
" " +
|
|
(await response2.text())
|
|
);
|
|
}
|
|
}
|
|
|
|
pastJobs = await getCompletedJobsCount(NON_TEST_TAGS);
|
|
|
|
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 +
|
|
" " +
|
|
(await response.text())
|
|
);
|
|
}
|
|
const uuids = await response.json();
|
|
const end_create = Date.now();
|
|
const create_duration = end_create - start_create;
|
|
console.log(
|
|
`Jobs successfully added to the queue in ${create_duration / 1000
|
|
}s. Windmill will start pulling them\n`
|
|
);
|
|
let start = Date.now();
|
|
|
|
let completedJobs = 0;
|
|
let lastElapsed = 0;
|
|
let lastCompletedJobs = 0;
|
|
|
|
// Timeout: 10 minutes for the polling loop to prevent hanging forever
|
|
// (e.g. if WAC suspend/resume fails or jobs get stuck)
|
|
const POLL_TIMEOUT_MS = 10 * 60 * 1000;
|
|
let didStart = false;
|
|
while (completedJobs < jobsSent) {
|
|
const loopStart = Date.now();
|
|
if (!didStart) {
|
|
const actual_queue = await getQueueCount(NON_TEST_TAGS);
|
|
if (actual_queue < jobsSent) {
|
|
start = Date.now();
|
|
didStart = true;
|
|
}
|
|
} else {
|
|
const elapsed = start ? Date.now() - start : 0;
|
|
if (elapsed > POLL_TIMEOUT_MS) {
|
|
console.error(`\nTimeout: benchmark did not complete within ${POLL_TIMEOUT_MS / 1000}s (${completedJobs}/${jobsSent} completed)`);
|
|
break;
|
|
}
|
|
completedJobs = await getCompletedJobsCount(NON_TEST_TAGS);
|
|
if (nStepsFlow > 0) {
|
|
completedJobs = Math.floor(completedJobs / (nStepsFlow + 1));
|
|
}
|
|
const avgThr = ((completedJobs / elapsed) * 1000).toFixed(2);
|
|
const instThr =
|
|
lastElapsed > 0
|
|
? (
|
|
((completedJobs - lastCompletedJobs) / (elapsed - lastElapsed)) *
|
|
1000
|
|
).toFixed(2)
|
|
: 0;
|
|
|
|
lastElapsed = elapsed;
|
|
lastCompletedJobs = completedJobs;
|
|
|
|
await Deno.stdout.write(
|
|
enc(
|
|
`elapsed: ${(elapsed / 1000).toFixed(
|
|
2
|
|
)} | jobs executed: ${completedJobs}/${jobsSent} (thr: inst ${instThr} - avg ${avgThr}) | remaining: ${jobsSent - completedJobs
|
|
} \r`
|
|
)
|
|
);
|
|
}
|
|
const loopDuration = (Date.now() - loopStart) / 1000.0;
|
|
if (loopDuration < 0.05) {
|
|
await sleep(0.05 - loopDuration);
|
|
}
|
|
}
|
|
|
|
const total_duration_sec = (Date.now() - start) / 1000.0;
|
|
|
|
console.log(`\njobs: ${jobsSent}`);
|
|
console.log(`duration: ${total_duration_sec}s`);
|
|
console.log(`avg. throughput (jobs/time): ${jobsSent / total_duration_sec}`);
|
|
|
|
console.log("completed jobs", completedJobs);
|
|
console.log("queue length:", await getQueueCount(NON_TEST_TAGS));
|
|
|
|
if (
|
|
!noVerify &&
|
|
kind !== "noop" &&
|
|
kind !== "nativets" &&
|
|
kind !== "dedicated_nativets" &&
|
|
!kind.startsWith("flow:") &&
|
|
!kind.startsWith("script:") &&
|
|
!WAC_KINDS.includes(kind) &&
|
|
!FLOW_COMPARISON_KINDS.includes(kind)
|
|
) {
|
|
await verifyOutputs(uuids, config.workspace_id);
|
|
}
|
|
|
|
console.log("done");
|
|
|
|
return {
|
|
throughput: jobsSent / total_duration_sec,
|
|
};
|
|
}
|
|
|
|
if (import.meta.main) {
|
|
await new Command()
|
|
.name("wmillbench")
|
|
.description("Run Benchmark to measure throughput of windmill.")
|
|
.version(VERSION)
|
|
.option("--host <url:string>", "The windmill host to benchmark.", {
|
|
default: "http://127.0.0.1:8000",
|
|
})
|
|
.option("-e --email <email:string>", "The email to use to login.", {
|
|
default: "admin@windmill.dev",
|
|
})
|
|
.option(
|
|
"-p --password <password:string>",
|
|
"The password to use to login.",
|
|
{
|
|
default: "changeme",
|
|
}
|
|
)
|
|
.env(
|
|
"WM_TOKEN=<token:string>",
|
|
"The token to use when talking to the API server. Preferred over manual login."
|
|
)
|
|
.option(
|
|
"-t --token <token:string>",
|
|
"The token to use when talking to the API server. Preferred over manual login."
|
|
)
|
|
.env(
|
|
"WM_WORKSPACE=<workspace:string>",
|
|
"The workspace to spawn scripts from."
|
|
)
|
|
.option(
|
|
"-w --workspace <workspace:string>",
|
|
"The workspace to spawn scripts from.",
|
|
{ default: "admins" }
|
|
)
|
|
.option(
|
|
"--kind <kind:string>",
|
|
"Specifiy the benchmark kind among: deno, identity, python, go, bash, dedicated, bun, noop, 2steps, nativets, dedicated_nativets, wac_seq_2, wac_par_2, wac_seq_3, wac_inline_2, flow_seq_2_bun, flow_par_2_bun, flow_seq_3_bun",
|
|
{
|
|
required: true,
|
|
}
|
|
)
|
|
.option("-j --jobs <jobs:number>", "Number of jobs to create.", {
|
|
default: 10000,
|
|
})
|
|
.option("--no-verify", "Do not verify the output of the jobs.", {
|
|
default: false,
|
|
})
|
|
.action(main)
|
|
.command(
|
|
"upgrade",
|
|
new UpgradeCommand({
|
|
main: "main.ts",
|
|
args: [
|
|
"--allow-net",
|
|
"--allow-read",
|
|
"--allow-write",
|
|
"--allow-env",
|
|
"--unstable",
|
|
],
|
|
provider: new DenoLandProvider({ name: "wmillbench" }),
|
|
})
|
|
)
|
|
.parse();
|
|
}
|