mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-08-23 08:00:45 +00:00
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) <noreply@anthropic.com> * test: dedupe git repo resource helper, fail loudly on callback timeout Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: reserve the workspace prefix in the git-sync concurrency key cap Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * test: cover the concurrency-key prefix reservation and the pull lane Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * 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) <noreply@anthropic.com> Co-authored-by: windmill-internal-app[bot] <windmill-internal-app[bot]@users.noreply.github.com>
This commit is contained in:
@@ -1 +1 @@
|
||||
483513b70979aa9497cab869837108d948449984
|
||||
dff61d6da80d15f8327af99d322c00cc91f784ff
|
||||
|
||||
@@ -89,27 +89,33 @@ async fn setup_git_sync_config(db: &Pool<Postgres>, 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<Postgres>) -> anyhow::Result<()> {
|
||||
async fn create_git_repo_resource_at(db: &Pool<Postgres>, 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<Postgres>) -> 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<Postgres>, path: &str) -> anyhow::Result<i64> {
|
||||
@@ -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<Postgres>,
|
||||
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<Postgres>,
|
||||
script_path: &str,
|
||||
) -> anyhow::Result<Vec<String>> {
|
||||
let rows: Vec<(String,)> = sqlx::query_as(
|
||||
r#"
|
||||
SELECT COALESCE(cs.concurrency_key, '<no concurrency settings row>')
|
||||
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<Postgres>) -> 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<Postgres>) -> 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<Postgres>) -> 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<Postgres>) -> 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<_>>(),
|
||||
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<Postgres>) -> 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.
|
||||
|
||||
@@ -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<String>,
|
||||
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<String>,
|
||||
runnable_path: &Option<String>,
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user