Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| cd7c2c4383 | |||
| f1820fd104 | |||
| 822f356881 | |||
| 73c0bfc670 | |||
| a8aa0c231b |
@@ -9,6 +9,18 @@ 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
|
||||
@@ -16,3 +28,7 @@ services:
|
||||
depends_on:
|
||||
mockupstream:
|
||||
condition: service_started
|
||||
mockupstreama:
|
||||
condition: service_started
|
||||
mockupstreamb:
|
||||
condition: service_started
|
||||
|
||||
@@ -30,6 +30,18 @@ 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
|
||||
|
||||
|
||||
Binary file not shown.
BIN
Binary file not shown.
BIN
Binary file not shown.
BIN
Binary file not shown.
BIN
Binary file not shown.
BIN
Binary file not shown.
BIN
Binary file not shown.
@@ -0,0 +1,55 @@
|
||||
<?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>
|
||||
@@ -0,0 +1,105 @@
|
||||
//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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
# 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";
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
# Mock upstream B for the multi-base_url e2e (see a.conf).
|
||||
server {
|
||||
listen 80;
|
||||
location / {
|
||||
default_type text/plain;
|
||||
return 200 "UPSTREAM-B";
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,163 @@
|
||||
//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)
|
||||
}
|
||||
}
|
||||
@@ -57,7 +57,7 @@ func do(t *testing.T, h http.Handler, method, path, body string) int {
|
||||
}
|
||||
|
||||
func TestRemotesErrorPaths(t *testing.T) {
|
||||
h := NewRemotesHandler(closedDB(t), nil).Routes()
|
||||
h := NewRemotesHandler(closedDB(t), nil, nil).Routes()
|
||||
if c := do(t, h, "GET", "/", ""); c != 500 {
|
||||
t.Errorf("list with dead db = %d, want 500", c)
|
||||
}
|
||||
|
||||
@@ -1,8 +1,10 @@
|
||||
package v2
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
|
||||
"github.com/go-chi/chi/v5"
|
||||
@@ -17,15 +19,23 @@ type Primer interface {
|
||||
EnqueuePrime(remote models.Remote)
|
||||
}
|
||||
|
||||
// MetadataFlusher purges a remote's cached mutable metadata (repodata / Release
|
||||
// / APKINDEX freshness keys). *cache.Redis satisfies it.
|
||||
type MetadataFlusher interface {
|
||||
FlushRemote(ctx context.Context, remote string) error
|
||||
}
|
||||
|
||||
type RemotesHandler struct {
|
||||
db *database.DB
|
||||
cache MetadataFlusher
|
||||
primers map[models.PackageType]Primer
|
||||
}
|
||||
|
||||
// NewRemotesHandler wires the handler to the per-type metadata primers. primers
|
||||
// may be nil; a package type with no registered primer simply skips priming.
|
||||
func NewRemotesHandler(db *database.DB, primers map[models.PackageType]Primer) *RemotesHandler {
|
||||
return &RemotesHandler{db: db, primers: primers}
|
||||
// NewRemotesHandler wires the handler to the metadata cache and per-type
|
||||
// primers. cache may be nil (flush-on-backend-change is skipped); primers may
|
||||
// be nil (a package type with no registered primer simply skips priming).
|
||||
func NewRemotesHandler(db *database.DB, cache MetadataFlusher, primers map[models.PackageType]Primer) *RemotesHandler {
|
||||
return &RemotesHandler{db: db, cache: cache, primers: primers}
|
||||
}
|
||||
|
||||
func (h *RemotesHandler) Routes() chi.Router {
|
||||
@@ -78,6 +88,14 @@ 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
|
||||
@@ -102,14 +120,42 @@ 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)
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,96 @@
|
||||
package v2
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"git.unkin.net/unkin/artifactapi/internal/database"
|
||||
"git.unkin.net/unkin/artifactapi/pkg/models"
|
||||
)
|
||||
|
||||
// fakeFlusher records FlushRemote calls so a test can assert whether — and how
|
||||
// often — a remote's cached metadata was purged.
|
||||
type fakeFlusher struct {
|
||||
calls []string
|
||||
err error
|
||||
}
|
||||
|
||||
func (f *fakeFlusher) FlushRemote(_ context.Context, remote string) error {
|
||||
f.calls = append(f.calls, remote)
|
||||
return f.err
|
||||
}
|
||||
|
||||
func seedRemote(t *testing.T, db *database.DB, name, baseURL string) {
|
||||
t.Helper()
|
||||
err := db.CreateRemote(context.Background(), &models.Remote{
|
||||
Name: name,
|
||||
PackageType: models.PackageRPM,
|
||||
RepoType: models.RepoTypeRemote,
|
||||
BaseURL: baseURL,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("seed remote: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// A base_url change must flush the remote's cached metadata exactly once, while
|
||||
// an update that leaves base_url untouched must not flush at all.
|
||||
func TestUpdateFlushesCacheOnBaseURLChange(t *testing.T) {
|
||||
if testDSN == "" {
|
||||
t.Skip("Docker unavailable")
|
||||
}
|
||||
db, err := database.New(testDSN)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
const name = "rpm-flush-change"
|
||||
seedRemote(t, db, name, "https://old.example.com/repo")
|
||||
|
||||
ff := &fakeFlusher{}
|
||||
h := NewRemotesHandler(db, ff, nil).Routes()
|
||||
|
||||
if c := do(t, h, "PUT", "/"+name, `{"package_type":"rpm","repo_type":"remote","base_url":"https://new.example.com/repo"}`); c != 200 {
|
||||
t.Fatalf("update (backend change) = %d, want 200", c)
|
||||
}
|
||||
if len(ff.calls) != 1 || ff.calls[0] != name {
|
||||
t.Fatalf("flush calls = %v, want exactly one flush of %q", ff.calls, name)
|
||||
}
|
||||
|
||||
// Re-updating with the same (now current) base_url must not flush again.
|
||||
ff.calls = nil
|
||||
if c := do(t, h, "PUT", "/"+name, `{"package_type":"rpm","repo_type":"remote","base_url":"https://new.example.com/repo"}`); c != 200 {
|
||||
t.Fatalf("update (no backend change) = %d, want 200", c)
|
||||
}
|
||||
if len(ff.calls) != 0 {
|
||||
t.Fatalf("flush calls = %v, want no flush when base_url is unchanged", ff.calls)
|
||||
}
|
||||
}
|
||||
|
||||
// A flush error must be swallowed: the DB update already succeeded, so the
|
||||
// request still returns 200.
|
||||
func TestUpdateFlushFailureStillSucceeds(t *testing.T) {
|
||||
if testDSN == "" {
|
||||
t.Skip("Docker unavailable")
|
||||
}
|
||||
db, err := database.New(testDSN)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
const name = "rpm-flush-error"
|
||||
seedRemote(t, db, name, "https://old.example.com/repo")
|
||||
|
||||
ff := &fakeFlusher{err: errors.New("redis down")}
|
||||
h := NewRemotesHandler(db, ff, nil).Routes()
|
||||
|
||||
if c := do(t, h, "PUT", "/"+name, `{"package_type":"rpm","repo_type":"remote","base_url":"https://new.example.com/repo"}`); c != 200 {
|
||||
t.Fatalf("update with failing flush = %d, want 200", c)
|
||||
}
|
||||
if len(ff.calls) != 1 {
|
||||
t.Fatalf("flush calls = %v, want exactly one attempted flush", ff.calls)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
ORDER BY name, version, arch, file_path
|
||||
`, repoName)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
||||
@@ -99,6 +99,53 @@ 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")
|
||||
|
||||
@@ -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
|
||||
size, md5, sha256, created_at
|
||||
FROM deb_metadata
|
||||
WHERE repo_name = $1
|
||||
ORDER BY name, version, architecture
|
||||
ORDER BY name, version, architecture, file_path
|
||||
`, 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.Size, &m.MD5, &m.SHA256, &m.CreatedAt,
|
||||
); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -44,6 +44,8 @@ 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 '',
|
||||
@@ -124,6 +126,8 @@ 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;
|
||||
|
||||
@@ -6,7 +6,7 @@ import (
|
||||
"git.unkin.net/unkin/artifactapi/pkg/models"
|
||||
)
|
||||
|
||||
const remoteCols = `name, package_type, repo_type, base_url, description, username, password,
|
||||
const remoteCols = `name, package_type, repo_type, base_url, mirrorlist, mirror_strategy, description, username, password,
|
||||
immutable_ttl, mutable_ttl, check_mutable,
|
||||
patterns, blocklist, mutable_patterns, immutable_patterns,
|
||||
ban_tags_enabled, ban_tags,
|
||||
@@ -15,9 +15,18 @@ const remoteCols = `name, package_type, repo_type, base_url, description, userna
|
||||
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.Description, &r.Username, &r.Password,
|
||||
&r.Name, &r.PackageType, &r.RepoType, &r.BaseURL, &r.Mirrorlist, &r.MirrorStrategy, &r.Description, &r.Username, &r.Password,
|
||||
&r.ImmutableTTL, &r.MutableTTL, &r.CheckMutable,
|
||||
&r.Patterns, &r.Blocklist, &r.MutablePatterns, &r.ImmutablePatterns,
|
||||
&r.BanTagsEnabled, &r.BanTags,
|
||||
@@ -58,22 +67,24 @@ 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, description, username, password,
|
||||
name, package_type, repo_type, base_url, mirrorlist, 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
|
||||
) 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)
|
||||
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)
|
||||
`,
|
||||
r.Name, r.PackageType, r.RepoType, r.BaseURL, r.Description, r.Username, r.Password,
|
||||
r.Name, r.PackageType, r.RepoType, r.BaseURL, r.Mirrorlist, 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
|
||||
}
|
||||
@@ -81,13 +92,14 @@ 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, description=$5, username=$6, password=$7,
|
||||
package_type=$2, repo_type=$3, base_url=$4, mirrorlist=$25, 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
|
||||
`,
|
||||
@@ -98,6 +110,8 @@ 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
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package database
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
"git.unkin.net/unkin/artifactapi/internal/provider"
|
||||
)
|
||||
@@ -65,6 +66,7 @@ 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) {
|
||||
@@ -94,6 +96,7 @@ 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)
|
||||
@@ -112,10 +115,11 @@ 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
|
||||
requires, provides, conflicts, obsoletes, files, changelogs,
|
||||
created_at
|
||||
FROM rpm_metadata
|
||||
WHERE repo_name = $1
|
||||
ORDER BY name, epoch, version, release, arch
|
||||
ORDER BY name, epoch, version, release, arch, file_path
|
||||
`, repoName)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -131,6 +135,7 @@ 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
|
||||
}
|
||||
|
||||
@@ -15,6 +15,7 @@ import (
|
||||
"path"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"archive/tar"
|
||||
|
||||
@@ -376,7 +377,10 @@ func generateAPKIndex(metas []provider.AlpineMetadata) []byte {
|
||||
var tarBuf bytes.Buffer
|
||||
tw := tar.NewWriter(&tarBuf)
|
||||
body := idx.Bytes()
|
||||
tw.WriteHeader(&tar.Header{Name: "APKINDEX", Mode: 0o644, Size: int64(len(body)), Typeflag: tar.TypeReg})
|
||||
// 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.Write(body)
|
||||
tw.Close()
|
||||
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
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())
|
||||
}
|
||||
}
|
||||
@@ -395,7 +395,7 @@ func generateRelease(metas []provider.DebMetadata) []byte {
|
||||
arches := uniqueArches(metas)
|
||||
|
||||
var b bytes.Buffer
|
||||
fmt.Fprintf(&b, "Date: %s\n", time.Now().UTC().Format(time.RFC1123Z))
|
||||
fmt.Fprintf(&b, "Date: %s\n", releaseDate(metas).Format(time.RFC1123Z))
|
||||
fmt.Fprintf(&b, "Architectures: %s\n", strings.Join(arches, " "))
|
||||
b.WriteString("Acquire-By-Hash: no\n")
|
||||
|
||||
@@ -410,6 +410,21 @@ 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)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,172 @@
|
||||
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
|
||||
}
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"git.unkin.net/unkin/artifactapi/pkg/models"
|
||||
)
|
||||
@@ -113,6 +114,10 @@ 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
|
||||
@@ -185,6 +190,9 @@ 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 {
|
||||
|
||||
@@ -275,7 +275,7 @@ func (p *Provider) serveRepomd(w http.ResponseWriter, r *http.Request, reader pr
|
||||
filelistsHash := sha256Hex(filelists)
|
||||
otherHash := sha256Hex(other)
|
||||
|
||||
repomd := generateRepomd(primaryHash, len(primary), filelistsHash, len(filelists), otherHash, len(other))
|
||||
repomd := generateRepomd(repomdRevision(metas), primaryHash, len(primary), filelistsHash, len(filelists), otherHash, len(other))
|
||||
|
||||
w.Header().Set("Content-Type", "application/xml")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
@@ -315,8 +315,32 @@ func (p *Provider) serveOther(w http.ResponseWriter, r *http.Request, reader pro
|
||||
w.Write(generateOtherXMLGZ(metas))
|
||||
}
|
||||
|
||||
func generateRepomd(primaryHash string, primarySize int, filelistsHash string, filelistsSize int, otherHash string, otherSize int) []byte {
|
||||
ts := fmt.Sprintf("%d", time.Now().Unix())
|
||||
// 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 {
|
||||
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")
|
||||
@@ -359,7 +383,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", time.Now().Unix())
|
||||
fmt.Fprintf(&xmlBuf, " <time file=\"%d\" build=\"0\"/>\n", stableUnix(m.CreatedAt))
|
||||
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")
|
||||
@@ -484,6 +508,9 @@ 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()
|
||||
|
||||
@@ -0,0 +1,148 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -10,7 +10,10 @@ import (
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"git.unkin.net/unkin/artifactapi/internal/cache"
|
||||
@@ -35,6 +38,15 @@ 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 {
|
||||
@@ -222,7 +234,32 @@ 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)
|
||||
@@ -277,7 +314,33 @@ func (e *Engine) headUpstream(ctx context.Context, remote models.Remote, path st
|
||||
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)
|
||||
@@ -454,7 +517,33 @@ 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)
|
||||
@@ -649,3 +738,83 @@ 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 {
|
||||
sort.SliceStable(ordered, func(a, b int) bool {
|
||||
return e.inflightCounter(remote.Name, ordered[a]).Load() < e.inflightCounter(remote.Name, ordered[b]).Load()
|
||||
})
|
||||
}
|
||||
return ordered
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
@@ -0,0 +1,103 @@
|
||||
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])
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
@@ -0,0 +1,188 @@
|
||||
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")
|
||||
}
|
||||
}
|
||||
@@ -175,7 +175,7 @@ func (s *Server) routes() chi.Router {
|
||||
r.Mount("/api/v1", proxyHandler.Routes())
|
||||
r.Mount("/v2", proxyHandler.DockerV2Routes())
|
||||
|
||||
remotesHandler := v2.NewRemotesHandler(s.db, map[models.PackageType]v2.Primer{
|
||||
remotesHandler := v2.NewRemotesHandler(s.db, s.cache, map[models.PackageType]v2.Primer{
|
||||
models.PackageGitHubRPM: s.syncer,
|
||||
models.PackageGitHubDeb: s.debSyncer,
|
||||
models.PackageGitHubAlpine: s.alpineSyncer,
|
||||
|
||||
+80
-3
@@ -2,6 +2,7 @@ package models
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/url"
|
||||
"regexp"
|
||||
"time"
|
||||
)
|
||||
@@ -39,9 +40,17 @@ type Remote struct {
|
||||
PackageType PackageType `json:"package_type"`
|
||||
RepoType RepoType `json:"repo_type"`
|
||||
BaseURL string `json:"base_url"`
|
||||
Description string `json:"description,omitempty"`
|
||||
Username string `json:"-"`
|
||||
Password string `json:"-"`
|
||||
// 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:"-"`
|
||||
|
||||
ImmutableTTL int `json:"immutable_ttl"`
|
||||
MutableTTL int `json:"mutable_ttl"`
|
||||
@@ -72,6 +81,74 @@ 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.
|
||||
|
||||
+107
-1
@@ -1,6 +1,10 @@
|
||||
package models
|
||||
|
||||
import "testing"
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestRemote_ValidatePatterns(t *testing.T) {
|
||||
valid := &Remote{
|
||||
@@ -17,3 +21,105 @@ 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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
+13
-3
@@ -17,8 +17,8 @@ cleanup() {
|
||||
}
|
||||
trap cleanup EXIT
|
||||
|
||||
echo "==> building and starting stack (postgres, redis, minio, mockupstream, artifactapi)"
|
||||
"${COMPOSE[@]}" up -d --build postgres redis minio mockupstream artifactapi
|
||||
echo "==> building and starting stack (postgres, redis, minio, mockupstream(s), artifactapi)"
|
||||
"${COMPOSE[@]}" up -d --build postgres redis minio mockupstream mockupstreama mockupstreamb artifactapi
|
||||
|
||||
echo "==> waiting for artifactapi health at ${API_URL}"
|
||||
for i in $(seq 1 60); do
|
||||
@@ -34,7 +34,17 @@ for i in $(seq 1 60); do
|
||||
sleep 1
|
||||
done
|
||||
|
||||
echo "==> running dockerised e2e suite"
|
||||
# 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})"
|
||||
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/...
|
||||
|
||||
Reference in New Issue
Block a user