mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-09-05 00:03:08 +00:00
* feat(forks): strip operational state from triggers/schedules on git-sync export When the source workspace is a fork (`wm-fork-*`), the tarball export now omits `mode` from triggers and `enabled` from schedules. The trigger update handler also preserves the existing DB `mode` when both fields are absent from the request, instead of falling back to the BaseTriggerData default. This prevents a fork's git-sync round-trip from flipping the parent workspace's enabled/disabled state when a merge applies the fork's YAML back to main. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * feat(forks): opt-in fork_triggers flag clones triggers/schedules disabled Adds `workspace.fork_triggers` (default false) and a matching field on CreateWorkspaceFork. When the user opts in, fork creation also runs clone_triggers_and_schedules: every row in schedule and the ten *_trigger tables is copied to the fork with mode='disabled' / enabled=false. Listener identifiers (group_id, replication_slot_name, subscription_name, …) are copied verbatim — the runtime suffix that prevents the fork from competing with the parent ships in a follow-up PR. native_trigger is intentionally skipped: those triggers manage external webhook state we don't want duplicated. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * feat(forks): warn before enabling triggers/schedules that conflict with parent set_trigger_mode and schedule's set_enabled now check whether the parent workspace has the same path actively enabled. If so, the call is rejected with a `fork-conflict:<kind>:<parent_id>` error unless the request includes `force=true`. The frontend interprets the prefix to surface a confirm-to- proceed dialog. This is the placeholder safety net until the Phase 3 listener-suffix work removes the conflict for the namespaceable kinds (Kafka/MQTT/NATS/Postgres/ Azure/GCP-CreateNew). For SQS, GCP-Existing, and schedules — where there's no namespacing fix — the warning is the durable solution. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * feat(forks): UI: opt-in clone-triggers checkbox + confirm-on-fork-conflict Adds the user-facing surface for the fork-trigger work: - CreateWorkspaceInner: new "Clone triggers and schedules" toggle in the fork-creation dialog (default off). Sends fork_triggers in the request. - forkConflict utility: detects the `fork-conflict:<kind>:<parent_id>` error string from the backend, shows a confirm() dialog explaining why the action is blocked, retries with `force: true` if accepted. - Wires withForkConflictRetry into every trigger setMode and the schedule setEnabled call, both in the per-kind editor components and the +page.svelte list views (HTTP, websocket, kafka, NATS, SQS, MQTT, GCP, Azure, Postgres, email, schedule). OpenAPI spec gains the `force` field on each setmode/setenabled body. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * feat(forks): CLI --fork-triggers flag, fork-trigger docs, skill update - Adds --fork-triggers boolean to wmill workspace fork; passes fork_triggers through to the create_fork API call. - New docs/fork-triggers.md describing the model end-to-end (default, opt-in clone, merge-direction filter, conflict warning, future runtime-suffix work). - Updates the adding-a-trigger SKILL.md to mention the fork-export ignore-keys participation and the clone_triggers_and_schedules block that new trigger kinds must extend. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * chore: regenerate sqlx offline query cache for fork-trigger SQL * fix(forks): replace browser confirm() with ConfirmationModal for fork conflict The fork-conflict warning previously used the browser's native confirm() which doesn't match Windmill's design system. Switches to a singleton ConfirmationModal mounted at the (logged) layout root, driven by a new forkConflictModal store. The withForkConflictRetry helper now sets the store and awaits the user's choice via a Promise, instead of blocking on window.confirm. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * fix(forks): filter unchanged triggers in merge UI, add diff view, surface parent-only ones The fork merge UI listed every trigger from the fork as a deployable item regardless of whether it differed from the parent — so a fork created with fork_triggers=true (which clones triggers in disabled state, otherwise identical) showed every trigger as a "Fork-only" change. The 'Update current' tab also missed triggers newly created in the parent that the fork hadn't pulled yet. This refactor: - fetchAllTriggers now lists both fork and parent in parallel for each trigger kind, then merges by path. - Computes a per-trigger `changeKind` (new / modified / deleted-in-source) using a JSON comparison that strips runtime + fork-local fields (mode/enabled/server_id/last_server_ping/edited_at/edited_by/etc.) so the disabled-on-clone difference doesn't show up as a change. - Filters the trigger items in deployableItems by the current direction: Deploy mode shows fork-side new/modified, Update mode shows parent-side new/modified. - Replaces the always-on "Fork-only" badge with proper New/Modified badges and surfaces a Diff button (modal Drawer + Monaco DiffEditor) for modified triggers — the diff strips the same ignored fields so users see only the meaningful config differences. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * fix(forks): always clone triggers/schedules disabled, drop opt-in flag Disabled triggers and schedules are inert — no listener attaches, no cron fires — so cloning them by default is safe by construction. Drops the fork_triggers opt-in flag introduced earlier in this PR: - Drops workspace.fork_triggers column (migration removed) - Removes fork_triggers from CreateWorkspaceFork (API + OpenAPI) - Removes the conditional in create_workspace_fork — clone always runs - Removes the toggle from the fork-creation dialog - Removes --fork-triggers from `wmill workspace fork` - Updates docs/fork-triggers.md and adding-a-trigger SKILL.md The merge UI continues to exclude triggers from the deploy/update default selection, so a routine merge from a fork doesn't accidentally push trigger config the user hasn't intentionally changed. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * fix(http-triggers): scope route exists check by workspace, skip non-workspaced clones in forks The non-CLOUD branch of `route_path_key_exists` self-excluded by trigger path alone, which silently masked cross-workspace collisions once forks started cloning trigger rows verbatim. Tighten it to exclude only the exact `(workspace_id, path)` row. Fork creation also now skips non-workspaced HTTP triggers — their URL has no workspace prefix, so a clone collides with the parent at the matchit router (which silently drops one of two duplicates) and there is no namespacing escape hatch. The clone copies all rows when CLOUD_HOSTED or HTTP_ROUTE_WORKSPACED_ROUTE forces every route workspaced regardless. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * fix(forks-ui): silent cancel on enable conflict, clean up trigger rows in compare view forkConflict helper now returns undefined when the user dismisses the modal instead of throwing, so the redundant 'Cannot enable: undefined' toast no longer appears. CompareWorkspaces trigger rows now mirror the script row layout: drop the redundant Disabled badge and the Trash/Details buttons (both belong on the dedicated trigger pages, not in the deploy/compare view); pass triggerKind through so RowIcon picks the right kind-specific icon; move extraLabel into the summary line; replace the yellow Modified badge with the same green ↗ ahead / blue ↘ behind treatment scripts use. Trigger diff drawer: switch JSON → YAML for parity with DiffDrawer, fix zero-height monaco render with className=!h-full, drop the redundant Original/Modified label banner above the diff. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * fix(email-trigger): scope local_part exists check, skip non-workspaced clones in forks Mirrors the HTTP route fix for the email-trigger non-CLOUD `email_exists` check (in EE) which had the same path-only self-exclusion bug, and the fork clone of `email_trigger` rows which copied non-workspaced `local_part` verbatim. Skip non-workspaced rows in the clone unless the instance is CLOUD_HOSTED (where lookup is workspace-scoped natively). EE companion change in windmill-trigger-email/src/handler_ee.rs. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * chore: update ee-repo-ref to 78512dd73b4a1c9f70574cff863374179e3a621b This commit updates the EE repository reference after PR #554 was merged in windmill-ee-private. Previous ee-repo-ref: 1ac77f50747b58e720a11162dfd309bc252a24ab New ee-repo-ref: 78512dd73b4a1c9f70574cff863374179e3a621b Automated by sync-ee-ref workflow. * fix(forks): always-warn on parent row, kind-specific modal copy, cancel-aware toggles - Conflict check now fires whenever the parent has the path (regardless of parent's mode), since the cloned upstream identifier is shared by construction; closes the Postgres slot-takeover gap when the parent is disabled. Schedule's set_schedule_enabled gets the same treatment. - Skip the warning entirely for HTTP and Email via a new TriggerCrud::FORK_CONFLICT_ON_ENABLE const — both kinds are workspace- scoped at runtime so cloned rows can't collide with the parent. - Modal copy branches by failure family: split-events (Kafka/NATS/MQTT/SQS/ GCP/Azure), duplicate-firing (Websocket/Schedule), slot-takeover (Postgres). Generic fallback for unknown kinds. - withForkConflictRetry now returns boolean (true=committed, false= cancelled). TriggerModeToggle reuses its existing innerTriggerMode local state via a function binding for the regular Toggle, snapping back to the prop when onToggleMode signals a cancel — needed because the native bind:checked diverges from the parent's prop after a click and Svelte's reactivity won't re-push a same-valued prop down. Schedule list page uses {#key} on a reset version since it renders Toggle directly. - Editor inners revert mode = previousMode on cancel; list pages skip the re-fetch (loadTriggers/loadSchedules) on cancel to avoid pointless network traffic and the schedule "Job stats loading..." flash. - Drop withForkConflictRetry from HTTP and Email editors + list pages since the backend never emits the conflict for those kinds. * fix(forks-ui): widen onToggleMode types, scope schedule toggle reset by path - TriggerEditorToolbar and TriggerSuspendedJobsModal forwarded onToggleMode as `(mode) => void`, dropping the new boolean return so any caller wired through them would silently no-op the cancel-revert. Match the wider TriggerModeToggle signature. - Schedule list page used a single resetVersion counter for every row's {#key}, so cancelling on any one schedule remounted every <Toggle> on the page. Switch to a per-path Record<string, number> bumped only for the affected row. * chore: bump ee-repo-ref to c3a4553 (email FORK_CONFLICT_ON_ENABLE override) * fix(forks): include Suspended in conflict gate, use parent_workspace_id for fork detection Three fixes from the Claude review on PR #8976: - Suspended mode still attaches the listener (it just pauses auto-run of queued jobs); two suspended fork+parent listeners would still split Kafka events / share a PG slot. Gate set_trigger_mode on `mode != Disabled` instead of `mode == Enabled` so Suspended also surfaces the warning. - workspaces_export.rs::fork_*_ignore_keys keyed off the wm-fork-* prefix while set_trigger_mode and set_schedule_enabled key off parent_workspace_id. Switch the export filter to query parent_workspace_id once at the top of tarball_workspace and pass is_fork through. The column is the contract; the prefix is a creation-time naming convention that could in principle drift. - TriggerModeToggle's suspend-dropdown action reassigned the non-bindable `triggerMode` prop instead of the local `innerTriggerMode` mirror, leaking inconsistent state if the dispatch was cancelled. Now writes to innerTriggerMode like the Toggle's on:change handler does. * fix(cli): skip setScheduleEnabled when local YAML lacks `enabled` Tarball export from a fork strips `enabled` from schedules so the fork→parent git-sync round-trip can't flip the parent's operational state. The CLI's pushSchedule called setScheduleEnabled whenever `localSchedule.enabled != schedule.enabled`, which evaluates truthy when local is undefined (fork-pulled YAML) and remote is true/false — sending `{ enabled: undefined }` that serializes to `{}` and gets rejected by the backend (`SetEnabled.enabled` is required). Skip the call when `localSchedule.enabled === undefined` so a sync push of fork-pulled YAMLs preserves the target's existing enabled state instead of erroring out. Trigger updates were already safe — the backend's update_trigger preserves `mode` when the request omits it. * Revert "fix(cli): skip setScheduleEnabled when local YAML lacks `enabled`" This reverts commit23ba7e72fc. * feat(cli): --force flag and friendlier error on fork-conflict for schedule enable `wmill schedule enable foo/bar` against a fork whose parent has the same path used to surface the raw `fork-conflict:schedule:<parent>` error body. The CLI now: - accepts `--force` to bypass the warning (mirrors the API field and the UI's "Enable anyway" confirmation), - detects the `fork-conflict:` prefix on errors and prints a one-screen explanation pointing at --force instead of the raw body. Disable doesn't trigger the warning (the gate fires only on transitions to listener-attaching modes), so no flag there. Trigger enable/disable isn't exposed as a standalone CLI command — sync push goes through updateTrigger which has its own backend mode-preservation, so no fork-conflict surfaces from the CLI for those. * chore: regenerate cli-commands docs after adding --force to schedule enable * chore: update ee-repo-ref to 967f961f0a88b027d894aebd03977181129477a8 This commit updates the EE repository reference after PR #555 was merged in windmill-ee-private. Previous ee-repo-ref: c3a4553296473932e15392a06415dd7fb9aa6591 New ee-repo-ref: 967f961f0a88b027d894aebd03977181129477a8 Automated by sync-ee-ref workflow. * fix(forks): address CI dead-code, claude/cubic review feedback - backend: cfg-gate `fork_trigger_ignore_keys` to match its already-gated callsite. CI compiles with `-D warnings`, so the unused-fn under feature combos that disable all trigger crates was breaking check_oss/check_ee/ cargo_test/test-linux/test-windows. - cli: re-apply the `pushSchedule` undefined-skip (originally23ba7e7, reverted in4d172a1). Tarball export from forks strips `enabled`, so fork-pulled YAMLs that get sync-pushed back via `wmill schedule push` would otherwise serialize `{ enabled: undefined }` → `{}` and the backend's required `SetEnabled.enabled` rejects the body. Skipping preserves the target's existing flag, which is the round-trip-safe behavior. (`wmill workspace merge` extension to triggers/schedules is tracked in #9001 — until then sync push is the only CLI path.) - TriggerModeToggle suspend-dropdown action awaits onToggleMode and resets `innerTriggerMode = triggerMode` on cancel, matching the Toggle on:change handler. Without this, dismissing the fork-conflict modal on a Suspend transition leaves the toggle stuck in 'suspended'. - forkConflict: when a new modal opens with a previous resolver still pending, resolve the older promise to false. Avoids a dangling promise if the user clicks toggles on two rows in quick succession. - schedules list: bump `toggleResetVersions[path]` on the permission-denied branch so the Toggle re-mounts back to the prop's `enabled` value. Without this, a user without write permission could click the toggle and have it stick visually flipped. - docs/fork-triggers.md: switch the merge-direction filter description from `wm-fork-*` prefix to `parent_workspace_id IS NOT NULL` (matches the code after4dd38fe). Drop the misleading "merge-direction filter strips identifier columns too" line in Future Work — the runtime suffix is applied at listener attach, the stored column never carries it, so no export filtering is needed there. --------- Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com> Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
1260 lines
46 KiB
Rust
1260 lines
46 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 crate::db::ApiAuthed;
|
|
|
|
use crate::{apps::AppWithLastVersion, db::DB, folders::Folder};
|
|
|
|
#[cfg(any(
|
|
feature = "http_trigger",
|
|
feature = "websocket",
|
|
feature = "postgres_trigger",
|
|
feature = "mqtt_trigger",
|
|
all(
|
|
feature = "enterprise",
|
|
any(
|
|
feature = "kafka",
|
|
feature = "sqs_trigger",
|
|
feature = "gcp_trigger",
|
|
feature = "azure_trigger",
|
|
feature = "nats",
|
|
feature = "smtp",
|
|
),
|
|
feature = "private"
|
|
)
|
|
))]
|
|
use crate::triggers::TriggerCrud;
|
|
|
|
use axum::{
|
|
extract::{Extension, Path, Query},
|
|
response::IntoResponse,
|
|
};
|
|
|
|
use http::HeaderName;
|
|
use itertools::Itertools;
|
|
|
|
use windmill_common::runnable_settings::{ConcurrencySettings, DebouncingSettings};
|
|
use windmill_common::scripts::ScriptRunnableSettingsHandle;
|
|
use windmill_common::utils::require_admin;
|
|
use windmill_common::variables::decrypt;
|
|
use windmill_common::worker::WINDMILL_DIR;
|
|
use windmill_common::{
|
|
db::UserDB,
|
|
error::{to_anyhow, Error, Result},
|
|
flows::Flow,
|
|
schedule::Schedule,
|
|
scripts::{Schema, Script, ScriptLang},
|
|
variables::{build_crypt, ExportableListableVariable},
|
|
workspace_dependencies::WorkspaceDependencies,
|
|
};
|
|
|
|
use hyper::header;
|
|
use serde::{Deserialize, Serialize};
|
|
use serde_json::Value;
|
|
use tempfile::TempDir;
|
|
use tokio::fs::File;
|
|
use tokio_util::io::ReaderStream;
|
|
use windmill_store::resources::{Resource, ResourceType};
|
|
|
|
#[derive(Serialize)]
|
|
struct ScriptMetadata {
|
|
summary: String,
|
|
description: String,
|
|
schema: Option<Schema>,
|
|
lock: Option<String>,
|
|
kind: String,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
envs: Option<Vec<String>>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
cache_ttl: Option<i32>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
dedicated_worker: Option<bool>,
|
|
#[serde(skip_serializing_if = "is_none_or_false")]
|
|
ws_error_handler_muted: Option<bool>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
priority: Option<i16>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
tag: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
pub timeout: Option<i32>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
pub delete_after_secs: Option<i32>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
pub restart_unless_cancelled: Option<bool>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
pub visible_to_runner_only: Option<bool>,
|
|
// auto_kind is intentionally excluded from export — it is auto-detected by the
|
|
// parser at deploy time from the script content (workflow/task patterns for "wac",
|
|
// no main function for "lib").
|
|
#[serde(skip_serializing)]
|
|
#[allow(dead_code)]
|
|
pub auto_kind: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
pub codebase: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
pub has_preprocessor: Option<bool>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
pub on_behalf_of_email: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
pub modules: Option<std::collections::HashMap<String, windmill_common::scripts::ScriptModule>>,
|
|
#[serde(flatten)]
|
|
pub concurrency_settings: ConcurrencySettings,
|
|
#[serde(flatten)]
|
|
pub debouncing_settings: DebouncingSettings,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
pub labels: Option<Vec<String>>,
|
|
}
|
|
|
|
pub fn is_none_or_false(val: &Option<bool>) -> bool {
|
|
match val {
|
|
Some(val) => !val,
|
|
None => true,
|
|
}
|
|
}
|
|
|
|
/// Returns the keys to strip from trigger/schedule serialization when the
|
|
/// source workspace is a fork. Stripping these keys avoids propagating
|
|
/// fork-local operational state (enabled flag, runtime listener identifiers)
|
|
/// back to the parent workspace through the git-sync round-trip.
|
|
#[cfg(any(
|
|
feature = "http_trigger",
|
|
feature = "websocket",
|
|
feature = "postgres_trigger",
|
|
feature = "mqtt_trigger",
|
|
feature = "native_trigger",
|
|
all(
|
|
feature = "enterprise",
|
|
any(
|
|
feature = "kafka",
|
|
feature = "sqs_trigger",
|
|
feature = "gcp_trigger",
|
|
feature = "azure_trigger",
|
|
feature = "nats",
|
|
feature = "smtp",
|
|
),
|
|
feature = "private"
|
|
)
|
|
))]
|
|
fn fork_trigger_ignore_keys(is_fork: bool) -> Option<Vec<&'static str>> {
|
|
if is_fork {
|
|
Some(vec!["mode", "enabled"])
|
|
} else {
|
|
None
|
|
}
|
|
}
|
|
|
|
fn fork_schedule_ignore_keys(is_fork: bool) -> Option<Vec<&'static str>> {
|
|
if is_fork {
|
|
Some(vec!["enabled"])
|
|
} else {
|
|
None
|
|
}
|
|
}
|
|
|
|
enum ArchiveImpl {
|
|
#[cfg(feature = "zip")]
|
|
Zip(async_zip::tokio::write::ZipFileWriter<tokio::fs::File>),
|
|
Tar(tokio_tar::Builder<File>),
|
|
}
|
|
|
|
impl ArchiveImpl {
|
|
async fn write_to_archive(&mut self, content: &str, path: &str) -> Result<()> {
|
|
match self {
|
|
ArchiveImpl::Tar(t) => {
|
|
let bytes = content.as_bytes();
|
|
let mut header = tokio_tar::Header::new_gnu();
|
|
header.set_size(bytes.len() as u64);
|
|
header.set_mtime(0);
|
|
header.set_uid(0);
|
|
header.set_gid(0);
|
|
header.set_mode(0o777);
|
|
header.set_cksum();
|
|
t.append_data(&mut header, path, bytes).await?;
|
|
}
|
|
#[cfg(feature = "zip")]
|
|
ArchiveImpl::Zip(z) => {
|
|
let header =
|
|
async_zip::ZipEntryBuilder::new(path.into(), async_zip::Compression::Deflate)
|
|
.last_modification_date(Default::default())
|
|
.unix_permissions(0o777)
|
|
.build();
|
|
z.write_entry_whole(header, content.as_bytes())
|
|
.await
|
|
.map_err(to_anyhow)?;
|
|
}
|
|
}
|
|
Ok(())
|
|
}
|
|
async fn finish(self) -> Result<()> {
|
|
match self {
|
|
ArchiveImpl::Tar(t) => t.into_inner().await?,
|
|
#[cfg(feature = "zip")]
|
|
ArchiveImpl::Zip(z) => z.close().await.map_err(to_anyhow)?.into_inner(),
|
|
}
|
|
.sync_all()
|
|
.await?;
|
|
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
#[derive(Deserialize)]
|
|
pub(crate) struct ArchiveQueryParams {
|
|
archive_type: Option<String>,
|
|
plain_secret: Option<bool>,
|
|
plain_secrets: Option<bool>,
|
|
skip_secrets: Option<bool>,
|
|
skip_variables: Option<bool>,
|
|
skip_resources: Option<bool>,
|
|
skip_resource_types: Option<bool>,
|
|
include_schedules: Option<bool>,
|
|
include_triggers: Option<bool>,
|
|
include_users: Option<bool>,
|
|
include_groups: Option<bool>,
|
|
include_settings: Option<bool>,
|
|
include_key: Option<bool>,
|
|
include_workspace_dependencies: Option<bool>,
|
|
default_ts: Option<String>,
|
|
/// Settings format version: "v1" (default) returns legacy flat format, "v2" returns grouped format
|
|
settings_version: Option<String>,
|
|
}
|
|
|
|
#[inline]
|
|
pub fn to_string_without_metadata<T>(
|
|
value: &T,
|
|
preserve_extra_perms: bool,
|
|
ignore_keys: Option<Vec<&str>>,
|
|
) -> Result<String>
|
|
where
|
|
T: ?Sized + Serialize,
|
|
{
|
|
let mut value = serde_json::to_value(value).map_err(to_anyhow)?;
|
|
value
|
|
.as_object_mut()
|
|
.map(|obj| {
|
|
let keys = [
|
|
vec![
|
|
"workspace_id",
|
|
"path",
|
|
"name",
|
|
"versions",
|
|
"id",
|
|
"created_at",
|
|
"updated_at",
|
|
"created_by",
|
|
"updated_by",
|
|
"edited_at",
|
|
"edited_by",
|
|
"permissioned_as",
|
|
"archived",
|
|
"has_draft",
|
|
"error",
|
|
"last_server_ping",
|
|
"server_id",
|
|
"raw_app",
|
|
],
|
|
ignore_keys.unwrap_or(vec![]),
|
|
]
|
|
.concat();
|
|
|
|
for key in keys {
|
|
if obj.contains_key(key) {
|
|
obj.remove(key);
|
|
}
|
|
}
|
|
|
|
if let Some(o2) = obj.get_mut("policy").and_then(|x| x.as_object_mut()) {
|
|
o2.remove("on_behalf_of");
|
|
o2.remove("on_behalf_of_email");
|
|
}
|
|
if !preserve_extra_perms && obj.contains_key("extra_perms") {
|
|
obj.remove("extra_perms");
|
|
}
|
|
if obj
|
|
.get("default_permissioned_as")
|
|
.and_then(|v| v.as_array())
|
|
.is_some_and(|a| a.is_empty())
|
|
{
|
|
obj.remove("default_permissioned_as");
|
|
}
|
|
|
|
serde_json::to_string_pretty(&obj).ok()
|
|
})
|
|
.flatten()
|
|
.ok_or_else(|| Error::BadRequest("Impossible to serialize value".to_string()))
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
struct SimplifiedUser {
|
|
username: String,
|
|
role: String,
|
|
disabled: bool,
|
|
email: String,
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
struct SimplifiedGroup {
|
|
name: String,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
summary: Option<String>,
|
|
members: Vec<String>,
|
|
admins: Vec<String>,
|
|
}
|
|
|
|
// V2 format: New grouped format
|
|
#[derive(Serialize)]
|
|
struct SimplifiedSettings {
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
auto_invite: Option<Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
webhook: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
deploy_to: Option<String>,
|
|
// Always serialize (including as `null`) so that `wmill sync pull` emits
|
|
// these fields in settings.yaml unconditionally. Makes round-trip
|
|
// bijective: YAML is the source of truth, absence/null = "clear remote",
|
|
// mirroring every other workspace setting.
|
|
error_handler: Option<Value>,
|
|
success_handler: Option<Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
ai_config: Option<serde_json::Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
large_file_storage: Option<Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
git_sync: Option<Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
default_app: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
default_scripts: Option<Value>,
|
|
name: String,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
mute_critical_alerts: Option<bool>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
color: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
operator_settings: Option<serde_json::Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
datatable: Option<Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
slack_team_id: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
slack_name: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
slack_command_script: Option<String>,
|
|
// Always serialize (see note above on error_handler / success_handler).
|
|
slack_oauth_client_id: Option<String>,
|
|
slack_oauth_client_secret: Option<String>,
|
|
}
|
|
|
|
// V1 format: Legacy flat format for backward compatibility (matches main branch exactly)
|
|
#[derive(Serialize)]
|
|
struct SimplifiedSettingsLegacy {
|
|
auto_invite_enabled: bool,
|
|
auto_invite_as: String,
|
|
auto_invite_mode: String,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
webhook: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
deploy_to: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
error_handler: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
error_handler_extra_args: Option<Value>,
|
|
error_handler_muted_on_cancel: bool,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
ai_config: Option<serde_json::Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
large_file_storage: Option<Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
git_sync: Option<Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
default_app: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
default_scripts: Option<Value>,
|
|
name: String,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
mute_critical_alerts: Option<bool>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
color: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
operator_settings: Option<serde_json::Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
datatable: Option<Value>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
slack_team_id: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
slack_name: Option<String>,
|
|
#[serde(skip_serializing_if = "Option::is_none")]
|
|
slack_command_script: Option<String>,
|
|
}
|
|
|
|
// Internal struct for querying database
|
|
#[derive(sqlx::FromRow)]
|
|
struct SettingsRow {
|
|
auto_invite: Option<Value>,
|
|
webhook: Option<String>,
|
|
deploy_to: Option<String>,
|
|
error_handler: Option<Value>,
|
|
success_handler: Option<Value>,
|
|
ai_config: Option<serde_json::Value>,
|
|
large_file_storage: Option<Value>,
|
|
git_sync: Option<Value>,
|
|
default_app: Option<String>,
|
|
default_scripts: Option<Value>,
|
|
name: Option<String>,
|
|
mute_critical_alerts: Option<bool>,
|
|
color: Option<String>,
|
|
operator_settings: Option<serde_json::Value>,
|
|
datatable: Option<Value>,
|
|
slack_team_id: Option<String>,
|
|
slack_name: Option<String>,
|
|
slack_command_script: Option<String>,
|
|
slack_oauth_client_id: Option<String>,
|
|
slack_oauth_client_secret: Option<String>,
|
|
}
|
|
|
|
pub(crate) async fn tarball_workspace(
|
|
authed: ApiAuthed,
|
|
Extension(user_db): Extension<UserDB>,
|
|
Extension(db): Extension<DB>,
|
|
Path(w_id): Path<String>,
|
|
Query(ArchiveQueryParams {
|
|
archive_type,
|
|
plain_secret,
|
|
plain_secrets,
|
|
skip_resources,
|
|
skip_resource_types,
|
|
skip_secrets,
|
|
skip_variables,
|
|
include_schedules,
|
|
include_triggers,
|
|
include_users,
|
|
include_groups,
|
|
include_settings,
|
|
include_key,
|
|
include_workspace_dependencies,
|
|
default_ts,
|
|
settings_version,
|
|
}): Query<ArchiveQueryParams>,
|
|
) -> Result<([(HeaderName, String); 2], impl IntoResponse)> {
|
|
tracing::info!(
|
|
"tarball_workspace called for workspace {}: include_workspace_dependencies={:?}, skip_variables={:?}, skip_resources={:?}",
|
|
w_id,
|
|
include_workspace_dependencies,
|
|
skip_variables,
|
|
skip_resources
|
|
);
|
|
|
|
let mut tx = user_db.begin(&authed).await?;
|
|
|
|
// Source-of-truth check for fork-ness: the workspace's parent_workspace_id
|
|
// column. The wm-fork-* prefix is a creation-time naming convention that
|
|
// could in principle drift (rename, manual SQL); the column is the
|
|
// contract that matches what the conflict-warning gates read.
|
|
let is_fork: bool = sqlx::query_scalar!(
|
|
"SELECT parent_workspace_id IS NOT NULL FROM workspace WHERE id = $1",
|
|
&w_id
|
|
)
|
|
.fetch_optional(&mut *tx)
|
|
.await?
|
|
.flatten()
|
|
.unwrap_or(false);
|
|
|
|
let tmp_dir = TempDir::new_in(&*WINDMILL_DIR)?;
|
|
|
|
let name = match archive_type.as_deref() {
|
|
Some("tar") | None => Ok(format!("windmill-{w_id}.tar")),
|
|
Some("zip") => Ok(format!("windmill-{w_id}.zip")),
|
|
Some(t) => Err(Error::BadRequest(format!("Invalid Archive Type {t}"))),
|
|
}?;
|
|
let file_path = tmp_dir.path().join(&name);
|
|
let mut archive = match archive_type.as_deref() {
|
|
Some("tar") | None => {
|
|
let file = File::create(&file_path).await?;
|
|
Ok(ArchiveImpl::Tar(tokio_tar::Builder::new(file)))
|
|
}
|
|
#[cfg(feature = "zip")]
|
|
Some("zip") => {
|
|
let file = tokio::fs::File::create(&file_path).await?;
|
|
Ok(ArchiveImpl::Zip(
|
|
async_zip::tokio::write::ZipFileWriter::with_tokio(file),
|
|
))
|
|
}
|
|
Some(t) => Err(Error::BadRequest(format!("Invalid Archive Type {t}"))),
|
|
}?;
|
|
{
|
|
let folders = sqlx::query_as::<_, Folder>("SELECT * FROM folder WHERE workspace_id = $1")
|
|
.bind(&w_id)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
for folder in folders {
|
|
archive
|
|
.write_to_archive(
|
|
&to_string_without_metadata(&folder, true, None).unwrap(),
|
|
&format!("f/{}/folder.meta.json", folder.name),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
{
|
|
let scripts = sqlx::query_as::<_, Script<ScriptRunnableSettingsHandle>>(
|
|
"SELECT * FROM script as o WHERE workspace_id = $1 AND archived = false
|
|
AND (draft_only IS NULL OR draft_only = false)
|
|
AND created_at = (select max(created_at) from script where path = o.path AND \
|
|
workspace_id = $1)",
|
|
)
|
|
.bind(&w_id)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
for script in scripts {
|
|
let script = windmill_common::scripts::prefetch_cached_script(script, &db).await?;
|
|
let ext = match script.language {
|
|
ScriptLang::Python3 => "py",
|
|
ScriptLang::Deno => {
|
|
if default_ts.as_ref().is_some_and(|x| x == "bun") {
|
|
"deno.ts"
|
|
} else {
|
|
"ts"
|
|
}
|
|
}
|
|
ScriptLang::Go => "go",
|
|
ScriptLang::Bash => "sh",
|
|
ScriptLang::Powershell => "ps1",
|
|
ScriptLang::Postgresql => "pg.sql",
|
|
ScriptLang::Mysql => "my.sql",
|
|
ScriptLang::Bigquery => "bq.sql",
|
|
ScriptLang::Snowflake => "sf.sql",
|
|
ScriptLang::Mssql => "ms.sql",
|
|
ScriptLang::DuckDb => "duckdb.sql",
|
|
ScriptLang::Graphql => "gql",
|
|
ScriptLang::Nativets => "fetch.ts",
|
|
ScriptLang::Bun | ScriptLang::Bunnative => {
|
|
if default_ts.as_ref().is_some_and(|x| x == "bun") {
|
|
"ts"
|
|
} else {
|
|
"bun.ts"
|
|
}
|
|
}
|
|
ScriptLang::Php => "php",
|
|
ScriptLang::Rust => "rs",
|
|
ScriptLang::Ansible => "playbook.yml",
|
|
ScriptLang::CSharp => "cs",
|
|
ScriptLang::Nu => "nu",
|
|
ScriptLang::OracleDB => "odb.sql",
|
|
ScriptLang::Java => "java",
|
|
ScriptLang::Ruby => "rb",
|
|
ScriptLang::Rlang => "r",
|
|
// for related places search: ADD_NEW_LANG
|
|
};
|
|
archive
|
|
.write_to_archive(&script.content, &format!("{}.{}", script.path, ext))
|
|
.await?;
|
|
|
|
let metadata = ScriptMetadata {
|
|
summary: script.summary,
|
|
description: script.description,
|
|
schema: script.schema,
|
|
kind: script.kind.to_string(),
|
|
lock: script.lock,
|
|
envs: script.envs,
|
|
concurrency_settings: script.runnable_settings.concurrency_settings,
|
|
debouncing_settings: script.runnable_settings.debouncing_settings,
|
|
cache_ttl: script.cache_ttl,
|
|
dedicated_worker: script.dedicated_worker,
|
|
ws_error_handler_muted: script.ws_error_handler_muted,
|
|
priority: script.priority,
|
|
tag: script.tag,
|
|
timeout: script.timeout,
|
|
delete_after_secs: script.delete_after_secs,
|
|
restart_unless_cancelled: script.restart_unless_cancelled,
|
|
visible_to_runner_only: script.visible_to_runner_only,
|
|
auto_kind: script.auto_kind,
|
|
codebase: script.codebase,
|
|
has_preprocessor: script.has_preprocessor,
|
|
on_behalf_of_email: script.on_behalf_of_email,
|
|
modules: script.modules,
|
|
labels: script.labels,
|
|
};
|
|
let metadata_str = serde_json::to_string_pretty(&metadata).unwrap();
|
|
archive
|
|
.write_to_archive(&metadata_str, &format!("{}.script.json", script.path))
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
if !skip_resources.unwrap_or(false) {
|
|
let resources = sqlx::query_as!(
|
|
Resource,
|
|
"SELECT * FROM resource WHERE workspace_id = $1 AND resource_type != 'state' AND resource_type != 'cache'",
|
|
&w_id
|
|
)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
for resource in resources {
|
|
let resource_str = &to_string_without_metadata(&resource, false, None).unwrap();
|
|
archive
|
|
.write_to_archive(&resource_str, &format!("{}.resource.json", resource.path))
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
if !skip_resource_types.unwrap_or(false) {
|
|
let resource_types = sqlx::query_as!(
|
|
ResourceType,
|
|
"SELECT * FROM resource_type WHERE workspace_id = $1",
|
|
&w_id
|
|
)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
for resource_type in resource_types {
|
|
let resource_str = &to_string_without_metadata(&resource_type, false, None).unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&resource_str,
|
|
&format!("{}.resource-type.json", resource_type.name),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
{
|
|
let flows = sqlx::query_as::<_, Flow>(
|
|
"SELECT flow.workspace_id, flow.path, flow.summary, flow.description, flow.archived, flow.extra_perms, flow.draft_only, 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.id = flow.versions[array_upper(flow.versions, 1)]
|
|
WHERE flow.workspace_id = $1 AND flow.archived = false AND (flow.draft_only IS NULL OR flow.draft_only = false)",
|
|
)
|
|
.bind(&w_id)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
for flow in flows {
|
|
let flow_str = &to_string_without_metadata(&flow, false, None).unwrap();
|
|
archive
|
|
.write_to_archive(&flow_str, &format!("{}.flow.json", flow.path))
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
if !skip_variables.unwrap_or(false) {
|
|
let variables =
|
|
sqlx::query_as::<_, ExportableListableVariable>(if !skip_secrets.unwrap_or(false) {
|
|
"SELECT * FROM variable WHERE workspace_id = $1 AND expires_at IS NULL"
|
|
} else {
|
|
"SELECT * FROM variable WHERE workspace_id = $1 AND is_secret = false AND expires_at IS NULL"
|
|
})
|
|
.bind(&w_id)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
let mc = build_crypt(&db, &w_id).await?;
|
|
|
|
for mut var in variables {
|
|
if plain_secret.or(plain_secrets).unwrap_or(false)
|
|
&& var.value.is_some()
|
|
&& var.is_secret
|
|
{
|
|
var.value = Some(decrypt(&mc, var.value.unwrap()).map_err(|e| {
|
|
Error::internal_err(format!("Error decrypting variable {}: {}", var.path, e))
|
|
})?);
|
|
}
|
|
let var_str = &to_string_without_metadata(&var, false, None).unwrap();
|
|
archive
|
|
.write_to_archive(&var_str, &format!("{}.variable.json", var.path))
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
{
|
|
let apps = sqlx::query_as::<_, AppWithLastVersion>(
|
|
"SELECT app.id, app.path, app.summary, app.versions, app.policy, app.custom_path,
|
|
app.extra_perms, app_version.value,
|
|
app_version.created_at, app_version.created_by, app_version.raw_app, app.labels from app, app_version
|
|
WHERE app.workspace_id = $1 AND app_version.id = app.versions[array_upper(app.versions, 1)]
|
|
AND (app.draft_only IS NULL OR app.draft_only = false)",
|
|
)
|
|
.bind(&w_id)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
for app in apps {
|
|
let app_str = &to_string_without_metadata(&app, false, None).unwrap();
|
|
let kind = if app.raw_app { "raw_app" } else { "app" };
|
|
archive
|
|
.write_to_archive(&app_str, &format!("{}.{}.json", app.path, kind))
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
if include_workspace_dependencies.unwrap_or(false)
|
|
&& require_admin(authed.is_admin, &authed.username).is_ok()
|
|
{
|
|
tracing::info!("Including workspace dependencies in tarball export");
|
|
let workspace_dependencies = WorkspaceDependencies::list(&w_id, &db).await?;
|
|
tracing::info!(
|
|
"Found {} workspace dependencies",
|
|
workspace_dependencies.len()
|
|
);
|
|
for dep in workspace_dependencies {
|
|
// let dep_str = &to_string_without_metadata(&dep, false, None).unwrap();
|
|
let filename = WorkspaceDependencies::to_path(&dep.name, dep.language)?;
|
|
tracing::info!(
|
|
"Adding workspace dependency: name={:?}, language={:?}, filename={}",
|
|
dep.name,
|
|
dep.language,
|
|
filename
|
|
);
|
|
archive.write_to_archive(&dep.content, &filename).await?;
|
|
}
|
|
} else {
|
|
tracing::info!(
|
|
"Skipping workspace dependencies: include_workspace_dependencies={:?}",
|
|
include_workspace_dependencies
|
|
);
|
|
}
|
|
|
|
if include_schedules.unwrap_or(false) {
|
|
let schedules = sqlx::query_as::<_, Schedule>(
|
|
"SELECT * FROM schedule
|
|
WHERE workspace_id = $1",
|
|
)
|
|
.bind(&w_id)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
let schedule_ignore_keys = fork_schedule_ignore_keys(is_fork);
|
|
for schedule in schedules {
|
|
let app_str =
|
|
&to_string_without_metadata(&schedule, false, schedule_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(&app_str, &format!("{}.schedule.json", schedule.path))
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
if include_triggers.unwrap_or(false) {
|
|
#[cfg(any(
|
|
feature = "http_trigger",
|
|
feature = "websocket",
|
|
feature = "postgres_trigger",
|
|
feature = "mqtt_trigger",
|
|
feature = "native_trigger",
|
|
all(
|
|
feature = "enterprise",
|
|
any(
|
|
feature = "kafka",
|
|
feature = "sqs_trigger",
|
|
feature = "gcp_trigger",
|
|
feature = "azure_trigger",
|
|
feature = "nats",
|
|
feature = "smtp",
|
|
),
|
|
feature = "private"
|
|
)
|
|
))]
|
|
let trigger_ignore_keys = fork_trigger_ignore_keys(is_fork);
|
|
|
|
#[cfg(feature = "http_trigger")]
|
|
{
|
|
use crate::triggers::http::HttpTrigger;
|
|
let handler = HttpTrigger;
|
|
let http_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in http_triggers {
|
|
let trigger_str =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.http_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "websocket")]
|
|
{
|
|
use crate::triggers::websocket::WebsocketTrigger;
|
|
let handler = WebsocketTrigger;
|
|
let websocket_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in websocket_triggers {
|
|
let trigger_str =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.websocket_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(all(feature = "enterprise", feature = "kafka", feature = "private"))]
|
|
{
|
|
use crate::triggers::kafka::KafkaTrigger;
|
|
let handler = KafkaTrigger;
|
|
let kafka_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in kafka_triggers {
|
|
let trigger_str =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.kafka_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(all(feature = "enterprise", feature = "sqs_trigger", feature = "private"))]
|
|
{
|
|
use crate::triggers::sqs::SqsTrigger;
|
|
let handler = SqsTrigger;
|
|
let sqs_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in sqs_triggers {
|
|
let trigger_str =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.sqs_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(all(feature = "enterprise", feature = "gcp_trigger", feature = "private"))]
|
|
{
|
|
use crate::triggers::gcp::GcpTrigger;
|
|
let handler = GcpTrigger;
|
|
let gcp_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in gcp_triggers {
|
|
let trigger_str =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.gcp_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(all(feature = "enterprise", feature = "azure_trigger", feature = "private"))]
|
|
{
|
|
use crate::triggers::azure::AzureTrigger;
|
|
let handler = AzureTrigger;
|
|
let azure_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in azure_triggers {
|
|
let trigger_str =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.azure_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(all(feature = "enterprise", feature = "nats", feature = "private"))]
|
|
{
|
|
use crate::triggers::nats::NatsTrigger;
|
|
let handler = NatsTrigger;
|
|
let nats_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in nats_triggers {
|
|
let trigger_str: &String =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.nats_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "postgres_trigger")]
|
|
{
|
|
use crate::triggers::postgres::PostgresTrigger;
|
|
let handler = PostgresTrigger;
|
|
let postgres_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in postgres_triggers {
|
|
let trigger_str =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.postgres_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "mqtt_trigger")]
|
|
{
|
|
use crate::triggers::mqtt::MqttTrigger;
|
|
let handler = MqttTrigger;
|
|
let mqtt_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in mqtt_triggers {
|
|
let trigger_str =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.mqtt_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(all(feature = "enterprise", feature = "smtp", feature = "private"))]
|
|
{
|
|
use crate::triggers::email::EmailTrigger;
|
|
let handler = EmailTrigger;
|
|
let email_triggers = handler.list_triggers(&mut *tx, &w_id, None).await?;
|
|
|
|
for trigger in email_triggers {
|
|
let trigger_str =
|
|
&to_string_without_metadata(&trigger, false, trigger_ignore_keys.clone())
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!("{}.email_trigger.json", trigger.base.path),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
#[cfg(feature = "native_trigger")]
|
|
{
|
|
use crate::native_triggers::{list_native_triggers, ServiceName};
|
|
use strum::IntoEnumIterator;
|
|
|
|
for service_name in ServiceName::iter() {
|
|
let native_triggers =
|
|
list_native_triggers(&mut *tx, &w_id, service_name, None, None, None, None)
|
|
.await?;
|
|
|
|
let mut native_ignore_keys = vec!["webhook_token_hash"];
|
|
if let Some(ref extra) = trigger_ignore_keys {
|
|
native_ignore_keys.extend_from_slice(extra);
|
|
}
|
|
|
|
for trigger in native_triggers {
|
|
let trigger_str = &to_string_without_metadata(
|
|
&trigger,
|
|
false,
|
|
Some(native_ignore_keys.clone()),
|
|
)
|
|
.unwrap();
|
|
archive
|
|
.write_to_archive(
|
|
&trigger_str,
|
|
&format!(
|
|
"{}.{}.{}.{}_native_trigger.json",
|
|
trigger.script_path,
|
|
if trigger.is_flow { "flow" } else { "script" },
|
|
trigger.external_id,
|
|
service_name.as_str()
|
|
),
|
|
)
|
|
.await?;
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
if include_users.unwrap_or(false) {
|
|
let users = sqlx::query!(
|
|
"SELECT * FROM usr
|
|
WHERE workspace_id = $1",
|
|
&w_id
|
|
)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
for user in users {
|
|
let user = SimplifiedUser {
|
|
username: user.username,
|
|
role: if user.is_admin {
|
|
"admin".to_string()
|
|
} else if user.operator {
|
|
"operator".to_string()
|
|
} else {
|
|
"developer".to_string()
|
|
},
|
|
disabled: user.disabled,
|
|
email: user.email,
|
|
};
|
|
let user_str = &to_string_without_metadata(&user, false, Some(vec!["email"])).unwrap();
|
|
archive
|
|
.write_to_archive(&user_str, &format!("users/{}.user.json", user.email))
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
if include_groups.unwrap_or(false) {
|
|
let groups = sqlx::query!(
|
|
r#"SELECT g_.workspace_id, name, summary, extra_perms, array_agg(u2g.usr) filter (where u2g.usr is not null) as members
|
|
FROM usr u
|
|
JOIN usr_to_group u2g ON u2g.usr = u.username AND u2g.workspace_id = u.workspace_id
|
|
RIGHT JOIN group_ g_ ON g_.workspace_id = u.workspace_id AND g_.name = u2g.group_
|
|
WHERE g_.workspace_id = $1 AND g_.name != 'all'
|
|
GROUP BY g_.workspace_id, name, summary, extra_perms"#,
|
|
&w_id
|
|
)
|
|
.fetch_all(&mut *tx)
|
|
.await?;
|
|
|
|
for group in groups {
|
|
let extra_perms: HashMap<String, bool> = serde_json::from_value(group.extra_perms)
|
|
.map_err(|e| {
|
|
Error::internal_err(format!(
|
|
"Error parsing extra_perms for group {}: {}",
|
|
group.name, e
|
|
))
|
|
})?;
|
|
tracing::info!("{:?}", extra_perms);
|
|
let members = group.members.unwrap_or(vec![]);
|
|
let admins: Vec<String> = extra_perms
|
|
.iter()
|
|
.filter_map(|(k, v)| {
|
|
// only consider extra_perms that concern actual members of the group
|
|
if members.contains(&k[2..].to_string()) && *v {
|
|
Some(k.clone())
|
|
} else {
|
|
None
|
|
}
|
|
})
|
|
.sorted()
|
|
.collect();
|
|
let group = SimplifiedGroup {
|
|
name: group.name,
|
|
summary: group.summary,
|
|
members: members
|
|
.iter()
|
|
.filter_map(|x| {
|
|
// remove members that are also admins as they are already in the admins list
|
|
let full_name = format!("u/{}", x);
|
|
if !admins.contains(&full_name) {
|
|
Some(full_name)
|
|
} else {
|
|
None
|
|
}
|
|
})
|
|
.collect(),
|
|
admins,
|
|
};
|
|
|
|
let group_str = &to_string_without_metadata(&group, true, None).unwrap();
|
|
archive
|
|
.write_to_archive(&group_str, &format!("groups/{}.group.json", group.name))
|
|
.await?;
|
|
}
|
|
}
|
|
|
|
if include_settings.unwrap_or(false) {
|
|
let row = sqlx::query_as::<_, SettingsRow>(
|
|
r#"SELECT
|
|
auto_invite,
|
|
webhook,
|
|
deploy_to,
|
|
error_handler,
|
|
success_handler,
|
|
ai_config,
|
|
large_file_storage,
|
|
git_sync,
|
|
default_app,
|
|
default_scripts,
|
|
workspace.name as name,
|
|
mute_critical_alerts,
|
|
color,
|
|
operator_settings,
|
|
datatable,
|
|
slack_team_id,
|
|
slack_name,
|
|
slack_command_script,
|
|
slack_oauth_client_id,
|
|
slack_oauth_client_secret
|
|
FROM workspace_settings
|
|
LEFT JOIN workspace ON workspace.id = workspace_settings.workspace_id
|
|
WHERE workspace_id = $1"#,
|
|
)
|
|
.bind(&w_id)
|
|
.fetch_one(&mut *tx)
|
|
.await?;
|
|
|
|
// Use v2 format only if explicitly requested, otherwise use v1 (legacy) for backward compatibility
|
|
let settings_str = if settings_version.as_deref() == Some("v2") {
|
|
let settings = SimplifiedSettings {
|
|
auto_invite: row.auto_invite,
|
|
webhook: row.webhook,
|
|
deploy_to: row.deploy_to,
|
|
error_handler: row.error_handler,
|
|
success_handler: row.success_handler,
|
|
ai_config: row.ai_config,
|
|
large_file_storage: row.large_file_storage,
|
|
git_sync: row.git_sync,
|
|
default_app: row.default_app,
|
|
default_scripts: row.default_scripts,
|
|
name: row.name.clone().unwrap_or_default(),
|
|
mute_critical_alerts: row.mute_critical_alerts,
|
|
color: row.color.clone(),
|
|
operator_settings: row.operator_settings.clone(),
|
|
datatable: row.datatable.clone(),
|
|
slack_team_id: row.slack_team_id.clone(),
|
|
slack_name: row.slack_name.clone(),
|
|
slack_command_script: row.slack_command_script.clone(),
|
|
slack_oauth_client_id: row.slack_oauth_client_id.clone(),
|
|
// Mirror the non-admin redaction in `get_settings`: the OAuth
|
|
// client secret is admin-only and must not leak via tarball.
|
|
slack_oauth_client_secret: if authed.is_admin {
|
|
row.slack_oauth_client_secret.clone()
|
|
} else {
|
|
None
|
|
},
|
|
};
|
|
serde_json::to_value(settings)
|
|
.map(|v| serde_json::to_string_pretty(&v).ok())
|
|
.ok()
|
|
.flatten()
|
|
} else {
|
|
// V1 (legacy) format: convert JSONB to flat fields (matches main branch exactly)
|
|
let (auto_invite_enabled, auto_invite_as, auto_invite_mode) =
|
|
if let Some(ref ai) = row.auto_invite {
|
|
let enabled = ai.get("enabled").and_then(|v| v.as_bool()).unwrap_or(false);
|
|
let operator = ai
|
|
.get("operator")
|
|
.and_then(|v| v.as_bool())
|
|
.unwrap_or(false);
|
|
let mode = ai.get("mode").and_then(|v| v.as_str()).unwrap_or("invite");
|
|
(
|
|
enabled,
|
|
if operator {
|
|
"operator".to_string()
|
|
} else {
|
|
"developer".to_string()
|
|
},
|
|
mode.to_string(),
|
|
)
|
|
} else {
|
|
(false, "developer".to_string(), "invite".to_string())
|
|
};
|
|
|
|
let (error_handler, error_handler_extra_args, error_handler_muted_on_cancel) =
|
|
if let Some(ref eh) = row.error_handler {
|
|
let path = eh.get("path").and_then(|v| v.as_str()).map(String::from);
|
|
let extra_args = eh.get("extra_args").cloned();
|
|
let muted_on_cancel = eh
|
|
.get("muted_on_cancel")
|
|
.and_then(|v| v.as_bool())
|
|
.unwrap_or(false);
|
|
(path, extra_args, muted_on_cancel)
|
|
} else {
|
|
(None, None, false)
|
|
};
|
|
|
|
let settings = SimplifiedSettingsLegacy {
|
|
auto_invite_enabled,
|
|
auto_invite_as,
|
|
auto_invite_mode,
|
|
webhook: row.webhook,
|
|
deploy_to: row.deploy_to,
|
|
error_handler,
|
|
error_handler_extra_args,
|
|
error_handler_muted_on_cancel,
|
|
ai_config: row.ai_config,
|
|
large_file_storage: row.large_file_storage,
|
|
git_sync: row.git_sync,
|
|
default_app: row.default_app,
|
|
default_scripts: row.default_scripts,
|
|
name: row.name.unwrap_or_default(),
|
|
mute_critical_alerts: row.mute_critical_alerts,
|
|
color: row.color,
|
|
operator_settings: row.operator_settings,
|
|
datatable: row.datatable,
|
|
slack_team_id: row.slack_team_id,
|
|
slack_name: row.slack_name,
|
|
slack_command_script: row.slack_command_script,
|
|
};
|
|
serde_json::to_value(settings)
|
|
.map(|v| serde_json::to_string_pretty(&v).ok())
|
|
.ok()
|
|
.flatten()
|
|
}
|
|
.ok_or_else(|| Error::internal_err("Error serializing settings".to_string()))?;
|
|
|
|
archive
|
|
.write_to_archive(&settings_str, "settings.json")
|
|
.await?;
|
|
}
|
|
|
|
if include_key.unwrap_or(false) {
|
|
require_admin(authed.is_admin, &authed.username)?;
|
|
|
|
let key = sqlx::query_scalar!(
|
|
"SELECT key FROM workspace_key WHERE workspace_id = $1",
|
|
&w_id
|
|
)
|
|
.fetch_one(&mut *tx)
|
|
.await?;
|
|
|
|
let key_json = serde_json::to_value(key)
|
|
.map(|v| serde_json::to_string_pretty(&v).ok())
|
|
.ok()
|
|
.flatten()
|
|
.ok_or_else(|| Error::internal_err("Error serializing enryption key".to_string()))?;
|
|
archive
|
|
.write_to_archive(&key_json, "encryption_key.json")
|
|
.await?;
|
|
}
|
|
|
|
archive.finish().await?;
|
|
|
|
let file = tokio::fs::File::open(&file_path).await?;
|
|
|
|
let stream = ReaderStream::new(file);
|
|
let body = axum::body::Body::from_stream(stream);
|
|
|
|
let headers = [
|
|
(header::CONTENT_TYPE, "application/x-tar".to_string()),
|
|
(
|
|
header::CONTENT_DISPOSITION,
|
|
format!("attachment; filename=\"{name}\""),
|
|
),
|
|
];
|
|
Ok((headers, body))
|
|
}
|