/// /// import { sleep } from "https://deno.land/x/sleep@v1.2.1/sleep.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; per_worker_throughput: number; useFlows: boolean; flowPattern: string; scriptPattern: string; continous: boolean; max_per_worker: number; custom: Action | undefined; server: string; token: string; hideProgress: boolean; }>((resolve, _reject) => { self.onmessage = (evt) => { const sharedConfig = evt.data; windmill.setClient(sharedConfig.token, sharedConfig.server); const config = { workspace_id: sharedConfig.workspace_id, per_worker_throughput: sharedConfig.per_worker_throughput, useFlows: sharedConfig.useFlows, flowPattern: sharedConfig.flowPattern, scriptPattern: sharedConfig.scriptPattern, continous: sharedConfig.continous, max_per_worker: sharedConfig.max_per_worker, custom: sharedConfig.custom, server: sharedConfig.server, token: sharedConfig.token, hideProgress: sharedConfig.hideProgress, }; self.name = "Worker " + sharedConfig.i; resolve(config); self.onmessage = null; }; }); const config = await promise; const outstanding: string[] = []; let cont = true; let total_spawned = 0; const start_time: number = Date.now(); // let complete_timeout = Infinity; self.onmessage = (evt) => { cont = false; // complete_timeout = evt.data; }; const updateStatusInterval = setInterval(() => { self.postMessage({ type: "jobs_sent", jobs_sent: total_spawned }); }, 100); while (cont) { try { const queue_length = await getQueueCount(); if (queue_length > 2500) { console.log( `queue length: ${queue_length} > 2500. waiting... ` ); await sleep(0.5); continue; } if ( (total_spawned * 1000) / (Date.now() - start_time) > config.per_worker_throughput ) { console.log("at maximum throughput. waiting..."); await sleep(0.1); continue; } total_spawned++; if (total_spawned > config.max_per_worker) { break; } let uuid: string; if (config.custom) { await evaluate(config.custom); continue; } else if (config.useFlows) { const payload = getFlowPayload(config.flowPattern); uuid = await windmill.JobService.runFlowPreview({ workspace: config.workspace_id, requestBody: payload, }); } else { try { if (config.scriptPattern === "identity") { uuid = await windmill.JobService.runScriptPreview({ workspace: config.workspace_id, requestBody: { path: "identity", kind: api.Preview.kind.IDENTITY, args: { identity: "itsme", }, }, }); } else { uuid = await windmill.JobService.runScriptByPath({ workspace: config.workspace_id, path: "f/benchmarks/" + (config.scriptPattern || "deno"), requestBody: {}, }); } } catch (e) { console.error("error running script: " + e.body); Deno.exit(1); } } if (!config.continous) outstanding.push(uuid); } catch (e) { console.log( `error while sending job: ${e} ` ); await sleep(0.5); continue; } } clearInterval(updateStatusInterval); self.postMessage({ type: "done", jobs_sent: total_spawned, });