From 200c85e98a47353e594b99a42368220bbcbf8ee5 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Fri, 21 Aug 2026 02:08:39 +0500 Subject: [PATCH] fix(queue): make scale-set reconciliation reachable --- CHANGELOG.md | 4 + config/example-runner-1.yaml | 2 +- config/example-runner-2.yaml | 2 +- config/example-runner-3.yaml | 2 +- config/example-runner-4.yaml | 2 +- config/example-services.yaml | 2 +- config/garm-derivative.yaml | 7 +- scripts/build-garm-nddev.sh | 6 +- ...achable-scale-set-job-reconciliation.patch | 128 ++++++++++++++++++ 9 files changed, 146 insertions(+), 9 deletions(-) create mode 100644 third_party/garm/patches/0020-reachable-scale-set-job-reconciliation.patch diff --git a/CHANGELOG.md b/CHANGELOG.md index 44b0f090..aef4d1f9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -46,6 +46,10 @@ Versioning. started instead of losing exact lifecycle correlation. A uniquely verified still-queued DB job also rehydrates its expired provisional intent, so an acknowledged JobAssigned message cannot strand valid old work forever. +- Added GARM derivative `v0.2.1-nddev.48` after live `.47` acceptance proved + the existing entity query excludes every scale-set row with + `workflow_job_id=0`. A dedicated scale-set-only SQL query makes cleanup and + rehydration reachable without exposing those rows to webhook pool consumers. ## [0.1.1] - 2026-08-16 diff --git a/config/example-runner-1.yaml b/config/example-runner-1.yaml index ebfac278..cf2ffa24 100644 --- a/config/example-runner-1.yaml +++ b/config/example-runner-1.yaml @@ -6,7 +6,7 @@ platform: control_plane: manager: garm - manager_version: v0.2.1-nddev.47 + manager_version: v0.2.1-nddev.48 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.52 diff --git a/config/example-runner-2.yaml b/config/example-runner-2.yaml index eef756e9..9dbbb5f0 100644 --- a/config/example-runner-2.yaml +++ b/config/example-runner-2.yaml @@ -6,7 +6,7 @@ platform: control_plane: manager: garm - manager_version: v0.2.1-nddev.47 + manager_version: v0.2.1-nddev.48 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.52 diff --git a/config/example-runner-3.yaml b/config/example-runner-3.yaml index 832aabc6..bf182e70 100644 --- a/config/example-runner-3.yaml +++ b/config/example-runner-3.yaml @@ -6,7 +6,7 @@ platform: control_plane: manager: garm - manager_version: v0.2.1-nddev.47 + manager_version: v0.2.1-nddev.48 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.52 diff --git a/config/example-runner-4.yaml b/config/example-runner-4.yaml index 4daec3be..d45f261e 100644 --- a/config/example-runner-4.yaml +++ b/config/example-runner-4.yaml @@ -6,7 +6,7 @@ platform: control_plane: manager: garm - manager_version: v0.2.1-nddev.47 + manager_version: v0.2.1-nddev.48 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.52 diff --git a/config/example-services.yaml b/config/example-services.yaml index d3421347..4ceb6e53 100644 --- a/config/example-services.yaml +++ b/config/example-services.yaml @@ -24,7 +24,7 @@ platform: control_plane: manager: garm - manager_version: v0.2.1-nddev.47 + manager_version: v0.2.1-nddev.48 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.52 diff --git a/config/garm-derivative.yaml b/config/garm-derivative.yaml index 9d0d1160..e7e69054 100644 --- a/config/garm-derivative.yaml +++ b/config/garm-derivative.yaml @@ -1,6 +1,6 @@ schema_version: 1 artifact: garm -derivative_version: v0.2.1-nddev.47 +derivative_version: v0.2.1-nddev.48 upstream: repository: https://github.com/cloudbase/garm release: v0.2.1 @@ -64,6 +64,9 @@ patches: - path: third_party/garm/patches/0019-authoritative-live-job-rehydration.patch sha256: e6eefa3cc56acf049161f8f020ae796f30aacf51cb55c7b14d8db452e53a544b purpose: Rehydrate one bounded queue intent when GitHub authoritatively confirms an acknowledged JobAssigned DB row is still queued, preventing valid old work from deadlocking after its provisional TTL. + - path: third_party/garm/patches/0020-reachable-scale-set-job-reconciliation.patch + sha256: 24661113a5fa3db00fc12f57da9764b1c5dc2c112e0a14b50eb39ffa9f6079ae + purpose: Query scale-set protocol rows through a dedicated bounded SQL path so authoritative cleanup and live-job rehydration are reachable without leaking those rows into ordinary webhook pool consumption. overlays: - path: third_party/garm/overlay/workers/scaleset/queue_intent.go sha256: 889086c63a7d3244efb37cbd5bfbd72db1b5b1ce2078397b69f29cc5f6333792 @@ -91,7 +94,7 @@ build: - sqlite_omit_load_extension reproducible_rebuilds: 2 maximum_required_glibc: "2.34" - binary_sha256: 47de0e484a1e8bdba3399925ac0e08d3fe0fefa429582209777f292bf32f3297 + binary_sha256: d0b7f57714756fe1283d815dc0ec9b624993c92aae1537a4b4638ae2e2c83ed3 runtime_contract: event_driven_scale_set_wake: true event_driven_instance_wake: true diff --git a/scripts/build-garm-nddev.sh b/scripts/build-garm-nddev.sh index 6a5bac11..d990e6ac 100755 --- a/scripts/build-garm-nddev.sh +++ b/scripts/build-garm-nddev.sh @@ -19,7 +19,7 @@ set -Eeuo pipefail # Every value below is the manifest's. Editing one here detaches the build # from the provenance it is reviewed against, which is why the region is # regenerated and compared rather than maintained. -readonly derivative_version="v0.2.1-nddev.47" +readonly derivative_version="v0.2.1-nddev.48" readonly upstream_repository="https://github.com/cloudbase/garm" readonly upstream_commit="154638445c3949c1958b01812f69d9a1e4d82684" readonly build_image="docker.io/library/golang@sha256:116d58cbd88c1297624acc6e967a060012422bacf9930927e23fb719189c6f36" @@ -32,7 +32,7 @@ readonly build_module_mode="vendor" readonly build_tags="osusergo,netgo,sqlite_omit_load_extension" readonly build_reproducible_rebuilds="2" readonly build_maximum_required_glibc="2.34" -readonly expected_binary_sha256="47de0e484a1e8bdba3399925ac0e08d3fe0fefa429582209777f292bf32f3297" +readonly expected_binary_sha256="d0b7f57714756fe1283d815dc0ec9b624993c92aae1537a4b4638ae2e2c83ed3" readonly patch_paths=( "third_party/garm/patches/0001-event-driven-reconciliation.patch" "third_party/garm/patches/0002-central-queue-admission.patch" @@ -53,6 +53,7 @@ readonly patch_paths=( "third_party/garm/patches/0017-authoritative-queue-intent-reconciliation.patch" "third_party/garm/patches/0018-oldest-first-stale-job-reconciliation.patch" "third_party/garm/patches/0019-authoritative-live-job-rehydration.patch" + "third_party/garm/patches/0020-reachable-scale-set-job-reconciliation.patch" ) readonly patch_sha256s=( "2f0571f141e7388d6ea0cb0341549ba5bf5dab26d0006382a71b76655e272d34" @@ -74,6 +75,7 @@ readonly patch_sha256s=( "fb67643be9a2ddce1eab86182cf844bceda7f6d40b3e8386fc7c5d4fd2caa5ad" "7e8822d4bd13dcab7990e15df38e828211609f2c20afb4d9721f524a617b0cb2" "e6eefa3cc56acf049161f8f020ae796f30aacf51cb55c7b14d8db452e53a544b" + "24661113a5fa3db00fc12f57da9764b1c5dc2c112e0a14b50eb39ffa9f6079ae" ) readonly overlay_paths=( "third_party/garm/overlay/workers/scaleset/queue_intent.go" diff --git a/third_party/garm/patches/0020-reachable-scale-set-job-reconciliation.patch b/third_party/garm/patches/0020-reachable-scale-set-job-reconciliation.patch new file mode 100644 index 00000000..e3e2acf4 --- /dev/null +++ b/third_party/garm/patches/0020-reachable-scale-set-job-reconciliation.patch @@ -0,0 +1,128 @@ +diff --git a/database/sql/jobs.go b/database/sql/jobs.go +index bb8b93b..2b10895 100644 +--- a/database/sql/jobs.go ++++ b/database/sql/jobs.go +@@ -481,6 +481,44 @@ func (s *sqlDatabase) ListEntityJobsByStatus(_ context.Context, entityType param + return ret, nil + } + ++// ListEntityScaleSetJobsByStatus lists only scale-set protocol jobs. The ++// ordinary entity query deliberately excludes these rows so webhook pool ++// consumers cannot create a second runner for work owned by a scale set. ++func (s *sqlDatabase) ListEntityScaleSetJobsByStatus(_ context.Context, entityType params.ForgeEntityType, entityID string, status params.JobStatus) ([]params.Job, error) { ++ u, err := uuid.Parse(entityID) ++ if err != nil { ++ return nil, err ++ } ++ var jobs []WorkflowJob ++ query := s.conn.Model(&WorkflowJob{}).Preload("Instance"). ++ Where("status = ?", status). ++ Where("workflow_job_id = 0"). ++ Where("scale_set_job_id != ?", "") ++ switch entityType { ++ case params.ForgeEntityTypeOrganization: ++ query = query.Where("org_id = ?", u) ++ case params.ForgeEntityTypeRepository: ++ query = query.Where("repo_id = ?", u) ++ case params.ForgeEntityTypeEnterprise: ++ query = query.Where("enterprise_id = ?", u) ++ } ++ if err := query.Find(&jobs); err.Error != nil { ++ if errors.Is(err.Error, gorm.ErrRecordNotFound) { ++ return []params.Job{}, nil ++ } ++ return nil, err.Error ++ } ++ ret := make([]params.Job, len(jobs)) ++ for index, job := range jobs { ++ converted, convertErr := sqlWorkflowJobToParamsJob(job) ++ if convertErr != nil { ++ return nil, fmt.Errorf("error converting scale-set job: %w", convertErr) ++ } ++ ret[index] = converted ++ } ++ return ret, nil ++} ++ + func (s *sqlDatabase) ListAllJobs(_ context.Context) ([]params.Job, error) { + var jobs []WorkflowJob + query := s.conn.Model(&WorkflowJob{}) +diff --git a/database/sql/jobs_test.go b/database/sql/jobs_test.go +index 7abaab0..510eadf 100644 +--- a/database/sql/jobs_test.go ++++ b/database/sql/jobs_test.go +@@ -20,6 +20,7 @@ import ( + "testing" + "time" + ++ "github.com/google/uuid" + "github.com/stretchr/testify/suite" + + dbCommon "github.com/cloudbase/garm/database/common" +@@ -57,6 +58,34 @@ func TestJobsTestSuite(t *testing.T) { + suite.Run(t, new(JobsTestSuite)) + } + ++func (s *JobsTestSuite) TestListEntityScaleSetJobsDoesNotLeakIntoWebhookQuery() { ++ endpoint := garmTesting.CreateDefaultGithubEndpoint(s.adminCtx, s.Store, s.T()) ++ credentials := garmTesting.CreateTestGithubCredentials(s.adminCtx, "scale-job-creds", s.Store, s.T(), endpoint) ++ repository, err := s.Store.CreateRepository( ++ s.adminCtx, "example-owner", "example-repo", credentials, "webhook-secret", ++ params.PoolBalancerTypeRoundRobin, false, ++ ) ++ s.Require().NoError(err) ++ entity, err := repository.GetEntity() ++ s.Require().NoError(err) ++ repositoryID := uuid.MustParse(entity.ID) ++ webhook := params.Job{WorkflowJobID: 42, RunID: 100, Status: string(params.JobStatusQueued), Name: "webhook", RepoID: &repositoryID} ++ scaleSet := params.Job{ScaleSetJobID: "00000000-0000-4000-8000-000000000404", RunID: 101, Status: string(params.JobStatusQueued), Name: "scale", RepoID: &repositoryID} ++ _, err = s.Store.CreateOrUpdateJob(s.adminCtx, webhook) ++ s.Require().NoError(err) ++ _, err = s.Store.CreateOrUpdateJob(s.adminCtx, scaleSet) ++ s.Require().NoError(err) ++ ++ ordinary, err := s.Store.ListEntityJobsByStatus(s.adminCtx, params.ForgeEntityTypeRepository, entity.ID, params.JobStatusQueued) ++ s.Require().NoError(err) ++ s.Require().Len(ordinary, 1) ++ s.Equal(int64(42), ordinary[0].WorkflowJobID) ++ rows, err := s.Store.(*sqlDatabase).ListEntityScaleSetJobsByStatus(s.adminCtx, params.ForgeEntityTypeRepository, entity.ID, params.JobStatusQueued) ++ s.Require().NoError(err) ++ s.Require().Len(rows, 1) ++ s.Equal(scaleSet.ScaleSetJobID, rows[0].ScaleSetJobID) ++} ++ + // TestDeleteInactionableJobs verifies the deletion logic for jobs + func (s *JobsTestSuite) TestDeleteInactionableJobs() { + db := s.Store.(*sqlDatabase) +diff --git a/runner/pool/pool.go b/runner/pool/pool.go +index 55c84d0..49d405b 100644 +--- a/runner/pool/pool.go ++++ b/runner/pool/pool.go +@@ -53,6 +53,10 @@ import ( + var nddevRemoveAuthoritativeQueueIntent = scaleSetWorker.NDDevRemoveQueueIntent + var nddevEnsureAuthoritativeQueueIntent = scaleSetWorker.NDDevEnsureQueueIntent + ++type entityScaleSetJobsStore interface { ++ ListEntityScaleSetJobsByStatus(context.Context, params.ForgeEntityType, string, params.JobStatus) ([]params.Job, error) ++} ++ + var ( + poolIDLabelprefix = "runner-pool-id" + controllerLabelPrefix = "runner-controller-id" +@@ -1987,6 +1991,15 @@ func (r *basePoolManager) reconcileStaleJobs() error { + } + return fmt.Errorf("error listing queued jobs: %w", err) + } ++ scaleSetStore, ok := r.store.(entityScaleSetJobsStore) ++ if !ok { ++ return fmt.Errorf("database does not expose bounded scale-set job reconciliation") ++ } ++ scaleSetQueued, err := scaleSetStore.ListEntityScaleSetJobsByStatus(r.ctx, r.entity.EntityType, r.entity.ID, params.JobStatusQueued) ++ if err != nil { ++ return fmt.Errorf("error listing queued scale-set jobs: %w", err) ++ } ++ queued = append(queued, scaleSetQueued...) + + queued = oldestQueuedJobs(queued) + checked := 0 +