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
15 changed files with 2789 additions and 29 deletions
+55
View File
@@ -3,7 +3,10 @@
package e2edocker
import (
"archive/tar"
"bytes"
"compress/gzip"
"io"
"net/http"
"strings"
"testing"
@@ -136,3 +139,55 @@ func TestLocalDebRepo(t *testing.T) {
t.Fatalf("deb content mismatch")
}
}
// TestLocalAlpineIndex uploads an .apk to an alpine local repo and validates
// that a per-arch APKINDEX.tar.gz is generated automatically from the parsed
// .PKGINFO (the apk-local analog of rpm repodata / deb Packages generation).
func TestLocalAlpineIndex(t *testing.T) {
createRepo(t, `{"name":"local-alpine","package_type":"alpine","repo_type":"local"}`)
defer deleteRepo(t, "local-alpine")
apk := testsupport.MinimalApk("e2e-testpkg", "1.0-r0", "x86_64")
uploadFile(t, "local-alpine", "x86_64/e2e-testpkg-1.0-r0.apk", apk, "application/vnd.android.package-archive")
// The index is generated asynchronously after upload; poll for it.
resp, body := getEventually(t, api("/api/v1/local/local-alpine/x86_64/APKINDEX.tar.gz"), 15*time.Second)
if resp.StatusCode != http.StatusOK {
t.Fatalf("APKINDEX: status %d: %s", resp.StatusCode, body)
}
zr, err := gzip.NewReader(bytes.NewReader(body))
if err != nil {
t.Fatalf("APKINDEX not gzip: %v", err)
}
tarBytes, _ := io.ReadAll(zr)
tr := tar.NewReader(bytes.NewReader(tarBytes))
var index string
for {
hdr, err := tr.Next()
if err == io.EOF {
break
}
if err != nil {
t.Fatalf("APKINDEX not tar: %v", err)
}
if hdr.Name == "APKINDEX" {
b, _ := io.ReadAll(tr)
index = string(b)
}
}
for _, want := range []string{"P:e2e-testpkg", "V:1.0-r0", "A:x86_64", "C:Q1", "S:", "I:"} {
if !strings.Contains(index, want) {
t.Fatalf("APKINDEX missing %q:\n%s", want, index)
}
}
// The .apk downloads back byte-identical from its arch path.
resp, body = doRequest(t, http.MethodGet, api("/api/v1/local/local-alpine/x86_64/e2e-testpkg-1.0-r0.apk"), nil, "")
if resp.StatusCode != http.StatusOK {
t.Fatalf("download apk: status %d: %s", resp.StatusCode, body)
}
if !bytes.Equal(body, apk) {
t.Fatalf("apk content mismatch")
}
}
+71
View File
@@ -0,0 +1,71 @@
package database
import (
"context"
"errors"
"time"
"github.com/jackc/pgx/v5"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// ListGitHubAlpineRemotes returns every github_alpine remote so the syncer can
// sweep them on each poll tick.
func (db *DB) ListGitHubAlpineRemotes(ctx context.Context) ([]models.Remote, error) {
rows, err := db.Pool.Query(ctx, `SELECT `+remoteCols+` FROM remotes WHERE package_type = $1 ORDER BY name`, models.PackageGitHubAlpine)
if err != nil {
return nil, err
}
defer rows.Close()
var remotes []models.Remote
for rows.Next() {
var r models.Remote
if err := scanRemote(rows, &r); err != nil {
return nil, err
}
remotes = append(remotes, r)
}
return remotes, rows.Err()
}
// ClaimGitHubAlpineSyncLease atomically claims the per-remote sync lease. It
// succeeds only when the remote is due (never synced, or synced longer than
// freshness ago) and no live lease is held by another replica. A zero freshness
// (prime scans) ignores the recency gate. The returned etag is the stored
// releases-list ETag, shared across replicas.
func (db *DB) ClaimGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (bool, string, error) {
row := db.Pool.QueryRow(ctx, `
INSERT INTO github_alpine_sync_state AS s (remote_name, sync_lease_owner, sync_lease_expires)
VALUES ($1, $2, now() + make_interval(secs => $4))
ON CONFLICT (remote_name) DO UPDATE
SET sync_lease_owner = $2,
sync_lease_expires = now() + make_interval(secs => $4)
WHERE (s.last_synced_at IS NULL OR s.last_synced_at < now() - make_interval(secs => $3))
AND (s.sync_lease_expires IS NULL OR s.sync_lease_expires < now())
RETURNING s.etag
`, remoteName, owner, freshness.Seconds(), lease.Seconds())
var etag string
if err := row.Scan(&etag); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return false, "", nil
}
return false, "", err
}
return true, etag, nil
}
// ReleaseGitHubAlpineSyncLease records the completed scan and frees the lease.
// Only the owning replica may release; last_synced_at advances so the next poll
// waits a full freshness window, and etag is persisted for the next conditional
// request.
func (db *DB) ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error {
_, err := db.Pool.Exec(ctx, `
UPDATE github_alpine_sync_state
SET last_synced_at = $3, etag = $4, sync_lease_owner = '', sync_lease_expires = NULL
WHERE remote_name = $1 AND sync_lease_owner = $2
`, remoteName, owner, syncedAt, etag)
return err
}
+70
View File
@@ -0,0 +1,70 @@
package database
import (
"context"
"strings"
"git.unkin.net/unkin/artifactapi/internal/provider"
)
func (db *DB) InsertAlpineMetadata(ctx context.Context, meta *provider.AlpineMetadata) error {
_, err := db.Pool.Exec(ctx, `
INSERT INTO alpine_metadata (
repo_name, file_path, content_hash, checksum,
name, version, arch, download_size, installed_size,
description, url, license, origin, maintainer,
build_time, commit_hash, provider_priority,
depends, provides, install_if
) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20)
ON CONFLICT (repo_name, file_path) DO NOTHING
`,
meta.RepoName, meta.FilePath, meta.ContentHash, meta.Checksum,
meta.Name, meta.Version, meta.Arch, meta.DownloadSize, meta.InstalledSize,
meta.Description, meta.URL, meta.License, meta.Origin, meta.Maintainer,
meta.BuildTime, meta.Commit, meta.ProviderPriority,
strings.Join(meta.Depends, " "), strings.Join(meta.Provides, " "), strings.Join(meta.InstallIf, " "),
)
return err
}
func (db *DB) DeleteAlpineMetadata(ctx context.Context, repoName, filePath string) error {
_, err := db.Pool.Exec(ctx, `DELETE FROM alpine_metadata WHERE repo_name = $1 AND file_path = $2`, repoName, filePath)
return err
}
func (db *DB) ListAlpineMetadataEntries(ctx context.Context, repoName string) ([]provider.AlpineMetadata, error) {
rows, err := db.Pool.Query(ctx, `
SELECT repo_name, file_path, content_hash, checksum,
name, version, arch, download_size, installed_size,
description, url, license, origin, maintainer,
build_time, commit_hash, provider_priority,
depends, provides, install_if
FROM alpine_metadata
WHERE repo_name = $1
ORDER BY name, version, arch
`, repoName)
if err != nil {
return nil, err
}
defer rows.Close()
var result []provider.AlpineMetadata
for rows.Next() {
var m provider.AlpineMetadata
var depends, provides, installIf string
if err := rows.Scan(
&m.RepoName, &m.FilePath, &m.ContentHash, &m.Checksum,
&m.Name, &m.Version, &m.Arch, &m.DownloadSize, &m.InstalledSize,
&m.Description, &m.URL, &m.License, &m.Origin, &m.Maintainer,
&m.BuildTime, &m.Commit, &m.ProviderPriority,
&depends, &provides, &installIf,
); err != nil {
return nil, err
}
m.Depends = strings.Fields(depends)
m.Provides = strings.Fields(provides)
m.InstallIf = strings.Fields(installIf)
result = append(result, m)
}
return result, rows.Err()
}
+37
View File
@@ -182,6 +182,35 @@ func (db *DB) migrate() error {
CREATE INDEX IF NOT EXISTS idx_deb_metadata_repo ON deb_metadata(repo_name);
CREATE TABLE IF NOT EXISTS alpine_metadata (
id BIGSERIAL PRIMARY KEY,
repo_name TEXT NOT NULL,
file_path TEXT NOT NULL,
content_hash TEXT NOT NULL,
checksum TEXT NOT NULL,
name TEXT NOT NULL,
version TEXT NOT NULL,
arch TEXT NOT NULL,
download_size BIGINT DEFAULT 0,
installed_size BIGINT DEFAULT 0,
description TEXT DEFAULT '',
url TEXT DEFAULT '',
license TEXT DEFAULT '',
origin TEXT DEFAULT '',
maintainer TEXT DEFAULT '',
build_time BIGINT DEFAULT 0,
commit_hash TEXT DEFAULT '',
provider_priority TEXT DEFAULT '',
depends TEXT DEFAULT '',
provides TEXT DEFAULT '',
install_if TEXT DEFAULT '',
created_at TIMESTAMPTZ DEFAULT NOW(),
UNIQUE(repo_name, file_path)
);
CREATE INDEX IF NOT EXISTS idx_alpine_metadata_repo ON alpine_metadata(repo_name);
CREATE INDEX IF NOT EXISTS idx_alpine_metadata_repo_arch ON alpine_metadata(repo_name, arch);
CREATE TABLE IF NOT EXISTS github_rpm_sync_state (
remote_name TEXT PRIMARY KEY,
etag TEXT DEFAULT '',
@@ -198,6 +227,14 @@ func (db *DB) migrate() error {
sync_lease_expires TIMESTAMPTZ
);
CREATE TABLE IF NOT EXISTS github_alpine_sync_state (
remote_name TEXT PRIMARY KEY,
etag TEXT DEFAULT '',
last_synced_at TIMESTAMPTZ,
sync_lease_owner TEXT DEFAULT '',
sync_lease_expires TIMESTAMPTZ
);
CREATE TABLE IF NOT EXISTS signing_keys (
purpose TEXT PRIMARY KEY,
private_key_armor TEXT NOT NULL,
+357
View File
@@ -1,12 +1,26 @@
package alpine
import (
"bufio"
"bytes"
"compress/gzip"
"context"
"crypto/sha1"
"encoding/base64"
"errors"
"fmt"
"io"
"log/slog"
"net/http"
"path"
"strconv"
"strings"
"archive/tar"
"git.unkin.net/unkin/artifactapi/internal/auth"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/internal/storage"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
@@ -46,3 +60,346 @@ func (p *Provider) RewriteResponse(_ []byte, _ models.Remote, _ string) ([]byte,
func (p *Provider) AuthHeaders(_ context.Context, remote models.Remote) (http.Header, error) {
return auth.BasicHeaders(remote), nil
}
// --- LocalUploader: hosting real .apk packages -----------------------------
// ValidateUpload accepts any *.apk and preserves the client-supplied directory
// (the arch prefix) as the storage path, since arch cannot be parsed from the
// filename alone and the generic uploader hands us only the path. apk clients
// fetch packages at <arch>/<file>.apk, so publishers upload to that same path;
// AfterUpload records the true arch (from .PKGINFO) for index filtering.
func (p *Provider) ValidateUpload(filePath string) (storagePath, contentType string, err error) {
clean := strings.TrimPrefix(path.Clean("/"+filePath), "/")
filename := clean
if i := strings.LastIndex(clean, "/"); i >= 0 {
filename = clean[i+1:]
}
if !strings.HasSuffix(strings.ToLower(filename), ".apk") {
return "", "", fmt.Errorf("file must be a .apk package")
}
return clean, "application/vnd.android.package-archive", nil
}
func (p *Provider) UploadResponse(storagePath, contentHash string, sizeBytes int64) map[string]any {
filename := storagePath
if i := strings.LastIndex(storagePath, "/"); i >= 0 {
filename = storagePath[i+1:]
}
return map[string]any{
"filename": filename,
"content_hash": contentHash,
"size_bytes": sizeBytes,
}
}
func (p *Provider) AfterUpload(ctx context.Context, repoName, storagePath, contentHash string, blobs provider.BlobReader, db provider.MetadataStore) {
s3Key := storage.BlobKey(strings.TrimPrefix(contentHash, "sha256:"))
reader, blobSize, err := blobs.Download(ctx, s3Key)
if err != nil {
slog.Error("alpine metadata: download failed", "repo", repoName, "path", storagePath, "error", err)
return
}
defer reader.Close()
raw, err := io.ReadAll(reader)
if err != nil {
slog.Error("alpine metadata: read failed", "repo", repoName, "path", storagePath, "error", err)
return
}
meta, err := parseApk(raw)
if err != nil {
slog.Error("alpine metadata: parse failed", "repo", repoName, "path", storagePath, "error", err)
return
}
meta.RepoName = repoName
meta.FilePath = storagePath
meta.ContentHash = contentHash
meta.DownloadSize = blobSize
if meta.Name == "" || meta.Arch == "" {
slog.Error("alpine metadata: .PKGINFO missing pkgname/arch", "repo", repoName, "path", storagePath)
return
}
store, ok := db.(provider.AlpineMetadataStore)
if !ok {
slog.Error("alpine metadata: store does not support alpine metadata", "repo", repoName)
return
}
if err := store.InsertAlpineMetadata(ctx, meta); err != nil {
slog.Error("alpine metadata: insert failed", "repo", repoName, "path", storagePath, "error", err)
return
}
slog.Info("alpine metadata: parsed", "repo", repoName, "name", meta.Name, "version", meta.Version, "arch", meta.Arch)
}
func (p *Provider) AfterDelete(ctx context.Context, repoName, storagePath string, db provider.MetadataDeleter) error {
deleter, ok := db.(provider.AlpineMetadataDeleter)
if !ok {
return nil
}
if err := deleter.DeleteAlpineMetadata(ctx, repoName, storagePath); err != nil {
slog.Error("alpine metadata: delete failed", "repo", repoName, "path", storagePath, "error", err)
return err
}
slog.Info("alpine metadata: deleted", "repo", repoName, "path", storagePath)
return nil
}
// --- LocalIndexer: generating a per-arch APKINDEX.tar.gz -------------------
// normalizeIndexPath collapses apk's dot-segment prefix: an /etc/apk/repositories
// line of "<url>/api/v1/local/<name>" makes apk request "./<arch>/APKINDEX.tar.gz".
// Mirrors deb's flat-repo normalization.
func normalizeIndexPath(p string) string {
return strings.TrimPrefix(path.Clean("/"+p), "/")
}
func (p *Provider) ServeLocalIndex(w http.ResponseWriter, r *http.Request, files provider.FileStore, repoName, reqPath string) bool {
clean := normalizeIndexPath(reqPath)
if !strings.HasSuffix(clean, "APKINDEX.tar.gz") {
return false
}
arch := strings.TrimSuffix(clean, "APKINDEX.tar.gz")
arch = strings.Trim(arch, "/")
if arch == "" || strings.Contains(arch, "/") {
http.Error(w, "APKINDEX must be requested per-arch: <arch>/APKINDEX.tar.gz", http.StatusNotFound)
return true
}
reader, ok := files.(provider.AlpineMetadataReader)
if !ok {
http.Error(w, "alpine metadata not available", http.StatusInternalServerError)
return true
}
metas, err := reader.ListAlpineMetadataEntries(r.Context(), repoName)
if err != nil {
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
slog.Warn("alpine: metadata read canceled", "repo", repoName, "error", err)
http.Error(w, "metadata read canceled", http.StatusServiceUnavailable)
return true
}
http.Error(w, err.Error(), http.StatusInternalServerError)
return true
}
var filtered []provider.AlpineMetadata
for _, m := range metas {
if m.Arch == arch {
filtered = append(filtered, m)
}
}
w.Header().Set("Content-Type", "application/gzip")
w.WriteHeader(http.StatusOK)
w.Write(generateAPKIndex(filtered))
return true
}
func (p *Provider) GenerateLocalIndex(ctx context.Context, files provider.FileStore, repoName, path string) ([]byte, error) {
return nil, fmt.Errorf("alpine local index generation for virtual repos not supported")
}
// --- pure-Go .apk parsing --------------------------------------------------
// parseApk reads an .apk (up to three concatenated, independently gzipped tar
// streams: optional signature, control, data). It locates the control stream by
// its .PKGINFO member, computes the apk pull checksum C: = "Q1" +
// base64(sha1(<control gzip stream bytes>)), and reads the .PKGINFO fields.
func parseApk(raw []byte) (*provider.AlpineMetadata, error) {
members, err := gzipMembers(raw)
if err != nil {
return nil, err
}
for _, m := range members {
pkginfo, ok := pkginfoFromTar(m.tar)
if !ok {
continue
}
meta := parsePkginfo(pkginfo)
sum := sha1.Sum(m.raw)
meta.Checksum = "Q1" + base64.StdEncoding.EncodeToString(sum[:])
return meta, nil
}
return nil, errors.New("no .PKGINFO found in any .apk gzip stream")
}
type gzMember struct {
raw []byte // the raw bytes of this gzip stream (for the Q1 checksum)
tar []byte // the decompressed tar payload
}
// gzipMembers splits the concatenated gzip streams, returning each stream's raw
// bytes alongside its decompressed tar. It relies on bytes.Reader being an
// io.ByteReader (so compress/gzip does not over-read past a member's trailer)
// to recover exact stream boundaries via Multistream(false)+Reset.
func gzipMembers(data []byte) ([]gzMember, error) {
br := bytes.NewReader(data)
zr, err := gzip.NewReader(br)
if err != nil {
return nil, err
}
var members []gzMember
prev := 0
for {
zr.Multistream(false)
out, err := io.ReadAll(zr)
if err != nil {
return nil, err
}
end := len(data) - br.Len()
members = append(members, gzMember{raw: data[prev:end], tar: out})
prev = end
if err := zr.Reset(br); err != nil {
if err == io.EOF {
break
}
return nil, err
}
}
return members, nil
}
func pkginfoFromTar(tarBytes []byte) (string, bool) {
tr := tar.NewReader(bytes.NewReader(tarBytes))
for {
hdr, err := tr.Next()
if err != nil {
return "", false
}
if strings.TrimPrefix(hdr.Name, "./") == ".PKGINFO" {
b, err := io.ReadAll(tr)
if err != nil {
return "", false
}
return string(b), true
}
}
}
// parsePkginfo reads the "key = value" .PKGINFO text, collecting the repeated
// depend/provides/install_if keys into slices.
func parsePkginfo(text string) *provider.AlpineMetadata {
m := &provider.AlpineMetadata{}
sc := bufio.NewScanner(strings.NewReader(text))
sc.Buffer(make([]byte, 0, 64*1024), 1024*1024)
for sc.Scan() {
line := strings.TrimSpace(sc.Text())
if line == "" || strings.HasPrefix(line, "#") {
continue
}
idx := strings.Index(line, "=")
if idx < 0 {
continue
}
key := strings.TrimSpace(line[:idx])
val := strings.TrimSpace(line[idx+1:])
switch key {
case "pkgname":
m.Name = val
case "pkgver":
m.Version = val
case "arch":
m.Arch = val
case "pkgdesc":
m.Description = val
case "url":
m.URL = val
case "license":
m.License = val
case "origin":
m.Origin = val
case "maintainer":
m.Maintainer = val
case "builddate":
if n, err := strconv.ParseInt(val, 10, 64); err == nil {
m.BuildTime = n
}
case "commit":
m.Commit = val
case "size":
if n, err := strconv.ParseInt(val, 10, 64); err == nil {
m.InstalledSize = n
}
case "provider_priority":
m.ProviderPriority = val
case "depend":
if val != "" {
m.Depends = append(m.Depends, val)
}
case "provides":
if val != "" {
m.Provides = append(m.Provides, val)
}
case "install_if":
if val != "" {
m.InstallIf = append(m.InstallIf, val)
}
}
}
return m
}
// generateAPKIndex builds the APKINDEX.tar.gz = gzip(tar(APKINDEX)) for the
// given (already arch-filtered) rows. Records are blank-line separated; fields
// follow the canonical C/P/V/A/S/I/T/U/L/o/m/t/c/k/D/p/i order and empties are
// omitted. Unsigned (clients use --allow-untrusted), matching rpm gpgcheck=0.
func generateAPKIndex(metas []provider.AlpineMetadata) []byte {
var idx bytes.Buffer
for i, m := range metas {
if i > 0 {
idx.WriteString("\n")
}
writeField(&idx, "C", m.Checksum)
writeField(&idx, "P", m.Name)
writeField(&idx, "V", m.Version)
writeField(&idx, "A", m.Arch)
writeField(&idx, "S", intField(m.DownloadSize))
writeField(&idx, "I", intField(m.InstalledSize))
writeField(&idx, "T", m.Description)
writeField(&idx, "U", m.URL)
writeField(&idx, "L", m.License)
writeField(&idx, "o", m.Origin)
writeField(&idx, "m", m.Maintainer)
writeField(&idx, "t", intField(m.BuildTime))
writeField(&idx, "c", m.Commit)
writeField(&idx, "k", m.ProviderPriority)
writeField(&idx, "D", strings.Join(m.Depends, " "))
writeField(&idx, "p", strings.Join(m.Provides, " "))
writeField(&idx, "i", strings.Join(m.InstallIf, " "))
}
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})
tw.Write(body)
tw.Close()
var gzBuf bytes.Buffer
gz := gzip.NewWriter(&gzBuf)
gz.Write(tarBuf.Bytes())
gz.Close()
return gzBuf.Bytes()
}
func writeField(b *bytes.Buffer, key, val string) {
if val == "" {
return
}
b.WriteString(key)
b.WriteString(":")
b.WriteString(val)
b.WriteString("\n")
}
func intField(n int64) string {
if n == 0 {
return ""
}
return strconv.FormatInt(n, 10)
}
@@ -0,0 +1,324 @@
package alpine
import (
"archive/tar"
"bytes"
"compress/gzip"
"context"
"crypto/sha1"
"encoding/base64"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/internal/testsupport"
)
type fakeBlobReader struct{ data []byte }
func (f fakeBlobReader) Download(_ context.Context, _ string) (io.ReadCloser, int64, error) {
return io.NopCloser(bytes.NewReader(f.data)), int64(len(f.data)), nil
}
type errBlobReader struct{}
func (errBlobReader) Download(_ context.Context, _ string) (io.ReadCloser, int64, error) {
return nil, 0, io.ErrUnexpectedEOF
}
// fakeAlpineStore satisfies provider.MetadataStore (shared) and
// provider.AlpineMetadataStore, recording the row AfterUpload writes.
type fakeAlpineStore struct{ inserted *provider.AlpineMetadata }
func (f *fakeAlpineStore) InsertRPMMetadata(context.Context, *provider.RPMMetadata) error { return nil }
func (f *fakeAlpineStore) InsertDebMetadata(context.Context, *provider.DebMetadata) error { return nil }
func (f *fakeAlpineStore) InsertAlpineMetadata(_ context.Context, m *provider.AlpineMetadata) error {
f.inserted = m
return nil
}
// fakeAlpineDeleter satisfies provider.MetadataDeleter and AlpineMetadataDeleter.
type fakeAlpineDeleter struct{ deleted bool }
func (f *fakeAlpineDeleter) DeleteRPMMetadata(context.Context, string, string) error { return nil }
func (f *fakeAlpineDeleter) DeleteDebMetadata(context.Context, string, string) error { return nil }
func (f *fakeAlpineDeleter) DeleteAlpineMetadata(context.Context, string, string) error {
f.deleted = true
return nil
}
// fakeAlpineReader is a FileStore that also serves alpine metadata rows.
type fakeAlpineReader struct{ metas []provider.AlpineMetadata }
func (f fakeAlpineReader) ListAlpineMetadataEntries(context.Context, string) ([]provider.AlpineMetadata, error) {
return f.metas, nil
}
func (f fakeAlpineReader) ListFilesByPrefix(context.Context, string, string) ([]provider.FileEntry, error) {
return nil, nil
}
func (f fakeAlpineReader) ListPackages(context.Context, string) ([]string, error) { return nil, nil }
type errAlpineReader struct{}
func (errAlpineReader) ListAlpineMetadataEntries(context.Context, string) ([]provider.AlpineMetadata, error) {
return nil, io.ErrUnexpectedEOF
}
func (errAlpineReader) ListFilesByPrefix(context.Context, string, string) ([]provider.FileEntry, error) {
return nil, nil
}
func (errAlpineReader) ListPackages(context.Context, string) ([]string, error) { return nil, nil }
func TestAlpineValidateUpload(t *testing.T) {
p := &Provider{}
sp, ct, err := p.ValidateUpload("x86_64/foo-1.0-r0.apk")
if err != nil || sp != "x86_64/foo-1.0-r0.apk" || ct != "application/vnd.android.package-archive" {
t.Errorf("sp=%q ct=%q err=%v", sp, ct, err)
}
// Dot-segment prefix is normalized away.
if sp, _, err := p.ValidateUpload("./aarch64/bar-2.0-r1.apk"); err != nil || sp != "aarch64/bar-2.0-r1.apk" {
t.Errorf("dot-seg: sp=%q err=%v", sp, err)
}
if _, _, err := p.ValidateUpload("foo.rpm"); err == nil {
t.Error("expected error for non-apk")
}
resp := p.UploadResponse("x86_64/foo-1.0-r0.apk", "sha256:abc", 42)
if resp["filename"] != "foo-1.0-r0.apk" || resp["content_hash"] != "sha256:abc" || resp["size_bytes"] != int64(42) {
t.Errorf("upload response %v", resp)
}
}
func TestAlpineAfterUpload(t *testing.T) {
data := testsupport.MinimalApk("hello", "1.0-r0", "x86_64")
store := &fakeAlpineStore{}
(&Provider{}).AfterUpload(context.Background(), "myrepo", "x86_64/hello-1.0-r0.apk",
"sha256:deadbeef", fakeBlobReader{data: data}, store)
m := store.inserted
if m == nil {
t.Fatal("no metadata inserted")
}
if m.Name != "hello" || m.Version != "1.0-r0" || m.Arch != "x86_64" {
t.Errorf("unexpected metadata: %+v", m)
}
if m.DownloadSize != int64(len(data)) {
t.Errorf("DownloadSize = %d, want %d", m.DownloadSize, len(data))
}
if m.InstalledSize != 4 {
t.Errorf("InstalledSize = %d, want 4", m.InstalledSize)
}
if m.License != "MIT" || m.Origin != "hello" || !strings.HasPrefix(m.Maintainer, "e2e") {
t.Errorf("scalar fields not parsed: %+v", m)
}
if len(m.Depends) != 1 || m.Depends[0] != "so:libc.musl-x86_64.so.1" {
t.Errorf("Depends = %v", m.Depends)
}
if len(m.Provides) != 1 || m.Provides[0] != "cmd:hello=1.0-r0" {
t.Errorf("Provides = %v", m.Provides)
}
// The Q1 checksum is the sha1 of the CONTROL gzip stream (the member whose
// tar carries .PKGINFO), not of the whole file.
controlRaw := controlStreamBytes(t, data)
sum := sha1.Sum(controlRaw)
want := "Q1" + base64.StdEncoding.EncodeToString(sum[:])
if m.Checksum != want {
t.Errorf("Checksum = %q, want %q (sha1 of control stream)", m.Checksum, want)
}
// And explicitly NOT the sha1 of the whole apk.
whole := sha1.Sum(data)
if m.Checksum == "Q1"+base64.StdEncoding.EncodeToString(whole[:]) {
t.Error("Checksum was computed over the whole file, not the control stream")
}
}
func TestAlpineAfterUploadErrors(t *testing.T) {
store := &fakeAlpineStore{}
(&Provider{}).AfterUpload(context.Background(), "r", "x86_64/p.apk", "sha256:x", errBlobReader{}, store)
if store.inserted != nil {
t.Error("no metadata should be inserted on download error")
}
store2 := &fakeAlpineStore{}
(&Provider{}).AfterUpload(context.Background(), "r", "x86_64/p.apk", "sha256:x", fakeBlobReader{data: []byte("not an apk")}, store2)
if store2.inserted != nil {
t.Error("no metadata should be inserted on parse error")
}
}
func TestAlpineAfterDelete(t *testing.T) {
d := &fakeAlpineDeleter{}
if err := (&Provider{}).AfterDelete(context.Background(), "r", "x86_64/p.apk", d); err != nil {
t.Fatalf("AfterDelete: %v", err)
}
if !d.deleted {
t.Error("DeleteAlpineMetadata not called")
}
}
func TestAlpineServeLocalIndex(t *testing.T) {
p := &Provider{}
reader := fakeAlpineReader{metas: []provider.AlpineMetadata{
{Name: "aaa", Version: "1.0-r0", Arch: "x86_64", Checksum: "Q1aaa", DownloadSize: 100, InstalledSize: 10,
Description: "pkg aaa", URL: "https://a", License: "MIT", Depends: []string{"so:libc"}, Provides: []string{"cmd:aaa"}},
{Name: "bbb", Version: "2.0-r0", Arch: "aarch64", Checksum: "Q1bbb", DownloadSize: 200, InstalledSize: 20},
}}
// x86_64 index contains only aaa, with its fields, and not bbb.
w := serveIndex(t, p, reader, "x86_64/APKINDEX.tar.gz")
if w.Code != 200 {
t.Fatalf("code %d", w.Code)
}
idx := untarIndex(t, w.Body.Bytes())
for _, want := range []string{"C:Q1aaa", "P:aaa", "V:1.0-r0", "A:x86_64", "S:100", "I:10", "T:pkg aaa", "U:https://a", "L:MIT", "D:so:libc", "p:cmd:aaa"} {
if !strings.Contains(idx, want) {
t.Errorf("x86_64 APKINDEX missing %q:\n%s", want, idx)
}
}
if strings.Contains(idx, "P:bbb") {
t.Errorf("x86_64 APKINDEX leaked aarch64 package:\n%s", idx)
}
// aarch64 index contains only bbb.
w = serveIndex(t, p, reader, "aarch64/APKINDEX.tar.gz")
idx = untarIndex(t, w.Body.Bytes())
if !strings.Contains(idx, "P:bbb") || strings.Contains(idx, "P:aaa") {
t.Errorf("aarch64 filtering wrong:\n%s", idx)
}
// Non-index and .apk paths are not owned by the indexer.
for _, path := range []string{"x86_64/foo-1.0-r0.apk", "x86_64/", "README"} {
w := httptest.NewRecorder()
r := httptest.NewRequest(http.MethodGet, "/"+path, nil)
if p.ServeLocalIndex(w, r, reader, "repo", path) {
t.Errorf("ServeLocalIndex should return false for %q", path)
}
}
}
// Empty fields are omitted from the record (bbb has no description/url).
func TestAlpineIndexOmitsEmptyFields(t *testing.T) {
p := &Provider{}
reader := fakeAlpineReader{metas: []provider.AlpineMetadata{
{Name: "bbb", Version: "2.0-r0", Arch: "x86_64", Checksum: "Q1bbb", DownloadSize: 200, InstalledSize: 20},
}}
idx := untarIndex(t, serveIndex(t, p, reader, "x86_64/APKINDEX.tar.gz").Body.Bytes())
for _, absent := range []string{"T:", "U:", "L:", "D:", "p:", "i:", "o:", "m:", "c:", "k:"} {
if strings.Contains(idx, absent) {
t.Errorf("empty field %q should be omitted:\n%s", absent, idx)
}
}
}
// apk requests "./<arch>/APKINDEX.tar.gz" for a bare repo base URL; the
// dot-segment must be collapsed and yield the same bytes as the plain path.
func TestAlpineServeLocalIndexDotSegment(t *testing.T) {
p := &Provider{}
reader := fakeAlpineReader{metas: []provider.AlpineMetadata{
{Name: "aaa", Version: "1.0-r0", Arch: "x86_64", Checksum: "Q1aaa", DownloadSize: 100, InstalledSize: 10},
}}
plain := untarIndex(t, serveIndex(t, p, reader, "x86_64/APKINDEX.tar.gz").Body.Bytes())
dotted := untarIndex(t, serveIndex(t, p, reader, "./x86_64/APKINDEX.tar.gz").Body.Bytes())
if plain != dotted {
t.Errorf("dot-segment path differs:\nplain=%q\ndotted=%q", plain, dotted)
}
}
func TestAlpineServeLocalIndexArchRequired(t *testing.T) {
p := &Provider{}
reader := fakeAlpineReader{}
w := httptest.NewRecorder()
r := httptest.NewRequest(http.MethodGet, "/APKINDEX.tar.gz", nil)
if !p.ServeLocalIndex(w, r, reader, "repo", "APKINDEX.tar.gz") {
t.Fatal("bare APKINDEX should be owned (and rejected) by the indexer")
}
if w.Code != http.StatusNotFound {
t.Errorf("bare APKINDEX code = %d, want 404", w.Code)
}
}
func TestAlpineServeMetadataError(t *testing.T) {
p := &Provider{}
w := httptest.NewRecorder()
r := httptest.NewRequest(http.MethodGet, "/x86_64/APKINDEX.tar.gz", nil)
p.ServeLocalIndex(w, r, errAlpineReader{}, "repo", "x86_64/APKINDEX.tar.gz")
if w.Code != 500 {
t.Errorf("failing reader code = %d, want 500", w.Code)
}
}
func TestAlpineGenerateLocalIndexUnsupported(t *testing.T) {
if _, err := (&Provider{}).GenerateLocalIndex(context.Background(), fakeAlpineReader{}, "r", "x86_64/APKINDEX.tar.gz"); err == nil {
t.Error("expected unsupported error")
}
}
func serveIndex(t *testing.T, p *Provider, files provider.FileStore, path string) *httptest.ResponseRecorder {
t.Helper()
w := httptest.NewRecorder()
r := httptest.NewRequest(http.MethodGet, "/"+path, nil)
if !p.ServeLocalIndex(w, r, files, "repo", path) {
t.Fatalf("ServeLocalIndex returned false for %q", path)
}
return w
}
// untarIndex un-gzips and un-tars an APKINDEX.tar.gz and returns the APKINDEX text.
func untarIndex(t *testing.T, gzTar []byte) string {
t.Helper()
zr, err := gzip.NewReader(bytes.NewReader(gzTar))
if err != nil {
t.Fatalf("APKINDEX not gzip: %v", err)
}
tarBytes, _ := io.ReadAll(zr)
tr := tar.NewReader(bytes.NewReader(tarBytes))
for {
hdr, err := tr.Next()
if err == io.EOF {
break
}
if err != nil {
t.Fatalf("APKINDEX not tar: %v", err)
}
if hdr.Name == "APKINDEX" {
b, _ := io.ReadAll(tr)
return string(b)
}
}
t.Fatal("no APKINDEX member in tarball")
return ""
}
// controlStreamBytes returns the raw bytes of the gzip stream whose tar carries
// .PKGINFO, so the test can independently compute the expected Q1 checksum.
func controlStreamBytes(t *testing.T, apk []byte) []byte {
t.Helper()
br := bytes.NewReader(apk)
zr, err := gzip.NewReader(br)
if err != nil {
t.Fatalf("gzip: %v", err)
}
prev := 0
for {
zr.Multistream(false)
out, _ := io.ReadAll(zr)
end := len(apk) - br.Len()
tr := tar.NewReader(bytes.NewReader(out))
for {
h, err := tr.Next()
if err != nil {
break
}
if strings.TrimPrefix(h.Name, "./") == ".PKGINFO" {
return apk[prev:end]
}
}
prev = end
if err := zr.Reset(br); err != nil {
break
}
}
t.Fatal("no control stream found")
return nil
}
+714
View File
@@ -0,0 +1,714 @@
package alpine
import (
"bytes"
"compress/gzip"
"context"
"crypto/sha1"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"net/http"
"net/url"
"regexp"
"strings"
"sync"
"time"
"golang.org/x/time/rate"
"git.unkin.net/unkin/artifactapi/internal/githubauth"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// gitHubProvider is the process-wide singleton for github_alpine. The background
// Syncer binds its shared rate limiter and work queue onto this instance so the
// request path and the syncer drive the same derive machinery.
var gitHubProvider = newGitHubProvider()
func init() {
provider.Register(gitHubProvider)
}
// Tuning knobs for the no-precache control fetch. An .apk is up to three
// concatenated gzip streams (optional signature, control, data); the control
// stream carrying .PKGINFO sits near the front, so a small prefix reliably
// covers it.
const (
defaultHeaderRangeInitial = 32 << 10 // 32 KiB — covers the control stream of almost every .apk
defaultHeaderRangeMax = 16 << 20 // 16 MiB — give up past this and skip the asset
defaultReleasePageCap = 10 // 100 releases/page * 10 pages
defaultScanTimeout = 10 * time.Minute
defaultServeTimeout = 30 * time.Second
defaultColdWait = 8 * time.Second
)
// GitHubProvider is a metadata-only remote: it scans a GitHub repo's releases
// for .apk assets, derives per-asset .PKGINFO metadata via a ranged prefix fetch
// (never downloading whole packages), synthesizes a per-arch APKINDEX from that
// cached metadata, and redirects package downloads to a backend "releases_remote"
// (the generic github.com remote) that serves the actual bytes.
type GitHubProvider struct {
client *http.Client
headerInitial int64
headerMax int64
pageCap int
scanTimeout time.Duration
serveTimeout time.Duration
coldWait time.Duration
limiter *rate.Limiter
syncer *Syncer
serverCred githubauth.Credential
mu sync.Mutex
scanning map[string]bool
lastScan map[string]time.Time
}
func newGitHubProvider() *GitHubProvider {
return &GitHubProvider{
client: &http.Client{},
headerInitial: defaultHeaderRangeInitial,
headerMax: defaultHeaderRangeMax,
pageCap: defaultReleasePageCap,
scanTimeout: defaultScanTimeout,
serveTimeout: defaultServeTimeout,
coldWait: defaultColdWait,
scanning: map[string]bool{},
lastScan: map[string]time.Time{},
}
}
func (p *GitHubProvider) limiterWait(ctx context.Context) error {
if p.limiter == nil {
return nil
}
return p.limiter.Wait(ctx)
}
func (p *GitHubProvider) Type() models.PackageType { return models.PackageGitHubAlpine }
func (p *GitHubProvider) Classify(path string) provider.Mutability {
if strings.HasSuffix(path, "APKINDEX.tar.gz") {
return provider.Mutable
}
return provider.Immutable
}
func (p *GitHubProvider) ContentType(path string) string {
switch {
case strings.HasSuffix(path, ".apk"):
return "application/vnd.android.package-archive"
case strings.HasSuffix(path, ".tar.gz"):
return "application/gzip"
}
return "application/octet-stream"
}
func (p *GitHubProvider) UpstreamURL(remote models.Remote, path string) string {
return strings.TrimRight(remote.BaseURL, "/") + "/" + strings.TrimLeft(path, "/")
}
func (p *GitHubProvider) RewriteResponse(_ []byte, _ models.Remote, _ string) ([]byte, error) {
return nil, nil
}
func (p *GitHubProvider) AuthHeaders(ctx context.Context, remote models.Remote) (http.Header, error) {
return p.githubHeaders(ctx, remote, false)
}
// ServeRemote answers a request against a github_alpine remote. It refreshes the
// derived metadata (bounded by mutable_ttl), serves a synthesized per-arch
// APKINDEX.tar.gz, and 302-redirects .apk downloads to the backend
// releases_remote. Returns false only for paths it does not own.
func (p *GitHubProvider) ServeRemote(w http.ResponseWriter, r *http.Request, remote models.Remote, reqPath, proxyBaseURL string, store provider.RemoteMetadataStore) bool {
p.onRequest(remote, store)
// apk requests the index at "./<arch>/APKINDEX.tar.gz"; collapse the
// dot-segment before matching, mirroring the local indexer.
path := normalizeIndexPath(reqPath)
if strings.HasSuffix(path, "APKINDEX.tar.gz") {
p.serveIndex(w, r, remote, path, store)
return true
}
if strings.HasSuffix(path, ".apk") {
if remote.ReleasesRemote == "" {
http.Error(w, "github_alpine remote has no releases_remote configured for downloads", http.StatusInternalServerError)
return true
}
p.serveApkRedirect(w, r, remote, path, proxyBaseURL, store)
return true
}
return false
}
// serveApkRedirect resolves an apk-reconstructed download path — apk builds
// "<arch>/<name>-<version>.apk" itself because APKINDEX carries no filename — to
// the real github-relative asset path stored on the metadata row, then redirects
// to the backend releases_remote. Passing the inbound path through verbatim would
// point at a nonexistent, allowlist-denied github.com path.
func (p *GitHubProvider) serveApkRedirect(w http.ResponseWriter, r *http.Request, remote models.Remote, path, proxyBaseURL string, store provider.RemoteMetadataStore) {
arch := strings.TrimSuffix(path[:strings.LastIndex(path, "/")+1], "/")
basename := path[strings.LastIndex(path, "/")+1:]
if arch == "" || strings.Contains(arch, "/") {
http.Error(w, "apk download must be requested per-arch: <arch>/<name>-<version>.apk", http.StatusNotFound)
return
}
reader, ok := store.(provider.AlpineMetadataReader)
if !ok {
http.Error(w, "alpine metadata not available", http.StatusInternalServerError)
return
}
sctx, cancel := context.WithTimeout(context.WithoutCancel(r.Context()), p.serveTimeout)
defer cancel()
rows, err := reader.ListAlpineMetadataEntries(sctx, remote.Name)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
for _, row := range rows {
if row.Arch == arch && row.Name+"-"+row.Version+".apk" == basename {
loc := strings.TrimRight(proxyBaseURL, "/") + "/api/v1/remote/" + remote.ReleasesRemote + "/" + strings.TrimLeft(row.FilePath, "/")
http.Redirect(w, r, loc, http.StatusFound)
return
}
}
http.Error(w, "package not found", http.StatusNotFound)
}
func (p *GitHubProvider) serveIndex(w http.ResponseWriter, r *http.Request, remote models.Remote, path string, store provider.RemoteMetadataStore) {
arch := strings.TrimSuffix(path, "APKINDEX.tar.gz")
arch = strings.Trim(arch, "/")
if arch == "" || strings.Contains(arch, "/") {
http.Error(w, "APKINDEX must be requested per-arch: <arch>/APKINDEX.tar.gz", http.StatusNotFound)
return
}
// Serve on a context detached from the inbound request so a client disconnect
// never cancels the metadata DB read and surfaces as a 500.
sctx, cancel := context.WithTimeout(context.WithoutCancel(r.Context()), p.serveTimeout)
defer cancel()
if p.syncer != nil && !p.ensurePrimed(sctx, remote, store) {
w.Header().Set("Retry-After", "5")
http.Error(w, "metadata is being prepared, retry shortly", http.StatusServiceUnavailable)
return
}
reader, ok := store.(provider.AlpineMetadataReader)
if !ok {
http.Error(w, "alpine metadata not available", http.StatusInternalServerError)
return
}
metas, err := reader.ListAlpineMetadataEntries(sctx, remote.Name)
if err != nil {
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
http.Error(w, "metadata read canceled", http.StatusServiceUnavailable)
return
}
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
var filtered []provider.AlpineMetadata
for _, m := range metas {
if m.Arch == arch {
filtered = append(filtered, m)
}
}
w.Header().Set("Content-Type", "application/gzip")
w.WriteHeader(http.StatusOK)
w.Write(generateAPKIndex(filtered))
}
// onRequest keeps a remote's derived metadata fresh off the request path.
func (p *GitHubProvider) onRequest(remote models.Remote, store provider.RemoteMetadataStore) {
if p.syncer != nil {
p.syncer.enqueue(remote, false)
return
}
p.refresh(remote, store)
}
// ensurePrimed returns true once the remote has at least one cached row. On an
// empty cache it enqueues a prime and polls briefly for it to land.
func (p *GitHubProvider) ensurePrimed(ctx context.Context, remote models.Remote, store provider.RemoteMetadataStore) bool {
if !p.cacheEmpty(ctx, store, remote.Name) {
return true
}
if p.syncer != nil {
p.syncer.enqueue(remote, true)
}
deadline := time.Now().Add(p.coldWait)
for time.Now().Before(deadline) {
select {
case <-ctx.Done():
return false
case <-time.After(400 * time.Millisecond):
}
if !p.cacheEmpty(ctx, store, remote.Name) {
return true
}
}
return false
}
func (p *GitHubProvider) cacheEmpty(ctx context.Context, store provider.RemoteMetadataStore, name string) bool {
reader, ok := store.(provider.AlpineMetadataReader)
if !ok {
return false
}
rows, err := reader.ListAlpineMetadataEntries(ctx, name)
if err != nil {
return false
}
return len(rows) == 0
}
// refresh brings the derived metadata up to date without coupling the scan to
// the inbound request (legacy inline path used without a syncer / in unit tests).
func (p *GitHubProvider) refresh(remote models.Remote, store provider.RemoteMetadataStore) {
ttl := time.Duration(remote.MutableTTL) * time.Second
if ttl <= 0 {
ttl = 5 * time.Minute
}
p.mu.Lock()
last, ok := p.lastScan[remote.Name]
fresh := ok && time.Since(last) < ttl
if fresh || p.scanning[remote.Name] {
p.mu.Unlock()
return
}
p.scanning[remote.Name] = true
p.mu.Unlock()
if p.cacheEmpty(context.Background(), store, remote.Name) {
p.runScan(remote, store)
return
}
go p.runScan(remote, store)
}
func (p *GitHubProvider) runScan(remote models.Remote, store provider.RemoteMetadataStore) {
defer func() {
p.mu.Lock()
delete(p.scanning, remote.Name)
p.mu.Unlock()
}()
ctx, cancel := context.WithTimeout(context.Background(), p.scanTimeout)
defer cancel()
if err := p.scan(ctx, remote, store); err != nil {
slog.Error("github_alpine: release scan failed", "remote", remote.Name, "error", err)
return
}
p.mu.Lock()
p.lastScan[remote.Name] = time.Now()
p.mu.Unlock()
}
// scan runs a full unconditional derive. Retained for the legacy inline refresh
// path and existing tests; the syncer uses scanWithState.
func (p *GitHubProvider) scan(ctx context.Context, remote models.Remote, store provider.RemoteMetadataStore) error {
_, _, err := p.scanWithState(ctx, remote, store, "")
return err
}
// scanWithState derives metadata incrementally. It sends the prior releases-list
// ETag as a conditional request: a 304 means nothing changed. On a 200 it diffs
// the release assets against the cache, derives only new/changed assets, prunes
// assets that disappeared, and returns the new ETag.
func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote, store provider.RemoteMetadataStore, etag string) (newEtag string, changed bool, err error) {
inserter, ok := store.(provider.AlpineMetadataStore)
if !ok {
return etag, false, errors.New("store does not support alpine metadata writes")
}
deleter, ok := store.(provider.AlpineMetadataDeleter)
if !ok {
return etag, false, errors.New("store does not support alpine metadata deletes")
}
reader, ok := store.(provider.AlpineMetadataReader)
if !ok {
return etag, false, errors.New("store does not support alpine metadata reads")
}
releases, newEtag, notModified, err := p.fetchReleases(ctx, remote, etag)
if err != nil {
return etag, false, err
}
if notModified {
return etag, false, nil
}
existing, err := reader.ListAlpineMetadataEntries(ctx, remote.Name)
if err != nil {
return newEtag, false, err
}
existingByPath := make(map[string]provider.AlpineMetadata, len(existing))
for _, m := range existing {
existingByPath[m.FilePath] = m
}
allow, err := compilePatterns(remote.Patterns)
if err != nil {
return newEtag, false, err
}
seen := map[string]bool{}
for _, rel := range releases {
if rel.Draft {
continue
}
for _, asset := range rel.Assets {
if !strings.HasSuffix(strings.ToLower(asset.Name), ".apk") {
continue
}
if !matchesAny(allow, asset.Name) {
continue
}
fp := assetPath(asset)
if fp == "" {
continue
}
seen[fp] = true
if cur, ok := existingByPath[fp]; ok {
if asset.Digest == "" || cur.ContentHash == asset.Digest {
continue
}
_ = deleter.DeleteAlpineMetadata(ctx, remote.Name, fp)
}
meta, err := p.deriveAsset(ctx, remote, asset, fp)
if err != nil {
slog.Warn("github_alpine: derive asset failed", "remote", remote.Name, "asset", asset.Name, "error", err)
continue
}
if err := inserter.InsertAlpineMetadata(ctx, meta); err != nil {
slog.Error("github_alpine: insert metadata failed", "remote", remote.Name, "asset", asset.Name, "error", err)
continue
}
slog.Info("github_alpine: derived asset", "remote", remote.Name, "name", meta.Name, "version", meta.Version, "arch", meta.Arch)
}
}
for fp := range existingByPath {
if !seen[fp] {
_ = deleter.DeleteAlpineMetadata(ctx, remote.Name, fp)
}
}
return newEtag, true, nil
}
type ghRelease struct {
TagName string `json:"tag_name"`
Draft bool `json:"draft"`
Assets []ghAsset `json:"assets"`
}
type ghAsset struct {
Name string `json:"name"`
Size int64 `json:"size"`
BrowserDownloadURL string `json:"browser_download_url"`
Digest string `json:"digest"`
}
// fetchReleases lists a repo's releases, sending the prior ETag as If-None-Match
// on page 1 so an unchanged repo short-circuits to notModified. Every call waits
// on the shared limiter first.
func (p *GitHubProvider) fetchReleases(ctx context.Context, remote models.Remote, etag string) (all []ghRelease, newEtag string, notModified bool, err error) {
base := strings.TrimRight(remote.BaseURL, "/") + "/releases"
for page := 1; page <= p.pageCap; page++ {
u := fmt.Sprintf("%s?per_page=100&page=%d", base, page)
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
if err != nil {
return nil, "", false, err
}
hdr, err := p.githubHeaders(ctx, remote, true)
if err != nil {
return nil, "", false, err
}
copyHeaders(req, hdr)
if page == 1 && etag != "" {
req.Header.Set("If-None-Match", etag)
}
if err := p.limiterWait(ctx); err != nil {
return nil, "", false, err
}
resp, err := p.client.Do(req)
if err != nil {
return nil, "", false, err
}
if page == 1 && resp.StatusCode == http.StatusNotModified {
io.Copy(io.Discard, resp.Body)
resp.Body.Close()
return nil, etag, true, nil
}
body, err := io.ReadAll(resp.Body)
respEtag := resp.Header.Get("ETag")
resp.Body.Close()
if err != nil {
return nil, "", false, err
}
if resp.StatusCode != http.StatusOK {
return nil, "", false, fmt.Errorf("github releases API %s: status %d", u, resp.StatusCode)
}
if page == 1 {
newEtag = respEtag
}
var releases []ghRelease
if err := json.Unmarshal(body, &releases); err != nil {
return nil, "", false, fmt.Errorf("decode releases: %w", err)
}
if len(releases) == 0 {
break
}
all = append(all, releases...)
if len(releases) < 100 {
break
}
}
return all, newEtag, false, nil
}
func (p *GitHubProvider) deriveAsset(ctx context.Context, remote models.Remote, asset ghAsset, fp string) (*provider.AlpineMetadata, error) {
meta, err := p.fetchPkginfo(ctx, remote, asset.BrowserDownloadURL)
if err != nil {
return nil, err
}
if meta.Name == "" || meta.Arch == "" {
return nil, errors.New(".PKGINFO missing pkgname/arch")
}
meta.RepoName = remote.Name
meta.FilePath = fp
// S: the on-disk .apk size comes straight from the releases API, so we never
// download the body just to size it.
meta.DownloadSize = asset.Size
// ContentHash records the GitHub asset digest (when present) purely so the
// next scan can detect a changed asset; unlike deb it is not the index
// checksum (that is the Q1 control-stream sum already set in fetchPkginfo).
if asset.Digest != "" {
meta.ContentHash = asset.Digest
}
return meta, nil
}
// fetchPkginfo pulls only the front of the .apk with a ranged GET and derives the
// .PKGINFO fields plus the apk pull checksum (C: = Q1 + base64(sha1(control gzip
// stream))). The control stream sits near the front, so a small prefix suffices;
// a prefix that truncates it doubles the range and retries.
func (p *GitHubProvider) fetchPkginfo(ctx context.Context, remote models.Remote, downloadURL string) (*provider.AlpineMetadata, error) {
n := p.headerInitial
for {
body, full, err := p.rangeGet(ctx, remote, downloadURL, n)
if err != nil {
return nil, err
}
meta, complete, perr := pkginfoFromPrefix(body)
if perr != nil {
return nil, fmt.Errorf("parse apk .PKGINFO: %w", perr)
}
if complete {
return meta, nil
}
if full || n >= p.headerMax {
return nil, fmt.Errorf(".PKGINFO not found within %d bytes of %s", n, downloadURL)
}
n *= 2
if n > p.headerMax {
n = p.headerMax
}
}
}
// pkginfoFromPrefix parses the concatenated gzip streams present in a front
// prefix of an .apk. It walks each fully-covered gzip member until it finds the
// control stream (the one whose tar carries .PKGINFO), computes the Q1 pull
// checksum from that stream's raw bytes, and reads the .PKGINFO fields. A prefix
// too short to fully cover the control stream returns complete=false so the
// caller can widen the range.
func pkginfoFromPrefix(prefix []byte) (meta *provider.AlpineMetadata, complete bool, err error) {
br := bytes.NewReader(prefix)
zr, zerr := gzip.NewReader(br)
if zerr != nil {
if zerr == io.EOF || zerr == io.ErrUnexpectedEOF {
return nil, false, nil
}
return nil, false, zerr
}
prev := 0
for {
zr.Multistream(false)
out, rerr := io.ReadAll(zr)
if rerr != nil {
// A member truncated by the range boundary is not an error — widen.
if rerr == io.ErrUnexpectedEOF || rerr == io.EOF {
return nil, false, nil
}
return nil, false, rerr
}
end := len(prefix) - br.Len()
raw := prefix[prev:end]
if pkginfo, ok := pkginfoFromTar(out); ok {
m := parsePkginfo(pkginfo)
sum := sha1.Sum(raw)
m.Checksum = "Q1" + base64.StdEncoding.EncodeToString(sum[:])
return m, true, nil
}
prev = end
if rsterr := zr.Reset(br); rsterr != nil {
if rsterr == io.EOF {
// No more complete members in the prefix; the control stream is
// either not covered yet or genuinely absent — let the caller
// decide by widening (or hitting the full-object guard).
return nil, false, nil
}
if rsterr == io.ErrUnexpectedEOF {
return nil, false, nil
}
return nil, false, rsterr
}
}
}
// rangeGet returns the first n bytes of downloadURL. full is true when the
// response body was shorter than n (i.e. we already have the whole object).
func (p *GitHubProvider) rangeGet(ctx context.Context, remote models.Remote, downloadURL string, n int64) ([]byte, bool, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, downloadURL, nil)
if err != nil {
return nil, false, err
}
hdr, err := p.githubHeaders(ctx, remote, false)
if err != nil {
return nil, false, err
}
copyHeaders(req, hdr)
req.Header.Set("Range", fmt.Sprintf("bytes=0-%d", n-1))
if err := p.limiterWait(ctx); err != nil {
return nil, false, err
}
resp, err := p.client.Do(req)
if err != nil {
return nil, false, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusPartialContent {
return nil, false, fmt.Errorf("range GET %s: status %d", downloadURL, resp.StatusCode)
}
body, err := io.ReadAll(io.LimitReader(resp.Body, n))
if err != nil {
return nil, false, err
}
full := int64(len(body)) < n
return body, full, nil
}
// assetPath is the package's location relative to github.com — the path the
// backend releases_remote (base https://github.com) proxies. It doubles as the
// alpine_metadata key and the redirect target, so an .apk download resolves back
// to this remote and redirects to the backend.
func assetPath(asset ghAsset) string {
u, err := url.Parse(asset.BrowserDownloadURL)
if err != nil {
return ""
}
return strings.TrimPrefix(u.Path, "/")
}
// githubHeaders builds the outbound headers for a GitHub request, attaching a
// bearer credential when one is available. A per-remote credential wins; absent
// that, the process-wide server credential is used; absent both, the request is
// unauthenticated.
func (p *GitHubProvider) githubHeaders(ctx context.Context, remote models.Remote, api bool) (http.Header, error) {
h := http.Header{}
if api {
h.Set("Accept", "application/vnd.github+json")
h.Set("X-GitHub-Api-Version", "2022-11-28")
}
tok, err := p.githubToken(ctx, remote)
if err != nil {
return nil, err
}
if tok != "" {
h.Set("Authorization", "Bearer "+tok)
}
return h, nil
}
// githubToken resolves the bearer token for a remote. Precedence: a per-remote
// credential (password, then username) overrides the server credential.
func (p *GitHubProvider) githubToken(ctx context.Context, remote models.Remote) (string, error) {
if remote.Password != "" {
return remote.Password, nil
}
if remote.Username != "" {
return remote.Username, nil
}
if c := p.serverCredential(); c != nil {
return c.Token(ctx)
}
return "", nil
}
func (p *GitHubProvider) serverCredential() githubauth.Credential {
if p.serverCred != nil {
return p.serverCred
}
return githubauth.Server()
}
func copyHeaders(req *http.Request, h http.Header) {
for k, vals := range h {
for _, v := range vals {
req.Header.Add(k, v)
}
}
}
func compilePatterns(patterns []string) ([]*regexp.Regexp, error) {
var out []*regexp.Regexp
for _, p := range patterns {
re, err := regexp.Compile(p)
if err != nil {
return nil, fmt.Errorf("invalid pattern %q: %w", p, err)
}
out = append(out, re)
}
return out, nil
}
func matchesAny(res []*regexp.Regexp, s string) bool {
if len(res) == 0 {
return true
}
for _, re := range res {
if re.MatchString(s) {
return true
}
}
return false
}
+497
View File
@@ -0,0 +1,497 @@
package alpine
import (
"archive/tar"
"bytes"
"compress/gzip"
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"io"
"net/http"
"net/http/httptest"
"strconv"
"strings"
"sync"
"testing"
"time"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/internal/testsupport"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// fakeStore is an in-memory provider.RemoteMetadataStore + AlpineMetadata
// store/reader/deleter keyed by file_path, mirroring the (repo_name, file_path)
// uniqueness of the real alpine_metadata table.
type fakeStore struct {
mu sync.Mutex
rows map[string]provider.AlpineMetadata
}
func newFakeStore() *fakeStore { return &fakeStore{rows: map[string]provider.AlpineMetadata{}} }
func (f *fakeStore) InsertAlpineMetadata(_ context.Context, m *provider.AlpineMetadata) error {
f.mu.Lock()
defer f.mu.Unlock()
if _, ok := f.rows[m.FilePath]; ok {
return nil // ON CONFLICT DO NOTHING
}
f.rows[m.FilePath] = *m
return nil
}
func (f *fakeStore) DeleteAlpineMetadata(_ context.Context, _, filePath string) error {
f.mu.Lock()
defer f.mu.Unlock()
delete(f.rows, filePath)
return nil
}
func (f *fakeStore) ListAlpineMetadataEntries(ctx context.Context, _ string) ([]provider.AlpineMetadata, error) {
if err := ctx.Err(); err != nil {
return nil, err
}
f.mu.Lock()
defer f.mu.Unlock()
out := make([]provider.AlpineMetadata, 0, len(f.rows))
for _, m := range f.rows {
out = append(out, m)
}
return out, nil
}
// The generic RemoteMetadataStore surface (rpm/deb) is unused by the alpine
// github provider but required to satisfy the interface passed to ServeRemote.
func (f *fakeStore) InsertRPMMetadata(context.Context, *provider.RPMMetadata) error { return nil }
func (f *fakeStore) DeleteRPMMetadata(context.Context, string, string) error { return nil }
func (f *fakeStore) ListRPMMetadataEntries(context.Context, string) ([]provider.RPMMetadata, error) {
return nil, nil
}
func (f *fakeStore) InsertDebMetadata(context.Context, *provider.DebMetadata) error { return nil }
func (f *fakeStore) DeleteDebMetadata(context.Context, string, string) error { return nil }
var _ provider.RemoteMetadataStore = (*fakeStore)(nil)
// githubFixture serves the releases API and the .apk asset downloads (with Range
// support) for a set of packages. digest controls whether the asset carries a
// sha256 digest (change-detection path) or not.
type githubFixture struct {
srv *httptest.Server
apkBytes map[string][]byte
rangeHit map[string]int
fullHit map[string]int
etag string
releasesHit int
notModHit int
releaseAuth string
assetAuth string
mu sync.Mutex
}
func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture {
t.Helper()
f := &githubFixture{
apkBytes: map[string][]byte{},
rangeHit: map[string]int{},
fullHit: map[string]int{},
}
f.apkBytes["demo-1.2.3-r0.apk"] = testsupport.MinimalApk("demo", "1.2.3-r0", "x86_64")
mux := http.NewServeMux()
mux.HandleFunc("/repos/acme/tools/releases", func(w http.ResponseWriter, r *http.Request) {
page := r.URL.Query().Get("page")
if page != "" && page != "1" {
w.Write([]byte("[]"))
return
}
f.mu.Lock()
f.releasesHit++
f.releaseAuth = r.Header.Get("Authorization")
etag := f.etag
if etag != "" && r.Header.Get("If-None-Match") == etag {
f.notModHit++
f.mu.Unlock()
w.WriteHeader(http.StatusNotModified)
return
}
f.mu.Unlock()
if etag != "" {
w.Header().Set("ETag", etag)
}
var assets []map[string]any
for name := range f.apkBytes {
a := map[string]any{
"name": name,
"size": len(f.apkBytes[name]),
"browser_download_url": f.srv.URL + "/acme/tools/releases/download/v1.2.3/" + name,
}
if withDigest {
sum := sha256.Sum256(f.apkBytes[name])
a["digest"] = "sha256:" + hex.EncodeToString(sum[:])
}
assets = append(assets, a)
}
rel := []map[string]any{{"tag_name": "v1.2.3", "draft": false, "assets": assets}}
json.NewEncoder(w).Encode(rel)
})
mux.HandleFunc("/acme/tools/releases/download/", func(w http.ResponseWriter, r *http.Request) {
name := r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]
body, ok := f.apkBytes[name]
if !ok {
http.Error(w, "not found", 404)
return
}
rng := r.Header.Get("Range")
f.mu.Lock()
f.assetAuth = r.Header.Get("Authorization")
if rng != "" {
f.rangeHit[name]++
} else {
f.fullHit[name]++
}
f.mu.Unlock()
if rng == "" {
w.WriteHeader(200)
w.Write(body)
return
}
var end int
fmt.Sscanf(rng, "bytes=0-%d", &end)
if end >= len(body)-1 {
end = len(body) - 1
}
w.Header().Set("Content-Range", fmt.Sprintf("bytes 0-%d/%d", end, len(body)))
w.Header().Set("Content-Length", strconv.Itoa(end+1))
w.WriteHeader(http.StatusPartialContent)
w.Write(body[:end+1])
})
f.srv = httptest.NewServer(mux)
t.Cleanup(f.srv.Close)
return f
}
func (f *githubFixture) remote() models.Remote {
return models.Remote{
Name: "acme-apk",
PackageType: models.PackageGitHubAlpine,
BaseURL: f.srv.URL + "/repos/acme/tools",
ReleasesRemote: "github",
MutableTTL: 3600,
}
}
func newTestProvider() *GitHubProvider {
p := newGitHubProvider()
p.headerInitial = 32 // force the ranged-fetch retry loop against the tiny fixture
p.headerMax = 1 << 20
return p
}
const demoPath = "acme/tools/releases/download/v1.2.3/demo-1.2.3-r0.apk"
func TestGitHubScanDerivesPkginfoFromPrefix(t *testing.T) {
fx := newGitHubFixture(t, true)
p := newTestProvider()
store := newFakeStore()
if err := p.scan(context.Background(), fx.remote(), store); err != nil {
t.Fatalf("scan: %v", err)
}
metas, _ := store.ListAlpineMetadataEntries(context.Background(), "acme-apk")
if len(metas) != 1 {
t.Fatalf("want 1 metadata row, got %d", len(metas))
}
m := metas[0]
if m.Name != "demo" || m.Version != "1.2.3-r0" || m.Arch != "x86_64" {
t.Fatalf("bad .PKGINFO fields: %+v", m)
}
if m.FilePath != demoPath {
t.Fatalf("FilePath = %q, want %q", m.FilePath, demoPath)
}
if int(m.DownloadSize) != len(fx.apkBytes["demo-1.2.3-r0.apk"]) {
t.Fatalf("DownloadSize = %d, want %d", m.DownloadSize, len(fx.apkBytes["demo-1.2.3-r0.apk"]))
}
if !strings.HasPrefix(m.Checksum, "Q1") {
t.Fatalf("Checksum not a Q1 pull checksum: %q", m.Checksum)
}
// The C: checksum must equal Q1 over the raw control gzip stream, matching the
// local-upload parser applied to the same bytes.
want, err := parseApk(fx.apkBytes["demo-1.2.3-r0.apk"])
if err != nil {
t.Fatalf("reference parseApk: %v", err)
}
if m.Checksum != want.Checksum {
t.Fatalf("Checksum = %q, want %q (Q1 of control stream)", m.Checksum, want.Checksum)
}
if fx.fullHit["demo-1.2.3-r0.apk"] != 0 {
t.Fatalf("expected no full download, got %d", fx.fullHit["demo-1.2.3-r0.apk"])
}
if fx.rangeHit["demo-1.2.3-r0.apk"] == 0 {
t.Fatalf("expected ranged .PKGINFO fetch")
}
}
func TestGitHubServeRemoteIndexAndRedirect(t *testing.T) {
fx := newGitHubFixture(t, true)
p := newTestProvider()
store := newFakeStore()
remote := fx.remote()
const proxyBase = "https://artifactapi.example"
// The per-arch index is served and triggers the initial scan.
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-apk/x86_64/APKINDEX.tar.gz", nil)
if !p.ServeRemote(rec, req, remote, "x86_64/APKINDEX.tar.gz", proxyBase, store) {
t.Fatal("ServeRemote did not handle APKINDEX")
}
if rec.Code != 200 {
t.Fatalf("APKINDEX bad: code=%d body=%s", rec.Code, rec.Body.String())
}
idx := readAPKIndex(t, rec.Body.Bytes())
if !strings.Contains(idx, "P:demo") || !strings.Contains(idx, "A:x86_64") {
t.Fatalf("APKINDEX missing package record: %s", idx)
}
if !strings.Contains(idx, "C:Q1") {
t.Fatalf("APKINDEX missing pull checksum: %s", idx)
}
// A different arch yields an empty (but valid) index.
rec = httptest.NewRecorder()
req = httptest.NewRequest(http.MethodGet, "/x", nil)
if !p.ServeRemote(rec, req, remote, "aarch64/APKINDEX.tar.gz", proxyBase, store) {
t.Fatal("ServeRemote did not handle aarch64 APKINDEX")
}
if rec.Code != 200 {
t.Fatalf("empty-arch index bad: %d", rec.Code)
}
if got := readAPKIndex(t, rec.Body.Bytes()); strings.Contains(got, "P:demo") {
t.Fatalf("aarch64 index should not carry the x86_64 package: %s", got)
}
// An .apk request arrives in apk's reconstructed shape
// "<arch>/<name>-<version>.apk" (APKINDEX carries no filename), NOT as the
// github-relative FilePath. ServeRemote must resolve it back to the stored
// FilePath before redirecting to the backend releases_remote.
rec = httptest.NewRecorder()
req = httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-apk/x86_64/demo-1.2.3-r0.apk", nil)
if !p.ServeRemote(rec, req, remote, "x86_64/demo-1.2.3-r0.apk", proxyBase, store) {
t.Fatal("ServeRemote did not handle .apk")
}
if rec.Code != http.StatusFound {
t.Fatalf("want 302, got %d", rec.Code)
}
wantLoc := proxyBase + "/api/v1/remote/github/" + demoPath
if got := rec.Header().Get("Location"); got != wantLoc {
t.Fatalf("Location = %q, want %q (must be the stored FilePath, not the inbound path)", got, wantLoc)
}
}
// An apk download whose reconstructed "<arch>/<name>-<version>.apk" matches no
// cached row must 404, never redirect to a bad path.
func TestGitHubServeRemoteApkRedirectNotFound(t *testing.T) {
fx := newGitHubFixture(t, true)
p := newTestProvider()
store := newFakeStore()
remote := fx.remote()
// Warm the cache so the store is populated but lacks the requested package.
if err := p.scan(context.Background(), remote, store); err != nil {
t.Fatalf("warm scan: %v", err)
}
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-apk/x86_64/nope-9.9.9.apk", nil)
if !p.ServeRemote(rec, req, remote, "x86_64/nope-9.9.9.apk", "https://x", store) {
t.Fatal("ServeRemote did not handle .apk")
}
if rec.Code != http.StatusNotFound {
t.Fatalf("want 404 for unknown package, got %d (Location=%q)", rec.Code, rec.Header().Get("Location"))
}
}
// apk requests the index at "./<arch>/APKINDEX.tar.gz"; ServeRemote must collapse
// the dot-segment and synthesize the same index as the un-prefixed request.
func TestGitHubServeRemoteApkDotSegment(t *testing.T) {
fx := newGitHubFixture(t, true)
p := newTestProvider()
store := newFakeStore()
remote := fx.remote()
const proxyBase = "https://artifactapi.example"
serve := func(path string) *httptest.ResponseRecorder {
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-apk/"+path, nil)
if !p.ServeRemote(rec, req, remote, path, proxyBase, store) {
t.Fatalf("ServeRemote did not handle %q", path)
}
return rec
}
plain, dotted := serve("x86_64/APKINDEX.tar.gz"), serve("./x86_64/APKINDEX.tar.gz")
if plain.Code != 200 || dotted.Code != 200 {
t.Fatalf("index: plain=%d dotted=%d, want 200/200", plain.Code, dotted.Code)
}
if !bytes.Equal(plain.Body.Bytes(), dotted.Body.Bytes()) {
t.Error("./<arch>/APKINDEX.tar.gz body differs from the un-prefixed body")
}
}
func TestGitHubServeRemoteRejectsNonPerArchIndex(t *testing.T) {
fx := newGitHubFixture(t, true)
p := newTestProvider()
store := newFakeStore()
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/x", nil)
if !p.ServeRemote(rec, req, fx.remote(), "APKINDEX.tar.gz", "https://x", store) {
t.Fatal("expected handled")
}
if rec.Code != http.StatusNotFound {
t.Fatalf("bare APKINDEX must 404 (per-arch required), got %d", rec.Code)
}
}
// A canceled inbound request must still serve the warm cache (detached context),
// not turn the metadata read into a 500.
func TestGitHubServeRemoteCanceledRequestServesCache(t *testing.T) {
fx := newGitHubFixture(t, true)
p := newTestProvider()
store := newFakeStore()
remote := fx.remote()
if err := p.scan(context.Background(), remote, store); err != nil {
t.Fatalf("warm scan: %v", err)
}
p.mu.Lock()
p.lastScan[remote.Name] = time.Now()
p.mu.Unlock()
ctx, cancel := context.WithCancel(context.Background())
cancel()
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-apk/x86_64/APKINDEX.tar.gz", nil).WithContext(ctx)
if !p.ServeRemote(rec, req, remote, "x86_64/APKINDEX.tar.gz", "https://x", store) {
t.Fatal("ServeRemote did not handle APKINDEX")
}
if rec.Code != http.StatusOK {
t.Fatalf("canceled request must serve cache, not error; got code=%d body=%s", rec.Code, rec.Body.String())
}
if got := readAPKIndex(t, rec.Body.Bytes()); !strings.Contains(got, "P:demo") {
t.Fatalf("expected index served from cache, got %s", got)
}
}
func TestGitHubServeRemoteRedirectRequiresReleasesRemote(t *testing.T) {
fx := newGitHubFixture(t, true)
p := newTestProvider()
store := newFakeStore()
remote := fx.remote()
remote.ReleasesRemote = ""
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/x", nil)
if !p.ServeRemote(rec, req, remote, demoPath, "https://x", store) {
t.Fatal("expected handled")
}
if rec.Code != http.StatusInternalServerError {
t.Fatalf("want 500 when releases_remote unset, got %d", rec.Code)
}
}
func TestGitHubScanPrunesRemovedAssets(t *testing.T) {
fx := newGitHubFixture(t, true)
p := newTestProvider()
store := newFakeStore()
if err := p.scan(context.Background(), fx.remote(), store); err != nil {
t.Fatalf("scan: %v", err)
}
if rows, _ := store.ListAlpineMetadataEntries(context.Background(), "acme-apk"); len(rows) != 1 {
t.Fatalf("want 1 row after first scan, got %d", len(rows))
}
delete(fx.apkBytes, "demo-1.2.3-r0.apk")
if err := p.scan(context.Background(), fx.remote(), store); err != nil {
t.Fatalf("rescan: %v", err)
}
if rows, _ := store.ListAlpineMetadataEntries(context.Background(), "acme-apk"); len(rows) != 0 {
t.Fatalf("want 0 rows after prune, got %d", len(rows))
}
}
func TestGitHubAssetPatternFilter(t *testing.T) {
fx := newGitHubFixture(t, true)
fx.apkBytes["other-9-r0.apk"] = testsupport.MinimalApk("other", "9-r0", "aarch64")
p := newTestProvider()
store := newFakeStore()
remote := fx.remote()
remote.Patterns = []string{`^demo-.*\.apk$`}
if err := p.scan(context.Background(), remote, store); err != nil {
t.Fatalf("scan: %v", err)
}
rows, _ := store.ListAlpineMetadataEntries(context.Background(), "acme-apk")
if len(rows) != 1 || rows[0].Name != "demo" {
t.Fatalf("pattern filter failed, rows=%+v", rows)
}
}
// Multi-arch: each asset's index record lands under its own arch bucket.
func TestGitHubServeRemotePerArchGrouping(t *testing.T) {
fx := newGitHubFixture(t, true)
fx.apkBytes["demo-1.2.3-r0-aarch64.apk"] = testsupport.MinimalApk("demo", "1.2.3-r0", "aarch64")
p := newTestProvider()
store := newFakeStore()
remote := fx.remote()
const proxyBase = "https://x"
if err := p.scan(context.Background(), remote, store); err != nil {
t.Fatalf("scan: %v", err)
}
serve := func(arch string) string {
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/x", nil)
if !p.ServeRemote(rec, req, remote, arch+"/APKINDEX.tar.gz", proxyBase, store) {
t.Fatalf("ServeRemote did not handle %s", arch)
}
return readAPKIndex(t, rec.Body.Bytes())
}
x86 := serve("x86_64")
if !strings.Contains(x86, "A:x86_64") || strings.Contains(x86, "A:aarch64") {
t.Fatalf("x86_64 index leaked another arch: %s", x86)
}
arm := serve("aarch64")
if !strings.Contains(arm, "A:aarch64") || strings.Contains(arm, "A:x86_64") {
t.Fatalf("aarch64 index leaked another arch: %s", arm)
}
}
func readAPKIndex(t *testing.T, gzBytes []byte) string {
t.Helper()
gz, err := gzip.NewReader(bytes.NewReader(gzBytes))
if err != nil {
t.Fatalf("gzip: %v", err)
}
tr := tar.NewReader(gz)
for {
hdr, err := tr.Next()
if err != nil {
t.Fatal("APKINDEX member missing from tar.gz")
}
if strings.TrimPrefix(hdr.Name, "./") == "APKINDEX" {
body, err := io.ReadAll(tr)
if err != nil {
t.Fatalf("read APKINDEX: %v", err)
}
return string(body)
}
}
}
+238
View File
@@ -0,0 +1,238 @@
package alpine
import (
"context"
"crypto/rand"
"encoding/hex"
"log/slog"
"os"
"sync"
"time"
"golang.org/x/time/rate"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
const (
syncLeaseDuration = 15 * time.Minute
defaultSyncFreshness = 5 * time.Minute
jobQueueDepth = 256
)
// SyncStore is the persistence surface the alpine syncer needs: the metadata
// cache it primes plus the shared sync-state coordination (remote enumeration
// and the per-remote lease). *database.DB satisfies it.
type SyncStore interface {
provider.RemoteMetadataStore
ListGitHubAlpineRemotes(ctx context.Context) ([]models.Remote, error)
ClaimGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (claimed bool, etag string, err error)
ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error
}
// SyncConfig tunes the shared syncer. Zero values fall back to safe defaults.
type SyncConfig struct {
RatePerSec float64
Burst int
Workers int
PollInterval time.Duration
}
type syncJob struct {
remote models.Remote
prime bool
}
// Syncer is the single per-process background worker that keeps every
// github_alpine remote's derived metadata fresh. It owns a deduped work queue, a
// pool of workers, and a global token-bucket rate limiter shared across all
// remotes and bound onto the github_alpine provider. Periodic checks are gated by
// a shared DB lease so, across replicas, only one performs each scan.
type Syncer struct {
store SyncStore
prov *GitHubProvider
limiter *rate.Limiter
cfg SyncConfig
owner string
jobs chan syncJob
mu sync.Mutex
active map[string]bool
}
// NewSyncer builds the syncer bound to the process-wide github_alpine provider
// singleton. Call Run to start it.
func NewSyncer(store SyncStore, cfg SyncConfig) *Syncer {
return newSyncer(store, gitHubProvider, cfg)
}
func newSyncer(store SyncStore, prov *GitHubProvider, cfg SyncConfig) *Syncer {
if cfg.RatePerSec <= 0 {
cfg.RatePerSec = 1
}
if cfg.Burst <= 0 {
cfg.Burst = 5
}
if cfg.Workers <= 0 {
cfg.Workers = 3
}
if cfg.PollInterval <= 0 {
cfg.PollInterval = 60 * time.Second
}
lim := rate.NewLimiter(rate.Limit(cfg.RatePerSec), cfg.Burst)
s := &Syncer{
store: store,
prov: prov,
limiter: lim,
cfg: cfg,
owner: leaseOwner(),
jobs: make(chan syncJob, jobQueueDepth),
active: map[string]bool{},
}
prov.limiter = lim
prov.syncer = s
return s
}
// Run starts the worker pool and the periodic scheduler and blocks until ctx is
// canceled, at which point it drains in-flight scans and returns.
func (s *Syncer) Run(ctx context.Context) {
slog.Info("github_alpine syncer started",
"rate_per_sec", s.cfg.RatePerSec, "burst", s.cfg.Burst,
"workers", s.cfg.Workers, "poll_interval", s.cfg.PollInterval, "owner", s.owner)
var wg sync.WaitGroup
for i := 0; i < s.cfg.Workers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
s.worker(ctx)
}()
}
ticker := time.NewTicker(s.cfg.PollInterval)
defer ticker.Stop()
s.schedule(ctx)
for {
select {
case <-ctx.Done():
wg.Wait()
slog.Info("github_alpine syncer stopped")
return
case <-ticker.C:
s.schedule(ctx)
}
}
}
// schedule enqueues a periodic check for every github_alpine remote. The DB lease
// enforces the per-remote mutable_ttl cadence and cross-replica coordination.
func (s *Syncer) schedule(ctx context.Context) {
remotes, err := s.store.ListGitHubAlpineRemotes(ctx)
if err != nil {
slog.Error("github_alpine syncer: list remotes", "error", err)
return
}
for _, r := range remotes {
s.enqueue(r, false)
}
}
// EnqueuePrime queues an immediate background prime for a freshly created remote.
func (s *Syncer) EnqueuePrime(remote models.Remote) {
if s == nil {
return
}
s.enqueue(remote, true)
}
// enqueue adds a job unless the remote is already queued or in-flight, coalescing
// duplicate requests down to one scan. It never blocks.
func (s *Syncer) enqueue(remote models.Remote, prime bool) {
s.mu.Lock()
if s.active[remote.Name] {
s.mu.Unlock()
return
}
s.active[remote.Name] = true
s.mu.Unlock()
select {
case s.jobs <- syncJob{remote: remote, prime: prime}:
default:
s.mu.Lock()
delete(s.active, remote.Name)
s.mu.Unlock()
}
}
func (s *Syncer) worker(ctx context.Context) {
for {
select {
case <-ctx.Done():
return
case job := <-s.jobs:
s.process(ctx, job)
}
}
}
// process claims the shared lease and, if won, runs an incremental scan. Losing
// the claim (another replica scanning, or not yet due) is a no-op.
func (s *Syncer) process(ctx context.Context, job syncJob) {
defer func() {
s.mu.Lock()
delete(s.active, job.remote.Name)
s.mu.Unlock()
}()
freshness := time.Duration(job.remote.MutableTTL) * time.Second
if freshness <= 0 {
freshness = defaultSyncFreshness
}
if job.prime {
freshness = 0
}
claimed, etag, err := s.store.ClaimGitHubAlpineSyncLease(ctx, job.remote.Name, s.owner, freshness, syncLeaseDuration)
if err != nil {
slog.Error("github_alpine syncer: claim lease", "remote", job.remote.Name, "error", err)
return
}
if !claimed {
return
}
scanCtx, cancel := context.WithTimeout(ctx, s.prov.scanTimeout)
defer cancel()
newEtag, changed, scanErr := s.prov.scanWithState(scanCtx, job.remote, s.store, etag)
releaseEtag := etag
if scanErr == nil {
releaseEtag = newEtag
} else {
slog.Error("github_alpine syncer: scan failed", "remote", job.remote.Name, "error", scanErr)
}
relCtx, relCancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
defer relCancel()
if err := s.store.ReleaseGitHubAlpineSyncLease(relCtx, job.remote.Name, s.owner, releaseEtag, time.Now()); err != nil {
slog.Warn("github_alpine syncer: release lease", "remote", job.remote.Name, "error", err)
}
if scanErr == nil && changed {
slog.Info("github_alpine syncer: refreshed", "remote", job.remote.Name, "prime", job.prime)
}
}
// leaseOwner is a per-replica identity for the lease: hostname plus a random
// suffix so restarts and colocated replicas never collide.
func leaseOwner() string {
host, _ := os.Hostname()
var b [6]byte
_, _ = rand.Read(b[:])
return host + "-" + hex.EncodeToString(b[:])
}
+300
View File
@@ -0,0 +1,300 @@
package alpine
import (
"context"
"net/http"
"net/http/httptest"
"sync"
"testing"
"time"
"golang.org/x/time/rate"
"git.unkin.net/unkin/artifactapi/internal/provider"
"git.unkin.net/unkin/artifactapi/internal/testsupport"
"git.unkin.net/unkin/artifactapi/pkg/models"
)
// fakeSyncStore is an in-memory SyncStore: the metadata cache (via the embedded
// fakeStore) plus the shared sync-state lease, whose claim mirrors the atomic
// semantics of the real SQL (recency gate AND no live lease).
type fakeSyncStore struct {
*fakeStore
mu sync.Mutex
remotes []models.Remote
leaseOwner map[string]string
leaseExp map[string]time.Time
lastSynced map[string]time.Time
etags map[string]string
}
func newFakeSyncStore() *fakeSyncStore {
return &fakeSyncStore{
fakeStore: newFakeStore(),
leaseOwner: map[string]string{},
leaseExp: map[string]time.Time{},
lastSynced: map[string]time.Time{},
etags: map[string]string{},
}
}
func (f *fakeSyncStore) ListGitHubAlpineRemotes(_ context.Context) ([]models.Remote, error) {
f.mu.Lock()
defer f.mu.Unlock()
return append([]models.Remote(nil), f.remotes...), nil
}
func (f *fakeSyncStore) ClaimGitHubAlpineSyncLease(_ context.Context, name, owner string, freshness, lease time.Duration) (bool, string, error) {
f.mu.Lock()
defer f.mu.Unlock()
now := time.Now()
ls, hasLS := f.lastSynced[name]
exp, hasExp := f.leaseExp[name]
freshOK := !hasLS || now.Sub(ls) >= freshness
leaseOK := !hasExp || exp.Before(now)
if freshOK && leaseOK {
f.leaseOwner[name] = owner
f.leaseExp[name] = now.Add(lease)
return true, f.etags[name], nil
}
return false, "", nil
}
func (f *fakeSyncStore) ReleaseGitHubAlpineSyncLease(_ context.Context, name, owner, etag string, syncedAt time.Time) error {
f.mu.Lock()
defer f.mu.Unlock()
if f.leaseOwner[name] != owner {
return nil
}
f.lastSynced[name] = syncedAt
f.etags[name] = etag
delete(f.leaseOwner, name)
delete(f.leaseExp, name)
return nil
}
func testSyncConfig() SyncConfig {
return SyncConfig{RatePerSec: 1000, Burst: 100, Workers: 1, PollInterval: time.Hour}
}
// (a) A 304 conditional response must derive nothing: no asset fetches and
// changed=false, so an unchanged repo is nearly free.
func TestSyncerConditionalNotModifiedSkipsDerive(t *testing.T) {
fx := newGitHubFixture(t, true)
fx.etag = `"v1"`
p := newTestProvider()
store := newFakeStore()
remote := fx.remote()
etag1, changed, err := p.scanWithState(context.Background(), remote, store, "")
if err != nil {
t.Fatalf("first scan: %v", err)
}
if !changed || etag1 != `"v1"` {
t.Fatalf("first scan changed=%v etag=%q, want true and \"v1\"", changed, etag1)
}
priorRange := fx.rangeHit["demo-1.2.3-r0.apk"]
if priorRange == 0 {
t.Fatal("first scan should have fetched the asset .PKGINFO")
}
etag2, changed2, err := p.scanWithState(context.Background(), remote, store, etag1)
if err != nil {
t.Fatalf("second scan: %v", err)
}
if changed2 {
t.Fatal("304 scan must report changed=false")
}
if etag2 != etag1 {
t.Fatalf("etag changed across 304: %q -> %q", etag1, etag2)
}
if fx.notModHit != 1 {
t.Fatalf("want exactly one 304 releases response, got %d", fx.notModHit)
}
if got := fx.rangeHit["demo-1.2.3-r0.apk"]; got != priorRange {
t.Fatalf("304 scan re-fetched asset .PKGINFO: %d -> %d", priorRange, got)
}
}
// (b) On a real change, only the newly added asset is derived.
func TestSyncerIncrementalDerivesOnlyNewAsset(t *testing.T) {
fx := newGitHubFixture(t, true)
fx.etag = `"v1"`
p := newTestProvider()
store := newFakeStore()
remote := fx.remote()
if _, _, err := p.scanWithState(context.Background(), remote, store, ""); err != nil {
t.Fatalf("first scan: %v", err)
}
demoRange := fx.rangeHit["demo-1.2.3-r0.apk"]
fx.apkBytes["other-9-r0.apk"] = testsupport.MinimalApk("other", "9-r0", "aarch64")
fx.etag = `"v2"`
if _, changed, err := p.scanWithState(context.Background(), remote, store, `"v1"`); err != nil || !changed {
t.Fatalf("second scan changed=%v err=%v", changed, err)
}
rows, _ := store.ListAlpineMetadataEntries(context.Background(), remote.Name)
if len(rows) != 2 {
t.Fatalf("want 2 cached rows after incremental derive, got %d", len(rows))
}
if got := fx.rangeHit["demo-1.2.3-r0.apk"]; got != demoRange {
t.Fatalf("already-cached asset was re-fetched: %d -> %d", demoRange, got)
}
if fx.rangeHit["other-9-r0.apk"] == 0 {
t.Fatal("newly added asset was not derived")
}
}
// (c) The shared limiter caps the request rate.
func TestRateLimiterCapsRequestRate(t *testing.T) {
fx := newGitHubFixture(t, true)
p := newTestProvider()
p.limiter = rate.NewLimiter(rate.Every(120*time.Millisecond), 1)
remote := fx.remote()
start := time.Now()
for i := 0; i < 3; i++ {
if _, _, _, err := p.fetchReleases(context.Background(), remote, ""); err != nil {
t.Fatalf("fetchReleases %d: %v", i, err)
}
}
if elapsed := time.Since(start); elapsed < 200*time.Millisecond {
t.Fatalf("rate limiter did not throttle: 3 calls took %v, want >= 200ms", elapsed)
}
}
// (d) Concurrent enqueues for the same remote coalesce to a single queued job.
func TestSyncerEnqueueDedup(t *testing.T) {
store := newFakeSyncStore()
p := newTestProvider()
s := newSyncer(store, p, testSyncConfig())
remote := models.Remote{Name: "acme-apk", PackageType: models.PackageGitHubAlpine, MutableTTL: 3600}
var wg sync.WaitGroup
for i := 0; i < 10; i++ {
wg.Add(1)
go func() { defer wg.Done(); s.enqueue(remote, false) }()
}
wg.Wait()
if got := len(s.jobs); got != 1 {
t.Fatalf("want exactly 1 coalesced job, got %d", got)
}
}
// (e) Prime-on-create enqueues a prime job.
func TestSyncerEnqueuePrime(t *testing.T) {
store := newFakeSyncStore()
p := newTestProvider()
s := newSyncer(store, p, testSyncConfig())
remote := models.Remote{Name: "acme-apk", PackageType: models.PackageGitHubAlpine, MutableTTL: 3600}
s.EnqueuePrime(remote)
select {
case job := <-s.jobs:
if !job.prime || job.remote.Name != "acme-apk" {
t.Fatalf("bad prime job: %+v", job)
}
default:
t.Fatal("EnqueuePrime did not enqueue a job")
}
}
// (f) A held lease prevents a second replica from scanning.
func TestSyncerLeasePreventsSecondReplica(t *testing.T) {
fx := newGitHubFixture(t, true)
fx.etag = `"v1"`
store := newFakeSyncStore()
p := newTestProvider()
s := newSyncer(store, p, testSyncConfig())
remote := fx.remote()
claimed, _, err := store.ClaimGitHubAlpineSyncLease(context.Background(), remote.Name, "replica-1", time.Duration(remote.MutableTTL)*time.Second, syncLeaseDuration)
if err != nil || !claimed {
t.Fatalf("replica-1 claim: claimed=%v err=%v", claimed, err)
}
s.process(context.Background(), syncJob{remote: remote})
if fx.releasesHit != 0 {
t.Fatalf("second replica scanned while lease held: %d releases calls", fx.releasesHit)
}
if rows, _ := store.ListAlpineMetadataEntries(context.Background(), remote.Name); len(rows) != 0 {
t.Fatalf("second replica derived metadata while lease held: %d rows", len(rows))
}
}
// With the syncer wired and the cache empty, an index request enqueues a prime
// and returns a retryable 503 when it has not landed within the cold wait.
func TestServeRemoteColdStartReturns503(t *testing.T) {
fx := newGitHubFixture(t, true)
store := newFakeSyncStore()
p := newTestProvider()
p.coldWait = 300 * time.Millisecond
_ = newSyncer(store, p, testSyncConfig()) // binds p.syncer, but no workers running
remote := fx.remote()
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-apk/x86_64/APKINDEX.tar.gz", nil)
if !p.ServeRemote(rec, req, remote, "x86_64/APKINDEX.tar.gz", "https://x", store) {
t.Fatal("ServeRemote did not handle APKINDEX")
}
if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("cold empty cache must return 503, got %d", rec.Code)
}
if rec.Header().Get("Retry-After") == "" {
t.Fatal("503 should carry Retry-After")
}
if got := len(p.syncer.jobs); got != 1 {
t.Fatalf("cold start did not enqueue a prime, jobs=%d", got)
}
}
// With the cache warm, the same request serves the index immediately (no 503).
func TestServeRemoteWarmCacheServesImmediately(t *testing.T) {
fx := newGitHubFixture(t, true)
store := newFakeSyncStore()
p := newTestProvider()
_ = newSyncer(store, p, testSyncConfig())
remote := fx.remote()
if err := p.scan(context.Background(), remote, store); err != nil {
t.Fatalf("warm scan: %v", err)
}
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-apk/x86_64/APKINDEX.tar.gz", nil)
if !p.ServeRemote(rec, req, remote, "x86_64/APKINDEX.tar.gz", "https://x", store) {
t.Fatal("ServeRemote did not handle APKINDEX")
}
if rec.Code != http.StatusOK {
t.Fatalf("warm cache must serve 200, got %d body=%s", rec.Code, rec.Body.String())
}
}
// A prime job (freshness 0) runs even right after a sync; a periodic job at the
// same moment is gated by the recency window.
func TestSyncerPrimeBypassesRecencyPeriodicDoesNot(t *testing.T) {
fx := newGitHubFixture(t, true)
fx.etag = `"v1"`
store := newFakeSyncStore()
p := newTestProvider()
s := newSyncer(store, p, testSyncConfig())
remote := fx.remote()
var _ provider.RemoteMetadataStore = store
s.process(context.Background(), syncJob{remote: remote, prime: true})
if rows, _ := store.ListAlpineMetadataEntries(context.Background(), remote.Name); len(rows) != 1 {
t.Fatalf("prime did not derive: %d rows", len(rows))
}
releasesAfterPrime := fx.releasesHit
s.process(context.Background(), syncJob{remote: remote, prime: false})
if fx.releasesHit != releasesAfterPrime {
t.Fatalf("periodic scan ran inside recency window: %d -> %d releases calls", releasesAfterPrime, fx.releasesHit)
}
}
+44
View File
@@ -115,6 +115,50 @@ type DebMetadata struct {
SHA256 string
}
// AlpineMetadataStore / AlpineMetadataDeleter / AlpineMetadataReader are the
// Alpine-specific persistence surfaces. They are kept separate from the shared
// RPM/Deb metadata interfaces so the apk provider can type-assert the generic
// MetadataStore/MetadataDeleter/FileStore it is handed without widening (and
// thus perturbing the test doubles of) the rpm and deb providers. *database.DB
// satisfies all three.
type AlpineMetadataStore interface {
InsertAlpineMetadata(ctx context.Context, meta *AlpineMetadata) error
}
type AlpineMetadataDeleter interface {
DeleteAlpineMetadata(ctx context.Context, repoName, filePath string) error
}
type AlpineMetadataReader interface {
ListAlpineMetadataEntries(ctx context.Context, repoName string) ([]AlpineMetadata, error)
}
// AlpineMetadata is the derived per-package metadata for an Alpine .apk, holding
// the fields an APKINDEX record carries plus the apk pull checksum (Q1…, the
// sha1 of the control gzip stream) and the download/installed sizes.
type AlpineMetadata struct {
RepoName string
FilePath string
ContentHash string
Checksum string // C: "Q1" + base64(sha1(control gzip stream))
Name string // P:
Version string // V:
Arch string // A:
DownloadSize int64 // S: on-disk .apk size
InstalledSize int64 // I: unpacked size from .PKGINFO
Description string // T:
URL string // U:
License string // L:
Origin string // o:
Maintainer string // m:
BuildTime int64 // t:
Commit string // c:
ProviderPriority string // k:
Depends []string // D:
Provides []string // p:
InstallIf []string // i:
}
type RPMMetadata struct {
RepoName string
FilePath string
+14 -3
View File
@@ -20,7 +20,7 @@ import (
"git.unkin.net/unkin/artifactapi/internal/database"
"git.unkin.net/unkin/artifactapi/internal/gc"
"git.unkin.net/unkin/artifactapi/internal/githubauth"
_ "git.unkin.net/unkin/artifactapi/internal/provider/alpine"
"git.unkin.net/unkin/artifactapi/internal/provider/alpine"
"git.unkin.net/unkin/artifactapi/internal/provider/deb"
_ "git.unkin.net/unkin/artifactapi/internal/provider/docker"
_ "git.unkin.net/unkin/artifactapi/internal/provider/generic"
@@ -52,6 +52,7 @@ type Server struct {
gc *gc.Collector
syncer *rpm.Syncer
debSyncer *deb.Syncer
alpineSyncer *alpine.Syncer
}
func New(cfg *config.Config, version string) (*Server, error) {
@@ -105,6 +106,12 @@ func New(cfg *config.Config, version string) (*Server, error) {
Workers: cfg.GitHubSyncWorkers,
PollInterval: time.Duration(cfg.GitHubSyncPollInterval) * time.Second,
})
alpineSyncer := alpine.NewSyncer(db, alpine.SyncConfig{
RatePerSec: cfg.GitHubSyncRatePerSec,
Burst: cfg.GitHubSyncBurst,
Workers: cfg.GitHubSyncWorkers,
PollInterval: time.Duration(cfg.GitHubSyncPollInterval) * time.Second,
})
// The terraform registry signs with a GPG key. A configured file wins (BYO
// key); otherwise artifactapi generates one on first start and persists it in
@@ -138,6 +145,7 @@ func New(cfg *config.Config, version string) (*Server, error) {
gc: collector,
syncer: syncer,
debSyncer: debSyncer,
alpineSyncer: alpineSyncer,
}
s.router = s.routes()
@@ -168,8 +176,9 @@ func (s *Server) routes() chi.Router {
r.Mount("/v2", proxyHandler.DockerV2Routes())
remotesHandler := v2.NewRemotesHandler(s.db, map[models.PackageType]v2.Primer{
models.PackageGitHubRPM: s.syncer,
models.PackageGitHubDeb: s.debSyncer,
models.PackageGitHubRPM: s.syncer,
models.PackageGitHubDeb: s.debSyncer,
models.PackageGitHubAlpine: s.alpineSyncer,
})
virtualsHandler := v2.NewVirtualsHandler(s.db)
healthHandler := v2.NewHealthHandler(s.db, s.cache, s.store)
@@ -239,6 +248,7 @@ func (s *Server) Run(ctx context.Context) error {
go s.gc.Run(ctx)
go s.syncer.Run(ctx)
go s.debSyncer.Run(ctx)
go s.alpineSyncer.Run(ctx)
httpServer := s.newHTTPServer()
@@ -261,6 +271,7 @@ func (s *Server) RunOnListener(ctx context.Context, ln net.Listener) error {
go s.gc.Run(ctx)
go s.syncer.Run(ctx)
go s.debSyncer.Run(ctx)
go s.alpineSyncer.Run(ctx)
httpServer := s.newHTTPServer()
+38
View File
@@ -0,0 +1,38 @@
package testsupport
import (
"bytes"
"fmt"
)
// MinimalApk builds a valid-enough Alpine package in pure Go (no committed
// binary fixture, no abuild): two concatenated, independently gzipped tar
// streams -- a control stream carrying .PKGINFO and a data stream carrying a
// single payload file. It mirrors MinimalDeb/MinimalRPM and is parseable by the
// alpine provider (which derives arch/name/version and the Q1 pull checksum from
// the control stream).
func MinimalApk(name, version, arch string) []byte {
pkginfo := fmt.Sprintf(
"# generated by testsupport\n"+
"pkgname = %s\n"+
"pkgver = %s\n"+
"arch = %s\n"+
"pkgdesc = minimal test package\n"+
"url = https://example.com/%s\n"+
"license = MIT\n"+
"origin = %s\n"+
"maintainer = e2e <e2e@example.com>\n"+
"builddate = 1700000000\n"+
"size = 4\n"+
"depend = so:libc.musl-x86_64.so.1\n"+
"provides = cmd:%s=%s\n",
name, version, arch, name, name, name, version)
control := gzipBytes(tarSingle(".PKGINFO", []byte(pkginfo)))
data := gzipBytes(tarSingle("usr/bin/"+name, []byte("body")))
var buf bytes.Buffer
buf.Write(control)
buf.Write(data)
return buf.Bytes()
}
+28 -26
View File
@@ -5,35 +5,37 @@ import "fmt"
type PackageType string
const (
PackageGeneric PackageType = "generic"
PackageDocker PackageType = "docker"
PackageHelm PackageType = "helm"
PackagePyPI PackageType = "pypi"
PackageNPM PackageType = "npm"
PackageRPM PackageType = "rpm"
PackageDeb PackageType = "deb"
PackageAlpine PackageType = "alpine"
PackagePuppet PackageType = "puppet"
PackageTerraform PackageType = "terraform"
PackageGoProxy PackageType = "goproxy"
PackageGitHubRPM PackageType = "github_rpm"
PackageGitHubDeb PackageType = "github_deb"
PackageGeneric PackageType = "generic"
PackageDocker PackageType = "docker"
PackageHelm PackageType = "helm"
PackagePyPI PackageType = "pypi"
PackageNPM PackageType = "npm"
PackageRPM PackageType = "rpm"
PackageDeb PackageType = "deb"
PackageAlpine PackageType = "alpine"
PackagePuppet PackageType = "puppet"
PackageTerraform PackageType = "terraform"
PackageGoProxy PackageType = "goproxy"
PackageGitHubRPM PackageType = "github_rpm"
PackageGitHubDeb PackageType = "github_deb"
PackageGitHubAlpine PackageType = "github_alpine"
)
var validPackageTypes = map[PackageType]bool{
PackageGeneric: true,
PackageDocker: true,
PackageHelm: true,
PackagePyPI: true,
PackageNPM: true,
PackageRPM: true,
PackageDeb: true,
PackageAlpine: true,
PackagePuppet: true,
PackageTerraform: true,
PackageGoProxy: true,
PackageGitHubRPM: true,
PackageGitHubDeb: true,
PackageGeneric: true,
PackageDocker: true,
PackageHelm: true,
PackagePyPI: true,
PackageNPM: true,
PackageRPM: true,
PackageDeb: true,
PackageAlpine: true,
PackagePuppet: true,
PackageTerraform: true,
PackageGoProxy: true,
PackageGitHubRPM: true,
PackageGitHubDeb: true,
PackageGitHubAlpine: true,
}
func (p PackageType) Valid() bool {
+2
View File
@@ -19,6 +19,8 @@ func TestPackageTypeValid(t *testing.T) {
models.PackageTerraform,
models.PackageGoProxy,
models.PackageGitHubRPM,
models.PackageGitHubDeb,
models.PackageGitHubAlpine,
}
for _, pt := range valid {
if !pt.Valid() {