From 53caecf1da8d76e246178dfb9b86d330f0ec52fd Mon Sep 17 00:00:00 2001 From: Diego Imbert <70353967+diegoimbert@users.noreply.github.com> Date: Wed, 4 Mar 2026 11:46:08 +0100 Subject: [PATCH 01/12] feat: Ducklake typechecker (#8118) * Typedchecked ducklake queries * Display script preview error as SQL error * Fix duplication * fix replacer * Revert "fix replacer" This reverts commit c5492033c850cabc8bf18a50151c089b83cd6826. * Don't recompile regex every call * nit OOB * avoid potential panic * Apply suggestions from code review Co-authored-by: claude[bot] <209825114+claude[bot]@users.noreply.github.com> * safety throw * Update backend/windmill-worker/src/duckdb_executor.rs Co-authored-by: claude[bot] <209825114+claude[bot]@users.noreply.github.com> * Try catch individual chunks in prepareDatatableQueries Co-authored-by: claude[bot] <209825114+claude[bot]@users.noreply.github.com> * format * nit comment * Revert "Try catch individual chunks in prepareDatatableQueries" This reverts commit ae64a8ad27deb7e5ddda10163c6db04a428827f2. * Correct try catch * better error messages * nit unused variable * comment * handle non describable queries * npm i --------- Co-authored-by: claude[bot] <209825114+claude[bot]@users.noreply.github.com> --- .../windmill-duckdb-ffi-internal/Cargo.lock | 1 + .../windmill-duckdb-ffi-internal/Cargo.toml | 1 + .../windmill-duckdb-ffi-internal/src/lib.rs | 255 +++++++++++++++--- .../windmill-worker/src/duckdb_executor.rs | 81 ++++++ frontend/src/lib/infer.svelte.ts | 174 ++++++++---- 5 files changed, 423 insertions(+), 89 deletions(-) diff --git a/backend/windmill-duckdb-ffi-internal/Cargo.lock b/backend/windmill-duckdb-ffi-internal/Cargo.lock index 07e428c633..559196a3c2 100644 --- a/backend/windmill-duckdb-ffi-internal/Cargo.lock +++ b/backend/windmill-duckdb-ffi-internal/Cargo.lock @@ -2164,6 +2164,7 @@ version = "0.1.0" dependencies = [ "chrono", "duckdb", + "regex", "rust_decimal", "serde", "serde_json", diff --git a/backend/windmill-duckdb-ffi-internal/Cargo.toml b/backend/windmill-duckdb-ffi-internal/Cargo.toml index 7043b33ee5..7eb6869ab9 100644 --- a/backend/windmill-duckdb-ffi-internal/Cargo.toml +++ b/backend/windmill-duckdb-ffi-internal/Cargo.toml @@ -6,6 +6,7 @@ edition = "2024" [dependencies] chrono = "0.4.41" duckdb = { version = "1.4.4", features = ["bundled"] } +regex = "1" rust_decimal = "1.37.2" serde = { version = "1.0", features = ["derive"] } serde_json = { version = "^1", features = ["preserve_order", "raw_value"] } diff --git a/backend/windmill-duckdb-ffi-internal/src/lib.rs b/backend/windmill-duckdb-ffi-internal/src/lib.rs index c5c819e60b..2701319d6e 100644 --- a/backend/windmill-duckdb-ffi-internal/src/lib.rs +++ b/backend/windmill-duckdb-ffi-internal/src/lib.rs @@ -1,12 +1,14 @@ use std::{ collections::HashMap, - ffi::{CStr, CString, c_char, c_uint}, + ffi::{c_char, c_uint, CStr, CString}, ptr::null_mut, + sync::LazyLock, }; -use duckdb::{Row, core::LogicalTypeId, params_from_iter, types::TimeUnit}; -use rust_decimal::{Decimal, prelude::FromPrimitive}; -use serde::Deserialize; +use duckdb::{core::LogicalTypeId, params_from_iter, types::TimeUnit, Row}; +use regex::Regex; +use rust_decimal::{prelude::FromPrimitive, Decimal}; +use serde::{Deserialize, Serialize}; use serde_json::value::RawValue; #[derive(Deserialize, Clone, Debug, PartialEq, Default)] @@ -96,6 +98,218 @@ pub extern "C" fn run_duckdb_ffi( }) } +#[derive(Serialize, Debug)] +struct PrepareQueryColumnInfo { + name: String, + #[serde(rename = "type")] + type_name: String, +} + +#[derive(Serialize, Debug)] +struct PrepareQueryResult { + #[serde(skip_serializing_if = "Option::is_none")] + columns: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + error: Option, +} + +fn is_setup_statement(query: &str) -> bool { + let trimmed = query.trim_start(); + let upper = trimmed.to_uppercase(); + upper.starts_with("ATTACH") + || upper.starts_with("USE") + || upper.starts_with("INSTALL") + || upper.starts_with("LOAD") + || upper.starts_with("SET") + || upper.starts_with("RESET") + || upper.starts_with("CREATE OR REPLACE SECRET") + || upper.starts_with("CREATE SECRET") +} + +/// Returns true if the query is expected to return a result set and can be wrapped with DESCRIBE. +fn is_describable_query(query: &str) -> bool { + let trimmed = query.trim_start(); + let upper = trimmed.to_uppercase(); + upper.starts_with("SELECT") + || upper.starts_with("WITH") + || upper.starts_with("VALUES") + || upper.starts_with("TABLE") + || upper.starts_with("FROM") +} + +static PARAM_RE: LazyLock = LazyLock::new(|| Regex::new(r"\$\d+").expect("invalid regex")); + +fn replace_params_with_null(query: &str) -> String { + PARAM_RE.replace_all(query, "NULL").to_string() +} + +#[unsafe(no_mangle)] +pub extern "C" fn prepare_duckdb_ffi( + query_block_list: *const *const c_char, + query_block_list_count: usize, + token: *const c_char, + base_internal_url: *const c_char, + w_id: *const c_char, +) -> *mut c_char { + let r = match convert_prepare_args( + query_block_list, + query_block_list_count, + token, + base_internal_url, + w_id, + ) + .and_then(|(query_block_list, token, base_internal_url, w_id)| { + prepare_duckdb_internal(query_block_list, token, base_internal_url, w_id) + }) { + Ok(result) => result, + Err(err) => { + let err = serde_json::to_string(&err) + .unwrap_or_else(|_| "Unknown error in duckdb ffi lib".to_string()); + format!("ERROR {}", err) + } + }; + + CString::new(r).map(|s| s.into_raw()).unwrap_or_else(|e| { + println!("Failed to allocate error string in duckdb ffi lib: {:?}", e); + null_mut() + }) +} + +fn setup_duckdb_connection( + conn: &duckdb::Connection, + token: &str, + base_internal_url: &str, + w_id: &str, +) -> Result<(), String> { + let (s3_access_key, s3_secret_key) = token.rsplit_once('.').unwrap_or(("", token)); + let (s3_endpoint_ssl, s3_endpoint) = base_internal_url + .split_once("://") + .unwrap_or(("http", &base_internal_url)); + let s3_endpoint_ssl = s3_endpoint_ssl == "https"; + + conn.execute_batch(&format!( + "INSTALL httpfs; LOAD httpfs; + INSTALL azure; LOAD azure; + CREATE OR REPLACE SECRET s3_secret ( + TYPE s3, + PROVIDER config, + KEY_ID '{s3_access_key}', + SECRET '{s3_secret_key}', + ENDPOINT '{s3_endpoint}/api/w/{w_id}/s3_proxy', + URL_STYLE path, + USE_SSL {s3_endpoint_ssl} + ); + CREATE OR REPLACE SECRET gcs_secret ( + TYPE gcs, + KEY_ID '{s3_access_key}', + SECRET '{s3_secret_key}', + ENDPOINT '{s3_endpoint}/api/w/{w_id}/s3_proxy', + USE_SSL {s3_endpoint_ssl} + ); + ", + )) + .map_err(|e| format!("Error setting up S3 secret: {}", e.to_string())) +} + +fn convert_prepare_args<'a>( + query_block_list: *const *const c_char, + query_block_list_count: usize, + token: *const c_char, + base_internal_url: *const c_char, + w_id: *const c_char, +) -> Result<(Vec<&'a str>, &'a str, &'a str, &'a str), String> { + let query_block_list = unsafe { + std::slice::from_raw_parts(query_block_list, query_block_list_count) + .iter() + .map(|q| { + CStr::from_ptr(*q).to_str().unwrap_or_else(|e| { + println!( + "Invalid query_block string pointer in duckdb ffi: {}", + e.to_string() + ); + "Invalid query_block string pointer in duckdb ffi" + }) + }) + .collect::>() + }; + let token = unsafe { CStr::from_ptr(token) } + .to_str() + .map_err(|e| format!("Invalid token string: {}", e.to_string()))?; + let base_internal_url = unsafe { CStr::from_ptr(base_internal_url) } + .to_str() + .map_err(|e| format!("Invalid base_internal_url string: {}", e.to_string()))?; + let w_id = unsafe { CStr::from_ptr(w_id) } + .to_str() + .map_err(|e| format!("Invalid w_id string: {}", e.to_string()))?; + Ok((query_block_list, token, base_internal_url, w_id)) +} + +fn prepare_duckdb_internal( + query_block_list: Vec<&str>, + token: &str, + base_internal_url: &str, + w_id: &str, +) -> Result { + let conn = duckdb::Connection::open_in_memory().map_err(|e| e.to_string())?; + + setup_duckdb_connection(&conn, token, base_internal_url, w_id)?; + + let mut results: Vec = vec![]; + + // IMPORTANT: Setup statements (ATTACH, USE, INSTALL, etc.) are executed but intentionally + // do not produce a PrepareQueryResult entry. The frontend prepends these as connection setup + // before the actual user queries, and mapPrepareResults expects results.length to equal the + // number of user queries (not setup statements). If a new setup-like statement is added to + // the connection flow (e.g. in setup_duckdb_connection or transform_attach_ducklake) without + // also being caught by is_setup_statement, the result count will mismatch and the frontend + // will throw. + for query_block in &query_block_list { + if is_setup_statement(query_block) { + conn.execute_batch(query_block) + .map_err(|e| format!("Error executing setup statement: {}", e.to_string()))?; + continue; + } + + let modified_query = replace_params_with_null(query_block); + // Validate the query parses correctly by preparing it + if let Err(e) = conn.prepare(&modified_query) { + results.push(PrepareQueryResult { columns: None, error: Some(e.to_string()) }); + continue; + } + + // DESCRIBE only works on queries that return result sets (SELECT, WITH, VALUES, TABLE, + // FROM). For non-returning statements (INSERT, UPDATE, DELETE, CREATE, DROP, ALTER, etc.) + // we skip DESCRIBE and assume no columns. + if !is_describable_query(&modified_query) { + results.push(PrepareQueryResult { columns: Some(vec![]), error: None }); + continue; + } + + // Note: We have to use a DESCRIBE statement and cannot simply use the + // methods returned by .prepare() because they panic if the statement was + // not executed at least once (which we specifically do not want to do). + let describe_query = format!("DESCRIBE {}", modified_query); + match conn.prepare(&describe_query).and_then(|mut stmt| { + let rows = stmt.query_map([], |row| { + Ok(PrepareQueryColumnInfo { + name: row.get::<_, String>(0)?, + type_name: row.get::<_, String>(1)?, + }) + })?; + rows.collect::, _>>() + }) { + Ok(columns) => { + results.push(PrepareQueryResult { columns: Some(columns), error: None }); + } + Err(e) => { + results.push(PrepareQueryResult { columns: None, error: Some(e.to_string()) }); + } + } + } + + serde_json::to_string(&results).map_err(|e| e.to_string()) +} + fn convert_args<'a>( query_block_list: *const *const c_char, query_block_list_count: usize, @@ -170,38 +384,7 @@ fn run_duckdb_internal<'a>( ) -> Result<(String, Option>), String> { let conn = duckdb::Connection::open_in_memory().map_err(|e| e.to_string())?; - let (s3_access_key, s3_secret_key) = token.split_at(token.rfind('.').unwrap_or(0)); - let s3_secret_key = &s3_secret_key[1..]; - let (s3_endpoint_ssl, s3_endpoint) = base_internal_url - .split_once("://") - .unwrap_or(("http", &base_internal_url)); - let s3_endpoint_ssl = match s3_endpoint_ssl { - "https" => true, - _ => false, - }; - - conn.execute_batch(&format!( - "INSTALL httpfs; LOAD httpfs; - INSTALL azure; LOAD azure; - CREATE OR REPLACE SECRET s3_secret ( - TYPE s3, - PROVIDER config, - KEY_ID '{s3_access_key}', - SECRET '{s3_secret_key}', - ENDPOINT '{s3_endpoint}/api/w/{w_id}/s3_proxy', - URL_STYLE path, - USE_SSL {s3_endpoint_ssl} - ); - CREATE OR REPLACE SECRET gcs_secret ( - TYPE gcs, - KEY_ID '{s3_access_key}', - SECRET '{s3_secret_key}', - ENDPOINT '{s3_endpoint}/api/w/{w_id}/s3_proxy', - USE_SSL {s3_endpoint_ssl} - ); - ", - )) - .map_err(|e| format!("Error setting up S3 secret: {}", e.to_string()))?; + setup_duckdb_connection(&conn, token, base_internal_url, w_id)?; let mut results: Vec>> = vec![]; let mut column_order = None; diff --git a/backend/windmill-worker/src/duckdb_executor.rs b/backend/windmill-worker/src/duckdb_executor.rs index 73136e1cc0..45e2f647a5 100644 --- a/backend/windmill-worker/src/duckdb_executor.rs +++ b/backend/windmill-worker/src/duckdb_executor.rs @@ -161,6 +161,22 @@ pub async fn do_duckdb( let base_internal_url = client.base_internal_url.clone(); let w_id = job.workspace_id.clone(); + if annotations.prepare { + let result = tokio::task::spawn_blocking(move || { + prepare_duckdb_ffi_safe( + query_block_list.iter().map(String::as_str), + &token, + &base_internal_url, + &w_id, + ) + }) + .await + .map_err(|e| Error::from(to_anyhow(e))) + .and_then(|r| r)?; + + return Ok(result); + } + let result = tokio::task::spawn_blocking(move || { run_duckdb_ffi_safe( query_block_list.iter().map(String::as_str), @@ -248,6 +264,18 @@ struct DuckDbFfiLib { collect_first_row_only: bool, ) -> *mut c_char, >, + prepare_duckdb_ffi: Option< + Symbol< + 'static, + unsafe extern "C" fn( + query_block_list: *const *const c_char, + query_block_list_count: usize, + token: *const c_char, + base_internal_url: *const c_char, + w_id: *const c_char, + ) -> *mut c_char, + >, + >, free_cstr: Symbol<'static, unsafe extern "C" fn(string: *mut c_char) -> ()>, } @@ -307,8 +335,11 @@ impl DuckDbFfiLib { } } + let prepare_duckdb_ffi = unsafe { lib.get(b"prepare_duckdb_ffi").ok() }; + Ok(DuckDbFfiLib { run_duckdb_ffi: unsafe { lib.get(b"run_duckdb_ffi").map_err(to_anyhow)? }, + prepare_duckdb_ffi, free_cstr: unsafe { lib.get(b"free_cstr").map_err(to_anyhow)? }, }) } @@ -388,6 +419,56 @@ fn run_duckdb_ffi_safe<'a>( } } +fn prepare_duckdb_ffi_safe<'a>( + query_block_list: impl Iterator, + token: &str, + base_internal_url: &str, + w_id: &str, +) -> Result> { + let query_block_list = query_block_list + .map(|s| { + CString::new(s).map_err(|e| { + Error::ExecutionErr(format!("Failed CString conversion: {}", e.to_string())) + }) + }) + .collect::>>()?; + let query_block_list = query_block_list + .iter() + .map(|s| s.as_ptr()) + .collect::>(); + + let token = CString::new(token).map_err(to_anyhow)?; + let base_internal_url = CString::new(base_internal_url).map_err(to_anyhow)?; + let w_id = CString::new(w_id).map_err(to_anyhow)?; + + let lib = DuckDbFfiLib::get_singleton()?; + let prepare_fn = lib.prepare_duckdb_ffi.as_ref().ok_or_else(|| { + Error::InternalErr( + "prepare_duckdb_ffi not available in duckdb ffi library. Please update to the latest windmill_duckdb_ffi_lib.".to_string(), + ) + })?; + let free_cstr = &lib.free_cstr; + + let result_str = unsafe { + let ptr = prepare_fn( + query_block_list.as_ptr(), + query_block_list.len(), + token.as_ptr(), + base_internal_url.as_ptr(), + w_id.as_ptr(), + ); + let str = CStr::from_ptr(ptr).to_string_lossy().to_string(); + free_cstr(ptr); + str + }; + + if result_str.starts_with("ERROR") { + Err(Error::ExecutionErr(result_str[6..].to_string())) + } else { + Ok(serde_json::value::RawValue::from_string(result_str).map_err(to_anyhow)?) + } +} + struct ParsedAttachDbResource<'a> { resource_path: &'a str, name: &'a str, diff --git a/frontend/src/lib/infer.svelte.ts b/frontend/src/lib/infer.svelte.ts index 419020631c..cbdfa95505 100644 --- a/frontend/src/lib/infer.svelte.ts +++ b/frontend/src/lib/infer.svelte.ts @@ -4,6 +4,13 @@ import { ChangeOnDeepInequality, MapResource } from './svelte5Utils.svelte' import { sqlDataTypeToJsTypeHeuristic } from './components/apps/components/display/dbtable/utils' import { chunkBy, clone, getQueryStmtCountHeuristic } from './utils' +function extractErrorMessage(e: unknown): string { + if (e != null && typeof e === 'object' && 'body' in e) { + return (e as any).body?.error?.message ?? JSON.stringify(e) + } + return e instanceof Error ? e.message : JSON.stringify(e) +} + function computeQueryKey(query: InferAssetsSqlQueryDetails, workspace?: string) { return `${query.source_kind}::${query.source_name}::${query.source_schema}::${workspace}::${query.query_string}` } @@ -21,66 +28,26 @@ export function usePreparedAssetSqlQueries( ), async (toFetch) => { let queries = Object.entries(clone(toFetch)) - // We only support preparing datatable source kinds for now. - queries = queries.filter(([_, q]) => q.source_kind === 'datatable') + queries = queries.filter( + ([_, q]) => q.source_kind === 'datatable' || q.source_kind === 'ducklake' + ) // We only support preparing single-statement queries for now. queries = queries.filter(([_, q]) => getQueryStmtCountHeuristic(q.query_string) === 1) if (!queries?.length) return {} - try { - // We chunk by source_name to minimize the number of requests. - // For example if we have 10 queries on the same data table, - // we can prepare them all with a single script. - queries.sort((a, b) => a[1].source_name.localeCompare(b[1].source_name)) - let results = ( - await Promise.all( - chunkBy(queries, ([key, q]) => q.source_name).map(async (chunk) => { - console.log( - 'Preparing chunk of queries:', - chunk.map(([_, q]) => q) - ) - let queryContent = chunk - .flatMap(([key, q]) => [ - q.source_schema ? `SET search_path TO ${q.source_schema};` : 'RESET search_path;', - q.query_string + (q.query_string.trim().endsWith(';') ? '' : ';') - ]) - .join('\n') - queryContent = - '-- prepare\n--result_collection=all_statements_first_row\n' + queryContent + let datatableQueries = queries.filter(([_, q]) => q.source_kind === 'datatable') + let ducklakeQueries = queries.filter(([_, q]) => q.source_kind === 'ducklake') - let res = (await JobService.runScriptPreviewAndWaitResult({ - workspace: getWorkspace()!, - requestBody: { - language: 'postgresql', - content: queryContent, - args: { database: `datatable://${chunk[0][1]?.source_name}` } - } - })) as { error?: string; columns?: { name: string; type: string }[] }[] + let allResults: [string, PreparedAssetsSqlQuery][] = [] - console.log('Prepared query content:', res) - - let res2: [string, PreparedAssetsSqlQuery][] = res.map((r, i) => [ - chunk[i][0], - r.columns - ? { - columns: Object.fromEntries( - r.columns.map(({ name, type }) => [ - name, - sqlDataTypeToJsTypeHeuristic(type) - ]) - ) - } - : { error: r.error ?? "Couldn't prepare query " } - ]) - return res2 - }) - ) - ).flat() - - return Object.fromEntries(results) - } catch (e) { - throw e + if (datatableQueries.length) { + allResults.push(...(await prepareDatatableQueries(datatableQueries, getWorkspace))) } + if (ducklakeQueries.length) { + allResults.push(...(await prepareDucklakeQueries(ducklakeQueries, getWorkspace))) + } + + return Object.fromEntries(allResults) } ) @@ -96,3 +63,104 @@ export function usePreparedAssetSqlQueries( } } } + +type QueryEntry = [string, InferAssetsSqlQueryDetails] + +function mapPrepareResults( + res: { error?: string; columns?: { name: string; type: string }[] }[], + chunk: QueryEntry[] +): [string, PreparedAssetsSqlQuery][] { + if (res.length !== chunk.length) { + throw new Error(`Prepare results count mismatch: got ${res.length}, expected ${chunk.length}`) + } + return res.map((r, i) => [ + chunk[i]?.[0], + r.columns + ? { + columns: Object.fromEntries( + r.columns.map(({ name, type: t }) => [name, sqlDataTypeToJsTypeHeuristic(t)]) + ) + } + : { error: r.error ?? "Couldn't prepare query " } + ]) +} + +async function prepareDatatableQueries( + queries: QueryEntry[], + getWorkspace: () => string | undefined +): Promise<[string, PreparedAssetsSqlQuery][]> { + queries.sort((a, b) => a[1].source_name.localeCompare(b[1].source_name)) + let results = ( + await Promise.all( + chunkBy(queries, ([_, q]) => q.source_name).map(async (chunk) => { + let queryContent = chunk + .flatMap(([_, q]) => [ + q.source_schema ? `SET search_path TO ${q.source_schema};` : 'RESET search_path;', + q.query_string + (q.query_string.trim().endsWith(';') ? '' : ';') + ]) + .join('\n') + queryContent = '-- prepare\n--result_collection=all_statements_first_row\n' + queryContent + + try { + let res = (await JobService.runScriptPreviewAndWaitResult({ + workspace: getWorkspace()!, + requestBody: { + language: 'postgresql', + content: queryContent, + args: { database: `datatable://${chunk[0][1]?.source_name}` } + } + })) as { error?: string; columns?: { name: string; type: string }[] }[] + + return mapPrepareResults(res, chunk) + } catch (e) { + const error = extractErrorMessage(e) + return chunk.map(([key]) => [key, { error }] as [string, PreparedAssetsSqlQuery]) + } + }) + ) + ).flat() + return results +} + +async function prepareDucklakeQueries( + queries: QueryEntry[], + getWorkspace: () => string | undefined +): Promise<[string, PreparedAssetsSqlQuery][]> { + queries.sort((a, b) => a[1].source_name.localeCompare(b[1].source_name)) + let results = ( + await Promise.all( + chunkBy(queries, ([_, q]) => `${q.source_name}::${q.source_schema ?? ''}`).map( + async (chunk) => { + let sourceName = chunk[0][1].source_name + let sourceSchema = chunk[0][1].source_schema + let attachSetup = `ATTACH 'ducklake://${sourceName}' AS dl;\n` + attachSetup += sourceSchema ? `USE dl.${sourceSchema};\n` : `USE dl;\n` + + let queryContent = chunk + .map(([_, q]) => q.query_string + (q.query_string.trim().endsWith(';') ? '' : ';')) + .join('\n') + queryContent = + '-- prepare\n--result_collection=all_statements_first_row\n' + + attachSetup + + queryContent + + try { + let res = (await JobService.runScriptPreviewAndWaitResult({ + workspace: getWorkspace()!, + requestBody: { + language: 'duckdb', + content: queryContent, + args: {} + } + })) as { error?: string; columns?: { name: string; type: string }[] }[] + return mapPrepareResults(res, chunk) + } catch (e) { + const error = extractErrorMessage(e) + return chunk.map(([key]) => [key, { error }] as [string, PreparedAssetsSqlQuery]) + } + } + ) + ) + ).flat() + return results +} From 4bf827bea4d44aca8c5ff7aa67ad449dbcf00673 Mon Sep 17 00:00:00 2001 From: Diego Imbert <70353967+diegoimbert@users.noreply.github.com> Date: Wed, 4 Mar 2026 11:46:34 +0100 Subject: [PATCH 02/12] feat: persistent Db manager state in URI (#8134) * DB Manager state in URL * Fix state not saving * shorted uri params * infer db_type from prefix * Revert "infer db_type from prefix" This reverts commit 7415fbed3db0d570f321a0b86e0c5db6e876b430. * dbm syntax * infer database type * Omit main and public * remove legacy #dbmanager: * Preserve hash * nit * Fix remaining dbManagerDrawer objects --- .../src/lib/components/DBManagerDrawer.svelte | 77 +++---- .../src/lib/components/DatatablePicker.svelte | 5 +- .../src/lib/components/DucklakePicker.svelte | 5 +- .../lib/components/ExploreAssetButton.svelte | 7 +- .../src/lib/components/ResourcePicker.svelte | 4 +- frontend/src/lib/components/RunsPage.svelte | 4 +- .../lib/components/assets/AssetButtons.svelte | 3 - .../assets/AssetsDropdownButton.svelte | 4 +- .../components/assets/JobAssetsViewer.svelte | 4 +- .../components/dbManagerDrawerModel.svelte.ts | 207 ++++++++++++++++++ .../graph/renderers/nodes/AssetNode.svelte | 3 +- .../components/sidebar/FavoriteMenu.svelte | 2 +- .../CustomInstanceDbSelect.svelte | 4 - .../CustomInstanceDbWizardModal.svelte | 3 - .../DataTableSettings.svelte | 6 +- .../workspaceSettings/DucklakeSettings.svelte | 5 +- frontend/src/lib/stores.ts | 5 +- frontend/src/lib/svelte5UtilsKit.svelte.ts | 5 +- .../src/routes/(root)/(logged)/+layout.svelte | 28 +-- .../(root)/(logged)/assets/+page.svelte | 3 - .../(root)/(logged)/resources/+page.svelte | 6 +- 21 files changed, 266 insertions(+), 124 deletions(-) create mode 100644 frontend/src/lib/components/dbManagerDrawerModel.svelte.ts diff --git a/frontend/src/lib/components/DBManagerDrawer.svelte b/frontend/src/lib/components/DBManagerDrawer.svelte index c5986e18a4..5996d98075 100644 --- a/frontend/src/lib/components/DBManagerDrawer.svelte +++ b/frontend/src/lib/components/DBManagerDrawer.svelte @@ -6,27 +6,20 @@ import DrawerContent from './common/drawer/DrawerContent.svelte' import Select from './select/Select.svelte' import { ArrowLeft, Expand, LoaderCircle, Minimize, RefreshCcw } from 'lucide-svelte' - import type { DbInput } from './dbTypes' import DBManagerContent from './DBManagerContent.svelte' import { resource } from 'runed' + import { untrack } from 'svelte' + import type { DbManagerUriState } from './dbManagerDrawerModel.svelte' interface Props { + uriState: DbManagerUriState /** Z-index offset for the drawer, useful when opening from within modals */ offset?: number } - let { offset = 0 }: Props = $props() + let { uriState, offset = 0 }: Props = $props() - let input: DbInput | undefined = $state() - let open = $derived(!!input) - - // For datatable inputs, track the selected datatable separately - let selectedDatatable = $state(undefined) - - // Check if input is a datatable type - const isDatatableInput = $derived( - input?.type === 'database' && input.resourcePath.startsWith('datatable://') - ) + let open = $derived(uriState.open) // Load available datatables when drawer opens with datatable input const datatables = resource([], async () => { @@ -39,16 +32,6 @@ } }) - // Computed input that updates when selectedDatatable changes - const effectiveInput: DbInput | undefined = $derived.by(() => { - if (!input) return undefined - if (!isDatatableInput || !selectedDatatable) return input - return { - ...input, - resourcePath: `datatable://${selectedDatatable}` - } - }) - const datatableItems = $derived( datatables.current.map((dt) => ({ value: dt, @@ -56,32 +39,26 @@ })) ) - export function openDrawer(nInput: DbInput) { - input = nInput - if (isDatatableInput) { - datatables.refetch() + // Refetch datatables when switching to a datatable input + $effect(() => { + if (uriState.isDatatableInput) { + untrack(() => datatables.refetch()) } - // If it's a datatable input, extract the datatable name for the selector - if (nInput.type === 'database' && nInput.resourcePath.startsWith('datatable://')) { - selectedDatatable = nInput.resourcePath.replace('datatable://', '') - datatables.refetch() - } else { - selectedDatatable = undefined - } - } - export function closeDrawer() { - input = undefined - selectedDatatable = undefined + }) + + function handleClose() { + uriState.closeDrawer() dbManagerContent?.clearReplResult() - if (window.location.hash.startsWith('#dbmanager:')) - history.replaceState('', document.title, window.location.href.replace(/#dbmanager:.*$/, '')) } let windowWidth = $state(window.innerWidth) let expand = $state(false) $effect(() => { - if (!open) expand = false + if (!open) { + expand = false + uriState.closeDrawer() + } }) let dbManagerContent: DBManagerContent | undefined = $state() @@ -96,7 +73,7 @@ size={expand ? `${windowWidth}px` : '1200px'} preventEscape {offset} - on:close={closeDrawer} + on:close={handleClose} > - {#if effectiveInput && $workspaceStore} - {#key selectedDatatable} - + {#if uriState.effectiveInput && $workspaceStore} + {#key uriState.selectedDatatable} + {#snippet dbSelector()} - {#if isDatatableInput} + {#if uriState.isDatatableInput} {#if datatables.loading}
@@ -125,7 +108,7 @@ updateEnvValue(entry.key, e.currentTarget.value, 'string')} - disabled={noEditor} - class="input w-full" - placeholder="Variable value" - /> +
+ + updateEnvValue(entry.key, e.currentTarget.value, 'string')} + disabled={noEditor} + class="input w-full" + placeholder="Variable value" + /> + {#if !noEditor} +
+ {#if typeof entry.value === 'string' && entry.value.startsWith('$var:') && entry.value.length > 5} +
+ Linked to variable {entry.value.slice(5)} +
+ {/if} {/if} -
+ {/each} @@ -285,3 +367,20 @@ + + { + if (pickForKey) { + setVarPath(pickForKey, path) + pickForKey = undefined + } + }} + itemName="Variable" + extraField="path" + loadItems={async () => + (await VariableService.listVariable({ workspace: $workspaceStore ?? '' })).map((x) => ({ + name: x.path, + ...x + }))} +/> diff --git a/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte b/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte index 6910fabbb5..ead5ad9ad3 100644 --- a/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte +++ b/frontend/src/lib/components/flows/propPicker/PropPickerWrapper.svelte @@ -27,7 +27,7 @@ import type { PickableProperties } from '../previousResults' import AnimatedButton from '$lib/components/common/button/AnimatedButton.svelte' import type { PropPickerContext } from '$lib/components/prop_picker' - import type { FlowEditorContext } from '../types' + interface Props { pickableProperties: PickableProperties | undefined @@ -67,9 +67,8 @@ const { flowPropPickerConfig } = getContext('PropPickerContext') flowPropPickerConfig.set(undefined) - const { flowStore } = getContext('FlowEditorContext') - let flow_env = $derived(pickableProperties?.flow_env || flowStore.val.value.flow_env) + setContext('PropPickerWrapper', { propPickerConfig, inputMatches, @@ -156,7 +155,6 @@ {extraResults} {displayContext} {error} - {flow_env} previousId={pickableProperties?.previousId} {pickableProperties} allowCopy={!notSelectable && !$propPickerConfig} diff --git a/frontend/src/lib/components/propertyPicker/PropPicker.svelte b/frontend/src/lib/components/propertyPicker/PropPicker.svelte index d04683e226..6910750480 100644 --- a/frontend/src/lib/components/propertyPicker/PropPicker.svelte +++ b/frontend/src/lib/components/propertyPicker/PropPicker.svelte @@ -19,7 +19,6 @@ error?: boolean allowCopy?: boolean previousId?: string | undefined - flow_env?: Record | undefined result?: any | undefined extraResults?: any } @@ -30,7 +29,6 @@ error = false, allowCopy = false, previousId = undefined, - flow_env = undefined, result = undefined, extraResults = undefined }: Props = $props() @@ -39,7 +37,6 @@ let resources: Record = $state({}) let displayVariable = $state(false) let displayResources = $state(false) - let displayFlowEnv = $state(false) let allResultsCollapsed = $state(true) let collapsableInitialState: @@ -47,7 +44,6 @@ allResultsCollapsed: boolean displayVariable: boolean displayResources: boolean - displayFlowEnv: boolean } | undefined @@ -139,7 +135,9 @@ resultByIdFiltered = {} } if (!$inputMatches?.some((match) => match.word === 'flow_env')) { - flowEnvFiltered = {} + if (search === EMPTY_STRING) { + flowEnvFiltered = pickableProperties.flow_env + } } if ($inputMatches?.length == 1) { filteringFlowInputsOrResult = $inputMatches[0].value @@ -185,8 +183,7 @@ collapsableInitialState = { allResultsCollapsed, displayVariable, - displayResources, - displayFlowEnv + displayResources } } @@ -200,10 +197,6 @@ displayResources = true return } - if ($inputMatches[0].word === 'flow_env') { - displayFlowEnv = true - return - } if ($inputMatches[0].word === 'results') { allResultsCollapsed = false return @@ -214,8 +207,7 @@ if (!collapsableInitialState) { return } - ;({ allResultsCollapsed, displayVariable, displayResources, displayFlowEnv } = - collapsableInitialState) + ;({ allResultsCollapsed, displayVariable, displayResources } = collapsableInitialState) collapsableInitialState = undefined } @@ -279,6 +271,18 @@ /> {/if} + {#if flowEnvFiltered && Object.keys(flowEnvFiltered ?? {}).length > 0} + Flow Env Variables +
+ +
+ {/if} {#if error} Error
@@ -445,45 +449,6 @@ {/if}
{/if} - {#if flow_env && Object.keys(flow_env).length > 0 && $inputMatches?.some((match) => match.word === 'flow_env')} -
- Flow Env Variables: - - {#if displayFlowEnv} - - - {:else} - - {/if} -
- {/if} {/if} diff --git a/frontend/src/routes/(root)/(logged)/variables/+page.svelte b/frontend/src/routes/(root)/(logged)/variables/+page.svelte index 4082e5f8ff..870bf8e7c4 100644 --- a/frontend/src/routes/(root)/(logged)/variables/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/variables/+page.svelte @@ -40,7 +40,8 @@ EyeOff, Circle } from 'lucide-svelte' - import { untrack } from 'svelte' + import { onMount, untrack } from 'svelte' + import { page } from '$app/stores' type ListableVariableW = ListableVariable & { canWrite: boolean } @@ -202,6 +203,14 @@ loadContextualVariables() }, 5000) } + + onMount(() => { + let hash = $page.url.hash + if (hash.length > 1) { + let path = hash.slice(1) + variableEditor?.editVariable(path) + } + }) diff --git a/openflow.openapi.yaml b/openflow.openapi.yaml index 2b181a6c7b..5031c7e889 100644 --- a/openflow.openapi.yaml +++ b/openflow.openapi.yaml @@ -96,9 +96,8 @@ components: type: boolean flow_env: type: object - description: Environment variables available to all steps - additionalProperties: - type: string + description: "Environment variables available to all steps. Values can be strings, JSON values, or special references: '$var:path' (workspace variable) or '$res:path' (resource)." + additionalProperties: {} priority: type: number description: Execution priority (higher numbers run first) From 19c065bed5468c484c8e7a50a6b79ab90153cc0e Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Wed, 4 Mar 2026 14:44:33 +0000 Subject: [PATCH 09/12] fix: handle multipart stream errors gracefully instead of panicking (#8226) Co-authored-by: Claude Opus 4.6 --- backend/windmill-api/src/apps.rs | 18 ++++++++++++++---- 1 file changed, 14 insertions(+), 4 deletions(-) diff --git a/backend/windmill-api/src/apps.rs b/backend/windmill-api/src/apps.rs index 1cfbfa6e19..9846912192 100644 --- a/backend/windmill-api/src/apps.rs +++ b/backend/windmill-api/src/apps.rs @@ -993,9 +993,18 @@ macro_rules! process_app_multipart { let mut uploaded_js = false; let mut multipart = $multipart; - while let Some(field) = multipart.next_field().await.unwrap() { - let name = field.name().unwrap().to_string(); - let data = field.bytes().await.unwrap(); + while let Some(field) = multipart + .next_field() + .await + .map_err(|e| Error::BadRequest(format!("failed to read multipart field: {e}")))? + { + let name = field + .name() + .ok_or_else(|| Error::BadRequest("multipart field missing name".to_string()))? + .to_string(); + let data = field.bytes().await.map_err(|e| { + Error::BadRequest(format!("failed to read multipart stream: {e}")) + })?; if name == "app" { let app = serde_json::from_slice(&data).map_err(to_anyhow)?; let (ntx, npath, nid) = $internal_fn( @@ -2149,7 +2158,8 @@ async fn execute_component( (email.as_str(), permissioned_as) }; - let end_user_email = get_end_user_email(&db, opt_authed.as_ref(), tokened.token.as_deref()).await; + let end_user_email = + get_end_user_email(&db, opt_authed.as_ref(), tokened.token.as_deref()).await; let (uuid, mut tx) = push( &db, From 62382fd2869ea0190dd0c0b714f9cbd35ceddd7a Mon Sep 17 00:00:00 2001 From: hugocasa Date: Wed, 4 Mar 2026 15:53:56 +0100 Subject: [PATCH 10/12] fix: wrap set_encryption_key in a single database transaction (#8212) Prevent workspace corruption when re-encryption fails mid-loop by wrapping the key update and variable re-encryption in a single transaction. If any step fails, the entire operation rolls back. Co-authored-by: Claude Opus 4.6 --- backend/Cargo.lock | 1 + backend/windmill-api-workspaces/Cargo.toml | 1 + .../windmill-api-workspaces/src/workspaces.rs | 29 ++++++++++++++----- 3 files changed, 24 insertions(+), 7 deletions(-) diff --git a/backend/Cargo.lock b/backend/Cargo.lock index d26d5b76d7..e52bb12fa1 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -16383,6 +16383,7 @@ dependencies = [ "http 1.4.0", "hyper 1.8.1", "lazy_static", + "magic-crypt", "regex", "serde", "serde_json", diff --git a/backend/windmill-api-workspaces/Cargo.toml b/backend/windmill-api-workspaces/Cargo.toml index b698426d53..a03bb3a490 100644 --- a/backend/windmill-api-workspaces/Cargo.toml +++ b/backend/windmill-api-workspaces/Cargo.toml @@ -29,6 +29,7 @@ windmill-dep-map.workspace = true axum.workspace = true chrono.workspace = true hex.workspace = true +magic-crypt.workspace = true http.workspace = true hyper.workspace = true lazy_static.workspace = true diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index 8f3311978a..80e497ba97 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -31,7 +31,9 @@ use windmill_audit::audit_oss::{audit_log, AuditAuthorable}; use windmill_audit::ActionKind; use windmill_common::db::UserDB; use windmill_common::users::username_to_permissioned_as; -use windmill_common::variables::{build_crypt, decrypt, encrypt, WORKSPACE_CRYPT_CACHE}; +use windmill_common::variables::{ + build_crypt, decrypt, encrypt, SECRET_SALT, WORKSPACE_CRYPT_CACHE, +}; use windmill_common::worker::{to_raw_value, CLOUD_HOSTED}; #[cfg(feature = "enterprise")] use windmill_common::workspaces::GitRepositorySettings; @@ -2418,20 +2420,28 @@ async fn set_encryption_key( )); } + // Build the previous cipher before the transaction (reads from cache/pool) let previous_encryption_key = build_crypt(&db, w_id.as_str()).await?; + let mut tx = db.begin().await?; + sqlx::query!( "UPDATE workspace_key SET key = $1 WHERE workspace_id = $2", request.new_key.clone(), w_id ) - .execute(&db) + .execute(&mut *tx) .await?; - WORKSPACE_CRYPT_CACHE.remove(w_id.as_str()); - if !request.skip_reencrypt.unwrap_or(false) { - let new_encryption_key = build_crypt(&db, w_id.as_str()).await?; + // Build the new cipher directly from the key string, since the transaction + // hasn't committed yet and build_crypt() would read the old key from the pool. + let crypt_key = if let Some(ref salt) = SECRET_SALT.as_ref() { + format!("{}{}", request.new_key, salt) + } else { + request.new_key.clone() + }; + let new_encryption_key = magic_crypt::new_magic_crypt!(crypt_key, 256); let mut truncated_new_key = request.new_key.clone(); truncated_new_key.truncate(8); @@ -2445,7 +2455,7 @@ async fn set_encryption_key( "SELECT path, value, is_secret FROM variable WHERE workspace_id = $1", w_id ) - .fetch_all(&db) + .fetch_all(&mut *tx) .await?; for variable in all_variables { @@ -2466,11 +2476,16 @@ async fn set_encryption_key( w_id, variable.path ) - .execute(&db) + .execute(&mut *tx) .await?; } } + tx.commit().await?; + + // Invalidate the cache only after the transaction has committed + WORKSPACE_CRYPT_CACHE.remove(w_id.as_str()); + // Trigger git sync for encryption key changes handle_deployment_metadata( &authed.email, From 87ebeaa51d9ca22bca9a1deba2590f9c6e5d3f77 Mon Sep 17 00:00:00 2001 From: centdix <40307056+centdix@users.noreply.github.com> Date: Wed, 4 Mar 2026 16:09:42 +0100 Subject: [PATCH 11/12] chore: make rust-analyzer plugin opt-in via USE_RUST_PLUGIN env var (#8227) * feat: optionally enable rust-analyzer plugin in worktree settings When USE_RUST_PLUGIN env var is set, the worktree-env script now includes the rust-analyzer-lsp plugin in .claude/settings.local.json. Co-Authored-By: Claude Opus 4.6 * chore: remove rust-analyzer plugin from default settings The rust-analyzer plugin is now opt-in via USE_RUST_PLUGIN env var in worktree-env, so it no longer needs to be in the shared settings. Co-Authored-By: Claude Opus 4.6 * chore: add WM_CLONE_DB and USE_RUST_PLUGIN to wmdev startup envs Defaults both to false so they can be toggled per-worktree. Co-Authored-By: Claude Opus 4.6 * fix: use explicit truthy checks for WM_CLONE_DB and USE_RUST_PLUGIN Co-Authored-By: Claude Opus 4.6 --------- Co-authored-by: Claude Opus 4.6 --- .claude/settings.json | 1 - .wmdev.yaml | 2 ++ scripts/worktree-env | 11 +++++++++-- 3 files changed, 11 insertions(+), 3 deletions(-) diff --git a/.claude/settings.json b/.claude/settings.json index fcd49c3140..cf8bfdd284 100644 --- a/.claude/settings.json +++ b/.claude/settings.json @@ -110,7 +110,6 @@ ] }, "enabledPlugins": { - "rust-analyzer-lsp@claude-plugins-official": true, "typescript-lsp@claude-plugins-official": true, "code-review@claude-plugins-official": true } diff --git a/.wmdev.yaml b/.wmdev.yaml index c028c8f3bf..1a949c94e2 100644 --- a/.wmdev.yaml +++ b/.wmdev.yaml @@ -2,6 +2,8 @@ name: Windmill startupEnvs: CARGO_FEATURES: "quickjs" + WM_CLONE_DB: false + USE_RUST_PLUGIN: false services: - name: BE diff --git a/scripts/worktree-env b/scripts/worktree-env index ae83b09405..5fd8490bc2 100755 --- a/scripts/worktree-env +++ b/scripts/worktree-env @@ -61,7 +61,7 @@ if command -v psql &>/dev/null; then if psql "$db_conn/postgres" -tc "SELECT 1 FROM pg_database WHERE datname = '${db_name}'" 2>/dev/null | grep -q 1; then echo "Database $db_name already exists" else - if [[ -n "${WM_CLONE_DB:-}" ]]; then + if [[ "${WM_CLONE_DB:-}" == "1" || "${WM_CLONE_DB:-}" == "true" ]]; then # Terminate active connections so CREATE DATABASE ... TEMPLATE works psql "$db_conn/postgres" -c "SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = 'windmill' AND pid <> pg_backend_pid();" 2>/dev/null || true psql "$db_conn/postgres" -c "CREATE DATABASE ${db_name} TEMPLATE windmill" 2>/dev/null \ @@ -178,13 +178,20 @@ if [ -n "$ee_repo" ]; then if [ -d "$ee_worktree_dir" ]; then ee_rel=$(python3 -c "import os; print(os.path.relpath('$ee_worktree_dir', '$(pwd)'))" 2>/dev/null || echo "$ee_worktree_dir") mkdir -p .claude + rust_plugin="" + if [[ "${USE_RUST_PLUGIN:-}" == "1" || "${USE_RUST_PLUGIN:-}" == "true" ]]; then + rust_plugin=', + "enabledPlugins": { + "rust-analyzer-lsp@claude-plugins-official": true + }' + fi cat > .claude/settings.local.json < Date: Wed, 4 Mar 2026 16:12:00 +0100 Subject: [PATCH 12/12] feat: replace hub error toasts with warning alerts and add disable hub setting (#8225) * feat: replace hub error toasts with warning alerts and add disable hub setting Co-Authored-By: Claude Opus 4.6 * fix: guard hub script cache refresh when hub is disabled Co-Authored-By: Claude Opus 4.6 --------- Co-authored-by: Claude Opus 4.6 --- backend/windmill-api-settings/src/lib.rs | 5 +++-- .../windmill-common/src/global_settings.rs | 1 + .../windmill-common/src/instance_config.rs | 2 ++ .../flows/pickers/PickHubApp.svelte | 22 +++++++++++++++++-- .../flows/pickers/PickHubFlow.svelte | 22 +++++++++++++++++-- .../flows/pickers/PickHubScript.svelte | 11 +++++++++- .../flows/pickers/PickHubScriptQuick.svelte | 22 +++++++++++++------ .../src/lib/components/instanceSettings.ts | 10 +++++++++ frontend/src/lib/stores.ts | 1 + .../src/routes/(root)/(logged)/+layout.svelte | 7 ++++++ 10 files changed, 89 insertions(+), 14 deletions(-) diff --git a/backend/windmill-api-settings/src/lib.rs b/backend/windmill-api-settings/src/lib.rs index bf493f9420..6b408724fd 100644 --- a/backend/windmill-api-settings/src/lib.rs +++ b/backend/windmill-api-settings/src/lib.rs @@ -43,8 +43,8 @@ use windmill_common::{ get_database_url, global_settings::{ APP_WORKSPACED_ROUTE_SETTING, AUTOMATE_USERNAME_CREATION_SETTING, - CRITICAL_ALERT_MUTE_UI_SETTING, DEFAULT_TAGS_WORKSPACES_SETTING, EMAIL_DOMAIN_SETTING, - ENV_SETTINGS, HUB_ACCESSIBLE_URL_SETTING, HUB_BASE_URL_SETTING, + CRITICAL_ALERT_MUTE_UI_SETTING, DEFAULT_TAGS_WORKSPACES_SETTING, DISABLE_HUB_SETTING, + EMAIL_DOMAIN_SETTING, ENV_SETTINGS, HUB_ACCESSIBLE_URL_SETTING, HUB_BASE_URL_SETTING, }, instance_config::{self, ApplyMode, InstanceConfig}, server::Smtp, @@ -519,6 +519,7 @@ pub async fn get_global_setting( && key != DEFAULT_TAGS_WORKSPACES_SETTING && key != HUB_BASE_URL_SETTING && key != HUB_ACCESSIBLE_URL_SETTING + && key != DISABLE_HUB_SETTING && key != EMAIL_DOMAIN_SETTING && key != APP_WORKSPACED_ROUTE_SETTING { diff --git a/backend/windmill-common/src/global_settings.rs b/backend/windmill-common/src/global_settings.rs index 3347127303..d4d8163ff2 100644 --- a/backend/windmill-common/src/global_settings.rs +++ b/backend/windmill-common/src/global_settings.rs @@ -44,6 +44,7 @@ pub const HUB_API_SECRET_SETTING: &str = "hub_api_secret"; pub const AUTOMATE_USERNAME_CREATION_SETTING: &str = "automate_username_creation"; pub const HUB_BASE_URL_SETTING: &str = "hub_base_url"; pub const HUB_ACCESSIBLE_URL_SETTING: &str = "hub_accessible_url"; +pub const DISABLE_HUB_SETTING: &str = "disable_hub"; pub const CRITICAL_ERROR_CHANNELS_SETTING: &str = "critical_error_channels"; pub const CRITICAL_ALERT_MUTE_UI_SETTING: &str = "critical_alert_mute_ui"; pub const CRITICAL_ALERTS_ON_DB_OVERSIZE_SETTING: &str = "critical_alerts_on_db_oversize"; diff --git a/backend/windmill-common/src/instance_config.rs b/backend/windmill-common/src/instance_config.rs index 4de90be5c5..4ef88f5b78 100644 --- a/backend/windmill-common/src/instance_config.rs +++ b/backend/windmill-common/src/instance_config.rs @@ -230,6 +230,8 @@ pub struct GlobalSettings { pub no_default_maven: Option, #[serde(skip_serializing_if = "Option::is_none")] pub default_tags_per_workspace: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub disable_hub: Option, // String settings #[serde(skip_serializing_if = "Option::is_none")] diff --git a/frontend/src/lib/components/flows/pickers/PickHubApp.svelte b/frontend/src/lib/components/flows/pickers/PickHubApp.svelte index 7f534a924a..234f053f7b 100644 --- a/frontend/src/lib/components/flows/pickers/PickHubApp.svelte +++ b/frontend/src/lib/components/flows/pickers/PickHubApp.svelte @@ -7,6 +7,8 @@ import RowIcon from '$lib/components/common/table/RowIcon.svelte' import { loadHubApps } from '$lib/hub' import TextInput from '$lib/components/text_input/TextInput.svelte' + import { Alert } from '$lib/components/common' + import { disableHubStore } from '$lib/stores' interface Props { filter?: string @@ -30,11 +32,22 @@ const dispatch = createEventDispatcher() + let hubNotAvailable = $state(false) + onMount(async () => { - hubApps = await loadHubApps() + if ($disableHubStore) return + const result = await loadHubApps() + if (result === undefined) { + hubNotAvailable = true + } else { + hubApps = result + } }) +{#if $disableHubStore} + +{:else} -{#if hubApps} +{#if hubNotAvailable} + + Could not connect to the Windmill Hub. If you are in a closed environment, you can disable the Hub in the instance settings. + +{:else if hubApps} {#if filteredItems.length == 0} {:else} @@ -93,3 +110,4 @@ {/each} {/if} +{/if} diff --git a/frontend/src/lib/components/flows/pickers/PickHubFlow.svelte b/frontend/src/lib/components/flows/pickers/PickHubFlow.svelte index 2ab12aea36..1f3c659bf8 100644 --- a/frontend/src/lib/components/flows/pickers/PickHubFlow.svelte +++ b/frontend/src/lib/components/flows/pickers/PickHubFlow.svelte @@ -7,6 +7,8 @@ import RowIcon from '$lib/components/common/table/RowIcon.svelte' import { loadHubFlows } from '$lib/hub' import TextInput from '$lib/components/text_input/TextInput.svelte' + import { Alert } from '$lib/components/common' + import { disableHubStore } from '$lib/stores' interface Props { filter?: string @@ -30,11 +32,22 @@ const dispatch = createEventDispatcher() + let hubNotAvailable = $state(false) + onMount(async () => { - hubFlows = await loadHubFlows() + if ($disableHubStore) return + const result = await loadHubFlows() + if (result === undefined) { + hubNotAvailable = true + } else { + hubFlows = result + } }) +{#if $disableHubStore} + +{:else} -{#if hubFlows} +{#if hubNotAvailable} + + Could not connect to the Windmill Hub. If you are in a closed environment, you can disable the Hub in the instance settings. + +{:else if hubFlows} {#if filteredItems.length == 0} {:else} @@ -95,3 +112,4 @@ {/each} {/if} +{/if} diff --git a/frontend/src/lib/components/flows/pickers/PickHubScript.svelte b/frontend/src/lib/components/flows/pickers/PickHubScript.svelte index b5ba87327e..7ec0a8ccc3 100644 --- a/frontend/src/lib/components/flows/pickers/PickHubScript.svelte +++ b/frontend/src/lib/components/flows/pickers/PickHubScript.svelte @@ -8,6 +8,7 @@ import { IntegrationService, ScriptService, type HubScriptKind } from '$lib/gen' import { Loader2 } from 'lucide-svelte' import TextInput from '$lib/components/text_input/TextInput.svelte' + import { disableHubStore } from '$lib/stores' interface Props { kind?: HubScriptKind & string @@ -47,6 +48,7 @@ ) async function getAllApps(filterKind: typeof kind) { + if ($disableHubStore) return try { hubNotAvailable = false allApps = ( @@ -67,6 +69,7 @@ filterKind: typeof kind, appFilter: string | undefined ) { + if ($disableHubStore) return try { loading = true hubNotAvailable = false @@ -138,6 +141,9 @@ }) +{#if $disableHubStore} + +{:else}
{@render children?.()}
@@ -156,7 +162,9 @@
{#if hubNotAvailable} - + + Could not connect to the Windmill Hub. If you are in a closed environment, you can disable the Hub in the instance settings. + {:else if (items.length > 0 && apps.length > 0) || !loading} {#if items.length == 0} @@ -204,3 +212,4 @@ {/each} {/if} +{/if} diff --git a/frontend/src/lib/components/flows/pickers/PickHubScriptQuick.svelte b/frontend/src/lib/components/flows/pickers/PickHubScriptQuick.svelte index 2091f76340..3d929dead7 100644 --- a/frontend/src/lib/components/flows/pickers/PickHubScriptQuick.svelte +++ b/frontend/src/lib/components/flows/pickers/PickHubScriptQuick.svelte @@ -24,7 +24,7 @@ []) : undefined } catch (err) { - sendUserToast('Failed to fetch hub scripts: ' + err, 'error') + console.error('Failed to fetch hub scripts:', err) return undefined } }, @@ -44,9 +44,10 @@ import { Circle, ExternalLink } from 'lucide-svelte' import Popover from '$lib/components/Popover.svelte' import { usePromise } from '$lib/svelte5Utils.svelte' - import { hubBaseUrlStore, userStore } from '$lib/stores' + import { disableHubStore, hubBaseUrlStore, userStore } from '$lib/stores' import { get } from 'svelte/store' import Button from '$lib/components/common/button/Button.svelte' + import { Alert } from '$lib/components/common' let hubNotAvailable = $state(false) @@ -94,13 +95,14 @@ }) async function getAllApps(filterKind: typeof kind) { + if ($disableHubStore) return try { hubNotAvailable = false allApps = (await listHubIntegrationsCached({ kind: filterKind, refreshCount })).map( (x) => x.name ) } catch (err) { - sendUserToast('Failed to fetch hub integrations: ' + err, 'error') + console.error('Failed to fetch hub integrations:', err) allApps = [] hubNotAvailable = true } @@ -112,7 +114,9 @@ ) $effect(() => { ;[filter, kind, appFilter, refreshCount] - hubScriptsFilteredPromise.refresh() + if (!$disableHubStore) { + hubScriptsFilteredPromise.refresh() + } }) $effect(() => { loading = hubScriptsFilteredPromise.status === 'loading' @@ -175,9 +179,13 @@ -{#if hubNotAvailable} -
- Hub not available +{#if $disableHubStore} + +{:else if hubNotAvailable} +
+ + Could not connect to the Windmill Hub. If you are in a closed environment, you can disable the Hub in the instance settings. +
{:else if loading} {#each Array(15).fill(0) as _} diff --git a/frontend/src/lib/components/instanceSettings.ts b/frontend/src/lib/components/instanceSettings.ts index de4cd1528b..25955b74db 100644 --- a/frontend/src/lib/components/instanceSettings.ts +++ b/frontend/src/lib/components/instanceSettings.ts @@ -331,6 +331,16 @@ export const settings: Record = { storage: 'setting', ee_only: '', hiddenIfEmpty: true + }, + { + label: 'Disable Hub', + description: + 'Disable the Windmill Hub integration entirely. Enable this if your instance runs in a closed environment without internet access and you do not have a private hub setup.', + key: 'disable_hub', + fieldType: 'boolean', + storage: 'setting', + ee_only: '', + requiresReloadOnChange: true } ], SMTP: [ diff --git a/frontend/src/lib/stores.ts b/frontend/src/lib/stores.ts index 6a9a8e7805..fd0a66cfb3 100644 --- a/frontend/src/lib/stores.ts +++ b/frontend/src/lib/stores.ts @@ -83,6 +83,7 @@ export const superadmin = writable(undefined) export const devopsRole = writable(undefined) export const lspTokenStore = writable(undefined) export const hubBaseUrlStore = writable(DEFAULT_HUB_BASE_URL) +export const disableHubStore = writable(false) export const userWorkspaces: Readable> = derived( [usersWorkspaceStore, superadmin], ([store, superadmin]) => { diff --git a/frontend/src/routes/(root)/(logged)/+layout.svelte b/frontend/src/routes/(root)/(logged)/+layout.svelte index e9b1884df2..78de6c1535 100644 --- a/frontend/src/routes/(root)/(logged)/+layout.svelte +++ b/frontend/src/routes/(root)/(logged)/+layout.svelte @@ -26,6 +26,7 @@ type UserExt, defaultScripts, hubBaseUrlStore, + disableHubStore, usedTriggerKinds, devopsRole, whitelabelNameStore, @@ -157,6 +158,7 @@ loadUsage() syncTutorialsTodos() loadHubBaseUrl() + loadDisableHub() loadUsedTriggerKinds() } @@ -176,6 +178,11 @@ DEFAULT_HUB_BASE_URL } + async function loadDisableHub() { + $disableHubStore = + ((await SettingService.getGlobal({ key: 'disable_hub' })) as boolean) ?? false + } + async function loadFavorites() { const scripts = await ScriptService.listScripts({ workspace: $workspaceStore ?? '',