From e34c6570fc6bd443cb7e3b027a41fbbfff36d9fd Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Thu, 13 Aug 2026 07:10:02 +1000 Subject: [PATCH] remotes: flush cached metadata when a remote's base_url changes Switching a remote's backend (base_url) left the previously-cached mutable metadata (repodata / Release / APKINDEX) served until TTL expiry, pointing at the old upstream. cache.FlushRemote existed but was wired to nothing. - Inject a MetadataFlusher (satisfied by *cache.Redis) into RemotesHandler. - On update, read the existing remote first, then after a successful DB update flush the remote's cached metadata when base_url changed so the next request re-fetches fresh from the new upstream. - A flush failure is logged as a warning and does not fail the request; the DB update already landed. - Add tests: base_url change flushes exactly once, an unchanged base_url does not flush, and a flush error still returns 200. --- internal/api/v2/errorpaths_test.go | 2 +- internal/api/v2/remotes.go | 38 +++++++++-- internal/api/v2/remotes_flush_test.go | 96 +++++++++++++++++++++++++++ internal/server/server.go | 2 +- 4 files changed, 132 insertions(+), 6 deletions(-) create mode 100644 internal/api/v2/remotes_flush_test.go diff --git a/internal/api/v2/errorpaths_test.go b/internal/api/v2/errorpaths_test.go index aecf463..63c0948 100644 --- a/internal/api/v2/errorpaths_test.go +++ b/internal/api/v2/errorpaths_test.go @@ -57,7 +57,7 @@ func do(t *testing.T, h http.Handler, method, path, body string) int { } func TestRemotesErrorPaths(t *testing.T) { - h := NewRemotesHandler(closedDB(t), nil).Routes() + h := NewRemotesHandler(closedDB(t), nil, nil).Routes() if c := do(t, h, "GET", "/", ""); c != 500 { t.Errorf("list with dead db = %d, want 500", c) } diff --git a/internal/api/v2/remotes.go b/internal/api/v2/remotes.go index b8d1e36..14c25c7 100644 --- a/internal/api/v2/remotes.go +++ b/internal/api/v2/remotes.go @@ -1,8 +1,10 @@ package v2 import ( + "context" "encoding/json" "fmt" + "log/slog" "net/http" "github.com/go-chi/chi/v5" @@ -17,15 +19,23 @@ type Primer interface { EnqueuePrime(remote models.Remote) } +// MetadataFlusher purges a remote's cached mutable metadata (repodata / Release +// / APKINDEX freshness keys). *cache.Redis satisfies it. +type MetadataFlusher interface { + FlushRemote(ctx context.Context, remote string) error +} + type RemotesHandler struct { db *database.DB + cache MetadataFlusher primers map[models.PackageType]Primer } -// NewRemotesHandler wires the handler to the per-type metadata primers. primers -// may be nil; a package type with no registered primer simply skips priming. -func NewRemotesHandler(db *database.DB, primers map[models.PackageType]Primer) *RemotesHandler { - return &RemotesHandler{db: db, primers: primers} +// NewRemotesHandler wires the handler to the metadata cache and per-type +// primers. cache may be nil (flush-on-backend-change is skipped); primers may +// be nil (a package type with no registered primer simply skips priming). +func NewRemotesHandler(db *database.DB, cache MetadataFlusher, primers map[models.PackageType]Primer) *RemotesHandler { + return &RemotesHandler{db: db, cache: cache, primers: primers} } func (h *RemotesHandler) Routes() chi.Router { @@ -106,10 +116,30 @@ func (h *RemotesHandler) update(w http.ResponseWriter, r *http.Request) { http.Error(w, err.Error(), http.StatusBadRequest) return } + // Capture the current backend before the update so we can tell whether the + // remote's base_url (its upstream) changed. A read failure just means we + // skip the freshness flush; it must not block the update. + oldBaseURL, oldKnown := "", false + if existing, err := h.db.GetRemote(r.Context(), name); err == nil { + oldBaseURL, oldKnown = existing.BaseURL, true + } if err := h.db.UpdateRemote(r.Context(), &remote); err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) return } + // Changing the backend invalidates any cached mutable metadata (repodata / + // Release / APKINDEX): purge it so the next request re-fetches from the new + // upstream instead of serving stale data until TTL expiry. A flush failure + // is logged but does not fail the request — the DB update already landed. + if oldKnown && oldBaseURL != remote.BaseURL && h.cache != nil { + if err := h.cache.FlushRemote(r.Context(), name); err != nil { + slog.Warn("flush cached metadata after base_url change failed", + "remote", name, "error", err) + } else { + slog.Info("flushed cached metadata after base_url change", + "remote", name, "old_base_url", oldBaseURL, "new_base_url", remote.BaseURL) + } + } writeJSON(w, http.StatusOK, remote) } diff --git a/internal/api/v2/remotes_flush_test.go b/internal/api/v2/remotes_flush_test.go new file mode 100644 index 0000000..09a25b5 --- /dev/null +++ b/internal/api/v2/remotes_flush_test.go @@ -0,0 +1,96 @@ +package v2 + +import ( + "context" + "errors" + "testing" + + "git.unkin.net/unkin/artifactapi/internal/database" + "git.unkin.net/unkin/artifactapi/pkg/models" +) + +// fakeFlusher records FlushRemote calls so a test can assert whether — and how +// often — a remote's cached metadata was purged. +type fakeFlusher struct { + calls []string + err error +} + +func (f *fakeFlusher) FlushRemote(_ context.Context, remote string) error { + f.calls = append(f.calls, remote) + return f.err +} + +func seedRemote(t *testing.T, db *database.DB, name, baseURL string) { + t.Helper() + err := db.CreateRemote(context.Background(), &models.Remote{ + Name: name, + PackageType: models.PackageRPM, + RepoType: models.RepoTypeRemote, + BaseURL: baseURL, + }) + if err != nil { + t.Fatalf("seed remote: %v", err) + } +} + +// A base_url change must flush the remote's cached metadata exactly once, while +// an update that leaves base_url untouched must not flush at all. +func TestUpdateFlushesCacheOnBaseURLChange(t *testing.T) { + if testDSN == "" { + t.Skip("Docker unavailable") + } + db, err := database.New(testDSN) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + const name = "rpm-flush-change" + seedRemote(t, db, name, "https://old.example.com/repo") + + ff := &fakeFlusher{} + h := NewRemotesHandler(db, ff, nil).Routes() + + if c := do(t, h, "PUT", "/"+name, `{"package_type":"rpm","repo_type":"remote","base_url":"https://new.example.com/repo"}`); c != 200 { + t.Fatalf("update (backend change) = %d, want 200", c) + } + if len(ff.calls) != 1 || ff.calls[0] != name { + t.Fatalf("flush calls = %v, want exactly one flush of %q", ff.calls, name) + } + + // Re-updating with the same (now current) base_url must not flush again. + ff.calls = nil + if c := do(t, h, "PUT", "/"+name, `{"package_type":"rpm","repo_type":"remote","base_url":"https://new.example.com/repo"}`); c != 200 { + t.Fatalf("update (no backend change) = %d, want 200", c) + } + if len(ff.calls) != 0 { + t.Fatalf("flush calls = %v, want no flush when base_url is unchanged", ff.calls) + } +} + +// A flush error must be swallowed: the DB update already succeeded, so the +// request still returns 200. +func TestUpdateFlushFailureStillSucceeds(t *testing.T) { + if testDSN == "" { + t.Skip("Docker unavailable") + } + db, err := database.New(testDSN) + if err != nil { + t.Fatal(err) + } + defer db.Close() + + const name = "rpm-flush-error" + seedRemote(t, db, name, "https://old.example.com/repo") + + ff := &fakeFlusher{err: errors.New("redis down")} + h := NewRemotesHandler(db, ff, nil).Routes() + + if c := do(t, h, "PUT", "/"+name, `{"package_type":"rpm","repo_type":"remote","base_url":"https://new.example.com/repo"}`); c != 200 { + t.Fatalf("update with failing flush = %d, want 200", c) + } + if len(ff.calls) != 1 { + t.Fatalf("flush calls = %v, want exactly one attempted flush", ff.calls) + } +} diff --git a/internal/server/server.go b/internal/server/server.go index dd0480b..18cf55d 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -175,7 +175,7 @@ func (s *Server) routes() chi.Router { r.Mount("/api/v1", proxyHandler.Routes()) r.Mount("/v2", proxyHandler.DockerV2Routes()) - remotesHandler := v2.NewRemotesHandler(s.db, map[models.PackageType]v2.Primer{ + remotesHandler := v2.NewRemotesHandler(s.db, s.cache, map[models.PackageType]v2.Primer{ models.PackageGitHubRPM: s.syncer, models.PackageGitHubDeb: s.debSyncer, models.PackageGitHubAlpine: s.alpineSyncer, -- 2.47.3