Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion config/example-runner-1.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion config/example-runner-2.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion config/example-runner-3.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion config/example-runner-4.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion config/example-services.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 5 additions & 2 deletions config/garm-derivative.yaml
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
6 changes: 4 additions & 2 deletions scripts/build-garm-nddev.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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"
Expand All @@ -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"
Expand All @@ -74,6 +75,7 @@ readonly patch_sha256s=(
"fb67643be9a2ddce1eab86182cf844bceda7f6d40b3e8386fc7c5d4fd2caa5ad"
"7e8822d4bd13dcab7990e15df38e828211609f2c20afb4d9721f524a617b0cb2"
"e6eefa3cc56acf049161f8f020ae796f30aacf51cb55c7b14d8db452e53a544b"
"24661113a5fa3db00fc12f57da9764b1c5dc2c112e0a14b50eb39ffa9f6079ae"
)
readonly overlay_paths=(
"third_party/garm/overlay/workers/scaleset/queue_intent.go"
Expand Down
Original file line number Diff line number Diff line change
@@ -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