From e16ded542b815a772581ca1b9e0a9a269141cb87 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Thu, 3 Sep 2026 18:33:22 +0500 Subject: [PATCH] Bind retry recovery to the current queue schema --- internal/providerretry/recovery_linux.go | 4 +++- internal/providerretry/recovery_linux_test.go | 10 ++++++---- 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/internal/providerretry/recovery_linux.go b/internal/providerretry/recovery_linux.go index 4f968e85..26b45078 100644 --- a/internal/providerretry/recovery_linux.go +++ b/internal/providerretry/recovery_linux.go @@ -15,6 +15,8 @@ import ( "strings" "syscall" "time" + + "github.com/NDDev-OpenNetwork/github-actions/internal/queueintent" ) const minimumTerminalRecoveryAge = time.Minute @@ -323,7 +325,7 @@ func requireActiveQueueJob(path, jobID string, now time.Time) error { ExpiresAt time.Time `json:"expires_at"` } `json:"intents"` } - if err := json.Unmarshal(data, &queue); err != nil || queue.SchemaVersion != 4 || queue.Intents == nil { + if err := json.Unmarshal(data, &queue); err != nil || queue.SchemaVersion != queueintent.SchemaVersion || queue.Intents == nil { return errors.New("queue journal identity is invalid") } for _, intent := range queue.Intents { diff --git a/internal/providerretry/recovery_linux_test.go b/internal/providerretry/recovery_linux_test.go index 205e7164..e3e092c3 100644 --- a/internal/providerretry/recovery_linux_test.go +++ b/internal/providerretry/recovery_linux_test.go @@ -11,6 +11,8 @@ import ( "strings" "testing" "time" + + "github.com/NDDev-OpenNetwork/github-actions/internal/queueintent" ) func TestRecoverTerminalUsesExactCASAndGeneration(t *testing.T) { @@ -114,8 +116,8 @@ func TestRecoverExactJobTerminalRequiresLiveQueueProofAndCAS(t *testing.T) { if err := os.WriteFile(journalPath, content, 0o600); err != nil { t.Fatal(err) } - queue := fmt.Sprintf(`{"schema_version":4,"intents":{"github-scale-set-job:v2:3:job-one":{"job_id":"job-one","state":"assigned","expires_at":%q}}}`, - time.Now().UTC().Add(time.Hour).Format(time.RFC3339Nano)) + queue := fmt.Sprintf(`{"schema_version":%d,"intents":{"github-scale-set-job:v2:3:job-one":{"job_id":"job-one","state":"assigned","expires_at":%q}}}`, + queueintent.SchemaVersion, time.Now().UTC().Add(time.Hour).Format(time.RFC3339Nano)) if err := os.WriteFile(queuePath, []byte(queue), 0o600); err != nil { t.Fatal(err) } @@ -142,8 +144,8 @@ func TestRecoverExactJobTerminalRejectsMissingOrInvalidQueueProof(t *testing.T) name string queue string }{ - {name: "job absent", queue: `{"schema_version":4,"intents":{}}`}, - {name: "job running", queue: `{"schema_version":4,"intents":{"job":{"job_id":"job-one","state":"running","expires_at":"2099-01-01T00:00:00Z"}}}`}, + {name: "job absent", queue: fmt.Sprintf(`{"schema_version":%d,"intents":{}}`, queueintent.SchemaVersion)}, + {name: "job running", queue: fmt.Sprintf(`{"schema_version":%d,"intents":{"job":{"job_id":"job-one","state":"running","expires_at":"2099-01-01T00:00:00Z"}}}`, queueintent.SchemaVersion)}, {name: "wrong schema", queue: `{"schema_version":3,"intents":{"job":{"job_id":"job-one","state":"assigned","expires_at":"2099-01-01T00:00:00Z"}}}`}, } { t.Run(test.name, func(t *testing.T) {