Compare commits

..

9 Commits

Author SHA1 Message Date
unkin-agent 492607a164 Retry failed GitHub release scans with backoff (#133)
ci/woodpecker/tag/docker Pipeline was successful
Releasing a GitHub sync lease always advanced `last_synced_at`, even after a failed scan (e.g. a rate-limit 403). A remote that failed once waited a full `mutable_ttl` before retrying, and a cold remote kept returning 503 until then.

- record scan outcomes in one shared lease helper for github_rpm/deb/alpine
- keep `last_synced_at` and the ETag on failure; retry from 60s with exponential backoff, capped at min(10m, ttl/4)
- honour `Retry-After` / `X-RateLimit-Reset`, clamped to `mutable_ttl`
- add `sync_failures` / `next_retry_at` columns (migration 0002)

Reviewed-on: #133
Co-authored-by: unkin-agent <unkin-agent@unkin.net>
Co-committed-by: unkin-agent <unkin-agent@unkin.net>
2026-10-09 23:15:55 +11:00
unkin-agent 1780b3d77c Push release images to artifactapi docker-internal (#134)
Gitea's registry rejects blob HEADs with 401, so the v3.13.0 tag pipeline cannot push. Release images move to the in-cluster artifactapi docker-internal registry, which trusts the runner without credentials.

- push docker-api to docker-internal/artifactapi and docker-web to docker-internal/artifactapi-ui
- use the CA-baked docker-internal/plugin-docker-buildx image
- drop the droneci credentials
- set kubernetes resources and service account on both steps

Reviewed-on: #134
Co-authored-by: unkin-agent <unkin-agent@unkin.net>
Co-committed-by: unkin-agent <unkin-agent@unkin.net>
2026-10-09 23:08:22 +11:00
unkin-agent 1837f6ef8c Add rpm virtual repositories (#131)
ci/woodpecker/tag/docker Pipeline failed
Virtual repos only merge helm and pypi, so several rpm repos (e.g. many github_rpm remotes) cannot be served as one yum repo.

- merge member primary/filelists/other into one repodata set; member order wins duplicate NEVRAs
- read member repodata through local, github_rpm and proxied remote paths
- prefix package locations (incl. xml:base under a member upstream) with the member name and 302 them to the member route
- reject absolute and dot-segment member paths; escape the redirect and keep its query
- return 502 when any member's repodata is unavailable
- reuse each virtual's merge for 60s behind singleflight; serve data files from its current and previous merge (per replica)
- add merger/engine unit tests and a dockerised dnf e2e case

Reviewed-on: #131
Co-authored-by: unkin-agent <unkin-agent@unkin.net>
Co-committed-by: unkin-agent <unkin-agent@unkin.net>
2026-10-09 22:07:20 +11:00
unkin-agent 80368986bb Add cargo sparse registry remote type (#132)
Rust crate builds (rpmbuilder CI) fetch from crates.io directly because artifactapi has no Cargo remote type.

- add `cargo` package type proxying the sparse registry protocol (RFC 2789)
- synthesize `config.json` so `dl` points crate downloads back at the remote
- treat index files as mutable and `crates/*/*.crate` as immutable, fetched from static.crates.io for crates.io
- add a `~/.cargo/config.toml` usage snippet to the UI and a cargo e2e caching case

Reviewed-on: #132
Co-authored-by: unkin-agent <unkin-agent@unkin.net>
Co-committed-by: unkin-agent <unkin-agent@unkin.net>
2026-10-09 21:44:17 +11:00
unkin-agent bda762b10e Fix objects test compile after Routes evictor change (#130)
ci/woodpecker/tag/docker Pipeline was successful
ObjectsHandler.Routes now takes an Evictor, but TestObjectsListPrefix still calls it with no arguments, so the api/v2 test package fails to compile on master.

- pass the existing fakeEvictor to Routes in objects_test.go

Reviewed-on: #130
Co-authored-by: unkin-agent <unkin-agent@unkin.net>
Co-committed-by: unkin-agent <unkin-agent@unkin.net>
2026-10-04 22:11:19 +11:00
unkin-agent a71d126239 Filter remote object listings by prefix (#129)
`GET /api/v2/remotes/{name}/objects` dropped the `prefix` query parameter, so a filtered listing returned the first page of every artifact in the remote.

- add a `prefix` argument to `ListArtifacts` that matches path prefixes
- pass the `prefix` query parameter through from the objects list handler

Reviewed-on: #129
Co-authored-by: unkin-agent <unkin-agent@unkin.net>
Co-committed-by: unkin-agent <unkin-agent@unkin.net>
2026-10-04 22:07:21 +11:00
unkin-agent 6a08539a78 Evict remote objects from every cache layer (#128)
Evicting a remote object only deleted its Postgres row. The S3 index and the Redis TTL/ETag keys survived, so stale mirror metadata such as EPEL repodata kept being served.

- evict the artifact row, S3 index and Redis keys via `Engine.Evict`
- evict a directory with `<dir>/*`; other wildcards return 400, unknown remotes 404
- wait on the per-path fetch lock for single-path evicts; return 503 on timeout or a Redis error
- wildcard evicts take no lock; a fetch already in flight may re-cache its path
- return S3 list errors from `DeletePrefix` instead of deleting nothing

Reviewed-on: #128
Co-authored-by: unkin-agent <unkin-agent@unkin.net>
Co-committed-by: unkin-agent <unkin-agent@unkin.net>
2026-10-04 22:07:14 +11:00
unkin-agent 0b159dad90 ci: move Go steps to the estate-built gobuilder image (#127)
Packer-built almalinux9-gobuilder is retiring. Moves test and
pre-commit steps to artifactapi docker-internal/gobuilder 0.1.2-alma9
(estate CA, go1.26.7, go-cache-plugin baked in).

- Repoint both steps at the new image, quoted (has a colon)
- Drop curl/sha256sum bootstrap for go-cache-plugin; it's on PATH now
- Set GOCACHEPROG as a static env var instead of a conditional export
- Verified build/test still pass with unreachable/invalid S3 creds —
  best-effort degrade is a plugin runtime property, unaffected

Reviewed-on: #127
Co-authored-by: unkin-agent <unkin-agent@unkin.net>
Co-committed-by: unkin-agent <unkin-agent@unkin.net>
2026-10-03 00:26:02 +10:00
unkin-agent 7a4f4054cc ci: back Go compiles with the shared S3 build cache (#126)
Every pipeline run recompiled the whole tree from scratch; the shared gocache S3 bucket already holds those objects.

- Bootstrap the go-cache-plugin GOCACHEPROG in the test and pre-commit steps, keyed `ci-artifactapi`.
- Best-effort: a failed plugin fetch leaves GOCACHEPROG unset, and S3 errors degrade to cache misses, so an outage never fails a build.
- Run tests on the gobuilder image, which trusts the internal CA the S3 endpoint presents.

Container builds still compile uncached; plumbing the cache into the Dockerfile is a follow-up.

Reviewed-on: #126
Co-authored-by: unkin-agent <unkin-agent@unkin.net>
Co-committed-by: unkin-agent <unkin-agent@unkin.net>
2026-09-27 17:33:19 +10:00
51 changed files with 2523 additions and 166 deletions
+26 -12
View File
@@ -4,31 +4,45 @@ when:
steps:
- name: docker-api
image: woodpeckerci/plugin-docker-buildx
image: artifactapi.k8s.syd1.au.unkin.net/docker-internal/plugin-docker-buildx:latest
settings:
registry: git.unkin.net
repo: git.unkin.net/unkin/artifactapi
registry: artifactapi.k8s.syd1.au.unkin.net
repo: artifactapi.k8s.syd1.au.unkin.net/docker-internal/artifactapi
build_args:
VERSION: ${CI_COMMIT_TAG}
username: droneci
password:
from_secret: DRONECI_PASSWORD
tags:
- ${CI_COMMIT_TAG}
- latest
backend_options:
kubernetes:
serviceAccountName: default
resources:
requests:
memory: 1Gi
cpu: 1
limits:
memory: 4Gi
cpu: 2
- name: docker-web
image: woodpeckerci/plugin-docker-buildx
image: artifactapi.k8s.syd1.au.unkin.net/docker-internal/plugin-docker-buildx:latest
settings:
registry: git.unkin.net
repo: git.unkin.net/unkin/artifactapi-ui
registry: artifactapi.k8s.syd1.au.unkin.net
repo: artifactapi.k8s.syd1.au.unkin.net/docker-internal/artifactapi-ui
dockerfile: ui/Dockerfile.ui
context: ui
build_args:
BASE_PATH: /ui
username: droneci
password:
from_secret: DRONECI_PASSWORD
tags:
- ${CI_COMMIT_TAG}
- latest
backend_options:
kubernetes:
serviceAccountName: default
resources:
requests:
memory: 1Gi
cpu: 1
limits:
memory: 4Gi
cpu: 2
+15 -1
View File
@@ -3,14 +3,28 @@ when:
steps:
- name: pre-commit
image: git.unkin.net/unkin/almalinux9-gobuilder:20260606
# gobuilder trusts the internal CA, which the S3 build cache endpoint needs.
# go-cache-plugin is baked into the image; S3 errors degrade to cache misses.
image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
commands:
- uvx pre-commit run --all-files
environment:
# golib lives on Gitea; skip the public proxy/sum db.
GOPRIVATE: git.unkin.net
GOCACHEPROG: "go-cache-plugin --cache-dir=/tmp/gocache"
GOCACHE_S3_BUCKET: gocache
# Explicit region skips a GetBucketLocation probe RGW handles poorly.
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-artifactapi
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
backend_options:
kubernetes:
serviceAccountName: default
resources:
requests:
memory: 512Mi
+24 -1
View File
@@ -3,9 +3,32 @@ when:
steps:
- name: test
image: golang:1.25
# gobuilder trusts the internal CA, which the S3 build cache endpoint needs.
# go-cache-plugin is baked into the image; S3 errors degrade to cache misses.
image: "artifactapi.k8s.syd1.au.unkin.net/docker-internal/gobuilder:0.1.2-alma9"
commands:
- go test -race -count=1 ./pkg/... ./internal/...
environment:
# golib lives on Gitea; skip the public proxy/sum db.
GOPRIVATE: git.unkin.net
GOCACHEPROG: "go-cache-plugin --cache-dir=/tmp/gocache"
GOCACHE_S3_BUCKET: gocache
# Explicit region skips a GetBucketLocation probe RGW handles poorly.
GOCACHE_S3_REGION: us-east-1
GOCACHE_S3_ENDPOINT_URL: "https://s3.ceph.unkin.net"
GOCACHE_S3_PATH_STYLE: "true"
GOCACHE_KEY_PREFIX: ci-artifactapi
AWS_ACCESS_KEY_ID:
from_secret: GOCACHE_AWS_ACCESS_KEY_ID
AWS_SECRET_ACCESS_KEY:
from_secret: GOCACHE_AWS_SECRET_ACCESS_KEY
backend_options:
kubernetes:
serviceAccountName: default
resources:
requests:
memory: 1Gi
cpu: 1
limits:
memory: 4Gi
cpu: 2
+2 -1
View File
@@ -1,6 +1,6 @@
# ArtifactAPI
Caching proxy for package repositories. Single Go binary, 10 package types, content-addressable storage, managed by Terraform.
Caching proxy for package repositories. Single Go binary, 11 package types, content-addressable storage, managed by Terraform.
## Quick Start
@@ -32,6 +32,7 @@ API: `http://localhost:8000` | Frontend: `http://localhost:5173`
| `puppet` | `v3/modules/*`, `v3/releases*` | `.tar.gz` |
| `terraform` | `*/versions` | `*/download/*/*` |
| `goproxy` | `@v/list`, `@latest` | `.info`, `.mod`, `.zip` |
| `cargo` | sparse index files (`config.json` synthesized) | `crates/*/*.crate` |
| `github_rpm` | `repodata/*` (synthesized) | `.rpm` (redirected) |
Providers classify paths automatically. Users only configure what to proxy and TTLs.
+1 -1
View File
@@ -23,7 +23,7 @@ already-running stack.
- **Repository lifecycle** — add / change / delete for remote, local and virtual repos.
- **Caching** — one immutable artifact per remote package type (generic, docker,
helm, pypi, npm, rpm, alpine, puppet, terraform, goproxy) proxied through the
helm, pypi, npm, rpm, alpine, puppet, terraform, goproxy, cargo) proxied through the
mock upstream: first fetch `X-Artifact-Source: remote`, second `cache`, bytes
verified against the origin fixture.
- **Local uploads** — generic (upload/download), pypi (wheel + generated `simple/`
+1
View File
@@ -29,6 +29,7 @@ func TestCachingPerProvider(t *testing.T) {
{"alpine", "alpine/x86_64/testpkg-1.0-r0.apk", "alpine/x86_64/testpkg-1.0-r0.apk"},
{"puppet", "puppet-releases/author-mod-1.0.0.tar.gz", "puppet-releases/author-mod-1.0.0.tar.gz"},
{"goproxy", "goproxy/example.com/mod/@v/v1.0.0.zip", "goproxy/example.com/mod/@v/v1.0.0.zip"},
{"cargo", "crates/mycrate/mycrate-1.0.0.crate", "crates/mycrate/mycrate-1.0.0.crate"},
{"terraform", "hashicorp/aws/download/pkg.zip", "v1/providers/hashicorp/aws/download/pkg.zip"},
{"docker", "library/testimg/blobs/blobdata", "v2/library/testimg/blobs/blobdata"},
}
+84
View File
@@ -3,9 +3,17 @@
package e2edocker
import (
"bytes"
"compress/gzip"
"encoding/xml"
"fmt"
"io"
"net/http"
"os"
"os/exec"
"strings"
"testing"
"time"
)
// TestVirtualPyPIMerge uploads different packages to two pypi locals and
@@ -52,3 +60,79 @@ func TestVirtualHelmMerge(t *testing.T) {
t.Fatalf("merged helm index missing a member chart (want alpha and beta): %s", s)
}
}
// TestVirtualRPMMerge merges a remote rpm repo and a local rpm repo carrying the
// same package: the first member wins the duplicate, its package download
// routes through the virtual, and a real dnf can consume the merged repo.
func TestVirtualRPMMerge(t *testing.T) {
createRepo(t, `{"name":"vrpm-remote","package_type":"rpm","repo_type":"remote","base_url":"`+mockUpstream()+`/rpm-mirror","stale_on_error":true}`)
createRepo(t, `{"name":"vrpm-local","package_type":"rpm","repo_type":"local"}`)
defer deleteRepo(t, "vrpm-remote")
defer deleteRepo(t, "vrpm-local")
pkg := fixtureBytes(t, "rpmrepo/Packages/e2e-testpkg-1.0-1.noarch.rpm")
uploadFile(t, "vrpm-local", "e2e-testpkg-1.0-1.noarch.rpm", pkg, "application/x-rpm")
createVirtual(t, `{"name":"vrpm","package_type":"rpm","members":["vrpm-remote","vrpm-local"]}`)
defer deleteVirtual(t, "vrpm")
resp, body := getEventually(t, api("/api/v1/virtual/vrpm/repodata/repomd.xml"), 15*time.Second)
if resp.StatusCode != http.StatusOK {
t.Fatalf("virtual repomd.xml: status %d: %s", resp.StatusCode, body)
}
var md struct {
Data []struct {
Type string `xml:"type,attr"`
Location struct {
Href string `xml:"href,attr"`
} `xml:"location"`
} `xml:"data"`
}
if err := xml.Unmarshal(body, &md); err != nil {
t.Fatalf("parse repomd.xml: %v\n%s", err, body)
}
var primary []byte
for _, d := range md.Data {
if d.Type == "primary" {
resp, gz := doRequest(t, http.MethodGet, api("/api/v1/virtual/vrpm/"+d.Location.Href), nil, "")
if resp.StatusCode != http.StatusOK {
t.Fatalf("primary: status %d", resp.StatusCode)
}
r, err := gzip.NewReader(bytes.NewReader(gz))
if err != nil {
t.Fatal(err)
}
primary, _ = io.ReadAll(r)
}
}
href := "vrpm-remote/Packages/e2e-testpkg-1.0-1.noarch.rpm"
if strings.Count(string(primary), "<name>e2e-testpkg</name>") != 1 || !strings.Contains(string(primary), `href="`+href+`"`) {
t.Fatalf("merged primary should hold one e2e-testpkg owned by the first member:\n%s", primary)
}
resp, got := doRequest(t, http.MethodGet, api("/api/v1/virtual/vrpm/"+href), nil, "")
if resp.StatusCode != http.StatusOK || !bytes.Equal(got, pkg) {
t.Fatalf("package via virtual: status %d, %d bytes (want %d)", resp.StatusCode, len(got), len(pkg))
}
network := os.Getenv("COMPOSE_NETWORK")
internal := os.Getenv("ARTIFACTAPI_INTERNAL")
if network == "" || internal == "" {
t.Skip("COMPOSE_NETWORK/ARTIFACTAPI_INTERNAL not set; skipping real dnf")
}
if _, err := exec.LookPath("docker"); err != nil {
t.Skip("docker not available on the test host")
}
repoConf := fmt.Sprintf("[vrpm]\nname=vrpm\nbaseurl=%s/api/v1/virtual/vrpm/\nenabled=1\ngpgcheck=0\nmetadata_expire=0\n", strings.TrimRight(internal, "/"))
script := "set -euo pipefail; " +
"printf '%s' \"$REPO\" > /etc/yum.repos.d/vrpm.repo; " +
"dnf -y --disablerepo='*' --enablerepo=vrpm makecache; " +
"dnf -y --disablerepo='*' --enablerepo=vrpm list --available; " +
"dnf -y --disablerepo='*' --enablerepo=vrpm install e2e-testpkg; " +
"rpm -q e2e-testpkg"
out, err := exec.Command("docker", "run", "--rm", "--network", network, "-e", "REPO="+repoConf,
"rockylinux:9", "bash", "-c", script).CombinedOutput()
if err != nil {
t.Fatalf("real dnf against rpm virtual failed: %v\n%s", err, out)
}
}
+1 -1
View File
@@ -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
+14
View File
@@ -162,7 +162,21 @@ func (h *ProxyHandler) handleVirtual(w http.ResponseWriter, r *http.Request) {
proxyBaseURL := fmt.Sprintf("%s://%s", scheme(r), r.Host)
loc, ok, err := h.virtualEngine.MemberRedirect(r.Context(), *virt, path, r.URL.RawQuery, proxyBaseURL)
if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
if ok {
http.Redirect(w, r, loc, http.StatusFound)
return
}
body, contentType, err := h.virtualEngine.Fetch(r.Context(), *virt, path, proxyBaseURL)
if errors.Is(err, virtual.ErrNotFound) {
http.Error(w, "not found", http.StatusNotFound)
return
}
if err != nil {
slog.Error("virtual fetch failed", "virtual", virtualName, "path", path, "error", err)
http.Error(w, "bad gateway", http.StatusBadGateway)
+21 -5
View File
@@ -1,6 +1,8 @@
package v2
import (
"context"
"errors"
"fmt"
"net/http"
"strconv"
@@ -8,8 +10,14 @@ import (
"github.com/go-chi/chi/v5"
"git.unkin.net/unkin/artifactapi/internal/database"
"git.unkin.net/unkin/artifactapi/internal/proxy"
)
// Evictor drops a remote path from every cache layer.
type Evictor interface {
Evict(ctx context.Context, remoteName, path string) error
}
type ObjectsHandler struct {
db *database.DB
}
@@ -18,10 +26,13 @@ func NewObjectsHandler(db *database.DB) *ObjectsHandler {
return &ObjectsHandler{db: db}
}
func (h *ObjectsHandler) Routes() chi.Router {
// Routes lists and evicts objects for remote repos; evictor serves the DELETE.
func (h *ObjectsHandler) Routes(evictor Evictor) chi.Router {
r := chi.NewRouter()
r.Get("/", h.list)
r.Delete("/*", h.evict)
r.Delete("/*", func(w http.ResponseWriter, r *http.Request) {
evict(w, r, evictor)
})
return r
}
@@ -51,7 +62,7 @@ func (h *ObjectsHandler) list(w http.ResponseWriter, r *http.Request) {
remoteName := chi.URLParam(r, "name")
limit, offset := pageBounds(r)
artifacts, err := h.db.ListArtifacts(r.Context(), remoteName, limit, offset)
artifacts, err := h.db.ListArtifacts(r.Context(), remoteName, r.URL.Query().Get("prefix"), limit, offset)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
@@ -82,11 +93,16 @@ func (h *ObjectsHandler) evictLocal(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusNoContent)
}
func (h *ObjectsHandler) evict(w http.ResponseWriter, r *http.Request) {
func evict(w http.ResponseWriter, r *http.Request, evictor Evictor) {
remoteName := chi.URLParam(r, "name")
path := chi.URLParam(r, "*")
if err := h.db.DeleteArtifact(r.Context(), remoteName, path); err != nil {
if err := evictor.Evict(r.Context(), remoteName, path); err != nil {
var proxyErr *proxy.ProxyError
if errors.As(err, &proxyErr) {
http.Error(w, proxyErr.Message, proxyErr.Status)
return
}
http.Error(w, fmt.Sprintf("evict failed: %v", err), http.StatusInternalServerError)
return
}
+59
View File
@@ -0,0 +1,59 @@
package v2
import (
"context"
"errors"
"fmt"
"net/http"
"net/http/httptest"
"testing"
"github.com/go-chi/chi/v5"
"git.unkin.net/unkin/artifactapi/internal/proxy"
)
type fakeEvictor struct {
remote, path string
err error
}
func (f *fakeEvictor) Evict(_ context.Context, remote, path string) error {
f.remote, f.path = remote, path
return f.err
}
func deleteObject(ev Evictor, path string) int {
router := chi.NewRouter()
router.Route("/remotes/{name}/objects", func(r chi.Router) {
r.Delete("/*", NewObjectsHandler(nil).Routes(ev).ServeHTTP)
})
w := httptest.NewRecorder()
router.ServeHTTP(w, httptest.NewRequest("DELETE", "/remotes/epel/objects/"+path, nil))
return w.Code
}
func TestRemoteEvictDelegatesToEvictor(t *testing.T) {
for _, path := range []string{"8/Everything/x86_64/repodata/repomd.xml", "8/Everything/x86_64/repodata/*"} {
ev := &fakeEvictor{}
if code := deleteObject(ev, path); code != 204 || ev.remote != "epel" || ev.path != path {
t.Errorf("DELETE %s: code=%d evicted=%q/%q", path, code, ev.remote, ev.path)
}
}
}
func TestRemoteEvictMapsErrorStatus(t *testing.T) {
for name, tc := range map[string]struct {
err error
want int
}{
"bad wildcard": {&proxy.ProxyError{Status: http.StatusBadRequest, Message: "wildcard evict must be <dir>/*"}, http.StatusBadRequest},
"unknown remote": {&proxy.ProxyError{Status: http.StatusNotFound, Message: "remote not found"}, http.StatusNotFound},
"lock busy": {fmt.Errorf("wrapped: %w", &proxy.ProxyError{Status: http.StatusServiceUnavailable, Message: "retry"}), http.StatusServiceUnavailable},
"other error": {errors.New("db down"), http.StatusInternalServerError},
} {
if code := deleteObject(&fakeEvictor{err: tc.err}, "x"); code != tc.want {
t.Errorf("%s: code = %d, want %d", name, code, tc.want)
}
}
}
+68
View File
@@ -0,0 +1,68 @@
package v2
import (
"context"
"encoding/json"
"net/http/httptest"
"testing"
"github.com/go-chi/chi/v5"
"git.unkin.net/unkin/artifactapi/internal/database"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// TestObjectsListPrefix verifies the remote objects listing passes ?prefix=
// through to the database filter.
func TestObjectsListPrefix(t *testing.T) {
if testDSN == "" {
t.Skip("Docker unavailable")
}
ctx := context.Background()
db, err := database.New(testDSN)
if err != nil {
t.Fatal(err)
}
defer db.Close()
const remote = "generic-objs-prefix"
if err := db.CreateRemote(ctx, &models.Remote{
Name: remote, PackageType: models.PackageGeneric, RepoType: models.RepoTypeRemote,
BaseURL: "https://example.com", MutableTTL: 3600,
}); err != nil {
t.Fatal(err)
}
const hash = "sha256:bb22"
if err := db.UpsertBlob(ctx, hash, "blobs/bb/22", 10, "text/plain"); err != nil {
t.Fatal(err)
}
for _, p := range []string{"a/one.txt", "b/two.txt"} {
if err := db.UpsertArtifact(ctx, remote, p, hash, ""); err != nil {
t.Fatal(err)
}
}
router := chi.NewRouter()
router.Mount("/remotes/{name}/objects", NewObjectsHandler(db).Routes(&fakeEvictor{}))
list := func(query string) []models.Artifact {
t.Helper()
w := httptest.NewRecorder()
router.ServeHTTP(w, httptest.NewRequest("GET", "/remotes/"+remote+"/objects"+query, nil))
if w.Code != 200 {
t.Fatalf("list%s = %d, want 200", query, w.Code)
}
var got []models.Artifact
if err := json.Unmarshal(w.Body.Bytes(), &got); err != nil {
t.Fatalf("decode: %v", err)
}
return got
}
if got := list(""); len(got) != 2 {
t.Fatalf("unfiltered listing returned %d objects, want 2", len(got))
}
if got := list("?prefix=b/"); len(got) != 1 || got[0].Path != "b/two.txt" {
t.Fatalf("prefix=b/ listing = %+v, want only b/two.txt", got)
}
}
+69
View File
@@ -131,3 +131,72 @@ func TestFlushRemote(t *testing.T) {
t.Error("expected keys flushed")
}
}
func setPathKeys(t *testing.T, remote string, paths ...string) {
t.Helper()
ctx := context.Background()
for _, p := range paths {
if err := testRedis.SetTTL(ctx, remote, p, time.Minute); err != nil {
t.Fatal(err)
}
if err := testRedis.SetETag(ctx, remote, p, `"e"`, time.Minute); err != nil {
t.Fatal(err)
}
}
}
func pathKeysExist(t *testing.T, remote, path string) (ttl, etag bool) {
t.Helper()
ctx := context.Background()
ttl, _ = testRedis.CheckTTL(ctx, remote, path)
e, _ := testRedis.GetETag(ctx, remote, path)
return ttl, e != ""
}
func TestForgetPath(t *testing.T) {
requireRedis(t)
const meta = `repo/a*b?[c]\d.xml`
setPathKeys(t, "fp", meta, "repo/aXb.xml")
setPathKeys(t, "fp-other", meta)
if err := testRedis.ForgetPath(context.Background(), "fp", meta); err != nil {
t.Fatal(err)
}
if ttl, etag := pathKeysExist(t, "fp", meta); ttl || etag {
t.Errorf("forgotten path keys remain: ttl=%v etag=%v", ttl, etag)
}
for _, k := range [][2]string{{"fp", "repo/aXb.xml"}, {"fp-other", meta}} {
if ttl, etag := pathKeysExist(t, k[0], k[1]); !ttl || !etag {
t.Errorf("%s:%s lost keys: ttl=%v etag=%v", k[0], k[1], ttl, etag)
}
}
}
func TestForgetPrefix(t *testing.T) {
requireRedis(t)
const prefix = `r*[1]?\/`
under := []string{prefix + "repomd.xml", prefix + "sub/x.rpm"}
// Each would match the prefix if its glob metacharacters were left unescaped.
globMatches := []string{`rX11/repomd.xml`, `r[1]?\/x`, `r*1Z/x`}
setPathKeys(t, "fx", append(under, globMatches...)...)
setPathKeys(t, "fx-other", under...)
if err := testRedis.ForgetPrefix(context.Background(), "fx", prefix); err != nil {
t.Fatal(err)
}
for _, p := range under {
if ttl, etag := pathKeysExist(t, "fx", p); ttl || etag {
t.Errorf("%s keys remain: ttl=%v etag=%v", p, ttl, etag)
}
}
for _, p := range globMatches {
if ttl, etag := pathKeysExist(t, "fx", p); !ttl || !etag {
t.Errorf("%s outside prefix lost keys: ttl=%v etag=%v", p, ttl, etag)
}
}
for _, p := range under {
if ttl, etag := pathKeysExist(t, "fx-other", p); !ttl || !etag {
t.Errorf("other remote %s lost keys: ttl=%v etag=%v", p, ttl, etag)
}
}
}
+25
View File
@@ -3,6 +3,7 @@ package cache
import (
"context"
"fmt"
"strings"
"time"
"github.com/redis/go-redis/v9"
@@ -115,3 +116,27 @@ func (r *Redis) FlushRemote(ctx context.Context, remote string) error {
}
return iter.Err()
}
// ForgetPath drops the freshness and ETag keys of one cached path.
func (r *Redis) ForgetPath(ctx context.Context, remote, path string) error {
return r.client.Del(ctx, fmt.Sprintf("ttl:%s:%s", remote, path), fmt.Sprintf("etag:%s:%s", remote, path)).Err()
}
// ForgetPrefix drops the freshness and ETag keys of every cached path under prefix.
func (r *Redis) ForgetPrefix(ctx context.Context, remote, prefix string) error {
glob := globEscaper.Replace(remote + ":" + prefix)
for _, kind := range []string{"ttl:", "etag:"} {
iter := r.client.Scan(ctx, 0, kind+glob+"*", 100).Iterator()
for iter.Next(ctx) {
if err := r.client.Del(ctx, iter.Val()).Err(); err != nil {
return err
}
}
if err := iter.Err(); err != nil {
return err
}
}
return nil
}
var globEscaper = strings.NewReplacer(`\`, `\\`, `*`, `\*`, `?`, `\?`, `[`, `\[`, `]`, `\]`)
+8 -38
View File
@@ -2,11 +2,9 @@ package database
import (
"context"
"errors"
"time"
"github.com/jackc/pgx/v5"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
@@ -30,42 +28,14 @@ func (db *DB) ListGitHubAlpineRemotes(ctx context.Context) ([]models.Remote, err
return remotes, rows.Err()
}
// ClaimGitHubAlpineSyncLease atomically claims the per-remote sync lease. It
// succeeds only when the remote is due (never synced, or synced longer than
// freshness ago) and no live lease is held by another replica. A zero freshness
// (prime scans) ignores the recency gate. The returned etag is the stored
// releases-list ETag, shared across replicas.
// ClaimGitHubAlpineSyncLease atomically claims the per-remote github_alpine sync
// lease. See claimSyncLease.
func (db *DB) ClaimGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (bool, string, error) {
row := db.Pool.QueryRow(ctx, `
INSERT INTO github_alpine_sync_state AS s (remote_name, sync_lease_owner, sync_lease_expires)
VALUES ($1, $2, now() + make_interval(secs => $4))
ON CONFLICT (remote_name) DO UPDATE
SET sync_lease_owner = $2,
sync_lease_expires = now() + make_interval(secs => $4)
WHERE (s.last_synced_at IS NULL OR s.last_synced_at < now() - make_interval(secs => $3))
AND (s.sync_lease_expires IS NULL OR s.sync_lease_expires < now())
RETURNING s.etag
`, remoteName, owner, freshness.Seconds(), lease.Seconds())
var etag string
if err := row.Scan(&etag); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return false, "", nil
}
return false, "", err
}
return true, etag, nil
return db.claimSyncLease(ctx, "github_alpine_sync_state", remoteName, owner, freshness, lease)
}
// ReleaseGitHubAlpineSyncLease records the completed scan and frees the lease.
// Only the owning replica may release; last_synced_at advances so the next poll
// waits a full freshness window, and etag is persisted for the next conditional
// request.
func (db *DB) ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error {
_, err := db.Pool.Exec(ctx, `
UPDATE github_alpine_sync_state
SET last_synced_at = $3, etag = $4, sync_lease_owner = '', sync_lease_expires = NULL
WHERE remote_name = $1 AND sync_lease_owner = $2
`, remoteName, owner, syncedAt, etag)
return err
// ReleaseGitHubAlpineSyncLease records a github_alpine scan outcome and frees the
// lease. See releaseSyncLease.
func (db *DB) ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error {
return db.releaseSyncLease(ctx, "github_alpine_sync_state", remoteName, owner, res)
}
+12 -4
View File
@@ -65,7 +65,8 @@ func (db *DB) TouchArtifactAccess(ctx context.Context, remoteName, path string)
return err
}
func (db *DB) ListArtifacts(ctx context.Context, remoteName string, limit, offset int) ([]models.Artifact, error) {
// ListArtifacts pages a remote's artifacts whose path starts with prefix ("" lists all).
func (db *DB) ListArtifacts(ctx context.Context, remoteName, prefix string, limit, offset int) ([]models.Artifact, error) {
rows, err := db.Pool.Query(ctx, `
SELECT a.id, a.remote_name, a.path, a.content_hash, a.upstream_etag,
a.upstream_last_modified, a.first_seen_at, a.last_fetched_at,
@@ -73,10 +74,10 @@ func (db *DB) ListArtifacts(ctx context.Context, remoteName string, limit, offse
b.size_bytes, b.content_type
FROM artifacts a
JOIN blobs b ON a.content_hash = b.content_hash
WHERE a.remote_name = $1
WHERE a.remote_name = $1 AND left(a.path, length($2)) = $2
ORDER BY a.path
LIMIT $2 OFFSET $3
`, remoteName, limit, offset)
LIMIT $3 OFFSET $4
`, remoteName, prefix, limit, offset)
if err != nil {
return nil, err
}
@@ -103,6 +104,13 @@ func (db *DB) DeleteArtifact(ctx context.Context, remoteName, path string) error
return err
}
// DeleteArtifactsByPrefix removes every artifact row of a remote whose path
// starts with prefix.
func (db *DB) DeleteArtifactsByPrefix(ctx context.Context, remoteName, prefix string) error {
_, err := db.Pool.Exec(ctx, `DELETE FROM artifacts WHERE remote_name = $1 AND left(path, length($2)) = $2`, remoteName, prefix)
return err
}
func (db *DB) InsertAccessLog(ctx context.Context, remoteName, path string, cacheHit bool, sizeBytes int64, upstreamMS int, clientIP string) error {
_, err := db.Pool.Exec(ctx, `
INSERT INTO access_log (remote_name, path, cache_hit, size_bytes, upstream_ms, client_ip)
+49 -3
View File
@@ -168,10 +168,17 @@ func TestArtifactsAndBlobs(t *testing.T) {
if err := testDB.TouchArtifactAccess(ctx(), "r-art", "path/a.txt"); err != nil {
t.Fatal(err)
}
arts, err := testDB.ListArtifacts(ctx(), "r-art", 10, 0)
if err != nil || len(arts) != 1 {
if err := testDB.UpsertArtifact(ctx(), "r-art", "other/b.txt", hash, ""); err != nil {
t.Fatal(err)
}
arts, err := testDB.ListArtifacts(ctx(), "r-art", "", 10, 0)
if err != nil || len(arts) != 2 {
t.Fatalf("list artifacts: %v %v", len(arts), err)
}
arts, err = testDB.ListArtifacts(ctx(), "r-art", "path/", 10, 0)
if err != nil || len(arts) != 1 || arts[0].Path != "path/a.txt" {
t.Fatalf("list artifacts with prefix: %+v %v", arts, err)
}
if err := testDB.InsertAccessLog(ctx(), "r-art", "path/a.txt", true, 10, 5, "1.2.3.4"); err != nil {
t.Fatal(err)
}
@@ -188,6 +195,45 @@ func TestArtifactsAndBlobs(t *testing.T) {
}
}
func TestListArtifactsPrefix(t *testing.T) {
requireDB(t)
seedRemote(t, "r-prefix")
seedBlob(t, "prefixhash")
for _, p := range []string{"a%b/y", "a1b/y", "a_b/x", "aXb/x", "pkg/1", "pkg/2", "pkg/3", "pkgx/4"} {
if err := testDB.UpsertArtifact(ctx(), "r-prefix", p, "sha256:prefixhash", ""); err != nil {
t.Fatal(err)
}
}
paths := func(prefix string, limit, offset int) []string {
t.Helper()
arts, err := testDB.ListArtifacts(ctx(), "r-prefix", prefix, limit, offset)
if err != nil {
t.Fatalf("list %q: %v", prefix, err)
}
out := make([]string, len(arts))
for i, a := range arts {
out[i] = a.Path
}
return out
}
// LIKE wildcards in the prefix must match literally.
if got := paths("a_b/", 10, 0); len(got) != 1 || got[0] != "a_b/x" {
t.Fatalf("prefix a_b/ = %v, want [a_b/x]", got)
}
if got := paths("a%", 10, 0); len(got) != 1 || got[0] != "a%b/y" {
t.Fatalf("prefix a%% = %v, want [a%%b/y]", got)
}
// limit/offset page the filtered set, not the whole remote.
if got := paths("pkg/", 2, 0); len(got) != 2 || got[0] != "pkg/1" || got[1] != "pkg/2" {
t.Fatalf("prefix pkg/ page 1 = %v, want [pkg/1 pkg/2]", got)
}
if got := paths("pkg/", 2, 2); len(got) != 1 || got[0] != "pkg/3" {
t.Fatalf("prefix pkg/ page 2 = %v, want [pkg/3]", got)
}
}
func TestOrphanAndColdCleanup(t *testing.T) {
requireDB(t)
seedBlob(t, "orphanhash")
@@ -320,7 +366,7 @@ func TestDatabaseErrorPaths(t *testing.T) {
if _, err := bad.ListVirtuals(ctx); err == nil {
t.Error("ListVirtuals should error")
}
if _, err := bad.ListArtifacts(ctx, "r", 10, 0); err == nil {
if _, err := bad.ListArtifacts(ctx, "r", "", 10, 0); err == nil {
t.Error("ListArtifacts should error")
}
if _, err := bad.ListLocalFiles(ctx, "r", 10, 0); err == nil {
+8 -37
View File
@@ -2,11 +2,9 @@ package database
import (
"context"
"errors"
"time"
"github.com/jackc/pgx/v5"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
@@ -30,41 +28,14 @@ func (db *DB) ListGitHubDebRemotes(ctx context.Context) ([]models.Remote, error)
return remotes, rows.Err()
}
// ClaimGitHubDebSyncLease atomically claims the per-remote sync lease. It
// succeeds only when the remote is due (never synced, or synced longer than
// freshness ago) and no live lease is held by another replica. A zero freshness
// (prime scans) ignores the recency gate. The returned etag is the stored
// releases-list ETag, shared across replicas.
// ClaimGitHubDebSyncLease atomically claims the per-remote github_deb sync
// lease. See claimSyncLease.
func (db *DB) ClaimGitHubDebSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (bool, string, error) {
row := db.Pool.QueryRow(ctx, `
INSERT INTO github_deb_sync_state AS s (remote_name, sync_lease_owner, sync_lease_expires)
VALUES ($1, $2, now() + make_interval(secs => $4))
ON CONFLICT (remote_name) DO UPDATE
SET sync_lease_owner = $2,
sync_lease_expires = now() + make_interval(secs => $4)
WHERE (s.last_synced_at IS NULL OR s.last_synced_at < now() - make_interval(secs => $3))
AND (s.sync_lease_expires IS NULL OR s.sync_lease_expires < now())
RETURNING s.etag
`, remoteName, owner, freshness.Seconds(), lease.Seconds())
var etag string
if err := row.Scan(&etag); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return false, "", nil
}
return false, "", err
}
return true, etag, nil
return db.claimSyncLease(ctx, "github_deb_sync_state", remoteName, owner, freshness, lease)
}
// ReleaseGitHubDebSyncLease records the completed scan and frees the lease. Only
// the owning replica may release; last_synced_at advances so the next poll waits
// a full freshness window, and etag is persisted for the next conditional request.
func (db *DB) ReleaseGitHubDebSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error {
_, err := db.Pool.Exec(ctx, `
UPDATE github_deb_sync_state
SET last_synced_at = $3, etag = $4, sync_lease_owner = '', sync_lease_expires = NULL
WHERE remote_name = $1 AND sync_lease_owner = $2
`, remoteName, owner, syncedAt, etag)
return err
// ReleaseGitHubDebSyncLease records a github_deb scan outcome and frees the
// lease. See releaseSyncLease.
func (db *DB) ReleaseGitHubDebSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error {
return db.releaseSyncLease(ctx, "github_deb_sync_state", remoteName, owner, res)
}
+2 -1
View File
@@ -4,6 +4,7 @@ import (
"testing"
"time"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
@@ -44,7 +45,7 @@ func TestGitHubDebSyncLease(t *testing.T) {
t.Fatal("replica-2 claimed while replica-1 holds the lease")
}
if err := testDB.ReleaseGitHubDebSyncLease(ctx(), name, "replica-1", `"etag-1"`, time.Now()); err != nil {
if err := testDB.ReleaseGitHubDebSyncLease(ctx(), name, "replica-1", provider.SyncResult{Etag: `"etag-1"`}); err != nil {
t.Fatalf("release: %v", err)
}
+45 -17
View File
@@ -7,6 +7,8 @@ import (
"github.com/jackc/pgx/v5"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
@@ -30,22 +32,35 @@ func (db *DB) ListGitHubRPMRemotes(ctx context.Context) ([]models.Remote, error)
return remotes, rows.Err()
}
// ClaimGitHubSyncLease atomically claims the per-remote sync lease. It succeeds
// (claimed=true) only when the remote is due — never synced, or synced longer
// than freshness ago — and no live lease is held by another replica. This bounds
// total GitHub load to roughly one scan per freshness window regardless of how
// many replicas poll. The returned etag is the stored releases-list ETag, shared
// across replicas so a conditional request can short-circuit an unchanged repo.
// A zero freshness (used for prime scans) ignores the recency gate and claims
// whenever no live lease is held.
// ClaimGitHubSyncLease atomically claims the per-remote github_rpm sync lease.
// See claimSyncLease.
func (db *DB) ClaimGitHubSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (bool, string, error) {
return db.claimSyncLease(ctx, "github_rpm_sync_state", remoteName, owner, freshness, lease)
}
// ReleaseGitHubSyncLease records a github_rpm scan outcome and frees the lease.
// See releaseSyncLease.
func (db *DB) ReleaseGitHubSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error {
return db.releaseSyncLease(ctx, "github_rpm_sync_state", remoteName, owner, res)
}
// claimSyncLease atomically claims a per-remote sync lease in table. It succeeds
// (claimed=true) only when the remote is due and no live lease is held by
// another replica. Due means: a pending retry after a failed scan has reached
// next_retry_at, or, with no retry pending, the remote was never synced or was
// synced longer than freshness ago. A zero freshness (prime scans) ignores the
// recency gate but still honours a pending retry's backoff. The returned etag is
// the stored releases-list ETag, shared across replicas so a conditional request
// can short-circuit an unchanged repo.
func (db *DB) claimSyncLease(ctx context.Context, table, remoteName, owner string, freshness, lease time.Duration) (bool, string, error) {
row := db.Pool.QueryRow(ctx, `
INSERT INTO github_rpm_sync_state AS s (remote_name, sync_lease_owner, sync_lease_expires)
INSERT INTO `+table+` AS s (remote_name, sync_lease_owner, sync_lease_expires)
VALUES ($1, $2, now() + make_interval(secs => $4))
ON CONFLICT (remote_name) DO UPDATE
SET sync_lease_owner = $2,
sync_lease_expires = now() + make_interval(secs => $4)
WHERE (s.last_synced_at IS NULL OR s.last_synced_at < now() - make_interval(secs => $3))
WHERE (CASE WHEN s.next_retry_at IS NOT NULL THEN s.next_retry_at <= now()
ELSE s.last_synced_at IS NULL OR s.last_synced_at < now() - make_interval(secs => $3) END)
AND (s.sync_lease_expires IS NULL OR s.sync_lease_expires < now())
RETURNING s.etag
`, remoteName, owner, freshness.Seconds(), lease.Seconds())
@@ -60,14 +75,27 @@ func (db *DB) ClaimGitHubSyncLease(ctx context.Context, remoteName, owner string
return true, etag, nil
}
// ReleaseGitHubSyncLease records the completed scan and frees the lease. Only the
// owning replica may release; last_synced_at advances so the next poll waits a
// full freshness window, and etag is persisted for the next conditional request.
func (db *DB) ReleaseGitHubSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error {
// releaseSyncLease records a scan outcome and frees the lease; only the owning
// replica may release. Success advances last_synced_at, persists the etag and
// clears any retry. Failure leaves last_synced_at and etag untouched (the last
// good metadata keeps serving) and schedules next_retry_at with exponential
// backoff from the consecutive-failure count, never before res.RetryAt.
func (db *DB) releaseSyncLease(ctx context.Context, table, remoteName, owner string, res provider.SyncResult) error {
var retryAt *time.Time
if !res.RetryAt.IsZero() {
retryAt = &res.RetryAt
}
_, err := db.Pool.Exec(ctx, `
UPDATE github_rpm_sync_state
SET last_synced_at = $3, etag = $4, sync_lease_owner = '', sync_lease_expires = NULL
UPDATE `+table+`
SET last_synced_at = CASE WHEN $3 THEN last_synced_at ELSE now() END,
etag = CASE WHEN $3 THEN etag ELSE $4 END,
sync_failures = CASE WHEN $3 THEN sync_failures + 1 ELSE 0 END,
next_retry_at = CASE WHEN $3 THEN GREATEST(
now() + make_interval(secs => LEAST($5 * power(2, LEAST(sync_failures, 20)), $6)),
$7::timestamptz)
END,
sync_lease_owner = '', sync_lease_expires = NULL
WHERE remote_name = $1 AND sync_lease_owner = $2
`, remoteName, owner, syncedAt, etag)
`, remoteName, owner, res.Failed, res.Etag, res.Backoff.Seconds(), res.MaxBackoff.Seconds(), retryAt)
return err
}
+115 -1
View File
@@ -4,6 +4,7 @@ import (
"testing"
"time"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
@@ -47,7 +48,7 @@ func TestGitHubSyncLease(t *testing.T) {
}
// Replica 1 finishes: record the sync and persist an etag.
if err := testDB.ReleaseGitHubSyncLease(ctx(), name, "replica-1", `"etag-1"`, time.Now()); err != nil {
if err := testDB.ReleaseGitHubSyncLease(ctx(), name, "replica-1", provider.SyncResult{Etag: `"etag-1"`}); err != nil {
t.Fatalf("release: %v", err)
}
@@ -93,3 +94,116 @@ func TestListGitHubRPMRemotes(t *testing.T) {
t.Fatalf("seeded remote %q not returned", name)
}
}
// A failed scan must not count as a sync: the etag and last_synced_at stay put,
// every claim (prime included) is held off until the backoff elapses, the delay
// doubles per consecutive failure up to the cap, an upstream retry hint pushes
// it later, and a success clears the retry state.
func TestGitHubSyncLeaseFailureBackoff(t *testing.T) {
requireDB(t)
name := "gh-retry-" + time.Now().Format("150405.000000")
seedGitHubRPMRemote(t, name)
const lease = 15 * time.Minute
freshness := time.Hour
claim := func(f time.Duration) (bool, string) {
t.Helper()
ok, etag, err := testDB.ClaimGitHubSyncLease(ctx(), name, "r1", f, lease)
if err != nil {
t.Fatalf("claim: %v", err)
}
return ok, etag
}
release := func(res provider.SyncResult) {
t.Helper()
if err := testDB.ReleaseGitHubSyncLease(ctx(), name, "r1", res); err != nil {
t.Fatalf("release: %v", err)
}
}
state := func() (failures int, retryIn time.Duration, synced *time.Time, etag string) {
t.Helper()
var next *time.Time
var dbNow time.Time
if err := testDB.Pool.QueryRow(ctx(), `SELECT sync_failures, next_retry_at, last_synced_at, etag, now() FROM github_rpm_sync_state WHERE remote_name = $1`, name).
Scan(&failures, &next, &synced, &etag, &dbNow); err != nil {
t.Fatalf("read state: %v", err)
}
if next != nil {
retryIn = next.Sub(dbNow)
}
return
}
elapse := func() {
t.Helper()
if _, err := testDB.Pool.Exec(ctx(), `UPDATE github_rpm_sync_state SET next_retry_at = now() - interval '1 second' WHERE remote_name = $1`, name); err != nil {
t.Fatalf("elapse backoff: %v", err)
}
}
failed := provider.SyncResult{Failed: true, Etag: `"ignored"`, Backoff: time.Minute, MaxBackoff: 3 * time.Minute}
// Last good sync.
if ok, _ := claim(freshness); !ok {
t.Fatal("first claim failed")
}
release(provider.SyncResult{Etag: `"good"`})
_, _, goodSynced, _ := state()
// First failure: 60s backoff, sync time and etag untouched.
if ok, _ := claim(0); !ok {
t.Fatal("prime claim failed")
}
release(failed)
n, in, synced, etag := state()
if n != 1 || in < 55*time.Second || in > 65*time.Second {
t.Fatalf("after 1st failure: failures=%d retry in %v, want 1 and ~60s", n, in)
}
if etag != `"good"` || synced == nil || !synced.Equal(*goodSynced) {
t.Fatalf("failed scan overwrote sync state: etag=%q synced=%v", etag, synced)
}
if ok, _ := claim(0); ok {
t.Fatal("prime claimed inside the retry backoff")
}
// Backoff elapsed: claimable despite being inside the freshness window,
// and the second failure doubles the delay.
elapse()
if ok, _ := claim(freshness); !ok {
t.Fatal("periodic claim refused after backoff elapsed")
}
release(failed)
if n, in, _, _ := state(); n != 2 || in < 115*time.Second || in > 125*time.Second {
t.Fatalf("after 2nd failure: failures=%d retry in %v, want 2 and ~120s", n, in)
}
// Third failure caps at MaxBackoff (3m, not 4m).
elapse()
claim(freshness)
release(failed)
if n, in, _, _ := state(); n != 3 || in < 175*time.Second || in > 185*time.Second {
t.Fatalf("after 3rd failure: failures=%d retry in %v, want 3 and ~180s", n, in)
}
// An upstream rate-limit reset later than the backoff wins.
elapse()
claim(freshness)
hinted := failed
hinted.RetryAt = time.Now().Add(20 * time.Minute)
release(hinted)
if _, in, _, _ := state(); in < 19*time.Minute || in > 21*time.Minute {
t.Fatalf("retry hint ignored: retry in %v, want ~20m", in)
}
// Recovery: success persists the new etag and clears the retry, so the
// normal freshness window applies again.
elapse()
if ok, etag := claim(freshness); !ok || etag != `"good"` {
t.Fatalf("recovery claim ok=%v etag=%q, want last good etag", ok, etag)
}
release(provider.SyncResult{Etag: `"new"`})
if n, in, _, etag := state(); n != 0 || in != 0 || etag != `"new"` {
t.Fatalf("after success: failures=%d retry in %v etag=%q", n, in, etag)
}
if ok, _ := claim(freshness); ok {
t.Fatal("claimed inside freshness window after a successful sync")
}
}
+1 -1
View File
@@ -470,7 +470,7 @@ func (p *GitHubProvider) fetchReleases(ctx context.Context, remote models.Remote
return nil, "", false, err
}
if resp.StatusCode != http.StatusOK {
return nil, "", false, fmt.Errorf("github releases API %s: status %d", u, resp.StatusCode)
return nil, "", false, provider.NewUpstreamStatusError("github releases API "+u, resp)
}
if page == 1 {
newEtag = respEtag
+7 -9
View File
@@ -28,7 +28,7 @@ type SyncStore interface {
provider.RemoteMetadataStore
ListGitHubAlpineRemotes(ctx context.Context) ([]models.Remote, error)
ClaimGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (claimed bool, etag string, err error)
ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error
ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error
}
// SyncConfig tunes the shared syncer. Zero values fall back to safe defaults.
@@ -189,10 +189,11 @@ func (s *Syncer) process(ctx context.Context, job syncJob) {
s.mu.Unlock()
}()
freshness := time.Duration(job.remote.MutableTTL) * time.Second
if freshness <= 0 {
freshness = defaultSyncFreshness
ttl := time.Duration(job.remote.MutableTTL) * time.Second
if ttl <= 0 {
ttl = defaultSyncFreshness
}
freshness := ttl
if job.prime {
freshness = 0
}
@@ -210,16 +211,13 @@ func (s *Syncer) process(ctx context.Context, job syncJob) {
defer cancel()
newEtag, changed, scanErr := s.prov.scanWithState(scanCtx, job.remote, s.store, etag)
releaseEtag := etag
if scanErr == nil {
releaseEtag = newEtag
} else {
if scanErr != nil {
slog.Error("github_alpine syncer: scan failed", "remote", job.remote.Name, "error", scanErr)
}
relCtx, relCancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
defer relCancel()
if err := s.store.ReleaseGitHubAlpineSyncLease(relCtx, job.remote.Name, s.owner, releaseEtag, time.Now()); err != nil {
if err := s.store.ReleaseGitHubAlpineSyncLease(relCtx, job.remote.Name, s.owner, provider.NewSyncResult(newEtag, scanErr, ttl)); err != nil {
slog.Warn("github_alpine syncer: release lease", "remote", job.remote.Name, "error", err)
}
+13 -3
View File
@@ -27,6 +27,7 @@ type fakeSyncStore struct {
leaseExp map[string]time.Time
lastSynced map[string]time.Time
etags map[string]string
retryAt map[string]time.Time
}
func newFakeSyncStore() *fakeSyncStore {
@@ -36,6 +37,7 @@ func newFakeSyncStore() *fakeSyncStore {
leaseExp: map[string]time.Time{},
lastSynced: map[string]time.Time{},
etags: map[string]string{},
retryAt: map[string]time.Time{},
}
}
@@ -52,6 +54,9 @@ func (f *fakeSyncStore) ClaimGitHubAlpineSyncLease(_ context.Context, name, owne
ls, hasLS := f.lastSynced[name]
exp, hasExp := f.leaseExp[name]
freshOK := !hasLS || now.Sub(ls) >= freshness
if ra, pending := f.retryAt[name]; pending {
freshOK = !now.Before(ra)
}
leaseOK := !hasExp || exp.Before(now)
if freshOK && leaseOK {
f.leaseOwner[name] = owner
@@ -61,14 +66,19 @@ func (f *fakeSyncStore) ClaimGitHubAlpineSyncLease(_ context.Context, name, owne
return false, "", nil
}
func (f *fakeSyncStore) ReleaseGitHubAlpineSyncLease(_ context.Context, name, owner, etag string, syncedAt time.Time) error {
func (f *fakeSyncStore) ReleaseGitHubAlpineSyncLease(_ context.Context, name, owner string, res provider.SyncResult) error {
f.mu.Lock()
defer f.mu.Unlock()
if f.leaseOwner[name] != owner {
return nil
}
f.lastSynced[name] = syncedAt
f.etags[name] = etag
if res.Failed {
f.retryAt[name] = time.Now().Add(res.Backoff)
} else {
f.lastSynced[name] = time.Now()
f.etags[name] = res.Etag
delete(f.retryAt, name)
}
delete(f.leaseOwner, name)
delete(f.leaseExp, name)
return nil
+92
View File
@@ -0,0 +1,92 @@
// Package cargo proxies a Cargo sparse registry (RFC 2789): mutable index
// files plus immutable .crate downloads, with config.json synthesized so
// cargo fetches crates back through artifactapi.
package cargo
import (
"context"
"encoding/json"
"net/http"
"net/url"
"strings"
"git.unkin.net/unkin/artifactapi/internal/auth"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
func init() {
provider.Register(&Provider{})
}
const (
crateDir = "crates/"
cratesIOIndex = "index.crates.io"
cratesIODownload = "https://static.crates.io"
)
type Provider struct{}
func (p *Provider) Type() models.PackageType { return models.PackageCargo }
func (p *Provider) Classify(path string) provider.Mutability {
if isCrate(path) {
return provider.Immutable
}
return provider.Mutable
}
func (p *Provider) ContentType(path string) string {
if isCrate(path) {
return "application/gzip"
}
if path == "config.json" {
return "application/json"
}
return "text/plain"
}
// UpstreamURL maps crates/{crate}/{crate}-{version}.crate to the download host
// and everything else to the index. crates.io splits the two across hosts;
// any other upstream is expected to serve both.
func (p *Provider) UpstreamURL(remote models.Remote, path string) string {
path = strings.TrimLeft(path, "/")
base := strings.TrimRight(remote.BaseURL, "/")
if isCrate(path) {
if u, err := url.Parse(base); err == nil && u.Host == cratesIOIndex {
base = cratesIODownload
}
}
return base + "/" + path
}
func (p *Provider) RewriteResponse(_ []byte, _ models.Remote, _ string) ([]byte, error) {
return nil, nil
}
func (p *Provider) AuthHeaders(_ context.Context, remote models.Remote) (http.Header, error) {
return auth.BasicHeaders(remote), nil
}
// ServeRemote answers config.json itself: dl must name the requesting host,
// which a cached upstream body cannot.
func (p *Provider) ServeRemote(w http.ResponseWriter, r *http.Request, remote models.Remote, path, proxyBaseURL string, _ provider.RemoteMetadataStore) bool {
if strings.TrimLeft(path, "/") != "config.json" {
return false
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(Config(remote, proxyBaseURL))
return true
}
// Config builds the sparse registry config.json pointing downloads at proxyBaseURL.
func Config(remote models.Remote, proxyBaseURL string) map[string]string {
return map[string]string{
"dl": strings.TrimRight(proxyBaseURL, "/") + "/api/v1/remote/" + remote.Name + "/" + crateDir + "{crate}/{crate}-{version}.crate",
}
}
func isCrate(path string) bool {
path = strings.TrimLeft(path, "/")
return strings.HasPrefix(path, crateDir) && strings.HasSuffix(path, ".crate")
}
+75
View File
@@ -0,0 +1,75 @@
package cargo
import (
"encoding/json"
"net/http/httptest"
"testing"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
func TestClassify(t *testing.T) {
p := &Provider{}
tests := []struct {
path string
want provider.Mutability
}{
{"crates/serde/serde-1.0.200.crate", provider.Immutable},
{"se/rd/serde", provider.Mutable},
{"1/a", provider.Mutable},
{"3/a/anyhow", provider.Mutable},
{"config.json", provider.Mutable},
{"cr/at/crates", provider.Mutable},
}
for _, tt := range tests {
if got := p.Classify(tt.path); got != tt.want {
t.Errorf("Classify(%q) = %v, want %v", tt.path, got, tt.want)
}
}
}
func TestUpstreamURL(t *testing.T) {
p := &Provider{}
cratesIO := models.Remote{BaseURL: "https://index.crates.io/"}
mirror := models.Remote{BaseURL: "http://mirror.example"}
tests := []struct {
remote models.Remote
path, want string
}{
{cratesIO, "se/rd/serde", "https://index.crates.io/se/rd/serde"},
{cratesIO, "crates/serde/serde-1.0.200.crate", "https://static.crates.io/crates/serde/serde-1.0.200.crate"},
{mirror, "crates/serde/serde-1.0.200.crate", "http://mirror.example/crates/serde/serde-1.0.200.crate"},
{mirror, "/2/ab", "http://mirror.example/2/ab"},
}
for _, tt := range tests {
if got := p.UpstreamURL(tt.remote, tt.path); got != tt.want {
t.Errorf("UpstreamURL(%q, %q) = %q, want %q", tt.remote.BaseURL, tt.path, got, tt.want)
}
}
}
func TestServeRemoteConfig(t *testing.T) {
p := &Provider{}
remote := models.Remote{Name: "crates-io", BaseURL: "https://index.crates.io"}
w := httptest.NewRecorder()
if !p.ServeRemote(w, httptest.NewRequest("GET", "/", nil), remote, "config.json", "https://aa.example/", nil) {
t.Fatal("config.json not served")
}
var cfg map[string]string
if err := json.Unmarshal(w.Body.Bytes(), &cfg); err != nil {
t.Fatal(err)
}
want := "https://aa.example/api/v1/remote/crates-io/crates/{crate}/{crate}-{version}.crate"
if cfg["dl"] != want {
t.Errorf("dl = %q, want %q", cfg["dl"], want)
}
if _, ok := cfg["api"]; ok {
t.Error("api must be omitted: publishing through the proxy is unsupported")
}
if p.ServeRemote(httptest.NewRecorder(), httptest.NewRequest("GET", "/", nil), remote, "se/rd/serde", "https://aa.example", nil) {
t.Error("index files must fall through to the proxy engine")
}
}
+1 -1
View File
@@ -431,7 +431,7 @@ func (p *GitHubProvider) fetchReleases(ctx context.Context, remote models.Remote
return nil, "", false, err
}
if resp.StatusCode != http.StatusOK {
return nil, "", false, fmt.Errorf("github releases API %s: status %d", u, resp.StatusCode)
return nil, "", false, provider.NewUpstreamStatusError("github releases API "+u, resp)
}
if page == 1 {
newEtag = respEtag
+7 -9
View File
@@ -28,7 +28,7 @@ type SyncStore interface {
provider.RemoteMetadataStore
ListGitHubDebRemotes(ctx context.Context) ([]models.Remote, error)
ClaimGitHubDebSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (claimed bool, etag string, err error)
ReleaseGitHubDebSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error
ReleaseGitHubDebSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error
}
// SyncConfig tunes the shared syncer. Zero values fall back to safe defaults.
@@ -189,10 +189,11 @@ func (s *Syncer) process(ctx context.Context, job syncJob) {
s.mu.Unlock()
}()
freshness := time.Duration(job.remote.MutableTTL) * time.Second
if freshness <= 0 {
freshness = defaultSyncFreshness
ttl := time.Duration(job.remote.MutableTTL) * time.Second
if ttl <= 0 {
ttl = defaultSyncFreshness
}
freshness := ttl
if job.prime {
freshness = 0
}
@@ -210,16 +211,13 @@ func (s *Syncer) process(ctx context.Context, job syncJob) {
defer cancel()
newEtag, changed, scanErr := s.prov.scanWithState(scanCtx, job.remote, s.store, etag)
releaseEtag := etag
if scanErr == nil {
releaseEtag = newEtag
} else {
if scanErr != nil {
slog.Error("github_deb syncer: scan failed", "remote", job.remote.Name, "error", scanErr)
}
relCtx, relCancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
defer relCancel()
if err := s.store.ReleaseGitHubDebSyncLease(relCtx, job.remote.Name, s.owner, releaseEtag, time.Now()); err != nil {
if err := s.store.ReleaseGitHubDebSyncLease(relCtx, job.remote.Name, s.owner, provider.NewSyncResult(newEtag, scanErr, ttl)); err != nil {
slog.Warn("github_deb syncer: release lease", "remote", job.remote.Name, "error", err)
}
+13 -3
View File
@@ -27,6 +27,7 @@ type fakeSyncStore struct {
leaseExp map[string]time.Time
lastSynced map[string]time.Time
etags map[string]string
retryAt map[string]time.Time
}
func newFakeSyncStore() *fakeSyncStore {
@@ -36,6 +37,7 @@ func newFakeSyncStore() *fakeSyncStore {
leaseExp: map[string]time.Time{},
lastSynced: map[string]time.Time{},
etags: map[string]string{},
retryAt: map[string]time.Time{},
}
}
@@ -52,6 +54,9 @@ func (f *fakeSyncStore) ClaimGitHubDebSyncLease(_ context.Context, name, owner s
ls, hasLS := f.lastSynced[name]
exp, hasExp := f.leaseExp[name]
freshOK := !hasLS || now.Sub(ls) >= freshness
if ra, pending := f.retryAt[name]; pending {
freshOK = !now.Before(ra)
}
leaseOK := !hasExp || exp.Before(now)
if freshOK && leaseOK {
f.leaseOwner[name] = owner
@@ -61,14 +66,19 @@ func (f *fakeSyncStore) ClaimGitHubDebSyncLease(_ context.Context, name, owner s
return false, "", nil
}
func (f *fakeSyncStore) ReleaseGitHubDebSyncLease(_ context.Context, name, owner, etag string, syncedAt time.Time) error {
func (f *fakeSyncStore) ReleaseGitHubDebSyncLease(_ context.Context, name, owner string, res provider.SyncResult) error {
f.mu.Lock()
defer f.mu.Unlock()
if f.leaseOwner[name] != owner {
return nil
}
f.lastSynced[name] = syncedAt
f.etags[name] = etag
if res.Failed {
f.retryAt[name] = time.Now().Add(res.Backoff)
} else {
f.lastSynced[name] = time.Now()
f.etags[name] = res.Etag
delete(f.retryAt, name)
}
delete(f.leaseOwner, name)
delete(f.leaseExp, name)
return nil
+1 -1
View File
@@ -447,7 +447,7 @@ func (p *GitHubProvider) fetchReleases(ctx context.Context, remote models.Remote
return nil, "", false, err
}
if resp.StatusCode != http.StatusOK {
return nil, "", false, fmt.Errorf("github releases API %s: status %d", u, resp.StatusCode)
return nil, "", false, provider.NewUpstreamStatusError("github releases API "+u, resp)
}
if page == 1 {
newEtag = respEtag
+9
View File
@@ -78,6 +78,7 @@ type githubFixture struct {
notModHit int // releases-list requests answered 304
releaseAuth string // Authorization header seen on the last releases request
assetAuth string // Authorization header seen on the last asset request
failStatus int // when set, the releases list answers this status (rate-limit style)
mu sync.Mutex
}
@@ -100,6 +101,14 @@ func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture {
f.mu.Lock()
f.releasesHit++
f.releaseAuth = r.Header.Get("Authorization")
if f.failStatus != 0 {
status := f.failStatus
f.mu.Unlock()
w.Header().Set("X-RateLimit-Remaining", "0")
w.Header().Set("X-RateLimit-Reset", strconv.FormatInt(time.Now().Add(30*time.Second).Unix(), 10))
http.Error(w, "API rate limit exceeded", status)
return
}
etag := f.etag
if etag != "" && r.Header.Get("If-None-Match") == etag {
f.notModHit++
+8 -10
View File
@@ -35,7 +35,7 @@ type SyncStore interface {
provider.RemoteMetadataStore
ListGitHubRPMRemotes(ctx context.Context) ([]models.Remote, error)
ClaimGitHubSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (claimed bool, etag string, err error)
ReleaseGitHubSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error
ReleaseGitHubSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error
}
// SyncConfig tunes the shared syncer. Zero values fall back to safe defaults.
@@ -205,10 +205,11 @@ func (s *Syncer) process(ctx context.Context, job syncJob) {
s.mu.Unlock()
}()
freshness := time.Duration(job.remote.MutableTTL) * time.Second
if freshness <= 0 {
freshness = defaultSyncFreshness
ttl := time.Duration(job.remote.MutableTTL) * time.Second
if ttl <= 0 {
ttl = defaultSyncFreshness
}
freshness := ttl
if job.prime {
freshness = 0 // prime ignores the recency gate but still respects a live lease
}
@@ -226,18 +227,15 @@ func (s *Syncer) process(ctx context.Context, job syncJob) {
defer cancel()
newEtag, changed, scanErr := s.prov.scanWithState(scanCtx, job.remote, s.store, etag)
releaseEtag := etag
if scanErr == nil {
releaseEtag = newEtag
} else {
if scanErr != nil {
slog.Error("github_rpm syncer: scan failed", "remote", job.remote.Name, "error", scanErr)
}
// Release on a detached context so a clean shutdown mid-scan still frees the
// lease and advances last_synced_at (otherwise it simply expires).
// lease and records the outcome (otherwise the lease simply expires).
relCtx, relCancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
defer relCancel()
if err := s.store.ReleaseGitHubSyncLease(relCtx, job.remote.Name, s.owner, releaseEtag, time.Now()); err != nil {
if err := s.store.ReleaseGitHubSyncLease(relCtx, job.remote.Name, s.owner, provider.NewSyncResult(newEtag, scanErr, ttl)); err != nil {
slog.Warn("github_rpm syncer: release lease", "remote", job.remote.Name, "error", err)
}
+92 -3
View File
@@ -27,6 +27,8 @@ type fakeSyncStore struct {
leaseExp map[string]time.Time
lastSynced map[string]time.Time
etags map[string]string
retryAt map[string]time.Time
lastResult provider.SyncResult
}
func newFakeSyncStore() *fakeSyncStore {
@@ -36,6 +38,7 @@ func newFakeSyncStore() *fakeSyncStore {
leaseExp: map[string]time.Time{},
lastSynced: map[string]time.Time{},
etags: map[string]string{},
retryAt: map[string]time.Time{},
}
}
@@ -52,6 +55,9 @@ func (f *fakeSyncStore) ClaimGitHubSyncLease(_ context.Context, name, owner stri
ls, hasLS := f.lastSynced[name]
exp, hasExp := f.leaseExp[name]
freshOK := !hasLS || now.Sub(ls) >= freshness
if ra, pending := f.retryAt[name]; pending {
freshOK = !now.Before(ra)
}
leaseOK := !hasExp || exp.Before(now)
if freshOK && leaseOK {
f.leaseOwner[name] = owner
@@ -61,14 +67,20 @@ func (f *fakeSyncStore) ClaimGitHubSyncLease(_ context.Context, name, owner stri
return false, "", nil
}
func (f *fakeSyncStore) ReleaseGitHubSyncLease(_ context.Context, name, owner, etag string, syncedAt time.Time) error {
func (f *fakeSyncStore) ReleaseGitHubSyncLease(_ context.Context, name, owner string, res provider.SyncResult) error {
f.mu.Lock()
defer f.mu.Unlock()
if f.leaseOwner[name] != owner {
return nil
}
f.lastSynced[name] = syncedAt
f.etags[name] = etag
f.lastResult = res
if res.Failed {
f.retryAt[name] = time.Now().Add(res.Backoff)
} else {
f.lastSynced[name] = time.Now()
f.etags[name] = res.Etag
delete(f.retryAt, name)
}
delete(f.leaseOwner, name)
delete(f.leaseExp, name)
return nil
@@ -310,3 +322,80 @@ func TestSyncerPrimeBypassesRecencyPeriodicDoesNot(t *testing.T) {
t.Fatalf("periodic scan ran inside recency window: %d -> %d releases calls", releasesAfterPrime, fx.releasesHit)
}
}
// A rate-limited scan is not recorded as a sync: the last good repodata keeps
// serving, the retry honours the upstream reset hint and backoff, and once the
// upstream recovers the next poll after the backoff re-syncs without waiting
// out mutable_ttl.
func TestSyncerFailedScanRetriesAfterBackoff(t *testing.T) {
fx := newGitHubFixture(t, true)
fx.etag = `"v1"`
store := newFakeSyncStore()
p := newTestProvider()
s := newSyncer(store, p, testSyncConfig())
remote := fx.remote()
bg := context.Background()
s.process(bg, syncJob{remote: remote, prime: true})
if rows, _ := store.ListRPMMetadataEntries(bg, remote.Name); len(rows) != 1 {
t.Fatalf("prime did not derive: %d rows", len(rows))
}
// mutable_ttl elapses while GitHub is rate-limiting.
fx.mu.Lock()
fx.failStatus = http.StatusForbidden
fx.etag = `"v2"`
fx.mu.Unlock()
store.mu.Lock()
store.lastSynced[remote.Name] = time.Now().Add(-2 * time.Hour)
store.mu.Unlock()
s.process(bg, syncJob{remote: remote})
res := store.lastResult
if !res.Failed {
t.Fatal("403 scan was recorded as a successful sync")
}
if until := time.Until(res.RetryAt); until < 20*time.Second || until > 40*time.Second {
t.Fatalf("retry hint from X-RateLimit-Reset not carried: retry in %v", until)
}
if res.Backoff != time.Minute || res.MaxBackoff >= time.Duration(remote.MutableTTL)*time.Second {
t.Fatalf("backoff %v cap %v, want 1m first retry capped below mutable_ttl", res.Backoff, res.MaxBackoff)
}
if store.etags[remote.Name] != `"v1"` {
t.Fatalf("failed scan replaced the etag: %q", store.etags[remote.Name])
}
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-rpm/repodata/repomd.xml", nil)
p.ServeRemote(rec, req, remote, "repodata/repomd.xml", "https://x", store)
if rec.Code != http.StatusOK {
t.Fatalf("last good repodata not served during failure: %d %s", rec.Code, rec.Body.String())
}
if rows, _ := store.ListRPMMetadataEntries(bg, remote.Name); len(rows) != 1 {
t.Fatalf("failed scan dropped cached metadata: %d rows", len(rows))
}
// Inside the backoff nothing is retried.
hits := fx.releasesHit
s.process(bg, syncJob{remote: remote})
if fx.releasesHit != hits {
t.Fatal("scan retried inside the backoff window")
}
// Upstream recovers and the backoff elapses: the next poll re-syncs.
fx.mu.Lock()
fx.failStatus = 0
fx.rpmBytes["other-9-9.aarch64.rpm"] = testsupport.MinimalRPM("other", "9", "9", "aarch64")
fx.mu.Unlock()
store.mu.Lock()
store.retryAt[remote.Name] = time.Now().Add(-time.Second)
store.mu.Unlock()
s.process(bg, syncJob{remote: remote})
if store.lastResult.Failed || store.etags[remote.Name] != `"v2"` {
t.Fatalf("recovery scan not recorded: failed=%v etag=%q", store.lastResult.Failed, store.etags[remote.Name])
}
if rows, _ := store.ListRPMMetadataEntries(bg, remote.Name); len(rows) != 2 {
t.Fatalf("recovery did not refresh metadata: %d rows", len(rows))
}
}
+75
View File
@@ -0,0 +1,75 @@
package provider
import (
"errors"
"fmt"
"net/http"
"strconv"
"time"
)
const (
syncRetryBase = time.Minute
syncRetryMax = 10 * time.Minute
)
// UpstreamStatusError is a non-success upstream response. RetryAt is the
// upstream's own retry hint (Retry-After, or X-RateLimit-Reset once the quota
// is exhausted); zero when it gave none.
type UpstreamStatusError struct {
URL string
Status int
RetryAt time.Time
}
func (e *UpstreamStatusError) Error() string {
return fmt.Sprintf("%s: status %d", e.URL, e.Status)
}
// NewUpstreamStatusError wraps a non-success response, capturing its retry hint.
func NewUpstreamStatusError(url string, resp *http.Response) *UpstreamStatusError {
e := &UpstreamStatusError{URL: url, Status: resp.StatusCode}
if ra := resp.Header.Get("Retry-After"); ra != "" {
if secs, err := strconv.Atoi(ra); err == nil {
e.RetryAt = time.Now().Add(time.Duration(secs) * time.Second)
} else if t, err := http.ParseTime(ra); err == nil {
e.RetryAt = t
}
} else if resp.Header.Get("X-RateLimit-Remaining") == "0" {
if reset, err := strconv.ParseInt(resp.Header.Get("X-RateLimit-Reset"), 10, 64); err == nil {
e.RetryAt = time.Unix(reset, 0)
}
}
return e
}
// SyncResult is a background scan's outcome, recorded when its sync lease is
// released. A failed scan keeps the prior sync time and ETag and schedules a
// retry after Backoff, doubled per consecutive failure up to MaxBackoff, and
// never earlier than RetryAt.
type SyncResult struct {
Etag string
Failed bool
RetryAt time.Time
Backoff time.Duration
MaxBackoff time.Duration
}
// NewSyncResult builds the result for a scan against a remote with the given
// mutable_ttl. The retry cap stays well below ttl, and an upstream hint is
// clamped to ttl so a bogus reset can never stall the remote longer than a
// normal sync interval would.
func NewSyncResult(etag string, scanErr error, ttl time.Duration) SyncResult {
if scanErr == nil {
return SyncResult{Etag: etag}
}
res := SyncResult{Failed: true, Backoff: syncRetryBase, MaxBackoff: min(syncRetryMax, max(syncRetryBase, ttl/4))}
var se *UpstreamStatusError
if errors.As(scanErr, &se) && !se.RetryAt.IsZero() {
res.RetryAt = se.RetryAt
if limit := time.Now().Add(ttl); res.RetryAt.After(limit) {
res.RetryAt = limit
}
}
return res
}
+63
View File
@@ -0,0 +1,63 @@
package provider
import (
"errors"
"fmt"
"net/http"
"strconv"
"testing"
"time"
)
func TestNewUpstreamStatusErrorRetryHint(t *testing.T) {
reset := time.Now().Add(15 * time.Minute).Truncate(time.Second)
cases := []struct {
name string
header map[string]string
want time.Duration // 0 = no hint
}{
{"retry-after seconds", map[string]string{"Retry-After": "90"}, 90 * time.Second},
{"retry-after date", map[string]string{"Retry-After": reset.UTC().Format(http.TimeFormat)}, 15 * time.Minute},
{"quota exhausted", map[string]string{"X-RateLimit-Remaining": "0", "X-RateLimit-Reset": strconv.FormatInt(reset.Unix(), 10)}, 15 * time.Minute},
{"quota left is not a hint", map[string]string{"X-RateLimit-Remaining": "12", "X-RateLimit-Reset": strconv.FormatInt(reset.Unix(), 10)}, 0},
{"no headers", nil, 0},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
resp := &http.Response{StatusCode: http.StatusForbidden, Header: http.Header{}}
for k, v := range c.header {
resp.Header.Set(k, v)
}
e := NewUpstreamStatusError("u", resp)
if c.want == 0 {
if !e.RetryAt.IsZero() {
t.Fatalf("unexpected hint %v", e.RetryAt)
}
return
}
if d := time.Until(e.RetryAt) - c.want; d < -2*time.Second || d > 2*time.Second {
t.Fatalf("retry in %v, want ~%v", time.Until(e.RetryAt), c.want)
}
})
}
}
func TestNewSyncResult(t *testing.T) {
if r := NewSyncResult(`"e"`, nil, time.Hour); r.Failed || r.Etag != `"e"` {
t.Fatalf("success result = %+v", r)
}
r := NewSyncResult(`"e"`, errors.New("boom"), time.Hour)
if !r.Failed || r.Backoff != time.Minute || r.MaxBackoff != 10*time.Minute || !r.RetryAt.IsZero() {
t.Fatalf("plain failure = %+v, want 1m backoff capped at 10m, no hint", r)
}
if r := NewSyncResult("", errors.New("boom"), 5*time.Minute); r.MaxBackoff != time.Minute+15*time.Second {
t.Fatalf("short ttl cap = %v, want ttl/4", r.MaxBackoff)
}
far := &UpstreamStatusError{Status: 403, RetryAt: time.Now().Add(3 * time.Hour)}
r = NewSyncResult("", fmt.Errorf("scan: %w", far), time.Hour)
if until := time.Until(r.RetryAt); until > time.Hour || until < 59*time.Minute {
t.Fatalf("wrapped hint not clamped to ttl: retry in %v", until)
}
}
+71
View File
@@ -16,6 +16,8 @@ import (
"sync/atomic"
"time"
"github.com/jackc/pgx/v5"
"git.unkin.net/unkin/artifactapi/internal/cache"
"git.unkin.net/unkin/artifactapi/internal/database"
"git.unkin.net/unkin/artifactapi/internal/provider"
@@ -47,6 +49,8 @@ type Engine struct {
// mirror strategy to prefer the mirror currently handling the fewest
// requests. Per-replica and approximate, which is fine.
inflight sync.Map
// evictLockWait bounds how long Evict waits on a held fetch lock.
evictLockWait time.Duration
}
func NewEngine(db *database.DB, c *cache.Redis, s *storage.S3) *Engine {
@@ -57,6 +61,8 @@ func NewEngine(db *database.DB, c *cache.Redis, s *storage.S3) *Engine {
cas: storage.NewCAS(s),
circuit: NewCircuitBreaker(c),
accessLog: make(chan database.AccessLogEntry, accessLogBufferSize),
evictLockWait: fetchLockTTL,
}
go e.runAccessLogWriter()
return e
@@ -208,6 +214,71 @@ func (e *Engine) Fetch(ctx context.Context, remote models.Remote, path string, p
return result, nil
}
// Evict drops path from every cache layer (artifact row, index object, Redis
// freshness and ETag keys) so the next request refetches from upstream. A
// trailing "/*" evicts every path under that directory.
func (e *Engine) Evict(ctx context.Context, remoteName, path string) error {
prefix, wildcard := strings.CutSuffix(path, "*")
if wildcard && !strings.HasSuffix(prefix, "/") {
return &ProxyError{Status: http.StatusBadRequest, Message: "wildcard evict must be <dir>/*"}
}
if _, err := e.db.GetRemote(ctx, remoteName); errors.Is(err, pgx.ErrNoRows) {
return &ProxyError{Status: http.StatusNotFound, Message: fmt.Sprintf("remote %q not found", remoteName)}
} else if err != nil {
return fmt.Errorf("get remote: %w", err)
}
if !wildcard {
if err := e.waitForLock(ctx, remoteName, path); err != nil {
return err
}
defer func() { _ = e.cache.ReleaseLock(context.WithoutCancel(ctx), remoteName, path) }()
if err := e.db.DeleteArtifact(ctx, remoteName, path); err != nil {
return fmt.Errorf("delete artifact: %w", err)
}
if err := e.store.Delete(ctx, storage.IndexKey(remoteName, path)); err != nil {
return fmt.Errorf("delete index: %w", err)
}
return e.cache.ForgetPath(ctx, remoteName, path)
}
// ponytail: no lock for wildcards; a Fetch already in flight under the
// prefix can re-cache its path after the evict. Per-path locks over a
// directory would close it if that ever matters.
if err := e.db.DeleteArtifactsByPrefix(ctx, remoteName, prefix); err != nil {
return fmt.Errorf("delete artifacts: %w", err)
}
if err := e.store.DeletePrefix(ctx, storage.IndexKey(remoteName, prefix)); err != nil {
return fmt.Errorf("delete indexes: %w", err)
}
return e.cache.ForgetPrefix(ctx, remoteName, prefix)
}
// waitForLock takes the per-path fetch lock so an in-flight Fetch cannot
// re-set TTL/ETag keys after an evict. It fails with a 503 when the lock
// cannot be taken within evictLockWait or Redis errors.
func (e *Engine) waitForLock(ctx context.Context, remoteName, path string) error {
deadline := time.Now().Add(e.evictLockWait)
for {
ok, err := e.cache.AcquireLock(ctx, remoteName, path, fetchLockTTL)
if ok {
return nil
}
if ctx.Err() != nil {
return ctx.Err()
}
if err != nil {
return &ProxyError{Status: http.StatusServiceUnavailable, Message: fmt.Sprintf("fetch lock: %v", err)}
}
if time.Now().After(deadline) {
return &ProxyError{Status: http.StatusServiceUnavailable, Message: "fetch in progress, retry evict"}
}
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(50 * time.Millisecond):
}
}
}
// HeadResult carries artifact metadata for a HEAD request. There is no body.
type HeadResult struct {
ContentType string
+271
View File
@@ -0,0 +1,271 @@
package proxy
import (
"context"
"errors"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"time"
_ "git.unkin.net/unkin/artifactapi/internal/provider/rpm"
"git.unkin.net/unkin/artifactapi/internal/storage"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// changingUpstream serves every path with the current revision, as a mirror
// does after a sync replaces its repodata. It answers 304 to a matching
// If-None-Match and counts conditional requests.
func changingUpstream(t *testing.T) (*httptest.Server, *atomic.Value, *atomic.Int32) {
t.Helper()
var rev atomic.Value
var conditional atomic.Int32
rev.Store("rev1")
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
v := rev.Load().(string)
etag := `"` + v + `"`
w.Header().Set("ETag", etag)
if inm := r.Header.Get("If-None-Match"); inm != "" {
conditional.Add(1)
if inm == etag {
w.WriteHeader(http.StatusNotModified)
return
}
}
_, _ = w.Write([]byte(v + ":" + r.URL.Path))
}))
t.Cleanup(srv.Close)
return srv, &rev, &conditional
}
func fetchBody(t *testing.T, r models.Remote, path string) string {
t.Helper()
res, err := testEngine.Fetch(context.Background(), r, path, prov(t, models.PackageRPM))
if err != nil {
t.Fatalf("fetch %s: %v", path, err)
}
return readAll(t, res)
}
func rpmRemote(t *testing.T, name, baseURL string) models.Remote {
return seed(t, models.Remote{Name: name, PackageType: models.PackageRPM, RepoType: models.RepoTypeRemote, BaseURL: baseURL, MutableTTL: 7200, CheckMutable: true})
}
// cached reports which cache layers hold path: artifact row, index object,
// Redis TTL key, Redis ETag key.
type cached struct{ row, index, ttl, etag bool }
// Mutable indexes live in the S3 index; immutable blobs get an artifact row.
var (
indexCached = cached{index: true, ttl: true, etag: true}
blobCached = cached{row: true, ttl: true, etag: true}
)
func layers(t *testing.T, remote, path string) cached {
t.Helper()
ctx := context.Background()
var c cached
_, err := testDB.GetArtifact(ctx, remote, path)
c.row = err == nil
c.index, err = testEngine.store.Exists(ctx, storage.IndexKey(remote, path))
if err != nil {
t.Fatalf("stat index %s: %v", path, err)
}
c.ttl, _ = testCache.CheckTTL(ctx, remote, path)
etag, _ := testCache.GetETag(ctx, remote, path)
c.etag = etag != ""
return c
}
func TestEvictMutableIndexRefetches(t *testing.T) {
requireStack(t)
srv, rev, conditional := changingUpstream(t)
r := rpmRemote(t, "evict-idx", srv.URL)
const path = "8/Everything/x86_64/repodata/repomd.xml"
if got := fetchBody(t, r, path); got != "rev1:/"+path {
t.Fatalf("initial fetch = %q", got)
}
if c := layers(t, r.Name, path); c != indexCached {
t.Fatalf("before evict = %+v, want %+v", c, indexCached)
}
rev.Store("rev2")
if got := fetchBody(t, r, path); got != "rev1:/"+path {
t.Fatalf("within TTL = %q, want cached rev1", got)
}
if err := testEngine.Evict(context.Background(), r.Name, path); err != nil {
t.Fatalf("evict: %v", err)
}
if c := layers(t, r.Name, path); c != (cached{}) {
t.Fatalf("after evict = %+v, want every layer gone", c)
}
conditional.Store(0)
if got := fetchBody(t, r, path); got != "rev2:/"+path {
t.Fatalf("after evict = %q, want rev2", got)
}
if n := conditional.Load(); n != 0 {
t.Errorf("post-evict fetch revalidated with the evicted ETag (%d conditional requests)", n)
}
}
func TestEvictImmutableBlobDropsRow(t *testing.T) {
requireStack(t)
srv, _, _ := changingUpstream(t)
r := rpmRemote(t, "evict-blob", srv.URL)
const path = "8/Everything/x86_64/Packages/a/a-1.0-1.el8.x86_64.rpm"
fetchBody(t, r, path)
if c := layers(t, r.Name, path); c != blobCached {
t.Fatalf("before evict = %+v, want %+v", c, blobCached)
}
if err := testEngine.Evict(context.Background(), r.Name, path); err != nil {
t.Fatalf("evict: %v", err)
}
if c := layers(t, r.Name, path); c != (cached{}) {
t.Fatalf("after evict = %+v, want every layer gone", c)
}
}
func TestEvictRejectsNonDirectoryWildcard(t *testing.T) {
requireStack(t)
srv, _, _ := changingUpstream(t)
r := rpmRemote(t, "evict-bare", srv.URL)
const path = "8/Everything/x86_64/repodata/repomd.xml"
fetchBody(t, r, path)
for _, bad := range []string{"*", "8*", "8/Every*"} {
var pe *ProxyError
if err := testEngine.Evict(context.Background(), r.Name, bad); !errors.As(err, &pe) || pe.Status != http.StatusBadRequest {
t.Errorf("evict %s = %v, want 400", bad, err)
}
}
if c := layers(t, r.Name, path); c != (indexCached) {
t.Errorf("rejected wildcard evicted layers: %+v", c)
}
}
func TestEvictUnknownRemote(t *testing.T) {
requireStack(t)
var pe *ProxyError
if err := testEngine.Evict(context.Background(), "evict-no-such-remote", "a/b"); !errors.As(err, &pe) || pe.Status != http.StatusNotFound {
t.Fatalf("evict unknown remote = %v, want 404", err)
}
}
func holdLock(t *testing.T, remote, path string) {
t.Helper()
ctx := context.Background()
if ok, err := testCache.AcquireLock(ctx, remote, path, time.Minute); !ok || err != nil {
t.Fatalf("acquire: %v %v", ok, err)
}
t.Cleanup(func() { _ = testCache.ReleaseLock(ctx, remote, path) })
}
func TestEvictWaitsForFetchLock(t *testing.T) {
requireStack(t)
srv, _, _ := changingUpstream(t)
r := rpmRemote(t, "evict-lock", srv.URL)
ctx := context.Background()
const path = "8/Everything/x86_64/repodata/repomd.xml"
holdLock(t, r.Name, path)
done := make(chan error, 1)
go func() { done <- testEngine.Evict(ctx, r.Name, path) }()
select {
case err := <-done:
t.Fatalf("evict returned while fetch lock held: %v", err)
case <-time.After(200 * time.Millisecond):
}
_ = testCache.ReleaseLock(ctx, r.Name, path)
if err := <-done; err != nil {
t.Fatalf("evict: %v", err)
}
ok, err := testCache.AcquireLock(ctx, r.Name, path, time.Second)
if !ok || err != nil {
t.Fatalf("lock still held after evict: %v %v", ok, err)
}
}
func TestEvictCancelledWhileWaitingDeletesNothing(t *testing.T) {
requireStack(t)
srv, _, _ := changingUpstream(t)
r := rpmRemote(t, "evict-cancel", srv.URL)
const path = "8/Everything/x86_64/repodata/repomd.xml"
fetchBody(t, r, path)
holdLock(t, r.Name, path)
ctx, cancel := context.WithCancel(context.Background())
done := make(chan error, 1)
go func() { done <- testEngine.Evict(ctx, r.Name, path) }()
time.Sleep(100 * time.Millisecond)
cancel()
select {
case err := <-done:
if err == nil {
t.Fatal("cancelled evict returned nil")
}
case <-time.After(2 * time.Second):
t.Fatal("cancelled evict did not return")
}
if c := layers(t, r.Name, path); c != (indexCached) {
t.Errorf("cancelled evict deleted layers: %+v", c)
}
}
func TestEvictLockTimeoutIs503(t *testing.T) {
requireStack(t)
srv, _, _ := changingUpstream(t)
r := rpmRemote(t, "evict-timeout", srv.URL)
const path = "8/Everything/x86_64/repodata/repomd.xml"
fetchBody(t, r, path)
holdLock(t, r.Name, path)
prev := testEngine.evictLockWait
testEngine.evictLockWait = 100 * time.Millisecond
t.Cleanup(func() { testEngine.evictLockWait = prev })
var pe *ProxyError
if err := testEngine.Evict(context.Background(), r.Name, path); !errors.As(err, &pe) || pe.Status != http.StatusServiceUnavailable {
t.Fatalf("evict = %v, want 503", err)
}
if c := layers(t, r.Name, path); c != (indexCached) {
t.Errorf("timed-out evict deleted layers: %+v", c)
}
}
func TestEvictWildcardClearsPrefixOnly(t *testing.T) {
requireStack(t)
srv, rev, _ := changingUpstream(t)
r := rpmRemote(t, "evict-wild", srv.URL)
const (
repomd = "8/Everything/x86_64/repodata/repomd.xml"
rpm = "8/Everything/x86_64/Packages/a/a-1.0-1.el8.x86_64.rpm"
other = "9/Everything/x86_64/repodata/repomd.xml"
)
for _, p := range []string{repomd, rpm, other} {
fetchBody(t, r, p)
}
if c := layers(t, r.Name, rpm); c != blobCached {
t.Fatalf("%s before evict = %+v, want %+v", rpm, c, blobCached)
}
rev.Store("rev2")
if err := testEngine.Evict(context.Background(), r.Name, "8/Everything/x86_64/*"); err != nil {
t.Fatalf("evict: %v", err)
}
for _, p := range []string{repomd, rpm} {
if c := layers(t, r.Name, p); c != (cached{}) {
t.Errorf("%s after evict = %+v, want every layer gone", p, c)
}
}
if c := layers(t, r.Name, other); c != (indexCached) {
t.Errorf("%s outside prefix = %+v, want %+v", other, c, indexCached)
}
if got := fetchBody(t, r, repomd); got != "rev2:/"+repomd {
t.Errorf("index under prefix = %q, want rev2", got)
}
if got := fetchBody(t, r, rpm); got != "rev2:/"+rpm {
t.Errorf("artifact under prefix = %q, want rev2", got)
}
if got := fetchBody(t, r, other); got != "rev1:/"+other {
t.Errorf("path outside prefix = %q, want cached rev1", got)
}
}
+3 -2
View File
@@ -21,6 +21,7 @@ import (
"git.unkin.net/unkin/artifactapi/internal/gc"
"git.unkin.net/unkin/artifactapi/internal/githubauth"
"git.unkin.net/unkin/artifactapi/internal/provider/alpine"
_ "git.unkin.net/unkin/artifactapi/internal/provider/cargo"
"git.unkin.net/unkin/artifactapi/internal/provider/deb"
_ "git.unkin.net/unkin/artifactapi/internal/provider/docker"
_ "git.unkin.net/unkin/artifactapi/internal/provider/generic"
@@ -196,8 +197,8 @@ func (s *Server) routes() chi.Router {
r.Route("/remotes/{name}/objects", func(r chi.Router) {
objHandler := v2.NewObjectsHandler(s.db)
r.Get("/", objHandler.Routes().ServeHTTP)
r.Delete("/*", objHandler.Routes().ServeHTTP)
r.Get("/", objHandler.Routes(s.engine).ServeHTTP)
r.Delete("/*", objHandler.Routes(s.engine).ServeHTTP)
})
r.Route("/locals/{name}/objects", func(r chi.Router) {
+24
View File
@@ -99,6 +99,30 @@ func (s *S3) Stat(ctx context.Context, key string) (*minio.ObjectInfo, error) {
return &info, nil
}
// DeletePrefix removes every object whose key starts with prefix. It lists
// first so a failed list returns its error instead of deleting nothing.
func (s *S3) DeletePrefix(ctx context.Context, prefix string) error {
var objs []minio.ObjectInfo
for obj := range s.client.ListObjects(ctx, s.bucket, minio.ListObjectsOptions{Prefix: prefix, Recursive: true}) {
if obj.Err != nil {
return obj.Err
}
objs = append(objs, obj)
}
ch := make(chan minio.ObjectInfo, len(objs))
for _, obj := range objs {
ch <- obj
}
close(ch)
var err error
for res := range s.client.RemoveObjects(ctx, s.bucket, ch, minio.RemoveObjectsOptions{}) {
if res.Err != nil && err == nil {
err = res.Err
}
}
return err
}
// ListStaleObjects returns keys under prefix last modified before cutoff. Used
// by the GC to reap abandoned staging objects (e.g. cancelled docker pushes).
func (s *S3) ListStaleObjects(ctx context.Context, prefix string, cutoff time.Time) ([]string, error) {
+27
View File
@@ -4,11 +4,15 @@ import (
"bytes"
"context"
"io"
"net/http"
"net/http/httptest"
"os"
"strings"
"testing"
"time"
"github.com/minio/minio-go/v7"
"git.unkin.net/unkin/artifactapi/internal/testsupport"
)
@@ -158,3 +162,26 @@ func TestCASStore(t *testing.T) {
t.Errorf("stored content mismatch: %q", got)
}
}
// The fake endpoint fails every list but accepts every delete, as an S3 that
// tolerates deleting an empty key would.
func TestDeletePrefixReturnsListError(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodPost {
_, _ = io.WriteString(w, `<DeleteResult></DeleteResult>`)
return
}
w.WriteHeader(http.StatusInternalServerError)
_, _ = io.WriteString(w, `<Error><Code>InternalError</Code><Message>list failed</Message></Error>`)
}))
defer srv.Close()
client, err := minio.New(strings.TrimPrefix(srv.URL, "http://"), &minio.Options{Region: "us-east-1", MaxRetries: 1})
if err != nil {
t.Fatal(err)
}
s := &S3{client: client, bucket: "bucket"}
err = s.DeletePrefix(context.Background(), "indexes/r/")
if err == nil {
t.Fatal("DeletePrefix on a failing list returned nil")
}
}
+262 -1
View File
@@ -2,27 +2,58 @@ package virtual
import (
"context"
"errors"
"fmt"
"io"
"log/slog"
"net/http"
"net/http/httptest"
"net/url"
"slices"
"strings"
"sync"
"time"
"git.unkin.net/unkin/artifactapi/internal/database"
"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 {
db *database.DB
proxyEngine *proxy.Engine
getRemote func(context.Context, string) (*models.Remote, error)
rpmMember func(context.Context, string) (*RPMMember, error)
mergeTTL time.Duration
sf singleflight.Group
mu sync.Mutex
rpm map[string]*rpmGen
}
// 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
}
const rpmMergeTTL = 60 * time.Second
func NewEngine(db *database.DB, proxyEngine *proxy.Engine) *Engine {
return &Engine{db: db, proxyEngine: proxyEngine}
e := &Engine{db: db, proxyEngine: proxyEngine, getRemote: db.GetRemote, mergeTTL: rpmMergeTTL}
e.rpmMember = e.fetchRPMMember
return e
}
func (e *Engine) Fetch(ctx context.Context, virt models.Virtual, path string, proxyBaseURL string) ([]byte, string, error) {
if virt.PackageType == models.PackageRPM {
return e.fetchRPM(ctx, virt, path)
}
merger, err := GetMerger(virt.PackageType)
if err != nil {
return nil, "", fmt.Errorf("unsupported virtual type %q: %w", virt.PackageType, err)
@@ -133,3 +164,233 @@ func (e *Engine) fetchLocalIndex(ctx context.Context, remote models.Remote, path
return indexer.GenerateLocalIndex(ctx, e.db, remote.Name, path)
}
var (
ErrNotFound = errors.New("not found")
ErrBadPath = errors.New("invalid path")
)
// MemberRedirect maps an rpm virtual package path (<member>/<href>, as written
// by MergeRPM) to the owning member's route. ok is false for any other path;
// ErrBadPath means the path is absolute or escapes the member.
func (e *Engine) MemberRedirect(ctx context.Context, virt models.Virtual, path, rawQuery, proxyBaseURL string) (string, bool, error) {
if virt.PackageType != models.PackageRPM || strings.HasPrefix(path, "repodata/") {
return "", false, nil
}
name, rest, found := strings.Cut(path, "/")
if !found || !slices.Contains(virt.Members, name) {
return "", false, nil
}
segs, err := memberSegments(rest)
if err != nil {
return "", false, err
}
remote, err := e.getRemote(ctx, name)
if err != nil {
return "", false, nil
}
loc := fmt.Sprintf("%s/api/v1/%s/%s/%s", strings.TrimRight(proxyBaseURL, "/"), remote.RepoType, url.PathEscape(name), strings.Join(segs, "/"))
if rawQuery != "" {
loc += "?" + rawQuery
}
return loc, true, nil
}
// memberSegments decodes rest and returns its path-escaped segments, rejecting
// absolute paths and any "." or ".." segment.
func memberSegments(rest string) ([]string, error) {
dec, err := url.PathUnescape(rest)
if err != nil || dec == "" || strings.HasPrefix(dec, "/") {
return nil, ErrBadPath
}
segs := strings.Split(dec, "/")
for i, s := range segs {
if s == "." || s == ".." || strings.Contains(s, "\\") {
return nil, ErrBadPath
}
segs[i] = url.PathEscape(s)
}
return segs, nil
}
func (e *Engine) fetchRPM(ctx context.Context, virt models.Virtual, path string) ([]byte, string, error) {
if !strings.HasPrefix(path, "repodata/") {
return nil, "", ErrNotFound
}
if path != "repodata/repomd.xml" {
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
for i, name := range virt.Members {
wg.Add(1)
go func() {
defer wg.Done()
m, err := e.rpmMember(ctx, name)
if err != nil {
errs[i] = fmt.Errorf("member %q: %w", name, err)
return
}
members[i] = *m
}()
}
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: %v", virt.Name, err)
}
repo, err := MergeRPM(members)
if err != nil {
return nil, fmt.Errorf("merge rpm repodata: %w", err)
}
return repo, nil
}
func (e *Engine) cachedRPMFile(virt, path string) ([]byte, bool) {
e.mu.Lock()
defer e.mu.Unlock()
g := e.rpm[virt]
if g == nil {
return nil, false
}
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) {
remote, err := e.getRemote(ctx, name)
if err != nil {
return nil, fmt.Errorf("remote %q: %w", name, err)
}
repomd, err := e.fetchMemberPath(ctx, *remote, "repodata/repomd.xml")
if err != nil {
return nil, err
}
locs, ts, err := parseRepomd(repomd)
if err != nil {
return nil, err
}
m := &RPMMember{RemoteName: name, Bases: remote.UpstreamPool(), Timestamp: ts, Data: map[string][]byte{}}
for t, href := range locs {
raw, err := e.fetchMemberPath(ctx, *remote, href)
if err != nil {
return nil, err
}
if m.Data[t], err = decompress(href, raw); err != nil {
return nil, fmt.Errorf("decompress %s: %w", href, err)
}
}
return m, nil
}
// fetchMemberPath reads one path from a member the same way its own route
// would serve it: local index, synthesized remote (e.g. github_rpm), or proxy.
func (e *Engine) fetchMemberPath(ctx context.Context, remote models.Remote, path string) ([]byte, error) {
prov, err := provider.Get(remote.PackageType)
if err != nil {
return nil, fmt.Errorf("provider %q: %w", remote.PackageType, err)
}
var serve func(w http.ResponseWriter, r *http.Request) bool
if remote.RepoType == models.RepoTypeLocal {
indexer, ok := prov.(provider.LocalIndexer)
if !ok {
return nil, fmt.Errorf("provider %q does not serve a local index", remote.PackageType)
}
serve = func(w http.ResponseWriter, r *http.Request) bool {
return indexer.ServeLocalIndex(w, r, e.db, remote.Name, path)
}
} else if rs, ok := prov.(provider.RemoteServer); ok {
serve = func(w http.ResponseWriter, r *http.Request) bool {
return rs.ServeRemote(w, r, remote, path, "", e.db)
}
}
if serve != nil {
rec := httptest.NewRecorder()
if serve(rec, httptest.NewRequestWithContext(ctx, http.MethodGet, "/"+path, nil)) {
if rec.Code != http.StatusOK {
return nil, fmt.Errorf("%s/%s: status %d", remote.Name, path, rec.Code)
}
return rec.Body.Bytes(), nil
}
if remote.RepoType == models.RepoTypeLocal {
return nil, fmt.Errorf("%s/%s: %w", remote.Name, path, ErrNotFound)
}
}
res, err := e.proxyEngine.Fetch(ctx, remote, path, prov)
if err != nil {
return nil, fmt.Errorf("fetch %s/%s: %w", remote.Name, path, err)
}
defer res.Reader.Close()
return io.ReadAll(res.Reader)
}
+158
View File
@@ -0,0 +1,158 @@
package virtual
import (
"context"
"errors"
"fmt"
"regexp"
"sync"
"sync/atomic"
"testing"
"time"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
func fakeEngine(members map[string]*RPMMember) *Engine {
return &Engine{
getRemote: func(_ context.Context, name string) (*models.Remote, error) {
return &models.Remote{Name: name, RepoType: models.RepoTypeRemote}, nil
},
rpmMember: func(_ context.Context, name string) (*RPMMember, error) {
if m := members[name]; m != nil {
return m, nil
}
return nil, errors.New("upstream down")
},
}
}
func rpmVirt(members ...string) models.Virtual {
return models.Virtual{Name: "v", PackageType: models.PackageRPM, Members: members}
}
func TestFetchRPMFailsClosed(t *testing.T) {
e := fakeEngine(map[string]*RPMMember{
"a": {RemoteName: "a", Data: map[string][]byte{"primary": primaryXML(primaryPkgXML("foo", "1", "aaa", "foo.rpm"))}},
})
_, _, err := e.Fetch(context.Background(), rpmVirt("a", "b"), "repodata/repomd.xml", "")
if err == nil || errors.Is(err, ErrNotFound) {
t.Fatalf("want upstream error with a member down, got %v", err)
}
}
func TestFetchRPMDataSurvivesMemberChange(t *testing.T) {
m := &RPMMember{RemoteName: "a", Data: map[string][]byte{"primary": primaryXML(primaryPkgXML("foo", "1", "aaa", "foo.rpm"))}}
e := fakeEngine(map[string]*RPMMember{"a": m})
virt := rpmVirt("a")
repomd, _, err := e.Fetch(context.Background(), virt, "repodata/repomd.xml", "")
if err != nil {
t.Fatal(err)
}
m.Data = map[string][]byte{"primary": primaryXML(primaryPkgXML("foo", "2", "bbb", "foo2.rpm"))}
href := regexp.MustCompile(`repodata/[0-9a-f]+-primary\.xml\.gz`).Find(repomd)
body, _, err := e.Fetch(context.Background(), virt, string(href), "")
if err != nil {
t.Fatalf("data from previous repomd must stay resolvable: %v", err)
}
if got := string(gunzip(t, body)); !regexp.MustCompile(`ver="1"`).MatchString(got) {
t.Fatalf("served wrong primary:\n%s", got)
}
newMD, _, err := e.Fetch(context.Background(), virt, "repodata/repomd.xml", "")
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) {
e := fakeEngine(nil)
virt := rpmVirt("gh", "other")
for _, tc := range []struct {
path, query, want string
ok bool
err error
}{
{path: "gh/o/r/releases/download/v1/c++-1.rpm", query: "x=1", ok: true,
want: "http://h/api/v1/remote/gh/o/r/releases/download/v1/c++-1.rpm?x=1"},
{path: "gh/a%20b/c%3Fd.rpm", ok: true, want: "http://h/api/v1/remote/gh/a%20b/c%3Fd.rpm"},
{path: "gh/../../remote/other/x", err: ErrBadPath},
{path: "gh/%2e%2e/%2e%2e/remote/other/x", err: ErrBadPath},
{path: "gh/a/./b", err: ErrBadPath},
{path: "gh//etc/passwd", err: ErrBadPath},
{path: "repodata/repomd.xml"},
{path: "nope/x.rpm"},
} {
got, ok, err := e.MemberRedirect(context.Background(), virt, tc.path, tc.query, "http://h/")
if got != tc.want || ok != tc.ok || !errors.Is(err, tc.err) {
t.Errorf("%s: got (%q, %v, %v), want (%q, %v, %v)", tc.path, got, ok, err, tc.want, tc.ok, tc.err)
}
}
}
+298
View File
@@ -0,0 +1,298 @@
package virtual
import (
"bytes"
"compress/bzip2"
"compress/gzip"
"crypto/sha256"
"encoding/hex"
"encoding/xml"
"errors"
"fmt"
"io"
"regexp"
"strings"
"github.com/klauspost/compress/zstd"
"github.com/ulikunitz/xz"
)
var rpmDataTypes = []string{"primary", "filelists", "other"}
var rpmRoots = map[string]string{
"primary": `<metadata xmlns="http://linux.duke.edu/metadata/common" xmlns:rpm="http://linux.duke.edu/metadata/rpm" packages="%d">`,
"filelists": `<filelists xmlns="http://linux.duke.edu/metadata/filelists" packages="%d">`,
"other": `<otherdata xmlns="http://linux.duke.edu/metadata/other" packages="%d">`,
}
var rpmRootClose = map[string]string{"primary": "</metadata>", "filelists": "</filelists>", "other": "</otherdata>"}
var locationRe = regexp.MustCompile(`<location\b[^>]*/>`)
// RPMMember is one member repo's decompressed primary/filelists/other XML,
// keyed by repomd data type. Missing types contribute no packages.
type RPMMember struct {
RemoteName string
Bases []string
Timestamp int64
Data map[string][]byte
}
// RPMRepo is a merged yum repo: repomd.xml plus the files it references, keyed
// by their path under the repo root.
type RPMRepo struct {
Repomd []byte
Files map[string][]byte
}
type repomdDoc struct {
Revision string `xml:"revision"`
Data []repomdData `xml:"data"`
}
type repomdData struct {
Type string `xml:"type,attr"`
Location struct {
Href string `xml:"href,attr"`
} `xml:"location"`
Timestamp int64 `xml:"timestamp"`
}
type rpmPkgVersion struct {
Epoch string `xml:"epoch,attr"`
Ver string `xml:"ver,attr"`
Rel string `xml:"rel,attr"`
}
type primaryPkg struct {
Name string `xml:"name"`
Arch string `xml:"arch"`
Version rpmPkgVersion `xml:"version"`
Checksum string `xml:"checksum"`
Location struct {
Href string `xml:"href,attr"`
Base string `xml:"http://www.w3.org/XML/1998/namespace base,attr"`
} `xml:"location"`
}
type pkgidPkg struct {
PkgID string `xml:"pkgid,attr"`
}
// parseRepomd returns the member's primary/filelists/other locations and its
// newest data timestamp.
func parseRepomd(body []byte) (map[string]string, int64, error) {
var doc repomdDoc
if err := xml.Unmarshal(body, &doc); err != nil {
return nil, 0, fmt.Errorf("parse repomd.xml: %w", err)
}
locs := map[string]string{}
var ts int64
for _, d := range doc.Data {
for _, t := range rpmDataTypes {
if d.Type == t && d.Location.Href != "" {
locs[t] = d.Location.Href
ts = max(ts, d.Timestamp)
}
}
}
if locs["primary"] == "" {
return nil, 0, errors.New("repomd.xml has no primary data")
}
return locs, ts, nil
}
func decompress(href string, body []byte) ([]byte, error) {
var r io.Reader
switch {
case strings.HasSuffix(href, ".gz"):
gz, err := gzip.NewReader(bytes.NewReader(body))
if err != nil {
return nil, err
}
r = gz
case strings.HasSuffix(href, ".xz"):
x, err := xz.NewReader(bytes.NewReader(body))
if err != nil {
return nil, err
}
r = x
case strings.HasSuffix(href, ".zst"):
z, err := zstd.NewReader(bytes.NewReader(body))
if err != nil {
return nil, err
}
defer z.Close()
r = z
case strings.HasSuffix(href, ".bz2"):
r = bzip2.NewReader(bytes.NewReader(body))
default:
return body, nil
}
return io.ReadAll(r)
}
// splitPackages returns the raw bytes of each top-level <package> element.
func splitPackages(doc []byte) ([][]byte, error) {
d := xml.NewDecoder(bytes.NewReader(doc))
var pkgs [][]byte
depth := 0
for {
start := d.InputOffset()
tok, err := d.Token()
if err == io.EOF {
return pkgs, nil
}
if err != nil {
return nil, err
}
switch t := tok.(type) {
case xml.StartElement:
if depth == 1 && t.Name.Local == "package" {
if err := d.Skip(); err != nil {
return nil, err
}
pkgs = append(pkgs, doc[start:d.InputOffset()])
continue
}
depth++
case xml.EndElement:
depth--
}
}
}
// memberHref returns href relative to the member root. An xml:base under one of
// the member's upstream bases is folded into the href; any other xml:base is a
// host the member can't serve, so ok is false and the location is left as is.
func memberHref(m RPMMember, base, href string) (string, bool) {
href = strings.TrimLeft(href, "/")
if base == "" {
return href, true
}
base = strings.TrimRight(base, "/") + "/"
for _, b := range m.Bases {
b = strings.TrimRight(b, "/") + "/"
if strings.HasPrefix(base, b) {
return strings.TrimPrefix(base, b) + href, true
}
}
return "", false
}
func xmlEscape(s string) string {
var b bytes.Buffer
_ = xml.EscapeText(&b, []byte(s))
return b.String()
}
// MergeRPM merges members into one repo. Members are in priority order: the
// first member to carry a NEVRA wins. Package locations are prefixed with the
// owning member's name so the virtual can route downloads back to it.
func MergeRPM(members []RPMMember) (*RPMRepo, error) {
kept := map[string]bool{}
seen := map[string]bool{}
out := map[string][][]byte{}
for i, m := range members {
pkgs, err := splitPackages(m.Data["primary"])
if err != nil {
return nil, fmt.Errorf("member %q primary: %w", m.RemoteName, err)
}
for _, raw := range pkgs {
var p primaryPkg
if err := xml.Unmarshal(raw, &p); err != nil {
return nil, fmt.Errorf("member %q primary package: %w", m.RemoteName, err)
}
epoch := p.Version.Epoch
if epoch == "" {
epoch = "0"
}
nevra := fmt.Sprintf("%s-%s:%s-%s.%s", p.Name, epoch, p.Version.Ver, p.Version.Rel, p.Arch)
if seen[nevra] {
continue
}
seen[nevra] = true
kept[fmt.Sprintf("%d/%s", i, p.Checksum)] = true
if href, ok := memberHref(m, p.Location.Base, p.Location.Href); ok {
loc := `<location href="` + xmlEscape(m.RemoteName+"/"+href) + `"/>`
raw = locationRe.ReplaceAll(raw, []byte(loc))
}
out["primary"] = append(out["primary"], raw)
}
}
for _, t := range rpmDataTypes[1:] {
emitted := map[string]bool{}
for i, m := range members {
pkgs, err := splitPackages(m.Data[t])
if err != nil {
return nil, fmt.Errorf("member %q %s: %w", m.RemoteName, t, err)
}
for _, raw := range pkgs {
var p pkgidPkg
if err := xml.Unmarshal(raw, &p); err != nil {
return nil, fmt.Errorf("member %q %s package: %w", m.RemoteName, t, err)
}
key := fmt.Sprintf("%d/%s", i, p.PkgID)
if kept[key] && !emitted[key] {
emitted[key] = true
out[t] = append(out[t], raw)
}
}
}
}
var ts int64
for _, m := range members {
ts = max(ts, m.Timestamp)
}
repo := &RPMRepo{Files: map[string][]byte{}}
var md bytes.Buffer
md.WriteString(xml.Header)
md.WriteString(`<repomd xmlns="http://linux.duke.edu/metadata/repo" xmlns:rpm="http://linux.duke.edu/metadata/rpm">` + "\n")
fmt.Fprintf(&md, " <revision>%d</revision>\n", ts)
for _, t := range rpmDataTypes {
var doc bytes.Buffer
doc.WriteString(xml.Header)
fmt.Fprintf(&doc, rpmRoots[t]+"\n", len(out[t]))
for _, raw := range out[t] {
doc.Write(raw)
doc.WriteString("\n")
}
doc.WriteString(rpmRootClose[t] + "\n")
gz := gzipDeterministic(doc.Bytes())
sum := sha256Hex(gz)
href := fmt.Sprintf("repodata/%s-%s.xml.gz", sum, t)
repo.Files[href] = gz
fmt.Fprintf(&md, " <data type=\"%s\">\n", t)
fmt.Fprintf(&md, " <checksum type=\"sha256\">%s</checksum>\n", sum)
fmt.Fprintf(&md, " <open-checksum type=\"sha256\">%s</open-checksum>\n", sha256Hex(doc.Bytes()))
fmt.Fprintf(&md, " <location href=\"%s\"/>\n", href)
fmt.Fprintf(&md, " <timestamp>%d</timestamp>\n", ts)
fmt.Fprintf(&md, " <size>%d</size>\n", len(gz))
fmt.Fprintf(&md, " <open-size>%d</open-size>\n", doc.Len())
md.WriteString(" </data>\n")
}
md.WriteString("</repomd>\n")
repo.Repomd = md.Bytes()
return repo, nil
}
func gzipDeterministic(data []byte) []byte {
var buf bytes.Buffer
gz := gzip.NewWriter(&buf)
gz.Header = gzip.Header{OS: 255}
_, _ = gz.Write(data)
_ = gz.Close()
return buf.Bytes()
}
func sha256Hex(data []byte) string {
h := sha256.Sum256(data)
return hex.EncodeToString(h[:])
}
+274
View File
@@ -0,0 +1,274 @@
package virtual
import (
"bytes"
"compress/gzip"
"crypto/sha256"
"encoding/hex"
"encoding/xml"
"fmt"
"io"
"strings"
"testing"
"github.com/klauspost/compress/zstd"
"github.com/ulikunitz/xz"
)
func primaryXML(pkgs ...string) []byte {
return []byte(`<?xml version="1.0" encoding="UTF-8"?>
<metadata xmlns="http://linux.duke.edu/metadata/common" xmlns:rpm="http://linux.duke.edu/metadata/rpm" packages="` + fmt.Sprint(len(pkgs)) + `">
` + strings.Join(pkgs, "\n") + `
</metadata>`)
}
func primaryPkgXML(name, ver, pkgid, href string) string {
return `<package type="rpm">
<name>` + name + `</name>
<arch>x86_64</arch>
<version epoch="0" ver="` + ver + `" rel="1"/>
<checksum type="sha256" pkgid="YES">` + pkgid + `</checksum>
<location href="` + href + `"/>
<format>
<rpm:provides><rpm:entry name="` + name + `"/></rpm:provides>
</format>
</package>`
}
func filelistsXML(pkgids ...string) []byte {
var b strings.Builder
b.WriteString(`<?xml version="1.0"?><filelists xmlns="http://linux.duke.edu/metadata/filelists" packages="1">`)
for _, id := range pkgids {
b.WriteString(`<package pkgid="` + id + `" name="x" arch="x86_64"><version epoch="0" ver="1" rel="1"/><file>/usr/bin/` + id + `</file></package>`)
}
b.WriteString(`</filelists>`)
return []byte(b.String())
}
func gunzip(t *testing.T, b []byte) []byte {
t.Helper()
r, err := gzip.NewReader(bytes.NewReader(b))
if err != nil {
t.Fatal(err)
}
out, err := io.ReadAll(r)
if err != nil {
t.Fatal(err)
}
return out
}
func mergedFile(t *testing.T, repo *RPMRepo, dtype string) []byte {
t.Helper()
for href, body := range repo.Files {
if strings.HasSuffix(href, "-"+dtype+".xml.gz") {
return gunzip(t, body)
}
}
t.Fatalf("no %s file in merged repo", dtype)
return nil
}
func TestMergeRPMDedupePriority(t *testing.T) {
members := []RPMMember{
{RemoteName: "first", Timestamp: 100, Data: map[string][]byte{
"primary": primaryXML(primaryPkgXML("foo", "1.0", "aaa", "Packages/foo-1.0.rpm")),
"filelists": filelistsXML("aaa"),
}},
{RemoteName: "second", Timestamp: 200, Data: map[string][]byte{
"primary": primaryXML(
primaryPkgXML("foo", "1.0", "bbb", "Packages/foo-1.0-rebuilt.rpm"),
primaryPkgXML("bar", "2.0", "ccc", "Packages/bar-2.0.rpm"),
),
"filelists": filelistsXML("bbb", "ccc"),
}},
}
repo, err := MergeRPM(members)
if err != nil {
t.Fatal(err)
}
primary := string(mergedFile(t, repo, "primary"))
if strings.Count(primary, "<name>foo</name>") != 1 {
t.Fatalf("duplicate NEVRA not deduped:\n%s", primary)
}
if !strings.Contains(primary, `href="first/Packages/foo-1.0.rpm"`) || strings.Contains(primary, "foo-1.0-rebuilt") {
t.Fatalf("first member should win duplicate NEVRA:\n%s", primary)
}
if !strings.Contains(primary, `href="second/Packages/bar-2.0.rpm"`) {
t.Fatalf("unique package from second member missing:\n%s", primary)
}
if !strings.Contains(primary, `packages="2"`) {
t.Fatalf("package count wrong:\n%s", primary)
}
filelists := string(mergedFile(t, repo, "filelists"))
if !strings.Contains(filelists, "/usr/bin/aaa") || !strings.Contains(filelists, "/usr/bin/ccc") || strings.Contains(filelists, "/usr/bin/bbb") {
t.Fatalf("filelists must follow the winning primary entries:\n%s", filelists)
}
other := string(mergedFile(t, repo, "other"))
if !strings.Contains(other, `packages="0"`) {
t.Fatalf("members without other data should yield an empty other.xml:\n%s", other)
}
if !strings.Contains(string(repo.Repomd), "<revision>200</revision>") {
t.Fatalf("revision should be the newest member timestamp:\n%s", repo.Repomd)
}
}
func TestMergeRPMOutputIsWellFormed(t *testing.T) {
repo, err := MergeRPM([]RPMMember{{RemoteName: "a", Data: map[string][]byte{
"primary": primaryXML(primaryPkgXML("foo", "1.0", "aaa", "Packages/foo.rpm")),
}}})
if err != nil {
t.Fatal(err)
}
var doc struct {
XMLName xml.Name
Packages []struct {
Name string `xml:"name"`
Provides []struct {
Name string `xml:"name,attr"`
} `xml:"format>provides>entry"`
} `xml:"package"`
}
if err := xml.Unmarshal(mergedFile(t, repo, "primary"), &doc); err != nil {
t.Fatal(err)
}
if doc.XMLName.Space != "http://linux.duke.edu/metadata/common" || len(doc.Packages) != 1 {
t.Fatalf("unexpected merged primary: %+v", doc)
}
if len(doc.Packages[0].Provides) != 1 || doc.Packages[0].Provides[0].Name != "foo" {
t.Fatalf("rpm: namespaced format lost: %+v", doc.Packages[0])
}
}
func TestMergeRPMHrefRewriting(t *testing.T) {
based := func(name, base string) string {
return `<package type="rpm"><name>` + name + `</name><arch>noarch</arch><version epoch="0" ver="1" rel="1"/>` +
`<checksum type="sha256" pkgid="YES">` + name + `</checksum><location xml:base="` + base + `" href="Packages/` + name + `.rpm"/></package>`
}
repo, err := MergeRPM([]RPMMember{{RemoteName: "gh", Bases: []string{"https://up.example/el9", "https://mirror.example/el9/"}, Data: map[string][]byte{
"primary": primaryXML(
primaryPkgXML("foo", "1.0", "aaa", "storytold/photocraft/releases/download/v1.0/foo&amp;bar.rpm"),
primaryPkgXML("lead", "1.0", "bbb", "/Packages/lead.rpm"),
based("root", "https://up.example/el9/"),
based("sub", "https://mirror.example/el9/extra"),
based("ext", "https://other.example/"),
),
}}})
if err != nil {
t.Fatal(err)
}
primary := string(mergedFile(t, repo, "primary"))
for _, want := range []string{
`<location href="gh/storytold/photocraft/releases/download/v1.0/foo&amp;bar.rpm"/>`,
`<location href="gh/Packages/lead.rpm"/>`,
`<location href="gh/Packages/root.rpm"/>`,
`<location href="gh/extra/Packages/sub.rpm"/>`,
`<location xml:base="https://other.example/" href="Packages/ext.rpm"/>`,
} {
if !strings.Contains(primary, want) {
t.Errorf("missing %s in:\n%s", want, primary)
}
}
}
func TestMergeRPMChecksums(t *testing.T) {
repo, err := MergeRPM([]RPMMember{{RemoteName: "a", Timestamp: 42, Data: map[string][]byte{
"primary": primaryXML(primaryPkgXML("foo", "1.0", "aaa", "Packages/foo.rpm")),
}}})
if err != nil {
t.Fatal(err)
}
var md struct {
Data []struct {
Type string `xml:"type,attr"`
Checksum string `xml:"checksum"`
OpenChecksum string `xml:"open-checksum"`
Location struct {
Href string `xml:"href,attr"`
} `xml:"location"`
Timestamp int64 `xml:"timestamp"`
Size int `xml:"size"`
OpenSize int `xml:"open-size"`
} `xml:"data"`
}
if err := xml.Unmarshal(repo.Repomd, &md); err != nil {
t.Fatal(err)
}
if len(md.Data) != 3 {
t.Fatalf("want primary/filelists/other, got %d entries", len(md.Data))
}
for _, d := range md.Data {
body, ok := repo.Files[d.Location.Href]
if !ok {
t.Fatalf("%s: location %q not served", d.Type, d.Location.Href)
}
sum := sha256.Sum256(body)
if hex.EncodeToString(sum[:]) != d.Checksum || len(body) != d.Size {
t.Errorf("%s: checksum/size do not match served bytes", d.Type)
}
open := gunzip(t, body)
osum := sha256.Sum256(open)
if hex.EncodeToString(osum[:]) != d.OpenChecksum || len(open) != d.OpenSize {
t.Errorf("%s: open-checksum/open-size do not match decompressed bytes", d.Type)
}
if d.Timestamp != 42 {
t.Errorf("%s: timestamp %d, want 42", d.Type, d.Timestamp)
}
}
again, _ := MergeRPM([]RPMMember{{RemoteName: "a", Timestamp: 42, Data: map[string][]byte{
"primary": primaryXML(primaryPkgXML("foo", "1.0", "aaa", "Packages/foo.rpm")),
}}})
if !bytes.Equal(repo.Repomd, again.Repomd) {
t.Error("merge must be deterministic so repomd checksums stay valid across requests")
}
}
func TestParseRepomd(t *testing.T) {
locs, ts, err := parseRepomd([]byte(`<repomd xmlns="http://linux.duke.edu/metadata/repo">
<data type="primary"><location href="repodata/p-primary.xml.zst"/><timestamp>10</timestamp></data>
<data type="primary_db"><location href="repodata/p.sqlite.bz2"/><timestamp>99</timestamp></data>
<data type="other"><location href="repodata/o-other.xml.gz"/><timestamp>20</timestamp></data>
</repomd>`))
if err != nil {
t.Fatal(err)
}
if locs["primary"] != "repodata/p-primary.xml.zst" || locs["other"] != "repodata/o-other.xml.gz" || len(locs) != 2 {
t.Fatalf("unexpected locations: %v", locs)
}
if ts != 20 {
t.Fatalf("timestamp %d, want 20 (sqlite entries ignored)", ts)
}
if _, _, err := parseRepomd([]byte(`<repomd><data type="other"><location href="x"/></data></repomd>`)); err == nil {
t.Fatal("repomd without primary should error")
}
}
func TestDecompress(t *testing.T) {
plain := []byte("<metadata/>")
var xzBuf bytes.Buffer
xw, _ := xz.NewWriter(&xzBuf)
_, _ = xw.Write(plain)
_ = xw.Close()
zw, _ := zstd.NewWriter(nil)
zst := zw.EncodeAll(plain, nil)
for href, body := range map[string][]byte{
"p.xml": plain,
"p.xml.gz": gzipDeterministic(plain),
"p.xml.xz": xzBuf.Bytes(),
"p.xml.zst": zst,
} {
got, err := decompress(href, body)
if err != nil || !bytes.Equal(got, plain) {
t.Errorf("%s: got %q, %v", href, got, err)
}
}
}
+9
View File
@@ -0,0 +1,9 @@
-- Failed GitHub release scans retry with backoff instead of waiting out a full
-- mutable_ttl: sync_failures counts consecutive failures, next_retry_at gates
-- the next claim while a retry is pending.
ALTER TABLE github_rpm_sync_state ADD COLUMN IF NOT EXISTS sync_failures INT NOT NULL DEFAULT 0;
ALTER TABLE github_rpm_sync_state ADD COLUMN IF NOT EXISTS next_retry_at TIMESTAMPTZ;
ALTER TABLE github_deb_sync_state ADD COLUMN IF NOT EXISTS sync_failures INT NOT NULL DEFAULT 0;
ALTER TABLE github_deb_sync_state ADD COLUMN IF NOT EXISTS next_retry_at TIMESTAMPTZ;
ALTER TABLE github_alpine_sync_state ADD COLUMN IF NOT EXISTS sync_failures INT NOT NULL DEFAULT 0;
ALTER TABLE github_alpine_sync_state ADD COLUMN IF NOT EXISTS next_retry_at TIMESTAMPTZ;
+2
View File
@@ -16,6 +16,7 @@ const (
PackagePuppet PackageType = "puppet"
PackageTerraform PackageType = "terraform"
PackageGoProxy PackageType = "goproxy"
PackageCargo PackageType = "cargo"
PackageGitHubRPM PackageType = "github_rpm"
PackageGitHubDeb PackageType = "github_deb"
PackageGitHubAlpine PackageType = "github_alpine"
@@ -33,6 +34,7 @@ var validPackageTypes = map[PackageType]bool{
PackagePuppet: true,
PackageTerraform: true,
PackageGoProxy: true,
PackageCargo: true,
PackageGitHubRPM: true,
PackageGitHubDeb: true,
PackageGitHubAlpine: true,
+1
View File
@@ -18,6 +18,7 @@ func TestPackageTypeValid(t *testing.T) {
models.PackagePuppet,
models.PackageTerraform,
models.PackageGoProxy,
models.PackageCargo,
models.PackageGitHubRPM,
models.PackageGitHubDeb,
models.PackageGitHubAlpine,
+15
View File
@@ -228,6 +228,21 @@ go mod download`,
},
];
case 'cargo':
return [
{
title: 'Point crates.io at this sparse registry',
language: 'toml',
code: `# ~/.cargo/config.toml
[source.crates-io]
replace-with = "${name}"
[registries.${name}]
index = "sparse+${proxy}/"`,
note: 'config.json is served by artifactapi so crate downloads also go through this remote.',
},
];
case 'puppet':
return [
{
+1
View File
@@ -17,6 +17,7 @@ const typeColors: Record<string, 'blue' | 'green' | 'yellow' | 'red' | 'default'
puppet: 'yellow',
terraform: 'blue',
goproxy: 'green',
cargo: 'yellow',
};
export function Locals() {
+1
View File
@@ -17,6 +17,7 @@ const typeColors: Record<string, 'blue' | 'green' | 'yellow' | 'red' | 'default'
puppet: 'yellow',
terraform: 'blue',
goproxy: 'green',
cargo: 'yellow',
};
export function Remotes() {