mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-09 08:02:25 +00:00
fix: workspace specfic tags compatibility with forked workspaces (#8850)
* fix: workspace specfic tags compatibility with forked workspaces * Rename _db to db and use saved WM_FORK_PREFIX * Add ttl cache for mapping fork id to parent workspace id * Change second option to just have a -fork suffix
This commit is contained in:
+21
-13
@@ -44,8 +44,9 @@ use windmill_common::{
|
||||
CRITICAL_ERROR_CHANNELS_SETTING, CUSTOM_TAGS_SETTING, DEFAULT_TAGS_PER_WORKSPACE_SETTING,
|
||||
DEFAULT_TAGS_WORKSPACES_SETTING, EMAIL_DOMAIN_SETTING, ENV_SETTINGS,
|
||||
EXPOSE_DEBUG_METRICS_SETTING, EXPOSE_METRICS_SETTING, EXTRA_PIP_INDEX_URL_SETTING,
|
||||
HTTP_ROUTE_WORKSPACED_ROUTE_SETTING, HUB_API_SECRET_SETTING, HUB_BASE_URL_SETTING,
|
||||
INDEXER_SETTING, INSTANCE_EVENTS_WEBHOOK_SETTING, INSTANCE_PYTHON_VERSION_SETTING,
|
||||
FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING, HTTP_ROUTE_WORKSPACED_ROUTE_SETTING,
|
||||
HUB_API_SECRET_SETTING, HUB_BASE_URL_SETTING, INDEXER_SETTING,
|
||||
INSTANCE_EVENTS_WEBHOOK_SETTING, INSTANCE_PYTHON_VERSION_SETTING,
|
||||
JOB_DEFAULT_TIMEOUT_SECS_SETTING, JOB_ISOLATION_SETTING, JWT_SECRET_SETTING,
|
||||
KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, MAVEN_REPOS_SETTING, MAVEN_SETTINGS_XML_SETTING,
|
||||
MONITOR_LOGS_ON_OBJECT_STORE_SETTING, NO_DEFAULT_MAVEN_SETTING,
|
||||
@@ -99,17 +100,17 @@ use windmill_worker::{
|
||||
};
|
||||
|
||||
use crate::monitor::{
|
||||
initial_load, load_keep_job_dir, load_metrics_debug_enabled, load_preview_tags_override,
|
||||
load_require_preexisting_user, load_tag_per_workspace_enabled,
|
||||
load_tag_per_workspace_workspaces, monitor_db, reload_app_workspaced_route_setting,
|
||||
reload_audit_log_retention_days_setting, reload_base_url_setting,
|
||||
reload_bunfig_install_scopes_setting, reload_critical_alert_mute_ui_setting,
|
||||
reload_critical_alerts_on_token_expiry_setting, reload_critical_error_channels_setting,
|
||||
reload_extra_pip_index_url_setting, reload_http_route_workspaced_route_setting,
|
||||
reload_hub_api_secret_setting, reload_hub_base_url_setting,
|
||||
reload_instance_events_webhook_setting, reload_job_default_timeout_setting,
|
||||
reload_job_isolation_setting, reload_jwt_secret_setting, reload_license_key,
|
||||
reload_npm_config_registry_setting, reload_otel_tracing_proxy_setting,
|
||||
initial_load, load_fork_workspace_tag_append_fork_suffix, load_keep_job_dir,
|
||||
load_metrics_debug_enabled, load_preview_tags_override, load_require_preexisting_user,
|
||||
load_tag_per_workspace_enabled, load_tag_per_workspace_workspaces, monitor_db,
|
||||
reload_app_workspaced_route_setting, reload_audit_log_retention_days_setting,
|
||||
reload_base_url_setting, reload_bunfig_install_scopes_setting,
|
||||
reload_critical_alert_mute_ui_setting, reload_critical_alerts_on_token_expiry_setting,
|
||||
reload_critical_error_channels_setting, reload_extra_pip_index_url_setting,
|
||||
reload_http_route_workspaced_route_setting, reload_hub_api_secret_setting,
|
||||
reload_hub_base_url_setting, reload_instance_events_webhook_setting,
|
||||
reload_job_default_timeout_setting, reload_job_isolation_setting, reload_jwt_secret_setting,
|
||||
reload_license_key, reload_npm_config_registry_setting, reload_otel_tracing_proxy_setting,
|
||||
reload_pip_index_url_setting, reload_retention_period_setting, reload_scim_token_setting,
|
||||
reload_smtp_config, reload_uv_index_strategy_setting, reload_worker_config, MonitorIteration,
|
||||
};
|
||||
@@ -1699,6 +1700,13 @@ async fn process_notify_event(
|
||||
);
|
||||
}
|
||||
}
|
||||
FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING => {
|
||||
if let Err(e) = load_fork_workspace_tag_append_fork_suffix(db).await {
|
||||
tracing::error!(
|
||||
"Error loading fork workspace tag append fork suffix: {e:#}"
|
||||
);
|
||||
}
|
||||
}
|
||||
PREVIEW_TAGS_OVERRIDE_SETTING => {
|
||||
if let Err(e) = load_preview_tags_override(db).await {
|
||||
tracing::error!("Error loading preview tags override: {e:#}");
|
||||
|
||||
+30
-11
@@ -57,15 +57,15 @@ use windmill_common::{
|
||||
CRITICAL_ALERT_MUTE_UI_SETTING, CRITICAL_ERROR_CHANNELS_SETTING,
|
||||
DEFAULT_TAGS_PER_WORKSPACE_SETTING, DEFAULT_TAGS_WORKSPACES_SETTING,
|
||||
EXPOSE_DEBUG_METRICS_SETTING, EXPOSE_METRICS_SETTING, EXTRA_PIP_INDEX_URL_SETTING,
|
||||
HUB_API_SECRET_SETTING, HUB_BASE_URL_SETTING, INSTANCE_PYTHON_VERSION_SETTING,
|
||||
JOB_DEFAULT_TIMEOUT_SECS_SETTING, JOB_ISOLATION_SETTING, JWT_SECRET_SETTING,
|
||||
KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, MONITOR_LOGS_ON_OBJECT_STORE_SETTING,
|
||||
NPMRC_SETTING, NPM_CONFIG_REGISTRY_SETTING, NUGET_CONFIG_SETTING, OTEL_SETTING,
|
||||
OTEL_TRACING_PROXY_SETTING, PIP_INDEX_URL_SETTING, POWERSHELL_REPO_PAT_SETTING,
|
||||
POWERSHELL_REPO_URL_SETTING, PREVIEW_TAGS_OVERRIDE_SETTING, REQUEST_SIZE_LIMIT_SETTING,
|
||||
REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING, RETENTION_PERIOD_SECS_SETTING,
|
||||
SAML_METADATA_SETTING, SCIM_TOKEN_SETTING, TIMEOUT_WAIT_RESULT_SETTING,
|
||||
UV_INDEX_STRATEGY_SETTING,
|
||||
FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING, HUB_API_SECRET_SETTING,
|
||||
HUB_BASE_URL_SETTING, INSTANCE_PYTHON_VERSION_SETTING, JOB_DEFAULT_TIMEOUT_SECS_SETTING,
|
||||
JOB_ISOLATION_SETTING, JWT_SECRET_SETTING, KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING,
|
||||
MONITOR_LOGS_ON_OBJECT_STORE_SETTING, NPMRC_SETTING, NPM_CONFIG_REGISTRY_SETTING,
|
||||
NUGET_CONFIG_SETTING, OTEL_SETTING, OTEL_TRACING_PROXY_SETTING, PIP_INDEX_URL_SETTING,
|
||||
POWERSHELL_REPO_PAT_SETTING, POWERSHELL_REPO_URL_SETTING, PREVIEW_TAGS_OVERRIDE_SETTING,
|
||||
REQUEST_SIZE_LIMIT_SETTING, REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING,
|
||||
RETENTION_PERIOD_SECS_SETTING, SAML_METADATA_SETTING, SCIM_TOKEN_SETTING,
|
||||
TIMEOUT_WAIT_RESULT_SETTING, UV_INDEX_STRATEGY_SETTING,
|
||||
},
|
||||
indexer::load_indexer_config,
|
||||
jwt::JWT_SECRET,
|
||||
@@ -79,8 +79,9 @@ use windmill_common::{
|
||||
load_periodic_bash_script_interval_from_env, load_whitelist_env_vars_from_env,
|
||||
load_worker_config, reload_custom_tags_setting, store_pull_query,
|
||||
store_suspended_pull_query, Connection, WorkerConfig, DEFAULT_TAGS_PER_WORKSPACE,
|
||||
DEFAULT_TAGS_WORKSPACES, INDEXER_CONFIG, PREVIEW_TAGS_OVERRIDE, SCRIPT_TOKEN_EXPIRY,
|
||||
SMTP_CONFIG, WINDMILL_DIR, WORKER_CONFIG, WORKER_GROUP,
|
||||
DEFAULT_TAGS_WORKSPACES, FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX, INDEXER_CONFIG,
|
||||
PREVIEW_TAGS_OVERRIDE, SCRIPT_TOKEN_EXPIRY, SMTP_CONFIG, WINDMILL_DIR, WORKER_CONFIG,
|
||||
WORKER_GROUP,
|
||||
},
|
||||
KillpillSender, AUDIT_LOG_RETENTION_DAYS, BASE_URL, CRITICAL_ALERTS_ON_DB_OVERSIZE,
|
||||
CRITICAL_ALERTS_ON_TOKEN_EXPIRY, CRITICAL_ALERT_MUTE_UI_ENABLED, CRITICAL_ERROR_CHANNELS, DB,
|
||||
@@ -236,6 +237,10 @@ pub async fn initial_load(
|
||||
tracing::error!("Error loading default tag per workpsace workspaces: {e:#}");
|
||||
}
|
||||
|
||||
if let Err(e) = load_fork_workspace_tag_append_fork_suffix(db).await {
|
||||
tracing::error!("Error loading fork workspace tag append fork suffix: {e:#}");
|
||||
}
|
||||
|
||||
if let Err(e) = load_preview_tags_override(db).await {
|
||||
tracing::error!("Error loading preview tags override: {e:#}");
|
||||
}
|
||||
@@ -510,6 +515,20 @@ pub async fn load_preview_tags_override(db: &DB) -> error::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn load_fork_workspace_tag_append_fork_suffix(db: &DB) -> error::Result<()> {
|
||||
let value =
|
||||
load_value_from_global_settings(db, FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING).await;
|
||||
|
||||
match value {
|
||||
Ok(Some(serde_json::Value::Bool(t))) => {
|
||||
FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX.store(t, Ordering::Relaxed)
|
||||
}
|
||||
Ok(None) => FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX.store(false, Ordering::Relaxed),
|
||||
_ => (),
|
||||
};
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn reload_critical_alert_mute_ui_setting(conn: &Connection) -> error::Result<()> {
|
||||
if let Ok(Some(serde_json::Value::Bool(t))) =
|
||||
load_value_from_global_settings_with_conn(conn, CRITICAL_ALERT_MUTE_UI_SETTING, true).await
|
||||
|
||||
@@ -41,7 +41,7 @@ use windmill_common::workspaces::WorkspaceDeploymentUISettings;
|
||||
use windmill_common::workspaces::{
|
||||
check_user_against_rule, get_datatable_resource_from_db_unchecked, DataTable,
|
||||
DataTableCatalogResourceType, DataTableForkBehavior, ProtectionRuleKind, ProtectionRules,
|
||||
ProtectionRuleset, RuleCheckResult, WorkspaceGitSyncSettings,
|
||||
ProtectionRuleset, RuleCheckResult, WorkspaceGitSyncSettings, WM_FORK_PREFIX,
|
||||
};
|
||||
use windmill_common::workspaces::{Ducklake, DucklakeCatalogResourceType};
|
||||
use windmill_common::PgDatabase;
|
||||
@@ -1107,7 +1107,6 @@ async fn edit_deploy_to() -> Result<String> {
|
||||
}
|
||||
|
||||
pub const BANNED_DOMAINS: &str = include_str!("../../windmill-api/banned_domains.txt");
|
||||
pub const WM_FORK_PREFIX: &str = "wm-fork-";
|
||||
pub const MAX_CUSTOM_PROMPT_LENGTH: usize = 5000;
|
||||
|
||||
async fn is_allowed_auto_domain(ApiAuthed { email, .. }: ApiAuthed) -> JsonResult<bool> {
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
use std::collections::HashMap;
|
||||
|
||||
use windmill_api_auth::{require_super_admin, ApiAuthed};
|
||||
use windmill_common::workspaces::WM_FORK_PREFIX;
|
||||
use windmill_common::DB;
|
||||
|
||||
use crate::workspaces::{
|
||||
archive_workspace_impl, check_w_id_conflict, CREATE_WORKSPACE_REQUIRE_SUPERADMIN,
|
||||
WM_FORK_PREFIX,
|
||||
};
|
||||
|
||||
use axum::extract::Query;
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
pub const CUSTOM_TAGS_SETTING: &str = "custom_tags";
|
||||
pub const DEFAULT_TAGS_PER_WORKSPACE_SETTING: &str = "default_tags_per_workspace";
|
||||
pub const DEFAULT_TAGS_WORKSPACES_SETTING: &str = "default_tags_workspaces";
|
||||
pub const FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING: &str =
|
||||
"fork_workspace_tag_append_fork_suffix";
|
||||
pub const PREVIEW_TAGS_OVERRIDE_SETTING: &str = "preview_tags_override";
|
||||
pub const BASE_URL_SETTING: &str = "base_url";
|
||||
pub const WS_BASE_URL_SETTING: &str = "ws_base_url";
|
||||
|
||||
@@ -208,6 +208,7 @@ lazy_static::lazy_static! {
|
||||
|
||||
pub static ref DEFAULT_TAGS_PER_WORKSPACE: AtomicBool = AtomicBool::new(false);
|
||||
pub static ref DEFAULT_TAGS_WORKSPACES: arc_swap::ArcSwap<Option<Vec<String>>> = arc_swap::ArcSwap::from_pointee(None);
|
||||
pub static ref FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX: AtomicBool = AtomicBool::new(false);
|
||||
pub static ref PREVIEW_TAGS_OVERRIDE: AtomicBool = AtomicBool::new(false);
|
||||
|
||||
pub static ref MAX_TIMEOUT: u64 = std::env::var("TIMEOUT")
|
||||
|
||||
@@ -159,6 +159,10 @@ pub enum ObjectType {
|
||||
|
||||
pub const LATEST_GIT_SYNC_SCRIPT_PATH: &str = "hub/28191/sync-script-to-git-repo-windmill";
|
||||
|
||||
/// Prefix used to identify fork workspaces. A workspace whose id starts with this string is a
|
||||
/// fork of another workspace.
|
||||
pub const WM_FORK_PREFIX: &str = "wm-fork-";
|
||||
|
||||
#[derive(Serialize, Deserialize, Debug)]
|
||||
pub struct GitRepositorySettings {
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
|
||||
@@ -3307,16 +3307,15 @@ pub async fn pull(
|
||||
&& !(job.kind.is_preview()
|
||||
&& PREVIEW_TAGS_OVERRIDE.load(std::sync::atomic::Ordering::Relaxed))
|
||||
{
|
||||
let per_workspace = per_workspace_tag(&job.workspace_id).await;
|
||||
let effective_ws = per_workspace_tag(&job.workspace_id, db).await;
|
||||
let base_tag = if job.is_flow() {
|
||||
"flow".to_string()
|
||||
} else {
|
||||
"dependency".to_string()
|
||||
};
|
||||
let tag = if per_workspace {
|
||||
format!("{}-{}", base_tag, job.workspace_id)
|
||||
} else {
|
||||
base_tag
|
||||
let tag = match &effective_ws {
|
||||
Some(ws) => format!("{}-{}", base_tag, ws),
|
||||
None => base_tag,
|
||||
};
|
||||
sqlx::query!(
|
||||
"UPDATE v2_job_queue SET tag = $1, running = false WHERE id = $2",
|
||||
@@ -4469,7 +4468,7 @@ use crate::cloud_usage::increment_usage_async;
|
||||
// Without this, the ~13KB future of push_inner is inlined into every caller's state machine,
|
||||
// causing stack overflows in deeply nested async call chains (e.g. flow execution).
|
||||
pub async fn push<'c, 'd>(
|
||||
_db: &Pool<Postgres>,
|
||||
db: &Pool<Postgres>,
|
||||
tx: PushIsolationLevel<'c>,
|
||||
workspace_id: &str,
|
||||
job_payload: JobPayload,
|
||||
@@ -4499,7 +4498,7 @@ pub async fn push<'c, 'd>(
|
||||
suspended_mode: Option<bool>,
|
||||
) -> Result<(Uuid, Transaction<'c, Postgres>), Error> {
|
||||
Box::pin(push_inner(
|
||||
_db,
|
||||
db,
|
||||
tx,
|
||||
workspace_id,
|
||||
job_payload,
|
||||
@@ -4533,7 +4532,7 @@ pub async fn push<'c, 'd>(
|
||||
|
||||
// #[instrument(level = "trace", skip_all)]
|
||||
async fn push_inner<'c, 'd>(
|
||||
_db: &Pool<Postgres>,
|
||||
db: &Pool<Postgres>,
|
||||
mut tx: PushIsolationLevel<'c>,
|
||||
workspace_id: &str,
|
||||
job_payload: JobPayload,
|
||||
@@ -4565,7 +4564,7 @@ async fn push_inner<'c, 'd>(
|
||||
#[cfg(feature = "cloud")]
|
||||
if *CLOUD_HOSTED {
|
||||
let team_plan_status =
|
||||
windmill_common::workspaces::get_team_plan_status(_db, workspace_id).await?;
|
||||
windmill_common::workspaces::get_team_plan_status(db, workspace_id).await?;
|
||||
// we track only non flow steps
|
||||
let (workspace_usage, user_usage) = if !matches!(
|
||||
job_payload,
|
||||
@@ -4574,11 +4573,11 @@ async fn push_inner<'c, 'd>(
|
||||
// Check current usage with SELECT (fast, no row locks)
|
||||
// Only check user usage for non-premium workspaces
|
||||
let (current_workspace_usage, current_user_usage) =
|
||||
check_usage_limits(_db, workspace_id, email, !team_plan_status.premium).await?;
|
||||
check_usage_limits(db, workspace_id, email, !team_plan_status.premium).await?;
|
||||
|
||||
// Spawn async task to update usage counters in the background
|
||||
increment_usage_async(
|
||||
_db.clone(),
|
||||
db.clone(),
|
||||
workspace_id.to_string(),
|
||||
if !team_plan_status.premium {
|
||||
Some(email.to_string())
|
||||
@@ -4601,7 +4600,7 @@ async fn push_inner<'c, 'd>(
|
||||
};
|
||||
|
||||
if !team_plan_status.premium || team_plan_status.is_past_due {
|
||||
let is_super_admin = is_superadmin_cached(_db, email).await?;
|
||||
let is_super_admin = is_superadmin_cached(db, email).await?;
|
||||
|
||||
#[cfg(feature = "private")]
|
||||
let recovery_email = crate::jobs_ee::SCHEDULE_RECOVERY_HANDLER_USER_EMAIL;
|
||||
@@ -4628,7 +4627,7 @@ async fn push_inner<'c, 'd>(
|
||||
AND id = $1",
|
||||
email
|
||||
)
|
||||
.fetch_optional(_db)
|
||||
.fetch_optional(db)
|
||||
.await?
|
||||
.flatten()
|
||||
.unwrap_or(1)
|
||||
@@ -4648,7 +4647,7 @@ async fn push_inner<'c, 'd>(
|
||||
"SELECT COUNT(j.id) FROM v2_job_queue q JOIN v2_job j USING (id) WHERE j.permissioned_as_email = $1",
|
||||
email
|
||||
)
|
||||
.fetch_one(_db)
|
||||
.fetch_one(db)
|
||||
.await?
|
||||
.unwrap_or(0);
|
||||
|
||||
@@ -4662,7 +4661,7 @@ async fn push_inner<'c, 'd>(
|
||||
"SELECT COUNT(j.id) FROM v2_job_queue q JOIN v2_job j USING (id) WHERE q.running = true AND j.permissioned_as_email = $1",
|
||||
email
|
||||
)
|
||||
.fetch_one(_db)
|
||||
.fetch_one(db)
|
||||
.await?
|
||||
.unwrap_or(0);
|
||||
|
||||
@@ -4684,7 +4683,7 @@ async fn push_inner<'c, 'd>(
|
||||
AND id = $1",
|
||||
workspace_id
|
||||
)
|
||||
.fetch_optional(_db)
|
||||
.fetch_optional(db)
|
||||
.await?
|
||||
.flatten()
|
||||
.unwrap_or(1)
|
||||
@@ -4713,7 +4712,7 @@ async fn push_inner<'c, 'd>(
|
||||
"SELECT COUNT(id) FROM v2_job_queue WHERE workspace_id = $1",
|
||||
workspace_id
|
||||
)
|
||||
.fetch_one(_db)
|
||||
.fetch_one(db)
|
||||
.await?
|
||||
.unwrap_or(0);
|
||||
|
||||
@@ -4727,7 +4726,7 @@ async fn push_inner<'c, 'd>(
|
||||
"SELECT COUNT(id) FROM v2_job_queue WHERE running = true AND workspace_id = $1",
|
||||
workspace_id
|
||||
)
|
||||
.fetch_one(_db)
|
||||
.fetch_one(db)
|
||||
.await?
|
||||
.unwrap_or(0);
|
||||
|
||||
@@ -4829,7 +4828,7 @@ async fn push_inner<'c, 'd>(
|
||||
..Default::default()
|
||||
},
|
||||
JobPayload::FlowNode { id, path } => {
|
||||
let data = cache::flow::fetch_flow(_db, id).await?;
|
||||
let data = cache::flow::fetch_flow(db, id).await?;
|
||||
let value = data.value();
|
||||
let status = Some(FlowStatus::new(value));
|
||||
// Keep inserting `value` if not all workers are updated.
|
||||
@@ -4875,7 +4874,7 @@ async fn push_inner<'c, 'd>(
|
||||
}
|
||||
|
||||
let hub_script =
|
||||
get_full_hub_script_by_path(StripPath(path.clone()), &HTTP_CLIENT, Some(_db))
|
||||
get_full_hub_script_by_path(StripPath(path.clone()), &HTTP_CLIENT, Some(db))
|
||||
.await?;
|
||||
|
||||
JobPayloadUntagged {
|
||||
@@ -5003,7 +5002,7 @@ async fn push_inner<'c, 'd>(
|
||||
Some(restarted_from_val) => {
|
||||
let (_, _, _, step_n, truncated_modules, user_states, cleanup_module) =
|
||||
restarted_flows_resolution(
|
||||
_db,
|
||||
db,
|
||||
workspace_id,
|
||||
restarted_from_val.flow_job_id,
|
||||
restarted_from_val.step_id.as_str(),
|
||||
@@ -5323,7 +5322,7 @@ async fn push_inner<'c, 'd>(
|
||||
user_states,
|
||||
cleanup_module,
|
||||
) = restarted_flows_resolution(
|
||||
_db,
|
||||
db,
|
||||
workspace_id,
|
||||
completed_job_id,
|
||||
step_id.as_str(),
|
||||
@@ -5509,7 +5508,7 @@ async fn push_inner<'c, 'd>(
|
||||
}
|
||||
|
||||
let interpolated_tag = tag.map(|x| interpolate_args(x, &args, workspace_id));
|
||||
let per_workspace = per_workspace_tag(&workspace_id).await;
|
||||
let effective_ws = per_workspace_tag(&workspace_id, db).await;
|
||||
|
||||
let default = || {
|
||||
let ntag = if job_kind.is_flow()
|
||||
@@ -5527,10 +5526,9 @@ async fn push_inner<'c, 'd>(
|
||||
} else {
|
||||
"deno".to_string()
|
||||
};
|
||||
if per_workspace {
|
||||
format!("{}-{}", ntag, workspace_id)
|
||||
} else {
|
||||
ntag
|
||||
match &effective_ws {
|
||||
Some(ws) => format!("{}-{}", ntag, ws),
|
||||
None => ntag,
|
||||
}
|
||||
};
|
||||
|
||||
@@ -5538,10 +5536,9 @@ async fn push_inner<'c, 'd>(
|
||||
if job_kind.is_preview()
|
||||
&& PREVIEW_TAGS_OVERRIDE.load(std::sync::atomic::Ordering::Relaxed)
|
||||
{
|
||||
if per_workspace {
|
||||
format!("preview-{}", workspace_id)
|
||||
} else {
|
||||
"preview".to_string()
|
||||
match &effective_ws {
|
||||
Some(ws) => format!("preview-{}", ws),
|
||||
None => "preview".to_string(),
|
||||
}
|
||||
} else {
|
||||
language
|
||||
@@ -5556,10 +5553,9 @@ async fn push_inner<'c, 'd>(
|
||||
} else {
|
||||
x.as_str()
|
||||
};
|
||||
if per_workspace {
|
||||
format!("{}-{}", tag_lang, workspace_id)
|
||||
} else {
|
||||
tag_lang.to_string()
|
||||
match &effective_ws {
|
||||
Some(ws) => format!("{}-{}", tag_lang, ws),
|
||||
None => tag_lang.to_string(),
|
||||
}
|
||||
})
|
||||
.unwrap_or_else(default)
|
||||
@@ -5710,10 +5706,10 @@ async fn push_inner<'c, 'd>(
|
||||
|
||||
let runnable_settings_handle = windmill_common::runnable_settings::insert_rs(
|
||||
RunnableSettings {
|
||||
debouncing_settings: debouncing_settings.insert_cached(_db).await?,
|
||||
concurrency_settings: concurrency_settings.insert_cached(_db).await?,
|
||||
debouncing_settings: debouncing_settings.insert_cached(db).await?,
|
||||
concurrency_settings: concurrency_settings.insert_cached(db).await?,
|
||||
},
|
||||
_db,
|
||||
db,
|
||||
)
|
||||
.await?;
|
||||
|
||||
|
||||
@@ -1,11 +1,90 @@
|
||||
use windmill_common::worker::{DEFAULT_TAGS_PER_WORKSPACE, DEFAULT_TAGS_WORKSPACES};
|
||||
use sqlx::{Pool, Postgres};
|
||||
use windmill_common::worker::{
|
||||
DEFAULT_TAGS_PER_WORKSPACE, DEFAULT_TAGS_WORKSPACES, FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX,
|
||||
};
|
||||
use windmill_common::workspaces::WM_FORK_PREFIX;
|
||||
|
||||
pub async fn per_workspace_tag(workspace_id: &str) -> bool {
|
||||
const FORK_PARENT_CACHE_TTL_SECS: u64 = 300;
|
||||
|
||||
lazy_static::lazy_static! {
|
||||
// Cache of fork workspace id -> (parent_workspace_id, cached_at).
|
||||
// `parent_workspace_id` is essentially immutable once a fork is created, so a multi-minute TTL
|
||||
// is safe. `None` means the lookup found no parent (or the DB call failed); we still cache it
|
||||
// briefly so that forks missing a parent do not hammer the DB.
|
||||
static ref FORK_PARENT_CACHE: quick_cache::sync::Cache<String, (Option<String>, std::time::Instant)> =
|
||||
quick_cache::sync::Cache::new(500);
|
||||
}
|
||||
|
||||
/// Returns `Some(effective_workspace_tag_id)` if jobs of `workspace_id` should use workspace-
|
||||
/// specific tags, where `effective_workspace_tag_id` is the string embedded in the tag. For forks,
|
||||
/// this is always the parent workspace id, optionally suffixed with `-fork` (controlled by the
|
||||
/// `FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX` instance setting) so admins can route fork jobs to
|
||||
/// dedicated workers. Returns `None` when default (non-workspaced) tags should be used.
|
||||
pub async fn per_workspace_tag(workspace_id: &str, db: &Pool<Postgres>) -> Option<String> {
|
||||
// Fast path: global toggle off -> no workspacing at all.
|
||||
if !DEFAULT_TAGS_PER_WORKSPACE.load(std::sync::atomic::Ordering::Relaxed) {
|
||||
return None;
|
||||
}
|
||||
|
||||
let is_fork = workspace_id.starts_with(WM_FORK_PREFIX);
|
||||
|
||||
// For forks, always resolve to the parent workspace id; regular workspaces avoid the lookup.
|
||||
let effective_ws_id: String = if is_fork {
|
||||
lookup_fork_parent(workspace_id, db)
|
||||
.await
|
||||
.unwrap_or_else(|| workspace_id.to_string()) // no parent found -> fall back to fork's own id
|
||||
} else {
|
||||
workspace_id.to_string()
|
||||
};
|
||||
|
||||
// Whitelist check is against the resolved (parent) id so that including a parent in the
|
||||
// whitelist transparently covers all of its forks.
|
||||
let per_workspace_workspaces = DEFAULT_TAGS_WORKSPACES.load();
|
||||
DEFAULT_TAGS_PER_WORKSPACE.load(std::sync::atomic::Ordering::Relaxed)
|
||||
&& (per_workspace_workspaces.is_none()
|
||||
|| (**per_workspace_workspaces)
|
||||
.as_ref()
|
||||
.unwrap()
|
||||
.contains(&workspace_id.to_string()))
|
||||
let whitelisted = per_workspace_workspaces.is_none()
|
||||
|| (**per_workspace_workspaces)
|
||||
.as_ref()
|
||||
.unwrap()
|
||||
.contains(&effective_ws_id);
|
||||
|
||||
if !whitelisted {
|
||||
return None;
|
||||
}
|
||||
|
||||
// For forks, optionally append a `-fork` suffix so all forks of a parent share a common
|
||||
// dedicated tag (e.g. `python3-{parent_id}-fork`).
|
||||
let append_fork_suffix =
|
||||
is_fork && FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX.load(std::sync::atomic::Ordering::Relaxed);
|
||||
|
||||
Some(if append_fork_suffix {
|
||||
format!("{}-fork", effective_ws_id)
|
||||
} else {
|
||||
effective_ws_id
|
||||
})
|
||||
}
|
||||
|
||||
/// Returns the parent workspace id for a fork, or `None` if the fork has no parent set (or the
|
||||
/// DB lookup failed). Backed by a short-TTL cache to avoid a DB round-trip per job push.
|
||||
async fn lookup_fork_parent(fork_id: &str, db: &Pool<Postgres>) -> Option<String> {
|
||||
if let Some((parent, cached_at)) = FORK_PARENT_CACHE.get(fork_id) {
|
||||
if cached_at.elapsed().as_secs() < FORK_PARENT_CACHE_TTL_SECS {
|
||||
return parent;
|
||||
}
|
||||
}
|
||||
|
||||
let parent = match sqlx::query_scalar!(
|
||||
"SELECT parent_workspace_id FROM workspace WHERE id = $1",
|
||||
fork_id
|
||||
)
|
||||
.fetch_optional(db)
|
||||
.await
|
||||
{
|
||||
Ok(Some(Some(parent))) => Some(parent),
|
||||
_ => None,
|
||||
};
|
||||
|
||||
FORK_PARENT_CACHE.insert(
|
||||
fork_id.to_string(),
|
||||
(parent.clone(), std::time::Instant::now()),
|
||||
);
|
||||
parent
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
import {
|
||||
DEFAULT_TAGS_PER_WORKSPACE_SETTING,
|
||||
DEFAULT_TAGS_WORKSPACES_SETTING,
|
||||
FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING,
|
||||
PREVIEW_TAGS_OVERRIDE_SETTING
|
||||
} from '$lib/consts'
|
||||
import Toggle from './Toggle.svelte'
|
||||
@@ -14,6 +15,8 @@
|
||||
import { safeSelectItems } from './select/utils.svelte'
|
||||
import Badge from './common/badge/Badge.svelte'
|
||||
import Section from './Section.svelte'
|
||||
import ToggleButtonGroup from './common/toggleButton-v2/ToggleButtonGroup.svelte'
|
||||
import ToggleButton from './common/toggleButton-v2/ToggleButton.svelte'
|
||||
interface Props {
|
||||
defaultTagPerWorkspace?: boolean | undefined
|
||||
defaultTagWorkspaces?: string[]
|
||||
@@ -27,18 +30,21 @@
|
||||
let defaultTags = $state<string[] | undefined>(undefined)
|
||||
let limitToWorkspaces = $state(false)
|
||||
let previewTagsOverride = $state(false)
|
||||
let forkAppendForkSuffix = $state(false)
|
||||
|
||||
// Change detection
|
||||
let originalDefaultTagPerWorkspace = $state<boolean | undefined>(defaultTagPerWorkspace)
|
||||
let originalDefaultTagWorkspaces = $state<string[]>(defaultTagWorkspaces)
|
||||
let originalPreviewTagsOverride = $state(false)
|
||||
let originalForkAppendForkSuffix = $state(false)
|
||||
|
||||
// Detect changes
|
||||
let hasChanges = $derived(
|
||||
originalDefaultTagPerWorkspace !== defaultTagPerWorkspace ||
|
||||
JSON.stringify($state.snapshot(originalDefaultTagWorkspaces)?.sort() || []) !==
|
||||
JSON.stringify($state.snapshot(defaultTagWorkspaces)?.sort() || []) ||
|
||||
originalPreviewTagsOverride !== previewTagsOverride
|
||||
originalPreviewTagsOverride !== previewTagsOverride ||
|
||||
originalForkAppendForkSuffix !== forkAppendForkSuffix
|
||||
)
|
||||
|
||||
let workspaces: string[] = $state([])
|
||||
@@ -59,6 +65,11 @@
|
||||
key: PREVIEW_TAGS_OVERRIDE_SETTING
|
||||
})) as any) ?? false
|
||||
originalPreviewTagsOverride = previewTagsOverride
|
||||
const forkSetting = (await SettingService.getGlobal({
|
||||
key: FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING
|
||||
})) as any
|
||||
forkAppendForkSuffix = forkSetting ?? false
|
||||
originalForkAppendForkSuffix = forkAppendForkSuffix
|
||||
} catch (err) {
|
||||
sendUserToast(`Could not load default tags: ${err}`, true)
|
||||
}
|
||||
@@ -86,11 +97,18 @@
|
||||
value: previewTagsOverride
|
||||
}
|
||||
})
|
||||
await SettingService.setGlobal({
|
||||
key: FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING,
|
||||
requestBody: {
|
||||
value: forkAppendForkSuffix
|
||||
}
|
||||
})
|
||||
|
||||
// Update original state after save
|
||||
originalDefaultTagPerWorkspace = defaultTagPerWorkspace
|
||||
originalDefaultTagWorkspaces = [...(defaultTagWorkspaces || [])]
|
||||
originalPreviewTagsOverride = previewTagsOverride
|
||||
originalForkAppendForkSuffix = forkAppendForkSuffix
|
||||
|
||||
loadDefaultTags()
|
||||
sendUserToast('Saved')
|
||||
@@ -150,6 +168,29 @@
|
||||
/>
|
||||
</div>
|
||||
{#if defaultTagPerWorkspace}
|
||||
<div class="flex flex-col gap-1">
|
||||
<span class="text-xs font-semibold">Fork workspace behavior</span>
|
||||
<ToggleButtonGroup
|
||||
selected={forkAppendForkSuffix ? 'parent-fork' : 'parent'}
|
||||
onSelected={(v) => (forkAppendForkSuffix = v === 'parent-fork')}
|
||||
disabled={!$enterpriseLicense}
|
||||
>
|
||||
{#snippet children({ item })}
|
||||
<ToggleButton
|
||||
value="parent"
|
||||
label="Use parent workspace tag"
|
||||
tooltip={'Fork jobs are tagged with the parent workspace id (e.g. python3-{parent_id}), so they are picked up by workers assigned to the parent workspace.'}
|
||||
{item}
|
||||
/>
|
||||
<ToggleButton
|
||||
value="parent-fork"
|
||||
label="Use parent workspace tag + -fork suffix"
|
||||
tooltip={'Fork jobs are tagged with the parent workspace id and a "-fork" suffix (e.g. python3-{parent_id}-fork), so all forks of a given parent share a dedicated tag. Route these jobs to workers provisioned specifically for forks of the parent.'}
|
||||
{item}
|
||||
/>
|
||||
{/snippet}
|
||||
</ToggleButtonGroup>
|
||||
</div>
|
||||
<Toggle
|
||||
bind:checked={limitToWorkspaces}
|
||||
options={{ right: 'only for some workspaces' }}
|
||||
|
||||
@@ -41,6 +41,7 @@ export const WORKER_S3_BUCKET_SYNC_SETTING = 'worker_s3_bucket_sync'
|
||||
export const CUSTOM_TAGS_SETTING = 'custom_tags'
|
||||
export const DEFAULT_TAGS_PER_WORKSPACE_SETTING = 'default_tags_per_workspace'
|
||||
export const DEFAULT_TAGS_WORKSPACES_SETTING = 'default_tags_workspaces'
|
||||
export const FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING = 'fork_workspace_tag_append_fork_suffix'
|
||||
export const PREVIEW_TAGS_OVERRIDE_SETTING = 'preview_tags_override'
|
||||
|
||||
export const WORKSPACE_SLACK_BOT_TOKEN_PATH = 'f/slack_bot/bot_token'
|
||||
|
||||
Reference in New Issue
Block a user