fix: hash long dedicated worker tags (#7914)

* fix: hash long dedicated worker tags

* Update frontend/src/lib/components/dedicated_worker.ts

Co-authored-by: claude[bot] <209825114+claude[bot]@users.noreply.github.com>

---------

Co-authored-by: claude[bot] <209825114+claude[bot]@users.noreply.github.com>
This commit is contained in:
hugocasa
2026-02-11 20:29:44 +01:00
committed by GitHub
parent 22f22c2661
commit aaa1b92300
9 changed files with 225 additions and 82 deletions

View File

@@ -12,7 +12,10 @@ use windmill_api_auth::{
check_scopes, maybe_refresh_folders, require_owner_of_path, ApiAuthed,
};
use windmill_common::{
utils::{BulkDeleteRequest, WithStarredInfoQuery, HTTP_CLIENT}, webhook::{WebhookMessage, WebhookShared}, workspaces::{check_user_against_rule, ProtectionRuleKind, RuleCheckResult}, DB
utils::{BulkDeleteRequest, WithStarredInfoQuery, HTTP_CLIENT},
webhook::{WebhookMessage, WebhookShared},
workspaces::{check_user_against_rule, ProtectionRuleKind, RuleCheckResult},
DB,
};
use windmill_queue::schedule::clear_schedule;
@@ -859,10 +862,13 @@ async fn create_script_internal<'c>(
}
};
let runnable_settings_handle = windmill_common::runnable_settings::insert_rs(RunnableSettings {
debouncing_settings: ns.debouncing_settings.insert_cached(&db).await?,
concurrency_settings: ns.concurrency_settings.insert_cached(&db).await?,
}, &db)
let runnable_settings_handle = windmill_common::runnable_settings::insert_rs(
RunnableSettings {
debouncing_settings: ns.debouncing_settings.insert_cached(&db).await?,
concurrency_settings: ns.concurrency_settings.insert_cached(&db).await?,
},
&db,
)
.await?;
let (
@@ -1072,7 +1078,9 @@ async fn create_script_internal<'c>(
}
if needs_lock_gen {
let tag = if ns.dedicated_worker.is_some_and(|x| x) {
Some(format!("{}:{}", &w_id, &ns.path,))
Some(windmill_common::worker::dedicated_worker_tag(
&w_id, &ns.path,
))
} else if ns.tag.as_ref().is_some_and(|x| x.contains("$args[")) {
None
} else {
@@ -1839,7 +1847,9 @@ async fn get_script_by_hash(
tx.commit().await?;
Ok(Json(windmill_common::scripts::prefetch_cached_script_with_starred(r, &db).await?))
Ok(Json(
windmill_common::scripts::prefetch_cached_script_with_starred(r, &db).await?,
))
}
async fn raw_script_by_hash(
@@ -2042,7 +2052,9 @@ async fn archive_script_by_hash(
WebhookMessage::DeleteScript { workspace: w_id, hash: hash.to_string() },
);
Ok(Json(windmill_common::scripts::prefetch_cached_script(script, &db).await?))
Ok(Json(
windmill_common::scripts::prefetch_cached_script(script, &db).await?,
))
}
async fn delete_script_by_hash(
@@ -2097,7 +2109,9 @@ async fn delete_script_by_hash(
WebhookMessage::DeleteScript { workspace: w_id, hash: hash.to_string() },
);
Ok(Json(windmill_common::scripts::prefetch_cached_script(script, &db).await?))
Ok(Json(
windmill_common::scripts::prefetch_cached_script(script, &db).await?,
))
}
#[derive(Deserialize)]

View File

@@ -2135,9 +2135,7 @@ pub async fn resume_suspended_flow_as_owner(
// Check approval conditions (self-approval, required groups, etc.)
if let Some(ref flow_status_value) = flow.flow_status {
if let Ok(flow_status) =
serde_json::from_value::<FlowStatus>(flow_status_value.clone())
{
if let Ok(flow_status) = serde_json::from_value::<FlowStatus>(flow_status_value.clone()) {
let trigger_email = flow.email.as_deref().unwrap_or("");
conditionally_require_authed_user(Some(authed.clone()), flow_status, trigger_email)?;
}
@@ -5224,7 +5222,7 @@ async fn add_batch_jobs(
let tag = if let Some(dedicated_worker) = dedicated_worker {
if dedicated_worker && path.is_some() {
format!("{}:{}", w_id, path.clone().unwrap())
windmill_common::worker::dedicated_worker_tag(&w_id, &path.clone().unwrap())
} else {
format!("{}", language.as_str())
}

View File

@@ -1560,6 +1560,24 @@ pub async fn update_worker_ping_main_loop_query(
// occupancy_rate = $3, memory_usage = $4, wm_memory_usage = $5, vcpus = COALESCE($7, vcpus),
// memory = COALESCE($8, memory), occupancy_rate_15s = $9, occupancy_rate_5m = $10, occupancy_rate_30m = $11 WHERE worker = $6",
const MAX_TAG_LEN: usize = 50;
const HASH_SUFFIX_LEN: usize = 16;
pub fn dedicated_worker_tag(workspace_id: &str, path: &str) -> String {
let full_tag = format!("{}:{}", workspace_id, path);
if full_tag.len() <= MAX_TAG_LEN {
return full_tag;
}
let hash = <sha2::Sha256 as sha2::Digest>::digest(full_tag.as_bytes());
let hex_hash = hex::encode(hash);
let prefix_len = MAX_TAG_LEN - 1 - HASH_SUFFIX_LEN;
format!(
"{}#{}",
&full_tag[..prefix_len],
&hex_hash[..HASH_SUFFIX_LEN]
)
}
pub async fn load_worker_config(
db: &DB,
killpill_tx: KillpillSender,
@@ -1671,7 +1689,7 @@ pub async fn load_worker_config(
if let Some(ref dws) = dedicated_workers.as_ref() {
let mut dedi_tags: Vec<String> = dws
.iter()
.map(|dw| format!("{}:{}", dw.workspace_id, dw.path))
.map(|dw| dedicated_worker_tag(&dw.workspace_id, &dw.path))
.collect();
if std::env::var("ADD_FLOW_TAG").is_ok() {
dedi_tags.push("flow".to_string());
@@ -1679,9 +1697,9 @@ pub async fn load_worker_config(
Some(dedi_tags)
} else if let Some(ref dedicated_worker) = dedicated_worker.as_ref() {
// Fallback to single dedicated worker for backward compatibility
let mut dedi_tags = vec![format!(
"{}:{}",
dedicated_worker.workspace_id, dedicated_worker.path
let mut dedi_tags = vec![dedicated_worker_tag(
&dedicated_worker.workspace_id,
&dedicated_worker.path,
)];
if std::env::var("ADD_FLOW_TAG").is_ok() {
dedi_tags.push("flow".to_string());
@@ -2178,4 +2196,58 @@ mod tests {
result.sort();
assert_eq!(result, vec!["foo", "legacy(^ws1^ws2)", "urgent(ws1+ws2)"]);
}
#[test]
fn test_dedicated_worker_tag_short() {
let tag = dedicated_worker_tag("demo", "u/alice/script");
assert_eq!(tag, "demo:u/alice/script");
assert!(tag.len() <= MAX_TAG_LEN);
}
#[test]
fn test_dedicated_worker_tag_exactly_50() {
// 50 chars exactly should not be hashed
let workspace = "ws";
let path = "a".repeat(50 - workspace.len() - 1); // -1 for ':'
let tag = dedicated_worker_tag(workspace, &path);
assert_eq!(tag.len(), 50);
assert!(!tag.contains('#'));
}
#[test]
fn test_dedicated_worker_tag_long_is_hashed() {
let tag = dedicated_worker_tag(
"my_workspace",
"u/engineering/team/automation/critical_workflow_script_v2",
);
assert_eq!(tag.len(), MAX_TAG_LEN);
assert_eq!(tag, "my_workspace:u/engineering/team/a#5bc26db79926d4f0");
}
#[test]
fn test_dedicated_worker_tag_deterministic() {
let a = dedicated_worker_tag(
"ws",
"some/very/long/path/that/exceeds/the/fifty/char/limit/easily",
);
assert_eq!(a, "ws:some/very/long/path/that/excee#bbb038d4268a0b41");
let b = dedicated_worker_tag(
"ws",
"some/very/long/path/that/exceeds/the/fifty/char/limit/easily",
);
assert_eq!(a, b);
}
#[test]
fn test_dedicated_worker_tag_different_paths_differ() {
let a = dedicated_worker_tag(
"ws",
"some/very/long/path/that/exceeds/the/fifty/char/limit/easily_a",
);
let b = dedicated_worker_tag(
"ws",
"some/very/long/path/that/exceeds/the/fifty/char/limit/easily_b",
);
assert_ne!(a, b);
}
}

View File

@@ -5274,16 +5274,17 @@ async fn push_inner<'c, 'd>(
.unwrap_or_else(|| (None, None));
let tag = if dedicated_worker.is_some_and(|x| x) {
format!(
"{}:{}{}",
workspace_id,
if job_kind == JobKind::Flow || job_kind == JobKind::FlowDependencies {
"flow/"
} else {
""
},
let flow_prefix = if job_kind == JobKind::Flow || job_kind == JobKind::FlowDependencies {
"flow/"
} else {
""
};
let full_path = format!(
"{}{}",
flow_prefix,
runnable_path.clone().expect("dedicated script has a path")
)
);
windmill_common::worker::dedicated_worker_tag(workspace_id, &full_path)
} else {
if tag == Some("".to_string()) {
tag = None;