remotes: flush cached metadata when a remote's base_url changes #120
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user