bound rpm virtual merges, scope data cache per virtual, 502 on member failure
This commit is contained in:
@@ -18,6 +18,7 @@ require (
|
|||||||
github.com/testcontainers/testcontainers-go/modules/redis v0.42.0
|
github.com/testcontainers/testcontainers-go/modules/redis v0.42.0
|
||||||
github.com/ulikunitz/xz v0.5.16
|
github.com/ulikunitz/xz v0.5.16
|
||||||
golang.org/x/crypto v0.54.0
|
golang.org/x/crypto v0.54.0
|
||||||
|
golang.org/x/sync v0.22.0
|
||||||
golang.org/x/time v0.15.0
|
golang.org/x/time v0.15.0
|
||||||
gopkg.in/yaml.v3 v3.0.1
|
gopkg.in/yaml.v3 v3.0.1
|
||||||
)
|
)
|
||||||
@@ -101,7 +102,6 @@ require (
|
|||||||
go.uber.org/atomic v1.11.0 // indirect
|
go.uber.org/atomic v1.11.0 // indirect
|
||||||
go.yaml.in/yaml/v3 v3.0.4 // indirect
|
go.yaml.in/yaml/v3 v3.0.4 // indirect
|
||||||
golang.org/x/net v0.56.0 // 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/sys v0.47.0 // indirect
|
||||||
golang.org/x/text v0.40.0 // indirect
|
golang.org/x/text v0.40.0 // indirect
|
||||||
gopkg.in/ini.v1 v1.67.2 // indirect
|
gopkg.in/ini.v1 v1.67.2 // indirect
|
||||||
|
|||||||
+86
-44
@@ -18,6 +18,7 @@ import (
|
|||||||
"git.unkin.net/unkin/artifactapi/internal/provider"
|
"git.unkin.net/unkin/artifactapi/internal/provider"
|
||||||
"git.unkin.net/unkin/artifactapi/internal/proxy"
|
"git.unkin.net/unkin/artifactapi/internal/proxy"
|
||||||
"git.unkin.net/unkin/artifactapi/pkg/models"
|
"git.unkin.net/unkin/artifactapi/pkg/models"
|
||||||
|
"golang.org/x/sync/singleflight"
|
||||||
)
|
)
|
||||||
|
|
||||||
type Engine struct {
|
type Engine struct {
|
||||||
@@ -26,21 +27,24 @@ type Engine struct {
|
|||||||
getRemote func(context.Context, string) (*models.Remote, error)
|
getRemote func(context.Context, string) (*models.Remote, error)
|
||||||
rpmMember func(context.Context, string) (*RPMMember, error)
|
rpmMember func(context.Context, string) (*RPMMember, error)
|
||||||
|
|
||||||
mu sync.Mutex
|
mergeTTL time.Duration
|
||||||
rpmFiles map[string]rpmFile
|
|
||||||
|
sf singleflight.Group
|
||||||
|
mu sync.Mutex
|
||||||
|
rpm map[string]*rpmGen
|
||||||
}
|
}
|
||||||
|
|
||||||
type rpmFile struct {
|
// rpmGen holds a virtual's current merge plus the previous one, so a client
|
||||||
body []byte
|
// holding the previous repomd can still fetch its content-hashed data files.
|
||||||
expires time.Time
|
type rpmGen struct {
|
||||||
|
cur, prev *RPMRepo
|
||||||
|
at time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
// rpmFileTTL keeps content-hashed data files resolvable after repomd moves on,
|
const rpmMergeTTL = 60 * time.Second
|
||||||
// so a client that fetched the previous repomd can still fetch its data.
|
|
||||||
const rpmFileTTL = 5 * time.Minute
|
|
||||||
|
|
||||||
func NewEngine(db *database.DB, proxyEngine *proxy.Engine) *Engine {
|
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
|
e.rpmMember = e.fetchRPMMember
|
||||||
return e
|
return e
|
||||||
}
|
}
|
||||||
@@ -214,11 +218,67 @@ func (e *Engine) fetchRPM(ctx context.Context, virt models.Virtual, path string)
|
|||||||
return nil, "", ErrNotFound
|
return nil, "", ErrNotFound
|
||||||
}
|
}
|
||||||
if path != "repodata/repomd.xml" {
|
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
|
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))
|
members := make([]RPMMember, len(virt.Members))
|
||||||
errs := make([]error, len(virt.Members))
|
errs := make([]error, len(virt.Members))
|
||||||
var wg sync.WaitGroup
|
var wg sync.WaitGroup
|
||||||
@@ -235,52 +295,34 @@ func (e *Engine) fetchRPM(ctx context.Context, virt models.Virtual, path string)
|
|||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
wg.Wait()
|
wg.Wait()
|
||||||
|
// %v, not %w: a member failure is a 502 whatever its cause, never a 404.
|
||||||
if err := errors.Join(errs...); err != nil {
|
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)
|
repo, err := MergeRPM(members)
|
||||||
if err != nil {
|
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)
|
return repo, nil
|
||||||
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
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// ponytail: per-replica cache; behind a non-sticky LB a data request landing on
|
func (e *Engine) cachedRPMFile(virt, path string) ([]byte, bool) {
|
||||||
// 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()
|
|
||||||
e.mu.Lock()
|
e.mu.Lock()
|
||||||
defer e.mu.Unlock()
|
defer e.mu.Unlock()
|
||||||
if e.rpmFiles == nil {
|
g := e.rpm[virt]
|
||||||
e.rpmFiles = map[string]rpmFile{}
|
if g == nil {
|
||||||
}
|
|
||||||
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) {
|
|
||||||
return nil, false
|
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) {
|
func (e *Engine) fetchRPMMember(ctx context.Context, name string) (*RPMMember, error) {
|
||||||
|
|||||||
@@ -3,8 +3,12 @@ package virtual
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
"fmt"
|
||||||
"regexp"
|
"regexp"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
"git.unkin.net/unkin/artifactapi/pkg/models"
|
"git.unkin.net/unkin/artifactapi/pkg/models"
|
||||||
)
|
)
|
||||||
@@ -61,6 +65,71 @@ func TestFetchRPMDataSurvivesMemberChange(t *testing.T) {
|
|||||||
if err != nil || string(newMD) == string(repomd) {
|
if err != nil || string(newMD) == string(repomd) {
|
||||||
t.Fatalf("repomd should reflect the member change (err %v)", err)
|
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) {
|
func TestMemberRedirect(t *testing.T) {
|
||||||
|
|||||||
Reference in New Issue
Block a user