diff --git a/backend/src/main.rs b/backend/src/main.rs index a7e3a8313f..39d25c3ee5 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -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:#}"); diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index e403ced388..66afaef6f9 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -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 diff --git a/backend/windmill-api-workspaces/src/workspaces.rs b/backend/windmill-api-workspaces/src/workspaces.rs index 3d48caaaa1..b51a3af62f 100644 --- a/backend/windmill-api-workspaces/src/workspaces.rs +++ b/backend/windmill-api-workspaces/src/workspaces.rs @@ -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 { } 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 { diff --git a/backend/windmill-api-workspaces/src/workspaces_extra.rs b/backend/windmill-api-workspaces/src/workspaces_extra.rs index 11f3a42e1b..3bfc45039a 100644 --- a/backend/windmill-api-workspaces/src/workspaces_extra.rs +++ b/backend/windmill-api-workspaces/src/workspaces_extra.rs @@ -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; diff --git a/backend/windmill-common/src/global_settings.rs b/backend/windmill-common/src/global_settings.rs index 77cb3c0202..72a94d805d 100644 --- a/backend/windmill-common/src/global_settings.rs +++ b/backend/windmill-common/src/global_settings.rs @@ -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"; diff --git a/backend/windmill-common/src/worker.rs b/backend/windmill-common/src/worker.rs index b77a3dce60..98e22455c1 100644 --- a/backend/windmill-common/src/worker.rs +++ b/backend/windmill-common/src/worker.rs @@ -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>> = 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") diff --git a/backend/windmill-common/src/workspaces.rs b/backend/windmill-common/src/workspaces.rs index 5909c481fa..0db5e133ee 100644 --- a/backend/windmill-common/src/workspaces.rs +++ b/backend/windmill-common/src/workspaces.rs @@ -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")] diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index f801113300..de91d25f2e 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -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, + db: &Pool, tx: PushIsolationLevel<'c>, workspace_id: &str, job_payload: JobPayload, @@ -4499,7 +4498,7 @@ pub async fn push<'c, 'd>( suspended_mode: Option, ) -> 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, + db: &Pool, 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?; diff --git a/backend/windmill-queue/src/tags.rs b/backend/windmill-queue/src/tags.rs index 469ad4eefd..a0cc3a305e 100644 --- a/backend/windmill-queue/src/tags.rs +++ b/backend/windmill-queue/src/tags.rs @@ -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, 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) -> Option { + // 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) -> Option { + 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 } diff --git a/frontend/src/lib/components/DefaultTagsInner.svelte b/frontend/src/lib/components/DefaultTagsInner.svelte index e57afd9659..9e890696ec 100644 --- a/frontend/src/lib/components/DefaultTagsInner.svelte +++ b/frontend/src/lib/components/DefaultTagsInner.svelte @@ -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(undefined) let limitToWorkspaces = $state(false) let previewTagsOverride = $state(false) + let forkAppendForkSuffix = $state(false) // Change detection let originalDefaultTagPerWorkspace = $state(defaultTagPerWorkspace) let originalDefaultTagWorkspaces = $state(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 @@ /> {#if defaultTagPerWorkspace} +
+ Fork workspace behavior + (forkAppendForkSuffix = v === 'parent-fork')} + disabled={!$enterpriseLicense} + > + {#snippet children({ item })} + + + {/snippet} + +