From ed2ff6c5e7755fb32bb7c6fdb015456d8088c701 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Wed, 19 Aug 2026 21:25:02 +0200 Subject: [PATCH] fix: scope git-sync concurrency key per repository (#10767) * [ee] fix: scope git-sync concurrency key per repository Co-Authored-By: Claude Opus 5 (1M context) * test: dedupe git repo resource helper, fail loudly on callback timeout Co-Authored-By: Claude Opus 5 (1M context) * fix: reserve the workspace prefix in the git-sync concurrency key cap Co-Authored-By: Claude Opus 5 (1M context) * test: cover the concurrency-key prefix reservation and the pull lane Co-Authored-By: Claude Opus 5 (1M context) * chore: update ee-repo-ref to dff61d6da80d15f8327af99d322c00cc91f784ff This commit updates the EE repository reference after PR #734 was merged in windmill-ee-private. Previous ee-repo-ref: e50a7eca7d7f8771979485f654831b15de59ec25 New ee-repo-ref: dff61d6da80d15f8327af99d322c00cc91f784ff Automated by sync-ee-ref workflow. --------- Co-authored-by: Claude Opus 5 (1M context) Co-authored-by: windmill-internal-app[bot] --- backend/ee-repo-ref.txt | 2 +- .../tests/workspace_dependencies_git_sync.rs | 313 ++++++++++++++++-- backend/windmill-queue/src/jobs.rs | 83 ++++- 3 files changed, 365 insertions(+), 33 deletions(-) diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index 1b9eb3fe0e..f6575c95db 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -483513b70979aa9497cab869837108d948449984 +dff61d6da80d15f8327af99d322c00cc91f784ff diff --git a/backend/windmill-api-integration-tests/tests/workspace_dependencies_git_sync.rs b/backend/windmill-api-integration-tests/tests/workspace_dependencies_git_sync.rs index a5af874f8c..4e25be1fa5 100644 --- a/backend/windmill-api-integration-tests/tests/workspace_dependencies_git_sync.rs +++ b/backend/windmill-api-integration-tests/tests/workspace_dependencies_git_sync.rs @@ -89,27 +89,33 @@ async fn setup_git_sync_config(db: &Pool, sync_script_path: &str) -> a Ok(()) } -/// Create a git repository resource for testing +/// Create a git repository resource at an arbitrary path. #[allow(dead_code)] -async fn create_git_repo_resource(db: &Pool) -> anyhow::Result<()> { +async fn create_git_repo_resource_at(db: &Pool, path: &str) -> anyhow::Result<()> { sqlx::query( r#" INSERT INTO resource (workspace_id, path, value, resource_type, extra_perms, created_by) - VALUES ('test-workspace', 'u/test-user/test_git_repo', $1::jsonb, 'git_repository', '{}'::jsonb, 'test-user') + VALUES ('test-workspace', $2, $1::jsonb, 'git_repository', '{}'::jsonb, 'test-user') ON CONFLICT (workspace_id, path) DO NOTHING "#, ) .bind(json!({ - "url": "https://github.com/test/test.git", + "url": format!("https://github.com/test/{}.git", path.rsplit('/').next().unwrap_or("test")), "branch": "main", "token": "test-token" })) + .bind(path) .execute(db) .await?; - Ok(()) } +/// Create a git repository resource for testing +#[allow(dead_code)] +async fn create_git_repo_resource(db: &Pool) -> anyhow::Result<()> { + create_git_repo_resource_at(db, "u/test-user/test_git_repo").await +} + /// Create a dummy sync script for testing (with version >= 28103 for debouncing support) #[allow(dead_code)] async fn create_sync_script(db: &Pool, path: &str) -> anyhow::Result { @@ -591,24 +597,58 @@ async fn test_promotion_individual_branch_debounces_per_path( Ok(()) } +/// Poll until `expected` callbacks for `script_path` are queued, or time out. +#[allow(dead_code)] +async fn wait_for_callback_jobs( + db: &Pool, + script_path: &str, + expected: usize, + timeout: Duration, +) -> anyhow::Result<()> { + let deadline = tokio::time::Instant::now() + timeout; + loop { + let jobs = + get_deployment_callback_jobs(db, script_path, Duration::from_millis(200)).await?; + if jobs.len() >= expected { + return Ok(()); + } + if tokio::time::Instant::now() >= deadline { + anyhow::bail!( + "timed out waiting for {expected} deployment callbacks on {script_path}, got {}", + jobs.len() + ); + } + tokio::time::sleep(Duration::from_millis(100)).await; + } +} + +/// The concurrency key each queued deployment callback was pushed with, read +/// back through the runnable-settings handle the push stored it under. +#[allow(dead_code)] +async fn get_concurrency_keys( + db: &Pool, + script_path: &str, +) -> anyhow::Result> { + let rows: Vec<(String,)> = sqlx::query_as( + r#" + SELECT COALESCE(cs.concurrency_key, '') + FROM v2_job j + JOIN v2_job_queue q ON q.id = j.id + LEFT JOIN runnable_settings rs ON rs.hash = q.runnable_settings_handle + LEFT JOIN concurrency_settings cs ON cs.hash = rs.concurrency_settings + WHERE j.runnable_path = $1 AND j.kind = 'deploymentcallback' + "#, + ) + .bind(script_path) + .fetch_all(db) + .await?; + Ok(rows.into_iter().map(|(k,)| k).collect()) +} + /// Create a second git repository resource for multi-repo tests. #[allow(dead_code)] async fn create_second_git_repo_resource(db: &Pool) -> anyhow::Result<()> { - sqlx::query( - r#" - INSERT INTO resource (workspace_id, path, value, resource_type, extra_perms, created_by) - VALUES ('test-workspace', 'u/test-user/test_git_repo_2', $1::jsonb, 'git_repository', '{}'::jsonb, 'test-user') - ON CONFLICT (workspace_id, path) DO NOTHING - "#, - ) - .bind(json!({ - "url": "https://github.com/test/test2.git", - "branch": "main", - "token": "test-token-2" - })) - .execute(db) - .await?; - Ok(()) + create_git_repo_resource_at(db, "u/test-user/test_git_repo_2").await } /// Configure git sync with TWO promotion-mode repositories pointing at distinct @@ -747,6 +787,239 @@ async fn test_two_promotion_repos_both_enqueue_callback(db: Pool) -> a Ok(()) } +/// Two repositories, first in workspace-wide mode and then in promotion mode: +/// each repo's sync must get its own concurrency lane in both. Sharing +/// `{workspace}:git_sync` serialises unrelated remotes, and a concurrency-limited +/// job is re-queued to an estimate derived from the key's average duration with no +/// wake-up when the slot frees, so one slow repo delays every other repo by +/// multiples of its own runtime. +#[cfg(all(feature = "enterprise", feature = "private"))] +#[sqlx::test(migrations = "../migrations", fixtures("base"))] +async fn test_two_repos_get_distinct_concurrency_lanes(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + + create_folder(&db, "28103").await?; + create_folder(&db, "target").await?; + // Concurrency settings are written through a filesystem-backed cache keyed by + // their own hash: a hash the cache has seen is assumed to be in the database + // already, so a key any earlier run inserted never lands in this test's fresh + // database and `get_concurrency_keys` cannot read it back. Repo paths unique to + // the run keep every key new. + let run_id: u32 = rand::random(); + let repo_a = format!("u/test-user/lane_repo_a_{run_id}"); + let repo_b = format!("u/test-user/lane_repo_b_{run_id}"); + create_git_repo_resource_at(&db, &repo_a).await?; + create_git_repo_resource_at(&db, &repo_b).await?; + let sync_script_path = "f/28103/test_sync_two_workspace_wide_repos"; + create_sync_script(&db, sync_script_path).await?; + + let git_sync_config = json!({ + "include_type": ["script"], + "include_path": ["**"], + "repositories": [ + { + "script_path": sync_script_path, + "git_repo_resource_path": format!("$res:{repo_a}"), + "use_individual_branch": false, + "group_by_folder": false + }, + { + "script_path": sync_script_path, + "git_repo_resource_path": format!("$res:{repo_b}"), + "use_individual_branch": false, + "group_by_folder": false + } + ] + }); + sqlx::query!( + "UPDATE workspace_settings SET git_sync = $1 WHERE workspace_id = $2", + git_sync_config, + "test-workspace" + ) + .execute(&db) + .await?; + + let (client, _port, _server) = init_client(db.clone()).await; + + create_test_script(&client, "f/target/alpha").await?; + wait_for_callback_jobs(&db, sync_script_path, 2, Duration::from_secs(5)).await?; + + let mut conc_keys = get_concurrency_keys(&db, sync_script_path).await?; + conc_keys.sort(); + assert_eq!( + conc_keys, + vec![ + format!("test-workspace:git_sync:{repo_a}"), + format!("test-workspace:git_sync:{repo_b}"), + ], + "each repo must push on its own concurrency lane" + ); + + // Promotion mode keys per branch, but two repos deploying the same object name + // the same branch, so the repo has to be in the key there too. + let promo_script_path = "f/28103/test_sync_two_promotion_lanes"; + create_sync_script(&db, promo_script_path).await?; + let promo_config = json!({ + "include_type": ["script"], + "include_path": ["**"], + "repositories": [ + { + "script_path": promo_script_path, + "git_repo_resource_path": format!("$res:{repo_a}"), + "use_individual_branch": true, + "group_by_folder": false + }, + { + "script_path": promo_script_path, + "git_repo_resource_path": format!("$res:{repo_b}"), + "use_individual_branch": true, + "group_by_folder": false + } + ] + }); + sqlx::query!( + "UPDATE workspace_settings SET git_sync = $1 WHERE workspace_id = $2", + promo_config, + "test-workspace" + ) + .execute(&db) + .await?; + + create_test_script(&client, "f/target/beta").await?; + wait_for_callback_jobs(&db, promo_script_path, 2, Duration::from_secs(5)).await?; + + let mut promo_keys = get_concurrency_keys(&db, promo_script_path).await?; + promo_keys.sort(); + assert_eq!( + promo_keys, + vec![ + format!("test-workspace:git_sync:{repo_a}:script:f/target/beta"), + format!("test-workspace:git_sync:{repo_b}:script:f/target/beta"), + ], + "each repo must push its branch on its own concurrency lane" + ); + + Ok(()) +} + +/// Pulls must NOT follow pushes into a per-repo lane. A pull writes workspace +/// objects, so the workspace-wide lane is what stops two repos applying creates, +/// updates and deletes to the same scripts and flows at once — the safety argument +/// for per-repo push lanes rests on this staying put. +#[cfg(all(feature = "enterprise", feature = "private"))] +#[sqlx::test(migrations = "../migrations", fixtures("base"))] +async fn test_pull_stays_on_the_workspace_lane(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + + create_git_repo_resource_at(&db, "u/test-user/pull_repo").await?; + + let repo: windmill_common::workspaces::GitRepositorySettings = serde_json::from_value(json!({ + "git_repo_resource_path": "$res:u/test-user/pull_repo", + "use_individual_branch": false, + "group_by_folder": false + }))?; + + windmill_git_sync::enqueue_git_pull_job( + &db, + "test-workspace", + &repo, + None, + false, + None, + None, + None, + ) + .await?; + + let keys: Vec<(String,)> = sqlx::query_as( + r#" + SELECT ck.key + FROM v2_job j + JOIN concurrency_key ck ON ck.job_id = j.id + WHERE j.workspace_id = 'test-workspace' AND j.kind = 'deploymentcallback' + "#, + ) + .fetch_all(&db) + .await?; + assert_eq!( + keys.iter().map(|(k,)| k.as_str()).collect::>(), + vec!["test-workspace:git_sync"], + "a pull must share the workspace-wide lane with every other repo's pull" + ); + + Ok(()) +} + +/// Putting the repo in the lane made the key long enough to matter: +/// `concurrency_key.key` is VARCHAR(255) and its INSERT runs inside `push`, so an +/// overflowing key fails the push and the sync job is never created at all. A repo +/// path that does not fit must be hashed down instead. +/// +/// A cloud build narrows the budget further — `resolve_concurrency_key` prepends +/// `{workspace_id}/` there — but `cloud` is a separate cargo feature this test +/// binary does not enable, so this covers the un-prefixed budget only. +#[cfg(all(feature = "enterprise", feature = "private"))] +#[sqlx::test(migrations = "../migrations", fixtures("base"))] +async fn test_long_repo_path_still_enqueues_callback(db: Pool) -> anyhow::Result<()> { + initialize_tracing().await; + + create_folder(&db, "28103").await?; + create_folder(&db, "target").await?; + // Long enough that `{workspace}:git_sync:{repo}` alone exceeds 255. + let long_repo = format!("u/test-user/{}", "x".repeat(228)); + create_git_repo_resource_at(&db, &long_repo).await?; + let sync_script_path = "f/28103/test_sync_long_repo_path"; + create_sync_script(&db, sync_script_path).await?; + + let git_sync_config = json!({ + "include_type": ["script"], + "include_path": ["**"], + "repositories": [{ + "script_path": sync_script_path, + "git_repo_resource_path": format!("$res:{long_repo}"), + "use_individual_branch": false, + "group_by_folder": false + }] + }); + sqlx::query!( + "UPDATE workspace_settings SET git_sync = $1 WHERE workspace_id = $2", + git_sync_config, + "test-workspace" + ) + .execute(&db) + .await?; + + let (client, _port, _server) = init_client(db.clone()).await; + + create_test_script(&client, "f/target/alpha").await?; + wait_for_callback_jobs(&db, sync_script_path, 1, Duration::from_secs(5)).await?; + + // The row exists only if the INSERT inside `push` accepted the key. + let keys: Vec<(String,)> = sqlx::query_as( + r#" + SELECT ck.key + FROM v2_job j + JOIN concurrency_key ck ON ck.job_id = j.id + WHERE j.runnable_path = $1 AND j.kind = 'deploymentcallback' + "#, + ) + .bind(sync_script_path) + .fetch_all(&db) + .await?; + assert_eq!( + keys.len(), + 1, + "expected a stored concurrency key, got {keys:?}" + ); + assert!( + keys[0].0.len() <= 255, + "concurrency key must fit the column: {} chars", + keys[0].0.len() + ); + + Ok(()) +} + /// Promotion mode with group_by_folder: items destined for the same per-folder /// branch must share one debounce key so they accumulate into a single sync /// job; scripts in different folders must get distinct keys. diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index b286082b1c..386416458f 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -4489,6 +4489,35 @@ pub async fn custom_debounce_key( .await } +/// Concurrency lane for a git-sync callback: the workspace's lane, narrowed by +/// `suffix` (the repo, plus its branch where one is knowable). +/// +/// Both destination columns are VARCHAR(255) and `concurrency_key.key` is written +/// inside `push`, so an over-long key aborts the push and the sync job is never +/// created — the suffix is hashed rather than allowed to overflow. +/// `reserved_prefix_len` is what a caller-side prefix will consume afterwards. +fn git_sync_concurrency_key( + workspace_id: &str, + suffix: Option, + reserved_prefix_len: usize, +) -> String { + const MAX_CONCURRENCY_KEY_LEN: usize = 255; + let max_key_len = MAX_CONCURRENCY_KEY_LEN.saturating_sub(reserved_prefix_len); + match suffix { + Some(suffix) => { + let full = format!("{workspace_id}:git_sync:{suffix}"); + if full.len() <= max_key_len { + full + } else { + // SHA-256 hex (64) over the whole suffix, so distinct repos stay + // distinct and the result fits any workspace id (VARCHAR(50)). + format!("{workspace_id}:git_sync:{}", calculate_hash(&suffix)) + } + } + None => format!("{workspace_id}:git_sync"), + } +} + pub fn resolve_debounce_key<'b>( unresolved_debounce_key: Option, runnable_path: &Option, @@ -6399,18 +6428,15 @@ async fn push_inner<'c, 'd>( } } JobPayload::DeploymentCallback { path, debouncing_settings, concurrency_key_append } => { - const MAX_CONCURRENCY_KEY_LEN: usize = 255; - let concurrency_key = match concurrency_key_append { - Some(suffix) => { - let full = format!("{workspace_id}:git_sync:{suffix}"); - if full.len() <= MAX_CONCURRENCY_KEY_LEN { - full - } else { - format!("{workspace_id}:git_sync:{}", calculate_hash(&suffix)) - } - } - None => format!("{workspace_id}:git_sync"), - }; + // `resolve_concurrency_key` prepends `{workspace_id}/` on cloud builds + // (compiled into every EE build), so that is what the key must leave room + // for here. + #[cfg(feature = "cloud")] + let reserved_prefix_len = workspace_id.len() + 1; + #[cfg(not(feature = "cloud"))] + let reserved_prefix_len = 0; + let concurrency_key = + git_sync_concurrency_key(workspace_id, concurrency_key_append, reserved_prefix_len); JobPayloadUntagged { runnable_path: Some(path.clone()), job_kind: JobKind::DeploymentCallback, @@ -7840,3 +7866,36 @@ pub async fn get_same_worker_job( )) }) } + +#[cfg(test)] +mod git_sync_concurrency_key_tests { + use super::git_sync_concurrency_key; + + /// The reservation is what keeps `{workspace_id}/` + the key inside the + /// VARCHAR(255) column on cloud builds. Without it, a suffix that fits the + /// bare budget is emitted verbatim and the prefixed insert fails. + #[test] + fn hashes_a_suffix_that_only_fits_before_the_prefix() { + let ws = "some-workspace"; + let suffix: String = format!("u/user/{}", "x".repeat(220)); + let bare = git_sync_concurrency_key(ws, Some(suffix.clone()), 0); + assert!(bare.len() <= 255 && bare.ends_with(&suffix)); + + let reserved = git_sync_concurrency_key(ws, Some(suffix), ws.len() + 1); + assert!( + ws.len() + 1 + reserved.len() <= 255, + "prefixed key must fit the column, got {}", + ws.len() + 1 + reserved.len() + ); + } + + #[test] + fn keeps_distinct_repos_distinct_when_hashed() { + let ws = "w"; + let long = "x".repeat(300); + let a = git_sync_concurrency_key(ws, Some(format!("u/user/a{long}")), 0); + let b = git_sync_concurrency_key(ws, Some(format!("u/user/b{long}")), 0); + assert_ne!(a, b); + assert!(a.len() <= 255 && b.len() <= 255); + } +}