diff --git a/config/example-runner-1.yaml b/config/example-runner-1.yaml index 829e230c..2e67bfc8 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.83 + manager_version: v0.2.1-nddev.84 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.80 diff --git a/config/example-runner-2.yaml b/config/example-runner-2.yaml index 9fa6bb85..7e038671 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.83 + manager_version: v0.2.1-nddev.84 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.80 diff --git a/config/example-runner-3.yaml b/config/example-runner-3.yaml index 6219848b..03d0ba14 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.83 + manager_version: v0.2.1-nddev.84 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.80 diff --git a/config/example-runner-4.yaml b/config/example-runner-4.yaml index 5d38b498..55ac4c54 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.83 + manager_version: v0.2.1-nddev.84 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.80 diff --git a/config/example-services.yaml b/config/example-services.yaml index dd918d7f..4be42aa1 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.83 + manager_version: v0.2.1-nddev.84 scheduling_mode: scale-set provider: incus provider_version: v0.1.5-nddev.80 diff --git a/config/garm-derivative.yaml b/config/garm-derivative.yaml index e52c8390..0f62c805 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.83 +derivative_version: v0.2.1-nddev.84 upstream: repository: https://github.com/cloudbase/garm release: v0.2.1 @@ -96,11 +96,11 @@ overlays: sha256: 2f7d98f63033d73d972a9e91ca7cca0934e09abbc62fa971b6ed99e328f531b9 purpose: Deterministic terminal-redelivery suppression, disabled-scale-set isolation, provisional bootstrap, authoritative repository enrichment, orphan-start rehydration, phase-entry and expiry, priority, fairness, width, idempotency, acknowledgement and concurrent-selection tests. - path: third_party/garm/overlay/workers/provider/nddev_create_retry.go - sha256: 611d67a74ff626f04cc0240409a22848525a06c2ae4dc7180c4b75001fb4b599 - purpose: Fsync-backed schema-v2 instance-to-intent reservations and per-job attempt leases, pruning inactive terminal jobs while retaining active state, and selecting the oldest eligible exact intent so one terminal or deferred job cannot block newer work in the same scale set. + sha256: 78cf49829a3c4d6a0ff4c0173fd0ca044cf34949a41565619a60e3cbcda907db + purpose: Fsync-backed schema-v2 instance-to-intent reservations and per-job attempt leases, pruning inactive terminal jobs while retaining active state, and selecting the oldest eligible tenant-qualified retry domain so same-named scale sets cannot alias across forge owners. - path: third_party/garm/overlay/workers/provider/nddev_create_retry_test.go - sha256: 9379d69b8a40295e176fdfb0587b299fabda69de9708e63672c41546fd9a3508 - purpose: Prove unique-owner reconstruction, stable per-job budgets, selective terminal pruning, blocked-intent preservation, next-eligible reservation, and fail-closed behavior when every exact intent is blocked or the queue journal is invalid. + sha256: 882e4cc3f13e124ddcc91404d437da69db37857ba3185d650befe8d3b180605b + purpose: Prove unique-owner reconstruction, tenant-qualified shared-capacity ownership, stable per-job budgets, selective terminal pruning, blocked-intent preservation, next-eligible reservation, and fail-closed behavior when every exact intent is blocked or the queue journal is invalid. build: container_image: docker.io/library/golang@sha256:116d58cbd88c1297624acc6e967a060012422bacf9930927e23fb719189c6f36 go_version: go1.26.6 @@ -115,7 +115,7 @@ build: - sqlite_omit_load_extension reproducible_rebuilds: 2 maximum_required_glibc: "2.34" - binary_sha256: f9594a5ff67643a3b1f9e3357db6a96b45e1ec83797ff9ad0189302a264bb238 + binary_sha256: 2c8dcf0f80071a5b488a3c20ac347f18975c7559bdae504091385cf48ac85534 runtime_contract: queue_intent_schema_version: 5 event_driven_scale_set_wake: true diff --git a/scripts/build-garm-nddev.sh b/scripts/build-garm-nddev.sh index 1bd13b6b..81c59a12 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.83" +readonly derivative_version="v0.2.1-nddev.84" 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="f9594a5ff67643a3b1f9e3357db6a96b45e1ec83797ff9ad0189302a264bb238" +readonly expected_binary_sha256="2c8dcf0f80071a5b488a3c20ac347f18975c7559bdae504091385cf48ac85534" readonly patch_paths=( "third_party/garm/patches/0001-event-driven-reconciliation.patch" "third_party/garm/patches/0002-central-queue-admission.patch" @@ -100,8 +100,8 @@ readonly overlay_paths=( readonly overlay_sha256s=( "31486d264f0a2b8dc63bb902ee49163c32e4c8774a6ad5b6d29efad92a228a03" "2f7d98f63033d73d972a9e91ca7cca0934e09abbc62fa971b6ed99e328f531b9" - "611d67a74ff626f04cc0240409a22848525a06c2ae4dc7180c4b75001fb4b599" - "9379d69b8a40295e176fdfb0587b299fabda69de9708e63672c41546fd9a3508" + "78cf49829a3c4d6a0ff4c0173fd0ca044cf34949a41565619a60e3cbcda907db" + "882e4cc3f13e124ddcc91404d437da69db37857ba3185d650befe8d3b180605b" ) readonly overlay_targets=( "workers/scaleset/queue_intent.go" diff --git a/third_party/garm/overlay/workers/provider/nddev_create_retry.go b/third_party/garm/overlay/workers/provider/nddev_create_retry.go index 73072508..b93442a3 100644 --- a/third_party/garm/overlay/workers/provider/nddev_create_retry.go +++ b/third_party/garm/overlay/workers/provider/nddev_create_retry.go @@ -48,6 +48,7 @@ type nddevRetryRecord struct { ProbeOwner string `json:"probe_owner,omitempty"` WakeReason string `json:"wake_reason,omitempty"` ScaleSetName string `json:"scale_set_name,omitempty"` + Owner string `json:"owner,omitempty"` } type nddevRetryJournal struct { @@ -244,6 +245,10 @@ func NDDevScaleSetCreateAllowed(ctx context.Context, scaleSet params.ScaleSet, e return nil } return nddevUpdateRetryJournal(ctx, func(journal *nddevRetryJournal, now time.Time) error { + if record, exists := journal.Records[key]; exists { + record.Owner = nddevRetryOwner(entity) + journal.Records[key] = record + } if err := nddevSharedCapacityCreateAllowed(journal, now, key, scaleSet.Name, false); err != nil { return err } @@ -358,7 +363,11 @@ func nddevRecordProviderCreateFailure(ctx context.Context, key string, providerE record.TerminalUntil = now.Add(nddevRetryExecutionTTL) } else if record.LastErrorClass == "capacity" { record.NextAllowedAt = now.Add(nddevCapacityRetryDelay(key, record.Attempts)) - nddevRecordSharedCapacityFailure(journal, now, record.ScaleSetName) + owner := record.Owner + if owner == "" { + owner = journal.Records[nddevRetryDomainKey(key)].Owner + } + nddevRecordSharedCapacityFailure(journal, now, record.ScaleSetName, owner) } else { // A missing intent is a cancellation race, not load. Keep it cheap // and responsive without teaching it capacity's accumulating delay. @@ -515,7 +524,8 @@ func nddevSharedCapacityCreateAllowed(journal *nddevRetryJournal, now time.Time, return err } owner := shared.ProbeOwner - if owner != "" && filterActive && (shared.ScaleSetName == "" || !activeScaleSets[shared.ScaleSetName]) { + ownerRecord := journal.Records[nddevRetryDomainKey(owner)] + if owner != "" && filterActive && !activeScaleSets[nddevOwnerScaleSetKey(ownerRecord.Owner, ownerRecord.ScaleSetName)] { owner = "" shared.ProbeOwner = "" shared.ScaleSetName = "" @@ -536,6 +546,7 @@ func nddevSharedCapacityCreateAllowed(journal *nddevRetryJournal, now time.Time, if owner != "" && reserve { shared.ProbeOwner = owner shared.ScaleSetName = journal.Records[owner].ScaleSetName + shared.Owner = journal.Records[owner].Owner } } ownerDomain := nddevRetryDomainKey(owner) @@ -555,6 +566,7 @@ func nddevSharedCapacityCreateAllowed(journal *nddevRetryJournal, now time.Time, } shared.ProbeOwner = key shared.ScaleSetName = scaleSetName + shared.Owner = journal.Records[domainKey].Owner shared.WakeReason = "probe-leased" shared.UpdatedAt = now shared.NextAllowedAt = now.Add(nddevRetryAttemptLease) @@ -562,7 +574,7 @@ func nddevSharedCapacityCreateAllowed(journal *nddevRetryJournal, now time.Time, return nil } -func nddevRecordSharedCapacityFailure(journal *nddevRetryJournal, now time.Time, scaleSetName string) { +func nddevRecordSharedCapacityFailure(journal *nddevRetryJournal, now time.Time, scaleSetName, owner string) { shared := journal.Records[nddevCapacityDomainKey] shared.JobID = nddevCapacityDomainKey if shared.Attempts < nddevRetryMaximum { @@ -574,6 +586,7 @@ func nddevRecordSharedCapacityFailure(journal *nddevRetryJournal, now time.Time, shared.TerminalUntil = time.Time{} shared.ProbeOwner = "" shared.ScaleSetName = scaleSetName + shared.Owner = owner shared.WakeReason = "capacity-refused" journal.Records[nddevCapacityDomainKey] = shared } @@ -590,6 +603,11 @@ func nddevGrantSharedCapacityProbe(journal *nddevRetryJournal, now time.Time, ow shared.TerminalUntil = time.Time{} shared.ProbeOwner = owner shared.ScaleSetName = scaleSetName + if owner != "" { + shared.Owner = journal.Records[nddevRetryDomainKey(owner)].Owner + } else { + shared.Owner = "" + } shared.WakeReason = reason journal.Records[nddevCapacityDomainKey] = shared } @@ -600,7 +618,7 @@ func nddevOldestEligibleCapacityDomain(journal *nddevRetryJournal, activeScaleSe if key == nddevCapacityDomainKey || nddevRetryDomainKey(key) != key || record.LastErrorClass != "capacity" { continue } - if filterActive && (record.ScaleSetName == "" || !activeScaleSets[record.ScaleSetName]) { + if filterActive && (record.ScaleSetName == "" || !activeScaleSets[nddevOwnerScaleSetKey(record.Owner, record.ScaleSetName)]) { continue } if selected == "" || record.UpdatedAt.Before(journal.Records[selected].UpdatedAt) || @@ -665,8 +683,8 @@ func nddevReadActiveQueueInventory(now time.Time) (nddevActiveQueueInventory, bo if !intent.ExpiresAt.After(now) { continue } - if name := strings.TrimSpace(intent.ScaleSetName); name != "" { - inventory.ScaleSets[name] = true + if name, owner := strings.TrimSpace(intent.ScaleSetName), strings.TrimSpace(intent.Owner); name != "" { + inventory.ScaleSets[nddevOwnerScaleSetKey(owner, name)] = true } jobID := strings.TrimSpace(intent.JobID) if jobID == "" { @@ -685,6 +703,10 @@ func nddevReadActiveQueueInventory(now time.Time) (nddevActiveQueueInventory, bo return inventory, true, nil } +func nddevOwnerScaleSetKey(owner, scaleSetName string) string { + return strings.TrimSpace(owner) + "\x00" + strings.TrimSpace(scaleSetName) +} + func nddevReadQueueRetryIntentMap() (map[string]nddevQueueRetryIntent, bool, error) { path := strings.TrimSpace(os.Getenv(nddevQueueIntentFileEnv)) if path == "" { @@ -874,6 +896,9 @@ func (j nddevRetryJournal) Validate() error { strings.ContainsAny(record.ScaleSetName, "\r\n\t") { return fmt.Errorf("provider retry record %q has invalid scale-set name", key) } + if len(record.Owner) > 128 || strings.TrimSpace(record.Owner) != record.Owner || strings.ContainsAny(record.Owner, "\r\n\t") { + return fmt.Errorf("provider retry record %q has invalid owner", key) + } } claimed := make(map[string]string, len(j.Reservations)) for instanceName, reservation := range j.Reservations { diff --git a/third_party/garm/overlay/workers/provider/nddev_create_retry_test.go b/third_party/garm/overlay/workers/provider/nddev_create_retry_test.go index 5c9b0cf6..2624c278 100644 --- a/third_party/garm/overlay/workers/provider/nddev_create_retry_test.go +++ b/third_party/garm/overlay/workers/provider/nddev_create_retry_test.go @@ -704,6 +704,53 @@ func TestNDDevCapacityRetrySerializesEveryScaleSetInTheMeasuredFleet(t *testing. } } +func TestNDDevSharedCapacitySkipsOwnerlessAliasFromAnotherTenant(t *testing.T) { + now := time.Date(2026, 8, 26, 19, 40, 0, 0, time.UTC) + originalNow := nddevRetryNow + nddevRetryNow = func() time.Time { return now } + t.Cleanup(func() { nddevRetryNow = originalNow }) + directory := t.TempDir() + retryPath := filepath.Join(directory, "retry.json") + t.Setenv(nddevRetryFileEnv, retryPath) + t.Setenv(nddevRetryLockEnv, filepath.Join(directory, "retry.lock")) + queuePath := filepath.Join(directory, "queue.json") + t.Setenv(nddevQueueIntentFileEnv, queuePath) + queue := fmt.Sprintf(`{"schema_version":5,"intents":{"active":{"key":"active","job_id":"job-active","scale_set_id":5,"scale_set_name":"nddev-linux-integration","owner":"active-org","state":"assigned","queue_time":%q,"expires_at":%q}}}`, + now.Add(-time.Minute).Format(time.RFC3339Nano), now.Add(time.Hour).Format(time.RFC3339Nano)) + if err := os.WriteFile(queuePath, []byte(queue), 0o600); err != nil { + t.Fatal(err) + } + staleDomain := "scale-set:stale-entity:5" + activeDomain := "scale-set:active-entity:5" + journal := nddevRetryJournal{SchemaVersion: nddevRetrySchemaVersion, Records: map[string]nddevRetryRecord{ + staleDomain: {JobID: staleDomain, Attempts: 1, LastErrorClass: "capacity", UpdatedAt: now.Add(-time.Hour), NextAllowedAt: now.Add(-time.Minute), ScaleSetName: "nddev-linux-integration"}, + activeDomain: {JobID: activeDomain, Attempts: 1, LastErrorClass: "capacity", UpdatedAt: now.Add(-time.Minute), NextAllowedAt: now.Add(-time.Second), ScaleSetName: "nddev-linux-integration"}, + nddevCapacityDomainKey: {JobID: nddevCapacityDomainKey, Attempts: 1, LastErrorClass: "capacity", UpdatedAt: now.Add(-time.Minute), NextAllowedAt: now, ScaleSetName: "nddev-linux-integration", WakeReason: "capacity-refused"}, + }, Reservations: map[string]nddevRetryReservation{}} + if err := nddevWriteRetryJournal(retryPath, journal); err != nil { + t.Fatal(err) + } + scaleSet := params.ScaleSet{ScaleSetID: 5, Name: "nddev-linux-integration"} + entity := params.ForgeEntity{ID: "active-entity", Owner: "active-org"} + if err := NDDevScaleSetCreateAllowed(context.Background(), scaleSet, entity); err != nil { + t.Fatalf("exact active tenant was blocked by stale alias: %v", err) + } + concrete := activeDomain + ":job:job-active" + if err := nddevBeforeProviderCreate(context.Background(), concrete, scaleSet.Name); err != nil { + t.Fatalf("exact active tenant did not acquire shared probe: %v", err) + } + updated, err := nddevReadRetryJournal(retryPath) + if err != nil { + t.Fatal(err) + } + if owner := updated.Records[nddevCapacityDomainKey].ProbeOwner; owner != concrete { + t.Fatalf("shared probe owner=%q, want %q", owner, concrete) + } + if updated.Records[activeDomain].Owner != "active-org" || updated.Records[staleDomain].Owner != "" { + t.Fatalf("tenant qualification drifted: active=%#v stale=%#v", updated.Records[activeDomain], updated.Records[staleDomain]) + } +} + func TestNDDevExpiredConcreteProbeOwnerWithoutFailureIsReleased(t *testing.T) { now := time.Date(2026, 8, 23, 8, 18, 54, 0, time.UTC) oldOwner := "scale-set:entity-one:5:instance:failed"