From 88b5f89952c562a13272a96e342d09b595cd5f06 Mon Sep 17 00:00:00 2001 From: dawn Date: Wed, 2 Sep 2026 21:51:09 +0900 Subject: [PATCH] spindle/engines/microvm: fix nix cache upload quota publication semantics Transition the reservation to publishing before anything can reach the destination: recovery charges publishing rows and drops reserved ones, so publishing after a confirmed upload leaves crash-window bytes unaccounted. A failed publishing transition now releases (or re-queues) the reservation and clears it from the staged nar instead of leaking a reserved row, and the narinfo retry re-reserves. A denied late reservation skips the object and answers 200 like the nar upload path does instead of failing the build with a 503. Publish failures that may have already landed keep the charge: transport uncertainty, 5xx from the destination, and anything after a confirmed nar PUT. Only definite 4xx refusals release it. Both reservations key off one canonical nar path so two spellings cannot double-charge. Signed-off-by: dawn --- spindle/engines/microvm/upload_quota_test.go | 132 +++++++++++++++++-- spindle/engines/microvm/upload_staging.go | 88 ++++++++++--- 2 files changed, 189 insertions(+), 31 deletions(-) diff --git a/spindle/engines/microvm/upload_quota_test.go b/spindle/engines/microvm/upload_quota_test.go index a955980ed..0a3e08b35 100644 --- a/spindle/engines/microvm/upload_quota_test.go +++ b/spindle/engines/microvm/upload_quota_test.go @@ -336,16 +336,16 @@ func TestQuotaHTTPFailureRetention(t *testing.T) { wrapper.ServeHTTP(rec, req) store.mu.Lock() - if len(store.releaseCalls) != 0 { + if len(store.releaseCalls) != 1 || store.releaseCalls[0] != "res-abc" { store.mu.Unlock() - t.Fatalf("HTTP refusal released reservations: %v", store.releaseCalls) + t.Fatalf("HTTP refusal release calls: %v", store.releaseCalls) } store.mu.Unlock() sw := wrapper.(*stagingWrapper) sw.mu.Lock() defer sw.mu.Unlock() - if len(sw.pendingCommits) != 1 || sw.pendingCommits[0] != "res-abc" { - t.Errorf("pending commits = %v, want [res-abc]", sw.pendingCommits) + if len(sw.pendingCommits) != 0 { + t.Errorf("pending commits = %v, want none", sw.pendingCommits) } } @@ -376,10 +376,17 @@ func TestQuotaUncertainFailureRetention(t *testing.T) { wrapper.ServeHTTP(rec, req) store.mu.Lock() - defer store.mu.Unlock() if len(store.releaseCalls) != 0 { + store.mu.Unlock() t.Errorf("expected 0 release calls on uncertain failure, got %v", store.releaseCalls) } + store.mu.Unlock() + sw := wrapper.(*stagingWrapper) + sw.mu.Lock() + defer sw.mu.Unlock() + if len(sw.pendingCommits) != 1 || sw.pendingCommits[0] != "res-abc" { + t.Errorf("pending commits = %v, want [res-abc] so the possibly published bytes stay charged", sw.pendingCommits) + } } func TestQuotaIdentityPropagation(t *testing.T) { @@ -531,8 +538,8 @@ func TestQuotaNixStoreSecondStageFailure(t *testing.T) { store.mu.Lock() defer store.mu.Unlock() - if len(store.releaseCalls) != 0 { - t.Errorf("expected no release calls on nix import failure, got %v", store.releaseCalls) + if len(store.releaseCalls) != 1 || store.releaseCalls[0] != "res-abc" { + t.Errorf("expected release on nix import failure, got %v", store.releaseCalls) } } @@ -782,14 +789,24 @@ func TestQuotaBeginCommitFailure(t *testing.T) { req = httptest.NewRequest(http.MethodPut, validQuotaNarinfoPath, strings.NewReader(validQuotaNarinfo)) wrapper.ServeHTTP(rec, req) - if rec.Code != http.StatusInternalServerError { - t.Fatalf("expected 500, got %d", rec.Code) + if rec.Code != http.StatusServiceUnavailable { + t.Fatalf("expected 503, got %d", rec.Code) } store.mu.Lock() - defer store.mu.Unlock() if len(store.releaseCalls) != 1 || store.releaseCalls[0] != "res-abc" { - t.Fatalf("expected release call with res-abc, got %v", store.releaseCalls) + store.mu.Unlock() + t.Fatalf("failed publishing transition leaked the reservation: %v", store.releaseCalls) + } + store.mu.Unlock() + + sw := wrapper.(*stagingWrapper) + sw.stageMu.Lock() + defer sw.stageMu.Unlock() + for path, staged := range sw.stagedNARs { + if staged.reservationID != "" { + t.Fatalf("staged nar %s kept reservation %q after a failed publishing transition", path, staged.reservationID) + } } } @@ -797,10 +814,9 @@ func TestQuotaTransientReleaseFailure(t *testing.T) { staging := t.TempDir() store := newMockQuotaStore(true, "") store.releaseErr = errors.New("transient release failure") - store.beginCommitErr = errors.New("begin publish failed") srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - w.WriteHeader(http.StatusOK) + w.WriteHeader(http.StatusForbidden) })) defer srv.Close() @@ -869,6 +885,96 @@ func TestQuotaTransientReleaseFailure(t *testing.T) { } } +func TestQuotaServerErrorPublishKeepsCharge(t *testing.T) { + staging := t.TempDir() + store := newMockQuotaStore(true, "") + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusBadGateway) + })) + defer srv.Close() + + u, _ := url.Parse(srv.URL) + backend := newHTTPUploadProxyBackend(u, nil, slog.Default()) + wrapper := NewStagingWrapper(backend, staging, store, "owner-1", "repo-1", srv.URL, slog.Default(), nil) + + rec := httptest.NewRecorder() + wrapper.ServeHTTP(rec, httptest.NewRequest(http.MethodPut, "/nar/abc.nar.xz", strings.NewReader("data"))) + + rec = httptest.NewRecorder() + wrapper.ServeHTTP(rec, httptest.NewRequest(http.MethodPut, validQuotaNarinfoPath, strings.NewReader(validQuotaNarinfo))) + if rec.Code != http.StatusBadGateway { + t.Fatalf("narinfo status = %d, want 502", rec.Code) + } + + store.mu.Lock() + if len(store.releaseCalls) != 0 { + store.mu.Unlock() + t.Fatalf("expected no release after a 5xx publish failure, got %v", store.releaseCalls) + } + store.mu.Unlock() + + sw := wrapper.(*stagingWrapper) + sw.mu.Lock() + defer sw.mu.Unlock() + if len(sw.pendingCommits) != 1 || sw.pendingCommits[0] != "res-abc" { + t.Fatalf("pending commits = %v, want [res-abc]", sw.pendingCommits) + } +} + +func TestQuotaLateReservationDenialSkipsUpload(t *testing.T) { + staging := t.TempDir() + store := newMockQuotaStore(true, "") + store.beginCommitErr = errors.New("begin publish failed") + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusOK) + })) + defer srv.Close() + + u, _ := url.Parse(srv.URL) + backend := newHTTPUploadProxyBackend(u, nil, slog.Default()) + wrapper := NewStagingWrapper(backend, staging, store, "owner-1", "repo-1", srv.URL, slog.Default(), nil) + + rec := httptest.NewRecorder() + wrapper.ServeHTTP(rec, httptest.NewRequest(http.MethodPut, "/nar/abc.nar.xz", strings.NewReader("data"))) + if rec.Code != http.StatusOK { + t.Fatalf("nar upload status = %d, want 200", rec.Code) + } + + rec = httptest.NewRecorder() + wrapper.ServeHTTP(rec, httptest.NewRequest(http.MethodPut, validQuotaNarinfoPath, strings.NewReader(validQuotaNarinfo))) + if rec.Code != http.StatusServiceUnavailable { + t.Fatalf("failed publishing transition status = %d, want 503", rec.Code) + } + + store.mu.Lock() + store.beginCommitErr = nil + store.allowed = false + store.reason = quota.ReasonUserLimit + store.mu.Unlock() + + rec = httptest.NewRecorder() + wrapper.ServeHTTP(rec, httptest.NewRequest(http.MethodPut, validQuotaNarinfoPath, strings.NewReader(validQuotaNarinfo))) + if rec.Code != http.StatusOK { + t.Fatalf("denied late reservation status = %d, want 200", rec.Code) + } + + if _, err := os.Stat(filepath.Join(staging, "nar", "abc.nar.xz")); !errors.Is(err, os.ErrNotExist) { + t.Errorf("staged nar survived a denied late reservation: %v", err) + } + if _, err := os.Stat(filepath.Join(staging, "00000000000000000000000000000000.narinfo")); !errors.Is(err, os.ErrNotExist) { + t.Errorf("staged narinfo survived a denied late reservation: %v", err) + } + + sw := wrapper.(*stagingWrapper) + sw.stageMu.Lock() + defer sw.stageMu.Unlock() + if !sw.skippedNARs["nar/abc.nar.xz"] { + t.Fatal("denied late reservation did not mark the nar skipped") + } +} + func TestQuotaCommitFailureRetriesOnNextRequest(t *testing.T) { staging := t.TempDir() store := newMockQuotaStore(true, "") diff --git a/spindle/engines/microvm/upload_staging.go b/spindle/engines/microvm/upload_staging.go index 7560e6ab3..39b359aae 100644 --- a/spindle/engines/microvm/upload_staging.go +++ b/spindle/engines/microvm/upload_staging.go @@ -220,6 +220,14 @@ func (s *stagingWrapper) handlePutNar(w http.ResponseWriter, r *http.Request, re w.WriteHeader(http.StatusOK) } +func canonicalNarKey(raw string) (string, error) { + key := strings.TrimPrefix(raw, "/") + if strings.Contains(key, "..") || !isNarObjectPath(key) { + return "", fmt.Errorf("invalid nar URL %q", raw) + } + return key, nil +} + func (s *stagingWrapper) handlePutNarinfo(w http.ResponseWriter, r *http.Request, relPath string) { body, err := io.ReadAll(io.LimitReader(r.Body, maxNarinfoSize+1)) if err != nil { @@ -251,19 +259,20 @@ func (s *stagingWrapper) handlePutNarinfo(w http.ResponseWriter, r *http.Request http.Error(w, "narinfo filename does not match StorePath hash", http.StatusBadRequest) return } - if !isNarObjectPath(info.URL) { - s.logger.Warn("narinfo references invalid nar URL", "path", relPath, "url", info.URL) + narKey, err := canonicalNarKey(info.URL) + if err != nil { + s.logger.Warn("narinfo references invalid nar URL", "path", relPath, "url", info.URL, "error", err) http.Error(w, "invalid nar URL", http.StatusBadRequest) return } - if s.skippedNARs[info.URL] { - delete(s.skippedNARs, info.URL) + if s.skippedNARs[narKey] { + delete(s.skippedNARs, narKey) w.WriteHeader(http.StatusOK) return } - narPath, err := s.stagingObjectPath(info.URL) + narPath, err := s.stagingObjectPath(narKey) if err != nil { s.logger.Warn("narinfo references unsafe nar URL", "path", relPath, "url", info.URL, "error", err) http.Error(w, "invalid nar URL", http.StatusBadRequest) @@ -296,6 +305,38 @@ func (s *stagingWrapper) handlePutNarinfo(w http.ResponseWriter, r *http.Request http.Error(w, "internal error", http.StatusInternalServerError) return } + if s.quotaStore != nil && staged.reservationID == "" { + reservation, reserveErr := s.quotaStore.Reserve(r.Context(), quota.ReserveRequest{ + Kind: quota.KindNixCache, + Key: narKey, + Identity: quota.Identity{ + OwnerDID: s.ownerDID, + RepoDID: s.repoDID, + }, + Resources: quota.Resources{ + quota.ResourceCacheStorageBytes: narSize, + }, + }) + if reserveErr != nil { + s.logger.Error("quota reservation failed", "error", reserveErr) + _ = os.Remove(dst) + http.Error(w, "quota reservation error", http.StatusServiceUnavailable) + return + } + if s.metrics != nil { + s.metrics.RecordQuotaDecision(string(quota.KindNixCache), quota.ResourceCacheStorageBytes, reservation.Allowed, reservation.Temporary, reservation.Reason) + } + if !reservation.Allowed { + s.removeStagedNar(narPath) + _ = os.Remove(dst) + s.skippedNARs[narKey] = true + w.WriteHeader(http.StatusOK) + return + } + staged.reservationID = reservation.ID + s.stagedNARs[narPath] = staged + } + backendName := "http" if _, ok := s.backend.(*NixStoreUploadBackend); ok { backendName = "nix_store" @@ -319,8 +360,11 @@ func (s *stagingWrapper) handlePutNarinfo(w http.ResponseWriter, r *http.Request if s.quotaStore != nil { if err := s.quotaStore.BeginCommit(r.Context(), resID.ID); err != nil { s.logger.Error("quota publishing transition failed", "error", err) - s.abortCleanup(resID.ID, narPath, dst) - http.Error(w, "internal error", http.StatusInternalServerError) + s.releaseOrTrack(resID.ID) + staged.reservationID = "" + s.stagedNARs[narPath] = staged + _ = os.Remove(dst) + http.Error(w, "quota commit failed", http.StatusServiceUnavailable) return } } @@ -328,10 +372,10 @@ func (s *stagingWrapper) handlePutNarinfo(w http.ResponseWriter, r *http.Request var pubErr error var narPublished bool if usesNixStore { - narPublished = true pubErr = nixBackend.importStorePath(r.Context(), info.StorePath) + narPublished = pubErr == nil } else { - pubErr = s.publishHTTP(r.Context(), info.URL, fNar, narSize) + pubErr = s.publishHTTP(r.Context(), narKey, fNar, narSize) if pubErr == nil { narPublished = true pubErr = s.publishHTTP(r.Context(), relPath, bytes.NewReader(body), int64(len(body))) @@ -339,13 +383,13 @@ func (s *stagingWrapper) handlePutNarinfo(w http.ResponseWriter, r *http.Request } if pubErr != nil { - s.logger.Warn("backend publish failed", "error", pubErr) - if !narPublished && !s.isUncertain(pubErr) { - s.abortCleanup(resID.ID, narPath, dst) - } else { + s.logger.Warn("backend publish failed", "error", pubErr, "nar_published", narPublished) + if narPublished || s.isUncertain(pubErr) { s.removeStagedNar(narPath) _ = os.Remove(dst) s.trackCommit(resID.ID) + } else { + s.abortCleanup(resID.ID, narPath, dst) } if metrics != nil { metrics.RecordCacheUpload(backendName, "failed") @@ -374,8 +418,6 @@ func (s *stagingWrapper) handlePutNarinfo(w http.ResponseWriter, r *http.Request w.WriteHeader(http.StatusOK) } -var errHTTPPublishUncertain = errors.New("http upload result is uncertain") - func (s *stagingWrapper) publishHTTP(ctx context.Context, relPath string, body io.Reader, contentLength int64) error { targetURL, err := url.Parse(s.uploadURL) if err != nil { @@ -412,18 +454,28 @@ func (s *stagingWrapper) publishHTTP(ctx context.Context, relPath string, body i if resp.StatusCode < 200 || resp.StatusCode >= 300 { respBody, _ := io.ReadAll(io.LimitReader(resp.Body, 1024)) - return fmt.Errorf("%w: status %d: %s", errHTTPPublishUncertain, resp.StatusCode, string(respBody)) + return publishStatusError{status: resp.StatusCode, body: string(respBody)} } return nil } +type publishStatusError struct { + status int + body string +} + +func (e publishStatusError) Error() string { + return fmt.Sprintf("upload destination returned status %d: %s", e.status, e.body) +} + func (s *stagingWrapper) isUncertain(err error) bool { if err == nil { return false } - if errors.Is(err, errHTTPPublishUncertain) { - return true + var statusErr publishStatusError + if errors.As(err, &statusErr) { + return statusErr.status >= 500 } if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) { return true -- 2.51.2