From 2ca6be11d28125e26293435e394502ed4cb0f22f Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Fri, 9 Oct 2026 20:21:22 +1100 Subject: [PATCH] bound rpm virtual merges, scope data cache per virtual, 502 on member failure --- go.mod | 2 +- internal/virtual/engine.go | 130 +++++++++++++++++++++----------- internal/virtual/engine_test.go | 69 +++++++++++++++++ 3 files changed, 156 insertions(+), 45 deletions(-) diff --git a/go.mod b/go.mod index 6333d15..6583252 100644 --- a/go.mod +++ b/go.mod @@ -18,6 +18,7 @@ require ( github.com/testcontainers/testcontainers-go/modules/redis v0.42.0 github.com/ulikunitz/xz v0.5.16 golang.org/x/crypto v0.54.0 + golang.org/x/sync v0.22.0 golang.org/x/time v0.15.0 gopkg.in/yaml.v3 v3.0.1 ) @@ -101,7 +102,6 @@ require ( go.uber.org/atomic v1.11.0 // indirect go.yaml.in/yaml/v3 v3.0.4 // indirect golang.org/x/net v0.56.0 // indirect - golang.org/x/sync v0.22.0 // indirect golang.org/x/sys v0.47.0 // indirect golang.org/x/text v0.40.0 // indirect gopkg.in/ini.v1 v1.67.2 // indirect diff --git a/internal/virtual/engine.go b/internal/virtual/engine.go index 0c680f2..d4c52d0 100644 --- a/internal/virtual/engine.go +++ b/internal/virtual/engine.go @@ -18,6 +18,7 @@ import ( "git.unkin.net/unkin/artifactapi/internal/provider" "git.unkin.net/unkin/artifactapi/internal/proxy" "git.unkin.net/unkin/artifactapi/pkg/models" + "golang.org/x/sync/singleflight" ) type Engine struct { @@ -26,21 +27,24 @@ type Engine struct { getRemote func(context.Context, string) (*models.Remote, error) rpmMember func(context.Context, string) (*RPMMember, error) - mu sync.Mutex - rpmFiles map[string]rpmFile + mergeTTL time.Duration + + sf singleflight.Group + mu sync.Mutex + rpm map[string]*rpmGen } -type rpmFile struct { - body []byte - expires time.Time +// rpmGen holds a virtual's current merge plus the previous one, so a client +// holding the previous repomd can still fetch its content-hashed data files. +type rpmGen struct { + cur, prev *RPMRepo + at time.Time } -// rpmFileTTL keeps content-hashed data files resolvable after repomd moves on, -// so a client that fetched the previous repomd can still fetch its data. -const rpmFileTTL = 5 * time.Minute +const rpmMergeTTL = 60 * time.Second func NewEngine(db *database.DB, proxyEngine *proxy.Engine) *Engine { - e := &Engine{db: db, proxyEngine: proxyEngine, getRemote: db.GetRemote} + e := &Engine{db: db, proxyEngine: proxyEngine, getRemote: db.GetRemote, mergeTTL: rpmMergeTTL} e.rpmMember = e.fetchRPMMember return e } @@ -214,11 +218,67 @@ func (e *Engine) fetchRPM(ctx context.Context, virt models.Virtual, path string) return nil, "", ErrNotFound } if path != "repodata/repomd.xml" { - if body, ok := e.cachedRPMFile(path); ok { + if body, ok := e.cachedRPMFile(virt.Name, path); ok { return body, "application/gzip", nil } } + repo, err := e.mergedRPM(ctx, virt) + if err != nil { + return nil, "", err + } + if path == "repodata/repomd.xml" { + return repo.Repomd, "application/xml", nil + } + if body, ok := repo.Files[path]; ok { + return body, "application/gzip", nil + } + return nil, "", ErrNotFound +} + +// mergedRPM returns the virtual's merge, reusing it for mergeTTL and +// collapsing concurrent merges of the same virtual into one. +// ponytail: per-replica cache; behind a non-sticky LB a data request landing on +// another replica re-merges and 404s if a member changed in between. Move to the +// shared redis cache if that shows up. +func (e *Engine) mergedRPM(ctx context.Context, virt models.Virtual) (*RPMRepo, error) { + e.mu.Lock() + g := e.rpm[virt.Name] + if g != nil && time.Since(g.at) < e.mergeTTL { + e.mu.Unlock() + return g.cur, nil + } + e.mu.Unlock() + + v, err, _ := e.sf.Do(virt.Name, func() (any, error) { + repo, err := e.mergeRPM(context.WithoutCancel(ctx), virt) + if err != nil { + return nil, err + } + e.mu.Lock() + defer e.mu.Unlock() + if e.rpm == nil { + e.rpm = map[string]*rpmGen{} + } + g := e.rpm[virt.Name] + switch { + case g == nil: + e.rpm[virt.Name] = &rpmGen{cur: repo, at: time.Now()} + case string(g.cur.Repomd) == string(repo.Repomd): + g.at = time.Now() + repo = g.cur + default: + g.prev, g.cur, g.at = g.cur, repo, time.Now() + } + return repo, nil + }) + if err != nil { + return nil, err + } + return v.(*RPMRepo), nil +} + +func (e *Engine) mergeRPM(ctx context.Context, virt models.Virtual) (*RPMRepo, error) { members := make([]RPMMember, len(virt.Members)) errs := make([]error, len(virt.Members)) var wg sync.WaitGroup @@ -235,52 +295,34 @@ func (e *Engine) fetchRPM(ctx context.Context, virt models.Virtual, path string) }() } wg.Wait() + // %v, not %w: a member failure is a 502 whatever its cause, never a 404. if err := errors.Join(errs...); err != nil { - return nil, "", fmt.Errorf("virtual %q: %w", virt.Name, err) + return nil, fmt.Errorf("virtual %q: %v", virt.Name, err) } repo, err := MergeRPM(members) if err != nil { - return nil, "", fmt.Errorf("merge rpm repodata: %w", err) + return nil, fmt.Errorf("merge rpm repodata: %w", err) } - e.storeRPMFiles(repo.Files) - if path == "repodata/repomd.xml" { - return repo.Repomd, "application/xml", nil - } - if body, ok := repo.Files[path]; ok { - return body, "application/gzip", nil - } - return nil, "", ErrNotFound + return repo, nil } -// ponytail: per-replica cache; behind a non-sticky LB a data request landing on -// another replica re-merges and 404s if a member changed in between. Move to the -// shared redis cache if that shows up. -func (e *Engine) storeRPMFiles(files map[string][]byte) { - now := time.Now() +func (e *Engine) cachedRPMFile(virt, path string) ([]byte, bool) { e.mu.Lock() defer e.mu.Unlock() - if e.rpmFiles == nil { - e.rpmFiles = map[string]rpmFile{} - } - for k, f := range e.rpmFiles { - if now.After(f.expires) { - delete(e.rpmFiles, k) - } - } - for k, body := range files { - e.rpmFiles[k] = rpmFile{body: body, expires: now.Add(rpmFileTTL)} - } -} - -func (e *Engine) cachedRPMFile(path string) ([]byte, bool) { - e.mu.Lock() - defer e.mu.Unlock() - f, ok := e.rpmFiles[path] - if !ok || time.Now().After(f.expires) { + g := e.rpm[virt] + if g == nil { return nil, false } - return f.body, true + for _, r := range []*RPMRepo{g.cur, g.prev} { + if r == nil { + continue + } + if body, ok := r.Files[path]; ok { + return body, true + } + } + return nil, false } func (e *Engine) fetchRPMMember(ctx context.Context, name string) (*RPMMember, error) { diff --git a/internal/virtual/engine_test.go b/internal/virtual/engine_test.go index 0571a56..d1c28c9 100644 --- a/internal/virtual/engine_test.go +++ b/internal/virtual/engine_test.go @@ -3,8 +3,12 @@ package virtual import ( "context" "errors" + "fmt" "regexp" + "sync" + "sync/atomic" "testing" + "time" "git.unkin.net/unkin/artifactapi/pkg/models" ) @@ -61,6 +65,71 @@ func TestFetchRPMDataSurvivesMemberChange(t *testing.T) { if err != nil || string(newMD) == string(repomd) { t.Fatalf("repomd should reflect the member change (err %v)", err) } + if _, _, err := e.Fetch(context.Background(), virt, string(href), ""); err != nil { + t.Fatalf("previous generation must survive one re-merge: %v", err) + } +} + +func TestFetchRPMLocalMemberWithoutRepodataIsUpstreamError(t *testing.T) { + e := fakeEngine(map[string]*RPMMember{ + "a": {RemoteName: "a", Data: map[string][]byte{"primary": primaryXML(primaryPkgXML("foo", "1", "aaa", "foo.rpm"))}}, + }) + e.rpmMember = func(_ context.Context, name string) (*RPMMember, error) { + if name == "local" { + return nil, fmt.Errorf("local/repodata/repomd.xml: %w", ErrNotFound) + } + return &RPMMember{RemoteName: name, Data: map[string][]byte{"primary": primaryXML(primaryPkgXML("foo", "1", "aaa", "foo.rpm"))}}, nil + } + _, _, err := e.Fetch(context.Background(), rpmVirt("a", "local"), "repodata/repomd.xml", "") + if err == nil || errors.Is(err, ErrNotFound) { + t.Fatalf("member failure must not surface as not-found, got %v", err) + } +} + +func TestFetchRPMMergeIsShared(t *testing.T) { + var calls atomic.Int32 + e := fakeEngine(nil) + e.mergeTTL = time.Minute + e.rpmMember = func(_ context.Context, name string) (*RPMMember, error) { + calls.Add(1) + time.Sleep(20 * time.Millisecond) + return &RPMMember{RemoteName: name, Data: map[string][]byte{"primary": primaryXML(primaryPkgXML("foo", "1", "aaa", "foo.rpm"))}}, nil + } + var wg sync.WaitGroup + for range 10 { + wg.Add(1) + go func() { + defer wg.Done() + if _, _, err := e.Fetch(context.Background(), rpmVirt("a"), "repodata/repomd.xml", ""); err != nil { + t.Error(err) + } + }() + } + wg.Wait() + if _, _, err := e.Fetch(context.Background(), rpmVirt("a"), "repodata/repomd.xml", ""); err != nil { + t.Fatal(err) + } + if n := calls.Load(); n != 1 { + t.Fatalf("member fetched %d times, want 1", n) + } +} + +func TestFetchRPMDataScopedToVirtual(t *testing.T) { + e := fakeEngine(map[string]*RPMMember{ + "a": {RemoteName: "a", Data: map[string][]byte{"primary": primaryXML(primaryPkgXML("foo", "1", "aaa", "foo.rpm"))}}, + "b": {RemoteName: "b", Data: map[string][]byte{"primary": primaryXML(primaryPkgXML("bar", "1", "bbb", "bar.rpm"))}}, + }) + virtA := models.Virtual{Name: "va", PackageType: models.PackageRPM, Members: []string{"a"}} + virtB := models.Virtual{Name: "vb", PackageType: models.PackageRPM, Members: []string{"b"}} + + repomd, _, err := e.Fetch(context.Background(), virtA, "repodata/repomd.xml", "") + if err != nil { + t.Fatal(err) + } + href := regexp.MustCompile(`repodata/[0-9a-f]+-primary\.xml\.gz`).Find(repomd) + if _, _, err := e.Fetch(context.Background(), virtB, string(href), ""); !errors.Is(err, ErrNotFound) { + t.Fatalf("virtual vb served va's data file (err %v)", err) + } } func TestMemberRedirect(t *testing.T) {