Files
windmill/backend/windmill-api-flows/src/flows.rs
T
Ruben Fiszel 76a9523009 feat: use derived username instead of email for non-member superadmins (#9857)
* feat: use derived username instead of email for non-member superadmins

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: address review - drop redundant username cache, guard whoami membership by email

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* refactor: use explicit non_member boolean instead of role string for superadmin banner

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: resolve email from password table for non-member superadmin permissioned_as

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: resolve non-member superadmin drafts via shared username->email resolver

Adds resolve_username_to_email (usr, then super_admin password fallback for both derived-username and email modes) and uses it in get_email_from_permissioned_as and the drafts get/list endpoints, so a non-member superadmin's drafts resolve and no email leaks into the drafts payload.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* test: superadmin-not-in-workspace schedule uses derived username as permissioned_as

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: resolve non-member superadmin identity in draft owner-circles, username_to_email, and home filter

Applies the password-fallback username resolution to the script/flow/app/draft owner-circle subqueries and the username_to_email endpoint (was an admins-workspace 'username == email' hack), and switches the home items-list user-folder filter to the non_member flag instead of the now-broken username-contains-@ heuristic.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: backfill non-member superadmin favorites from email to derived username

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: propagate DB errors in username resolution instead of leaking email (CI review)

Addresses cubic-dev-ai P2: get_instance_username_or_fallback_to_email now returns Result and only falls back to the email for a genuine 'no derived username'; a query error propagates so callers fail closed rather than leaking the raw email as the acting username.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: clarify non-member superadmin popover (username used + admin permissions)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: keep username_to_email endpoint member-only to not disclose non-member superadmin email (CI review)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix: forbid disabling automate_username_creation once usernames assigned (CI review)

Makes the setting effectively one-way once instance-wide usernames exist, so the global-uniqueness invariant that keeps stored u/<username> identities (schedules/triggers/drafts/superadmin ownership) unambiguous can never be dropped back to workspace-local uniqueness. Re-saving false on an already-disabled instance stays a no-op.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-01 14:54:07 +00:00

2256 lines
70 KiB
Rust

/*
* Author: Ruben Fiszel
* Copyright: Windmill Labs, Inc 2022
* This file and its contents are licensed under the AGPLv3 License.
* Please see the included NOTICE for copyright information and
* LICENSE-AGPL for a copy of the license.
*/
use std::collections::HashMap;
use axum::response::IntoResponse;
use axum::{
extract::{Extension, Path, Query},
routing::{delete, get, post},
Json, Router,
};
use windmill_api_auth::{
auth::{list_tokens_internal, TruncatedTokenWithEmail},
build_scope_path_predicate, check_scopes, maybe_refresh_folders, require_owner_of_path,
ApiAuthed,
};
use windmill_common::workspaces::{check_deploy_rules, RuleCheckResult};
use windmill_common::{
user_drafts::{overlay_or_draft_only, DraftUserRef, UserDraftItemKind, WithDraftOverlay},
utils::HTTP_CLIENT,
webhook::{WebhookMessage, WebhookShared},
DB,
};
use windmill_queue::schedule::clear_schedule;
use hyper::StatusCode;
use serde::{Deserialize, Serialize};
use sql_builder::prelude::*;
use sqlx::{FromRow, Postgres, Transaction};
use windmill_audit::audit_oss::{audit_log, AuditAuthorable};
use windmill_audit::ActionKind;
use windmill_common::assets::{clear_static_asset_usage, AssetUsageKind};
use windmill_common::flows::FlowModule;
use windmill_common::min_version::{
MIN_VERSION_SUPPORTS_DEBOUNCING, MIN_VERSION_SUPPORTS_DEBOUNCING_V2,
MIN_VERSION_SUPPORTS_NODE_DEBOUNCING,
};
use windmill_common::runnable_settings::RunnableSettingsTrait;
use windmill_common::utils::query_elems_from_hub;
use windmill_common::worker::{to_raw_value, CLOUD_HOSTED};
use windmill_common::HUB_BASE_URL;
use windmill_common::{
db::UserDB,
error::{self, to_anyhow, Error, JsonResult, Result},
flows::{Flow, FlowWithStarred, ListFlowQuery, ListableFlow, NewFlow},
jobs::JobPayload,
schedule::Schedule,
utils::{http_get_from_hub, not_found_if_none, paginate, Pagination, RunnableKind, StripPath},
};
use windmill_dep_map::scoped_dependency_map::ScopedDependencyMap;
use windmill_git_sync::{handle_deployment_metadata, DeployedObject};
use windmill_queue::WMDEBUG_FORCE_NO_LEGACY_DEBOUNCING_COMPAT;
use windmill_queue::{push, schedule::push_scheduled_job, PushIsolationLevel};
pub fn workspaced_service() -> Router {
Router::new()
.route("/list", get(list_flows))
.route("/list_search", get(list_search_flows))
.route("/create", post(create_flow))
.route("/update/{*path}", post(update_flow))
.route("/archive/{*path}", post(archive_flow_by_path))
.route("/delete/{*path}", delete(delete_flow_by_path))
.route("/list_tokens/{*path}", get(list_tokens))
.route("/get/{*path}", get(get_flow_by_path))
.route("/deployment_status/p/{*path}", get(get_deployment_status))
.route("/exists/{*path}", get(exists_flow_by_path))
.route("/list_paths", get(list_paths))
.route("/history/p/{*path}", get(get_flow_history))
.route("/get_latest_version/{*path}", get(get_latest_version))
.route(
"/list_paths_from_workspace_runnable/{runnable_kind}/{*path}",
get(list_paths_from_workspace_runnable),
)
.route("/history_update/v/{version}", post(update_flow_history))
.route("/get/v/{version}", get(get_flow_version_by_id))
.route("/get/v/{version}/p/{*path}", get(get_flow_version))
.route(
"/toggle_workspace_error_handler/{*path}",
post(toggle_workspace_error_handler),
)
}
pub fn global_service() -> Router {
Router::new()
.route("/hub/list", get(list_hub_flows))
.route("/hub/get/{id}", get(get_hub_flow_by_id))
}
#[derive(Serialize, FromRow)]
pub struct SearchFlow {
path: String,
value: sqlx::types::Json<Box<serde_json::value::RawValue>>,
}
async fn list_search_flows(
authed: ApiAuthed,
Path(w_id): Path<String>,
Extension(user_db): Extension<UserDB>,
) -> JsonResult<Vec<SearchFlow>> {
#[cfg(feature = "enterprise")]
let n = 1000;
#[cfg(not(feature = "enterprise"))]
let n = 3;
let mut tx = user_db.begin(&authed).await?;
let allowed = build_scope_path_predicate(&authed, "flows", "read");
let rows = sqlx::query_as::<_, SearchFlow>(
"SELECT flow.path, flow_version.value
FROM flow
LEFT JOIN flow_version ON flow_version.id = flow.versions[array_upper(flow.versions, 1)]
WHERE flow.workspace_id = $1 LIMIT $2",
)
.bind(&w_id)
.bind(n)
.fetch_all(&mut *tx)
.await?
.into_iter()
.filter(|r| allowed(&r.path))
.collect::<Vec<_>>();
tx.commit().await?;
Ok(Json(rows))
}
async fn list_flows(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Extension(db): Extension<DB>,
Path(w_id): Path<String>,
Query(pagination): Query<Pagination>,
Query(lq): Query<ListFlowQuery>,
) -> JsonResult<Vec<ListableFlow>> {
let (per_page, offset) = paginate(pagination);
let mut sqlb = SqlBuilder::select_from("flow as o")
.fields(&[
"o.workspace_id",
"o.path",
"summary",
if !lq.without_description.unwrap_or(false) {
"description"
} else {
"NULL as description"
},
"fv.created_by as edited_by",
"fv.created_at as edited_at",
"archived",
"extra_perms",
"favorite.path IS NOT NULL as starred",
"ws_error_handler_muted",
"o.labels",
"draft.email IS NOT NULL as is_draft",
// Per-path draft owners as a JSON array; see scripts.rs for the rationale
// (non-member superadmin identity fallback via `password`, legacy NULL-email row).
"(SELECT json_agg(json_build_object('username', COALESCE(u.username, p.username, CASE WHEN p.email IS NOT NULL THEN d.email END)) ORDER BY COALESCE(u.username, p.username, CASE WHEN p.email IS NOT NULL THEN d.email END) NULLS LAST) \
FROM draft d \
LEFT JOIN usr u ON u.workspace_id = d.workspace_id AND u.email = d.email \
LEFT JOIN password p ON p.email = d.email AND p.super_admin = true \
WHERE d.workspace_id = o.workspace_id AND d.path = o.path AND d.typ = 'flow') as draft_users",
"folder_labels(o.workspace_id, o.path) as inherited_labels"
])
.left()
.join("favorite")
.on(
"favorite.favorite_kind = 'flow' AND favorite.workspace_id = o.workspace_id AND favorite.path = o.path AND favorite.usr = ?"
.bind(&authed.username),
)
.left()
.join("draft")
.on(
"draft.path = o.path AND draft.workspace_id = o.workspace_id AND draft.typ = 'flow' AND draft.email = ?"
.bind(&authed.email),
)
.left()
.join("flow_version fv")
.on(
"fv.id = o.versions[array_upper(o.versions, 1)]"
)
.order_desc("favorite.path IS NOT NULL")
.order_by("fv.created_at", lq.order_desc.unwrap_or(true))
.and_where("o.workspace_id = ?".bind(&w_id))
.offset(offset)
.limit(per_page)
.clone();
sqlb.and_where_eq("archived", lq.show_archived.unwrap_or(false));
if let Some(ps) = &lq.path_start {
sqlb.and_where_like_left("o.path", ps);
}
if let Some(p) = &lq.path_exact {
sqlb.and_where_eq("o.path", "?".bind(p));
}
if let Some(cb) = &lq.edited_by {
sqlb.and_where_eq("fv.created_by", "?".bind(cb));
}
if lq.starred_only.unwrap_or(false) {
sqlb.and_where_is_not_null("favorite.path");
}
if let Some(dw) = &lq.dedicated_worker {
sqlb.and_where_eq("dedicated_worker", dw);
}
if let Some(label) = &lq.label {
for l in label.split(',') {
sqlb.and_where(
"(o.labels @> ARRAY[?] OR folder_labels(o.workspace_id, o.path) @> ARRAY[?])"
.bind(&l.trim())
.bind(&l.trim()),
);
}
}
if lq.with_deployment_msg.unwrap_or(false) {
sqlb.join("deployment_metadata dm")
.left()
.on("dm.flow_version = o.versions[array_upper(o.versions, 1)]")
.fields(&["dm.deployment_msg"]);
}
let sql = sqlb.sql().map_err(|e| Error::internal_err(e.to_string()))?;
let mut tx = user_db.begin(&authed).await?;
let allowed = build_scope_path_predicate(&authed, "flows", "read");
let mut rows = sqlx::query_as::<_, ListableFlow>(&sql)
.fetch_all(&mut *tx)
.await?
.into_iter()
.filter(|r| allowed(&r.path))
.collect::<Vec<_>>();
tx.commit().await?;
// Append the authed user's drafts at paths with no deployed flow; see scripts.rs.
if lq.include_draft_only.unwrap_or(false)
&& !authed.is_operator
&& offset == 0
&& lq.path_start.is_none()
&& lq.path_exact.is_none()
&& lq.edited_by.is_none()
&& lq.dedicated_worker.is_none()
&& lq.label.is_none()
&& !lq.starred_only.unwrap_or(false)
&& !lq.show_archived.unwrap_or(false)
{
// `(email = $2 OR email IS NULL)` + `DISTINCT ON (path)` ordered NULL-last; see scripts.rs.
let draft_only_rows = sqlx::query!(
r#"SELECT DISTINCT ON (path)
path,
value as "value!: sqlx::types::Json<Box<serde_json::value::RawValue>>",
created_at
FROM draft
WHERE workspace_id = $1
AND typ = 'flow'
AND (email = $2 OR email IS NULL)
AND NOT EXISTS (
SELECT 1 FROM flow f
WHERE f.workspace_id = draft.workspace_id
AND f.path = draft.path
)
ORDER BY path, (email IS NULL)"#,
&w_id,
&authed.email,
)
.fetch_all(&db)
.await?;
for row in draft_only_rows {
let v: serde_json::Value =
serde_json::from_str(row.value.0.get()).unwrap_or(serde_json::Value::Null);
// The Path widget binds `$pathStore` one-way (`flow.path → $pathStore`),
// so the editor writes a separate `draft_path` field only when the typed
// path differs from the deployed one. `None` = unchanged.
let draft_path = v
.get("draft_path")
.and_then(|s| s.as_str())
.filter(|s| !s.is_empty() && *s != row.path.as_str())
.map(|s| s.to_string());
rows.push(ListableFlow {
workspace_id: w_id.clone(),
path: row.path,
summary: v
.get("summary")
.and_then(|s| s.as_str())
.unwrap_or("")
.to_string(),
description: v
.get("description")
.and_then(|s| s.as_str())
.map(|s| s.to_string()),
edited_by: Some(authed.email.clone()),
edited_at: Some(row.created_at),
archived: false,
extra_perms: serde_json::Value::Object(serde_json::Map::new()),
starred: false,
draft_only: Some(true),
ws_error_handler_muted: None,
deployment_msg: None,
labels: None,
// No deployed row to inherit folder labels from.
inherited_labels: None,
is_draft: true,
draft_path,
// Synthesized rows are the authed user's own draft.
draft_users: Some(sqlx::types::Json(vec![DraftUserRef {
username: Some(authed.username.clone()),
}])),
});
}
}
Ok(Json(rows))
}
async fn list_hub_flows(Extension(db): Extension<DB>) -> impl IntoResponse {
let (status_code, headers, response) = query_elems_from_hub(
&HTTP_CLIENT,
&format!("{}/searchFlowData?approved=true", **HUB_BASE_URL.load()),
None,
&db,
)
.await?;
Ok::<_, Error>((status_code, headers, response))
}
async fn list_paths(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path(w_id): Path<String>,
) -> JsonResult<Vec<String>> {
let mut tx = user_db.begin(&authed).await?;
let flows = sqlx::query_scalar!(
"SELECT distinct(path) FROM flow WHERE workspace_id = $1",
w_id
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
Ok(Json(flows))
}
pub async fn get_hub_flow_by_id(
Path(id): Path<i32>,
Extension(db): Extension<DB>,
) -> JsonResult<Box<serde_json::value::RawValue>> {
let value = http_get_from_hub(
&HTTP_CLIENT,
&format!("{}/flows/{}/json", **HUB_BASE_URL.load(), id),
false,
None,
Some(&db),
)
.await?
.json()
.await
.map_err(to_anyhow)?;
Ok(Json(value))
}
#[derive(Deserialize)]
pub struct ToggleWorkspaceErrorHandler {
#[cfg(feature = "enterprise")]
pub muted: Option<bool>,
}
#[cfg(not(feature = "enterprise"))]
async fn toggle_workspace_error_handler(
_authed: ApiAuthed,
Extension(_user_db): Extension<UserDB>,
Path((_w_id, _path)): Path<(String, StripPath)>,
Json(_req): Json<ToggleWorkspaceErrorHandler>,
) -> Result<String> {
return Err(Error::BadRequest(
"Muting the error handler for certain flow is only available in enterprise version"
.to_string(),
));
}
#[cfg(feature = "enterprise")]
async fn toggle_workspace_error_handler(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
Json(req): Json<ToggleWorkspaceErrorHandler>,
) -> Result<String> {
let mut tx = user_db.begin(&authed).await?;
let error_handler_maybe: Option<String> = sqlx::query_scalar!(
r#"
SELECT
error_handler->>'path'
FROM
workspace_settings
WHERE
workspace_id = $1
"#,
w_id
)
.fetch_optional(&mut *tx)
.await?
.unwrap_or(None);
let response = match error_handler_maybe {
Some(_) => {
sqlx::query_scalar!(
r#"
UPDATE
flow
SET
ws_error_handler_muted = $3
WHERE
path = $1 AND
workspace_id = $2
"#,
path.to_path(),
w_id,
req.muted,
)
.execute(&mut *tx)
.await?;
Ok("".to_string())
}
None => Err(Error::BadRequest(
"Workspace error handler needs to be defined".to_string(),
)),
};
tx.commit().await?;
return response;
}
async fn check_path_conflict<'c>(
tx: &mut Transaction<'c, Postgres>,
w_id: &str,
path: &str,
) -> Result<()> {
let exists = sqlx::query_scalar!(
"SELECT EXISTS(SELECT 1 FROM flow WHERE path = $1 AND workspace_id = $2)",
path,
w_id
)
.fetch_one(&mut **tx)
.await?
.unwrap_or(false);
if exists {
return Err(Error::BadRequest(format!("Flow {} already exists", path)));
}
return Ok(());
}
#[derive(Deserialize)]
struct ListPathsFromWorkspaceRunnableQuery {
match_path_start: Option<bool>,
}
async fn list_paths_from_workspace_runnable(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, runnable_kind, path)): Path<(String, RunnableKind, StripPath)>,
Query(query): Query<ListPathsFromWorkspaceRunnableQuery>,
) -> JsonResult<Vec<String>> {
let path = path.to_path();
check_scopes(&authed, || {
format!("flows:read:{}", format!("{}/{}", runnable_kind, path))
})?;
let mut tx = user_db.begin(&authed).await?;
let runnables = if query.match_path_start.unwrap_or(false) {
sqlx::query_scalar!(
r#"SELECT DISTINCT f.path
FROM workspace_runnable_dependencies wru
JOIN flow f
ON wru.flow_path = f.path AND wru.workspace_id = f.workspace_id
WHERE wru.runnable_path LIKE $1 || '%' AND wru.runnable_is_flow = $2 AND wru.workspace_id = $3"#,
path,
matches!(runnable_kind, RunnableKind::Flow),
w_id
)
.fetch_all(&mut *tx)
.await?
} else {
sqlx::query_scalar!(
r#"SELECT f.path
FROM workspace_runnable_dependencies wru
JOIN flow f
ON wru.flow_path = f.path AND wru.workspace_id = f.workspace_id
WHERE wru.runnable_path = $1 AND wru.runnable_is_flow = $2 AND wru.workspace_id = $3"#,
path,
matches!(runnable_kind, RunnableKind::Flow),
w_id
)
.fetch_all(&mut *tx)
.await?
};
tx.commit().await?;
Ok(Json(runnables))
}
async fn validate_flow(new_flow: &NewFlow) -> error::Result<()> {
#[cfg(not(feature = "enterprise"))]
if new_flow.ws_error_handler_muted.is_some_and(|val| val) {
return Err(Error::BadRequest(
"Muting the error handler for certain flow is only available in enterprise version"
.to_string(),
));
}
guard_flow_from_debounce_data(new_flow).await?;
return Ok(());
}
async fn create_flow(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Extension(user_db): Extension<UserDB>,
Extension(webhook): Extension<WebhookShared>,
Path(w_id): Path<String>,
Json(mut nf): Json<NewFlow>,
) -> Result<(StatusCode, String)> {
if authed.is_operator {
return Err(Error::NotAuthorized(
"Operators cannot create flows for security reasons".to_string(),
));
}
check_scopes(&authed, || format!("flows:write:{}", nf.path))?;
if let RuleCheckResult::Blocked(msg) = check_deploy_rules(
&w_id,
AuditAuthorable::username(&authed),
&authed.groups,
authed.is_admin,
&db,
)
.await?
{
return Err(Error::PermissionDenied(msg));
}
validate_flow(&nf).await?;
if *CLOUD_HOSTED {
let nb_flows =
sqlx::query_scalar!("SELECT COUNT(*) FROM flow WHERE workspace_id = $1", &w_id)
.fetch_one(&db)
.await?;
if nb_flows.unwrap_or(0) >= 1000 {
return Err(Error::BadRequest(
"You have reached the maximum number of flows (1000) on cloud. Check your usage in Workspace Settings > General > Cloud Quotas. Contact support@windmill.dev to increase the limit"
.to_string(),
));
}
if nf.summary.len() > 300 {
return Err(Error::BadRequest(
"Summary must be less than 300 characters on cloud".to_string(),
));
}
if nf
.description
.as_ref()
.is_some_and(|desc| desc.len() > 3000)
{
return Err(Error::BadRequest(
"Description must be less than 3000 characters on cloud".to_string(),
));
}
}
// cron::Schedule::from_str(&ns.schedule).map_err(|e| error::Error::BadRequest(e.to_string()))?;
let authed = maybe_refresh_folders(&nf.path, &w_id, authed, &db).await;
// Apply folder default_permissioned_as on create when the caller did not
// explicitly preserve a value and the user can preserve.
let explicit_preserve = nf.on_behalf_of_email.is_some()
&& nf.preserve_on_behalf_of.unwrap_or(false)
&& windmill_common::can_preserve_on_behalf_of(&authed);
if !explicit_preserve && windmill_common::can_preserve_on_behalf_of(&authed) {
if let Some(default_email) =
windmill_common::folders::resolve_folder_default_on_behalf_of_email(
&db, &w_id, &nf.path,
)
.await?
{
nf.on_behalf_of_email = Some(default_email);
nf.preserve_on_behalf_of = Some(true);
}
}
let mut tx = user_db.clone().begin(&authed).await?;
check_path_conflict(&mut tx, &w_id, &nf.path).await?;
check_schedule_conflict(&mut tx, &w_id, &nf.path).await?;
let schema_str = nf.schema.and_then(|x| serde_json::to_string(&x.0).ok());
sqlx::query!(
r#"INSERT INTO flow (
workspace_id, path, summary, description,
dependency_job, lock_error_logs, tag,
dedicated_worker, visible_to_runner_only, on_behalf_of_email,
ws_error_handler_muted,
value, schema, edited_by, edited_at, labels
) VALUES (
$1, $2, $3, $4,
NULL, '', $5,
$6, $7, $8,
$9,
$10, $11::text::json, $12, now(), $13
)"#,
w_id,
nf.path,
nf.summary,
nf.description.as_deref().unwrap_or(""),
nf.tag,
nf.dedicated_worker,
nf.visible_to_runner_only.unwrap_or(false),
windmill_common::resolve_on_behalf_of_email(
nf.on_behalf_of_email.as_deref(),
nf.preserve_on_behalf_of.unwrap_or(false),
&authed,
),
nf.ws_error_handler_muted.unwrap_or(false),
sqlx::types::Json(&nf.value) as _,
schema_str,
&authed.username,
nf.labels.as_deref() as Option<&[String]>,
)
.execute(&mut *tx)
.await?;
let version = sqlx::query_scalar!(
"INSERT INTO flow_version (workspace_id, path, value, schema, created_by)
VALUES ($1, $2, $3, $4::text::json, $5)
RETURNING id",
w_id,
nf.path,
sqlx::types::Json(nf.value) as _,
schema_str,
&authed.username,
)
.fetch_one(&mut *tx)
.await?;
sqlx::query!(
"UPDATE flow SET versions = array_append(versions, $1) WHERE path = $2 AND workspace_id = $3",
version,
nf.path,
w_id
).execute(&mut *tx).await?;
// CLI / git-sync deploys ask us to preserve any existing user draft at this
// path instead of wiping it as part of the deploy. Only wipe the deployer's
// own draft (plus the legacy NULL-email row); see scripts.rs.
if !nf.skip_draft_deletion.unwrap_or(false) {
sqlx::query!(
"DELETE FROM draft WHERE path = $1 AND workspace_id = $2 AND typ = 'flow' \
AND (email = $3 OR email IS NULL)",
nf.path,
&w_id,
&authed.email,
)
.execute(&mut *tx)
.await?;
}
audit_log(
&mut *tx,
&authed,
"flows.create",
ActionKind::Create,
&w_id,
Some(&nf.path.to_string()),
Some(
[Some(("flow", nf.path.as_str()))]
.into_iter()
.flatten()
.collect(),
),
)
.await?;
if let Some(on_behalf_of) = windmill_common::check_on_behalf_of_preservation(
nf.on_behalf_of_email.as_deref(),
nf.preserve_on_behalf_of.unwrap_or(false),
&authed,
&authed.email,
) {
audit_log(
&mut *tx,
&authed,
"flows.on_behalf_of",
ActionKind::Create,
&w_id,
Some(&nf.path),
Some(
[
("on_behalf_of", on_behalf_of.as_str()),
("action", "create"),
]
.into(),
),
)
.await?;
}
let mut args: HashMap<String, Box<serde_json::value::RawValue>> = HashMap::new();
if let Some(dm) = nf.deployment_message {
args.insert("deployment_message".to_string(), to_raw_value(&dm));
}
let tx = PushIsolationLevel::Transaction(tx);
let (dependency_job_uuid, mut new_tx) = push(
&db,
tx,
&w_id,
JobPayload::FlowDependencies {
path: nf.path.clone(),
dedicated_worker: nf.dedicated_worker,
version: version,
debouncing_settings: Default::default(),
},
windmill_queue::PushArgs { args: &args, extra: None },
&authed.username,
&authed.email,
windmill_common::users::username_to_permissioned_as(&authed.username),
authed.token_prefix.as_deref(),
None,
None,
None,
None,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
Some(&authed.clone().into()),
false,
None,
None,
None,
)
.await?;
sqlx::query!(
"UPDATE flow SET dependency_job = $1 WHERE path = $2 AND workspace_id = $3",
dependency_job_uuid,
nf.path,
w_id
)
.execute(&mut *new_tx)
.await?;
// Store the job_id in deployment_metadata for this flow deployment
sqlx::query!(
"INSERT INTO deployment_metadata (workspace_id, path, flow_version, job_id)
VALUES ($1, $2, $3, $4)
ON CONFLICT (workspace_id, path, flow_version) WHERE flow_version IS NOT NULL
DO UPDATE SET job_id = EXCLUDED.job_id",
w_id,
nf.path,
version,
dependency_job_uuid
)
.execute(&mut *new_tx)
.await?;
new_tx.commit().await?;
webhook.send_message(
w_id.clone(),
WebhookMessage::CreateFlow { workspace: w_id.clone(), path: nf.path.clone() },
);
Ok((StatusCode::CREATED, nf.path.to_string()))
}
async fn check_schedule_conflict<'c>(
tx: &mut Transaction<'c, Postgres>,
w_id: &str,
path: &str,
) -> error::Result<()> {
let exists_flow = sqlx::query_scalar!(
"SELECT EXISTS (SELECT 1 FROM schedule WHERE path = $1 AND workspace_id = $2 AND path != \
script_path)",
path,
w_id
)
.fetch_one(&mut **tx)
.await?
.unwrap_or(false);
if exists_flow {
return Err(error::Error::BadConfig(format!(
"A flow cannot have the same path as a schedule if the schedule does not trigger that \
same flow: {path}",
)));
};
Ok(())
}
#[derive(Serialize)]
pub struct FlowVersion {
pub id: i64,
pub created_at: chrono::DateTime<chrono::Utc>,
#[serde(skip_serializing_if = "Option::is_none")]
pub deployment_msg: Option<String>,
}
async fn get_flow_history(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<Vec<FlowVersion>> {
let path = path.to_path();
check_scopes(&authed, || format!("flows:read:{}", path))?;
let mut tx = user_db.begin(&authed).await?;
let flows = sqlx::query_as!(
FlowVersion,
"SELECT flow_version.id, flow_version.created_at, deployment_metadata.deployment_msg FROM flow_version
LEFT JOIN deployment_metadata ON flow_version.id = deployment_metadata.flow_version
WHERE flow_version.path = $1 AND flow_version.workspace_id = $2
ORDER BY flow_version.created_at DESC",
path,
w_id
)
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
Ok(Json(flows))
}
async fn get_latest_version(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<Option<FlowVersion>> {
let path = path.to_path();
check_scopes(&authed, || format!("flows:read:{}", path))?;
let mut tx = user_db.begin(&authed).await?;
let version = sqlx::query_as!(
FlowVersion,
"SELECT flow_version.id, flow_version.created_at, deployment_metadata.deployment_msg FROM flow_version
LEFT JOIN deployment_metadata ON flow_version.id = deployment_metadata.flow_version
WHERE flow_version.path = $1 AND flow_version.workspace_id = $2
ORDER BY flow_version.created_at DESC",
path,
w_id
)
.fetch_optional(&mut *tx)
.await?;
tx.commit().await?;
Ok(Json(version))
}
async fn get_flow_version(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, version, path)): Path<(String, i64, StripPath)>,
) -> JsonResult<Flow> {
let path = path.to_path();
check_scopes(&authed, || format!("flows:read:{}", path))?;
let mut tx = user_db.begin(&authed).await?;
let flow = sqlx::query_as::<_, Flow>(
"SELECT flow.workspace_id, flow.path, flow.summary, flow.description, flow.archived, flow.extra_perms, flow.dedicated_worker, flow.tag, flow.ws_error_handler_muted, flow.timeout, flow.visible_to_runner_only, flow.on_behalf_of_email, flow.labels, flow_version.schema, flow_version.value, flow_version.created_at as edited_at, flow_version.created_by as edited_by
FROM flow
LEFT JOIN flow_version ON flow_version.path = flow.path AND flow_version.workspace_id = flow.workspace_id
WHERE flow.path = $1 AND flow.workspace_id = $2 AND flow_version.id = $3",
)
.bind(path)
.bind(w_id)
.bind(version)
.fetch_optional(&mut *tx)
.await?;
tx.commit().await?;
let flow = not_found_if_none(flow, "Flow version", version.to_string())?;
Ok(Json(flow))
}
async fn get_flow_version_by_id(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, version)): Path<(String, i64)>,
) -> JsonResult<Flow> {
let mut tx = user_db.begin(&authed).await?;
// First, fetch the path to perform authorization check early
let path: Option<String> =
sqlx::query_scalar("SELECT path FROM flow_version WHERE id = $1 AND workspace_id = $2")
.bind(version)
.bind(&w_id)
.fetch_optional(&mut *tx)
.await?;
let path = not_found_if_none(
path,
"Flow version",
format!("{} in workspace {}", version, w_id),
)?;
// Perform authorization check before fetching full data
check_scopes(&authed, || format!("flows:read:{}", path))?;
// Now fetch the full flow data with INNER JOIN to ensure flow exists
let flow = sqlx::query_as::<_, Flow>(
"SELECT
flow.workspace_id,
flow.path,
flow.summary,
flow.description,
flow.archived,
flow.extra_perms,
flow.dedicated_worker,
flow.tag,
flow.ws_error_handler_muted,
flow.timeout,
flow.visible_to_runner_only,
flow.on_behalf_of_email,
flow.labels,
flow_version.schema,
flow_version.value,
flow_version.created_at as edited_at,
flow_version.created_by as edited_by
FROM flow
INNER JOIN flow_version
ON flow_version.path = flow.path
AND flow_version.workspace_id = flow.workspace_id
WHERE flow_version.id = $1 AND flow.workspace_id = $2",
)
.bind(version)
.bind(&w_id)
.fetch_optional(&mut *tx)
.await?;
tx.commit().await?;
let flow = not_found_if_none(
flow,
"Flow",
format!("for version {} (flow may have been deleted)", version),
)?;
Ok(Json(flow))
}
#[derive(Deserialize)]
pub struct FlowHistoryUpdate {
pub deployment_msg: String,
}
async fn update_flow_history(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Path((w_id, version)): Path<(String, i64)>,
Json(history_update): Json<FlowHistoryUpdate>,
) -> Result<()> {
let mut tx = user_db.begin(&authed).await?;
// Fetch path and perform authorization check early
let path: Option<String> =
sqlx::query_scalar("SELECT path FROM flow_version WHERE workspace_id = $1 AND id = $2")
.bind(&w_id)
.bind(version)
.fetch_optional(&mut *tx)
.await?;
let path = not_found_if_none(
path,
"Flow version",
format!("{} in workspace {}", version, w_id),
)?;
// Perform authorization check before any modifications
check_scopes(&authed, || format!("flows:write:{}", path))?;
// Insert or update deployment metadata
sqlx::query!(
"INSERT INTO deployment_metadata (workspace_id, path, flow_version, deployment_msg) VALUES ($1, $2, $3, $4) ON CONFLICT (workspace_id, path, flow_version) WHERE flow_version IS NOT NULL DO UPDATE SET deployment_msg = EXCLUDED.deployment_msg",
&w_id,
path,
version,
history_update.deployment_msg,
)
.fetch_optional(&mut *tx)
.await?;
tx.commit().await?;
Ok(())
}
async fn update_flow(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Extension(db): Extension<DB>,
Extension(webhook): Extension<WebhookShared>,
Path((w_id, flow_path)): Path<(String, StripPath)>,
Json(nf): Json<NewFlow>,
) -> Result<String> {
if authed.is_operator {
return Err(Error::NotAuthorized(
"Operators cannot update flows for security reasons".to_string(),
));
}
let flow_path = flow_path.to_path();
check_scopes(&authed, || format!("flows:write:{}", flow_path))?;
if let RuleCheckResult::Blocked(msg) = check_deploy_rules(
&w_id,
AuditAuthorable::username(&authed),
&authed.groups,
authed.is_admin,
&db,
)
.await?
{
return Err(Error::PermissionDenied(msg));
}
validate_flow(&nf).await?;
let authed = maybe_refresh_folders(&flow_path, &w_id, authed, &db).await;
let mut tx = user_db.clone().begin(&authed).await?;
check_schedule_conflict(&mut tx, &w_id, flow_path).await?;
let schema = nf.schema.map(|x| x.0);
let old_dep_job = sqlx::query_scalar!(
"SELECT dependency_job FROM flow WHERE path = $1 AND workspace_id = $2",
flow_path,
w_id
)
.fetch_optional(&mut *tx)
.await?;
let old_dep_job = not_found_if_none(old_dep_job, "Flow", flow_path)?;
let is_new_path = nf.path != flow_path;
let schema_str = schema.and_then(|x| serde_json::to_string(&x).ok());
sqlx::query!(
"
UPDATE
flow
SET
path = $1,
summary = $2,
description = $3,
dependency_job = NULL,
lock_error_logs = '',
tag = $4,
dedicated_worker = $5,
visible_to_runner_only = $6,
on_behalf_of_email = $7,
ws_error_handler_muted = $8,
value = $9,
schema = $10::text::json,
edited_by = $11,
edited_at = now(),
labels = COALESCE($14, labels)
WHERE
path = $12 AND workspace_id = $13",
if is_new_path { flow_path } else { &nf.path },
nf.summary,
nf.description.as_deref().unwrap_or(""),
nf.tag,
nf.dedicated_worker,
nf.visible_to_runner_only.unwrap_or(false),
windmill_common::resolve_on_behalf_of_email(
nf.on_behalf_of_email.as_deref(),
nf.preserve_on_behalf_of.unwrap_or(false),
&authed,
),
nf.ws_error_handler_muted.unwrap_or(false),
sqlx::types::Json(&nf.value) as _,
schema_str,
authed.username,
flow_path,
w_id,
nf.labels.as_deref() as Option<&[String]>,
)
.execute(&mut *tx)
.await
.map_err(|e| {
error::Error::internal_err(format!("Error updating flow due to flow update: {e:#}"))
})?;
if is_new_path {
// if new path, must clone flow to new path and delete old flow for flow_version foreign key constraint
sqlx::query!(
"INSERT INTO flow
(workspace_id, path, summary, description, archived, extra_perms, dependency_job, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, on_behalf_of_email, concurrency_key, versions, value, schema, edited_by, edited_at, labels)
SELECT workspace_id, $1, summary, description, archived, extra_perms, dependency_job, tag, ws_error_handler_muted, dedicated_worker, timeout, visible_to_runner_only, on_behalf_of_email, concurrency_key, versions, value, schema, edited_by, edited_at, labels
FROM flow
WHERE path = $2 AND workspace_id = $3",
nf.path,
flow_path,
w_id
)
.execute(&mut *tx)
.await
.map_err(|e| {
error::Error::internal_err(format!("Error updating flow due to create new flow: {e:#}"))
})?;
sqlx::query!(
"UPDATE flow_version SET path = $1 WHERE path = $2 AND workspace_id = $3",
nf.path,
flow_path,
w_id
)
.execute(&mut *tx)
.await
.map_err(|e| {
error::Error::internal_err(format!(
"Error updating flow due to updating flow history path: {e:#}"
))
})?;
sqlx::query!(
"DELETE FROM flow WHERE path = $1 AND workspace_id = $2",
flow_path,
w_id
)
.execute(&mut *tx)
.await
.map_err(|e| {
error::Error::internal_err(format!(
"Error updating flow due to deleting old flow: {e:#}"
))
})?;
sqlx::query!(
"UPDATE capture_config SET path = $1 WHERE path = $2 AND workspace_id = $3 AND is_flow IS TRUE",
nf.path,
flow_path,
w_id
)
.execute(&mut *tx)
.await?;
sqlx::query!(
"UPDATE capture SET path = $1 WHERE path = $2 AND workspace_id = $3 AND is_flow IS TRUE",
nf.path,
flow_path,
w_id
)
.execute(&mut *tx)
.await?;
// Update ci_test_reference when a tested flow is renamed
sqlx::query!(
"UPDATE ci_test_reference SET tested_item_path = $1 WHERE tested_item_path = $2 AND workspace_id = $3 AND tested_item_kind = 'flow'",
nf.path,
flow_path,
w_id
)
.execute(&mut *tx)
.await?;
}
// tracing::error!("Updating flow: {:?}", nf.value.get());
// This will lock anyone who is trying to iterate on flow_versions with given path and parameters.
let version = sqlx::query_scalar!(
"INSERT INTO flow_version (workspace_id, path, value, schema, created_by) VALUES ($1, $2, $3, $4::text::json, $5) RETURNING id",
w_id,
nf.path,
sqlx::types::Json(nf.value) as _,
schema_str,
&authed.username,
)
.fetch_one(&mut *tx)
.await
.map_err(|e| {
error::Error::internal_err(format!(
"Error updating flow due to flow history insert: {e:#}"
))
})?;
// TODO: This should happen only after we are done with dependency job.
sqlx::query!(
"UPDATE flow SET versions = array_append(versions, $1) WHERE path = $2 AND workspace_id = $3",
version, nf.path, w_id
).execute(&mut *tx).await?;
if is_new_path {
check_schedule_conflict(&mut tx, &w_id, &nf.path).await?;
if !authed.is_admin {
require_owner_of_path(&authed, flow_path)?;
}
}
let mut schedulables: Vec<Schedule> = sqlx::query_as::<_, Schedule>(
"UPDATE schedule SET script_path = $1 WHERE script_path = $2 AND path != $2 AND workspace_id = $3 AND is_flow IS true RETURNING *")
.bind(&nf.path)
.bind(&flow_path)
.bind(&w_id)
.fetch_all(&mut *tx)
.await.map_err(|e| error::Error::internal_err(format!("Error updating flow due to related schedules update: {e:#}")))?;
let schedule = sqlx::query_as::<_, Schedule>(
"UPDATE schedule SET path = $1, script_path = $1 WHERE path = $2 AND workspace_id = $3 AND is_flow IS true RETURNING *")
.bind(&nf.path)
.bind(&flow_path)
.bind(&w_id)
.fetch_optional(&mut *tx)
.await.map_err(|e| error::Error::internal_err(format!("Error updating flow due to related schedule update: {e:#}")))?;
if let Some(schedule) = schedule {
clear_schedule(&mut tx, &flow_path, &w_id).await?;
schedulables.push(schedule);
}
for schedule in schedulables.into_iter() {
clear_schedule(&mut tx, &schedule.path, &w_id).await?;
if schedule.enabled {
tx = push_scheduled_job(&db, tx, &schedule, None, None).await?;
}
}
if is_new_path {
windmill_common::triggers::update_triggers_script_path(
&mut tx, &nf.path, &flow_path, &w_id, true,
)
.await
.map_err(|e| {
error::Error::internal_err(format!(
"Error updating triggers due to runnable path change: {e:#}"
))
})?;
}
// CLI / git-sync deploys ask us to preserve any existing user draft at this
// path instead of wiping it as part of the deploy. Only wipe the deployer's
// own draft (plus the legacy NULL-email row); see scripts.rs.
if !nf.skip_draft_deletion.unwrap_or(false) {
sqlx::query!(
"DELETE FROM draft WHERE path = $1 AND workspace_id = $2 AND typ = 'flow' \
AND (email = $3 OR email IS NULL)",
flow_path,
&w_id,
&authed.email,
)
.execute(&mut *tx)
.await?;
}
audit_log(
&mut *tx,
&authed,
"flows.update",
ActionKind::Create,
&w_id,
Some(&nf.path.to_string()),
Some(
[Some(("flow", nf.path.as_str()))]
.into_iter()
.flatten()
.collect(),
),
)
.await?;
if let Some(on_behalf_of) = windmill_common::check_on_behalf_of_preservation(
nf.on_behalf_of_email.as_deref(),
nf.preserve_on_behalf_of.unwrap_or(false),
&authed,
&authed.email,
) {
audit_log(
&mut *tx,
&authed,
"flows.on_behalf_of",
ActionKind::Update,
&w_id,
Some(&nf.path),
Some(
[
("on_behalf_of", on_behalf_of.as_str()),
("action", "update"),
]
.into(),
),
)
.await?;
}
webhook.send_message(
w_id.clone(),
WebhookMessage::UpdateFlow {
workspace: w_id.clone(),
old_path: flow_path.to_owned(),
new_path: nf.path.clone(),
},
);
let tx = PushIsolationLevel::Transaction(tx);
let mut args: HashMap<String, Box<serde_json::value::RawValue>> = HashMap::new();
if let Some(dm) = nf.deployment_message {
args.insert("deployment_message".to_string(), to_raw_value(&dm));
}
args.insert("parent_path".to_string(), to_raw_value(&flow_path));
let (dependency_job_uuid, mut new_tx) = push(
&db,
tx,
&w_id,
JobPayload::FlowDependencies {
path: nf.path.clone(),
dedicated_worker: nf.dedicated_worker,
version,
debouncing_settings: Default::default(),
},
windmill_queue::PushArgs { args: &args, extra: None },
&authed.username,
&authed.email,
windmill_common::users::username_to_permissioned_as(&authed.username),
authed.token_prefix.as_deref(),
None,
None,
None,
None,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
Some(&authed.clone().into()),
false,
None,
None,
None,
)
.await?;
sqlx::query!(
"UPDATE flow SET dependency_job = $1 WHERE path = $2 AND workspace_id = $3",
dependency_job_uuid,
nf.path,
w_id
)
.execute(&mut *new_tx)
.await
.map_err(|e| {
error::Error::internal_err(format!(
"Error updating flow due to updating dependency job field: {e:#}"
))
})?;
// Store the job_id in deployment_metadata for this flow deployment
sqlx::query!(
"INSERT INTO deployment_metadata (workspace_id, path, flow_version, job_id)
VALUES ($1, $2, $3, $4)
ON CONFLICT (workspace_id, path, flow_version) WHERE flow_version IS NOT NULL
DO UPDATE SET job_id = EXCLUDED.job_id",
w_id,
nf.path,
version,
dependency_job_uuid
)
.execute(&mut *new_tx)
.await
.map_err(|e| {
error::Error::internal_err(format!(
"Error updating deployment_metadata with job_id: {e:#}"
))
})?;
if let Some(old_dep_job) = old_dep_job {
sqlx::query!(
"UPDATE v2_job_queue SET
canceled_by = $2,
canceled_reason = 're-deployment'
WHERE id = $1",
old_dep_job,
&authed.username
)
.execute(&mut *new_tx)
.await
.map_err(|e| {
error::Error::internal_err(format!(
"Error updating flow due to cancelling dependency job: {e:#}"
))
})?;
}
new_tx.commit().await?;
// Trigger CI tests for items that reference this flow
{
let db2 = db.clone();
let w_id2 = w_id.clone();
let flow_path2 = nf.path.clone();
let email2 = authed.email.clone();
let username2 = authed.username.clone();
tokio::spawn(async move {
if let Err(e) = windmill_dep_map::ci_tests::trigger_ci_tests_for_item(
&db2,
&w_id2,
&flow_path2,
"flow",
&email2,
&username2,
)
.await
{
tracing::error!(%e, "error triggering CI tests after flow deploy");
}
});
}
Ok(nf.path.to_string())
}
async fn list_tokens(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<Vec<TruncatedTokenWithEmail>> {
let path = path.to_path();
check_scopes(&authed, || format!("flows:read:{}", path))?;
list_tokens_internal(&db, &w_id, &path, true).await
}
#[derive(Serialize)]
struct DeploymentStatus {
lock_error_logs: Option<String>,
job_id: Option<sqlx::types::Uuid>,
}
async fn get_deployment_status(
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<DeploymentStatus> {
let path = path.to_path();
let mut tx = db.begin().await?;
let status_o = sqlx::query!(
"SELECT f.lock_error_logs, dm.job_id
FROM flow f
LEFT JOIN deployment_metadata dm ON f.versions[array_upper(f.versions, 1)] = dm.flow_version
AND f.workspace_id = dm.workspace_id AND f.path = dm.path
WHERE f.path = $1 AND f.workspace_id = $2",
path,
w_id,
)
.fetch_optional(&mut *tx)
.await?;
let status = not_found_if_none(status_o, "DeploymentStatus", path)?;
let deployment_status =
DeploymentStatus { lock_error_logs: status.lock_error_logs, job_id: status.job_id };
tx.commit().await?;
Ok(Json(deployment_status))
}
// Fields inlined rather than flattened (axum query bool quirk); see GetScriptByPathQuery in scripts.rs.
#[derive(Deserialize)]
struct GetFlowByPathQuery {
with_starred_info: Option<bool>,
#[serde(default)]
get_draft: bool,
}
async fn get_flow_by_path(
authed: ApiAuthed,
Extension(user_db): Extension<UserDB>,
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
Query(query): Query<GetFlowByPathQuery>,
) -> JsonResult<WithDraftOverlay> {
let path = path.to_path();
check_scopes(&authed, || format!("flows:read:{}", path))?;
let mut tx = user_db.begin(&authed).await?;
let flow_o = if query.with_starred_info.unwrap_or(false) {
sqlx::query_as::<_, FlowWithStarred>(
r#"
SELECT
flow.workspace_id,
flow.path,
flow.lock_error_logs,
flow.summary,
flow.description,
flow.archived,
flow.extra_perms,
flow.dedicated_worker,
flow.tag,
flow.ws_error_handler_muted,
flow.timeout,
flow.visible_to_runner_only,
flow.on_behalf_of_email,
flow.labels,
folder_labels(flow.workspace_id, flow.path) AS inherited_labels,
flow_version.id AS version_id,
flow_version.schema,
flow_version.value,
flow_version.created_at AS edited_at,
flow_version.created_by AS edited_by,
favorite.path IS NOT NULL AS starred
FROM flow
LEFT JOIN favorite
ON favorite.favorite_kind = 'flow'
AND favorite.workspace_id = flow.workspace_id
AND favorite.path = flow.path
AND favorite.usr = $3
LEFT JOIN flow_version
ON flow_version.id = flow.versions[array_upper(flow.versions, 1)]
WHERE flow.path = $1 AND flow.workspace_id = $2
"#,
)
.bind(path)
.bind(&w_id)
.bind(&authed.username)
.fetch_optional(&mut *tx)
.await?
} else {
sqlx::query_as::<_, FlowWithStarred>(
r#"
SELECT
flow.workspace_id,
flow.path,
flow.lock_error_logs,
flow.summary,
flow.description,
flow.archived,
flow.extra_perms,
flow.dedicated_worker,
flow.tag,
flow.ws_error_handler_muted,
flow.timeout,
flow.visible_to_runner_only,
flow.on_behalf_of_email,
flow.labels,
folder_labels(flow.workspace_id, flow.path) AS inherited_labels,
flow_version.id AS version_id,
flow_version.schema,
flow_version.value,
flow_version.created_at AS edited_at,
flow_version.created_by AS edited_by,
NULL AS starred
FROM flow
LEFT JOIN flow_version
ON flow_version.id = flow.versions[array_upper(flow.versions, 1)]
WHERE flow.path = $1 AND flow.workspace_id = $2
"#,
)
.bind(path)
.bind(&w_id)
.fetch_optional(&mut *tx)
.await?
};
tx.commit().await?;
// No deployed row + `get_draft`: fall back to the draft table; see scripts.rs.
let overlay = overlay_or_draft_only(
&db,
&w_id,
&authed.email,
UserDraftItemKind::Flow,
path,
query.get_draft,
flow_o,
|| windmill_common::error::Error::NotFound(format!("Flow not found at path {path}")),
)
.await?;
Ok(Json(overlay))
}
async fn exists_flow_by_path(
Extension(db): Extension<DB>,
Path((w_id, path)): Path<(String, StripPath)>,
) -> JsonResult<bool> {
let path = path.to_path();
let exists = sqlx::query_scalar!(
"SELECT EXISTS(SELECT 1 FROM flow WHERE path = $1 AND workspace_id = $2)",
path,
w_id
)
.fetch_one(&db)
.await?
.unwrap_or(false);
Ok(Json(exists))
}
#[derive(Deserialize)]
struct Archived {
archived: Option<bool>,
}
async fn archive_flow_by_path(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Extension(user_db): Extension<UserDB>,
Extension(webhook): Extension<WebhookShared>,
Path((w_id, path)): Path<(String, StripPath)>,
Json(archived): Json<Archived>,
) -> Result<String> {
let path = path.to_path();
check_scopes(&authed, || format!("flows:write:{}", path))?;
if let RuleCheckResult::Blocked(msg) = check_deploy_rules(
&w_id,
AuditAuthorable::username(&authed),
&authed.groups,
authed.is_admin,
&db,
)
.await?
{
return Err(Error::PermissionDenied(msg));
}
let mut tx = user_db.begin(&authed).await?;
sqlx::query!(
"UPDATE flow SET archived = $1 WHERE path = $2 AND workspace_id = $3",
archived.archived.unwrap_or(true),
path,
&w_id
)
.execute(&mut *tx)
.await?;
clear_static_asset_usage(&mut *tx, &w_id, path, AssetUsageKind::Flow).await?;
audit_log(
&mut *tx,
&authed,
"flows.archive",
ActionKind::Delete,
&w_id,
Some(path),
Some([("workspace", w_id.as_str())].into()),
)
.await?;
ScopedDependencyMap::clear_map_for_item(path, &w_id, "flow", tx, &None)
.await
.commit()
.await?;
handle_deployment_metadata(
&authed.email,
&authed.username,
&db,
&w_id,
DeployedObject::Flow {
path: path.to_string(),
parent_path: Some(path.to_string()),
version: 0, // dummy version as it will not get inserted in db
},
Some(format!(
"Flow '{}' {}",
path,
if archived.archived.unwrap_or(true) {
"archived"
} else {
"unarchived"
}
)),
true,
None,
)
.await?;
webhook.send_message(
w_id.clone(),
WebhookMessage::ArchiveFlow { workspace: w_id, path: path.to_owned() },
);
Ok(format!("Flow {path} archived"))
}
/// Validates that flow debouncing configuration is supported by all workers
/// Returns an error if debouncing is configured but workers are behind required version
async fn guard_flow_from_debounce_data(nf: &NewFlow) -> Result<()> {
let flow_value = nf.parse_flow_value()?;
if !MIN_VERSION_SUPPORTS_DEBOUNCING.met().await && !flow_value.debouncing_settings.is_default()
{
tracing::warn!(
"Flow debouncing configuration rejected: workers are behind minimum required version for debouncing feature"
);
return Err(Error::WorkersAreBehind {
feature: "Debouncing".into(),
min_version: "1.566.0".into(),
});
}
if !MIN_VERSION_SUPPORTS_DEBOUNCING_V2.met().await
&& !flow_value.debouncing_settings.is_legacy_compatible()
&& !*WMDEBUG_FORCE_NO_LEGACY_DEBOUNCING_COMPAT
{
tracing::warn!(
"Flow debouncing configuration rejected: workers are behind minimum required version for debouncing feature"
);
return Err(Error::WorkersAreBehind {
feature: "V2 Debouncing".into(),
min_version: "1.597.0".into(),
});
}
// Check node-level debouncing on all modules (including nested branches/loops)
let mut has_node_debouncing = false;
let check_result = FlowModule::traverse_modules(&flow_value.modules, &mut |m| {
if m.debouncing
.as_ref()
.is_some_and(|d| d.debounce_delay_s.is_some_and(|s| s > 0))
{
has_node_debouncing = true;
}
Ok(())
});
if let Err(e) = check_result {
tracing::warn!("Failed to traverse flow modules for debounce guard: {e}");
}
if has_node_debouncing && !MIN_VERSION_SUPPORTS_NODE_DEBOUNCING.met().await {
return Err(Error::WorkersAreBehind {
feature: "Flow node debouncing".into(),
min_version: "1.658.0".into(),
});
}
Ok(())
}
#[derive(Deserialize)]
struct DeleteFlowQuery {
keep_captures: Option<bool>,
}
async fn delete_flow_by_path(
authed: ApiAuthed,
Extension(db): Extension<DB>,
Extension(user_db): Extension<UserDB>,
Extension(webhook): Extension<WebhookShared>,
Path((w_id, path)): Path<(String, StripPath)>,
Query(query): Query<DeleteFlowQuery>,
) -> Result<String> {
let path = path.to_path();
check_scopes(&authed, || format!("flows:write:{}", path))?;
if let RuleCheckResult::Blocked(msg) = check_deploy_rules(
&w_id,
AuditAuthorable::username(&authed),
&authed.groups,
authed.is_admin,
&db,
)
.await?
{
return Err(Error::PermissionDenied(msg));
}
let mut tx = user_db.begin(&authed).await?;
// Capture all related data for trashbin before deleting (CASCADE will remove flow_version, flow_node)
let trash_flow: Option<serde_json::Value> =
sqlx::query_scalar("SELECT to_jsonb(t) FROM flow t WHERE path = $1 AND workspace_id = $2")
.bind(path)
.bind(&w_id)
.fetch_optional(&mut *tx)
.await?;
let trash_flow_versions: Vec<serde_json::Value> = sqlx::query_scalar(
"SELECT to_jsonb(t) FROM flow_version t WHERE path = $1 AND workspace_id = $2",
)
.bind(path)
.bind(&w_id)
.fetch_all(&mut *tx)
.await?;
let trash_flow_nodes: Vec<serde_json::Value> = sqlx::query_scalar(
"SELECT to_jsonb(t) FROM flow_node t WHERE path = $1 AND workspace_id = $2",
)
.bind(path)
.bind(&w_id)
.fetch_all(&mut *tx)
.await?;
let trash_drafts: Vec<serde_json::Value> = sqlx::query_scalar(
"SELECT to_jsonb(t) FROM draft t WHERE path = $1 AND workspace_id = $2 AND typ = 'flow'",
)
.bind(path)
.bind(&w_id)
.fetch_all(&mut *tx)
.await?;
sqlx::query!(
"DELETE FROM draft WHERE path = $1 AND workspace_id = $2 AND typ = 'flow'",
path,
&w_id
)
.execute(&mut *tx)
.await?;
sqlx::query!(
"DELETE FROM flow WHERE path = $1 AND workspace_id = $2",
path,
&w_id
)
.execute(&mut *tx)
.await?;
if let Some(flow_data) = trash_flow {
let mut trash_data = serde_json::json!({"row": flow_data});
if !trash_flow_versions.is_empty() {
trash_data["flow_versions"] = serde_json::Value::Array(trash_flow_versions);
}
if !trash_flow_nodes.is_empty() {
trash_data["flow_nodes"] = serde_json::Value::Array(trash_flow_nodes);
}
if !trash_drafts.is_empty() {
trash_data["drafts"] = serde_json::Value::Array(trash_drafts);
}
windmill_common::trashbin::move_to_trash(
&mut *tx,
&w_id,
"flow",
path,
trash_data,
&authed.username,
)
.await?;
}
if !query.keep_captures.unwrap_or(false) {
sqlx::query!(
"DELETE FROM capture_config WHERE path = $1 AND workspace_id = $2 AND is_flow IS TRUE",
path,
&w_id
)
.execute(&mut *tx)
.await?;
sqlx::query!(
"DELETE FROM capture WHERE path = $1 AND workspace_id = $2 AND is_flow IS TRUE",
path,
&w_id
)
.execute(&mut *tx)
.await?;
}
audit_log(
&mut *tx,
&authed,
"flows.delete",
ActionKind::Delete,
&w_id,
Some(path),
Some([("workspace", w_id.as_str())].into()),
)
.await?;
tx.commit().await?;
handle_deployment_metadata(
&authed.email,
&authed.username,
&db,
&w_id,
DeployedObject::Flow {
path: path.to_string(),
parent_path: Some(path.to_string()),
version: 0, // dummy version as it will not get inserted in db
},
Some(format!("Flow '{}' deleted", path)),
true,
None,
)
.await?;
sqlx::query!(
"DELETE FROM deployment_metadata WHERE path = $1 AND workspace_id = $2 AND script_hash IS NULL and app_version IS NULL",
path,
w_id
)
.execute(&db)
.await
.map_err(|e| {
Error::internal_err(format!(
"error deleting deployment metadata for script with path {path} in workspace {w_id}: {e:#}"
))
})?;
webhook.send_message(
w_id.clone(),
WebhookMessage::DeleteFlow { workspace: w_id, path: path.to_owned() },
);
Ok(format!("Flow {path} deleted"))
}
#[cfg(test)]
mod tests {
use std::{collections::HashMap, time::Duration};
use windmill_common::{
flows::{
ConstantDelay, ExponentialDelay, FlowModule, FlowModuleValue, FlowValue,
InputTransform, Retry, StopAfterIf,
},
runnable_settings::{
ConcurrencySettings, ConcurrencySettingsWithCustom, DebouncingSettings,
},
scripts,
};
const SECOND: Duration = Duration::from_secs(1);
#[test]
fn flowmodule_serde() {
let fv = FlowValue {
modules: vec![
FlowModule {
id: "a".to_string(),
value: windmill_common::worker::to_raw_value(&FlowModuleValue::Script {
path: "test".to_string(),
input_transforms: [(
"test".to_string(),
InputTransform::Static {
value: windmill_common::worker::to_raw_value(&"test2".to_string()),
},
)]
.into(),
hash: None,
tag_override: None,
is_trigger: None,
pass_flow_input_directly: None,
}),
stop_after_if: None,
stop_after_all_iters_if: None,
summary: None,
suspend: Default::default(),
retry: None,
sleep: None,
cache_ttl: None,
cache_ignore_s3_path: None,
mock: None,
timeout: None,
priority: None,
delete_after_use: None,
delete_after_secs: None,
continue_on_error: None,
skip_if: None,
apply_preprocessor: None,
pass_flow_input_directly: None,
debouncing: None,
},
FlowModule {
id: "b".to_string(),
value: windmill_common::worker::to_raw_value(&FlowModuleValue::RawScript {
input_transforms: HashMap::new(),
content: "test".to_string(),
language: scripts::ScriptLang::Deno,
path: None,
lock: None,
tag: None,
is_trigger: None,
assets: None,
concurrency_settings: ConcurrencySettingsWithCustom::default(),
}),
stop_after_if: Some(StopAfterIf {
expr: "foo = 'bar'".to_string(),
..Default::default()
}),
stop_after_all_iters_if: None,
summary: None,
suspend: Default::default(),
retry: None,
sleep: None,
cache_ttl: None,
cache_ignore_s3_path: None,
mock: None,
timeout: None,
priority: None,
delete_after_use: None,
delete_after_secs: None,
continue_on_error: None,
skip_if: None,
apply_preprocessor: None,
pass_flow_input_directly: None,
debouncing: None,
},
FlowModule {
id: "c".to_string(),
value: windmill_common::worker::to_raw_value(&FlowModuleValue::ForloopFlow {
iterator: InputTransform::Static {
value: windmill_common::worker::to_raw_value(&[1, 2, 3]),
},
modules: vec![],
modules_node: None,
skip_failures: true,
parallel: false,
parallelism: None,
squash: None,
}),
stop_after_if: Some(StopAfterIf {
expr: "previous.isEmpty()".to_string(),
..Default::default()
}),
stop_after_all_iters_if: None,
summary: None,
suspend: Default::default(),
retry: None,
sleep: None,
cache_ttl: None,
cache_ignore_s3_path: None,
mock: None,
timeout: None,
priority: None,
delete_after_use: None,
delete_after_secs: None,
continue_on_error: None,
skip_if: None,
apply_preprocessor: None,
pass_flow_input_directly: None,
debouncing: None,
},
],
failure_module: Some(Box::new(FlowModule {
id: "d".to_string(),
value: FlowModuleValue::Script {
path: "test".to_string(),
input_transforms: HashMap::new(),
hash: None,
tag_override: None,
is_trigger: None,
pass_flow_input_directly: None,
}
.into(),
stop_after_if: Some(StopAfterIf {
expr: "previous.isEmpty()".to_string(),
..Default::default()
}),
stop_after_all_iters_if: None,
summary: None,
suspend: Default::default(),
retry: None,
sleep: None,
cache_ttl: None,
cache_ignore_s3_path: None,
mock: None,
timeout: None,
priority: None,
delete_after_use: None,
delete_after_secs: None,
continue_on_error: None,
skip_if: None,
apply_preprocessor: None,
pass_flow_input_directly: None,
debouncing: None,
})),
preprocessor_module: None,
same_worker: false,
preserve_step_tags: false,
skip_expr: None,
cache_ttl: None,
cache_ignore_s3_path: None,
priority: None,
early_return: None,
chat_input_enabled: None,
flow_env: None,
delete_after_use: None,
delete_after_secs: None,
concurrency_settings: ConcurrencySettings::default(),
debouncing_settings: DebouncingSettings::default(),
};
let expect = serde_json::json!({
"modules": [
{
"id": "a",
"value": {
"input_transforms": {
"test": {
"type": "static",
"value": "test2"
}
},
"type": "script",
"path": "test",
},
},
{
"id": "b",
"value": {
"input_transforms": {},
"type": "rawscript",
"content": "test",
"language": "deno"
},
"stop_after_if": {
"expr": "foo = 'bar'",
"skip_if_stopped": false
}
},
{
"id": "c",
"value": {
"type": "forloopflow",
"iterator": {
"type": "static",
"value": [
1,
2,
3
]
},
"parallel": false,
"skip_failures": true,
"modules": []
},
"stop_after_if": {
"expr": "previous.isEmpty()",
"skip_if_stopped": false
}
}
],
"failure_module": {
"id": "d",
"value": {
"input_transforms": {},
"type": "script",
"path": "test",
},
"stop_after_if": {
"expr": "previous.isEmpty()",
"skip_if_stopped": false
}
},
});
assert_eq!(dbg!(serde_json::json!(fv)), dbg!(expect));
}
#[test]
fn retry_serde() {
assert_eq!(Retry::default(), serde_json::from_str(r#"{}"#).unwrap());
assert_eq!(
Retry::default(),
serde_json::from_str(
r#"
{
"constant": {
"seconds": 0
},
"exponential": {
"multiplier": 1,
"seconds": 0
},
"retry_if": null
}
"#
)
.unwrap()
);
assert_eq!(
Retry {
constant: Default::default(),
exponential: ExponentialDelay {
attempts: 0,
multiplier: 1,
seconds: 123,
random_factor: None
},
retry_if: None
},
serde_json::from_str(
r#"
{
"constant": {},
"exponential": { "seconds": 123 },
"retry_if" : null
}
"#
)
.unwrap()
);
}
#[test]
fn retry_exponential() {
let retry = Retry {
constant: ConstantDelay::default(),
exponential: ExponentialDelay {
attempts: 3,
multiplier: 4,
seconds: 3,
random_factor: None,
},
retry_if: None,
};
assert_eq!(
vec![
Some(12 * SECOND),
Some(36 * SECOND),
Some(108 * SECOND),
None
],
(0..4)
.map(|previous_attempts| retry.interval(previous_attempts, false))
.collect::<Vec<_>>()
);
assert_eq!(Some(108 * SECOND), retry.max_interval());
}
#[test]
fn retry_both() {
let retry = Retry {
constant: ConstantDelay { attempts: 2, seconds: 4 },
exponential: ExponentialDelay {
attempts: 2,
multiplier: 1,
seconds: 3,
random_factor: None,
},
retry_if: None,
};
assert_eq!(
vec![
Some(4 * SECOND),
Some(4 * SECOND),
Some(27 * SECOND),
Some(81 * SECOND),
None,
],
(0..5)
.map(|previous_attempts| retry.interval(previous_attempts, false))
.collect::<Vec<_>>()
);
assert_eq!(Some(81 * SECOND), retry.max_interval());
}
}