diff --git a/CHANGELOG.md b/CHANGELOG.md index 98af86545..c38745e7c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,7 @@ ## Unreleased +- Bound Nomad's finite control-plane requests while retaining durable recovery identity for uncertain registration, preserving caller cancellation and keeping established exec streams outside the request ceiling. [PR 1916](https://github.com/openclaw/crabbox/pull/1916). Thanks @SebTardif. - Preserve GCP capacity fallback when a bounded error summary omits retry evidence, while keeping user-visible diagnostics redacted and bounded. [PR 1987](https://github.com/openclaw/crabbox/pull/1987). Thanks @steipete. - AWS: add administrator-only legacy cleanup recovery backed by authenticated original CloudTrail allocation evidence, preserving remaining key and access cleanup and recording an atomic scope-recovery audit without force-success. [PR 1975](https://github.com/openclaw/crabbox/pull/1975). - AWS: complete cleanup after a verified empty instance response without skipping owned keys, bind new leases to the original account and Region, and retain unresolved historical cleanup when that authority is missing. [PR 1904](https://github.com/openclaw/crabbox/pull/1904). Thanks @vincentkoc. diff --git a/docs/providers/nomad.md b/docs/providers/nomad.md index 8f7818823..2d3d4b8b2 100644 --- a/docs/providers/nomad.md +++ b/docs/providers/nomad.md @@ -70,6 +70,29 @@ region and namespace values. A reachable ACL-disabled cluster passes without a token; an anonymous `401`/`403` reports the missing token environment variable. Every check prints `mutation=false`. +Finite JSON calls (`agent.self`, regions, namespace information, job registration, +job information, job allocations, evaluation information and deregistration) +have a two-minute request ceiling, preserving earlier caller deadlines and +cancellation. Regions queries carry the caller context and retain sorted results. +Internally created HTTP transports also bound response-header waits to 30 seconds; +injected clients retain their settings. No whole-request client timeout is added +to established allocation exec streams. Exec startup's HTTP node discovery can +encounter the header deadline before the WebSocket connection is established. + +Before registration, Crabbox durably records the exact job and lease identity. +If registration remains uncertain after reconciliation, Crabbox retains that +recovery claim; a single missing-job response does not prove registration was +rejected. `status` and `list` show +`registration-pending` while a submitted job is absent, and `stop`/`cleanup` +retain the claim with an unknown-outcome diagnostic. A prepared attempt that was +never submitted can be removed locally. Matching observed jobs use the existing +ownership-checked removal path; unexpected state is not adopted or deleted. +Run failures expose a kept recovery session only for the exact retained claim, +and warmup failures retain recovery identifiers in their diagnostic. Registration +is not automatically retried. Existing claims without registration markers keep +their confirmed-job behavior. Setup-failure rollback remains adapter-owned and +does not become a retained successful allocation merely because `--keep` is set. + `warmup` creates a Nomad job and local Crabbox claim. The job stays running until explicit `stop` or `cleanup`, even if `--keep` is omitted. A `run` without `--id` creates a fresh job and deletes it after the command unless `--keep` or diff --git a/internal/cli/claim.go b/internal/cli/claim.go index b8516a3b2..b0098ef21 100644 --- a/internal/cli/claim.go +++ b/internal/cli/claim.go @@ -1636,6 +1636,10 @@ func replaceLeaseClaimIfUnchangedDurableReturning(leaseID string, current, repla return replaceLeaseClaimIfUnchangedWithWrite(leaseID, current, replacement, writeLeaseClaimAtomicDurable) } +func replaceLeaseClaimIfUnchangedDurableReturningContext(ctx context.Context, leaseID string, current, replacement leaseClaim) (leaseClaim, error) { + return replaceLeaseClaimTransactionContext(ctx, leaseID, current, replacement, nil, writeLeaseClaimAtomicDurable) +} + func replaceLeaseClaimIfUnchangedDurableAfter(leaseID string, current, replacement leaseClaim, action func() error) (leaseClaim, error) { // This entrypoint binds replacement identity; restore/replace retain the // supplied payload, including incomplete claims used by rollback. @@ -1648,7 +1652,12 @@ func replaceLeaseClaimIfUnchangedWithWrite(leaseID string, current, replacement } func replaceLeaseClaimTransaction(leaseID string, current, replacement leaseClaim, action func() error, write func(string, leaseClaim) error) (leaseClaim, error) { + return replaceLeaseClaimTransactionContext(context.Background(), leaseID, current, replacement, action, write) +} + +func replaceLeaseClaimTransactionContext(ctx context.Context, leaseID string, current, replacement leaseClaim, action func() error, write func(string, leaseClaim) error) (leaseClaim, error) { return transactLeaseClaim(leaseID, leaseClaimTransaction{ + context: ctx, guard: unchangedLeaseClaimGuard(leaseID, current, true), action: claimTransactionAction(action), revision: claimRevisionAfterMutation, diff --git a/internal/cli/claim_lock_test.go b/internal/cli/claim_lock_test.go index ccf8d2dc5..67d91cd24 100644 --- a/internal/cli/claim_lock_test.go +++ b/internal/cli/claim_lock_test.go @@ -269,7 +269,7 @@ func TestClaimSharedFinalizationPreservesReplacement(t *testing.T) { } func TestClaimFenceContextCancelsPublicationWithoutMutation(t *testing.T) { - for _, operation := range []string{"reuse", "publish", "finalize", "shared"} { + for _, operation := range []string{"reuse", "publish", "finalize", "shared", "durable replacement"} { t.Run(operation, func(t *testing.T) { t.Setenv("XDG_STATE_HOME", t.TempDir()) const id = "cbx_shared_publication" @@ -300,6 +300,11 @@ func TestClaimFenceContextCancelsPublicationWithoutMutation(t *testing.T) { return WithDurableLeaseClaimLockContext(ctx, id, func(*LeaseClaim, bool, func() error) error { return action() }) case "finalize": return CleanupLeaseClaimIfUnchangedAfterContext(ctx, id, claim, true, action) + case "durable replacement": + replacement := cloneLeaseClaim(claim) + replacement.Labels = map[string]string{"state": "submitting"} + _, err := ReplaceLeaseClaimIfUnchangedDurableReturningContext(ctx, id, claim, replacement) + return err default: return WithLeaseClaimUnchangedShared(ctx, id, claim, action) } diff --git a/internal/cli/claim_transaction_test.go b/internal/cli/claim_transaction_test.go index bb128d062..6fe439598 100644 --- a/internal/cli/claim_transaction_test.go +++ b/internal/cli/claim_transaction_test.go @@ -3,6 +3,7 @@ package cli import ( "context" "errors" + "fmt" "os" "path/filepath" "reflect" @@ -154,6 +155,44 @@ func TestClaimTransactionContractActionFailureAndRevision(t *testing.T) { } } +func TestClaimTransactionDurableReplacementContext(t *testing.T) { + for _, contextual := range []bool{false, true} { + for _, stale := range []bool{false, true} { + t.Run(fmt.Sprintf("context=%t/stale=%t", contextual, stale), func(t *testing.T) { + stored := seedClaimContract(t) + expected := cloneLeaseClaim(stored) + if stale { + expected.Revision = "older" + } + replacement := cloneLeaseClaim(stored) + replacement.Labels = map[string]string{"state": "submitting"} + var updated leaseClaim + var err error + if contextual { + updated, err = ReplaceLeaseClaimIfUnchangedDurableReturningContext(t.Context(), stored.LeaseID, expected, replacement) + } else { + updated, err = ReplaceLeaseClaimIfUnchangedDurableReturning(stored.LeaseID, expected, replacement) + } + if stale { + if err == nil || !strings.Contains(err.Error(), "claim changed") { + t.Fatalf("stale replacement err=%v", err) + } + assertClaimContractStored(t, stored.LeaseID, stored) + return + } + if err != nil || updated.Revision == "" || updated.Revision == stored.Revision { + t.Fatalf("updated=%#v err=%v", updated, err) + } + replacement.Revision = updated.Revision + if !reflect.DeepEqual(updated, replacement) { + t.Fatalf("replacement policy changed: got=%#v want=%#v", updated, replacement) + } + assertClaimContractStored(t, stored.LeaseID, updated) + }) + } + } +} + // The provider hook observes different revision phases for input endpoints and // action-produced endpoints. It must always run after the exact-claim guard. type claimContractProvider struct { diff --git a/internal/cli/provider_exports.go b/internal/cli/provider_exports.go index e00dc36e0..f7a4a4a2b 100644 --- a/internal/cli/provider_exports.go +++ b/internal/cli/provider_exports.go @@ -290,6 +290,12 @@ func ReplaceLeaseClaimIfUnchangedDurableReturning(leaseID string, current, repla return replaceLeaseClaimIfUnchangedDurableReturning(leaseID, current, replacement) } +// ReplaceLeaseClaimIfUnchangedDurableReturningContext also bounds waiting for +// the claim fence. Existing context-free replacement APIs remain detached. +func ReplaceLeaseClaimIfUnchangedDurableReturningContext(ctx context.Context, leaseID string, current, replacement LeaseClaim) (LeaseClaim, error) { + return replaceLeaseClaimIfUnchangedDurableReturningContext(ctx, leaseID, current, replacement) +} + func ReplaceLeaseClaimIfUnchangedDurableAfter(leaseID string, current, replacement LeaseClaim, action func() error) (LeaseClaim, error) { return replaceLeaseClaimIfUnchangedDurableAfter(leaseID, current, replacement, action) } diff --git a/internal/providers/nomad/claim.go b/internal/providers/nomad/claim.go index 721d36484..914aa854b 100644 --- a/internal/providers/nomad/claim.go +++ b/internal/providers/nomad/claim.go @@ -41,17 +41,6 @@ func claimScope(cfg Config) string { }, "|") } -func writeNomadClaim(cfg Config, leaseID, slug string, repo Repo, reclaim bool, ready allocationReadiness, expiresAt time.Time) (LeaseClaim, error) { - if err := claimLeaseForRepoProviderScopePond(leaseID, slug, providerName, claimScope(cfg), cfg.Pond, repo.Root, cfg.IdleTimeout, reclaim); err != nil { - return LeaseClaim{}, err - } - claim, err := readLeaseClaim(leaseID) - if err != nil { - return LeaseClaim{}, err - } - return updateLeaseClaimLabelsIfUnchanged(leaseID, claim, claimLabels(cfg, leaseID, slug, ready, expiresAt)) -} - func claimLabels(cfg Config, leaseID, slug string, ready allocationReadiness, expiresAt time.Time) map[string]string { labels := map[string]string{ "provider": providerName, @@ -96,6 +85,9 @@ func resolveNomadClaim(cfg Config, id string) (LeaseClaim, error) { } func authorizeClaimScope(cfg Config, claim LeaseClaim) error { + if _, err := registrationState(claim); err != nil { + return err + } if claim.Provider != "" && claim.Provider != providerName { return exit(2, "lease %s belongs to provider=%s, not %s", claim.LeaseID, claim.Provider, providerName) } @@ -118,6 +110,9 @@ func listNomadLeaseClaims() ([]LeaseClaim, error) { } func validateRemoteOwnership(cfg Config, claim LeaseClaim, job *nomadapi.Job) error { + if _, err := registrationState(claim); err != nil { + return err + } if job == nil { return exit(4, "nomad job for lease %s is missing or inaccessible", claim.LeaseID) } diff --git a/internal/providers/nomad/cleanup_ownership_test.go b/internal/providers/nomad/cleanup_ownership_test.go index c517972ad..5ea6def82 100644 --- a/internal/providers/nomad/cleanup_ownership_test.go +++ b/internal/providers/nomad/cleanup_ownership_test.go @@ -64,11 +64,8 @@ func TestNomadDestructionFencesValidationPurgeAndAbsence(t *testing.T) { expireClaim(t, claim, core.ClockNow(b.rt.Clock).Add(-time.Hour)) claim, _ = readLeaseClaim(claim.LeaseID) } - expectedJob := cloneJob(fake.jobs[claim.Labels[claimLabelJobID]]) if operation == "rollback" { - if err := core.RemoveLeaseClaimIfUnchanged(claim.LeaseID, claim); err != nil { - t.Fatal(err) - } + claim = markRegistrationClaim(t, claim, fake.jobs[claim.Labels[claimLabelJobID]], registrationConfirmed) } lockPath := filepath.Join(os.Getenv("XDG_STATE_HOME"), "crabbox", "claim-locks", claim.LeaseID+".json.lock") if err := os.MkdirAll(filepath.Dir(lockPath), 0o700); err != nil { @@ -111,9 +108,9 @@ func TestNomadDestructionFencesValidationPurgeAndAbsence(t *testing.T) { err = b.deleteOwnedRunJob(context.Background(), client, claim) case "rollback": cause := errors.New("setup failed") - err = b.cleanupUnclaimedJob(context.Background(), client, expectedJob, cause) - if err == cause { - err = nil + recovery, failure := b.rollbackRegistration(context.Background(), client, claim, cause) + if recovery != nil || !errors.Is(failure, cause) || !strings.Contains(failure.Error(), "rolled back") { + t.Fatalf("rollback result recovery=%#v err=%v", recovery, failure) } } if err != nil || purges != 1 || reads < 2 { @@ -320,23 +317,26 @@ func TestNomadKeptRunRefreshFencesClaimChangedDuringExecution(t *testing.T) { } } -func TestNomadSetupRollbackRetainsPublishedClaim(t *testing.T) { +func TestNomadSetupRollbackRetainsChangedClaim(t *testing.T) { for _, partial := range []bool{false, true} { t.Run(strconv.FormatBool(partial), func(t *testing.T) { fake := newLifecycleFakeClient() b, _, _ := testBackend(t, fake) claim := createClaim(t, b, "cbx_a44444444444", "setup-crab", "crabbox-a44444444444", "alloc-a") - expected := cloneJob(fake.jobs[claim.Labels[claimLabelJobID]]) + claim = markRegistrationClaim(t, claim, fake.jobs[claim.Labels[claimLabelJobID]], registrationConfirmed) + expected := claim + labels := maps.Clone(claim.Labels) + labels["owner_revision"] = "replacement" if partial { - var err error - claim, err = core.UpdateLeaseClaimLabelsIfUnchanged(claim.LeaseID, claim, nil) - if err != nil { - t.Fatal(err) - } + labels = nil + } + claim, err := core.UpdateLeaseClaimLabelsIfUnchanged(claim.LeaseID, claim, labels) + if err != nil { + t.Fatal(err) } cause := errors.New("publication failed") - err := b.cleanupUnclaimedJob(context.Background(), fake, expected, cause) - if !errors.Is(err, cause) || err == cause || len(fake.deregisters) != 0 { + recovery, err := b.rollbackRegistration(context.Background(), fake, expected, cause) + if recovery != nil || !errors.Is(err, cause) || err == cause || len(fake.deregisters) != 0 { t.Fatalf("err=%v purges=%v", err, fake.deregisters) } assertNomadClaimRetained(t, claim) diff --git a/internal/providers/nomad/client.go b/internal/providers/nomad/client.go index b6a1bf744..68abdc222 100644 --- a/internal/providers/nomad/client.go +++ b/internal/providers/nomad/client.go @@ -10,6 +10,7 @@ import ( "net/http" "net/url" "os" + "sort" "strings" "time" @@ -35,12 +36,29 @@ type liveClient struct { cfg Config } +const ( + defaultNomadControlRequestTimeout = 2 * time.Minute + nomadDefaultResponseHeaderTimeout = 30 * time.Second +) + var ( errNomadCrossOriginRedirect = errors.New("nomad refused cross-origin redirect") errNomadInvalidRedirect = errors.New("nomad refused invalid redirect") errNomadRedirectLimit = errors.New("nomad redirect stopped after 10 redirects") + nomadControlRequestTimeout = defaultNomadControlRequestTimeout ) +// controlRequestContext bounds finite JSON calls, not allocation exec. +func controlRequestContext(ctx context.Context) (context.Context, context.CancelFunc) { + if nomadControlRequestTimeout <= 0 { + return ctx, func() {} + } + if deadline, ok := ctx.Deadline(); ok && time.Until(deadline) <= nomadControlRequestTimeout { + return ctx, func() {} + } + return context.WithTimeout(ctx, nomadControlRequestTimeout) +} + func newNomadClient(cfg Config, rt Runtime) (Client, error) { apiConfig, err := newNomadAPIConfig(cfg, os.Getenv) if err != nil { @@ -76,6 +94,7 @@ func configureNomadHTTPClient(apiConfig *nomadapi.Config, source *http.Client) e source = cleanhttp.DefaultPooledClient() transport = source.Transport.(*http.Transport) } + transport.ResponseHeaderTimeout = nomadDefaultResponseHeaderTimeout transport.TLSHandshakeTimeout = 10 * time.Second transport.TLSClientConfig = &tls.Config{MinVersion: tls.VersionTLS12} transport.ForceAttemptHTTP2 = false @@ -160,6 +179,8 @@ func newNomadAPIConfig(cfg Config, lookup func(string) string) (*nomadapi.Config } func (c liveClient) AgentSelf(ctx context.Context) (*nomadapi.AgentSelf, error) { + ctx, cancel := controlRequestContext(ctx) + defer cancel() var value nomadapi.AgentSelf query := (&nomadapi.QueryOptions{}).WithContext(ctx) if _, err := c.client.Raw().Query("/v1/agent/self", &value, query); err != nil { @@ -168,17 +189,28 @@ func (c liveClient) AgentSelf(ctx context.Context) (*nomadapi.AgentSelf, error) return &value, nil } -func (c liveClient) Regions(context.Context) ([]string, error) { - value, err := c.client.Regions().List() - return value, sanitizeNomadClientError(err) +func (c liveClient) Regions(ctx context.Context) ([]string, error) { + ctx, cancel := controlRequestContext(ctx) + defer cancel() + var value []string + query := (&nomadapi.QueryOptions{}).WithContext(ctx) + if _, err := c.client.Raw().Query("/v1/regions", &value, query); err != nil { + return nil, sanitizeNomadClientError(err) + } + sort.Strings(value) + return value, nil } func (c liveClient) NamespaceInfo(ctx context.Context, namespace string) (*nomadapi.Namespace, error) { + ctx, cancel := controlRequestContext(ctx) + defer cancel() ns, _, err := c.client.Namespaces().Info(namespace, c.queryOptions(ctx)) return ns, sanitizeNomadClientError(err) } func (c liveClient) RegisterJob(ctx context.Context, job *nomadapi.Job) (string, error) { + ctx, cancel := controlRequestContext(ctx) + defer cancel() resp, _, err := c.client.Jobs().RegisterOpts(job, &nomadapi.RegisterOptions{ EnforceIndex: true, ModifyIndex: 0, @@ -193,21 +225,29 @@ func (c liveClient) RegisterJob(ctx context.Context, job *nomadapi.Job) (string, } func (c liveClient) JobInfo(ctx context.Context, jobID string) (*nomadapi.Job, error) { + ctx, cancel := controlRequestContext(ctx) + defer cancel() job, _, err := c.client.Jobs().Info(jobID, c.queryOptions(ctx)) return job, sanitizeNomadClientError(err) } func (c liveClient) JobAllocations(ctx context.Context, jobID string, all bool) ([]*nomadapi.AllocationListStub, error) { + ctx, cancel := controlRequestContext(ctx) + defer cancel() allocs, _, err := c.client.Jobs().Allocations(jobID, all, c.queryOptions(ctx)) return allocs, sanitizeNomadClientError(err) } func (c liveClient) EvaluationInfo(ctx context.Context, evalID string) (*nomadapi.Evaluation, error) { + ctx, cancel := controlRequestContext(ctx) + defer cancel() eval, _, err := c.client.Evaluations().Info(evalID, c.queryOptions(ctx)) return eval, sanitizeNomadClientError(err) } func (c liveClient) DeregisterJob(ctx context.Context, jobID string, purge bool) (string, error) { + ctx, cancel := controlRequestContext(ctx) + defer cancel() evalID, _, err := c.client.Jobs().Deregister(jobID, purge, c.writeOptions(ctx)) return evalID, sanitizeNomadClientError(err) } diff --git a/internal/providers/nomad/client_test.go b/internal/providers/nomad/client_test.go index 2af072b53..08eae1aa1 100644 --- a/internal/providers/nomad/client_test.go +++ b/internal/providers/nomad/client_test.go @@ -6,6 +6,7 @@ import ( "encoding/json" "errors" "fmt" + "io" "net" "net/http" "net/http/httptest" @@ -14,11 +15,209 @@ import ( "strings" "sync/atomic" "testing" + "time" nomadapi "github.com/hashicorp/nomad/api" core "github.com/openclaw/crabbox/internal/cli" + "github.com/openclaw/crabbox/internal/testutil" ) +func TestNomadFiniteControlRequestContexts(t *testing.T) { + for _, tc := range []struct { + name, method, endpoint, body string + call func(Client, context.Context) error + }{ + {"agent", http.MethodGet, "/v1/agent/self", `{}`, func(c Client, ctx context.Context) error { _, err := c.AgentSelf(ctx); return err }}, + {"regions", http.MethodGet, "/v1/regions", `["west","east"]`, func(c Client, ctx context.Context) error { + regions, err := c.Regions(ctx) + if err == nil && strings.Join(regions, ",") != "east,west" { + return fmt.Errorf("regions order=%v", regions) + } + return err + }}, + {"namespace", http.MethodGet, "/v1/namespace/team-a", `{}`, func(c Client, ctx context.Context) error { _, err := c.NamespaceInfo(ctx, "team-a"); return err }}, + {"job", http.MethodGet, "/v1/job/job-one", `{}`, func(c Client, ctx context.Context) error { _, err := c.JobInfo(ctx, "job-one"); return err }}, + {"allocations", http.MethodGet, "/v1/job/job-one/allocations", `[]`, func(c Client, ctx context.Context) error { + _, err := c.JobAllocations(ctx, "job-one", true) + return err + }}, + {"evaluation", http.MethodGet, "/v1/evaluation/eval-one", `{}`, func(c Client, ctx context.Context) error { _, err := c.EvaluationInfo(ctx, "eval-one"); return err }}, + {"register", http.MethodPut, "/v1/jobs", `{"EvalID":"eval-one"}`, func(c Client, ctx context.Context) error { + id := "job-one" + _, err := c.RegisterJob(ctx, &nomadapi.Job{ID: &id}) + return err + }}, + {"deregister", http.MethodDelete, "/v1/job/job-one", `{"EvalID":"eval-one"}`, func(c Client, ctx context.Context) error { _, err := c.DeregisterJob(ctx, "job-one", true); return err }}, + } { + for _, earlier := range []bool{false, true} { + name := "ceiling" + if earlier { + name = "earlier-caller" + } + t.Run(tc.name+"/"+name, func(t *testing.T) { + type contextKey struct{} + ctx := context.WithValue(context.Background(), contextKey{}, "read-request") + if earlier { + var cancel context.CancelFunc + ctx, cancel = context.WithTimeout(ctx, time.Minute) + defer cancel() + } + calls := 0 + var observed context.Context + httpClient := &http.Client{Transport: testutil.RoundTripFunc(func(req *http.Request) (*http.Response, error) { + calls++ + observed = req.Context() + if req.Method != tc.method || req.URL.Path != tc.endpoint { + t.Errorf("request=%s %s", req.Method, req.URL.Path) + } + if req.Context().Value(contextKey{}) != "read-request" { + t.Error("caller context value lost") + } + deadline, ok := req.Context().Deadline() + if !ok { + t.Error("finite control call has no deadline") + } else if earlier { + want, _ := ctx.Deadline() + if !deadline.Equal(want) { + t.Errorf("caller deadline changed: got=%s want=%s", deadline, want) + } + } else if remaining := time.Until(deadline); remaining <= time.Minute || remaining > 2*time.Minute { + t.Errorf("read ceiling remaining=%s, want within second minute", remaining) + } + return &http.Response{StatusCode: http.StatusOK, Header: make(http.Header), ContentLength: int64(len(tc.body)), Body: io.NopCloser(strings.NewReader(tc.body)), Request: req}, nil + })} + client, err := newNomadClient(nomadTestConfig("https://nomad.example.test"), Runtime{HTTP: httpClient}) + if err != nil { + t.Fatal(err) + } + if err := tc.call(client, ctx); err != nil { + t.Fatal(err) + } + if calls != 1 { + t.Fatalf("requests=%d, want 1", calls) + } + if !earlier && observed.Err() != context.Canceled { + t.Error("derived read context not canceled after return") + } + if ctx.Err() != nil { + t.Error("read canceled its caller") + } + }) + } + } +} + +func TestNomadControlTransportDeadlinePolicy(t *testing.T) { + config, err := newNomadAPIConfig(nomadTestConfig("https://nomad.example.test"), func(string) string { return "" }) + if err != nil { + t.Fatal(err) + } + if err := configureNomadHTTPClient(config, nil); err != nil { + t.Fatal(err) + } + if config.HttpClient.Timeout != 0 || config.HttpClient.Transport.(*http.Transport).ResponseHeaderTimeout != 30*time.Second { + t.Fatal("default transport must bound headers without a whole-request timeout") + } + injected := &http.Client{Transport: &http.Transport{ResponseHeaderTimeout: 7 * time.Second}} + if err := configureNomadHTTPClient(config, injected); err != nil { + t.Fatal(err) + } + if config.HttpClient.Timeout != 0 || config.HttpClient.Transport.(*http.Transport).ResponseHeaderTimeout != 7*time.Second || injected.Timeout != 0 { + t.Fatal("injected client deadline policy changed") + } +} + +func TestNomadRegionsStalledHostHonorsCallerDeadline(t *testing.T) { + server := stalledNomadJSONServer() + t.Cleanup(func() { + server.CloseClientConnections() + server.Close() + }) + + client, err := newNomadClient(nomadTestConfig(server.URL), Runtime{}) + if err != nil { + t.Fatal(err) + } + + ctx, cancel := context.WithTimeout(context.Background(), 40*time.Millisecond) + defer cancel() + started := time.Now() + _, err = client.Regions(ctx) + if err == nil || !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("Regions error=%v, want caller deadline", err) + } + if elapsed := time.Since(started); elapsed >= time.Second { + t.Fatalf("Regions returned after %s, want under 1s", elapsed) + } +} + +func TestNomadFallbackBoundsControlJSONWithoutClientTimeout(t *testing.T) { + const controlTimeout = 40 * time.Millisecond + old := nomadControlRequestTimeout + nomadControlRequestTimeout = controlTimeout + t.Cleanup(func() { nomadControlRequestTimeout = old }) + + server := stalledNomadJSONServer() + t.Cleanup(func() { + server.CloseClientConnections() + server.Close() + }) + + client, err := newNomadClient(nomadTestConfig(server.URL), Runtime{}) + if err != nil { + t.Fatal(err) + } + + started := time.Now() + _, err = client.JobInfo(context.Background(), "job-stall") + if err == nil || !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("JobInfo error=%v, want control deadline", err) + } + if elapsed := time.Since(started); elapsed >= time.Second { + t.Fatalf("JobInfo returned after %s, want under 1s", elapsed) + } + + apiConfig, err := newNomadAPIConfig(nomadTestConfig(server.URL), func(string) string { return "" }) + if err != nil { + t.Fatal(err) + } + if err := configureNomadHTTPClient(apiConfig, nil); err != nil { + t.Fatal(err) + } + if apiConfig.HttpClient.Timeout != 0 { + t.Fatalf("shared Nomad HTTP client Timeout=%s, want 0 so exec streams are not cut", apiConfig.HttpClient.Timeout) + } +} + +func TestNomadJobInfoHonorsCallerCancel(t *testing.T) { + server := stalledNomadJSONServer() + t.Cleanup(func() { + server.CloseClientConnections() + server.Close() + }) + + client, err := newNomadClient(nomadTestConfig(server.URL), Runtime{}) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + _, err = client.JobInfo(ctx, "job-cancel") + if err == nil || !errors.Is(err, context.Canceled) { + t.Fatalf("JobInfo error=%v, want canceled", err) + } +} + +func stalledNomadJSONServer() *httptest.Server { + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(`{"`)) + w.(http.Flusher).Flush() + <-r.Context().Done() + })) +} + func TestNomadClientRefusesCrossOriginRedirectBeforeTokenReplay(t *testing.T) { const token = "nomad-test-token" var sinkRequests atomic.Int32 diff --git a/internal/providers/nomad/core.go b/internal/providers/nomad/core.go index 6452666e9..37c578f9a 100644 --- a/internal/providers/nomad/core.go +++ b/internal/providers/nomad/core.go @@ -2,7 +2,6 @@ package nomad import ( "io" - "time" core "github.com/openclaw/crabbox/internal/cli" ) @@ -42,10 +41,6 @@ func allocateClaimLeaseSlug(leaseID, requested string) (string, error) { return core.AllocateClaimLeaseSlug(leaseID, requested) } -func claimLeaseForRepoProviderScopePond(leaseID, slug, provider, providerScope, pond, repoRoot string, idleTimeout time.Duration, reclaim bool) error { - return core.ClaimLeaseForRepoProviderScopePond(leaseID, slug, provider, providerScope, pond, repoRoot, idleTimeout, reclaim) -} - func readLeaseClaim(leaseID string) (LeaseClaim, error) { return core.ReadLeaseClaim(leaseID) } diff --git a/internal/providers/nomad/lifecycle.go b/internal/providers/nomad/lifecycle.go index 377cff7c3..99a0d5c1e 100644 --- a/internal/providers/nomad/lifecycle.go +++ b/internal/providers/nomad/lifecycle.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "maps" "net/http" "strings" "time" @@ -53,10 +54,11 @@ func (b *backend) Warmup(ctx context.Context, req WarmupRequest) error { if err != nil { return err } - leaseID, slug, ready, _, err := b.createJob(ctx, client, req.Repo, req.RequestedSlug, req.Reclaim) + ready, claim, _, err := b.createJob(ctx, client, req.Repo, req.RequestedSlug) if err != nil { return err } + leaseID, slug := claim.LeaseID, claim.Slug fmt.Fprintf(b.rt.Stdout, "leased %s slug=%s provider=%s job=%s allocation=%s task=%s workdir=%s\n", leaseID, slug, providerName, ready.JobID, ready.AllocationID, b.cfg.Nomad.Task, b.cfg.Nomad.Workdir) if !req.Keep { fmt.Fprintf(b.rt.Stderr, "warning: nomad warmup keeps the job until explicit stop\n") @@ -108,9 +110,10 @@ func (b *backend) Run(ctx context.Context, req RunRequest) (RunResult, error) { }, Acquire: func(ctx context.Context) (shared.DelegatedSandbox, error) { var err error - _, _, ready, claim, err = b.createJob(ctx, client, req.Repo, req.RequestedSlug, req.Reclaim) + var recovery *shared.DelegatedSandboxRecovery + ready, claim, recovery, err = b.createJob(ctx, client, req.Repo, req.RequestedSlug) if err != nil { - return shared.DelegatedSandbox{}, err + return shared.DelegatedSandbox{Recovery: recovery}, err } fmt.Fprintf(b.rt.Stderr, "leased %s slug=%s provider=%s job=%s allocation=%s task=%s\n", claim.LeaseID, claim.Slug, providerName, ready.JobID, ready.AllocationID, ready.Task) return bound(), nil @@ -121,14 +124,32 @@ func (b *backend) Run(ctx context.Context, req RunRequest) (RunResult, error) { if err != nil { return shared.DelegatedSandbox{}, err } + state, err := registrationState(claim) + if err != nil { + return shared.DelegatedSandbox{}, err + } + if state == registrationPrepared { + return shared.DelegatedSandbox{}, registrationRecoveryError(claim, exit(4, "registration was not submitted")) + } jobID := claim.Labels[claimLabelJobID] job, err := client.JobInfo(ctx, jobID) if err != nil { + if state == registrationSubmitting && isNotFoundError(err) { + return shared.DelegatedSandbox{}, unresolvedRegistrationError(claim) + } return shared.DelegatedSandbox{}, err } if err := validateRemoteOwnership(b.cfg, claim, job); err != nil { return shared.DelegatedSandbox{}, err } + if err := validateRegistrationMetadata(claim, job); err != nil { + return shared.DelegatedSandbox{}, err + } + if state == registrationSubmitting { + if err := advanceRegistration(ctx, &claim, registrationConfirmed, "", nil); err != nil { + return shared.DelegatedSandbox{}, err + } + } claimedWorkdir := strings.TrimSpace(claim.Labels[claimLabelWorkdir]) if claimedWorkdir != "" && claimedWorkdir != workdir { return shared.DelegatedSandbox{}, exit(2, "nomad lease %s uses workdir %s; requested workdir %s differs; stop the lease or rerun with the matching --nomad-workdir", claim.LeaseID, claimedWorkdir, workdir) @@ -145,10 +166,16 @@ func (b *backend) Run(ctx context.Context, req RunRequest) (RunResult, error) { if err != nil { return shared.DelegatedSandbox{}, err } - claim, err = updateLeaseClaimLabelsIfUnchanged(claim.LeaseID, updated, claimLabels(b.cfg, claim.LeaseID, claim.Slug, ready, claimExpiresAt(claim))) + labels := maps.Clone(updated.Labels) + maps.Copy(labels, claimLabels(b.cfg, claim.LeaseID, claim.Slug, ready, claimExpiresAt(claim))) + claim, err = updateLeaseClaimLabelsIfUnchanged(claim.LeaseID, updated, labels) if err != nil { return shared.DelegatedSandbox{}, err } + } else if registrationAwaitingReadiness(claim) { + if err := advanceRegistration(ctx, &claim, registrationConfirmed, "", &ready); err != nil { + return shared.DelegatedSandbox{}, err + } } return bound(), nil }, @@ -182,69 +209,39 @@ func (b *backend) Run(ctx context.Context, req RunRequest) (RunResult, error) { }) } -func (b *backend) createJob(ctx context.Context, client Client, repo Repo, requestedSlug string, reclaim bool) (string, string, allocationReadiness, LeaseClaim, error) { - leaseID, err := newLeaseID() +func (b *backend) createJob(ctx context.Context, client Client, repo Repo, requestedSlug string) (allocationReadiness, LeaseClaim, *shared.DelegatedSandboxRecovery, error) { + // Fresh acquisition always uses an absent random identity. Reclaim stays in Resolve. + prepared, err := b.prepareRegistration(ctx, repo, requestedSlug) if err != nil { - return "", "", allocationReadiness{}, LeaseClaim{}, err + return allocationReadiness{}, LeaseClaim{}, nil, err } - slug, err := allocateClaimLeaseSlug(leaseID, requestedSlug) + ready, err := b.submitRegistration(ctx, client, prepared) if err != nil { - return "", "", allocationReadiness{}, LeaseClaim{}, err - } - expiresAt := time.Time{} - if b.cfg.TTL > 0 { - expiresAt = core.ClockNow(b.rt.Clock).UTC().Add(b.cfg.TTL) - } - jobID := jobIDForLease(leaseID) - job, err := buildJobSpec(b.cfg, jobSpecInput{LeaseID: leaseID, Slug: slug, JobID: jobID, ExpiresAt: expiresAt}) - if err != nil { - return "", "", allocationReadiness{}, LeaseClaim{}, err - } - evalID, err := client.RegisterJob(ctx, job) - if err != nil { - return "", "", allocationReadiness{}, LeaseClaim{}, err - } - if evalID != "" { - if err := b.waitForEvaluation(ctx, client, evalID); err != nil { - return "", "", allocationReadiness{}, LeaseClaim{}, b.cleanupUnclaimedJob(ctx, client, job, err) + recovery, failure := b.rollbackRegistration(ctx, client, prepared.claim, err) + if recovery != nil { + return allocationReadiness{}, prepared.claim, recovery, registrationRecoveryError(prepared.claim, failure) } + return allocationReadiness{}, LeaseClaim{}, nil, failure } - ready, err := b.waitForAllocation(ctx, client, jobID, b.allocReadyTimeout()) - if err != nil { - return "", "", allocationReadiness{}, LeaseClaim{}, b.cleanupUnclaimedJob(ctx, client, job, err) - } - claim, err := writeNomadClaim(b.cfg, leaseID, slug, repo, reclaim, ready, expiresAt) - if err != nil { - return "", "", allocationReadiness{}, LeaseClaim{}, b.cleanupUnclaimedJob(ctx, client, job, err) - } - return leaseID, slug, ready, claim, nil + return ready, prepared.claim, nil, nil } -func (b *backend) cleanupUnclaimedJob(ctx context.Context, client Client, expected *nomadapi.Job, cause error) error { +func (b *backend) rollbackRegistration(ctx context.Context, client Client, claim LeaseClaim, cause error) (*shared.DelegatedSandboxRecovery, error) { cleanupCtx, cancel := b.cleanupContext(ctx) defer cancel() - jobID := stringValue(expected.ID) - leaseID := expected.Meta[metadataLeaseID] - cleanupErr := core.CleanupLeaseClaimIfUnchangedAfter(leaseID, LeaseClaim{}, false, func() error { - if err := cleanupCtx.Err(); err != nil { - return err - } - job, err := client.JobInfo(cleanupCtx, jobID) + _, cleanupErr := b.removeOwnedJob(cleanupCtx, client, claim, false) + if cleanupErr != nil { + failure := errors.Join(cause, fmt.Errorf("cleanup nomad job %s after setup failure: %w", claim.Labels[claimLabelJobID], cleanupErr)) + recovery, err := registrationRecovery(cleanupCtx, claim) if err != nil { - if isNotFoundError(err) { - return nil - } - return fmt.Errorf("inspect nomad job %s before setup cleanup: %w", jobID, err) + failure = errors.Join(failure, fmt.Errorf("verify retained nomad registration lease=%s: %w", claim.LeaseID, err)) + message := fmt.Sprintf("nomad setup lease=%s job=%s scope=%q: %v; recovery claim could not be verified", claim.LeaseID, claim.Labels[claimLabelJobID], claim.ProviderScope, failure) + return nil, shared.ExitErrorWithCause(core.ExitCodeForError(cause, 1), message, failure) } - if job == nil || stringValue(job.ID) != jobID || !metadataMatches(job.Meta, expected.Meta) { - return exit(4, "refusing cleanup of nomad job %s after setup failure: ownership changed", jobID) - } - return b.deregisterJobAndConfirmAbsent(cleanupCtx, client, jobID) - }) - if cleanupErr != nil { - return errors.Join(cause, fmt.Errorf("cleanup nomad job %s after unclaimed setup failure: %w", jobID, cleanupErr)) + return recovery, shared.ExitErrorWithCause(core.ExitCodeForError(cause, 1), failure.Error(), failure) } - return cause + message := fmt.Sprintf("nomad setup lease=%s job=%s rolled back: %v", claim.LeaseID, claim.Labels[claimLabelJobID], cause) + return nil, shared.ExitErrorWithCause(core.ExitCodeForError(cause, 1), message, cause) } func (b *backend) deregisterJobAndConfirmAbsent(ctx context.Context, client Client, jobID string) error { @@ -293,14 +290,20 @@ func (b *backend) removeOwnedJob(ctx context.Context, client Client, expected Le ctx, cancel := context.WithTimeout(ctx, b.evalTimeout()) defer cancel() missing := false - err := core.RemoveLeaseClaimIfUnchangedAfter(expected.LeaseID, expected, func() error { - // Lock admission is not cancelable; an expired waiter must do no work. + err := core.CleanupLeaseClaimIfUnchangedAfterContext(ctx, expected.LeaseID, expected, true, func() error { if err := ctx.Err(); err != nil { return err } if err := authorizeClaimScope(b.cfg, expected); err != nil { return err } + state, err := registrationState(expected) + if err != nil { + return err + } + if state == registrationPrepared { + return nil // Exact durable prepared state proves dispatch never began. + } jobID := expected.Labels[claimLabelJobID] if strings.TrimSpace(jobID) == "" { return exit(4, "nomad lease %s has no job ID", expected.LeaseID) @@ -308,6 +311,9 @@ func (b *backend) removeOwnedJob(ctx context.Context, client Client, expected Le job, err := client.JobInfo(ctx, jobID) if err != nil { if isNotFoundError(err) { + if state == registrationSubmitting { + return unresolvedRegistrationError(expected) + } missing = true return nil } @@ -316,6 +322,9 @@ func (b *backend) removeOwnedJob(ctx context.Context, client Client, expected Le if requireMissing { return exit(4, "refusing removal of nomad claim %s: job %s reappeared", expected.LeaseID, jobID) } + if err := validateRegistrationMetadata(expected, job); err != nil { + return err + } if err := validateRemoteOwnership(b.cfg, expected, job); err != nil { return err } @@ -468,12 +477,33 @@ func (b *backend) Cleanup(ctx context.Context, req CleanupRequest) error { return err } checked++ + state, err := registrationState(claim) + if err != nil { + return err + } jobID := claim.Labels[claimLabelJobID] + if state == registrationPrepared { + if req.DryRun { + if err := core.VerifyLeaseClaimUnchanged(claim.LeaseID, claim); err != nil { + return err + } + fmt.Fprintf(b.rt.Stdout, "would remove nomad claim lease=%s job=%s reason=registration_not_submitted\n", claim.LeaseID, jobID) + continue + } + if _, err := b.removeOwnedJob(ctx, client, claim, false); err != nil { + return err + } + removed++ + continue + } job, err := client.JobInfo(ctx, jobID) if err != nil { if !isNotFoundError(err) { return err } + if state == registrationSubmitting { + return unresolvedRegistrationError(claim) + } if req.DryRun { if err := core.VerifyLeaseClaimUnchanged(claim.LeaseID, claim); err != nil { return err @@ -496,6 +526,9 @@ func (b *backend) Cleanup(ctx context.Context, req CleanupRequest) error { if err := validateRemoteOwnership(b.cfg, claim, job); err != nil { return err } + if err := validateRegistrationMetadata(claim, job); err != nil { + return err + } if req.DryRun { if err := core.VerifyLeaseClaimUnchanged(claim.LeaseID, claim); err != nil { return err @@ -520,6 +553,10 @@ func (b *backend) statusFromClaim(ctx context.Context, client Client, claim Leas return StatusView{}, err } jobID := claim.Labels[claimLabelJobID] + state, err := registrationState(claim) + if err != nil { + return StatusView{}, err + } base := StatusView{ ID: claim.LeaseID, Slug: claim.Slug, @@ -530,10 +567,18 @@ func (b *backend) statusFromClaim(ctx context.Context, client Client, claim Leas Network: networkPublic, Labels: baseStatusLabels(b.cfg, claim, "not-ready"), } + if state == registrationPrepared { + base.State = "registration-prepared" + base.Labels[claimLabelState] = base.State + return base, nil + } job, err := client.JobInfo(ctx, jobID) if err != nil { if isNotFoundError(err) { base.State = "missing" + if state == registrationSubmitting { + base.State = "registration-pending" + } base.Labels[claimLabelState] = base.State base.Labels["reason"] = err.Error() return base, nil @@ -543,6 +588,9 @@ func (b *backend) statusFromClaim(ctx context.Context, client Client, claim Leas if err := validateRemoteOwnership(b.cfg, claim, job); err != nil { return StatusView{}, err } + if err := validateRegistrationMetadata(claim, job); err != nil { + return StatusView{}, err + } ready, err := b.currentAllocation(ctx, client, jobID) if err != nil { return StatusView{}, err @@ -693,6 +741,10 @@ func baseStatusLabels(cfg Config, claim LeaseClaim, state string) map[string]str if expiresAt := strings.TrimSpace(claim.Labels[claimLabelExpiresAt]); expiresAt != "" { labels[claimLabelExpiresAt] = expiresAt } + if version, exists := claim.Labels[registrationVersionLabel]; exists { + labels[registrationVersionLabel] = version + labels[registrationStateLabel] = claim.Labels[registrationStateLabel] + } return labels } diff --git a/internal/providers/nomad/lifecycle_test.go b/internal/providers/nomad/lifecycle_test.go index 6cb421f8b..507477b57 100644 --- a/internal/providers/nomad/lifecycle_test.go +++ b/internal/providers/nomad/lifecycle_test.go @@ -3,6 +3,8 @@ package nomad import ( "bytes" "context" + "crypto/sha256" + "encoding/hex" "encoding/json" "errors" "io" @@ -10,6 +12,7 @@ import ( "os" "os/exec" "path/filepath" + "sort" "strings" "testing" "time" @@ -394,7 +397,7 @@ func TestWarmupTimingJSONIncludesNomadLease(t *testing.T) { } } -func TestWarmupCleansRegisteredJobWhenReadinessFailsBeforeClaim(t *testing.T) { +func TestWarmupCleansRegisteredJobWhenReadinessFails(t *testing.T) { fake := newLifecycleFakeClient() fake.evalStatus = nomadapi.EvalStatusFailed b, _, _ := testBackend(t, fake) @@ -406,6 +409,10 @@ func TestWarmupCleansRegisteredJobWhenReadinessFailsBeforeClaim(t *testing.T) { if len(fake.deregisters) != 1 { t.Fatalf("deregisters=%v, want one cleanup", fake.deregisters) } + var displayed core.ExitError + if !core.AsExitError(err, &displayed) || !strings.Contains(displayed.Message, fake.deregisters[0]) || !strings.Contains(displayed.Message, "lease=") || !strings.Contains(displayed.Message, "rolled back") || strings.Contains(displayed.Message, "recover with") { + t.Fatalf("warmup rollback diagnostic lost identity or implied retention: %v", err) + } claims, err := listNomadLeaseClaims() if err != nil { t.Fatal(err) @@ -525,14 +532,11 @@ func TestStopRetainsClaimWhenRemovalCannotBeConfirmed(t *testing.T) { func TestSetupFailureRefusesCleanupAfterRemoteOwnershipChanges(t *testing.T) { fake := newLifecycleFakeClient() b, _, _ := testBackend(t, fake) - expected, err := buildJobSpec(b.cfg, jobSpecInput{LeaseID: "cbx_474747474747", Slug: "changed-crab", JobID: "crabbox-474747474747"}) - if err != nil { - t.Fatal(err) - } - fake.jobs[stringValue(expected.ID)] = cloneJob(expected) - fake.jobs[stringValue(expected.ID)].Meta[metadataLeaseID] = "cbx_someone_else" + claim := createClaim(t, b, "cbx_474747474747", "changed-crab", "crabbox-474747474747", "alloc-47") + claim = markRegistrationClaim(t, claim, fake.jobs[claim.Labels[claimLabelJobID]], registrationConfirmed) + fake.jobs[claim.Labels[claimLabelJobID]].Meta[metadataLeaseID] = "cbx_someone_else" - err = b.cleanupUnclaimedJob(context.Background(), fake, expected, errors.New("readiness failed")) + _, err := b.rollbackRegistration(context.Background(), fake, claim, errors.New("readiness failed")) if err == nil || !strings.Contains(err.Error(), "ownership changed") { t.Fatalf("err=%v, want ownership refusal", err) } @@ -541,6 +545,50 @@ func TestSetupFailureRefusesCleanupAfterRemoteOwnershipChanges(t *testing.T) { } } +// Legacy fixture construction intentionally has no registration-attempt markers. +func writeNomadClaim(cfg Config, leaseID, slug string, repo Repo, reclaim bool, ready allocationReadiness, expiresAt time.Time) (LeaseClaim, error) { + if err := core.ClaimLeaseForRepoProviderScopePond(leaseID, slug, providerName, claimScope(cfg), cfg.Pond, repo.Root, cfg.IdleTimeout, reclaim); err != nil { + return LeaseClaim{}, err + } + claim, err := readLeaseClaim(leaseID) + if err != nil { + return LeaseClaim{}, err + } + return updateLeaseClaimLabelsIfUnchanged(leaseID, claim, claimLabels(cfg, leaseID, slug, ready, expiresAt)) +} + +func markRegistrationClaim(t *testing.T, claim LeaseClaim, job *nomadapi.Job, state string) LeaseClaim { + t.Helper() + metadata, err := json.Marshal(job.Meta) + if err != nil { + t.Fatal(err) + } + labels := make(map[string]string, len(claim.Labels)+3) + for key, value := range claim.Labels { + labels[key] = value + } + labels[registrationVersionLabel] = "1" + labels[registrationStateLabel] = state + delete(labels, claimLabelAllocationID) + keys := make([]string, 0, len(job.Meta)) + for key := range job.Meta { + keys = append(keys, key) + } + sort.Strings(keys) + keyJSON, err := json.Marshal(keys) + if err != nil { + t.Fatal(err) + } + digest := sha256.Sum256(metadata) + labels[registrationMetaKeysLabel] = string(keyJSON) + labels[registrationMetaHashLabel] = hex.EncodeToString(digest[:]) + updated, err := core.UpdateLeaseClaimLabelsIfUnchanged(claim.LeaseID, claim, labels) + if err != nil { + t.Fatal(err) + } + return updated +} + func TestCleanupDryRunAndLiveOwnedExpiredClaims(t *testing.T) { fake := newLifecycleFakeClient() b, stdout, _ := testBackend(t, fake) diff --git a/internal/providers/nomad/registration.go b/internal/providers/nomad/registration.go new file mode 100644 index 000000000..c18de18ee --- /dev/null +++ b/internal/providers/nomad/registration.go @@ -0,0 +1,260 @@ +package nomad + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "maps" + "sort" + "strings" + "time" + + nomadapi "github.com/hashicorp/nomad/api" + core "github.com/openclaw/crabbox/internal/cli" + "github.com/openclaw/crabbox/internal/providers/shared" +) + +const ( + registrationVersionLabel = "nomad_registration_version" + registrationStateLabel = "nomad_registration_state" + registrationMetaKeysLabel = "nomad_registration_metadata_keys" + registrationMetaHashLabel = "nomad_registration_metadata_sha256" + registrationEvalLabel = "nomad_registration_evaluation" + registrationPrepared = "prepared" + registrationSubmitting = "submitting" + registrationConfirmed = "confirmed" +) + +type preparedRegistration struct { + job *nomadapi.Job + claim LeaseClaim +} + +func registrationState(claim LeaseClaim) (string, error) { + version, hasVersion := claim.Labels[registrationVersionLabel] + state, hasState := claim.Labels[registrationStateLabel] + if !hasVersion && !hasState { + if _, exists := claim.Labels[registrationMetaKeysLabel]; exists { + return "", exit(4, "nomad lease %s has incomplete registration markers", claim.LeaseID) + } + if _, exists := claim.Labels[registrationMetaHashLabel]; exists { + return "", exit(4, "nomad lease %s has incomplete registration markers", claim.LeaseID) + } + if _, exists := claim.Labels[registrationEvalLabel]; exists { + return "", exit(4, "nomad lease %s has incomplete registration markers", claim.LeaseID) + } + return registrationConfirmed, nil // Legacy claims were published after registration. + } + if version != "1" { + return "", exit(4, "nomad lease %s has an unsupported registration version", claim.LeaseID) + } + if _, _, err := registrationMetadataFingerprint(claim); err != nil { + return "", err + } + switch state { + case registrationPrepared, registrationSubmitting, registrationConfirmed: + return state, nil + default: + return "", exit(4, "nomad lease %s has an invalid registration state", claim.LeaseID) + } +} + +func unresolvedRegistrationError(claim LeaseClaim) error { + return exit(5, "nomad lease %s job %s registration outcome is unknown; claim retained", claim.LeaseID, claim.Labels[claimLabelJobID]) +} + +func registrationRecovery(ctx context.Context, claim LeaseClaim) (*shared.DelegatedSandboxRecovery, error) { + if claim.LeaseID == "" || claim.Provider != providerName || claim.Labels[claimLabelJobID] == "" { + return nil, exit(4, "nomad registration recovery has no bound claim") + } + if _, err := registrationState(claim); err != nil { + return nil, err + } + if err := core.WithLeaseClaimUnchangedContext(ctx, claim.LeaseID, claim, func() error { return nil }); err != nil { + return nil, err + } + return &shared.DelegatedSandboxRecovery{ + LeaseID: claim.LeaseID, Slug: claim.Slug, + CleanupCommand: "crabbox stop --provider nomad " + shellQuote(claim.LeaseID), + }, nil +} + +func registrationRecoveryError(claim LeaseClaim, cause error) error { + message := fmt.Sprintf("nomad registration lease=%s slug=%s job=%s scope=%q: %v; inspect with crabbox status --provider nomad --id %s and recover with crabbox stop --provider nomad %s", claim.LeaseID, claim.Slug, claim.Labels[claimLabelJobID], claim.ProviderScope, cause, claim.LeaseID, claim.LeaseID) + return shared.ExitErrorWithCause(core.ExitCodeForError(cause, 1), message, cause) +} + +func (b *backend) prepareRegistration(ctx context.Context, repo Repo, requestedSlug string) (*preparedRegistration, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + leaseID, err := newLeaseID() + if err != nil { + return nil, err + } + slug, err := allocateClaimLeaseSlug(leaseID, requestedSlug) + if err != nil { + return nil, err + } + expiresAt := time.Time{} + if b.cfg.TTL > 0 { + expiresAt = core.ClockNow(b.rt.Clock).UTC().Add(b.cfg.TTL) + } + jobID := jobIDForLease(leaseID) + job, err := buildJobSpec(b.cfg, jobSpecInput{LeaseID: leaseID, Slug: slug, JobID: jobID, ExpiresAt: expiresAt}) + if err != nil { + return nil, err + } + metadata, err := json.Marshal(job.Meta) + if err != nil { + return nil, err + } + labels := claimLabels(b.cfg, leaseID, slug, allocationReadiness{JobID: jobID}, expiresAt) + labels[registrationVersionLabel] = "1" + labels[registrationStateLabel] = registrationPrepared + keys := make([]string, 0, len(job.Meta)) + for key := range job.Meta { + keys = append(keys, key) + } + sort.Strings(keys) + keyJSON, err := json.Marshal(keys) + if err != nil { + return nil, err + } + digest := sha256.Sum256(metadata) + labels[registrationMetaKeysLabel] = string(keyJSON) + labels[registrationMetaHashLabel] = hex.EncodeToString(digest[:]) + if err := ctx.Err(); err != nil { + return nil, err + } + claim, err := core.ClaimLeaseForRepoProviderScopePondWithLabels(leaseID, slug, providerName, claimScope(b.cfg), b.cfg.Pond, repo.Root, b.cfg.IdleTimeout, labels) + if err != nil { + // Admission may be visible after a durability error; never submit after it. + return nil, registrationRecoveryError(LeaseClaim{LeaseID: leaseID, Slug: slug, ProviderScope: claimScope(b.cfg), Labels: labels}, err) + } + return &preparedRegistration{job: job, claim: claim}, nil +} + +func advanceRegistration(ctx context.Context, claim *LeaseClaim, state, evaluation string, ready *allocationReadiness) error { + previous, err := registrationState(*claim) + if err != nil { + return err + } + if claim.Labels[registrationVersionLabel] != "1" || + (previous == registrationPrepared && state != registrationSubmitting) || + (previous != registrationPrepared && state != registrationConfirmed) { + return exit(4, "nomad lease %s has an invalid registration transition", claim.LeaseID) + } + next := *claim + next.Labels = maps.Clone(claim.Labels) + next.Labels[registrationStateLabel] = state + if evaluation != "" { + next.Labels[registrationEvalLabel] = evaluation + } + if ready != nil { + applyReadinessLabels(next.Labels, *ready) + next.Labels[claimLabelState] = ready.State() + // Match the old ready-claim publication point, not preparation time. + next.ClaimedAt = time.Now().UTC().Format(time.RFC3339) + next.LastUsedAt = next.ClaimedAt + } + updated, err := core.ReplaceLeaseClaimIfUnchangedDurableReturningContext(ctx, claim.LeaseID, *claim, next) + if err != nil { + return err + } + *claim = updated + return nil +} + +func (b *backend) submitRegistration(ctx context.Context, client Client, prepared *preparedRegistration) (allocationReadiness, error) { + state, err := registrationState(prepared.claim) + if err != nil { + return allocationReadiness{}, err + } + if state != registrationPrepared { + return allocationReadiness{}, exit(4, "nomad lease %s registration was already submitted", prepared.claim.LeaseID) + } + if err := ctx.Err(); err != nil { + return allocationReadiness{}, err + } + if err := advanceRegistration(ctx, &prepared.claim, registrationSubmitting, "", nil); err != nil { + return allocationReadiness{}, err + } + if err := ctx.Err(); err != nil { + return allocationReadiness{}, err + } + evalID, err := client.RegisterJob(ctx, prepared.job) + if err != nil { + return allocationReadiness{}, err + } + // Acknowledged submission is distinct from allocation readiness. + recordCtx, cancel := b.cleanupContext(ctx) + err = advanceRegistration(recordCtx, &prepared.claim, registrationConfirmed, evalID, nil) + cancel() + if err != nil { + return allocationReadiness{}, err + } + if evalID != "" { + if err := b.waitForEvaluation(ctx, client, evalID); err != nil { + return allocationReadiness{}, err + } + } + ready, err := b.waitForAllocation(ctx, client, prepared.claim.Labels[claimLabelJobID], b.allocReadyTimeout()) + if err != nil { + return allocationReadiness{}, err + } + if err := advanceRegistration(ctx, &prepared.claim, registrationConfirmed, "", &ready); err != nil { + return allocationReadiness{}, err + } + return ready, nil +} + +func validateRegistrationMetadata(claim LeaseClaim, job *nomadapi.Job) error { + if !registrationAwaitingReadiness(claim) { + return nil + } + keys, expectedHash, err := registrationMetadataFingerprint(claim) + if err != nil { + return err + } + if job == nil || stringValue(job.ID) != claim.Labels[claimLabelJobID] { + return exit(4, "nomad job for lease %s registration ownership changed", claim.LeaseID) + } + selected := make(map[string]string, len(keys)) + for _, key := range keys { + // Match metadataMatches: added keys are ignored; missing empty values match. + selected[key] = strings.TrimSpace(job.Meta[key]) + } + data, err := json.Marshal(selected) + if err != nil { + return err + } + digest := sha256.Sum256(data) + if hex.EncodeToString(digest[:]) != expectedHash { + return exit(4, "nomad job for lease %s registration ownership changed", claim.LeaseID) + } + return nil +} + +// Allocation ID is committed only after successful readiness. Normal ready and +// legacy claims retain their existing reserved-ownership metadata contract. +func registrationAwaitingReadiness(claim LeaseClaim) bool { + return claim.Labels[registrationVersionLabel] != "" && claim.Labels[claimLabelAllocationID] == "" +} + +func registrationMetadataFingerprint(claim LeaseClaim) ([]string, string, error) { + var keys []string + hash := claim.Labels[registrationMetaHashLabel] + decoded, err := hex.DecodeString(hash) + if err != nil || len(decoded) != sha256.Size || json.Unmarshal([]byte(claim.Labels[registrationMetaKeysLabel]), &keys) != nil || len(keys) == 0 || !sort.StringsAreSorted(keys) { + return nil, "", exit(4, "nomad lease %s has invalid registration metadata fingerprint", claim.LeaseID) + } + for i := 1; i < len(keys); i++ { + if keys[i] == keys[i-1] { + return nil, "", exit(4, "nomad lease %s has duplicate registration metadata keys", claim.LeaseID) + } + } + return keys, hash, nil +} diff --git a/internal/providers/nomad/registration_test.go b/internal/providers/nomad/registration_test.go new file mode 100644 index 000000000..de7092d11 --- /dev/null +++ b/internal/providers/nomad/registration_test.go @@ -0,0 +1,403 @@ +package nomad + +import ( + "context" + "encoding/json" + "errors" + "maps" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "slices" + "strings" + "sync" + "testing" + "time" + + nomadapi "github.com/hashicorp/nomad/api" + core "github.com/openclaw/crabbox/internal/cli" +) + +type registrationHookClient struct { + Client + register func(context.Context, *nomadapi.Job) (string, error) + evaluation func(context.Context, string) (*nomadapi.Evaluation, error) +} + +func (c registrationHookClient) RegisterJob(ctx context.Context, job *nomadapi.Job) (string, error) { + return c.register(ctx, job) +} + +func (c registrationHookClient) EvaluationInfo(ctx context.Context, id string) (*nomadapi.Evaluation, error) { + if c.evaluation != nil { + return c.evaluation(ctx, id) + } + return c.Client.EvaluationInfo(ctx, id) +} + +func onlyRegistrationClaim(t *testing.T) LeaseClaim { + t.Helper() + claims, err := listNomadLeaseClaims() + if err != nil || len(claims) != 1 { + t.Fatalf("claims=%d err=%v, want one retained registration", len(claims), err) + } + return claims[0] +} + +func TestNomadRegistrationUncertaintyRetainsRecovery(t *testing.T) { + for _, primary := range []struct { + name string + cause error + code int + }{ + {"typed-four", exit(4, "registration response unavailable"), 4}, + {"typed-five", exit(5, "registration response unavailable"), 5}, + {"generic", errors.New("registration response unavailable"), 1}, + {"canceled", context.Canceled, 1}, + {"deadline", context.DeadlineExceeded, 1}, + } { + for _, mode := range []string{"warmup", "run", "run-keep"} { + t.Run(primary.name+"/"+mode, func(t *testing.T) { + fake := newLifecycleFakeClient() + b, _, stderr := testBackend(t, fake) + cause := primary.cause + calls := 0 + client := registrationHookClient{Client: fake, register: func(_ context.Context, job *nomadapi.Job) (string, error) { + calls++ + claim := onlyRegistrationClaim(t) + if state, err := registrationState(claim); err != nil || state != registrationSubmitting || claim.Labels[claimLabelJobID] != stringValue(job.ID) { + t.Fatalf("registration was not durably submitting: state=%s err=%v", state, err) + } + if claim.Labels[claimLabelAllocationID] != "" { + t.Fatal("pending registration fabricated an allocation") + } + return "", cause + }} + b.clientFactory = func(Config, Runtime) (Client, error) { return client, nil } + repo := Repo{Root: t.TempDir(), Name: "ordinary-repo"} + var err error + if mode == "warmup" { + err = b.Warmup(context.Background(), WarmupRequest{Repo: repo, Keep: true}) + } else { + var result RunResult + result, err = b.Run(context.Background(), RunRequest{Repo: repo, Command: []string{"true"}, NoSync: true, Keep: mode == "run-keep", TimingJSON: true}) + claim := onlyRegistrationClaim(t) + if result.Session == nil || !result.Session.Kept || result.Session.Reused || result.LeaseID != claim.LeaseID || result.Session.CleanupCommand == "" { + t.Fatalf("missing recovery session: %#v", result.Session) + } + if !strings.Contains(stderr.String(), claim.LeaseID) { + t.Fatal("timing/reporting omitted retained lease") + } + } + claim := onlyRegistrationClaim(t) + var displayed core.ExitError + if !errors.Is(err, cause) || !core.AsExitError(err, &displayed) || displayed.Code != primary.code { + t.Fatalf("primary error/code lost: code=%d want=%d: %v", displayed.Code, primary.code, err) + } + for _, value := range []string{claim.LeaseID, claim.Labels[claimLabelJobID], claim.ProviderScope, "crabbox stop --provider nomad"} { + if !strings.Contains(displayed.Message, value) { + t.Fatalf("CLI-selected error omitted %q: %s", value, displayed.Message) + } + } + view, err := b.Status(context.Background(), StatusRequest{ID: claim.LeaseID}) + if err != nil || view.State != "registration-pending" || view.Ready { + t.Fatalf("pending status=%#v err=%v", view, err) + } + for _, dryRun := range []bool{true, false} { + if err := b.Cleanup(context.Background(), CleanupRequest{DryRun: dryRun}); err == nil || !strings.Contains(err.Error(), "registration outcome is unknown; claim retained") { + t.Fatalf("cleanup dryRun=%t erased uncertainty: %v", dryRun, err) + } + } + if err := b.Stop(context.Background(), StopRequest{ID: claim.LeaseID}); err == nil { + t.Fatal("pending absence treated as completed stop") + } + assertNomadClaimRetained(t, claim) + if calls != 1 || len(fake.deregisters) != 0 || len(fake.execs) != 0 { + t.Fatalf("unexpected retry/cleanup/exec: register=%d purge=%v exec=%v", calls, fake.deregisters, fake.execs) + } + }) + } + } +} + +func TestNomadAcceptedRegistrationHTTPDeadlineRollsBack(t *testing.T) { + old := nomadControlRequestTimeout + nomadControlRequestTimeout = 200 * time.Millisecond + t.Cleanup(func() { nomadControlRequestTimeout = old }) + t.Setenv("NOMAD_TOKEN", "") + b, _, _ := testBackend(t, nil) + var mu sync.Mutex + var stored *nomadapi.Job + var submittedClaim LeaseClaim + var requests []string + registrationCanceled := make(chan struct{}) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + if r.Method == http.MethodPut && r.URL.Path == "/v1/jobs" { + claims, err := listNomadLeaseClaims() + if err != nil || len(claims) != 1 { + t.Errorf("registration reached HTTP without one durable claim: claims=%d err=%v", len(claims), err) + http.Error(w, "missing registration claim", http.StatusInternalServerError) + return + } + claim := claims[0] + if state, err := registrationState(claim); err != nil || state != registrationSubmitting { + t.Errorf("registration reached HTTP before submitting: state=%s err=%v", state, err) + } + var body struct{ Job *nomadapi.Job } + if err := json.NewDecoder(r.Body).Decode(&body); err != nil || body.Job == nil { + t.Errorf("decode submitted job: job=%v err=%v", body.Job, err) + http.Error(w, "invalid job", http.StatusBadRequest) + return + } + if stringValue(body.Job.ID) != claim.Labels[claimLabelJobID] { + t.Errorf("submitted job does not match durable claim") + } + mu.Lock() + stored, submittedClaim = body.Job, claim + requests = append(requests, "PUT") + mu.Unlock() + // Accept the unchanged job, but withhold the rest of the JSON acknowledgement. + _, _ = w.Write([]byte(`{"EvalID":"`)) + w.(http.Flusher).Flush() + <-r.Context().Done() + close(registrationCanceled) + return + } + mu.Lock() + defer mu.Unlock() + requests = append(requests, r.Method) + if submittedClaim.LeaseID == "" || r.URL.Path != "/v1/job/"+submittedClaim.Labels[claimLabelJobID] { + t.Errorf("unexpected request: %s %s", r.Method, r.URL.Path) + http.NotFound(w, r) + return + } + switch r.Method { + case http.MethodGet: + if stored == nil { + http.NotFound(w, r) + return + } + if err := json.NewEncoder(w).Encode(stored); err != nil { + t.Errorf("write stored job: %v", err) + } + case http.MethodDelete: + if stored == nil || r.URL.Query().Get("purge") != "true" { + t.Errorf("cleanup did not purge the accepted job") + } + stored = nil + _, _ = w.Write([]byte(`{}`)) + default: + t.Errorf("unexpected request method: %s", r.Method) + http.Error(w, "unexpected method", http.StatusMethodNotAllowed) + } + })) + t.Cleanup(func() { + server.CloseClientConnections() + server.Close() + }) + b.cfg.Nomad.Address = server.URL + b.clientFactory = newNomadClient + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + err := b.Warmup(ctx, WarmupRequest{Repo: Repo{Root: t.TempDir(), Name: "ordinary-repo"}, Keep: true}) + var displayed core.ExitError + if ctx.Err() != nil || !errors.Is(err, context.DeadlineExceeded) || !core.AsExitError(err, &displayed) || displayed.Code != 1 { + t.Fatalf("initiating request deadline/code lost: caller=%v code=%d err=%v", ctx.Err(), displayed.Code, err) + } + select { + case <-registrationCanceled: + case <-ctx.Done(): + t.Fatal("registration HTTP request was not canceled") + } + mu.Lock() + defer mu.Unlock() + if stored != nil || !slices.Equal(requests, []string{"PUT", "GET", "DELETE", "GET"}) { + t.Fatalf("accepted registration rollback sequence: stored=%v requests=%v", stored != nil, requests) + } + if !strings.Contains(displayed.Message, submittedClaim.LeaseID) || !strings.Contains(displayed.Message, submittedClaim.Labels[claimLabelJobID]) || !strings.Contains(displayed.Message, "rolled back") || strings.Contains(displayed.Message, "recover with") { + t.Fatalf("rollback diagnostic lost identity or implied retention: %s", displayed.Message) + } + if claims, err := listNomadLeaseClaims(); err != nil || len(claims) != 0 { + t.Fatalf("claim remains after confirmed HTTP rollback: claims=%d err=%v", len(claims), err) + } +} + +func TestNomadRegistrationAcknowledgedBeforeReadiness(t *testing.T) { + for _, keep := range []bool{false, true} { + t.Run(map[bool]string{false: "temporary", true: "keep"}[keep], func(t *testing.T) { + fake := newLifecycleFakeClient() + b, _, _ := testBackend(t, fake) + confirmed := false + client := registrationHookClient{Client: fake, register: fake.RegisterJob, evaluation: func(ctx context.Context, id string) (*nomadapi.Evaluation, error) { + if !strings.HasPrefix(id, "eval-deregister-") { + claim := onlyRegistrationClaim(t) + state, err := registrationState(claim) + if err != nil || state != registrationConfirmed || claim.Labels[registrationEvalLabel] != id || claim.Labels[claimLabelAllocationID] != "" { + t.Fatalf("acknowledgement not persisted before wait: state=%s err=%v", state, err) + } + confirmed = true + return &nomadapi.Evaluation{ID: id, Status: nomadapi.EvalStatusFailed}, nil + } + return fake.EvaluationInfo(ctx, id) + }} + b.clientFactory = func(Config, Runtime) (Client, error) { return client, nil } + result, err := b.Run(context.Background(), RunRequest{Repo: Repo{Root: t.TempDir(), Name: "repo"}, Command: []string{"true"}, NoSync: true, Keep: keep}) + if err == nil || !confirmed || len(fake.deregisters) != 1 || result.Session != nil { + t.Fatalf("readiness rollback changed Keep contract: err=%v confirmed=%t purges=%v session=%#v", err, confirmed, fake.deregisters, result.Session) + } + var displayed core.ExitError + if !core.AsExitError(err, &displayed) || !strings.Contains(displayed.Message, "rolled back") || !strings.Contains(displayed.Message, "lease=") || !strings.Contains(displayed.Message, "job=") || strings.Contains(displayed.Message, "recover with") { + t.Fatalf("successful rollback lost identity or implied retention: %v", err) + } + claims, err := listNomadLeaseClaims() + if err != nil || len(claims) != 0 { + t.Fatalf("claim remains after rollback: %d %v", len(claims), err) + } + }) + } +} + +func TestNomadPreparedRegistrationNeverSubmitsAfterRemoval(t *testing.T) { + fake := newLifecycleFakeClient() + b, _, _ := testBackend(t, fake) + prepared, err := b.prepareRegistration(context.Background(), Repo{Root: t.TempDir()}, "prepared") + if err != nil { + t.Fatal(err) + } + lookup := cleanupHookClient{Client: fake, jobInfo: func(context.Context, string) (*nomadapi.Job, error) { + t.Fatal("prepared cleanup must not query remote jobs") + return nil, nil + }} + if _, err := b.removeOwnedJob(context.Background(), lookup, prepared.claim, false); err != nil { + t.Fatal(err) + } + if _, err := b.submitRegistration(context.Background(), fake, prepared); err == nil || fake.registers != 0 { + t.Fatalf("removed attempt submitted: err=%v calls=%d", err, fake.registers) + } + if current, err := readLeaseClaim(prepared.claim.LeaseID); err != nil || current.LeaseID != "" { + t.Fatalf("removed claim recreated: %#v %v", current, err) + } +} + +func TestNomadRegistrationAdmissionFailureDoesNotDispatch(t *testing.T) { + fake := newLifecycleFakeClient() + b, _, _ := testBackend(t, fake) + notDirectory := filepath.Join(t.TempDir(), "file") + if err := os.WriteFile(notDirectory, []byte("ordinary fixture"), 0o600); err != nil { + t.Fatal(err) + } + t.Setenv("XDG_STATE_HOME", notDirectory) + _, _, recovery, err := b.createJob(context.Background(), fake, Repo{Root: t.TempDir()}, "no-dispatch") + if err == nil || recovery != nil || fake.registers != 0 { + t.Fatalf("admission failure dispatched: err=%v recovery=%#v calls=%d", err, recovery, fake.registers) + } +} + +func TestNomadObservedRegistrationErrorRollsBack(t *testing.T) { + fake := newLifecycleFakeClient() + b, _, _ := testBackend(t, fake) + cause := exit(5, "registration response unavailable") + client := registrationHookClient{Client: fake, register: func(ctx context.Context, job *nomadapi.Job) (string, error) { + _, err := fake.RegisterJob(ctx, job) + if err != nil { + t.Fatal(err) + } + return "", cause + }} + _, claim, recovery, err := b.createJob(context.Background(), client, Repo{Root: t.TempDir()}, "observed") + if !errors.Is(err, cause) || recovery != nil || claim.LeaseID != "" || fake.registers != 1 || len(fake.deregisters) != 1 { + t.Fatalf("observed registration rollback: claim=%s recovery=%#v calls=%d purge=%v err=%v", claim.LeaseID, recovery, fake.registers, fake.deregisters, err) + } +} + +func TestNomadAcknowledgementRecordedAfterCallerCancel(t *testing.T) { + fake := newLifecycleFakeClient() + b, _, _ := testBackend(t, fake) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + observedConfirmed := false + reads := cleanupHookClient{Client: fake, jobInfo: func(ctx context.Context, id string) (*nomadapi.Job, error) { + claim := onlyRegistrationClaim(t) + if state, err := registrationState(claim); err != nil || state != registrationConfirmed { + t.Fatalf("successful acknowledgement remained unknown: state=%s err=%v", state, err) + } + observedConfirmed = true + return fake.JobInfo(ctx, id) + }} + client := registrationHookClient{Client: reads, register: func(ctx context.Context, job *nomadapi.Job) (string, error) { + id, err := fake.RegisterJob(ctx, job) + cancel() + return id, err + }} + _, _, recovery, err := b.createJob(ctx, client, Repo{Root: t.TempDir()}, "canceled-after-ack") + if !errors.Is(err, context.Canceled) || recovery != nil || !observedConfirmed || len(fake.deregisters) != 1 { + t.Fatalf("ack/cancel recovery=%#v confirmed=%t purge=%v err=%v", recovery, observedConfirmed, fake.deregisters, err) + } +} + +func TestNomadRecoveryVerificationHonorsCancellation(t *testing.T) { + fake := newLifecycleFakeClient() + b, _, _ := testBackend(t, fake) + prepared, err := b.prepareRegistration(context.Background(), Repo{Root: t.TempDir()}, "verify-canceled") + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if recovery, err := registrationRecovery(ctx, prepared.claim); recovery != nil || !errors.Is(err, context.Canceled) { + t.Fatalf("canceled verification recovery=%#v err=%v", recovery, err) + } + assertNomadClaimRetained(t, prepared.claim) +} + +func TestNomadRegistrationMetadataAndLegacyBoundary(t *testing.T) { + fake := newLifecycleFakeClient() + b, _, _ := testBackend(t, fake) + claim := createClaim(t, b, "cbx_a77777777777", "metadata", "crabbox-a77777777777", "alloc-a") + job := fake.jobs[claim.Labels[claimLabelJobID]] + job.Meta["custom"] = "original-value" + job.Meta["empty"] = "" + claim = markRegistrationClaim(t, claim, job, registrationSubmitting) + if strings.Contains(claim.Labels[registrationMetaKeysLabel]+claim.Labels[registrationMetaHashLabel], "original-value") { + t.Fatal("arbitrary metadata value persisted") + } + job.Meta["server-added"] = "allowed" + delete(job.Meta, "empty") + if err := validateRegistrationMetadata(claim, job); err != nil { + t.Fatalf("additional metadata or missing expected-empty changed matching: %v", err) + } + job.Meta["custom"] = "changed" + if err := validateRegistrationMetadata(claim, job); err == nil { + t.Fatal("pending original metadata change accepted") + } + ready := claim + ready.Labels = maps.Clone(claim.Labels) + ready.Labels[registrationStateLabel] = registrationConfirmed + ready.Labels[claimLabelAllocationID] = "alloc-a" + if err := validateRemoteOwnership(b.cfg, ready, job); err != nil { + t.Fatalf("normal ready ownership semantics strengthened: %v", err) + } + if err := validateRegistrationMetadata(ready, job); err != nil { + t.Fatalf("ready claim retained setup-only full map guard: %v", err) + } + for _, labels := range []map[string]string{{registrationVersionLabel: "2"}, {registrationStateLabel: registrationSubmitting}, {registrationMetaHashLabel: "unknown"}} { + if _, err := registrationState(LeaseClaim{Labels: labels}); err == nil { + t.Fatal("unknown or incomplete registration schema accepted") + } + } + if state, err := registrationState(LeaseClaim{}); err != nil || state != registrationConfirmed { + t.Fatalf("legacy claim semantics changed: %s %v", state, err) + } + job.Meta["custom"] = "original-value" + if err := b.Stop(context.Background(), StopRequest{ID: claim.LeaseID}); err != nil { + t.Fatalf("later matching observation did not reconcile: %v", err) + } + if len(fake.deregisters) != 1 { + t.Fatalf("matching observed job purges=%v", fake.deregisters) + } + if _, err := fake.JobInfo(context.Background(), claim.Labels[claimLabelJobID]); !isNotFoundError(err) { + t.Fatalf("confirmed removal not observed: %v, want HTTP %d", err, http.StatusNotFound) + } +} diff --git a/internal/providers/shared/delegated.go b/internal/providers/shared/delegated.go index 39a424d1f..a8180a2a4 100644 --- a/internal/providers/shared/delegated.go +++ b/internal/providers/shared/delegated.go @@ -15,11 +15,22 @@ import ( // the resource. Unlock, when set, is held through cleanup and final reporting. // On acquisition failure the adapter owns partial-resource rollback; Unlock is // still called, but no session or permission to delete is inferred from an ID. +// Recovery explicitly reports a durable recovery claim after failed acquisition. type DelegatedSandbox struct { LeaseID string Slug string CleanupCommand string Unlock func() + Recovery *DelegatedSandboxRecovery +} + +// DelegatedSandboxRecovery reports an authorized durable claim left after the +// adapter's acquisition rollback. It is used only on Acquire error and grants +// no permission for shared cleanup, retention, execution, or readiness. +type DelegatedSandboxRecovery struct { + LeaseID string + Slug string + CleanupCommand string } // DelegatedSandboxCommand keeps command preparation (including credential @@ -151,6 +162,7 @@ func RunDelegatedSandbox(ctx context.Context, req core.RunRequest, lifecycle Del syncPhases = []core.TimingPhase{{Name: "sync", Skipped: true, Reason: "--no-sync"}} } acquired := req.ID == "" + acquireFailed := false reuseAdmitted := acquired || lifecycle.AdmitReuse == nil commandRan := false @@ -171,7 +183,7 @@ func RunDelegatedSandbox(ctx context.Context, req core.RunRequest, lifecycle Del cancel() appendFailure(closeErr, core.ExitCodeForError(closeErr, 1)) } - if result.Session != nil { + if result.Session != nil && !acquireFailed { shouldStop := acquired && !req.Keep if retErr != nil && reuseAdmitted { core.HandleDelegatedRunFailure(stderr, req, lifecycle.Provider, sandbox.LeaseID, sandbox.Slug, lifecycle.IdleTimeout, lifecycle.TTL, acquired, &shouldStop) @@ -238,6 +250,14 @@ func RunDelegatedSandbox(ctx context.Context, req core.RunRequest, lifecycle Del sandbox, err = lifecycle.Resolve(ctx) } if err != nil { + acquireFailed = acquired + if recovery := sandbox.Recovery; acquired && recovery != nil && recovery.LeaseID != "" { + result.LeaseID, result.Slug = recovery.LeaseID, recovery.Slug + result.Session = &core.RunSessionHandle{ + Provider: lifecycle.Provider, LeaseID: recovery.LeaseID, Slug: recovery.Slug, + Kept: true, CleanupCommand: recovery.CleanupCommand, + } + } return result, err } result.LeaseID, result.Slug = sandbox.LeaseID, sandbox.Slug diff --git a/internal/providers/shared/delegated_test.go b/internal/providers/shared/delegated_test.go index 2ae92b3be..4260bc7e7 100644 --- a/internal/providers/shared/delegated_test.go +++ b/internal/providers/shared/delegated_test.go @@ -492,6 +492,131 @@ func TestDelegatedSandboxCleanupDeadline(t *testing.T) { type sandboxFailingTimingWriter struct{} +func TestDelegatedSandboxAcquisitionRecovery(t *testing.T) { + failure := errors.New("acquisition unresolved") + for _, tc := range []struct { + name string + err error + keep, keepOnFailure, reuse, noCarrier, emptyCarrier bool + wantCode int + wantStatus core.RunStatus + wantKind core.RunErrorKind + }{ + {name: "default", err: failure, wantCode: 1, wantStatus: core.RunStatusFailed, wantKind: core.RunErrorProvider}, + {name: "keep", err: failure, keep: true, wantCode: 1, wantStatus: core.RunStatusFailed, wantKind: core.RunErrorProvider}, + {name: "keep on failure", err: failure, keepOnFailure: true, wantCode: 1, wantStatus: core.RunStatusFailed, wantKind: core.RunErrorProvider}, + {name: "typed primary", err: core.ExitError{Code: 4, Message: "pending claim"}, wantCode: 4, wantStatus: core.RunStatusFailed, wantKind: core.RunErrorProvider}, + {name: "canceled", err: context.Canceled, wantCode: 1, wantStatus: core.RunStatusCanceled, wantKind: core.RunErrorCanceled}, + {name: "deadline", err: context.DeadlineExceeded, wantCode: 1, wantStatus: core.RunStatusTimedOut, wantKind: core.RunErrorTimeout}, + {name: "ID alone", err: failure, noCarrier: true, wantCode: 1, wantStatus: core.RunStatusFailed, wantKind: core.RunErrorProvider}, + {name: "empty carrier", err: failure, emptyCarrier: true, wantCode: 1, wantStatus: core.RunStatusFailed, wantKind: core.RunErrorProvider}, + {name: "resolve ignores carrier", err: failure, reuse: true, wantCode: 1, wantStatus: core.RunStatusFailed, wantKind: core.RunErrorProvider}, + {name: "success ignores carrier", wantStatus: core.RunStatusSucceeded}, + } { + t.Run(tc.name, func(t *testing.T) { + var stderr bytes.Buffer + var calls []string + clock := &sandboxTestClock{current: time.Unix(0, 0)} + acquire := func(context.Context) (DelegatedSandbox, error) { + calls = append(calls, "acquire") + clock.Sleep(3 * time.Millisecond) + s := DelegatedSandbox{LeaseID: "ordinary", Slug: "ordinary-slug", CleanupCommand: "stop ordinary", Unlock: func() { calls = append(calls, "unlock") }} + if !tc.noCarrier { + s.Recovery = &DelegatedSandboxRecovery{LeaseID: "pending", Slug: "pending-slug", CleanupCommand: "stop pending"} + if tc.emptyCarrier { + s.Recovery.LeaseID = "" + } + } + return s, tc.err + } + step := func(name string) func(context.Context) error { + return func(context.Context) error { calls = append(calls, name); return nil } + } + req := core.RunRequest{NoSync: true, TimingJSON: true, Keep: tc.keep, KeepOnFailure: tc.keepOnFailure} + if tc.reuse { + req.ID = "ordinary" + } + result, err := RunDelegatedSandbox(t.Context(), req, DelegatedSandboxLifecycle{ + Provider: "fixture", Runtime: core.Runtime{Stderr: &stderr, Clock: clock}, + Acquire: acquire, Resolve: acquire, Setup: step("setup"), NoSync: step("workspace"), + Command: func(context.Context) (DelegatedSandboxCommand, error) { + calls = append(calls, "command") + return DelegatedSandboxCommand{Run: func(context.Context) (int, error) { calls = append(calls, "run"); return 0, nil }}, nil + }, + Cleanup: step("cleanup"), Retained: step("retained"), + }) + if result.ExitCode != tc.wantCode || result.Status != tc.wantStatus || result.ErrorKind != tc.wantKind || !errors.Is(err, tc.err) { + t.Fatalf("result=%#v err=%v", result, err) + } + wantCalls := []string{"acquire", "unlock"} + wantLease := "" + wantRecovery := tc.err != nil && !tc.reuse && !tc.noCarrier && !tc.emptyCarrier + if wantRecovery { + wantLease = "pending" + want := &core.RunSessionHandle{Provider: "fixture", LeaseID: "pending", Slug: "pending-slug", CleanupCommand: "stop pending", Kept: true} + if !reflect.DeepEqual(result.Session, want) || result.Slug != "pending-slug" { + t.Fatalf("recovery session=%#v result=%#v", result.Session, result) + } + } else if tc.err == nil { + wantLease = "ordinary" + wantCalls = []string{"acquire", "setup", "workspace", "command", "run", "cleanup", "unlock"} + if result.Session == nil || result.Session.LeaseID != "ordinary" || result.Session.Kept { + t.Fatalf("success session=%#v", result.Session) + } + } else if result.Session != nil { + t.Fatalf("unadmitted session=%#v", result.Session) + } + if !reflect.DeepEqual(calls, wantCalls) || result.LeaseID != wantLease { + t.Fatalf("calls=%v result=%#v", calls, result) + } + var report core.TimingReport + reports := 0 + for _, line := range strings.Split(strings.TrimSpace(stderr.String()), "\n") { + if strings.HasPrefix(line, "{") { + reports++ + if err := json.Unmarshal([]byte(line), &report); err != nil { + t.Fatal(err) + } + } else if tc.err != nil { + t.Fatalf("acquisition error ran normal reporting policy: %q", line) + } + } + if reports != 1 || report.LeaseID != wantLease || report.ExitCode != tc.wantCode || report.RunStatus != tc.wantStatus || report.ErrorKind != tc.wantKind || report.TotalMs != 3 || report.CommandMs != 0 { + t.Fatalf("timing=%#v", report) + } + }) + } +} + +func TestDelegatedSandboxAcquisitionRecoveryTimingFailureReleasesResources(t *testing.T) { + file, err := os.CreateTemp(t.TempDir(), "prepared-") + if err != nil { + t.Fatal(err) + } + primary := core.ExitError{Code: 4, Message: "registration unresolved"} + unlocks := 0 + result, err := RunDelegatedSandbox(t.Context(), core.RunRequest{TimingJSON: true, Keep: true}, DelegatedSandboxLifecycle{ + Provider: "fixture", Runtime: core.Runtime{Stderr: sandboxFailingTimingWriter{}}, + PrepareArchive: func(context.Context) (*core.PreparedArchive, error) { return &core.PreparedArchive{File: file}, nil }, + Acquire: func(context.Context) (DelegatedSandbox, error) { + return DelegatedSandbox{Recovery: &DelegatedSandboxRecovery{LeaseID: "pending"}, Unlock: func() { + unlocks++ + if _, err := os.Stat(file.Name()); !os.IsNotExist(err) { + t.Errorf("archive not closed before unlock: %v", err) + } + }}, primary + }, + Cleanup: func(context.Context) error { t.Fatal("shared acquisition cleanup"); return nil }, + Retained: func(context.Context) error { t.Fatal("shared acquisition retention"); return nil }, + }) + if !errors.Is(err, primary) || !errors.Is(err, io.ErrClosedPipe) || result.ExitCode != 4 || unlocks != 1 || result.Session == nil || !result.Session.Kept { + t.Fatalf("result=%#v err=%v unlocks=%d", result, err, unlocks) + } + if _, err := file.Stat(); !errors.Is(err, os.ErrClosed) { + t.Fatalf("archive descriptor remains open: %v", err) + } +} + func (sandboxFailingTimingWriter) Write(p []byte) (int, error) { return len(p), nil } func (sandboxFailingTimingWriter) WriteTimingReport(core.TimingReport) error { return io.ErrClosedPipe }