Files
windmill/typescript-client/client.ts
Ruben Fiszel a6d4390790 feat: workflow-as-code (WAC) v2 (#8172)
* feat: workflow-as-code v2 with @task decorator API

Replace ctx.step("name", "script") API with @task decorators where
functions are called directly. Users no longer need to pass WorkflowCtx
or use string-based step names/script paths.

Python: @task decorator with contextvars-based implicit context
TypeScript: task() wrapper with module-level context variable
Parsers: detect @task function calls instead of ctx.step() calls
Worker: updated wrappers to set implicit context

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* feat: WAC v2 checkpoint/replay with _executing_key child dispatch

- Rust-side orchestration: parent dispatches child jobs, suspends, resumes on completion
- _executing_key in checkpoint tells child which step to execute directly
- task() throws StepSuspend(mode="step_complete") after executing target step
- result_processor handles child completion and updates parent checkpoint
- WacGraph.svelte for runtime execution visualization
- Sequential and parallel workflows tested end-to-end

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: WAC v2 bundle cache, globalThis ctx sharing, description optional

- Disable bun bundle caching for WAC v2 scripts (wrapper needs
  windmill-client from node_modules, not available in bundle mode)
- Use Reflect.set/get(globalThis, "__wmill_wf_ctx") to share workflow
  context across dual module instances (wrapper vs user script)
- Never-resolving thenable for non-matching steps in child job mode
  prevents Promise.all race conditions
- Make description field optional in NewScript API (defaults to "")

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* feat: add step() primitive for inline checkpointed steps

step() executes a function inline (no child job) and persists the result
to the checkpoint. On replay, the cached value is returned — ensuring
deterministic behavior for non-deterministic operations like Date.now()
or Math.random().

- TypeScript: step(name, fn) — executes inline, throws StepSuspend with
  mode "inline_checkpoint" to persist before continuing
- Rust: InlineCheckpoint variant in WacOutput, saves to checkpoint and
  resets running=false for immediate re-pickup (no zombie wait)
- Shared step counter between task() and step() via _allocKey()

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* feat: add Python WAC v2 support with task(), step(), workflow()

- Python SDK: WorkflowCtx with _executing_key child mode, _alloc_key
  shared counter, _run_inline_step for step(), _execute_directly and
  _never_resolve for child mode, step() async function
- Python executor: WAC v2 detection, checkpoint.json writing, WAC
  wrapper.py generation calling _run_workflow(), post-execution hook
  into shared handle_wac_v2_output()
- Make handle_wac_v2_output pub so both bun and python executors share
  the same dispatch/suspend/inline-checkpoint logic
- 17 Python tests covering dispatch, replay, parallel, conditional,
  inline checkpoint, and child mode

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* chore: update sqlx prepared queries

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: WacGraph Tooltip→Popover, simplify wacToFlow parsers

- Fix type error: Tooltip doesn't accept text snippet, use Popover
- Extract shared helpers for task matching and block collection
- Replace linear tasks.find() with Map lookups
- Remove mutable module-level counter

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: Box::pin WAC v2 output handler to prevent stack overflow

handle_python_job's async state machine was too large when combined
with handle_wac_v2_output. Box::pin heap-allocates the future.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: merge WAC v1 and v2 task decorators to preserve backward compat

The v2 @task decorator was shadowing the v1 one, breaking WAC v1
scripts that rely on HTTP-based dispatch via /workflow_as_code/ API.

The merged decorator handles three modes:
- v2: inside @workflow context → checkpoint/replay dispatch
- v1: WM_JOB_ID set, no @workflow → HTTP API dispatch + wait_job
- standalone: no Windmill env → execute function body directly

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: skip no_main_func detection for WAC v2 scripts in TS and Python parsers

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: prevent empty/noop dispatch causing infinite requeue loop

- Validate steps.len() > 0 in WAC dispatch handler (issue 3)
- Replace noop StepSuspend throw with never-resolving promise so it
  can't reach the backend as an empty dispatch (issue 4)

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: Python task wrapper now converts positional args to kwargs in v2 mode

Previously only **kwargs were passed to _next_step(), silently dropping
positional arguments. Extract shared _merge_args() helper used by both
v1 and v2 paths.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: replace unwrap() with proper error propagation in WAC arg serialization

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: add workspace_id filter to v2_job queries in WAC dispatch

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: prevent race condition in WAC child dispatch

Restructure dispatch to save checkpoint + suspend parent + seed child
checkpoints in a single transaction BEFORE pushing child jobs. This
ensures a fast child can't complete before the parent is suspended.

Also wrap InlineCheckpoint save + running reset in a transaction to
prevent corrupted state on crash.

Use ULID for pre-generated child job IDs (consistent with rest of API).

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: include step key and child job ID in WAC error propagation

Move step_key lookup before the success check so failed child errors
include which task failed, the child job ID, and the original error.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* docs: document WAC determinism contract and step dispatch semantics

- Document that workflow functions must be deterministic across replays
- Document that WacStepDispatch.script/args are metadata, not dispatch targets
- Add comments on counter-based key allocation

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: tighten WAC v2 detection to reduce false positives

Replace naive substring matching with line-aware checks that skip
comments and look for specific patterns:
- TS: import from "windmill-client" containing workflow/task
- Python: @workflow and @task decorators with wmill import

Extracted shared helpers in wac_executor.rs used by both executors.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: show failed steps in WacGraph when workflow completes with errors

When flowDone is true and a pending step isn't in completedSteps,
mark it as 'failed' instead of 'running'. The failed state CSS and
XCircle icon were already defined but never triggered.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: unsuspend and fail parent when WAC child push fails

Previously if a child push failed mid-batch, the parent remained
suspended with suspend = num_steps but fewer children, hanging until
the 14-day timeout. Now the push loop catches errors and unsuspends
the parent before returning the error.

Also adds source hash validation: if the script content changes between
replays, the job fails with a clear error instead of silently feeding
stale checkpoint data into wrong steps.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: clear suspend_until when unsuspending WAC parent

Set suspend_until = NULL alongside suspend = 0 in both the child
failure and all-children-complete paths, so the parent doesn't rely
on subtle pull query invariants to be re-picked-up.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test: add exhaustive edge case tests for WAC v2 SDK

fix: make TS task wrapper non-async to fix unawaited task flush

The async wrapper caused microtask-based thenable auto-resolution that
fired .then() and threw StepSuspend before _flushPending() could capture
unawaited steps — making the flush mechanism completely broken. Now the
thenable is returned directly without async wrapping. Backward compatible
with v1 (all code paths still return awaitables).

Tests added (59 TS + 66 Python) covering: full sequential lifecycle,
step after parallel, parallel after parallel, conditional on step result,
empty/single-task workflows, 10+ steps, falsy value preservation, inline
steps, mixed step/task, unawaited flush, child mode with parallel,
key determinism, large parallel groups, and complex mixed patterns.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: atomic checkpoint updates to prevent parallel child race condition

Replace read-modify-write pattern in handle_wac_child_completion with
atomic SQL operations:
- completed_steps merged via jsonb_set(... || jsonb_build_object(...))
  so concurrent children on different workers don't overwrite each other
- suspend counter decremented atomically with RETURNING to determine
  "all done" condition (instead of checking completed_steps in memory)
- suspend_until cleared in the same atomic decrement statement

Before this fix, two parallel children completing simultaneously could
both load the same checkpoint, each add their step, and save — the
second write would overwrite the first, silently losing a child result
and leaving the parent suspended forever.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: cancel already-pushed children on partial WAC dispatch failure

When pushing child jobs sequentially, if pushing child N fails, children
1..N-1 are already running. Previously the error handler only unsuspended
the parent, leaving orphaned children that would complete and corrupt the
checkpoint state (decrementing suspend on an already-unsuspended parent,
potentially causing duplicate step execution on re-run).

Now on partial failure:
1. Cancel all already-pushed children (prevents them from completing
   and corrupting checkpoint state)
2. Clear pending_steps from checkpoint (so parent doesn't think
   children are outstanding on re-run)
3. Then unsuspend parent (so the error propagates)

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: skip WAC duration write and child check for non-WAC parents

The duration write to workflow_as_code_status was running for every
non-flow child with a parent (error handlers, success handlers,
run_script children), even though it was only intended for WAC jobs.

Add WHERE workflow_as_code_status IS NOT NULL to skip non-WAC parents
entirely. Piggyback RETURNING pending_steps.job_ids on the same query
so WAC v2 child completion needs zero extra DB round-trips on the
success path.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: seed child checkpoint in same transaction as push

The child checkpoint insert was happening before the child job was
pushed, violating the FK constraint on v2_job_status. Move it into
the push transaction so the job row exists and the child can't be
picked up before its checkpoint is ready.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: set running=false when WAC parent suspends for child dispatch

The parent job kept running=true after suspending, so workers wouldn't
pick it up when children completed and suspend reached 0. The parent
only advanced when the zombie job detector reset it (~90s). Now the
dispatch suspend sets running=false so the parent is immediately
eligible for pickup.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: WAC parent suspend/unsuspend lifecycle

Keep running=true when suspending the parent so the normal pull query
(WHERE running=false) never picks it up. Keep suspend_until non-null
when decrementing suspend to 0 so the suspended pull query
(WHERE suspend_until IS NOT NULL AND suspend<=0) picks it up.

Previously: setting running=false caused infinite restart loops because
the normal pull query has no suspend check and would immediately re-pick
the parent. Clearing suspend_until on the last child prevented the
suspended pull from ever seeing it, requiring the 90s zombie detector.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* feat: add approval primitive, flow child completion, timeline fixes for WAC v2

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* feat: add error propagation, task options, sleep, and parallel for WAC v2

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test: fix python SDK tests to use name-based keys and add new test coverage

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: address WAC v2 review findings (sleep timing, error marker, atomicity)

- Fix sleep using suspend=1 instead of 0 to enforce actual delay
- Add approval/sleep resume injection to Python executor
- Fix TS SDK concurrency_limit mapping (was reading wrong property)
- Namespace error marker as __wmill_error to avoid user data collision
- Wrap child completion SQL in transaction for atomicity
- Decrement suspend even when step key is missing (prevents hang)
- Expand TASK_RE to handle export const, let, var, generics
- Validate step key uniqueness before dispatch
- Log warning on checkpoint deserialization failure
- Remove unimplemented delete_after_use from SDKs
- Add TaskError exception class to Python SDK with diagnostic context
- Fix extra positional args handling and add functools.wraps
- Improve getParamNames to handle typed/destructured params

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* sqlx

* sqlx

* test: add WAC v1 e2e integration tests for TS and Python

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix: revert fake test versions in typescript-client

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* refactor: remove unused WacGraph component and strip wacToFlow to isWorkflowAsCode

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* refactor: extract shared approval/sleep resume logic into wac_executor

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-09 19:39:24 +00:00

1869 lines
57 KiB
TypeScript

// Import only the services actually used in this file (not re-exported)
// This enables tree-shaking - importing setClient won't pull in all services
import {
ResourceService,
VariableService,
JobService,
HelpersService,
AppService,
MetricsService,
OidcService,
UserService,
} from "./services.gen";
import { OpenAPI } from "./core/OpenAPI";
// import type { DenoS3LightClientSettings } from "./index";
import {
DenoS3LightClientSettings,
S3ObjectRecord,
type S3Object,
} from "./s3Types";
export {
type S3Object,
type S3ObjectRecord,
type S3ObjectURI,
} from "./s3Types";
export {
datatable,
ducklake,
type SqlTemplateFunction,
type DatatableSqlTemplateFunction,
} from "./sqlUtils";
// Services are NOT re-exported here to enable tree-shaking
// Import services directly from "windmill-client" or use the default export
export type Sql = string;
export type Email = string;
export type Base64 = string;
export type Resource<S extends string> = any;
export const SHARED_FOLDER = "/shared";
let mockedApi: MockedApi | undefined = undefined;
/**
* Initialize the Windmill client with authentication token and base URL
* @param token - Authentication token (defaults to WM_TOKEN env variable)
* @param baseUrl - API base URL (defaults to BASE_INTERNAL_URL or BASE_URL env variable)
*/
export function setClient(token?: string, baseUrl?: string) {
if (baseUrl === undefined) {
baseUrl =
getEnv("BASE_INTERNAL_URL") ??
getEnv("BASE_URL") ??
"http://localhost:8000";
}
if (token === undefined) {
token = getEnv("WM_TOKEN") ?? "no_token";
}
OpenAPI.WITH_CREDENTIALS = true;
OpenAPI.TOKEN = token;
OpenAPI.BASE = baseUrl + "/api";
}
function getPublicBaseUrl(): string {
return getEnv("WM_BASE_URL") ?? "http://localhost:3000";
}
const getEnv = (key: string) => {
if (typeof window === "undefined") {
// node
return process?.env?.[key];
}
// browser
return window?.process?.env?.[key];
};
/**
* Create a client configuration from env variables
* @returns client configuration
*/
export function getWorkspace(): string {
return getEnv("WM_WORKSPACE") ?? "no_workspace";
}
/**
* Get a resource value by path
* @param path path of the resource, default to internal state path
* @param undefinedIfEmpty if the resource does not exist, return undefined instead of throwing an error
* @returns resource value
*/
export async function getResource(
path?: string,
undefinedIfEmpty?: boolean
): Promise<any> {
path = parseResourceSyntax(path) ?? path ?? getStatePath();
const mockedApi = await getMockedApi();
if (mockedApi) {
if (mockedApi.resources[path]) {
return mockedApi.resources[path];
} else {
console.log(
`MockedAPI present, but resource not found at ${path}, falling back to real API`
);
}
}
const workspace = getWorkspace();
try {
return await ResourceService.getResourceValueInterpolated({
workspace,
path,
});
} catch (e: any) {
if (undefinedIfEmpty && e.status === 404) {
return undefined;
} else {
throw Error(
`Resource not found at ${path} or not visible to you: ${e.body}`
);
}
}
}
/**
* Get the true root job id
* @param jobId job id to get the root job id from (default to current job)
* @returns root job id
*/
export async function getRootJobId(jobId?: string): Promise<string> {
const workspace = getWorkspace();
jobId = jobId ?? getEnv("WM_JOB_ID");
if (jobId === undefined) {
throw Error("Job ID not set");
}
return await JobService.getRootJobId({ workspace, id: jobId });
}
/**
* @deprecated Use runScriptByPath or runScriptByHash instead
*/
export async function runScript(
path: string | null = null,
hash_: string | null = null,
args: Record<string, any> | null = null,
verbose: boolean = false
): Promise<any> {
console.warn(
"runScript is deprecated. Use runScriptByPath or runScriptByHash instead."
);
if (path && hash_) {
throw new Error("path and hash_ are mutually exclusive");
}
return _runScriptInternal(path, hash_, args, verbose);
}
async function _runScriptInternal(
path: string | null = null,
hash_: string | null = null,
args: Record<string, any> | null = null,
verbose: boolean = false
): Promise<any> {
args = args || {};
if (verbose) {
if (path) {
console.info(`running \`${path}\` synchronously with args:`, args);
} else if (hash_) {
console.info(
`running script with hash \`${hash_}\` synchronously with args:`,
args
);
}
}
const jobId = await _runScriptAsyncInternal(path, hash_, args);
return await waitJob(jobId, verbose);
}
/**
* Run a script synchronously by its path and wait for the result
* @param path - Script path in Windmill
* @param args - Arguments to pass to the script
* @param verbose - Enable verbose logging
* @returns Script execution result
*/
export async function runScriptByPath(
path: string,
args: Record<string, any> | null = null,
verbose: boolean = false
): Promise<any> {
return _runScriptInternal(path, null, args, verbose);
}
/**
* Run a script synchronously by its hash and wait for the result
* @param hash_ - Script hash in Windmill
* @param args - Arguments to pass to the script
* @param verbose - Enable verbose logging
* @returns Script execution result
*/
export async function runScriptByHash(
hash_: string,
args: Record<string, any> | null = null,
verbose: boolean = false
): Promise<any> {
return _runScriptInternal(null, hash_, args, verbose);
}
/**
* Append a text to the result stream
* @param text text to append to the result stream
*/
export function appendToResultStream(text: string) {
console.log("WM_STREAM: " + text.replace(/\n/g, "\\n"));
}
/**
* Stream to the result stream
* @param stream stream to stream to the result stream
*/
export async function streamResult(stream: AsyncIterable<string>) {
for await (const text of stream) {
appendToResultStream(text);
}
}
/**
* Run a flow synchronously by its path and wait for the result
* @param path - Flow path in Windmill
* @param args - Arguments to pass to the flow
* @param verbose - Enable verbose logging
* @returns Flow execution result
*/
export async function runFlow(
path: string | null = null,
args: Record<string, any> | null = null,
verbose: boolean = false
): Promise<any> {
args = args || {};
if (verbose) {
console.info(`running \`${path}\` synchronously with args:`, args);
}
const jobId = await runFlowAsync(path, args, null, false);
return await waitJob(jobId, verbose);
}
/**
* Wait for a job to complete and return its result
* @param jobId - ID of the job to wait for
* @param verbose - Enable verbose logging
* @returns Job result when completed
*/
export async function waitJob(
jobId: string,
verbose: boolean = false
): Promise<any> {
while (true) {
// Implement your HTTP request logic here to get job result
const resultRes = await getResultMaybe(jobId);
const started = resultRes.started;
const completed = resultRes.completed;
const success = resultRes.success;
if (!started && verbose) {
console.info(`job ${jobId} has not started yet`);
}
if (completed) {
const result = resultRes.result;
if (success) {
return result;
} else {
const error = result.error;
throw new Error(
`Job ${jobId} was not successful: ${JSON.stringify(error)}`
);
}
}
if (verbose) {
console.info(`sleeping 0.5 seconds for jobId: ${jobId}`);
}
await new Promise((resolve) => setTimeout(resolve, 500));
}
}
/**
* Get the result of a completed job
* @param jobId - ID of the completed job
* @returns Job result
*/
export async function getResult(jobId: string): Promise<any> {
const workspace = getWorkspace();
return await JobService.getCompletedJobResult({ workspace, id: jobId });
}
/**
* Get the result of a job if completed, or its current status
* @param jobId - ID of the job
* @returns Object with started, completed, success, and result properties
*/
export async function getResultMaybe(jobId: string): Promise<any> {
const workspace = getWorkspace();
return await JobService.getCompletedJobResultMaybe({ workspace, id: jobId });
}
const STRIP_COMMENTS =
/(\/\/.*$)|(\/\*[\s\S]*?\*\/)|(\s*=[^,\)]*(('(?:\\'|[^'\r\n])*')|("(?:\\"|[^"\r\n])*"))|(\s*=[^,\)]*))/gm;
function getParamNames(func: Function): string[] {
const fnStr = func.toString().replace(STRIP_COMMENTS, "");
// Find the matching closing paren for the parameter list, handling nesting
const openIdx = fnStr.indexOf("(");
if (openIdx === -1) return [];
let depth = 1;
let closeIdx = openIdx + 1;
for (; closeIdx < fnStr.length && depth > 0; closeIdx++) {
if (fnStr[closeIdx] === "(") depth++;
else if (fnStr[closeIdx] === ")") depth--;
}
const paramStr = fnStr.slice(openIdx + 1, closeIdx - 1).trim();
if (!paramStr) return [];
// Split on commas at depth 0 (skip nested parens, angle brackets, braces)
const params: string[] = [];
let current = "";
let d = 0;
for (const ch of paramStr) {
if ("(<{".includes(ch)) d++;
else if (")>}".includes(ch)) d--;
if (ch === "," && d === 0) {
params.push(current.trim());
current = "";
} else {
current += ch;
}
}
if (current.trim()) params.push(current.trim());
// Extract the parameter name from each param (strip type annotations, destructuring, rest)
return params.map((p) => {
// Remove rest operator
p = p.replace(/^\.\.\./, "");
// For destructured params like { url, depth }: Config, use a positional fallback
if (p.startsWith("{") || p.startsWith("[")) return "";
// Strip type annotation (e.g. "x: number" -> "x", "x?: string" -> "x")
const colonIdx = p.indexOf(":");
if (colonIdx !== -1) p = p.slice(0, colonIdx);
return p.replace(/\?$/, "").trim();
}).filter(Boolean);
}
/**
* @deprecated Use runScriptByPathAsync or runScriptByHashAsync instead
*/
export async function runScriptAsync(
path: string | null,
hash_: string | null,
args: Record<string, any> | null,
scheduledInSeconds: number | null = null
): Promise<string> {
console.warn(
"runScriptAsync is deprecated. Use runScriptByPathAsync or runScriptByHashAsync instead."
);
// Create a script job and return its job id.
if (path && hash_) {
throw new Error("path and hash_ are mutually exclusive");
}
return _runScriptAsyncInternal(path, hash_, args, scheduledInSeconds);
}
async function _runScriptAsyncInternal(
path: string | null = null,
hash_: string | null = null,
args: Record<string, any> | null = null,
scheduledInSeconds: number | null = null
): Promise<string> {
// Create a script job and return its job id.
args = args || {};
const params: Record<string, any> = {};
if (scheduledInSeconds) {
params["scheduled_in_secs"] = scheduledInSeconds;
}
let parentJobId = getEnv("WM_JOB_ID");
if (parentJobId !== undefined) {
params["parent_job"] = parentJobId;
}
let rootJobId = getEnv("WM_ROOT_FLOW_JOB_ID");
if (rootJobId != undefined && rootJobId != "") {
params["root_job"] = rootJobId;
}
let endpoint: string;
if (path) {
endpoint = `/w/${getWorkspace()}/jobs/run/p/${path}`;
} else if (hash_) {
endpoint = `/w/${getWorkspace()}/jobs/run/h/${hash_}`;
} else {
throw new Error("path or hash_ must be provided");
}
let url = new URL(OpenAPI.BASE + endpoint);
url.search = new URLSearchParams(params).toString();
return fetch(url, {
method: "POST",
headers: {
"Content-Type": "application/json",
Authorization: `Bearer ${OpenAPI.TOKEN}`,
},
body: JSON.stringify(args),
}).then((res) => res.text());
}
/**
* Run a script asynchronously by its path
* @param path - Script path in Windmill
* @param args - Arguments to pass to the script
* @param scheduledInSeconds - Schedule execution for a future time (in seconds)
* @returns Job ID of the created job
*/
export async function runScriptByPathAsync(
path: string,
args: Record<string, any> | null = null,
scheduledInSeconds: number | null = null
): Promise<string> {
return _runScriptAsyncInternal(path, null, args, scheduledInSeconds);
}
/**
* Run a script asynchronously by its hash
* @param hash_ - Script hash in Windmill
* @param args - Arguments to pass to the script
* @param scheduledInSeconds - Schedule execution for a future time (in seconds)
* @returns Job ID of the created job
*/
export async function runScriptByHashAsync(
hash_: string,
args: Record<string, any> | null = null,
scheduledInSeconds: number | null = null
): Promise<string> {
return _runScriptAsyncInternal(null, hash_, args, scheduledInSeconds);
}
/**
* Run a flow asynchronously by its path
* @param path - Flow path in Windmill
* @param args - Arguments to pass to the flow
* @param scheduledInSeconds - Schedule execution for a future time (in seconds)
* @param doNotTrackInParent - If false, tracks state in parent job (only use when fully awaiting the job)
* @returns Job ID of the created job
*/
export async function runFlowAsync(
path: string | null,
args: Record<string, any> | null,
scheduledInSeconds: number | null = null,
// can only be set to false if this the job will be fully await and not concurrent with any other job
// as otherwise the child flow and its own child will store their state in the parent job which will
// lead to incorrectness and failures
doNotTrackInParent: boolean = true
): Promise<string> {
// Create a script job and return its job id.
args = args || {};
const params: Record<string, any> = {};
if (scheduledInSeconds) {
params["scheduled_in_secs"] = scheduledInSeconds;
}
if (!doNotTrackInParent) {
let parentJobId = getEnv("WM_JOB_ID");
if (parentJobId !== undefined) {
params["parent_job"] = parentJobId;
}
let rootJobId = getEnv("WM_ROOT_FLOW_JOB_ID");
if (rootJobId != undefined && rootJobId != "") {
params["root_job"] = rootJobId;
}
}
let endpoint: string;
if (path) {
endpoint = `/w/${getWorkspace()}/jobs/run/f/${path}`;
} else {
throw new Error("path must be provided");
}
let url = new URL(OpenAPI.BASE + endpoint);
url.search = new URLSearchParams(params).toString();
return fetch(url, {
method: "POST",
headers: {
"Content-Type": "application/json",
Authorization: `Bearer ${OpenAPI.TOKEN}`,
},
body: JSON.stringify(args),
}).then((res) => res.text());
}
/**
* Resolve a resource value in case the default value was picked because the input payload was undefined
* @param obj resource value or path of the resource under the format `$res:path`
* @returns resource value
*/
export async function resolveDefaultResource(obj: any): Promise<any> {
if (typeof obj === "string" && obj.startsWith("$res:")) {
return await getResource(obj.substring(5), true);
} else {
return obj;
}
}
/**
* Get the state file path from environment variables
* @returns State path string
*/
export function getStatePath(): string {
const state_path = getEnv("WM_STATE_PATH_NEW") ?? getEnv("WM_STATE_PATH");
if (state_path === undefined) {
throw Error("State path not set");
}
return state_path;
}
/**
* Set a resource value by path
* @param path path of the resource to set, default to state path
* @param value new value of the resource to set
* @param initializeToTypeIfNotExist if the resource does not exist, initialize it with this type
*/
export async function setResource(
value: any,
path?: string,
initializeToTypeIfNotExist?: string
): Promise<void> {
path = parseResourceSyntax(path) ?? path ?? getStatePath();
const mockedApi = await getMockedApi();
if (mockedApi) {
mockedApi.resources[path] = value;
return;
}
const workspace = getWorkspace();
if (await ResourceService.existsResource({ workspace, path })) {
await ResourceService.updateResourceValue({
workspace,
path,
requestBody: { value },
});
} else if (initializeToTypeIfNotExist) {
await ResourceService.createResource({
workspace,
requestBody: { path, value, resource_type: initializeToTypeIfNotExist },
});
} else {
throw Error(
`Resource at path ${path} does not exist and no type was provided to initialize it`
);
}
}
/**
* Set the state
* @param state state to set
* @deprecated use setState instead
*/
export async function setInternalState(state: any): Promise<void> {
await setResource(state, undefined, "state");
}
/**
* Set the state
* @param state state to set
* @param path Optional state resource path override. Defaults to `getStatePath()`.
*/
export async function setState(state: any, path?: string): Promise<void> {
await setResource(state, path ?? getStatePath(), "state");
}
/**
* Set the progress
* Progress cannot go back and limited to 0% to 99% range
* @param percent Progress to set in %
* @param jobId? Job to set progress for
*/
export async function setProgress(percent: number, jobId?: any): Promise<void> {
const workspace = getWorkspace();
let flowId = getEnv("WM_FLOW_JOB_ID");
// If jobId specified we need to find if there is a parent/flow
if (jobId) {
const job = await JobService.getJob({
id: jobId ?? "NO_JOB_ID",
workspace,
noLogs: true,
});
// Could be actual flowId or undefined
flowId = job.parent_job;
}
await MetricsService.setJobProgress({
id: jobId ?? getEnv("WM_JOB_ID") ?? "NO_JOB_ID",
workspace,
requestBody: {
// In case user inputs float, it should be converted to int
percent: Math.floor(percent),
flow_job_id: flowId == "" ? undefined : flowId,
},
});
}
/**
* Get the progress
* @param jobId? Job to get progress from
* @returns Optional clamped between 0 and 100 progress value
*/
export async function getProgress(jobId?: any): Promise<number | null> {
// TODO: Delete or set to 100 completed job metrics
return await MetricsService.getJobProgress({
id: jobId ?? getEnv("WM_JOB_ID") ?? "NO_JOB_ID",
workspace: getWorkspace(),
});
}
/**
* Set a flow user state
* @param key key of the state
* @param value value of the state
*/
export async function setFlowUserState(
key: string,
value: any,
errorIfNotPossible?: boolean
): Promise<void> {
if (value === undefined) {
value = null;
}
const workspace = getWorkspace();
try {
await JobService.setFlowUserState({
workspace,
id: await getRootJobId(),
key,
requestBody: value,
});
} catch (e: any) {
if (errorIfNotPossible) {
throw Error(`Error setting flow user state at ${key}: ${e.body}`);
} else {
console.error(`Error setting flow user state at ${key}: ${e.body}`);
}
}
}
/**
* Get a flow user state
* @param path path of the variable
*/
export async function getFlowUserState(
key: string,
errorIfNotPossible?: boolean
): Promise<any> {
const workspace = getWorkspace();
try {
return await JobService.getFlowUserState({
workspace,
id: await getRootJobId(),
key,
});
} catch (e: any) {
if (errorIfNotPossible) {
throw Error(`Error setting flow user state at ${key}: ${e.body}`);
} else {
console.error(`Error setting flow user state at ${key}: ${e.body}`);
}
}
}
/**
* Get the internal state
* @deprecated use getState instead
*/
export async function getInternalState(): Promise<any> {
return await getResource(getStatePath(), true);
}
/**
* Get the state shared across executions
* @param path Optional state resource path override. Defaults to `getStatePath()`.
*/
export async function getState(path?: string): Promise<any> {
return await getResource(path ?? getStatePath(), true);
}
/**
* Get a variable by path
* @param path path of the variable
* @returns variable value
*/
export async function getVariable(path: string): Promise<string> {
path = parseVariableSyntax(path) ?? path;
const mockedApi = await getMockedApi();
if (mockedApi) {
if (mockedApi.variables[path]) {
return mockedApi.variables[path];
} else {
console.log(
`MockedAPI present, but variable not found at ${path}, falling back to real API`
);
}
}
const workspace = getWorkspace();
try {
return await VariableService.getVariableValue({ workspace, path });
} catch (e: any) {
throw Error(
`Variable not found at ${path} or not visible to you: ${e.body}`
);
}
}
/**
* Set a variable by path, create if not exist
* @param path path of the variable
* @param value value of the variable
* @param isSecretIfNotExist if the variable does not exist, create it as secret or not (default: false)
* @param descriptionIfNotExist if the variable does not exist, create it with this description (default: "")
*/
export async function setVariable(
path: string,
value: string,
isSecretIfNotExist?: boolean,
descriptionIfNotExist?: string
): Promise<void> {
path = parseVariableSyntax(path) ?? path;
const mockedApi = await getMockedApi();
if (mockedApi) {
mockedApi.variables[path] = value;
return;
}
const workspace = getWorkspace();
if (await VariableService.existsVariable({ workspace, path })) {
await VariableService.updateVariable({
workspace,
path,
requestBody: { value },
});
} else {
await VariableService.createVariable({
workspace,
requestBody: {
path,
value,
is_secret: isSecretIfNotExist ?? false,
description: descriptionIfNotExist ?? "",
},
});
}
}
/**
* Build a PostgreSQL connection URL from a database resource
* @param path - Path to the database resource
* @returns PostgreSQL connection URL string
*/
export async function databaseUrlFromResource(path: string): Promise<string> {
const resource = await getResource(path);
return `postgresql://${resource.user}:${resource.password}@${resource.host}:${resource.port}/${resource.dbname}?sslmode=${resource.sslmode}`;
}
// TODO(gb): need to investigate more how Polars and DuckDB work in TS
// export async function polarsConnectionSettings(s3_resource_path: string | undefined): Promise<any> {
// const workspace = getWorkspace();
// return await HelpersService.polarsConnectionSettingsV2({
// workspace: workspace,
// requestBody: {
// s3_resource_path: s3_resource_path
// }
// });
// }
// export async function duckdbConnectionSettings(s3_resource_path: string | undefined): Promise<any> {
// const workspace = getWorkspace();
// return await HelpersService.duckdbConnectionSettingsV2({
// workspace: workspace,
// requestBody: {
// s3_resource_path: s3_resource_path
// }
// });
// }
/**
* Get S3 client settings from a resource or workspace default
* @param s3_resource_path - Path to S3 resource (uses workspace default if undefined)
* @returns S3 client configuration settings
*/
export async function denoS3LightClientSettings(
s3_resource_path: string | undefined
): Promise<DenoS3LightClientSettings> {
const workspace = getWorkspace();
const s3Resource = await HelpersService.s3ResourceInfo({
workspace: workspace,
requestBody: {
s3_resource_path:
parseResourceSyntax(s3_resource_path) ?? s3_resource_path,
},
});
let settings: DenoS3LightClientSettings = {
...s3Resource,
};
return settings;
}
/**
* Load the content of a file stored in S3. If the s3ResourcePath is undefined, it will default to the workspace S3 resource.
*
* ```typescript
* let fileContent = await wmill.loadS3FileContent(inputFile)
* // if the file is a raw text file, it can be decoded and printed directly:
* const text = new TextDecoder().decode(fileContentStream)
* console.log(text);
* ```
*/
export async function loadS3File(
s3object: S3Object,
s3ResourcePath: string | undefined = undefined
): Promise<Uint8Array | undefined> {
const fileContentBlob = await loadS3FileStream(s3object, s3ResourcePath);
if (fileContentBlob === undefined) {
return undefined;
}
// we read the stream until completion and put the content in an Uint8Array
const reader = fileContentBlob.stream().getReader();
const chunks: Uint8Array[] = [];
while (true) {
const { value: chunk, done } = await reader.read();
if (done) {
break;
}
chunks.push(chunk);
}
let fileContentLength = 0;
chunks.forEach((item) => {
fileContentLength += item.length;
});
let fileContent = new Uint8Array(fileContentLength);
let offset = 0;
chunks.forEach((chunk) => {
fileContent.set(chunk, offset);
offset += chunk.length;
});
return fileContent;
}
/**
* Load the content of a file stored in S3 as a stream. If the s3ResourcePath is undefined, it will default to the workspace S3 resource.
*
* ```typescript
* let fileContentBlob = await wmill.loadS3FileStream(inputFile)
* // if the content is plain text, the blob can be read directly:
* console.log(await fileContentBlob.text());
* ```
*/
export async function loadS3FileStream(
s3object: S3Object,
s3ResourcePath: string | undefined = undefined
): Promise<Blob | undefined> {
let s3Obj = s3object && parseS3Object(s3object);
let params: Record<string, string> = {};
params["file_key"] = s3Obj.s3;
if (s3ResourcePath !== undefined) {
params["s3_resource_path"] = s3ResourcePath;
}
if (s3Obj.storage !== undefined) {
params["storage"] = s3Obj.storage;
}
const queryParams = new URLSearchParams(params);
// We use raw fetch here b/c OpenAPI generated client doesn't handle Blobs nicely
const response = await fetch(
`${
OpenAPI.BASE
}/w/${getWorkspace()}/job_helpers/download_s3_file?${queryParams}`,
{
method: "GET",
headers: {
Authorization: `Bearer ${OpenAPI.TOKEN}`,
},
}
);
// Check if the response was successful
if (!response.ok) {
const errorText = await response.text();
throw new Error(
`Failed to load S3 file: ${response.status} ${response.statusText} - ${errorText}`
);
}
return response.blob();
}
/**
* Persist a file to the S3 bucket. If the s3ResourcePath is undefined, it will default to the workspace S3 resource.
*
* ```typescript
* const s3object = await writeS3File(s3Object, "Hello Windmill!")
* const fileContentAsUtf8Str = (await s3object.toArray()).toString('utf-8')
* console.log(fileContentAsUtf8Str)
* ```
*/
export async function writeS3File(
s3object: S3Object | undefined,
fileContent: string | Blob,
s3ResourcePath: string | undefined = undefined,
contentType: string | undefined = undefined,
contentDisposition: string | undefined = undefined
): Promise<S3Object> {
let fileContentBlob: Blob;
if (typeof fileContent === "string") {
fileContentBlob = new Blob([fileContent as string], {
type: "text/plain",
});
} else {
fileContentBlob = fileContent as Blob;
}
let s3Obj = s3object && parseS3Object(s3object);
const response = await HelpersService.fileUpload({
workspace: getWorkspace(),
fileKey: s3Obj?.s3,
fileExtension: undefined,
s3ResourcePath: s3ResourcePath,
requestBody: fileContentBlob,
storage: s3Obj?.storage,
contentType,
contentDisposition,
});
return {
s3: response.file_key,
...(s3Obj?.storage && { storage: s3Obj?.storage }),
};
}
/**
* Sign S3 objects to be used by anonymous users in public apps
* @param s3objects s3 objects to sign
* @returns signed s3 objects
*/
export async function signS3Objects(
s3objects: S3Object[]
): Promise<S3Object[]> {
const signedKeys = await AppService.signS3Objects({
workspace: getWorkspace(),
requestBody: {
s3_objects: s3objects.map(parseS3Object),
},
});
return signedKeys;
}
/**
* Sign S3 object to be used by anonymous users in public apps
* @param s3object s3 object to sign
* @returns signed s3 object
*/
export async function signS3Object(s3object: S3Object): Promise<S3Object> {
const [signedObject] = await signS3Objects([s3object]);
return signedObject;
}
/**
* Generate a presigned public URL for an array of S3 objects.
* If an S3 object is not signed yet, it will be signed first.
* @param s3Objects s3 objects to sign
* @returns list of signed public URLs
*/
export async function getPresignedS3PublicUrls(
s3Objects: S3Object[],
{ baseUrl }: { baseUrl?: string } = {}
): Promise<string[]> {
baseUrl ??= getPublicBaseUrl();
const s3Objs = s3Objects.map(parseS3Object);
// Sign all S3 objects that need to be signed in one go
const s3ObjsToSign: (readonly [S3ObjectRecord, number])[] = s3Objs
.map((s3Obj, index) => [s3Obj, index] as const)
.filter(([s3Obj, _]) => s3Obj.presigned === undefined);
if (s3ObjsToSign.length > 0) {
const signedS3Objs = await signS3Objects(
s3ObjsToSign.map(([s3Obj, _]) => s3Obj)
);
for (let i = 0; i < s3ObjsToSign.length; i++) {
const [_, originalIndex] = s3ObjsToSign[i];
s3Objs[originalIndex] = parseS3Object(signedS3Objs[i]);
}
}
const signedUrls: string[] = [];
for (const s3Obj of s3Objs) {
const { s3, presigned, storage = "_default_" } = s3Obj;
const signedUrl = `${baseUrl}/api/w/${getWorkspace()}/s3_proxy/${storage}/${s3}?${presigned}`;
signedUrls.push(signedUrl);
}
return signedUrls;
}
/**
* Generate a presigned public URL for an S3 object. If the S3 object is not signed yet, it will be signed first.
* @param s3Object s3 object to sign
* @returns signed public URL
*/
export async function getPresignedS3PublicUrl(
s3Objects: S3Object,
{ baseUrl }: { baseUrl?: string } = {}
): Promise<string> {
const [s3Object] = await getPresignedS3PublicUrls([s3Objects], { baseUrl });
return s3Object;
}
/**
* Get URLs needed for resuming a flow after this step
* @param approver approver name
* @param flowLevel if true, generate resume URLs for the parent flow instead of the specific step.
* This allows pre-approvals that can be consumed by any later suspend step in the same flow.
* @returns approval page UI URL, resume and cancel API URLs for resuming the flow
*/
export async function getResumeUrls(
approver?: string,
flowLevel?: boolean
): Promise<{
approvalPage: string;
resume: string;
cancel: string;
}> {
const nonce = Math.floor(Math.random() * 4294967295);
const workspace = getWorkspace();
return await JobService.getResumeUrls({
workspace,
resumeId: nonce,
approver,
flowLevel,
id: getEnv("WM_JOB_ID") ?? "NO_JOB_ID",
});
}
/**
* @deprecated use getResumeUrls instead
*/
export function getResumeEndpoints(approver?: string): Promise<{
approvalPage: string;
resume: string;
cancel: string;
}> {
return getResumeUrls(approver);
}
/**
* Get an OIDC jwt token for auth to external services (e.g: Vault, AWS) (ee only)
* @param audience audience of the token
* @param expiresIn Optional number of seconds until the token expires
* @returns jwt token
*/
export async function getIdToken(
audience: string,
expiresIn?: number
): Promise<string> {
const workspace = getWorkspace();
return await OidcService.getOidcToken({
workspace,
audience,
expiresIn,
});
}
/**
* Convert a base64-encoded string to Uint8Array
* @param data - Base64-encoded string
* @returns Decoded Uint8Array
*/
export function base64ToUint8Array(data: string): Uint8Array {
return Uint8Array.from(atob(data), (c) => c.charCodeAt(0));
}
/**
* Convert a Uint8Array to base64-encoded string
* @param arrayBuffer - Uint8Array to encode
* @returns Base64-encoded string
*/
export function uint8ArrayToBase64(arrayBuffer: Uint8Array): string {
let base64 = "";
const encodings =
"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
const bytes = new Uint8Array(arrayBuffer);
const byteLength = bytes.byteLength;
const byteRemainder = byteLength % 3;
const mainLength = byteLength - byteRemainder;
let a, b, c, d;
let chunk;
// Main loop deals with bytes in chunks of 3
for (let i = 0; i < mainLength; i = i + 3) {
// Combine the three bytes into a single integer
chunk = (bytes[i] << 16) | (bytes[i + 1] << 8) | bytes[i + 2];
// Use bitmasks to extract 6-bit segments from the triplet
a = (chunk & 16515072) >> 18; // 16515072 = (2^6 - 1) << 18
b = (chunk & 258048) >> 12; // 258048 = (2^6 - 1) << 12
c = (chunk & 4032) >> 6; // 4032 = (2^6 - 1) << 6
d = chunk & 63; // 63 = 2^6 - 1
// Convert the raw binary segments to the appropriate ASCII encoding
base64 += encodings[a] + encodings[b] + encodings[c] + encodings[d];
}
// Deal with the remaining bytes and padding
if (byteRemainder == 1) {
chunk = bytes[mainLength];
a = (chunk & 252) >> 2; // 252 = (2^6 - 1) << 2
// Set the 4 least significant bits to zero
b = (chunk & 3) << 4; // 3 = 2^2 - 1
base64 += encodings[a] + encodings[b] + "==";
} else if (byteRemainder == 2) {
chunk = (bytes[mainLength] << 8) | bytes[mainLength + 1];
a = (chunk & 64512) >> 10; // 64512 = (2^6 - 1) << 10
b = (chunk & 1008) >> 4; // 1008 = (2^6 - 1) << 4
// Set the 2 least significant bits to zero
c = (chunk & 15) << 2; // 15 = 2^4 - 1
base64 += encodings[a] + encodings[b] + encodings[c] + "=";
}
return base64;
}
/**
* Get email from workspace username
* This method is particularly useful for apps that require the email address of the viewer.
* Indeed, in the viewer context, WM_USERNAME is set to the username of the viewer but WM_EMAIL is set to the email of the creator of the app.
* @param username
* @returns email address
*/
export async function usernameToEmail(username: string): Promise<string> {
const workspace = getWorkspace();
return await UserService.usernameToEmail({ username, workspace });
}
interface SlackApprovalOptions {
slackResourcePath: string;
channelId: string;
message?: string;
approver?: string;
defaultArgsJson?: Record<string, any>;
dynamicEnumsJson?: Record<string, any>;
resumeButtonText?: string;
cancelButtonText?: string;
}
interface TeamsApprovalOptions {
teamName: string;
channelName: string;
message?: string;
approver?: string;
defaultArgsJson?: Record<string, any>;
dynamicEnumsJson?: Record<string, any>;
}
/**
* Sends an interactive approval request via Slack, allowing optional customization of the message, approver, and form fields.
*
* **[Enterprise Edition Only]** To include form fields in the Slack approval request, go to **Advanced -> Suspend -> Form**
* and define a form. Learn more at [Windmill Documentation](https://www.windmill.dev/docs/flows/flow_approval#form).
*
* @param {Object} options - The configuration options for the Slack approval request.
* @param {string} options.slackResourcePath - The path to the Slack resource in Windmill.
* @param {string} options.channelId - The Slack channel ID where the approval request will be sent.
* @param {string} [options.message] - Optional custom message to include in the Slack approval request.
* @param {string} [options.approver] - Optional user ID or name of the approver for the request.
* @param {DefaultArgs} [options.defaultArgsJson] - Optional object defining or overriding the default arguments to a form field.
* @param {Enums} [options.dynamicEnumsJson] - Optional object overriding the enum default values of an enum form field.
* @param {string} [options.resumeButtonText] - Optional text for the resume button.
* @param {string} [options.cancelButtonText] - Optional text for the cancel button.
*
* @returns {Promise<void>} Resolves when the Slack approval request is successfully sent.
*
* @throws {Error} If the function is not called within a flow or flow preview.
* @throws {Error} If the `JobService.getSlackApprovalPayload` call fails.
*
* **Usage Example:**
* ```typescript
* await requestInteractiveSlackApproval({
* slackResourcePath: "/u/alex/my_slack_resource",
* channelId: "admins-slack-channel",
* message: "Please approve this request",
* approver: "approver123",
* defaultArgsJson: { key1: "value1", key2: 42 },
* dynamicEnumsJson: { foo: ["choice1", "choice2"], bar: ["optionA", "optionB"] },
* resumeButtonText: "Resume",
* cancelButtonText: "Cancel",
* });
* ```
*
* **Note:** This function requires execution within a Windmill flow or flow preview.
*/
export async function requestInteractiveSlackApproval({
slackResourcePath,
channelId,
message,
approver,
defaultArgsJson,
dynamicEnumsJson,
resumeButtonText,
cancelButtonText,
}: SlackApprovalOptions): Promise<void> {
const workspace = getWorkspace();
const flowJobId = getEnv("WM_FLOW_JOB_ID");
if (!flowJobId) {
throw new Error(
"You can't use this function in a standalone script or flow step preview. Please use it in a flow or a flow preview."
);
}
const flowStepId = getEnv("WM_FLOW_STEP_ID");
if (!flowStepId) {
throw new Error("This function can only be called as a flow step");
}
// Only include non-empty parameters
const params: {
approver?: string;
message?: string;
slackResourcePath: string;
channelId: string;
flowStepId: string;
defaultArgsJson?: string;
dynamicEnumsJson?: string;
resumeButtonText?: string;
cancelButtonText?: string;
} = {
slackResourcePath,
channelId,
flowStepId,
};
if (message) {
params.message = message;
}
if (approver) {
params.approver = approver;
}
if (defaultArgsJson) {
params.defaultArgsJson = JSON.stringify(defaultArgsJson);
}
if (dynamicEnumsJson) {
params.dynamicEnumsJson = JSON.stringify(dynamicEnumsJson);
}
if (resumeButtonText) {
params.resumeButtonText = resumeButtonText;
}
if (cancelButtonText) {
params.cancelButtonText = cancelButtonText;
}
await JobService.getSlackApprovalPayload({
workspace,
...params,
id: getEnv("WM_JOB_ID") ?? "NO_JOB_ID",
});
}
/**
* Sends an interactive approval request via Teams, allowing optional customization of the message, approver, and form fields.
*
* **[Enterprise Edition Only]** To include form fields in the Teams approval request, go to **Advanced -> Suspend -> Form**
* and define a form. Learn more at [Windmill Documentation](https://www.windmill.dev/docs/flows/flow_approval#form).
*
* @param {Object} options - The configuration options for the Teams approval request.
* @param {string} options.teamName - The Teams team name where the approval request will be sent.
* @param {string} options.channelName - The Teams channel name where the approval request will be sent.
* @param {string} [options.message] - Optional custom message to include in the Teams approval request.
* @param {string} [options.approver] - Optional user ID or name of the approver for the request.
* @param {DefaultArgs} [options.defaultArgsJson] - Optional object defining or overriding the default arguments to a form field.
* @param {Enums} [options.dynamicEnumsJson] - Optional object overriding the enum default values of an enum form field.
*
* @returns {Promise<void>} Resolves when the Teams approval request is successfully sent.
*
* @throws {Error} If the function is not called within a flow or flow preview.
* @throws {Error} If the `JobService.getTeamsApprovalPayload` call fails.
*
* **Usage Example:**
* ```typescript
* await requestInteractiveTeamsApproval({
* teamName: "admins-teams",
* channelName: "admins-teams-channel",
* message: "Please approve this request",
* approver: "approver123",
* defaultArgsJson: { key1: "value1", key2: 42 },
* dynamicEnumsJson: { foo: ["choice1", "choice2"], bar: ["optionA", "optionB"] },
* });
* ```
*
* **Note:** This function requires execution within a Windmill flow or flow preview.
*/
export async function requestInteractiveTeamsApproval({
teamName,
channelName,
message,
approver,
defaultArgsJson,
dynamicEnumsJson,
}: TeamsApprovalOptions): Promise<void> {
const workspace = getWorkspace();
const flowJobId = getEnv("WM_FLOW_JOB_ID");
if (!flowJobId) {
throw new Error(
"You can't use this function in a standalone script or flow step preview. Please use it in a flow or a flow preview."
);
}
const flowStepId = getEnv("WM_FLOW_STEP_ID");
if (!flowStepId) {
throw new Error("This function can only be called as a flow step");
}
// Only include non-empty parameters
const params: {
approver?: string;
message?: string;
teamName: string;
channelName: string;
flowStepId: string;
defaultArgsJson?: string;
dynamicEnumsJson?: string;
} = {
teamName,
channelName,
flowStepId,
};
if (message) {
params.message = message;
}
if (approver) {
params.approver = approver;
}
if (defaultArgsJson) {
params.defaultArgsJson = JSON.stringify(defaultArgsJson);
}
if (dynamicEnumsJson) {
params.dynamicEnumsJson = JSON.stringify(dynamicEnumsJson);
}
await JobService.getTeamsApprovalPayload({
workspace,
...params,
id: getEnv("WM_JOB_ID") ?? "NO_JOB_ID",
});
}
async function getMockedApi(): Promise<MockedApi | undefined> {
if (mockedApi) {
return mockedApi;
}
const mockedPath = getEnv("WM_MOCKED_API_FILE");
if (mockedPath) {
console.info("Using mocked API from", mockedPath);
} else {
return undefined;
}
try {
const fs = await import("node:fs/promises");
const file = await fs.readFile(mockedPath, "utf-8");
try {
mockedApi = JSON.parse(file) as MockedApi;
if (!mockedApi.variables) {
mockedApi.variables = {};
}
if (!mockedApi.resources) {
mockedApi.resources = {};
}
return mockedApi;
} catch {
console.warn("Error parsing mocked API file at path", mockedPath);
}
} catch {
console.warn("Error reading mocked API file at path", mockedPath);
}
if (!mockedApi) {
console.warn(
"No mocked API file path provided at env variable WM_MOCKED_API_FILE. Using empty mocked API."
);
mockedApi = {
variables: {},
resources: {},
};
return mockedApi;
}
}
interface MockedApi {
variables: Record<string, string>;
resources: Record<string, any>;
}
function parseResourceSyntax(s: string | undefined) {
if (s?.startsWith("$res:")) return s.substring(5);
if (s?.startsWith("res://")) return s.substring(6);
}
/**
* Parse an S3 object from URI string or record format
* @param s3Object - S3 object as URI string (s3://storage/key) or record
* @returns S3 object record with storage and s3 key
*/
export function parseS3Object(s3Object: S3Object): S3ObjectRecord {
if (typeof s3Object === "object") return s3Object;
const match = s3Object.match(/^s3:\/\/([^/]*)\/(.*)$/);
return { storage: match?.[1] || undefined, s3: match?.[2] ?? "" };
}
function parseVariableSyntax(s: string) {
if (s.startsWith("var://")) return s.substring(6);
}
// ── Workflow-as-Code SDK ──────────────────────────────────────────────
export class StepSuspend extends Error {
constructor(public dispatchInfo: Record<string, any>) {
super("__step_suspend__");
this.name = "StepSuspend";
}
}
export interface TaskOptions {
timeout?: number;
tag?: string;
cache_ttl?: number;
priority?: number;
concurrency_limit?: number;
concurrency_key?: string;
concurrency_time_window_s?: number;
}
export let _workflowCtx: WorkflowCtx | null = null;
export function setWorkflowCtx(ctx: WorkflowCtx | null) {
_workflowCtx = ctx;
Reflect.set(globalThis, "__wmill_wf_ctx", ctx);
}
export class WorkflowCtx {
private completed: Record<string, any>;
private counters: Record<string, number> = {};
private pending: Array<{
name: string;
script: string;
args: Record<string, any>;
key: string;
dispatch_type: string;
[k: string]: any;
}> = [];
private _suspended = false;
/** When set, the task matching this key executes its inner function directly */
_executingKey: string | null;
constructor(checkpoint: Record<string, any> = {}) {
this.completed = checkpoint?.completed_steps ?? {};
this._executingKey = checkpoint?._executing_key ?? null;
}
/** Name-based key: `double` for first call, `double_2`, `double_3` for subsequent. */
_allocKey(name: string): string {
const n = (this.counters[name] ?? 0) + 1;
this.counters[name] = n;
return n === 1 ? name : `${name}_${n}`;
}
_nextStep(
name: string,
script: string,
args: Record<string, any> = {},
dispatch_type: string = "inline",
options?: TaskOptions,
): PromiseLike<any> {
const key = this._allocKey(name || script || "step");
if (key in this.completed) {
const value = this.completed[key];
if (value && typeof value === "object" && (value as any).__wmill_error) {
const err = new Error((value as any).message || `Task '${name}' failed`);
(err as any).result = (value as any).result;
(err as any).step_key = (value as any).step_key;
(err as any).child_job_id = (value as any).child_job_id;
return { then: (_resolve: any, reject?: any) => { if (reject) reject(err); else throw err; } } as PromiseLike<any>;
}
return { then: (resolve: any) => resolve(value) };
}
// If this is a child job executing a specific step, return null to signal
// that the task wrapper should run the inner function directly
if (this._executingKey === key) {
return { then: (resolve: any) => resolve(null), _execute_directly: true } as any;
}
// In child job mode (_executingKey is set), non-matching uncompleted steps
// should never resolve or throw — the matching step will throw step_complete
// which terminates the workflow. Returning a never-resolving thenable prevents
// race conditions where a non-matching step's StepSuspend fires before step_complete.
if (this._executingKey !== null) {
return { then: () => new Promise(() => {}) };
}
const stepInfo: any = { name: name || key, script: script || key, args, key, dispatch_type };
if (options) {
if (options.timeout !== undefined) stepInfo.timeout = options.timeout;
if (options.tag !== undefined) stepInfo.tag = options.tag;
if (options.cache_ttl !== undefined) stepInfo.cache_ttl = options.cache_ttl;
if (options.priority !== undefined) stepInfo.priority = options.priority;
if (options.concurrency_limit !== undefined) stepInfo.concurrent_limit = options.concurrency_limit;
if (options.concurrency_key !== undefined) stepInfo.concurrency_key = options.concurrency_key;
if (options.concurrency_time_window_s !== undefined) stepInfo.concurrency_time_window_s = options.concurrency_time_window_s;
}
this.pending.push(stepInfo);
return {
then: (): never => {
// Only the first .then() call throws with all accumulated steps.
// Subsequent calls (e.g. from Promise.all resolving other thenables)
// also throw (they'll be caught by the same handler).
if (this._suspended) return new Promise(() => {}) as never;
this._suspended = true;
const steps = [...this.pending];
this.pending = [];
throw new StepSuspend({
mode: steps.length > 1 ? "parallel" : "sequential",
steps,
});
},
};
}
/** Return and clear any pending (unawaited) steps. */
_flushPending(): Array<{ name: string; script: string; args: Record<string, any>; key: string; dispatch_type: string }> {
const steps = [...this.pending];
this.pending = [];
return steps;
}
_waitForApproval(options?: {
timeout?: number;
form?: object;
}): PromiseLike<{ value: any; approver: string; approved: boolean }> {
const key = this._allocKey("approval");
if (key in this.completed) {
const value = this.completed[key];
return { then: (resolve: any) => resolve(value) };
}
// In child job mode, return never-resolving thenable (same as _nextStep)
if (this._executingKey !== null) {
return { then: () => new Promise(() => {}) };
}
// Throw immediately — approval is always a blocking step
throw new StepSuspend({
mode: "approval",
key,
timeout: options?.timeout ?? 1800,
form: options?.form,
steps: [],
});
}
_sleep(seconds: number): PromiseLike<void> {
const key = this._allocKey("sleep");
if (key in this.completed) {
return { then: (resolve: any) => resolve(undefined) };
}
if (this._executingKey !== null) {
return { then: () => new Promise(() => {}) };
}
throw new StepSuspend({
mode: "sleep",
key,
seconds: Math.max(1, Math.round(seconds)),
steps: [],
});
}
async _runInlineStep<T>(name: string, fn: () => T | Promise<T>): Promise<T> {
const key = this._allocKey(name || "step");
if (key in this.completed) {
const value = this.completed[key];
if (value && typeof value === "object" && (value as any).__wmill_error) {
const err = new Error((value as any).message || `Step '${name}' failed`);
(err as any).result = (value as any).result;
(err as any).step_key = (value as any).step_key;
(err as any).child_job_id = (value as any).child_job_id;
throw err;
}
return value as T;
}
if (this._executingKey !== null) {
return new Promise(() => {});
}
const result = await fn();
throw new StepSuspend({ mode: "inline_checkpoint", steps: [], key, result });
}
}
export async function sleep(seconds: number): Promise<void> {
const ctx: WorkflowCtx | null = _workflowCtx ?? Reflect.get(globalThis, "__wmill_wf_ctx");
if (ctx) {
return ctx._sleep(seconds) as Promise<void>;
}
// Outside workflow context, just wait locally
await new Promise((r) => setTimeout(r, seconds * 1000));
}
export async function step<T>(name: string, fn: () => T | Promise<T>): Promise<T> {
const ctx: WorkflowCtx | null = _workflowCtx ?? Reflect.get(globalThis, "__wmill_wf_ctx");
if (ctx) {
return ctx._runInlineStep(name, fn);
}
return fn();
}
/**
* Wrap an async function as a workflow task.
*
* @example
* const extract_data = task(async (url: string) => { ... });
* const run_external = task("f/external_script", async (x: number) => { ... });
*
* Inside a `workflow()`, calling a task dispatches it as a step.
* Outside a workflow, the function body executes directly.
*/
export function task<T extends (...args: any[]) => Promise<any>>(
fnOrPath: T | string,
maybeFnOrOptions?: T | TaskOptions,
maybeOptions?: TaskOptions,
): T {
let fn: T;
let taskPath: string | undefined;
let taskOptions: TaskOptions | undefined;
if (typeof fnOrPath === "string") {
taskPath = fnOrPath;
fn = maybeFnOrOptions as T;
taskOptions = maybeOptions;
} else {
fn = fnOrPath;
taskOptions = maybeFnOrOptions as TaskOptions | undefined;
}
const taskName = fn.name || taskPath || "";
// NOT async — in workflow context we return the thenable directly so that
// unawaited task calls leave the step in ctx.pending (for _flushPending).
// An async wrapper would auto-resolve the thenable in a microtask, calling
// .then() which throws StepSuspend and empties pending before the caller
// can flush.
const wrapper = function (...args: any[]) {
const ctx: WorkflowCtx | null = _workflowCtx ?? Reflect.get(globalThis, "__wmill_wf_ctx");
if (ctx) {
// Inside a workflow with checkpoint/replay context — dispatch as step
const script = taskPath ?? taskName;
const paramNames = getParamNames(fn);
const kwargs: Record<string, any> = {};
for (let i = 0; i < args.length; i++) {
if (paramNames[i]) {
kwargs[paramNames[i]] = args[i];
} else {
kwargs[`arg${i}`] = args[i];
}
}
const stepResult = ctx._nextStep(taskName, script, kwargs, "inline", taskOptions);
// If this step should execute directly (child job mode), run the inner function
// and throw StepSuspend with mode "step_complete" to signal that we're done
if ((stepResult as any)?._execute_directly) {
return (async () => {
const result = await fn(...args);
throw new StepSuspend({ mode: "step_complete", steps: [], result });
})();
}
return stepResult;
} else if (getEnv("WM_JOB_ID") && !getEnv("WM_FLOW_JOB_ID")) {
// Inside a Windmill root job without checkpoint context — v1 HTTP dispatch
// WM_FLOW_JOB_ID is set on child jobs, so we skip dispatch for those
return (async () => {
const paramNames = getParamNames(fn);
const kwargs: Record<string, any> = {};
args.forEach((x, i) => (kwargs[paramNames[i]] = x));
let req = await fetch(
`${OpenAPI.BASE}/w/${getWorkspace()}/jobs/run/workflow_as_code/${getEnv(
"WM_JOB_ID"
)}/${taskName}`,
{
method: "POST",
headers: {
"Content-Type": "application/json",
Authorization: `Bearer ${getEnv("WM_TOKEN")}`,
},
body: JSON.stringify({ args: kwargs }),
}
);
let jobId = await req.text();
console.log(`Started task ${taskName} as job ${jobId}`);
let r = await waitJob(jobId);
console.log(`Task ${taskName} (${jobId}) completed`);
return r;
})();
} else {
// Standalone — execute directly
return fn(...args);
}
} as unknown as T;
Object.defineProperty(wrapper, "name", { value: taskName });
(wrapper as any)._is_task = true;
(wrapper as any)._task_path = taskPath;
return wrapper;
}
/**
* Create a task that dispatches to a separate Windmill script.
*
* @example
* const extract = taskScript("f/data/extract");
* // inside workflow: await extract({ url: "https://..." })
*/
export function taskScript(path: string, options?: TaskOptions): (...args: any[]) => PromiseLike<any> {
const name = path.split("/").pop() || path;
const wrapper = function (...args: any[]) {
const ctx: WorkflowCtx | null = _workflowCtx ?? Reflect.get(globalThis, "__wmill_wf_ctx");
if (ctx) {
const kwargs = args.length === 1 && typeof args[0] === "object" && args[0] !== null
? args[0]
: args.reduce((acc, v, i) => { acc[`arg${i}`] = v; return acc; }, {} as Record<string, any>);
return ctx._nextStep(name, path, kwargs, "script", options);
}
throw new Error(`taskScript("${path}") can only be called inside a workflow()`);
};
Object.defineProperty(wrapper, "name", { value: name });
(wrapper as any)._is_task = true;
(wrapper as any)._task_path = path;
return wrapper;
}
/**
* Create a task that dispatches to a separate Windmill flow.
*
* @example
* const pipeline = taskFlow("f/etl/pipeline");
* // inside workflow: await pipeline({ input: data })
*/
export function taskFlow(path: string, options?: TaskOptions): (...args: any[]) => PromiseLike<any> {
const name = path.split("/").pop() || path;
const wrapper = function (...args: any[]) {
const ctx: WorkflowCtx | null = _workflowCtx ?? Reflect.get(globalThis, "__wmill_wf_ctx");
if (ctx) {
const kwargs = args.length === 1 && typeof args[0] === "object" && args[0] !== null
? args[0]
: args.reduce((acc, v, i) => { acc[`arg${i}`] = v; return acc; }, {} as Record<string, any>);
return ctx._nextStep(name, path, kwargs, "flow", options);
}
throw new Error(`taskFlow("${path}") can only be called inside a workflow()`);
};
Object.defineProperty(wrapper, "name", { value: name });
(wrapper as any)._is_task = true;
(wrapper as any)._task_path = path;
return wrapper;
}
/**
* Mark an async function as a workflow-as-code entry point.
*
* The function must be **deterministic**: given the same inputs it must call
* tasks in the same order on every replay. Branching on task results is fine
* (results are replayed from checkpoint), but branching on external state
* (current time, random values, external API calls) must use `step()` to
* checkpoint the value so replays see the same result.
*/
export function workflow<T>(fn: (...args: any[]) => Promise<T>) {
(fn as any)._is_workflow = true;
return fn;
}
/**
* Suspend the workflow and wait for an external approval.
*
* Use `getResumeUrls()` (wrapped in `step()`) to obtain resume/cancel/approvalPage
* URLs before calling this function.
*
* @example
* const urls = await step("urls", () => getResumeUrls());
* await step("notify", () => sendEmail(urls.approvalPage));
* const { value, approver } = await waitForApproval({ timeout: 3600 });
*/
export function waitForApproval(options?: {
timeout?: number;
form?: object;
}): PromiseLike<{ value: any; approver: string; approved: boolean }> {
const ctx: WorkflowCtx | null = _workflowCtx ?? Reflect.get(globalThis, "__wmill_wf_ctx");
if (!ctx) {
throw new Error("waitForApproval can only be called inside a workflow()");
}
return ctx._waitForApproval(options);
}
/**
* Process items in parallel with optional concurrency control.
*
* Each item is processed by calling `fn(item)`, which should be a task().
* Items are dispatched in batches of `concurrency` (default: all at once).
*
* @example
* const process = task(async (item: string) => { ... });
* const results = await parallel(items, process, { concurrency: 5 });
*/
export async function parallel<T, R>(
items: T[],
fn: (item: T) => PromiseLike<R> | R,
options?: { concurrency?: number },
): Promise<R[]> {
const concurrency = options?.concurrency ?? items.length;
if (concurrency <= 0 || items.length === 0) return [];
const results: R[] = [];
for (let i = 0; i < items.length; i += concurrency) {
const batch = items.slice(i, i + concurrency);
const batchResults = await Promise.all(batch.map((item) => fn(item)));
results.push(...batchResults);
}
return results;
}