Compare commits

..

3 Commits

Author SHA1 Message Date
unkin-agent 6c6ad3066e Fix github_alpine .apk redirect to resolve stored FilePath
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
apk reconstructs the download URL itself as <arch>/<name>-<version>.apk
because APKINDEX carries no filename field (unlike rpm's <location> or
deb's Filename:). ServeRemote forwarded that synthesized path verbatim
into the releases_remote redirect, pointing at a nonexistent,
allowlist-denied github.com path (404/403).

Look up the cached metadata row by arch plus the full reconstructed
filename (no hyphen-split, so -rN suffixes are preserved) and redirect
to the stored github-relative FilePath. Unknown packages now 404 instead
of redirecting to a bad path.
2026-08-12 01:28:25 +10:00
unkin-agent b1de05d3b4 Add github_alpine metadata-only package type
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
github_alpine is the Alpine/apk analog of github_deb/github_rpm: a
metadata-only remote that scans a GitHub repo's releases for .apk assets,
derives each package's .PKGINFO via a ranged prefix fetch (never
downloading whole packages), synthesizes a per-arch APKINDEX.tar.gz from
that cached metadata, and 302-redirects .apk downloads to a backend
releases_remote. It stacks on the apk-local work, reusing the alpine
provider's APKINDEX generator, .PKGINFO parser, Q1 checksum, and
AlpineMetadata store.

- pkg/models: add PackageGitHubAlpine to the enum + validators
- internal/provider/alpine/github.go: the github_alpine provider
  (ServeRemote per-arch index + .apk redirect, cold-start 503,
  scanWithState incremental derive, ranged .PKGINFO prefix fetch with
  range-doubling on truncation)
- internal/provider/alpine/syncer.go: parallel background Syncer
  (worker pool, shared limiter, deduped queue, DB lease)
- internal/database/alpine_github_sync.go + github_alpine_sync_state
  table: remote enumeration + per-remote sync lease
- internal/api/v2/remotes.go: primed on create via the shared Primer map
- internal/server/server.go: construct + Run the alpine syncer, register
  it in the Primer map
- tests mirror the deb github_test/syncer_test (scan/diff/prune, ranged
  .PKGINFO parse, per-arch ServeRemote routing, .apk 302, DB lease)
2026-08-12 01:14:20 +10:00
unkin-agent 58a24a15dd Add Alpine/apk local repository support
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
The alpine provider hosted only remote (proxy) repos; there was no way to
publish first-party .apk packages the way rpm-local and deb-local already
allow. This extends the existing alpine provider into a real apk repository:
uploaded .apk files are parsed in pure Go and a per-arch APKINDEX.tar.gz is
generated on demand, at parity with rpm repodata and deb Packages generation.

- Implement LocalUploader/LocalIndexer/PostUploadHook/PostDeleteHook on the
  alpine provider, keeping the remote proxy methods intact.
- Parse .apk (concatenated gzipped tar streams) in pure Go: read .PKGINFO from
  the control stream and compute the apk pull checksum C: = Q1+base64(sha1) over
  the raw control gzip stream (not the whole file).
- Generate an unsigned per-arch APKINDEX.tar.gz (clients use --allow-untrusted),
  applying the same dot-segment normalization as deb for ./<arch>/... requests.
- Add AlpineMetadata plus separate Alpine store/reader/deleter interfaces so the
  rpm/deb metadata interfaces are not widened.
- Add the alpine_metadata table and its Insert/Delete/List DB methods.
- Add testsupport.MinimalApk plus unit tests (parse, Q1 checksum, per-arch
  filtering, dot-segment handling, validate) and a dockere2e index test.
2026-08-12 00:59:15 +10:00
40 changed files with 37 additions and 2137 deletions
-16
View File
@@ -9,18 +9,6 @@ services:
# No host port needed: only the artifactapi container talks to it, and the
# tests compare served bytes against the on-disk fixtures.
# Two constant-body upstreams for the multi-base_url suite: each returns a
# distinct, upstream-identifying body for any path, so round-robin
# distribution across a two-mirror remote is directly observable.
mockupstreama:
image: nginx:alpine
volumes:
- ./e2e-docker/mirror-conf/a.conf:/etc/nginx/conf.d/default.conf:ro,z
mockupstreamb:
image: nginx:alpine
volumes:
- ./e2e-docker/mirror-conf/b.conf:/etc/nginx/conf.d/default.conf:ro,z
artifactapi:
# The host port is set via ARTIFACTAPI_PORT (see scripts/docker-e2e.sh),
# defaulting to 8000; the e2e run uses 8001 to avoid colliding with a
@@ -28,7 +16,3 @@ services:
depends_on:
mockupstream:
condition: service_started
mockupstreama:
condition: service_started
mockupstreamb:
condition: service_started
-12
View File
@@ -30,18 +30,6 @@ already-running stack.
index), rpm (real package + **automatic repodata** generation).
- **Virtual repositories** — pypi simple-index merge and helm `index.yaml` merge
across two members.
- **Mirrorlist** — an rpm remote with a `mirrorlist` of extra upstream mirrors
(pool = `base_url` + `mirrorlist`): round-robin distribution across both mirrors
(constant-body `mockupstreama` / `mockupstreamb`), failover past a dead primary,
no-mirrorlist regression, and a real `dnf` (stock `rockylinux:9` container)
`makecache` + `install` through a two-mirror rpm remote whose `base_url` is dead
— a dead mirror must not break the client.
- **Mirror strategy (`least_conn`)** — a `mirror_strategy: least_conn` rpm remote
over the two constant-body mirrors exercises the least-connections selection
path end-to-end (both mirrors serve, all requests succeed), plus a real `dnf`
install through a `least_conn` remote with a dead primary (failover unchanged).
The precise least-loaded pick is asserted deterministically in the proxy unit
test, since an in-flight-skew assertion over HTTP is timing-sensitive.
## Fixtures
@@ -1,55 +0,0 @@
<?xml version="1.0" encoding="UTF-8"?>
<repomd xmlns="http://linux.duke.edu/metadata/repo" xmlns:rpm="http://linux.duke.edu/metadata/rpm">
<revision>1786573032</revision>
<data type="primary">
<checksum type="sha256">d82f717e4da1afe96b8e7857de9e852f5785c0d75734e50d0f8afdcbc6261b08</checksum>
<open-checksum type="sha256">3345bb631380ae6c0620fe2a29cf7dff2ed4e28c5cb06c91e8a1628ba1979bcb</open-checksum>
<location href="repodata/d82f717e4da1afe96b8e7857de9e852f5785c0d75734e50d0f8afdcbc6261b08-primary.xml.gz"/>
<timestamp>1786573032</timestamp>
<size>631</size>
<open-size>1192</open-size>
</data>
<data type="filelists">
<checksum type="sha256">daa313cc5eeb7df556e1d4885d7701b10b9f012f436ef239fa46827f966222be</checksum>
<open-checksum type="sha256">648bd0ce00fda09abbc6e9c3ff3278518a76f258576ac24cd10c12e41e0e5bd7</open-checksum>
<location href="repodata/daa313cc5eeb7df556e1d4885d7701b10b9f012f436ef239fa46827f966222be-filelists.xml.gz"/>
<timestamp>1786573032</timestamp>
<size>256</size>
<open-size>338</open-size>
</data>
<data type="other">
<checksum type="sha256">8510c74a6f288828bbc92abee5d0d8ae9687d3c31a2579ea95e31a4c3a320d85</checksum>
<open-checksum type="sha256">c42cfd3843e9c53a60ad84bded44aa46c65faae7b3da99a099a9ca0b018872a9</open-checksum>
<location href="repodata/8510c74a6f288828bbc92abee5d0d8ae9687d3c31a2579ea95e31a4c3a320d85-other.xml.gz"/>
<timestamp>1786573032</timestamp>
<size>296</size>
<open-size>399</open-size>
</data>
<data type="primary_db">
<checksum type="sha256">f6bd7755da13d9726381048f467992869104a4c5521338ef740dc35eb85b9b71</checksum>
<open-checksum type="sha256">c45c85d12ccb0f8172b7bfae466362c08ac1867559574a9fb9cb2118c04daddb</open-checksum>
<location href="repodata/f6bd7755da13d9726381048f467992869104a4c5521338ef740dc35eb85b9b71-primary.sqlite.bz2"/>
<timestamp>1786573032</timestamp>
<size>1740</size>
<open-size>106496</open-size>
<database_version>10</database_version>
</data>
<data type="filelists_db">
<checksum type="sha256">be3c6e4c7a13ece48bd5d6a4d6d5e6a2395fe006ef9c5f217f5b87144f465e57</checksum>
<open-checksum type="sha256">1ccfa3dff532d782ce3225aae807506a4ce4534291386f1c47455dcc6b70cfd6</open-checksum>
<location href="repodata/be3c6e4c7a13ece48bd5d6a4d6d5e6a2395fe006ef9c5f217f5b87144f465e57-filelists.sqlite.bz2"/>
<timestamp>1786573032</timestamp>
<size>764</size>
<open-size>28672</open-size>
<database_version>10</database_version>
</data>
<data type="other_db">
<checksum type="sha256">ba593cd8ab5ec1e127888707c1fd882920996f1ce173fd7a589d647918fd7da4</checksum>
<open-checksum type="sha256">5d4d38380f0e359bfc0a50033d4faa84c75c5ae8fe2ef11c1c82d64748e9b8e2</open-checksum>
<location href="repodata/ba593cd8ab5ec1e127888707c1fd882920996f1ce173fd7a589d647918fd7da4-other.sqlite.bz2"/>
<timestamp>1786573032</timestamp>
<size>738</size>
<open-size>24576</open-size>
<database_version>10</database_version>
</data>
</repomd>
-105
View File
@@ -1,105 +0,0 @@
//go:build dockere2e
package e2edocker
import (
"fmt"
"net/http"
"os"
"os/exec"
"strings"
"testing"
)
// TestLeastConnMultiBaseURL configures an rpm remote with mirror_strategy =
// least_conn over a two-mirror pool and drives distinct cache-miss paths through
// it, asserting every request succeeds and both mirrors serve traffic. This
// exercises the least-connections selection path (leastConnOrder + the in-flight
// gauge inc/dec around each upstream call) end-to-end through a real HTTP client.
// A precise least-loaded assertion is timing-sensitive over HTTP and is covered
// deterministically by the proxy unit test (TestLeastConnPicksLeastLoaded);
// here, with requests issued serially, in-flight counts return to zero between
// them so equal-load mirrors are spread by the round-robin tie-break.
func TestLeastConnMultiBaseURL(t *testing.T) {
name := "e2e-leastconn"
createRepo(t, fmt.Sprintf(`{
"name": %q,
"package_type": "rpm",
"repo_type": "remote",
"base_url": %q,
"mirrorlist": [%q],
"mirror_strategy": "least_conn",
"stale_on_error": false
}`, name, mockUpstreamA(), mockUpstreamB()))
defer deleteRepo(t, name)
seenA, seenB := false, false
const n = 12
for i := 0; i < n; i++ {
url := api(fmt.Sprintf("/api/v1/remote/%s/lc/%d", name, i))
resp, body := doRequest(t, http.MethodGet, url, nil, "")
if resp.StatusCode != http.StatusOK {
t.Fatalf("request %d: status %d: %s", i, resp.StatusCode, body)
}
switch strings.TrimSpace(string(body)) {
case "UPSTREAM-A":
seenA = true
case "UPSTREAM-B":
seenB = true
default:
t.Fatalf("request %d: unexpected body %q", i, body)
}
}
if !seenA || !seenB {
t.Fatalf("least_conn remote did not reach both upstreams: A=%v B=%v", seenA, seenB)
}
}
// TestLeastConnDnfInstall drives a real dnf (stock rockylinux container) at a
// two-mirror rpm remote configured with mirror_strategy = least_conn whose
// primary base_url is dead: makecache + install must succeed via the live mirror.
// This proves a real package-manager client installs correctly through a
// least_conn remote and that failover semantics are unchanged under the new
// strategy. Requires the compose network exported by scripts/docker-e2e.sh.
func TestLeastConnDnfInstall(t *testing.T) {
network := os.Getenv("COMPOSE_NETWORK")
internal := os.Getenv("ARTIFACTAPI_INTERNAL")
if network == "" || internal == "" {
t.Skip("COMPOSE_NETWORK/ARTIFACTAPI_INTERNAL not set; run via scripts/docker-e2e.sh")
}
if _, err := exec.LookPath("docker"); err != nil {
t.Skip("docker not available on the test host")
}
name := "e2e-leastconn-dnf"
createRepo(t, fmt.Sprintf(`{
"name": %q,
"package_type": "rpm",
"repo_type": "remote",
"base_url": "http://mockupstream-dead:80",
"mirrorlist": [%q],
"mirror_strategy": "least_conn",
"stale_on_error": false
}`, name, mockUpstream()))
defer deleteRepo(t, name)
repoURL := strings.TrimRight(internal, "/") + "/api/v1/remote/" + name + "/rpm-mirror"
repoConf := fmt.Sprintf("[dnflc]\nname=dnflc\nbaseurl=%s\nenabled=1\ngpgcheck=0\nsslverify=0\nmetadata_expire=0\n", repoURL)
script := "set -euo pipefail; " +
"printf '%s' \"$REPO\" > /etc/yum.repos.d/dnflc.repo; " +
"dnf -y --disablerepo='*' --enablerepo=dnflc makecache; " +
"dnf -y --disablerepo='*' --enablerepo=dnflc install e2e-testpkg; " +
"rpm -q e2e-testpkg"
cmd := exec.Command("docker", "run", "--rm",
"--network", network,
"-e", "REPO="+repoConf,
"rockylinux:9", "bash", "-c", script)
out, err := cmd.CombinedOutput()
if err != nil {
t.Fatalf("real dnf install through a least_conn remote failed: %v\n%s", err, out)
}
if !strings.Contains(string(out), "e2e-testpkg-1.0-1") {
t.Fatalf("dnf did not install the expected package via least_conn remote; output:\n%s", out)
}
}
-9
View File
@@ -1,9 +0,0 @@
# Mock upstream A for the multi-base_url e2e: any path returns a constant,
# upstream-identifying body so round-robin distribution is observable.
server {
listen 80;
location / {
default_type text/plain;
return 200 "UPSTREAM-A";
}
}
-8
View File
@@ -1,8 +0,0 @@
# Mock upstream B for the multi-base_url e2e (see a.conf).
server {
listen 80;
location / {
default_type text/plain;
return 200 "UPSTREAM-B";
}
}
-163
View File
@@ -1,163 +0,0 @@
//go:build dockere2e
package e2edocker
import (
"fmt"
"net/http"
"os"
"os/exec"
"strings"
"testing"
)
// mockUpstreamA/B are the constant-body upstreams (see docker-compose.e2e.yml)
// that let the round-robin test observe which mirror served each request.
func mockUpstreamA() string {
if v := os.Getenv("MOCK_UPSTREAM_A_INTERNAL"); v != "" {
return strings.TrimRight(v, "/")
}
return "http://mockupstreama"
}
func mockUpstreamB() string {
if v := os.Getenv("MOCK_UPSTREAM_B_INTERNAL"); v != "" {
return strings.TrimRight(v, "/")
}
return "http://mockupstreamb"
}
// TestMultiBaseURLRoundRobin configures an rpm remote with base_url = mirror A
// and mirrorlist = [mirror B] and drives distinct paths through it, asserting
// both mirrors serve traffic. Each path is a cache miss, so every request reaches
// upstream and the round-robin cursor alternates mirrors.
func TestMultiBaseURLRoundRobin(t *testing.T) {
name := "e2e-rr"
createRepo(t, fmt.Sprintf(`{
"name": %q,
"package_type": "rpm",
"repo_type": "remote",
"base_url": %q,
"mirrorlist": [%q],
"stale_on_error": false
}`, name, mockUpstreamA(), mockUpstreamB()))
defer deleteRepo(t, name)
seenA, seenB := false, false
const n = 12
for i := 0; i < n; i++ {
url := api(fmt.Sprintf("/api/v1/remote/%s/rr/%d", name, i))
resp, body := doRequest(t, http.MethodGet, url, nil, "")
if resp.StatusCode != http.StatusOK {
t.Fatalf("request %d: status %d: %s", i, resp.StatusCode, body)
}
switch strings.TrimSpace(string(body)) {
case "UPSTREAM-A":
seenA = true
case "UPSTREAM-B":
seenB = true
default:
t.Fatalf("request %d: unexpected body %q", i, body)
}
}
if !seenA || !seenB {
t.Fatalf("round-robin did not reach both upstreams: A=%v B=%v", seenA, seenB)
}
}
// TestMultiBaseURLFailover points a two-mirror remote at a dead primary and a
// healthy secondary and asserts every request still succeeds via the secondary.
func TestMultiBaseURLFailover(t *testing.T) {
name := "e2e-failover"
createRepo(t, fmt.Sprintf(`{
"name": %q,
"package_type": "rpm",
"repo_type": "remote",
"base_url": "http://mockupstream-dead:80",
"mirrorlist": [%q],
"stale_on_error": false
}`, name, mockUpstreamB()))
defer deleteRepo(t, name)
for i := 0; i < 6; i++ {
url := api(fmt.Sprintf("/api/v1/remote/%s/fo/%d", name, i))
resp, body := doRequest(t, http.MethodGet, url, nil, "")
if resp.StatusCode != http.StatusOK {
t.Fatalf("request %d: dead primary broke fetch: status %d: %s", i, resp.StatusCode, body)
}
if got := strings.TrimSpace(string(body)); got != "UPSTREAM-B" {
t.Fatalf("request %d: body %q, want UPSTREAM-B (served via failover)", i, got)
}
}
}
// TestSingleBaseURLRegression asserts a remote with no mirrorlist works exactly
// as before the mirrorlist change.
func TestSingleBaseURLRegression(t *testing.T) {
name := "e2e-single"
createRepo(t, fmt.Sprintf(`{
"name": %q,
"package_type": "rpm",
"repo_type": "remote",
"base_url": %q,
"stale_on_error": false
}`, name, mockUpstreamA()))
defer deleteRepo(t, name)
resp, body := doRequest(t, http.MethodGet, api("/api/v1/remote/"+name+"/solo/0"), nil, "")
if resp.StatusCode != http.StatusOK {
t.Fatalf("single-url fetch: status %d: %s", resp.StatusCode, body)
}
if got := strings.TrimSpace(string(body)); got != "UPSTREAM-A" {
t.Fatalf("single-url body %q, want UPSTREAM-A", got)
}
}
// TestMultiBaseURLDnfFailover drives a real dnf (stock rockylinux container) at
// a two-mirror rpm remote whose primary is dead: makecache + install must
// succeed via the live secondary mirror, proving a dead mirror does not break a
// real package-manager client. Requires the compose network and internal API
// URL exported by scripts/docker-e2e.sh; skipped when run standalone.
func TestMultiBaseURLDnfFailover(t *testing.T) {
network := os.Getenv("COMPOSE_NETWORK")
internal := os.Getenv("ARTIFACTAPI_INTERNAL")
if network == "" || internal == "" {
t.Skip("COMPOSE_NETWORK/ARTIFACTAPI_INTERNAL not set; run via scripts/docker-e2e.sh")
}
if _, err := exec.LookPath("docker"); err != nil {
t.Skip("docker not available on the test host")
}
name := "e2e-dnf-failover"
// Primary base_url is dead; the live mirror serves the real yum repo under
// fixtures/rpm-mirror via the shared mock upstream.
createRepo(t, fmt.Sprintf(`{
"name": %q,
"package_type": "rpm",
"repo_type": "remote",
"base_url": "http://mockupstream-dead:80",
"mirrorlist": [%q],
"stale_on_error": false
}`, name, mockUpstream()))
defer deleteRepo(t, name)
repoURL := strings.TrimRight(internal, "/") + "/api/v1/remote/" + name + "/rpm-mirror"
repoConf := fmt.Sprintf("[dnffo]\nname=dnffo\nbaseurl=%s\nenabled=1\ngpgcheck=0\nsslverify=0\nmetadata_expire=0\n", repoURL)
script := "set -euo pipefail; " +
"printf '%s' \"$REPO\" > /etc/yum.repos.d/dnffo.repo; " +
"dnf -y --disablerepo='*' --enablerepo=dnffo makecache; " +
"dnf -y --disablerepo='*' --enablerepo=dnffo install e2e-testpkg; " +
"rpm -q e2e-testpkg"
cmd := exec.Command("docker", "run", "--rm",
"--network", network,
"-e", "REPO="+repoConf,
"rockylinux:9", "bash", "-c", script)
out, err := cmd.CombinedOutput()
if err != nil {
t.Fatalf("real dnf install through a dead primary mirror failed: %v\n%s", err, out)
}
if !strings.Contains(string(out), "e2e-testpkg-1.0-1") {
t.Fatalf("dnf did not install the expected package via failover; output:\n%s", out)
}
}
+1 -1
View File
@@ -57,7 +57,7 @@ func do(t *testing.T, h http.Handler, method, path, body string) int {
}
func TestRemotesErrorPaths(t *testing.T) {
h := NewRemotesHandler(closedDB(t), nil, nil).Routes()
h := NewRemotesHandler(closedDB(t), nil).Routes()
if c := do(t, h, "GET", "/", ""); c != 500 {
t.Errorf("list with dead db = %d, want 500", c)
}
+4 -50
View File
@@ -1,10 +1,8 @@
package v2
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"net/http"
"github.com/go-chi/chi/v5"
@@ -19,23 +17,15 @@ type Primer interface {
EnqueuePrime(remote models.Remote)
}
// MetadataFlusher purges a remote's cached mutable metadata (repodata / Release
// / APKINDEX freshness keys). *cache.Redis satisfies it.
type MetadataFlusher interface {
FlushRemote(ctx context.Context, remote string) error
}
type RemotesHandler struct {
db *database.DB
cache MetadataFlusher
primers map[models.PackageType]Primer
}
// NewRemotesHandler wires the handler to the metadata cache and per-type
// primers. cache may be nil (flush-on-backend-change is skipped); primers may
// be nil (a package type with no registered primer simply skips priming).
func NewRemotesHandler(db *database.DB, cache MetadataFlusher, primers map[models.PackageType]Primer) *RemotesHandler {
return &RemotesHandler{db: db, cache: cache, primers: primers}
// NewRemotesHandler wires the handler to the per-type metadata primers. primers
// may be nil; a package type with no registered primer simply skips priming.
func NewRemotesHandler(db *database.DB, primers map[models.PackageType]Primer) *RemotesHandler {
return &RemotesHandler{db: db, primers: primers}
}
func (h *RemotesHandler) Routes() chi.Router {
@@ -88,14 +78,6 @@ func (h *RemotesHandler) create(w http.ResponseWriter, r *http.Request) {
http.Error(w, "base_url is required for remote repositories", http.StatusBadRequest)
return
}
if err := remote.ValidateMirrorlist(); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
if err := remote.ValidateMirrorStrategy(); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
if err := remote.ValidatePatterns(); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
@@ -120,42 +102,14 @@ func (h *RemotesHandler) update(w http.ResponseWriter, r *http.Request) {
return
}
remote.Name = name
if err := remote.ValidateMirrorlist(); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
if err := remote.ValidateMirrorStrategy(); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
if err := remote.ValidatePatterns(); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
// Capture the current backend before the update so we can tell whether the
// remote's base_url (its upstream) changed. A read failure just means we
// skip the freshness flush; it must not block the update.
oldBaseURL, oldKnown := "", false
if existing, err := h.db.GetRemote(r.Context(), name); err == nil {
oldBaseURL, oldKnown = existing.BaseURL, true
}
if err := h.db.UpdateRemote(r.Context(), &remote); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Changing the backend invalidates any cached mutable metadata (repodata /
// Release / APKINDEX): purge it so the next request re-fetches from the new
// upstream instead of serving stale data until TTL expiry. A flush failure
// is logged but does not fail the request — the DB update already landed.
if oldKnown && oldBaseURL != remote.BaseURL && h.cache != nil {
if err := h.cache.FlushRemote(r.Context(), name); err != nil {
slog.Warn("flush cached metadata after base_url change failed",
"remote", name, "error", err)
} else {
slog.Info("flushed cached metadata after base_url change",
"remote", name, "old_base_url", oldBaseURL, "new_base_url", remote.BaseURL)
}
}
writeJSON(w, http.StatusOK, remote)
}
-96
View File
@@ -1,96 +0,0 @@
package v2
import (
"context"
"errors"
"testing"
"git.unkin.net/unkin/artifactapi/internal/database"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// fakeFlusher records FlushRemote calls so a test can assert whether — and how
// often — a remote's cached metadata was purged.
type fakeFlusher struct {
calls []string
err error
}
func (f *fakeFlusher) FlushRemote(_ context.Context, remote string) error {
f.calls = append(f.calls, remote)
return f.err
}
func seedRemote(t *testing.T, db *database.DB, name, baseURL string) {
t.Helper()
err := db.CreateRemote(context.Background(), &models.Remote{
Name: name,
PackageType: models.PackageRPM,
RepoType: models.RepoTypeRemote,
BaseURL: baseURL,
})
if err != nil {
t.Fatalf("seed remote: %v", err)
}
}
// A base_url change must flush the remote's cached metadata exactly once, while
// an update that leaves base_url untouched must not flush at all.
func TestUpdateFlushesCacheOnBaseURLChange(t *testing.T) {
if testDSN == "" {
t.Skip("Docker unavailable")
}
db, err := database.New(testDSN)
if err != nil {
t.Fatal(err)
}
defer db.Close()
const name = "rpm-flush-change"
seedRemote(t, db, name, "https://old.example.com/repo")
ff := &fakeFlusher{}
h := NewRemotesHandler(db, ff, nil).Routes()
if c := do(t, h, "PUT", "/"+name, `{"package_type":"rpm","repo_type":"remote","base_url":"https://new.example.com/repo"}`); c != 200 {
t.Fatalf("update (backend change) = %d, want 200", c)
}
if len(ff.calls) != 1 || ff.calls[0] != name {
t.Fatalf("flush calls = %v, want exactly one flush of %q", ff.calls, name)
}
// Re-updating with the same (now current) base_url must not flush again.
ff.calls = nil
if c := do(t, h, "PUT", "/"+name, `{"package_type":"rpm","repo_type":"remote","base_url":"https://new.example.com/repo"}`); c != 200 {
t.Fatalf("update (no backend change) = %d, want 200", c)
}
if len(ff.calls) != 0 {
t.Fatalf("flush calls = %v, want no flush when base_url is unchanged", ff.calls)
}
}
// A flush error must be swallowed: the DB update already succeeded, so the
// request still returns 200.
func TestUpdateFlushFailureStillSucceeds(t *testing.T) {
if testDSN == "" {
t.Skip("Docker unavailable")
}
db, err := database.New(testDSN)
if err != nil {
t.Fatal(err)
}
defer db.Close()
const name = "rpm-flush-error"
seedRemote(t, db, name, "https://old.example.com/repo")
ff := &fakeFlusher{err: errors.New("redis down")}
h := NewRemotesHandler(db, ff, nil).Routes()
if c := do(t, h, "PUT", "/"+name, `{"package_type":"rpm","repo_type":"remote","base_url":"https://new.example.com/repo"}`); c != 200 {
t.Fatalf("update with failing flush = %d, want 200", c)
}
if len(ff.calls) != 1 {
t.Fatalf("flush calls = %v, want exactly one attempted flush", ff.calls)
}
}
+1 -1
View File
@@ -41,7 +41,7 @@ func (db *DB) ListAlpineMetadataEntries(ctx context.Context, repoName string) ([
depends, provides, install_if
FROM alpine_metadata
WHERE repo_name = $1
ORDER BY name, version, arch, file_path
ORDER BY name, version, arch
`, repoName)
if err != nil {
return nil, err
-47
View File
@@ -99,53 +99,6 @@ func TestRemotesCRUD(t *testing.T) {
}
}
func TestRemoteMirrorlistRoundTrip(t *testing.T) {
requireDB(t)
mirrors := []string{"https://b.example", "https://c.example"}
if err := testDB.CreateRemote(ctx(), &models.Remote{
Name: "r-mirror", PackageType: models.PackageRPM, RepoType: models.RepoTypeRemote,
BaseURL: "https://a.example", Mirrorlist: mirrors, MutableTTL: 3600,
}); err != nil {
t.Fatalf("create mirrorlist remote: %v", err)
}
defer testDB.DeleteRemote(ctx(), "r-mirror")
got, err := testDB.GetRemote(ctx(), "r-mirror")
if err != nil {
t.Fatalf("get: %v", err)
}
if got.BaseURL != "https://a.example" {
t.Fatalf("BaseURL = %q, want https://a.example", got.BaseURL)
}
if len(got.Mirrorlist) != 2 || got.Mirrorlist[0] != mirrors[0] || got.Mirrorlist[1] != mirrors[1] {
t.Fatalf("Mirrorlist round-trip = %v, want %v", got.Mirrorlist, mirrors)
}
// An unset strategy is stored as the round_robin default.
if got.MirrorStrategy != models.MirrorStrategyRoundRobin {
t.Fatalf("MirrorStrategy default = %q, want %q", got.MirrorStrategy, models.MirrorStrategyRoundRobin)
}
// Updating to least_conn round-trips.
got.MirrorStrategy = models.MirrorStrategyLeastConn
if err := testDB.UpdateRemote(ctx(), got); err != nil {
t.Fatalf("update to least_conn: %v", err)
}
got, _ = testDB.GetRemote(ctx(), "r-mirror")
if got.MirrorStrategy != models.MirrorStrategyLeastConn {
t.Fatalf("MirrorStrategy after update = %q, want least_conn", got.MirrorStrategy)
}
// Clearing the mirrorlist on update persists an empty list.
got.Mirrorlist = nil
if err := testDB.UpdateRemote(ctx(), got); err != nil {
t.Fatalf("update clearing mirrorlist: %v", err)
}
got, _ = testDB.GetRemote(ctx(), "r-mirror")
if len(got.Mirrorlist) != 0 {
t.Fatalf("mirrorlist after clear = %v, want empty", got.Mirrorlist)
}
}
func TestArtifactsAndBlobs(t *testing.T) {
requireDB(t)
seedRemote(t, "r-art")
+3 -3
View File
@@ -31,10 +31,10 @@ func (db *DB) ListDebMetadataEntries(ctx context.Context, repoName string) ([]pr
rows, err := db.Pool.Query(ctx, `
SELECT repo_name, file_path, content_hash,
name, version, architecture, control,
size, md5, sha256, created_at
size, md5, sha256
FROM deb_metadata
WHERE repo_name = $1
ORDER BY name, version, architecture, file_path
ORDER BY name, version, architecture
`, repoName)
if err != nil {
return nil, err
@@ -47,7 +47,7 @@ func (db *DB) ListDebMetadataEntries(ctx context.Context, repoName string) ([]pr
if err := rows.Scan(
&m.RepoName, &m.FilePath, &m.ContentHash,
&m.Name, &m.Version, &m.Architecture, &m.Control,
&m.Size, &m.MD5, &m.SHA256, &m.CreatedAt,
&m.Size, &m.MD5, &m.SHA256,
); err != nil {
return nil, err
}
-4
View File
@@ -44,8 +44,6 @@ func (db *DB) migrate() error {
package_type TEXT NOT NULL,
repo_type TEXT DEFAULT 'remote',
base_url TEXT NOT NULL DEFAULT '',
mirrorlist TEXT[] DEFAULT '{}',
mirror_strategy TEXT NOT NULL DEFAULT 'round_robin',
description TEXT DEFAULT '',
username TEXT DEFAULT '',
password TEXT DEFAULT '',
@@ -126,8 +124,6 @@ func (db *DB) migrate() error {
CREATE INDEX IF NOT EXISTS idx_access_log_remote_time ON access_log(remote_name, created_at);
ALTER TABLE remotes ADD COLUMN IF NOT EXISTS repo_type TEXT DEFAULT 'remote';
ALTER TABLE remotes ADD COLUMN IF NOT EXISTS mirrorlist TEXT[] DEFAULT '{}';
ALTER TABLE remotes ADD COLUMN IF NOT EXISTS mirror_strategy TEXT NOT NULL DEFAULT 'round_robin';
ALTER TABLE remotes ADD COLUMN IF NOT EXISTS upstream_dial_timeout INTEGER DEFAULT 0;
ALTER TABLE remotes ADD COLUMN IF NOT EXISTS upstream_tls_timeout INTEGER DEFAULT 0;
ALTER TABLE remotes ADD COLUMN IF NOT EXISTS upstream_response_header_timeout INTEGER DEFAULT 0;
+7 -21
View File
@@ -6,7 +6,7 @@ import (
"git.unkin.net/unkin/artifactapi/pkg/models"
)
const remoteCols = `name, package_type, repo_type, base_url, mirrorlist, mirror_strategy, description, username, password,
const remoteCols = `name, package_type, repo_type, base_url, description, username, password,
immutable_ttl, mutable_ttl, check_mutable,
patterns, blocklist, mutable_patterns, immutable_patterns,
ban_tags_enabled, ban_tags,
@@ -15,18 +15,9 @@ const remoteCols = `name, package_type, repo_type, base_url, mirrorlist, mirror_
upstream_dial_timeout, upstream_tls_timeout, upstream_response_header_timeout,
created_at, updated_at`
// normalizeMirrorStrategy maps an empty strategy to the round_robin default so
// the NOT NULL mirror_strategy column always stores a canonical value.
func normalizeMirrorStrategy(s string) string {
if s == "" {
return models.MirrorStrategyRoundRobin
}
return s
}
func scanRemote(scanner interface{ Scan(...any) error }, r *models.Remote) error {
return scanner.Scan(
&r.Name, &r.PackageType, &r.RepoType, &r.BaseURL, &r.Mirrorlist, &r.MirrorStrategy, &r.Description, &r.Username, &r.Password,
&r.Name, &r.PackageType, &r.RepoType, &r.BaseURL, &r.Description, &r.Username, &r.Password,
&r.ImmutableTTL, &r.MutableTTL, &r.CheckMutable,
&r.Patterns, &r.Blocklist, &r.MutablePatterns, &r.ImmutablePatterns,
&r.BanTagsEnabled, &r.BanTags,
@@ -67,24 +58,22 @@ func (db *DB) ListRemotes(ctx context.Context) ([]models.Remote, error) {
func (db *DB) CreateRemote(ctx context.Context, r *models.Remote) error {
_, err := db.Pool.Exec(ctx, `
INSERT INTO remotes (
name, package_type, repo_type, base_url, mirrorlist, description, username, password,
name, package_type, repo_type, base_url, description, username, password,
immutable_ttl, mutable_ttl, check_mutable,
patterns, blocklist, mutable_patterns, immutable_patterns,
ban_tags_enabled, ban_tags,
quarantine_enabled, quarantine_days, stale_on_error,
releases_remote, managed_by,
upstream_dial_timeout, upstream_tls_timeout, upstream_response_header_timeout,
mirror_strategy
) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20,$21,$22,$23,$24,$25,$26)
upstream_dial_timeout, upstream_tls_timeout, upstream_response_header_timeout
) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20,$21,$22,$23,$24)
`,
r.Name, r.PackageType, r.RepoType, r.BaseURL, r.Mirrorlist, r.Description, r.Username, r.Password,
r.Name, r.PackageType, r.RepoType, r.BaseURL, r.Description, r.Username, r.Password,
r.ImmutableTTL, r.MutableTTL, r.CheckMutable,
r.Patterns, r.Blocklist, r.MutablePatterns, r.ImmutablePatterns,
r.BanTagsEnabled, r.BanTags,
r.QuarantineEnabled, r.QuarantineDays, r.StaleOnError,
r.ReleasesRemote, r.ManagedBy,
r.UpstreamDialTimeout, r.UpstreamTLSTimeout, r.UpstreamResponseHeaderTimeout,
normalizeMirrorStrategy(r.MirrorStrategy),
)
return err
}
@@ -92,14 +81,13 @@ func (db *DB) CreateRemote(ctx context.Context, r *models.Remote) error {
func (db *DB) UpdateRemote(ctx context.Context, r *models.Remote) error {
_, err := db.Pool.Exec(ctx, `
UPDATE remotes SET
package_type=$2, repo_type=$3, base_url=$4, mirrorlist=$25, description=$5, username=$6, password=$7,
package_type=$2, repo_type=$3, base_url=$4, description=$5, username=$6, password=$7,
immutable_ttl=$8, mutable_ttl=$9, check_mutable=$10,
patterns=$11, blocklist=$12, mutable_patterns=$13, immutable_patterns=$14,
ban_tags_enabled=$15, ban_tags=$16,
quarantine_enabled=$17, quarantine_days=$18, stale_on_error=$19,
releases_remote=$20, managed_by=$21,
upstream_dial_timeout=$22, upstream_tls_timeout=$23, upstream_response_header_timeout=$24,
mirror_strategy=$26,
updated_at=NOW()
WHERE name=$1
`,
@@ -110,8 +98,6 @@ func (db *DB) UpdateRemote(ctx context.Context, r *models.Remote) error {
r.QuarantineEnabled, r.QuarantineDays, r.StaleOnError,
r.ReleasesRemote, r.ManagedBy,
r.UpstreamDialTimeout, r.UpstreamTLSTimeout, r.UpstreamResponseHeaderTimeout,
r.Mirrorlist,
normalizeMirrorStrategy(r.MirrorStrategy),
)
return err
}
+2 -7
View File
@@ -3,7 +3,6 @@ package database
import (
"context"
"encoding/json"
"time"
"git.unkin.net/unkin/artifactapi/internal/provider"
)
@@ -66,7 +65,6 @@ type RPMMetadataRow struct {
Obsoletes json.RawMessage
Files json.RawMessage
Changelogs json.RawMessage
CreatedAt time.Time
}
func (db *DB) ListRPMMetadataEntries(ctx context.Context, repoName string) ([]provider.RPMMetadata, error) {
@@ -96,7 +94,6 @@ func (db *DB) ListRPMMetadataEntries(ctx context.Context, repoName string) ([]pr
SourceRPM: r.SourceRPM,
URL: r.URL,
Packager: r.Packager,
CreatedAt: r.CreatedAt,
}
json.Unmarshal(r.Requires, &meta.Requires)
json.Unmarshal(r.Provides, &meta.Provides)
@@ -115,11 +112,10 @@ func (db *DB) ListRPMMetadata(ctx context.Context, repoName string) ([]RPMMetada
name, epoch, version, release, arch,
summary, description, rpm_size, installed_size,
license, vendor, build_group, build_host, source_rpm, url, packager,
requires, provides, conflicts, obsoletes, files, changelogs,
created_at
requires, provides, conflicts, obsoletes, files, changelogs
FROM rpm_metadata
WHERE repo_name = $1
ORDER BY name, epoch, version, release, arch, file_path
ORDER BY name, epoch, version, release, arch
`, repoName)
if err != nil {
return nil, err
@@ -135,7 +131,6 @@ func (db *DB) ListRPMMetadata(ctx context.Context, repoName string) ([]RPMMetada
&r.Summary, &r.Description, &r.RPMSize, &r.InstalledSize,
&r.License, &r.Vendor, &r.Group, &r.BuildHost, &r.SourceRPM, &r.URL, &r.Packager,
&r.Requires, &r.Provides, &r.Conflicts, &r.Obsoletes, &r.Files, &r.Changelogs,
&r.CreatedAt,
); err != nil {
return nil, err
}
+1 -5
View File
@@ -15,7 +15,6 @@ import (
"path"
"strconv"
"strings"
"time"
"archive/tar"
@@ -377,10 +376,7 @@ func generateAPKIndex(metas []provider.AlpineMetadata) []byte {
var tarBuf bytes.Buffer
tw := tar.NewWriter(&tarBuf)
body := idx.Bytes()
// ModTime is pinned to the Unix epoch (never wall clock) so APKINDEX.tar.gz
// is byte-identical across replicas and regenerations (issue #117); apk
// clients ignore the tar mtime.
tw.WriteHeader(&tar.Header{Name: "APKINDEX", Mode: 0o644, Size: int64(len(body)), Typeflag: tar.TypeReg, ModTime: time.Unix(0, 0)})
tw.WriteHeader(&tar.Header{Name: "APKINDEX", Mode: 0o644, Size: int64(len(body)), Typeflag: tar.TypeReg})
tw.Write(body)
tw.Close()
@@ -1,72 +0,0 @@
package alpine
import (
"archive/tar"
"bytes"
"compress/gzip"
"io"
"testing"
"time"
"git.unkin.net/unkin/artifactapi/internal/provider"
)
func apkFixture() []provider.AlpineMetadata {
return []provider.AlpineMetadata{
{
RepoName: "r", FilePath: "x86_64/aaa-1.0-r0.apk", Checksum: "Q1aaa",
Name: "aaa", Version: "1.0-r0", Arch: "x86_64", DownloadSize: 100, InstalledSize: 10,
Description: "pkg aaa", URL: "https://a", License: "MIT",
Depends: []string{"so:libc"}, Provides: []string{"cmd:aaa"}, BuildTime: 1710000000,
},
{
RepoName: "r", FilePath: "x86_64/bbb-2.0-r0.apk", Checksum: "Q1bbb",
Name: "bbb", Version: "2.0-r0", Arch: "x86_64", DownloadSize: 200, InstalledSize: 20,
},
}
}
// TestAPKIndexDeterministic asserts APKINDEX.tar.gz is byte-identical across two
// generations separated by wall-clock time, so the two no-affinity replicas and
// every regeneration serve the same bytes (issue #117).
func TestAPKIndexDeterministic(t *testing.T) {
metas := apkFixture()
first := generateAPKIndex(metas)
time.Sleep(10 * time.Millisecond)
second := generateAPKIndex(metas)
if !bytes.Equal(first, second) {
t.Error("APKINDEX.tar.gz differs across generations")
}
}
// TestAPKIndexTarModTimePinned guards the tar header: its ModTime must be the
// pinned Unix epoch, never wall clock. Fails if a future edit stamps time.Now().
func TestAPKIndexTarModTimePinned(t *testing.T) {
metas := apkFixture()
zr, err := gzip.NewReader(bytes.NewReader(generateAPKIndex(metas)))
if err != nil {
t.Fatalf("gzip: %v", err)
}
if !zr.ModTime.IsZero() && zr.ModTime.Unix() != 0 {
t.Errorf("gzip header ModTime = %v, want zero/epoch", zr.ModTime)
}
tarBytes, err := io.ReadAll(zr)
if err != nil {
t.Fatalf("gunzip: %v", err)
}
tr := tar.NewReader(bytes.NewReader(tarBytes))
hdr, err := tr.Next()
if err != nil {
t.Fatalf("tar: %v", err)
}
if hdr.Name != "APKINDEX" {
t.Fatalf("tar entry = %q, want APKINDEX", hdr.Name)
}
if hdr.ModTime.Unix() != 0 {
t.Errorf("APKINDEX tar ModTime = %v (unix %d), want epoch (0)", hdr.ModTime, hdr.ModTime.Unix())
}
}
+1 -16
View File
@@ -395,7 +395,7 @@ func generateRelease(metas []provider.DebMetadata) []byte {
arches := uniqueArches(metas)
var b bytes.Buffer
fmt.Fprintf(&b, "Date: %s\n", releaseDate(metas).Format(time.RFC1123Z))
fmt.Fprintf(&b, "Date: %s\n", time.Now().UTC().Format(time.RFC1123Z))
fmt.Fprintf(&b, "Architectures: %s\n", strings.Join(arches, " "))
b.WriteString("Acquire-By-Hash: no\n")
@@ -410,21 +410,6 @@ func generateRelease(metas []provider.DebMetadata) []byte {
return b.Bytes()
}
// releaseDate derives the Release Date: from the newest package's persisted
// created_at (in UTC) so the file is byte-identical across the no-affinity
// replicas and across regenerations (issue #117); an empty repo falls back to
// the Unix epoch. This never uses wall clock, which also keeps Date: from
// running ahead of any Valid-Until logic.
func releaseDate(metas []provider.DebMetadata) time.Time {
newest := time.Unix(0, 0)
for _, m := range metas {
if m.CreatedAt.After(newest) {
newest = m.CreatedAt
}
}
return newest.UTC()
}
func writeReleaseEntry(b *bytes.Buffer, hash string, size int, name string) {
fmt.Fprintf(b, " %s %d %s\n", hash, size, name)
}
@@ -1,172 +0,0 @@
package deb
import (
"bytes"
"strconv"
"strings"
"testing"
"time"
"git.unkin.net/unkin/artifactapi/internal/provider"
)
// debFixture returns a fixed set of rows with persisted created_at values, in
// the total order ListDebMetadataEntries produces (name, version, arch,
// file_path), so the generators are exercised on a stable input.
func debFixture() []provider.DebMetadata {
t1 := time.Date(2026, 3, 1, 8, 30, 0, 0, time.UTC)
t2 := time.Date(2026, 4, 15, 12, 0, 0, 0, time.UTC) // newest
return []provider.DebMetadata{
{
RepoName: "r", FilePath: "pool/aaa_1.0_amd64.deb", ContentHash: "sha256:aa",
Name: "aaa", Version: "1.0", Architecture: "amd64",
Control: "Package: aaa\nVersion: 1.0\nArchitecture: amd64",
Size: 100, MD5: "d41d8cd98f00b204e9800998ecf8427e", SHA256: "aa", CreatedAt: t1,
},
{
RepoName: "r", FilePath: "pool/bbb_2.0_arm64.deb", ContentHash: "sha256:bb",
Name: "bbb", Version: "2.0", Architecture: "arm64",
Control: "Package: bbb\nVersion: 2.0\nArchitecture: arm64",
Size: 200, MD5: "0cc175b9c0f1b6a831c399e269772661", SHA256: "bb", CreatedAt: t2,
},
}
}
// TestDebGeneratorsDeterministic asserts the served bytes are a pure function of
// DB state: Packages, Packages.gz and Release are byte-identical across two
// generations separated by wall-clock time. Fails against the old
// time.Now()-stamped Release Date:.
func TestDebGeneratorsDeterministic(t *testing.T) {
metas := debFixture()
pkgs1 := generatePackages(metas)
rel1 := generateRelease(metas)
gz1 := gzipBytes(pkgs1)
time.Sleep(10 * time.Millisecond)
pkgs2 := generatePackages(metas)
rel2 := generateRelease(metas)
gz2 := gzipBytes(pkgs2)
if !bytes.Equal(pkgs1, pkgs2) {
t.Error("Packages differs across generations")
}
if !bytes.Equal(gz1, gz2) {
t.Error("Packages.gz differs across generations")
}
if !bytes.Equal(rel1, rel2) {
t.Errorf("Release differs across generations:\n--- first ---\n%s\n--- second ---\n%s", rel1, rel2)
}
}
// TestDebReleaseDateUsesPersistedCreatedAt pins the Release Date: to the newest
// persisted created_at (RFC1123Z, UTC), not wall clock. Fails against the old
// time.Now() code.
func TestDebReleaseDateUsesPersistedCreatedAt(t *testing.T) {
metas := debFixture()
want := time.Date(2026, 4, 15, 12, 0, 0, 0, time.UTC).Format(time.RFC1123Z)
rel := string(generateRelease(metas))
var got string
for _, line := range strings.Split(rel, "\n") {
if strings.HasPrefix(line, "Date:") {
got = strings.TrimSpace(strings.TrimPrefix(line, "Date:"))
break
}
}
if got != want {
t.Errorf("Release Date: = %q, want %q (newest created_at)", got, want)
}
}
// TestDebReleaseDateEmptyRepoIsEpoch guards the fallback: an empty repo yields a
// deterministic epoch Date: rather than wall clock.
func TestDebReleaseDateEmptyRepoIsEpoch(t *testing.T) {
want := time.Unix(0, 0).UTC().Format(time.RFC1123Z)
rel := string(generateRelease(nil))
if !strings.Contains(rel, "Date: "+want+"\n") {
t.Errorf("empty-repo Release missing epoch Date: %q\n%s", want, rel)
}
}
// TestDebReleaseChecksumsMatchServedBytes is the exact apt invariant: the
// sha256/size (and md5/size) advertised for Packages and Packages.gz in Release
// equal the sha256/size of the actual bytes ServeLocalIndex serves. apt rejects
// any mismatch.
func TestDebReleaseChecksumsMatchServedBytes(t *testing.T) {
metas := debFixture()
packages := generatePackages(metas)
packagesGz := gzipBytes(packages)
rel := string(generateRelease(metas))
wantSHA := map[string]struct {
hash string
size int
}{
"Packages": {sha256Hex(packages), len(packages)},
"Packages.gz": {sha256Hex(packagesGz), len(packagesGz)},
}
wantMD5 := map[string]struct {
hash string
size int
}{
"Packages": {md5Hex(packages), len(packages)},
"Packages.gz": {md5Hex(packagesGz), len(packagesGz)},
}
sha := parseReleaseSection(rel, "SHA256:")
md5s := parseReleaseSection(rel, "MD5Sum:")
for name, w := range wantSHA {
got, ok := sha[name]
if !ok {
t.Fatalf("Release SHA256 section missing %q", name)
}
if got.hash != w.hash || got.size != w.size {
t.Errorf("Release SHA256 %s = (%s, %d), served bytes are (%s, %d)", name, got.hash, got.size, w.hash, w.size)
}
}
for name, w := range wantMD5 {
got, ok := md5s[name]
if !ok {
t.Fatalf("Release MD5Sum section missing %q", name)
}
if got.hash != w.hash || got.size != w.size {
t.Errorf("Release MD5Sum %s = (%s, %d), served bytes are (%s, %d)", name, got.hash, got.size, w.hash, w.size)
}
}
}
type releaseEntry struct {
hash string
size int
}
// parseReleaseSection reads the indented " <hash> <size> <name>" lines that
// follow a "SHA256:" / "MD5Sum:" header until the next non-indented line.
func parseReleaseSection(release, header string) map[string]releaseEntry {
out := map[string]releaseEntry{}
lines := strings.Split(release, "\n")
in := false
for _, line := range lines {
if line == header {
in = true
continue
}
if !in {
continue
}
if !strings.HasPrefix(line, " ") {
break
}
fields := strings.Fields(line)
if len(fields) != 3 {
continue
}
size, _ := strconv.Atoi(fields[1])
out[fields[2]] = releaseEntry{hash: fields[0], size: size}
}
return out
}
-8
View File
@@ -5,7 +5,6 @@ import (
"fmt"
"io"
"net/http"
"time"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
@@ -114,10 +113,6 @@ type DebMetadata struct {
Size int64
MD5 string
SHA256 string
// CreatedAt is the persisted insert time; the Release Date: is derived from
// the newest value so the index is byte-identical across replicas and
// regenerations (issue #117) rather than stamped from wall clock.
CreatedAt time.Time
}
// AlpineMetadataStore / AlpineMetadataDeleter / AlpineMetadataReader are the
@@ -190,9 +185,6 @@ type RPMMetadata struct {
Obsoletes []RPMDep
Files []RPMFile
Changelogs []RPMChangelog
// CreatedAt is the persisted upload timestamp; used as a stable, replica-independent
// value for the repodata <time>/<revision> fields so generated indexes are deterministic.
CreatedAt time.Time
}
type RPMDep struct {
+4 -31
View File
@@ -275,7 +275,7 @@ func (p *Provider) serveRepomd(w http.ResponseWriter, r *http.Request, reader pr
filelistsHash := sha256Hex(filelists)
otherHash := sha256Hex(other)
repomd := generateRepomd(repomdRevision(metas), primaryHash, len(primary), filelistsHash, len(filelists), otherHash, len(other))
repomd := generateRepomd(primaryHash, len(primary), filelistsHash, len(filelists), otherHash, len(other))
w.Header().Set("Content-Type", "application/xml")
w.WriteHeader(http.StatusOK)
@@ -315,32 +315,8 @@ func (p *Provider) serveOther(w http.ResponseWriter, r *http.Request, reader pro
w.Write(generateOtherXMLGZ(metas))
}
// stableUnix maps a persisted timestamp to a fixed integer for repodata's
// informational <time>/<timestamp> fields. Zero times (unset) collapse to 0 so
// output stays byte-identical across replicas and requests. dnf does not
// validate these values.
func stableUnix(t time.Time) int64 {
if t.IsZero() {
return 0
}
return t.Unix()
}
// repomdRevision derives repomd.xml's <revision>/<timestamp> from persisted
// state: the newest package upload time in the repo. It changes only when the
// repo's package set does, and is identical on every replica reading the same
// rows, so repomd.xml is byte-stable.
func repomdRevision(metas []provider.RPMMetadata) string {
var max int64
for _, m := range metas {
if u := stableUnix(m.CreatedAt); u > max {
max = u
}
}
return fmt.Sprintf("%d", max)
}
func generateRepomd(ts string, primaryHash string, primarySize int, filelistsHash string, filelistsSize int, otherHash string, otherSize int) []byte {
func generateRepomd(primaryHash string, primarySize int, filelistsHash string, filelistsSize int, otherHash string, otherSize int) []byte {
ts := fmt.Sprintf("%d", time.Now().Unix())
var b bytes.Buffer
b.WriteString(xml.Header)
b.WriteString(`<repomd xmlns="http://linux.duke.edu/metadata/repo" xmlns:rpm="http://linux.duke.edu/metadata/rpm">` + "\n")
@@ -383,7 +359,7 @@ func generatePrimaryXMLGZ(metas []provider.RPMMetadata) []byte {
if m.URL != "" {
fmt.Fprintf(&xmlBuf, " <url>%s</url>\n", xmlEscape(m.URL))
}
fmt.Fprintf(&xmlBuf, " <time file=\"%d\" build=\"0\"/>\n", stableUnix(m.CreatedAt))
fmt.Fprintf(&xmlBuf, " <time file=\"%d\" build=\"0\"/>\n", time.Now().Unix())
fmt.Fprintf(&xmlBuf, " <size package=\"%d\" installed=\"%d\" archive=\"0\"/>\n", m.RPMSize, m.InstalledSize)
fmt.Fprintf(&xmlBuf, " <location href=\"%s\"/>\n", xmlEscape(m.FilePath))
fmt.Fprintf(&xmlBuf, " <format>\n")
@@ -508,9 +484,6 @@ func xmlEscape(s string) string {
func gzipBytes(data []byte) []byte {
var buf bytes.Buffer
gz := gzip.NewWriter(&buf)
// Pin every header field so the compressed bytes (and their sha256) depend
// only on the payload, never on wall-clock time or the Go version's gzip defaults.
gz.Header = gzip.Header{OS: 255}
gz.Write(data)
gz.Close()
return buf.Bytes()
@@ -1,148 +0,0 @@
package rpm
import (
"bytes"
"compress/gzip"
"encoding/xml"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"git.unkin.net/unkin/artifactapi/internal/provider"
)
func gunzip(t *testing.T, data []byte) string {
t.Helper()
zr, err := gzip.NewReader(bytes.NewReader(data))
if err != nil {
t.Fatalf("gzip reader: %v", err)
}
out, err := io.ReadAll(zr)
if err != nil {
t.Fatalf("gunzip: %v", err)
}
return string(out)
}
// sampleMetas returns a fixed two-package repo state whose upload timestamps are
// pinned, so any nondeterminism must come from the generators themselves.
func sampleMetas() []provider.RPMMetadata {
base := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
return []provider.RPMMetadata{
{
Name: "alpha", Version: "1.0", Release: "1", Arch: "x86_64",
Summary: "a", Description: "d", ContentHash: "sha256:aaa",
FilePath: "Packages/alpha-1.0-1.x86_64.rpm", RPMSize: 10, InstalledSize: 20,
Provides: []provider.RPMDep{{Name: "alpha"}},
Requires: []provider.RPMDep{{Name: "libc", Flags: "GE", Version: "2.0"}},
CreatedAt: base,
},
{
Name: "beta", Version: "2.0", Release: "3", Arch: "noarch",
Summary: "b", Description: "d2", ContentHash: "sha256:bbb",
FilePath: "Packages/beta-2.0-3.noarch.rpm", RPMSize: 30, InstalledSize: 40,
CreatedAt: base.Add(time.Hour),
},
}
}
// TestRepodataGeneratorsDeterministic is the direct regression guard for #117:
// generating each metadata document twice from identical state must yield
// byte-identical output (hence an identical sha256). The old code embedded
// time.Now() inside primary.xml.gz, so its bytes/hash drifted every second.
func TestRepodataGeneratorsDeterministic(t *testing.T) {
metas := sampleMetas()
gens := map[string]func([]provider.RPMMetadata) []byte{
"primary": generatePrimaryXMLGZ,
"filelists": generateFilelistsXMLGZ,
"other": generateOtherXMLGZ,
}
for name, gen := range gens {
a := gen(metas)
b := gen(metas)
if sha256Hex(a) != sha256Hex(b) {
t.Errorf("%s: sha256 differs between two generations (nondeterministic): %s != %s",
name, sha256Hex(a), sha256Hex(b))
}
}
// repomd.xml itself must also be byte-stable across regenerations.
r1 := generateRepomd(repomdRevision(metas), sha256Hex(generatePrimaryXMLGZ(metas)), 1, "f", 2, "o", 3)
r2 := generateRepomd(repomdRevision(metas), sha256Hex(generatePrimaryXMLGZ(metas)), 1, "f", 2, "o", 3)
if string(r1) != string(r2) {
t.Error("repomd.xml differs between two generations")
}
}
// TestPrimaryTimeUsesPersistedCreatedAt proves the <time> element is a pure
// function of the persisted upload timestamp, not the wall clock.
func TestPrimaryTimeUsesPersistedCreatedAt(t *testing.T) {
metas := sampleMetas()
out := gunzip(t, generatePrimaryXMLGZ(metas))
if want := `<time file="1767323045" build="0"/>`; !strings.Contains(out, want) {
t.Errorf("primary.xml missing persisted <time> %q; got:\n%s", want, out)
}
// A zero (unset) CreatedAt collapses to a fixed 0, never a live clock value.
metas[0].CreatedAt = time.Time{}
out = gunzip(t, generatePrimaryXMLGZ(metas))
if !strings.Contains(out, `<time file="0" build="0"/>`) {
t.Errorf("zero CreatedAt should emit file=\"0\"; got:\n%s", out)
}
}
type repomdDoc struct {
Revision string `xml:"revision"`
Data []struct {
Type string `xml:"type,attr"`
Checksum struct {
Value string `xml:",chardata"`
} `xml:"checksum"`
Location struct {
Href string `xml:"href,attr"`
} `xml:"location"`
} `xml:"data"`
}
// TestRepomdHashMatchesServedBytes asserts the exact invariant #117 violated:
// the sha256 advertised in repomd.xml equals the sha256 of the bytes the
// content-addressed serve* handler returns for the same repo state.
func TestRepomdHashMatchesServedBytes(t *testing.T) {
p := &Provider{}
reader := fakeRPMReader{metas: sampleMetas()}
serve := func(path string) *httptest.ResponseRecorder {
w := httptest.NewRecorder()
r := httptest.NewRequest(http.MethodGet, "/"+path, nil)
if !p.ServeLocalIndex(w, r, reader, "repo", path) {
t.Fatalf("ServeLocalIndex false for %q", path)
}
if w.Code != http.StatusOK {
t.Fatalf("%s: code %d", path, w.Code)
}
return w
}
var doc repomdDoc
if err := xml.Unmarshal(serve("repodata/repomd.xml").Body.Bytes(), &doc); err != nil {
t.Fatalf("parse repomd: %v", err)
}
if len(doc.Data) != 3 {
t.Fatalf("expected 3 <data> entries, got %d", len(doc.Data))
}
for _, d := range doc.Data {
// The advertised location is content-addressed: repodata/<sha256>-<type>.xml.gz.
body := serve("repodata/" + d.Location.Href[len("repodata/"):]).Body.Bytes()
got := sha256Hex(body)
if got != d.Checksum.Value {
t.Errorf("%s: repomd advertises %s but served bytes hash to %s (dnf would reject)",
d.Type, d.Checksum.Value, got)
}
if d.Location.Href != "repodata/"+d.Checksum.Value+"-"+d.Type+".xml.gz" {
t.Errorf("%s: location %q not addressed by its checksum %s", d.Type, d.Location.Href, d.Checksum.Value)
}
}
}
-162
View File
@@ -1,162 +0,0 @@
# Mirror-selection benchmarks
These benchmarks (`selection_bench_test.go`) isolate the **mirror load-balancing
selection overhead** — no network, no DB, no Redis. They build a zero-value
`Engine` and call `baseURLAttemptOrder` / `beginAttempt` / `endAttempt`
directly, the same way `leastconn_test.go` and `multibaseurl_test.go` do.
Goal: quantify how much latency the load-balancing strategy (`round_robin` vs
`least_conn`) adds versus a plain single-URL remote, and give a permanent
regression guard.
## Key context: selection is cache-miss-only
`baseURLAttemptOrder` is called from exactly three places — `headUpstream`,
`fetchFromUpstream`, and `checkUpstream` — all on the **upstream / cache-miss
path**. A cache hit returns `Source: "cache"` from `GetArtifact` / `store.Stat`
*before* any selection code runs. So none of the numbers below apply to the hot
cache-hit path: cache hits pay **zero** selection cost regardless of strategy.
The overhead here is paid once per upstream fetch, alongside a network round-trip
measured in milliseconds.
## How to run
```
go test -run=^$ -bench='BaseURLAttemptOrder|BeginEndAttempt' -benchmem \
-benchtime=1s -count=6 -cpu=8 ./internal/proxy/
```
## Results
Machine: AMD Ryzen 7 4700U (8 threads), linux/amd64, go1.26.5.
`-benchtime=1s -count=6`; figures below are the **median of 6 runs**.
### Sequential (single-goroutine)
| Benchmark | ns/op | B/op | allocs/op |
|----------------------------------|-------:|-----:|----------:|
| BaseURLAttemptOrder_SingleURL | ~133 | 16 | 1 |
| BaseURLAttemptOrder_RoundRobin/3 | ~462 | 120 | 4 |
| BaseURLAttemptOrder_RoundRobin/8 | ~682 | 280 | 4 |
| BaseURLAttemptOrder_LeastConn/3 | ~2690 | 474 | 22 |
| BaseURLAttemptOrder_LeastConn/8 | ~15200 | 2688 | 132 |
| BeginEndAttempt (gauge inc/dec) | ~514 | 104 | 4 |
### Parallel (`RunParallel`, GOMAXPROCS=8) — ns/op is wall-time across 8 cores
| Benchmark | ns/op | B/op | allocs/op |
|-------------------------------------------|------:|-----:|----------:|
| BaseURLAttemptOrder_RoundRobin_Parallel/3 | ~67.5 | 120 | 4 |
| BaseURLAttemptOrder_RoundRobin_Parallel/8 | ~130 | 280 | 4 |
| BaseURLAttemptOrder_LeastConn_Parallel/3 | ~292 | 474 | 22 |
| BaseURLAttemptOrder_LeastConn_Parallel/8 | ~1673 | 2688 | 132 |
| BeginEndAttempt_Parallel | ~71.5 | 104 | 4 |
## Reading the numbers
- **Single-URL is a near-no-op** (~133 ns, 1 alloc): the `len(urls) <= 1`
early return just returns the pool slice. Every non-mirrored remote takes this
path.
- **round_robin is cheap**: ~462 ns for a 3-mirror pool, ~682 ns for 8. Cost is
one atomic cursor increment plus building the rotated `[]string`. Allocs are
constant at 4 (the ordered slice + its backing string headers), size grows
with pool length.
- **least_conn is more expensive and scales super-linearly**: ~2.7 µs / 22
allocs at 3 mirrors, ~15 µs / 132 allocs at 8. The cost is the per-call
`sort.SliceStable`, whose comparator calls `inflightCounter` (a
`sync.Map.LoadOrStore` with a `remoteName\x00url` string-concat key plus a
speculative `new(atomic.Int64)`) O(n·log n) times. That is where the alloc
count and the time come from — not the sort itself. A future optimization
could snapshot each mirror's load once before sorting; out of scope for this
measurement PR.
- **beginAttempt/endAttempt** (~514 ns seq, ~72 ns parallel) is one
`LoadOrStore` + two atomic adds; it only runs for least_conn multi-mirror
remotes, once per upstream attempt.
- **Under concurrency the atomics/sync.Map do not collapse**: every parallel
variant reports *lower* ns/op than its sequential twin because work spreads
across 8 cores (RunParallel reports aggregate wall-time-per-op). No contention
cliff on the shared rrCounters cursor, the inflight `sync.Map`, or the
per-mirror `atomic.Int64` gauges.
## Verdict
At the per-request scale that matters (a cache-miss that is *already* doing a
multi-millisecond network fetch), even the worst case here — least_conn across 8
mirrors at ~15 µs — is <1% of a single upstream round-trip, and round_robin
(~0.5 µs) is negligible. The strategy adds no meaningful latency, and it adds
**exactly zero** to the cache-hit hot path because selection never runs there.
## Raw output (all 6 runs)
```
goos: linux
goarch: amd64
pkg: git.unkin.net/unkin/artifactapi/internal/proxy
cpu: AMD Ryzen 7 4700U with Radeon Graphics
BenchmarkBaseURLAttemptOrder_SingleURL-8 8635875 138.4 ns/op 16 B/op 1 allocs/op
BenchmarkBaseURLAttemptOrder_SingleURL-8 9673108 126.7 ns/op 16 B/op 1 allocs/op
BenchmarkBaseURLAttemptOrder_SingleURL-8 8001002 147.3 ns/op 16 B/op 1 allocs/op
BenchmarkBaseURLAttemptOrder_SingleURL-8 11303490 135.0 ns/op 16 B/op 1 allocs/op
BenchmarkBaseURLAttemptOrder_SingleURL-8 10125138 132.0 ns/op 16 B/op 1 allocs/op
BenchmarkBaseURLAttemptOrder_SingleURL-8 8025687 130.1 ns/op 16 B/op 1 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool3-8 2706013 453.2 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool3-8 2498718 444.4 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool3-8 2605516 471.6 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool3-8 2799928 487.6 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool3-8 2463375 418.6 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool3-8 2472265 474.3 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool8-8 1741789 689.4 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool8-8 1775467 602.1 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool8-8 1829398 688.4 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool8-8 1781149 679.3 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool8-8 1795680 599.4 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin/pool8-8 1739122 684.5 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool3-8 424184 2675 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool3-8 426796 2498 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool3-8 427116 2716 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool3-8 430540 2681 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool3-8 418156 2700 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool3-8 423009 2711 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool8-8 163642 14007 ns/op 2688 B/op 132 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool8-8 78501 15457 ns/op 2688 B/op 131 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool8-8 76380 14928 ns/op 2688 B/op 132 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool8-8 183634 16091 ns/op 2688 B/op 132 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool8-8 73809 15561 ns/op 2688 B/op 132 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn/pool8-8 74546 13962 ns/op 2688 B/op 132 allocs/op
BenchmarkBeginEndAttempt-8 2350348 519.7 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt-8 2321659 514.0 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt-8 2284635 438.0 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt-8 2287051 513.9 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt-8 2286481 520.3 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt-8 2837775 512.9 ns/op 104 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool3-8 16052568 67.78 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool3-8 17452791 66.07 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool3-8 17549858 69.25 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool3-8 18845167 64.55 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool3-8 16285608 69.81 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool3-8 17382639 67.12 ns/op 120 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool8-8 9734368 120.6 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool8-8 10154736 133.3 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool8-8 10061422 131.8 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool8-8 10212364 127.4 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool8-8 10259030 132.5 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel/pool8-8 10069576 122.4 ns/op 280 B/op 4 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool3-8 4288112 292.7 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool3-8 4009249 295.9 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool3-8 4176378 291.3 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool3-8 4104871 289.9 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool3-8 4245262 296.4 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool3-8 4079778 290.3 ns/op 474 B/op 22 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool8-8 748200 1636 ns/op 2688 B/op 132 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool8-8 763029 1653 ns/op 2688 B/op 132 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool8-8 663717 1772 ns/op 2688 B/op 132 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool8-8 739677 1676 ns/op 2688 B/op 132 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool8-8 763148 1669 ns/op 2688 B/op 132 allocs/op
BenchmarkBaseURLAttemptOrder_LeastConn_Parallel/pool8-8 610597 1684 ns/op 2688 B/op 132 allocs/op
BenchmarkBeginEndAttempt_Parallel-8 17266058 69.46 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt_Parallel-8 17151303 72.05 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt_Parallel-8 16919542 74.17 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt_Parallel-8 16948015 72.49 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt_Parallel-8 16918693 69.52 ns/op 104 B/op 4 allocs/op
BenchmarkBeginEndAttempt_Parallel-8 17376012 70.99 ns/op 104 B/op 4 allocs/op
```
-197
View File
@@ -10,10 +10,7 @@ import (
"io"
"log/slog"
"net/http"
"sort"
"strings"
"sync"
"sync/atomic"
"time"
"git.unkin.net/unkin/artifactapi/internal/cache"
@@ -38,15 +35,6 @@ type Engine struct {
cas *storage.CAS
circuit *CircuitBreaker
accessLog chan database.AccessLogEntry
// rrCounters holds a per-remote round-robin cursor (remoteName ->
// *atomic.Uint64) used to rotate the starting mirror across upstream base
// URLs. Distribution is per-replica and approximate, which is fine.
rrCounters sync.Map
// inflight holds a per-remote, per-upstream-URL in-flight request gauge
// (key "remoteName\x00baseURL" -> *atomic.Int64) used by the least_conn
// mirror strategy to prefer the mirror currently handling the fewest
// requests. Per-replica and approximate, which is fine.
inflight sync.Map
}
func NewEngine(db *database.DB, c *cache.Redis, s *storage.S3) *Engine {
@@ -234,32 +222,7 @@ func (e *Engine) Head(ctx context.Context, remote models.Remote, path string, pr
return e.headUpstream(ctx, remote, path, prov)
}
// headUpstream issues an upstream HEAD, load-balancing across the remote's base
// URLs and failing over to the next mirror on a network error or 5xx.
func (e *Engine) headUpstream(ctx context.Context, remote models.Remote, path string, prov provider.Provider) (*HeadResult, error) {
order := e.baseURLAttemptOrder(remote)
if len(order) == 0 {
return nil, &ProxyError{Status: http.StatusBadGateway, Message: "no upstream base_url configured"}
}
var lastErr error
for i, url := range order {
ctr := e.beginAttempt(remote, url)
result, err := e.headUpstreamOnce(ctx, withBaseURL(remote, url), path, prov)
endAttempt(ctr)
if err == nil {
return result, nil
}
lastErr = err
if i < len(order)-1 && shouldFailover(err) {
slog.Warn("upstream HEAD failed, failing over", "remote", remote.Name, "base_url", url, "error", err)
continue
}
return nil, err
}
return nil, lastErr
}
func (e *Engine) headUpstreamOnce(ctx context.Context, remote models.Remote, path string, prov provider.Provider) (*HeadResult, error) {
url := prov.UpstreamURL(remote, path)
authHeaders, err := prov.AuthHeaders(ctx, remote)
@@ -314,33 +277,7 @@ func (e *Engine) headUpstreamOnce(ctx context.Context, remote models.Remote, pat
return &HeadResult{ContentType: contentType, Size: resp.ContentLength, Source: "remote"}, nil
}
// fetchFromUpstream fetches an artifact from upstream, load-balancing across the
// remote's base URLs and failing over to the next mirror on a network error or
// 5xx before returning an error.
func (e *Engine) fetchFromUpstream(ctx context.Context, remote models.Remote, path string, prov provider.Provider, class Classification, ttl time.Duration, clientHeaders http.Header) (*FetchResult, error) {
order := e.baseURLAttemptOrder(remote)
if len(order) == 0 {
return nil, &ProxyError{Status: http.StatusBadGateway, Message: "no upstream base_url configured"}
}
var lastErr error
for i, url := range order {
ctr := e.beginAttempt(remote, url)
result, err := e.fetchFromUpstreamOnce(ctx, withBaseURL(remote, url), path, prov, class, ttl, clientHeaders)
endAttempt(ctr)
if err == nil {
return result, nil
}
lastErr = err
if i < len(order)-1 && shouldFailover(err) {
slog.Warn("upstream fetch failed, failing over", "remote", remote.Name, "base_url", url, "error", err)
continue
}
return nil, err
}
return nil, lastErr
}
func (e *Engine) fetchFromUpstreamOnce(ctx context.Context, remote models.Remote, path string, prov provider.Provider, class Classification, ttl time.Duration, clientHeaders http.Header) (*FetchResult, error) {
url := prov.UpstreamURL(remote, path)
authHeaders, err := prov.AuthHeaders(ctx, remote)
@@ -517,33 +454,7 @@ func (e *Engine) serveFromStore(ctx context.Context, remote models.Remote, path
}, nil
}
// checkUpstream issues a conditional upstream HEAD (If-None-Match), load
// balancing across the remote's base URLs and failing over to the next mirror on
// a network error or 5xx.
func (e *Engine) checkUpstream(ctx context.Context, remote models.Remote, path, etag string, prov provider.Provider) (bool, error) {
order := e.baseURLAttemptOrder(remote)
if len(order) == 0 {
return false, &ProxyError{Status: http.StatusBadGateway, Message: "no upstream base_url configured"}
}
var lastErr error
for i, url := range order {
ctr := e.beginAttempt(remote, url)
notModified, err := e.checkUpstreamOnce(ctx, withBaseURL(remote, url), path, etag, prov)
endAttempt(ctr)
if err == nil {
return notModified, nil
}
lastErr = err
if i < len(order)-1 && shouldFailover(err) {
slog.Warn("upstream revalidation failed, failing over", "remote", remote.Name, "base_url", url, "error", err)
continue
}
return false, err
}
return false, lastErr
}
func (e *Engine) checkUpstreamOnce(ctx context.Context, remote models.Remote, path, etag string, prov provider.Provider) (bool, error) {
url := prov.UpstreamURL(remote, path)
req, err := http.NewRequestWithContext(ctx, http.MethodHead, url, nil)
@@ -738,111 +649,3 @@ func isNetworkError(err error) bool {
var ue *UpstreamError
return errors.As(err, &ue)
}
// baseURLAttemptOrder returns the ordered upstream base URLs to try for a single
// request, drawn from the remote's pool ([base_url] + mirrorlist). A multi-mirror
// remote starts at a strategy-chosen position and advances linearly for
// failover; a remote with no mirrorlist yields exactly [base_url], preserving the
// original single-attempt behavior. The default (round_robin) rotates the
// starting mirror; least_conn starts with the mirror handling the fewest
// in-flight requests. Failover order after the first pick is unchanged.
func (e *Engine) baseURLAttemptOrder(remote models.Remote) []string {
urls := remote.UpstreamPool()
if len(urls) <= 1 {
return urls
}
// Rotate by the round-robin cursor first so equal-load mirrors still spread
// evenly; least_conn then stable-sorts this rotation by in-flight count.
v, _ := e.rrCounters.LoadOrStore(remote.Name, new(atomic.Uint64))
start := int(v.(*atomic.Uint64).Add(1) - 1)
ordered := make([]string, len(urls))
for i := range urls {
ordered[i] = urls[(start+i)%len(urls)]
}
if remote.MirrorStrategy == models.MirrorStrategyLeastConn {
// Snapshot each mirror's in-flight count once, then sort the snapshot.
// Reading the gauge inside the comparator would repeat an allocating
// sync.Map lookup on every comparison (O(n log n) lookups); this is O(n).
snap := make([]inflightSnapshot, len(ordered))
for i, url := range ordered {
snap[i] = inflightSnapshot{url: url, count: e.inflightCount(remote.Name, url)}
}
sort.SliceStable(snap, func(a, b int) bool {
return snap[a].count < snap[b].count
})
for i := range snap {
ordered[i] = snap[i].url
}
}
return ordered
}
// inflightSnapshot pairs a mirror URL with its sampled in-flight count so the
// least_conn sort compares plain ints instead of re-reading the gauge.
type inflightSnapshot struct {
url string
count int64
}
// inflightCount reads the in-flight request gauge for a given (remote, upstream
// URL) without creating it, returning 0 when the counter is absent. This keeps
// the selection read path allocation-free (plain Load, no LoadOrStore).
func (e *Engine) inflightCount(remoteName, url string) int64 {
v, ok := e.inflight.Load(remoteName + "\x00" + url)
if !ok {
return 0
}
return v.(*atomic.Int64).Load()
}
// inflightCounter returns the shared in-flight request gauge for a given
// (remote, upstream URL), creating it on first use.
func (e *Engine) inflightCounter(remoteName, url string) *atomic.Int64 {
v, _ := e.inflight.LoadOrStore(remoteName+"\x00"+url, new(atomic.Int64))
return v.(*atomic.Int64)
}
// beginAttempt increments the in-flight gauge for a least_conn multi-mirror
// remote before an upstream call and returns the counter to release; it is a
// no-op (returns nil) for round-robin remotes and single-URL pools.
func (e *Engine) beginAttempt(remote models.Remote, url string) *atomic.Int64 {
if remote.MirrorStrategy != models.MirrorStrategyLeastConn {
return nil
}
if len(remote.UpstreamPool()) <= 1 {
return nil
}
ctr := e.inflightCounter(remote.Name, url)
ctr.Add(1)
return ctr
}
// endAttempt decrements a gauge returned by beginAttempt, tolerating nil.
func endAttempt(ctr *atomic.Int64) {
if ctr != nil {
ctr.Add(-1)
}
}
// withBaseURL narrows a remote's active BaseURL to a single selected mirror so
// providers (UpstreamURL/AuthHeaders/RewriteResponse) operate on exactly that
// upstream for this attempt.
func withBaseURL(remote models.Remote, url string) models.Remote {
remote.BaseURL = url
remote.Mirrorlist = nil
return remote
}
// shouldFailover reports whether an upstream attempt error is worth retrying
// against the next mirror: network errors/timeouts and upstream 5xx responses.
// Definitive statuses (404/403/401/...) are returned to the caller unchanged.
func shouldFailover(err error) bool {
if isNetworkError(err) {
return true
}
var pe *ProxyError
if errors.As(err, &pe) {
return pe.Status >= 500
}
return false
}
-148
View File
@@ -1,148 +0,0 @@
package proxy
import (
"testing"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// baseURLAttemptOrder and the in-flight gauge only touch the engine's sync.Map
// fields, so these tests run against a zero-value Engine without a DB/S3/redis
// stack and can drive the gauge deterministically.
// TestLeastConnPicksLeastLoaded pre-loads one mirror's in-flight gauge and
// asserts a least_conn remote starts its attempt order with the idle mirror.
func TestLeastConnPicksLeastLoaded(t *testing.T) {
e := &Engine{}
r := models.Remote{
Name: "lc",
BaseURL: "https://a.example",
Mirrorlist: []string{"https://b.example"},
MirrorStrategy: models.MirrorStrategyLeastConn,
}
// Make A appear busy: least_conn must prefer B regardless of RR rotation.
e.inflightCounter(r.Name, "https://a.example").Add(3)
for i := 0; i < 5; i++ {
order := e.baseURLAttemptOrder(r)
if len(order) != 2 {
t.Fatalf("attempt %d: order len = %d, want 2", i, len(order))
}
if order[0] != "https://b.example" {
t.Fatalf("attempt %d: least_conn started with %q, want idle mirror https://b.example", i, order[0])
}
}
// Once B is the busier mirror, the starting pick flips to A.
e.inflightCounter(r.Name, "https://b.example").Add(10)
if order := e.baseURLAttemptOrder(r); order[0] != "https://a.example" {
t.Fatalf("after loading B, least_conn started with %q, want https://a.example", order[0])
}
}
// TestLeastConnStableTieBreak asserts that when every mirror carries equal
// in-flight load, least_conn falls back to the round-robin rotation: the
// snapshot sort is stable, so tied mirrors keep the RR-rotated order and the
// starting pick advances across the whole pool on successive calls.
func TestLeastConnStableTieBreak(t *testing.T) {
e := &Engine{}
r := models.Remote{
Name: "lc-tie",
BaseURL: "https://a.example",
Mirrorlist: []string{"https://b.example", "https://c.example"},
MirrorStrategy: models.MirrorStrategyLeastConn,
}
pool := r.UpstreamPool()
// Equal (zero) load on every mirror: order must equal the RR rotation.
starts := map[string]int{}
for i := 0; i < len(pool); i++ {
order := e.baseURLAttemptOrder(r)
if len(order) != len(pool) {
t.Fatalf("attempt %d: order len = %d, want %d", i, len(order), len(pool))
}
// A stable sort of an all-tied slice is a pure RR rotation: for the
// call whose cursor selects start s, order must be pool rotated by s.
start := indexOf(pool, order[0])
for j := range order {
if want := pool[(start+j)%len(pool)]; order[j] != want {
t.Fatalf("attempt %d: order[%d] = %q, want RR-rotated %q", i, j, order[j], want)
}
}
starts[order[0]]++
}
if len(starts) != len(pool) {
t.Fatalf("tied least_conn did not rotate across the whole pool: %v", starts)
}
}
func indexOf(s []string, v string) int {
for i := range s {
if s[i] == v {
return i
}
}
return -1
}
// TestRoundRobinDefaultUnchanged asserts an unset strategy still rotates the
// starting mirror across the pool and ignores the in-flight gauge.
func TestRoundRobinDefaultUnchanged(t *testing.T) {
e := &Engine{}
r := models.Remote{
Name: "rr",
BaseURL: "https://a.example",
Mirrorlist: []string{"https://b.example"},
}
// Even with A heavily loaded, round-robin must still rotate (not avoid A).
e.inflightCounter(r.Name, "https://a.example").Add(100)
starts := map[string]int{}
for i := 0; i < 4; i++ {
starts[e.baseURLAttemptOrder(r)[0]]++
}
if starts["https://a.example"] == 0 || starts["https://b.example"] == 0 {
t.Fatalf("round-robin did not rotate starting mirror: %v", starts)
}
}
// TestLeastConnSingleURLNoOp asserts a single-URL pool yields exactly [base_url]
// and beginAttempt is a no-op there and for round-robin remotes.
func TestLeastConnSingleURLNoOp(t *testing.T) {
e := &Engine{}
solo := models.Remote{Name: "solo", BaseURL: "https://a.example", MirrorStrategy: models.MirrorStrategyLeastConn}
if order := e.baseURLAttemptOrder(solo); len(order) != 1 || order[0] != "https://a.example" {
t.Fatalf("single-url order = %v, want [base_url]", order)
}
if ctr := e.beginAttempt(solo, "https://a.example"); ctr != nil {
t.Fatal("beginAttempt on single-url pool should be a no-op (nil)")
}
rr := models.Remote{Name: "rr2", BaseURL: "https://a.example", Mirrorlist: []string{"https://b.example"}}
if ctr := e.beginAttempt(rr, "https://a.example"); ctr != nil {
t.Fatal("beginAttempt on round-robin remote should be a no-op (nil)")
}
}
// TestBeginEndAttemptGauge asserts the gauge increments on begin and returns to
// zero after endAttempt, so it tracks live in-flight requests.
func TestBeginEndAttemptGauge(t *testing.T) {
e := &Engine{}
r := models.Remote{
Name: "g",
BaseURL: "https://a.example",
Mirrorlist: []string{"https://b.example"},
MirrorStrategy: models.MirrorStrategyLeastConn,
}
c1 := e.beginAttempt(r, "https://a.example")
c2 := e.beginAttempt(r, "https://a.example")
if got := e.inflightCounter(r.Name, "https://a.example").Load(); got != 2 {
t.Fatalf("gauge after two begins = %d, want 2", got)
}
endAttempt(c1)
endAttempt(c2)
if got := e.inflightCounter(r.Name, "https://a.example").Load(); got != 0 {
t.Fatalf("gauge after matching ends = %d, want 0", got)
}
endAttempt(nil) // tolerated
}
-188
View File
@@ -1,188 +0,0 @@
package proxy
import (
"context"
"fmt"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// TestFetchMultiBaseURLRoundRobin drives distinct artifact paths through a
// remote configured with two upstreams and asserts both receive traffic.
func TestFetchMultiBaseURLRoundRobin(t *testing.T) {
requireStack(t)
ctx := context.Background()
var hitsA, hitsB atomic.Int64
upA := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
hitsA.Add(1)
w.Write([]byte("A"))
}))
defer upA.Close()
upB := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
hitsB.Add(1)
w.Write([]byte("B"))
}))
defer upB.Close()
r := seed(t, models.Remote{
Name: "eng-rr",
PackageType: models.PackageGeneric,
RepoType: models.RepoTypeRemote,
BaseURL: upA.URL,
Mirrorlist: []string{upB.URL},
StaleOnError: true,
})
p := prov(t, models.PackageGeneric)
const n = 10
for i := 0; i < n; i++ {
res, err := testEngine.Fetch(ctx, r, fmt.Sprintf("rr-%d.bin", i), p)
if err != nil {
t.Fatalf("fetch %d: %v", i, err)
}
res.Reader.Close()
}
if hitsA.Load() == 0 || hitsB.Load() == 0 {
t.Fatalf("round-robin did not spread across both upstreams: A=%d B=%d", hitsA.Load(), hitsB.Load())
}
if total := hitsA.Load() + hitsB.Load(); total != n {
t.Fatalf("expected %d upstream hits total, got %d (A=%d B=%d)", n, total, hitsA.Load(), hitsB.Load())
}
}
// TestFetchMultiBaseURLFailover asserts that a dead/erroring primary mirror
// transparently fails over to a healthy secondary, for both a 5xx primary and a
// network-unreachable primary.
func TestFetchMultiBaseURLFailover(t *testing.T) {
requireStack(t)
ctx := context.Background()
var hitsB atomic.Int64
upB := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
hitsB.Add(1)
w.Write([]byte("served-by-B"))
}))
defer upB.Close()
up500 := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
}))
defer up500.Close()
p := prov(t, models.PackageGeneric)
// Primary returns 5xx: every request must still succeed via the secondary.
r5xx := seed(t, models.Remote{
Name: "eng-failover-5xx",
PackageType: models.PackageGeneric,
RepoType: models.RepoTypeRemote,
BaseURL: up500.URL,
Mirrorlist: []string{upB.URL},
})
for i := 0; i < 6; i++ {
res, err := testEngine.Fetch(ctx, r5xx, fmt.Sprintf("fo5-%d.bin", i), p)
if err != nil {
t.Fatalf("5xx failover fetch %d: %v", i, err)
}
if got := readAll(t, res); got != "served-by-B" {
t.Fatalf("5xx failover fetch %d body=%q, want served-by-B", i, got)
}
}
// Primary is network-unreachable: failover must still reach the secondary.
rNet := seed(t, models.Remote{
Name: "eng-failover-net",
PackageType: models.PackageGeneric,
RepoType: models.RepoTypeRemote,
BaseURL: "http://127.0.0.1:1",
Mirrorlist: []string{upB.URL},
})
res, err := testEngine.Fetch(ctx, rNet, "fonet.bin", p)
if err != nil {
t.Fatalf("network failover fetch: %v", err)
}
if got := readAll(t, res); got != "served-by-B" {
t.Fatalf("network failover body=%q, want served-by-B", got)
}
if hitsB.Load() == 0 {
t.Fatal("secondary upstream never served during failover")
}
}
// TestFetchDefinitiveStatusNoFailover asserts a definitive 404 from the first
// mirror is returned as-is (not failed over): a missing artifact is not a mirror
// outage. The remote is fresh so its round-robin cursor starts at index 0.
func TestFetchDefinitiveStatusNoFailover(t *testing.T) {
requireStack(t)
ctx := context.Background()
var hitsB atomic.Int64
up404 := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.NotFound(w, r)
}))
defer up404.Close()
upB := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
hitsB.Add(1)
w.Write([]byte("B"))
}))
defer upB.Close()
r := seed(t, models.Remote{
Name: "eng-no-failover-404",
PackageType: models.PackageGeneric,
RepoType: models.RepoTypeRemote,
BaseURL: up404.URL,
Mirrorlist: []string{upB.URL},
})
_, err := testEngine.Fetch(ctx, r, "missing.bin", prov(t, models.PackageGeneric))
var pe *ProxyError
if err == nil || !asProxyError(err, &pe) || pe.Status != http.StatusNotFound {
t.Fatalf("expected 404 ProxyError without failover, got %v", err)
}
if hitsB.Load() != 0 {
t.Fatalf("404 from primary must not fail over, but secondary was hit %d times", hitsB.Load())
}
}
// TestFetchSingleBaseURLUnchanged asserts a single-URL remote behaves exactly as
// before: one healthy URL succeeds, and one dead URL errors with no failover.
func TestFetchSingleBaseURLUnchanged(t *testing.T) {
requireStack(t)
ctx := context.Background()
upB := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Write([]byte("solo"))
}))
defer upB.Close()
p := prov(t, models.PackageGeneric)
rOK := seed(t, models.Remote{
Name: "eng-solo",
PackageType: models.PackageGeneric,
RepoType: models.RepoTypeRemote,
BaseURL: upB.URL,
})
res, err := testEngine.Fetch(ctx, rOK, "solo.bin", p)
if err != nil {
t.Fatalf("single-url fetch: %v", err)
}
if got := readAll(t, res); got != "solo" {
t.Fatalf("single-url body=%q, want solo", got)
}
rDead := seed(t, models.Remote{
Name: "eng-solo-dead",
PackageType: models.PackageGeneric,
RepoType: models.RepoTypeRemote,
BaseURL: "http://127.0.0.1:1",
})
if _, err := testEngine.Fetch(ctx, rDead, "x.bin", p); err == nil {
t.Fatal("single dead upstream should error, not succeed")
}
}
-156
View File
@@ -1,156 +0,0 @@
package proxy
import (
"fmt"
"testing"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// These benchmarks isolate the mirror-selection overhead only: they construct a
// zero-value Engine (no DB/S3/redis) and call baseURLAttemptOrder /
// beginAttempt / endAttempt directly, mirroring leastconn_test.go and
// multibaseurl_test.go. This quantifies how much latency the load-balancing
// strategy (round_robin vs least_conn) adds versus a single-URL remote. Note
// that in the live proxy this selection runs only on the cache-miss/upstream
// path; a cache hit never calls it.
// mirrorPool builds a remote with n upstreams (base_url + n-1 mirrorlist
// entries) under the given strategy.
func mirrorPool(name, strategy string, n int) models.Remote {
r := models.Remote{
Name: name,
BaseURL: "https://mirror0.example/repo",
MirrorStrategy: strategy,
}
for i := 1; i < n; i++ {
r.Mirrorlist = append(r.Mirrorlist, fmt.Sprintf("https://mirror%d.example/repo", i))
}
return r
}
// skewInflight sets an ascending in-flight load across the pool so least_conn's
// stable sort has real work to do (mirror0 busiest, last mirror idle).
func skewInflight(e *Engine, r models.Remote) {
pool := r.UpstreamPool()
for i, u := range pool {
e.inflightCounter(r.Name, u).Add(int64(len(pool) - i))
}
}
// BenchmarkBaseURLAttemptOrder_SingleURL measures the early-return no-op path
// (pool of 1): the branch that preserves original single-attempt behavior and
// must add effectively zero overhead. This is the same code the cache-miss path
// takes for every non-mirrored remote.
func BenchmarkBaseURLAttemptOrder_SingleURL(b *testing.B) {
e := &Engine{}
r := models.Remote{Name: "solo", BaseURL: "https://mirror0.example/repo"}
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
_ = e.baseURLAttemptOrder(r)
}
}
// BenchmarkBaseURLAttemptOrder_RoundRobin measures the default strategy: rotate
// the starting mirror by an atomic cursor and materialize the ordered slice. No
// in-flight sort.
func BenchmarkBaseURLAttemptOrder_RoundRobin(b *testing.B) {
for _, n := range []int{3, 8} {
b.Run(fmt.Sprintf("pool%d", n), func(b *testing.B) {
e := &Engine{}
r := mirrorPool("rr", models.MirrorStrategyRoundRobin, n)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
_ = e.baseURLAttemptOrder(r)
}
})
}
}
// BenchmarkBaseURLAttemptOrder_LeastConn measures the least_conn strategy: RR
// rotation plus a stable sort of the pool by the atomic in-flight gauges. Skew
// is preloaded so the sort compares distinct loads.
func BenchmarkBaseURLAttemptOrder_LeastConn(b *testing.B) {
for _, n := range []int{3, 8} {
b.Run(fmt.Sprintf("pool%d", n), func(b *testing.B) {
e := &Engine{}
r := mirrorPool("lc", models.MirrorStrategyLeastConn, n)
skewInflight(e, r)
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
_ = e.baseURLAttemptOrder(r)
}
})
}
}
// BenchmarkBeginEndAttempt measures the gauge inc/dec pair that brackets each
// least_conn upstream attempt (LoadOrStore + atomic add, then atomic add back).
func BenchmarkBeginEndAttempt(b *testing.B) {
e := &Engine{}
r := mirrorPool("g", models.MirrorStrategyLeastConn, 3)
url := r.UpstreamPool()[0]
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
ctr := e.beginAttempt(r, url)
endAttempt(ctr)
}
}
// BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel surfaces atomic-cursor
// contention on the shared rrCounters entry under concurrent selection.
func BenchmarkBaseURLAttemptOrder_RoundRobin_Parallel(b *testing.B) {
for _, n := range []int{3, 8} {
b.Run(fmt.Sprintf("pool%d", n), func(b *testing.B) {
e := &Engine{}
r := mirrorPool("rrp", models.MirrorStrategyRoundRobin, n)
b.ReportAllocs()
b.ResetTimer()
b.RunParallel(func(pb *testing.PB) {
for pb.Next() {
_ = e.baseURLAttemptOrder(r)
}
})
})
}
}
// BenchmarkBaseURLAttemptOrder_LeastConn_Parallel surfaces sync.Map read
// contention on the in-flight gauges plus the per-call sort under concurrency.
func BenchmarkBaseURLAttemptOrder_LeastConn_Parallel(b *testing.B) {
for _, n := range []int{3, 8} {
b.Run(fmt.Sprintf("pool%d", n), func(b *testing.B) {
e := &Engine{}
r := mirrorPool("lcp", models.MirrorStrategyLeastConn, n)
skewInflight(e, r)
b.ReportAllocs()
b.ResetTimer()
b.RunParallel(func(pb *testing.PB) {
for pb.Next() {
_ = e.baseURLAttemptOrder(r)
}
})
})
}
}
// BenchmarkBeginEndAttempt_Parallel exercises the gauge inc/dec pair under
// concurrency: all goroutines hammer the same atomic.Int64, the realistic
// hot-mirror case, to surface counter contention.
func BenchmarkBeginEndAttempt_Parallel(b *testing.B) {
e := &Engine{}
r := mirrorPool("gp", models.MirrorStrategyLeastConn, 3)
url := r.UpstreamPool()[0]
b.ReportAllocs()
b.ResetTimer()
b.RunParallel(func(pb *testing.PB) {
for pb.Next() {
ctr := e.beginAttempt(r, url)
endAttempt(ctr)
}
})
}
+1 -1
View File
@@ -175,7 +175,7 @@ func (s *Server) routes() chi.Router {
r.Mount("/api/v1", proxyHandler.Routes())
r.Mount("/v2", proxyHandler.DockerV2Routes())
remotesHandler := v2.NewRemotesHandler(s.db, s.cache, map[models.PackageType]v2.Primer{
remotesHandler := v2.NewRemotesHandler(s.db, map[models.PackageType]v2.Primer{
models.PackageGitHubRPM: s.syncer,
models.PackageGitHubDeb: s.debSyncer,
models.PackageGitHubAlpine: s.alpineSyncer,
+3 -80
View File
@@ -2,7 +2,6 @@ package models
import (
"fmt"
"net/url"
"regexp"
"time"
)
@@ -40,17 +39,9 @@ type Remote struct {
PackageType PackageType `json:"package_type"`
RepoType RepoType `json:"repo_type"`
BaseURL string `json:"base_url"`
// Mirrorlist holds additional upstream mirror base URLs. The effective
// upstream pool is [base_url] + mirrorlist, load-balanced round-robin with
// failover by the proxy engine. Only valid on remote rpm/deb/apk repos.
Mirrorlist []string `json:"mirrorlist,omitempty"`
// MirrorStrategy selects how the proxy engine picks the starting upstream
// from the pool: round_robin (default/empty) rotates, least_conn favors the
// mirror with the fewest in-flight requests. Failover order is unchanged.
MirrorStrategy string `json:"mirror_strategy,omitempty"`
Description string `json:"description,omitempty"`
Username string `json:"-"`
Password string `json:"-"`
Description string `json:"description,omitempty"`
Username string `json:"-"`
Password string `json:"-"`
ImmutableTTL int `json:"immutable_ttl"`
MutableTTL int `json:"mutable_ttl"`
@@ -81,74 +72,6 @@ type Remote struct {
UpdatedAt time.Time `json:"updated_at"`
}
// Mirror balancing strategies for MirrorStrategy. An empty value is treated as
// round_robin, so existing remotes keep their current behavior.
const (
MirrorStrategyRoundRobin = "round_robin"
MirrorStrategyLeastConn = "least_conn"
)
// mirrorlistPackageTypes are the package types for which a mirrorlist is
// allowed: OS package repos (rpm, deb, apk/alpine) that fetch many small files
// and benefit most from mirror load-balancing and failover.
var mirrorlistPackageTypes = map[PackageType]bool{
PackageRPM: true,
PackageDeb: true,
PackageAlpine: true,
}
// UpstreamPool returns the ordered upstream base URLs for this remote: the
// primary base_url first, followed by any mirrorlist entries. The proxy engine
// load-balances round-robin across the pool and fails over between them.
func (r Remote) UpstreamPool() []string {
pool := make([]string, 0, 1+len(r.Mirrorlist))
if r.BaseURL != "" {
pool = append(pool, r.BaseURL)
}
pool = append(pool, r.Mirrorlist...)
return pool
}
// ValidateMirrorlist enforces that a mirrorlist is only configured on remote
// rpm/deb/apk repositories and that every entry is a parseable http/https URL.
func (r *Remote) ValidateMirrorlist() error {
if len(r.Mirrorlist) == 0 {
return nil
}
if r.RepoType != RepoTypeRemote {
return fmt.Errorf("mirrorlist is only allowed on remote repositories")
}
if !mirrorlistPackageTypes[r.PackageType] {
return fmt.Errorf("mirrorlist is only allowed for rpm, deb and alpine package types, not %q", r.PackageType)
}
for _, u := range r.Mirrorlist {
parsed, err := url.ParseRequestURI(u)
if err != nil {
return fmt.Errorf("invalid mirrorlist url %q: %w", u, err)
}
if parsed.Scheme != "http" && parsed.Scheme != "https" {
return fmt.Errorf("mirrorlist url %q must be http or https", u)
}
}
return nil
}
// ValidateMirrorStrategy enforces that mirror_strategy is one of the allowed
// values and that a non-default strategy (least_conn) is only set alongside a
// non-empty mirrorlist, where balancing is meaningful. An empty strategy is
// accepted and behaves as round_robin.
func (r *Remote) ValidateMirrorStrategy() error {
switch r.MirrorStrategy {
case "", MirrorStrategyRoundRobin, MirrorStrategyLeastConn:
default:
return fmt.Errorf("invalid mirror_strategy %q: must be %q or %q", r.MirrorStrategy, MirrorStrategyRoundRobin, MirrorStrategyLeastConn)
}
if r.MirrorStrategy == MirrorStrategyLeastConn && len(r.Mirrorlist) == 0 {
return fmt.Errorf("mirror_strategy %q requires a non-empty mirrorlist", r.MirrorStrategy)
}
return nil
}
// ValidatePatterns ensures every configured regex compiles. Storing an
// invalid pattern would otherwise be silently dropped at match time, which
// for the blocklist is a fail-open: a mistyped deny rule becomes a no-op.
+1 -107
View File
@@ -1,10 +1,6 @@
package models
import (
"encoding/json"
"strings"
"testing"
)
import "testing"
func TestRemote_ValidatePatterns(t *testing.T) {
valid := &Remote{
@@ -21,105 +17,3 @@ func TestRemote_ValidatePatterns(t *testing.T) {
t.Fatal("expected error for invalid blocklist regex, got nil")
}
}
func TestRemoteMirrorlistJSON(t *testing.T) {
// base_url stays a plain string; mirrorlist round-trips as an array.
var r Remote
body := `{"name":"x","package_type":"rpm","repo_type":"remote","base_url":"https://a.example","mirrorlist":["https://b.example","https://c.example"]}`
if err := json.Unmarshal([]byte(body), &r); err != nil {
t.Fatal(err)
}
if r.BaseURL != "https://a.example" {
t.Errorf("BaseURL = %q, want https://a.example", r.BaseURL)
}
if len(r.Mirrorlist) != 2 || r.Mirrorlist[0] != "https://b.example" || r.Mirrorlist[1] != "https://c.example" {
t.Errorf("Mirrorlist = %v, want two entries", r.Mirrorlist)
}
out, err := json.Marshal(r)
if err != nil {
t.Fatal(err)
}
if !strings.Contains(string(out), `"base_url":"https://a.example"`) {
t.Errorf("marshal lost base_url: %s", out)
}
if !strings.Contains(string(out), `"mirrorlist":["https://b.example","https://c.example"]`) {
t.Errorf("marshal lost mirrorlist: %s", out)
}
}
func TestRemoteMirrorlistOmitempty(t *testing.T) {
out, err := json.Marshal(Remote{Name: "x", PackageType: PackageRPM, RepoType: RepoTypeRemote, BaseURL: "https://a.example"})
if err != nil {
t.Fatal(err)
}
if strings.Contains(string(out), "mirrorlist") {
t.Errorf("empty mirrorlist should be omitted: %s", out)
}
}
func TestUpstreamPool(t *testing.T) {
// base_url first, then mirrorlist.
r := Remote{BaseURL: "https://a.example", Mirrorlist: []string{"https://b.example", "https://c.example"}}
pool := r.UpstreamPool()
want := []string{"https://a.example", "https://b.example", "https://c.example"}
if strings.Join(pool, ",") != strings.Join(want, ",") {
t.Errorf("UpstreamPool = %v, want %v", pool, want)
}
// No mirrorlist ⇒ pool is just [base_url].
solo := Remote{BaseURL: "https://a.example"}
if got := solo.UpstreamPool(); len(got) != 1 || got[0] != "https://a.example" {
t.Errorf("solo UpstreamPool = %v, want [base_url]", got)
}
}
func TestValidateMirrorStrategy(t *testing.T) {
ml := []string{"https://m.example"}
cases := []struct {
name string
remote Remote
wantErr bool
}{
{"empty defaults ok", Remote{Mirrorlist: ml}, false},
{"explicit round_robin ok", Remote{MirrorStrategy: MirrorStrategyRoundRobin, Mirrorlist: ml}, false},
{"round_robin without mirrorlist ok", Remote{MirrorStrategy: MirrorStrategyRoundRobin}, false},
{"least_conn with mirrorlist ok", Remote{MirrorStrategy: MirrorStrategyLeastConn, Mirrorlist: ml}, false},
{"least_conn without mirrorlist rejected", Remote{MirrorStrategy: MirrorStrategyLeastConn}, true},
{"unknown strategy rejected", Remote{MirrorStrategy: "random", Mirrorlist: ml}, true},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
err := tc.remote.ValidateMirrorStrategy()
if (err != nil) != tc.wantErr {
t.Errorf("ValidateMirrorStrategy() err = %v, wantErr = %v", err, tc.wantErr)
}
})
}
}
func TestValidateMirrorlist(t *testing.T) {
cases := []struct {
name string
remote Remote
wantErr bool
}{
{"empty is ok on anything", Remote{RepoType: RepoTypeRemote, PackageType: PackageGeneric}, false},
{"rpm remote ok", Remote{RepoType: RepoTypeRemote, PackageType: PackageRPM, Mirrorlist: []string{"https://m.example"}}, false},
{"deb remote ok", Remote{RepoType: RepoTypeRemote, PackageType: PackageDeb, Mirrorlist: []string{"http://m.example"}}, false},
{"alpine remote ok", Remote{RepoType: RepoTypeRemote, PackageType: PackageAlpine, Mirrorlist: []string{"https://m.example"}}, false},
{"generic remote rejected", Remote{RepoType: RepoTypeRemote, PackageType: PackageGeneric, Mirrorlist: []string{"https://m.example"}}, true},
{"docker remote rejected", Remote{RepoType: RepoTypeRemote, PackageType: PackageDocker, Mirrorlist: []string{"https://m.example"}}, true},
{"local rpm rejected", Remote{RepoType: RepoTypeLocal, PackageType: PackageRPM, Mirrorlist: []string{"https://m.example"}}, true},
{"bad scheme rejected", Remote{RepoType: RepoTypeRemote, PackageType: PackageRPM, Mirrorlist: []string{"ftp://m.example"}}, true},
{"unparseable rejected", Remote{RepoType: RepoTypeRemote, PackageType: PackageRPM, Mirrorlist: []string{"://nope"}}, true},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
err := tc.remote.ValidateMirrorlist()
if (err != nil) != tc.wantErr {
t.Errorf("ValidateMirrorlist() err = %v, wantErr = %v", err, tc.wantErr)
}
})
}
}
+3 -13
View File
@@ -17,8 +17,8 @@ cleanup() {
}
trap cleanup EXIT
echo "==> building and starting stack (postgres, redis, minio, mockupstream(s), artifactapi)"
"${COMPOSE[@]}" up -d --build postgres redis minio mockupstream mockupstreama mockupstreamb artifactapi
echo "==> building and starting stack (postgres, redis, minio, mockupstream, artifactapi)"
"${COMPOSE[@]}" up -d --build postgres redis minio mockupstream artifactapi
echo "==> waiting for artifactapi health at ${API_URL}"
for i in $(seq 1 60); do
@@ -34,17 +34,7 @@ for i in $(seq 1 60); do
sleep 1
done
# Resolve the compose network the artifactapi container is attached to, so the
# real-package-manager test can launch a stock distro container on the same
# network and reach artifactapi by service name.
API_CID="$("${COMPOSE[@]}" ps -q artifactapi)"
COMPOSE_NETWORK="$(docker inspect -f '{{range $k,$_ := .NetworkSettings.Networks}}{{$k}}{{end}}' "${API_CID}" 2>/dev/null || true)"
echo "==> running dockerised e2e suite (compose network: ${COMPOSE_NETWORK:-unknown})"
echo "==> running dockerised e2e suite"
ARTIFACTAPI_URL="${API_URL}" \
MOCK_UPSTREAM_INTERNAL="${MOCK_UPSTREAM_INTERNAL:-http://mockupstream}" \
MOCK_UPSTREAM_A_INTERNAL="${MOCK_UPSTREAM_A_INTERNAL:-http://mockupstreama}" \
MOCK_UPSTREAM_B_INTERNAL="${MOCK_UPSTREAM_B_INTERNAL:-http://mockupstreamb}" \
ARTIFACTAPI_INTERNAL="${ARTIFACTAPI_INTERNAL:-http://artifactapi:8000}" \
COMPOSE_NETWORK="${COMPOSE_NETWORK}" \
go test -tags=dockere2e -count=1 -timeout=10m -v ./e2e-docker/...
+5 -35
View File
@@ -176,44 +176,14 @@ helm install <release> ${name}/<chart>`,
];
case 'alpine':
return isLocal
? [
{
title: 'Add the apk repo (real apk repo, APKINDEX auto-generated)',
language: 'bash',
code: `echo '${url}/api/v1/local/${name}' | sudo tee -a /etc/apk/repositories
sudo apk update --allow-untrusted
sudo apk add --allow-untrusted <package>`,
note: `Served unsigned (parity with the rpm repo's gpgcheck=0) — use --allow-untrusted, or install a signing key. apk fetches <arch>/APKINDEX.tar.gz under this base.`,
},
{
title: 'Publish a .apk (index regenerates automatically)',
language: 'bash',
code: `curl -fsSL --upload-file ./mypkg-1.0-r0.apk \\
${url}/api/v2/remotes/${name}/files/x86_64/mypkg-1.0-r0.apk`,
note: 'Upload each package at <arch>/<name>-<version>.apk — apk reconstructs that exact path from the index (APKINDEX carries no filename), so a mismatched path will 404 on install.',
},
]
: [
{
title: 'Add the APK repository',
language: 'bash',
code: `echo '${proxy}/' | sudo tee -a /etc/apk/repositories
sudo apk update
sudo apk add <package>`,
note: 'If the index is unsigned over the proxy, add --allow-untrusted or install the signing key into /etc/apk/keys.',
},
];
case 'github_alpine':
return [
{
title: 'Add the apk repo (metadata-only, from GitHub releases)',
title: 'Add the APK repository',
language: 'bash',
code: `echo '${proxy}' | sudo tee -a /etc/apk/repositories
sudo apk update --allow-untrusted
sudo apk add --allow-untrusted <package>`,
note: "The per-arch APKINDEX is synthesized from the configured GitHub repo's release .apk assets; package downloads are redirected to the backing releases remote. Served unsigned, so --allow-untrusted.",
code: `echo '${proxy}/' | sudo tee -a /etc/apk/repositories
sudo apk update
sudo apk add <package>`,
note: 'If the index is unsigned over the proxy, add --allow-untrusted or install the signing key into /etc/apk/keys.',
},
];