* feat: teach the AI the raw-app job bindings and the draft/deployed split Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: scope the raw-app deploy advice to the referenced item, and stop kind-conversion from stranding fields The draft/deployed guidance added in the previous commit was read as "deploy the app too": the agent asked for both the flow and the app and routed a one-item dependency through the review-and-deploy page. Only the referenced flow or script has to exist deployed — the preview runs the app's draft — so the prompts, the `write_app_runnable` warning and the testing rule now say to offer that one deploy and leave the app a draft. `buildPersistedRunnable` spread the existing runnable when rewriting it, so converting a path runnable to inline left `runType`/`path` behind (and the reverse left `inlineScript`). `isRunnableByName` matches the inline branch first, so an app "wired to a flow" silently ran stale inline code. `test_run_app_runnable` now fills ctx-bound inputs with `$ctx:<prop>` the way RawAppBackgroundRunner does, so a ctx argument no longer arrives missing. The SDK-reference rationale claimed WM_TOKEN may be unset, that a missing base URL falls back to localhost, and that a job token is scoped enough to 403 a hand-rolled REST call. None of the three is true, and it shipped to every write-script prompt; the text now only says the client configures itself. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: address review round on the raw-app AI instructions The eval case could pass on the exact answer it exists to reject. Every `requiredMentionsAnyOf` alternative but one was flow-agnostic, so "the app must be deployed" satisfied "must be deployed". All alternatives now name the flow, and a unit test pins that the app-only phrasing fails. `instanceLine` asserted "self-hosted Community Edition" outside the browser, where `isCloudHosted()` reads false and the license store is unset — so every global eval was told that regardless of what it pointed at. It is now emitted only under BROWSER. `assistantExpect.forbiddenMentions` defaulted a missing `assistantText` to "", which passes every entry forever on a mode whose runner does not report it. It now fails with that as the reason. `buildPersistedRunnable` carried `schema` across a retarget, so a path runnable pointed at a new flow kept the previous item's schema and `genWmillTs` typed `backend.<key>(args)` from the wrong inputs. It survives only while kind and path both match. The SDK header claimed "a function that is not listed below does not exist". `windmill-client` also exports the generated services, and the Python client exposes `Windmill.get`/`.post`, so an endpoint without a helper had no legal move. Each language now names its own escape hatch. `getAppInstructions` said the attached reference carries the TypeScript SDK even when `language: "python3"` had swapped in the Python one — on the very sentence telling the model to make that call. The kind-conversion comment claimed a hybrid runnable "silently runs stale inline code". It does not: `isRunnableByName`, `isRunnableByPath`, `convertPersistedToBackendRunnable` and `rawAppPolicy.processRunnable` all dispatch on `type` alone. The leftovers contradict the runnable's kind rather than override it, which is what the comment now says. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: address round-2 review nits on the raw-app AI instructions `flow is deployed` was satisfied both by "once the flow is deployed, the button works" and by a hallucinated "done — the flow is deployed", which eval mode makes impossible and the drafts-only judge cannot see. Every alternative now states an outstanding obligation, and two more real phrasings ("will need to be deployed") are accepted so a correct answer is not failed on wording. Condenses the three comment blocks that ran past the four-line limit in AGENTS.md, and drops two claims inside them that no longer hold: the `testRunAppRunnable` doc said it runs a runnable the way the app's own frontend does (it is the editor preview, which a deployed app's stored policy does not match), and `undeployedRunnableTargets` described its argument as the write tool's raw input when the call site passes the persisted runnable. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: report the real cause when a test run fails, and label the app-runnable card Driving `test_run_app_runnable` in a live session surfaced two defects the API-level check could not see. `executeTestRun` built its failure message from `error.message`, which the generated client leaves as the bare status text while the server's message sits in `body`. A path runnable aimed at an undeployed flow reported "Not Found" instead of "Not found: flow not found at name u/admin/current_time" — dropping the one diagnostic the run exists to produce. `formatToolError`, in the same file and written for exactly this, now does it. This also applies to test_run_script and test_run_flow, which had the same loss. The completion card read "Flow test completed successfully" for an app runnable, because `contextName` doubles as the jobs-tray kind and a path runnable pointing at a flow really does queue a flow job. A `completionName` override now names what ran without changing the kind. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * test: pin the deploy expectation against wrong answers, not just correct ones `deploying the flow` was satisfied by "done deploying the flow" — a deploy the agent only claims to have made, which eval mode makes impossible and the drafts-only judge cannot see. Replaced with the prospective forms, and dropped the same reading from the workflow variant. Three review rounds each found this same class of hole in the phrasing list, so the list is now exercised against the wrong answers themselves rather than eyeballed: naming the app as what needs deploying, claiming the deploy is already done, claiming to have deployed the flow, and saying nothing about deploying all have to fail, while four real correct phrasings have to pass. The test reads the case out of global.yaml, so a future edit to the alternatives is checked by it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * test: drop the tense-neutral deploy alternatives and cover completed claims A gerund after a preposition carries no tense, so `before`/`after`/`by deploying the flow` all match a deploy the agent only claims to have made ("after deploying the flow, I clicked the button and it returns the greeting") just as the bare gerund did. All three are gone rather than swapped for whichever reads least badly, and the two completed-deploy phrasings are now negative fixtures. The remaining alternatives are imperative or obligational, which a claim of having already deployed cannot satisfy. Condenses the two comments this list carries: the YAML block to four lines, and the test's rationale to the durable constraint about substring matching. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: encrypt sensitive inputs when test-running an app runnable `test_run_app_runnable` sent `force_viewer_static_fields` but not `force_viewer_sensitive_inputs`, which every other preview path derives from the runnable's `sensitive` user fields. That list is the only thing driving the encryption loop in apps.rs, so testing a runnable with a sensitive input wrote the real value into the job's args in plaintext, readable by anyone with run access to the workspace. Verified against a running EE instance. With the list, `api_key` is stored as `$encrypted:mvqtSRI9…` and the sentinel appears nowhere in the job record; without it, the sentinel is readable in run details. A non-sensitive field is left plaintext either way. The tool claims parity with the editor preview, so it uses that same filter (`type == 'user' && sensitive`) and omits the field entirely when empty. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-authored-by: Ruben Fiszel <ruben@windmill.dev>
27 KiB
Python SDK (wmill)
Import: import wmill
The client configures itself from the job's environment — base URL, token and credentials mode are all set before your code runs, so there is nothing to initialize and no reason to read WM_TOKEN or BASE_INTERNAL_URL and build an API URL yourself. Reconstructing that by hand only reintroduces details the client already handles. Call the SDK for anything Windmill, and use raw HTTP for third-party APIs.
The functions below are the surface to prefer. For an endpoint none of them covers, wmill.Windmill().get(endpoint) and .post(endpoint) issue an authenticated request against this instance. What does not exist is a function name you guessed at: if it is not listed below, do not call it.
To know who is running the script, read the contextual variables rather than calling the API:
os.environ.get("WM_END_USER_EMAIL") or os.environ.get("WM_EMAIL"). WM_END_USER_EMAIL is the app
viewer when the run was triggered from an app and empty otherwise (both variables are always
defined), WM_EMAIL is the user the job is permissioned as. WM_USERNAME is the matching username.
def worker_has_internal_server() -> bool
def get_mocked_api() -> Optional[dict]
Get the HTTP client instance.
Returns:
Configured httpx.Client for API requests
def get_client() -> httpx.Client
Make an HTTP GET request to the Windmill API.
Args:
endpoint: API endpoint path
raise_for_status: Whether to raise an exception on HTTP errors
**kwargs: Additional arguments passed to httpx.get
Returns:
HTTP response object
def get(endpoint, raise_for_status = True, **kwargs) -> httpx.Response
Make an HTTP POST request to the Windmill API.
Args:
endpoint: API endpoint path
raise_for_status: Whether to raise an exception on HTTP errors
**kwargs: Additional arguments passed to httpx.post
Returns:
HTTP response object
def post(endpoint, raise_for_status = True, **kwargs) -> httpx.Response
Create a new authentication token.
Args:
duration: Token validity duration (default: 1 day)
Returns:
New authentication token string
def create_token(duration = dt.timedelta(days=1)) -> str
Create a script job by path and return its job id.
def run_script_by_path_async(path: str, args: dict = None, scheduled_in_secs: int = None, tag: str = None) -> str
Create a script job by hash and return its job id.
def run_script_by_hash_async(hash_: str, args: dict = None, scheduled_in_secs: int = None, tag: str = None) -> str
Create a flow job and return its job id.
def run_flow_async(path: str, args: dict = None, scheduled_in_secs: int = None, do_not_track_in_parent: bool = True, tag: str = None) -> str
Run script by path synchronously and return its result.
def run_script_by_path(path: str, args: dict = None, timeout: dt.timedelta | int | float | None = None, verbose: bool = False, cleanup: bool = True, assert_result_is_not_none: bool = False, tag: str = None) -> Any
Run script by hash synchronously and return its result.
def run_script_by_hash(hash_: str, args: dict = None, timeout: dt.timedelta | int | float | None = None, verbose: bool = False, cleanup: bool = True, assert_result_is_not_none: bool = False, tag: str = None) -> Any
Run a script on the current worker without creating a job.
On agent workers (no internal server), falls back to running a normal
preview job and waiting for the result.
def run_inline_script_preview(content: str, language: str, args: dict = None) -> Any
Wait for a job to complete and return its result.
Args:
job_id: ID of the job to wait for
timeout: Maximum time to wait (seconds or timedelta)
verbose: Enable verbose logging
cleanup: Register cleanup handler to cancel job on exit
assert_result_is_not_none: Raise exception if result is None
Returns:
Job result when completed
Raises:
TimeoutError: If timeout is reached
Exception: If job fails
def wait_job(job_id, timeout: dt.timedelta | int | float | None = None, verbose: bool = False, cleanup: bool = True, assert_result_is_not_none: bool = False)
Cancel a specific job by ID.
Args:
job_id: UUID of the job to cancel
reason: Optional reason for cancellation
Returns:
Response message from the cancel endpoint
def cancel_job(job_id: str, reason: str = None) -> str
Cancel currently running executions of the same script.
def cancel_running() -> dict
Get job details by ID.
Args:
job_id: UUID of the job
Returns:
Job details dictionary
def get_job(job_id: str) -> dict
Get the root job ID for a flow hierarchy.
Args:
job_id: Job ID (defaults to current WM_JOB_ID)
Returns:
Root job ID
def get_root_job_id(job_id: str | None = None) -> dict
Get an OIDC JWT token for authentication to external services.
Args:
audience: Token audience (e.g., "vault", "aws")
expires_in: Optional expiration time in seconds
Returns:
JWT token string
def get_id_token(audience: str, expires_in: int | None = None) -> str
Get the status of a job.
Args:
job_id: UUID of the job
Returns:
Job status: "RUNNING", "WAITING", or "COMPLETED"
def get_job_status(job_id: str) -> JobStatus
Get the result of a completed job.
Args:
job_id: UUID of the completed job
assert_result_is_not_none: Raise exception if result is None
Returns:
Job result
def get_result(job_id: str, assert_result_is_not_none: bool = True) -> Any
Get a variable value by path.
Args:
path: Variable path in Windmill
Returns:
Variable value as string
def get_variable(path: str) -> str
Set a variable value by path, creating it if it doesn't exist.
Args:
path: Variable path in Windmill
value: Variable value to set
is_secret: Whether the variable should be secret (default: False)
def set_variable(path: str, value: str, is_secret: bool = False) -> None
Get a resource value by path.
Args:
path: Resource path in Windmill
none_if_undefined: Return None instead of raising if not found
interpolated: if variables and resources are fully unrolled
Returns:
Resource value dictionary or None
def get_resource(path: str, none_if_undefined: bool = False, interpolated: bool = True) -> dict | None
Set a resource value by path, creating it if it doesn't exist.
Args:
value: Resource value to set
path: Resource path in Windmill
resource_type: Resource type for creation
def set_resource(value: Any, path: str, resource_type: str)
List resources from Windmill workspace.
Args:
resource_type: Optional resource type to filter by (e.g., "postgresql", "mysql", "s3")
page: Optional page number for pagination
per_page: Optional number of results per page
Returns:
List of resource dictionaries
def list_resources(resource_type: str = None, page: int = None, per_page: int = None) -> list[dict]
Set the workflow state.
Args:
value: State value to set
path: Optional state resource path override.
def set_state(value: Any, path: str | None = None) -> None
Get the workflow state.
Args:
path: Optional state resource path override.
Returns:
State value or None if not set
def get_state(path: str | None = None) -> Any
Set job progress percentage (0-99).
Args:
value: Progress percentage
job_id: Job ID (defaults to current WM_JOB_ID)
def set_progress(value: int, job_id: Optional[str] = None)
Get job progress percentage.
Args:
job_id: Job ID (defaults to current WM_JOB_ID)
Returns:
Progress value (0-100) or None if not set
def get_progress(job_id: Optional[str] = None) -> Any
Set the user state of a flow at a given key
def set_flow_user_state(key: str, value: Any) -> None
Get the user state of a flow at a given key
def get_flow_user_state(key: str) -> Any
Get the Windmill server version.
Returns:
Version string
def version()
Convenient helpers that takes an S3 resource as input and returns the settings necessary to
initiate an S3 connection from DuckDB
def get_duckdb_connection_settings(s3_resource_path: str = '') -> DuckDbConnectionSettings | None
Convenient helpers that takes an S3 resource as input and returns the settings necessary to
initiate an S3 connection from Polars
def get_polars_connection_settings(s3_resource_path: str = '') -> PolarsConnectionSettings
Convenient helpers that takes an S3 resource as input and returns the settings necessary to
initiate an S3 connection using boto3
def get_boto3_connection_settings(s3_resource_path: str = '') -> Boto3ConnectionSettings
Load a file from the workspace s3 bucket and returns its content as bytes.
'''python
from wmill import S3Object
s3_obj = S3Object(s3="/path/to/my_file.txt")
my_obj_content = client.load_s3_file(s3_obj)
file_content = my_obj_content.decode("utf-8")
'''
def load_s3_file(s3object: S3Object | str, s3_resource_path: str | None) -> bytes
Load a file from the workspace s3 bucket and returns the bytes stream.
'''python
from wmill import S3Object
s3_obj = S3Object(s3="/path/to/my_file.txt")
with wmill.load_s3_file_reader(s3object, s3_resource_path) as file_reader:
print(file_reader.read())
'''
def load_s3_file_reader(s3object: S3Object | str, s3_resource_path: str | None) -> BufferedReader
Write a file to the workspace S3 bucket
'''python
from wmill import S3Object
s3_obj = S3Object(s3="/path/to/my_file.txt")
# for an in memory bytes array:
file_content = b'Hello Windmill!'
client.write_s3_file(s3_obj, file_content)
# for a file:
with open("my_file.txt", "rb") as my_file:
client.write_s3_file(s3_obj, my_file)
'''
def write_s3_file(s3object: S3Object | str | None, file_content: BufferedReader | bytes, s3_resource_path: str | None, content_type: str | None = None, content_disposition: str | None = None) -> S3Object
Permanently delete a file from the workspace S3 bucket.
'''python
from wmill import S3Object
s3_obj = S3Object(s3="/path/to/my_file.txt")
client.delete_s3_object(s3_obj)
'''
def delete_s3_object(s3object: S3Object | str, s3_resource_path: str | None = None) -> None
Sign S3 objects for use by anonymous users in public apps.
Args:
s3_objects: List of S3 objects to sign
Returns:
List of signed S3 objects
def sign_s3_objects(s3_objects: list[S3Object | str]) -> list[S3Object]
Sign a single S3 object for use by anonymous users in public apps.
Args:
s3_object: S3 object to sign
Returns:
Signed S3 object
def sign_s3_object(s3_object: S3Object | str) -> S3Object
Generate presigned public URLs for an array of S3 objects.
If an S3 object is not signed yet, it will be signed first.
Args:
s3_objects: List of S3 objects to sign
base_url: Optional base URL for the presigned URLs (defaults to WM_BASE_URL)
Returns:
List of signed public URLs
Example:
>>> s3_objs = [S3Object(s3="/path/to/file1.txt"), S3Object(s3="/path/to/file2.txt")]
>>> urls = client.get_presigned_s3_public_urls(s3_objs)
def get_presigned_s3_public_urls(s3_objects: list[S3Object | str], base_url: str | None = None) -> list[str]
Generate a presigned public URL for an S3 object.
If the S3 object is not signed yet, it will be signed first.
Args:
s3_object: S3 object to sign
base_url: Optional base URL for the presigned URL (defaults to WM_BASE_URL)
Returns:
Signed public URL
Example:
>>> s3_obj = S3Object(s3="/path/to/file.txt")
>>> url = client.get_presigned_s3_public_url(s3_obj)
def get_presigned_s3_public_url(s3_object: S3Object | str, base_url: str | None = None) -> str
Get the current user information.
Returns:
User details dictionary
def whoami() -> dict
Get the current user information (alias for whoami).
Returns:
User details dictionary
def user() -> dict
Get the state resource path from environment.
Returns:
State path string
def state_path() -> str
Get the workflow state.
Returns:
State value or None if not set
def state() -> Any
Set the state in the shared folder using pickle
def set_shared_state_pickle(value: Any, path: str = 'state.pickle') -> None
Get the state in the shared folder using pickle
def get_shared_state_pickle(path: str = 'state.pickle') -> Any
Set the state in the shared folder using pickle
def set_shared_state(value: Any, path: str = 'state.json') -> None
Get the state in the shared folder using pickle
def get_shared_state(path: str = 'state.json') -> None
Get URLs needed for resuming a flow after suspension.
Args:
approver: Optional approver name
flow_level: 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:
Dictionary with approvalPage, resume, and cancel URLs
def get_resume_urls(approver: str = None, flow_level: bool = None) -> dict
Get the resume URLs bound to one wait_for_approval step of this workflow.
Args:
step_key: Checkpoint key of the approval step, as passed to
wait_for_approval(key=...)
approver: Optional approver name
Returns:
Dictionary with approvalPage, resume, and cancel URLs
def get_approval_urls(step_key: str = 'approval', approver: str = None) -> dict
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, use the "Advanced -> Suspend -> Form" functionality.
Learn more at: https://www.windmill.dev/docs/flows/flow_approval#form
:param slack_resource_path: The path to the Slack resource in Windmill.
:type slack_resource_path: str
:param channel_id: The Slack channel ID where the approval request will be sent.
:type channel_id: str
:param message: Optional custom message to include in the Slack approval request.
:type message: str, optional
:param approver: Optional user ID or name of the approver for the request.
:type approver: str, optional
:param default_args_json: Optional dictionary defining or overriding the default arguments for form fields.
:type default_args_json: dict, optional
:param dynamic_enums_json: Optional dictionary overriding the enum default values of enum form fields.
:type dynamic_enums_json: dict, optional
:raises Exception: If the function is not called within a flow or flow preview.
:raises Exception: If the required flow job or flow step environment variables are not set.
:return: None
Usage Example:
>>> client.request_interactive_slack_approval(
... slack_resource_path="/u/alex/my_slack_resource",
... channel_id="admins-slack-channel",
... message="Please approve this request",
... approver="approver123",
... default_args_json={"key1": "value1", "key2": 42},
... dynamic_enums_json={"foo": ["choice1", "choice2"], "bar": ["optionA", "optionB"]},
... )
Notes:
- This function must be executed within a Windmill flow or flow preview.
- The function checks for required environment variables (WM_FLOW_JOB_ID, WM_FLOW_STEP_ID) to ensure it is run in the appropriate context.
def request_interactive_slack_approval(slack_resource_path: str, channel_id: str, message: str = None, approver: str = None, default_args_json: dict = None, dynamic_enums_json: dict = None) -> None
Send a message to a Microsoft Teams conversation with conversation_id, where success is used to style the message
def send_teams_message(conversation_id: str, text: str, success: bool = True, card_block: dict = None)
Get a DataTable client for SQL queries.
Args:
name: Database name (default: "main")
Returns:
DataTableClient instance
def datatable(name: str = 'main')
Get a DuckLake client for DuckDB queries.
Args:
name: Database name (default: "main")
Returns:
DucklakeClient instance
def ducklake(name: str = 'main')
def init_global_client(f)
def deprecate(in_favor_of: str)
Get the current workspace ID.
Returns:
Workspace ID string
def get_workspace() -> str
def get_version() -> str
Create a script job and return its job ID.
Args:
hash_or_path: Script hash or path (determined by presence of '/')
args: Script arguments
scheduled_in_secs: Delay before execution in seconds
tag: Override the worker tag the job runs on
Returns:
Job ID string
def run_script_async(hash_or_path: str, args: Dict[str, Any] = None, scheduled_in_secs: int = None, tag: str = None) -> str
Run a script synchronously by hash and return its result.
Args:
hash: Script hash
args: Script arguments
verbose: Enable verbose logging
assert_result_is_not_none: Raise exception if result is None
cleanup: Register cleanup handler to cancel job on exit
timeout: Maximum time to wait
tag: Override the worker tag the job runs on
Returns:
Script result
def run_script_sync(hash: str, args: Dict[str, Any] = None, verbose: bool = False, assert_result_is_not_none: bool = True, cleanup: bool = True, timeout: dt.timedelta = None, tag: str = None) -> Any
Run a script synchronously by path and return its result.
Args:
path: Script path
args: Script arguments
verbose: Enable verbose logging
assert_result_is_not_none: Raise exception if result is None
cleanup: Register cleanup handler to cancel job on exit
timeout: Maximum time to wait
tag: Override the worker tag the job runs on
Returns:
Script result
def run_script_by_path_sync(path: str, args: Dict[str, Any] = None, verbose: bool = False, assert_result_is_not_none: bool = True, cleanup: bool = True, timeout: dt.timedelta = None, tag: str = None) -> Any
Convenient helpers that takes an S3 resource as input and returns the settings necessary to
initiate an S3 connection from DuckDB
def duckdb_connection_settings(s3_resource_path: str = '') -> DuckDbConnectionSettings
Convenient helpers that takes an S3 resource as input and returns the settings necessary to
initiate an S3 connection from Polars
def polars_connection_settings(s3_resource_path: str = '') -> PolarsConnectionSettings
Convenient helpers that takes an S3 resource as input and returns the settings necessary to
initiate an S3 connection using boto3
def boto3_connection_settings(s3_resource_path: str = '') -> Boto3ConnectionSettings
Get the state resource path from environment.
Returns:
State path string
def get_state_path() -> str
Parse resource syntax from string.
def parse_resource_syntax(s: str) -> Optional[str]
Parse S3 object from a s3://<storage>/<key> URI string (s3:///<key>
for the default storage) or S3Object format. Any other string raises
rather than falling back to an auto-generated key: an auto key is
requested by omitting the object, and a fallback would silently misplace
the upload on any typo.
def parse_s3_object(s3_object: S3Object | str) -> S3Object
Parse variable syntax from string.
def parse_variable_syntax(s: str) -> Optional[str]
Append a text to the result stream.
Args:
text: text to append to the result stream
def append_to_result_stream(text: str) -> None
Stream to the result stream.
Args:
stream: stream to stream to the result stream
def stream_result(stream) -> None
Execute a SQL query against the DataTable.
Args:
sql: SQL query string with $1, $2, etc. placeholders
*args: Positional arguments to bind to query placeholders
Returns:
SqlQuery instance for fetching results
def query(sql: str, *args) -> SqlQuery
Idempotently materialize the rows of select_sql into ducklake
table for one partition (or the whole table when partition is
None). Client-side equivalent of the // materialize engine: with
unique_key it upserts within the slice (delete-by-key + insert);
without it, it replaces (whole table → CREATE OR REPLACE; partition →
delete the partition + insert). Re-running the same slice is safe — the
backfill / failure-recovery contract.
The partition value is bound as a DuckDB arg (never string-interpolated)
so it cannot inject SQL. select_sql is trusted (your own query).
def upsert_partition(table: str, select_sql: str, partition: str = None, unique_key: str = None, partition_col: str = '_wm_partition', schema: str = None)
INSERT-only materialization (no dedup / no replace) for an immutable
event-log table — for one partition, or the whole table when
partition is None. NOTE: unlike upsert_partition, re-running the same
slice duplicates rows — use only for append-only sources.
def append_partition(table: str, select_sql: str, partition: str = None, partition_col: str = '_wm_partition', schema: str = None)
Read a materialized ducklake table, optionally a single partition.
def read(table: str, partition: str = None, partition_col: str = '_wm_partition', schema: str = None)
Execute query and fetch results.
Args:
result_collection: Optional result collection mode
Returns:
Query results
def fetch(result_collection: str | None = None)
Execute query and fetch first row of results.
Returns:
First row of query results
def fetch_one()
Execute query and fetch first row of results. Return result as a scalar value.
Returns:
First row of query result as a scalar value
def fetch_one_scalar()
Execute query and don't return any results.
def execute()
DuckDB executor requires explicit argument types at declaration
These types exist in both DuckDB and Postgres
Check that the types exist if you plan to extend this function for other SQL engines.
def infer_sql_type(value) -> str
def parse_sql_client_name(name: str) -> tuple[str, Optional[str]]
Decorator that marks a function as a workflow task.
Works in both WAC v1 (sync, HTTP-based dispatch) and WAC v2
(async, checkpoint/replay) modes:
- v2 (inside @workflow): dispatches as a checkpoint step.
- v1 (WM_JOB_ID set, no @workflow): dispatches via HTTP API.
- Standalone: executes the function body directly.
A task runs as its own job, so its result is always encoded as JSON and
decoded back before the caller sees it: a datetime comes back as a
string, a tuple as a list.
Usage::
@task
async def extract_data(url: str): ...
@task(path="f/external_script", timeout=600, tag="gpu")
async def run_external(x: int): ...
def task(_func = None, path: Optional[str] = None, tag: Optional[str] = None, timeout: Optional[int] = None, cache_ttl: Optional[int] = None, priority: Optional[int] = None, concurrency_limit: Optional[int] = None, concurrency_key: Optional[str] = None, concurrency_time_window_s: Optional[int] = None)
Create a task that dispatches to a separate Windmill script.
Usage::
extract = task_script("f/data/extract", timeout=600)
@workflow
async def main():
data = await extract(url="https://...")
def task_script(path: str, timeout: Optional[int] = None, tag: Optional[str] = None, cache_ttl: Optional[int] = None, priority: Optional[int] = None, concurrency_limit: Optional[int] = None, concurrency_key: Optional[str] = None, concurrency_time_window_s: Optional[int] = None)
Create a task that dispatches to a separate Windmill flow.
Usage::
pipeline = task_flow("f/etl/pipeline", priority=10)
@workflow
async def main():
result = await pipeline(input=data)
def task_flow(path: str, timeout: Optional[int] = None, tag: Optional[str] = None, cache_ttl: Optional[int] = None, priority: Optional[int] = None, concurrency_limit: Optional[int] = None, concurrency_key: Optional[str] = None, concurrency_time_window_s: Optional[int] = None)
Decorator marking 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.
def workflow(func)
Execute fn inline and checkpoint the result.
On replay the cached value is returned without re-executing fn.
Use for lightweight deterministic operations (timestamps, random IDs,
config reads) that should not incur the overhead of a child job.
fn's result is encoded as JSON and decoded back before it is returned,
so the round that runs the body sees the same types every replay sees:
a datetime comes back as a string, a tuple as a list.
async def step(name: str, fn)
Server-side sleep — suspend the workflow for the given duration without holding a worker.
Inside a @workflow, the parent job suspends and auto-resumes after seconds.
Outside a workflow, falls back to asyncio.sleep.
async def sleep(seconds: int)
Suspend the workflow and wait for an external approval.
Pass key to name the step, then get_approval_urls(key) yields the URLs
that resume exactly this approval — route them through your own channel.
Without a key the steps are named approval, approval_2, ...
Returns a dict with value (form data), approver, and approved.
Args:
timeout: Approval timeout in seconds (default 1800).
form: Optional form schema for the approval page.
self_approval: Whether the user who triggered the flow can approve it (default True).
key: Optional checkpoint key naming this approval step.
Example::
urls = await step("urls", lambda: get_approval_urls("manager"))
await step("notify", lambda: send_email(urls["resume"], urls["cancel"]))
result = await wait_for_approval(key="manager", timeout=3600)
async def wait_for_approval(timeout: int = 1800, form: dict | None = None, self_approval: bool = True, key: str | None = None) -> dict
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::
@task
async def process(item: str):
...
results = await parallel(items, process, concurrency=5)
async def parallel(items, fn, concurrency: Optional[int] = None)
Commit Kafka offsets for a trigger with auto_commit disabled.
Args:
trigger_path: Path to the Kafka trigger (from event['wm_trigger']['trigger_path'])
topic: Kafka topic name (from event['topic'])
partition: Partition number (from event['partition'])
offset: Message offset to commit (from event['offset'])
def commit_kafka_offsets(trigger_path: str, topic: str, partition: int, offset: int) -> None