mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-09 16:05:42 +00:00
* feat(fork): merge a fork deletion on evidence, not on the counters `workspace_diff.ahead`/`.behind` record that a write happened on a side, not what it was or who made it. That leaves one row shape undecidable: an item the parent has and the fork does not can mean the parent added it, the fork deleted it, or a git-sync pull reverted a deploy that had just brought it in. #10467 kept every such row out of the merge direction, which killed the phantom but also dropped the only way to propagate a fork-side deletion and left a rename's old path behind in the parent. Record the evidence instead: - `workspace_diff` gains, per side, the last event's kind (`write` / `delete` / `rename_from`) and origin (`authored` / `sync`). Rows written before the migration have neither and keep #10467's behavior. - The kind is probed from whether the path still holds an item once the write has committed; an item kind the probe doesn't map records no evidence rather than a deletion. Create and update are not split — nothing at that point tells them apart for every kind, and the comparison already recomputes existence per side. - The origin comes from an `X-Windmill-Deploy-Origin` header the API scopes into a task-local for the request. It is the load-bearing half: recording `delete` alone would read a git-sync revert as a fork deletion and reproduce the original bug. Two clients set it — `wmill sync push` (which the git-sync auto-pull runs inside a job) and the compare page's parent→fork "Update fork". Merging the other way stays authored so a deletion keeps propagating up a fork chain. - The merge direction admits a parent-only row only when the fork's last event was an authored delete or rename-away. Such a row stays opt-in, never bulk-selected, and reads "Removes in <parent>"; the update direction keeps offering it back as "New". A fork deletion and a rename now merge into the parent, a rename leaves no duplicate behind, and a fork the parent also edited surfaces in both directions instead of the parent silently winning. Fixes WIN-2289 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(fork): address review — detached tallies, enum wire values, doc duplication Codex P1: a dependency job tallies its deploy whenever it happens to finish, and the event kind is probed from the state at that moment. If anything removed the path in between (a git-sync revert), the stale tally read that deletion as its own and filed it as authored — handing the merge exactly the removal this is meant to withhold. `tally_deployed_object_changes` now takes `Option<DeployOrigin>`; `None` bumps the counter and leaves the evidence columns as the last vouching tally left them, and the worker path passes it. Covered by extending the removal-origin test: a detached tally after the sync archive must not disturb `(delete, sync)`. Also from review: - `fork_removed_it` compares through `DeployOrigin::as_str()` / `DeployEventKind::as_str()` rather than repeating their wire values, so a renamed variant can't silently make the predicate always false. - `deploy_origin`'s module doc no longer claims `sync` is inert: it cannot make the merge propose a removal, but it does drop a row out of both sides of the `all_ahead_items_visible` comparison. - `WorkspaceDiffRow` says why only the fork half of the evidence is consumed. - The delete-vs-revert rationale is stated once (the migration) instead of restated in eight files. - `PATH_KEYED_TABLES` is swept by a test: its query is built at runtime, so a wrong table name is not a compile error and would only surface as a failed tally for that trigger kind in a fork. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(fork): let only a request task vouch for a deploy event Round 2 found the first fix incomplete. Detaching only the failed/cancelled dependency path left the common route untouched: a dependency job that succeeds calls `handle_deployment_metadata` from the worker, where `deploy_origin::current()` read as `Authored`. A sync archiving the script while its lock generation was pending then had its deletion probed on completion and refiled as authored — the same fabricated removal, on the path most deploys actually take. `current()` now returns `Option`, `Some` only inside the request scope the API always enters. Having no scope means "not the task that served this write", which is true of every worker-side call and needs no marking at the call site. The integration test drives the real `handle_deployment_metadata` off a request task instead of the tally directly, and fails without this. Two more from the same round: - The script dependency handler passed no `renamed_from`, unlike the flow and app handlers next to it. A lock-generating create has no earlier tally, so that was the only chance for the path a rename vacated to be recorded at all — renames of Python/TS scripts left the old path in the parent, which the bash-only manual check missed. - The tally now drops a `renamed_from` equal to the path itself. Callers pass the previous path whether or not the deploy moved the item, so an unfiltered one both counted the path twice and stamped it `rename_from` when nothing was renamed. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(fork): carry a deploy's origin into the dependency job it queues Round 3 caught the previous fix cutting too deep. Refusing a detached tally any claim also refused its rename evidence, and a lock-generating deploy has no other tally — so the `renamed_from` added alongside it was inert, and a renamed flow, app or Python script still left its old path in the parent with nothing to merge. Flows and apps always generate, so renames worked essentially nowhere. The two capabilities are now separate. `TallyEvidence` says whether the tallying task served the write (`Served`, may probe what the path holds now) or is reporting one that committed earlier (`Deferred`, may not), and each column is written only from a source that answers for it. The origin itself is a fact of the deploy either way, so the request stamps it into the dependency job's args and the worker re-enters the scope with it — the last place that knows it handing it to the only tally that will run. Also from round 3: `WorkspaceDiffRow`'s event fields skip serializing `None` rather than emitting `null`, matching what the schema declares (OpenAPI 3.0.3 ignores a `description` sibling of `$ref`, so those moved onto the shared schemas). Verified against a live worker: renaming a flow in a fork records `(rename_from, authored)` on the vacated path and the merge offers its removal, while the deployed path claims nothing. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(fork): mark the CLI's parent-to-fork merge as sync `wmill workspace merge --direction to-fork` is the CLI's "Update fork" and deletes items in the fork, but without the marker the compare page sets. Its deletions were recorded as authored fork decisions, so once the parent recreated such a path the merge would offer deleting it there. Also from review: an unrecognized deploy-origin arg now reads as no evidence rather than as authored — strict where a request header is lenient, since an unmarked request really is authored but an unreadable stored value is skew. Reading the arg moved next to `stamp_origin_arg`, the half that writes it, so the round trip a lock-generating deploy depends on is covered by one test. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: drop the imports the shared arg reader made unused CI compiles with `-D warnings`, so this was four red Backend jobs rather than a lint. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(fork): stop a stale deferred rename from restating a removed path Nothing orders these events. A tally that served the write made its claim inside its own commit, but a deferred one reports a write that landed at an unknown remove. So a lock-generating rename whose dependency job finished after a sync had removed the vacated path could overwrite `(delete, sync)` with `(rename_from, authored)` — the path is gone either way, so the merge would then offer removing it from the parent on the strength of the older event. A deferred claim now only writes where the side has none, which is the case it exists for: a vacated path that nothing else has spoken for. The regression asserts the ordering directly, and fails without the guard. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(fork): record a rename's vacated path from the request that made it The deferred mechanism could not be made correct, as round 7 showed: its guard protected an existing row, but that row is deleted as soon as the two workspaces agree on the path — so a rename job finishing after the reconciliation inserted fresh, and the stale claim reappeared against whatever the parent later recreated there. Ordering cannot be recovered outside the row, because the row is disposable. So the vacated path is now recorded by the request, which is inside its own commit and whose row shares the counter's lifetime. A deploy that hands its metadata to a dependency job — every flow and app, and any script needing a lock — calls `tally_rename_vacated_path` once its transaction has committed; scripts reach it through the post-commit hook they already had, which grew a second variant rather than new plumbing. That lets the whole deferred apparatus go: `TallyEvidence`, the origin job arg and its round trip. `deploy_origin::current` is `Some` only inside a request scope again, and `handle_deployment_metadata` hands `renamed_from` to the tally only when it can answer for it — git-sync still gets it either way, so the rename keeps naming itself in the commit message. The vacated path's kind now reads `delete` rather than `rename_from` for these deploys, since it is probed rather than declared. The merge treats the two alike; only the row's tooltip is less specific. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(fork): cover raw-app renames, and stop firing CI before the lock exists Two things the vacated-path call broke or missed: - `create_script` reads its third return value as "no lock generation needed" to decide whether the script is runnable now, and the new `VacatedPath` variant made that true for renames that do generate. Those fired dependent CI tests from the API against a version with no lockfile, and again from the dependency job. The variant now decides it explicitly. - Raw apps rename through `update_app_raw`, a separate route into `update_app_internal`, which the new call had not been attached to. Both routes now go through one helper. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * test(fork): assert the kind only an inline rename can record `rename_from` is what a deploy says when it knows it moved the item, which only the path that reports both halves from its own request can. Nothing pinned it, and that is the side the vacated-path change touched. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * chore: update ee-repo-ref to a45bec03922d305aad5893ed354dc029c7f97bb4 This commit updates the EE repository reference after PR #709 was merged in windmill-ee-private. Previous ee-repo-ref: 62f494b2a51de0dfc0cfa0c3530ff19a1d32667c New ee-repo-ref: a45bec03922d305aad5893ed354dc029c7f97bb4 Automated by sync-ee-ref workflow. --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
2380 lines
76 KiB
Rust
2380 lines
76 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::{EditFlow, Flow, FlowWithStarred, ListFlowQuery, ListableFlow, NewFlow},
|
|
jobs::JobPayload,
|
|
schedule::Schedule,
|
|
triggers::MovedNativeTrigger,
|
|
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>> {
|
|
let n = 1000;
|
|
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))?;
|
|
|
|
// A `<= 0` flow timeout is "unset", not a 0-second limit that kills every run instantly.
|
|
// (The concurrency settings inside the flow value are normalized on deserialization; see
|
|
// ConcurrencySettings.) Runtime guards also protect already-stored rows.
|
|
nf.timeout = windmill_common::runnable_settings::none_if_non_positive(nf.timeout);
|
|
|
|
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.on_behalf_of.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, default_permissioned_as)) =
|
|
windmill_common::folders::resolve_folder_default_on_behalf_of(&db, &w_id, &nf.path)
|
|
.await?
|
|
{
|
|
nf.on_behalf_of_email = Some(default_email);
|
|
nf.on_behalf_of = Some(default_permissioned_as);
|
|
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());
|
|
let resolved_on_behalf_of =
|
|
windmill_common::resolve_on_behalf_of(
|
|
nf.on_behalf_of_email.as_deref(),
|
|
nf.on_behalf_of.as_deref(),
|
|
nf.preserve_on_behalf_of.unwrap_or(false),
|
|
&authed,
|
|
&w_id,
|
|
&db,
|
|
)
|
|
.await?;
|
|
// Written beside the principal only while a worker that still reads it may be live.
|
|
let legacy_on_behalf_of_email =
|
|
windmill_common::legacy_on_behalf_of_email(resolved_on_behalf_of.as_deref(), &w_id, &db)
|
|
.await?;
|
|
sqlx::query!(
|
|
r#"INSERT INTO flow (
|
|
workspace_id, path, summary, description,
|
|
dependency_job, lock_error_logs, tag,
|
|
dedicated_worker, visible_to_runner_only,
|
|
ws_error_handler_muted,
|
|
value, schema, edited_by, edited_at, labels,
|
|
on_behalf_of, on_behalf_of_email
|
|
) VALUES (
|
|
$1, $2, $3, $4,
|
|
NULL, '', $5,
|
|
$6, $7,
|
|
$8,
|
|
$9, $10::text::json, $11, now(), $12,
|
|
$13, $14
|
|
)"#,
|
|
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),
|
|
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]>,
|
|
resolved_on_behalf_of,
|
|
legacy_on_behalf_of_email,
|
|
)
|
|
.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(
|
|
resolved_on_behalf_of.as_deref(),
|
|
nf.preserve_on_behalf_of.unwrap_or(false),
|
|
&authed,
|
|
&windmill_common::users::username_to_permissioned_as(&authed.username),
|
|
) {
|
|
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(),
|
|
authed.username_override.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))
|
|
}
|
|
|
|
/// `on_behalf_of_email` is derived rather than selected: the read paths fill it from the
|
|
/// principal so clients written against the address keep working. The column itself still
|
|
/// exists for the workers that read it — see `legacy_on_behalf_of_email`.
|
|
async fn derived_on_behalf_of_email(
|
|
db: &DB,
|
|
w_id: &str,
|
|
flow: &Flow,
|
|
) -> error::Result<Option<String>> {
|
|
let Some(permissioned_as) = flow.on_behalf_of.as_deref() else {
|
|
return Ok(None);
|
|
};
|
|
// Uncached, for the reason given on `prefetch_cached_script`: this pair is round-tripped.
|
|
Ok(Some(
|
|
windmill_common::users::get_email_from_permissioned_as_uncached(permissioned_as, w_id, db)
|
|
.await?,
|
|
))
|
|
}
|
|
|
|
async fn get_flow_version(
|
|
authed: ApiAuthed,
|
|
Extension(db): Extension<DB>,
|
|
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, 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 mut flow = not_found_if_none(flow, "Flow version", version.to_string())?;
|
|
flow.on_behalf_of_email = derived_on_behalf_of_email(&db, &w_id, &flow).await?;
|
|
|
|
Ok(Json(flow))
|
|
}
|
|
|
|
async fn get_flow_version_by_id(
|
|
authed: ApiAuthed,
|
|
Extension(db): Extension<DB>,
|
|
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,
|
|
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 mut flow = not_found_if_none(
|
|
flow,
|
|
"Flow",
|
|
format!("for version {} (flow may have been deleted)", version),
|
|
)?;
|
|
flow.on_behalf_of_email = derived_on_behalf_of_email(&db, &w_id, &flow).await?;
|
|
|
|
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(())
|
|
}
|
|
|
|
/// Re-point the webhooks of the native triggers a rename carried onto the new path.
|
|
///
|
|
/// Runs after the deploy transaction commits — repointing a webhook is not undoable — and off the
|
|
/// request, because it waits on a third-party service that may be slow or gone, and a deploy that
|
|
/// already committed must not look like it failed. The rename itself marked these rows
|
|
/// `REREGISTRATION_PENDING`, so nothing is lost silently if this never finishes.
|
|
fn reregister_moved_native_triggers(
|
|
db: &DB,
|
|
authed: &ApiAuthed,
|
|
w_id: &str,
|
|
moved: Vec<MovedNativeTrigger>,
|
|
) {
|
|
if moved.is_empty() {
|
|
return;
|
|
}
|
|
#[cfg(feature = "native_trigger")]
|
|
{
|
|
let (db, authed, w_id) = (db.clone(), authed.clone(), w_id.to_string());
|
|
tokio::spawn(async move {
|
|
windmill_native_triggers::rename::reregister_triggers_after_rename(
|
|
&db, &authed, &w_id, &moved,
|
|
)
|
|
.await;
|
|
});
|
|
}
|
|
#[cfg(not(feature = "native_trigger"))]
|
|
let _ = (db, authed, w_id, moved);
|
|
}
|
|
|
|
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(ef): Json<EditFlow>,
|
|
) -> 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();
|
|
// The URL identifies the flow being updated; the body path is only needed to rename.
|
|
let mut nf = ef.into_new_flow(flow_path);
|
|
// A `<= 0` flow timeout is "unset", not a 0-second limit (see create_flow).
|
|
nf.timeout = windmill_common::runnable_settings::none_if_non_positive(nf.timeout);
|
|
check_scopes(&authed, || format!("flows:write:{}", flow_path))?;
|
|
// A rename writes the destination as much as the source, so a path-scoped token needs both.
|
|
// Checking only the source would let it move a flow onto a path it has no say over — and
|
|
// everything that follows the rename, native triggers included, is then acting on a path this
|
|
// caller was never authorized for. `create_script` already scopes against its destination.
|
|
if nf.path != flow_path {
|
|
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?;
|
|
|
|
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());
|
|
let resolved_on_behalf_of =
|
|
windmill_common::resolve_on_behalf_of(
|
|
nf.on_behalf_of_email.as_deref(),
|
|
nf.on_behalf_of.as_deref(),
|
|
nf.preserve_on_behalf_of.unwrap_or(false),
|
|
&authed,
|
|
&w_id,
|
|
&db,
|
|
)
|
|
.await?;
|
|
// Written beside the principal only while a worker that still reads it may be live.
|
|
let legacy_on_behalf_of_email =
|
|
windmill_common::legacy_on_behalf_of_email(resolved_on_behalf_of.as_deref(), &w_id, &db)
|
|
.await?;
|
|
|
|
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,
|
|
ws_error_handler_muted = $7,
|
|
value = $8,
|
|
schema = $9::text::json,
|
|
edited_by = $10,
|
|
edited_at = now(),
|
|
labels = COALESCE($13, labels),
|
|
on_behalf_of = $14,
|
|
on_behalf_of_email = $15
|
|
WHERE
|
|
path = $11 AND workspace_id = $12",
|
|
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),
|
|
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]>,
|
|
resolved_on_behalf_of,
|
|
legacy_on_behalf_of_email,
|
|
)
|
|
.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, 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, 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?;
|
|
}
|
|
}
|
|
|
|
let mut moved_native_triggers = Vec::new();
|
|
if is_new_path {
|
|
moved_native_triggers = 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(
|
|
resolved_on_behalf_of.as_deref(),
|
|
nf.preserve_on_behalf_of.unwrap_or(false),
|
|
&authed,
|
|
&windmill_common::users::username_to_permissioned_as(&authed.username),
|
|
) {
|
|
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(),
|
|
authed.username_override.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?;
|
|
|
|
// See `tally_rename_vacated_path`.
|
|
if flow_path != nf.path {
|
|
if let Err(e) = windmill_git_sync::tally_rename_vacated_path(
|
|
&db,
|
|
&w_id,
|
|
DeployedObject::Flow { path: flow_path.to_string(), parent_path: None, version },
|
|
)
|
|
.await
|
|
{
|
|
tracing::error!(%e, "error tallying the path renamed away from");
|
|
}
|
|
}
|
|
|
|
reregister_moved_native_triggers(&db, &authed, &w_id, moved_native_triggers);
|
|
|
|
// 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,
|
|
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,
|
|
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?;
|
|
|
|
// The primary GET is what the CLI reads before a preserving push, so it must carry the
|
|
// derived address: without it the push sends neither half and the identity is cleared.
|
|
let mut flow_o = flow_o;
|
|
if let Some(fws) = flow_o.as_mut() {
|
|
fws.flow.on_behalf_of_email = derived_on_behalf_of_email(&db, &w_id, &fws.flow).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> {
|
|
if authed.is_operator {
|
|
return Err(Error::NotAuthorized(
|
|
"Operators cannot archive flows for security reasons".to_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> {
|
|
if authed.is_operator {
|
|
return Err(Error::NotAuthorized(
|
|
"Operators cannot delete flows for security reasons".to_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());
|
|
}
|
|
}
|