Add github_deb metadata-only package type #112
@@ -12,18 +12,20 @@ import (
|
||||
)
|
||||
|
||||
// Primer enqueues a background metadata prime for a newly created remote so the
|
||||
// create call never blocks on a derive. *rpm.Syncer satisfies it.
|
||||
// create call never blocks on a derive. *rpm.Syncer and *deb.Syncer satisfy it.
|
||||
type Primer interface {
|
||||
EnqueuePrime(remote models.Remote)
|
||||
}
|
||||
|
||||
type RemotesHandler struct {
|
||||
db *database.DB
|
||||
primer Primer
|
||||
db *database.DB
|
||||
primers map[models.PackageType]Primer
|
||||
}
|
||||
|
||||
func NewRemotesHandler(db *database.DB, primer Primer) *RemotesHandler {
|
||||
return &RemotesHandler{db: db, primer: primer}
|
||||
// NewRemotesHandler wires the handler to the per-type metadata primers. primers
|
||||
// may be nil; a package type with no registered primer simply skips priming.
|
||||
func NewRemotesHandler(db *database.DB, primers map[models.PackageType]Primer) *RemotesHandler {
|
||||
return &RemotesHandler{db: db, primers: primers}
|
||||
}
|
||||
|
||||
func (h *RemotesHandler) Routes() chi.Router {
|
||||
@@ -84,10 +86,10 @@ func (h *RemotesHandler) create(w http.ResponseWriter, r *http.Request) {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
// Prime a github_rpm remote's metadata in the background so its first
|
||||
// repodata request is served from cache instead of a cold on-demand derive.
|
||||
if h.primer != nil && remote.PackageType == models.PackageGitHubRPM {
|
||||
h.primer.EnqueuePrime(remote)
|
||||
// Prime a metadata-only remote (github_rpm/github_deb) in the background so
|
||||
// its first index request is served from cache instead of a cold derive.
|
||||
if primer := h.primers[remote.PackageType]; primer != nil {
|
||||
primer.EnqueuePrime(remote)
|
||||
}
|
||||
writeJSON(w, http.StatusCreated, remote)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,70 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
|
||||
"git.unkin.net/unkin/artifactapi/pkg/models"
|
||||
)
|
||||
|
||||
// ListGitHubDebRemotes returns every github_deb remote so the syncer can sweep
|
||||
// them on each poll tick.
|
||||
func (db *DB) ListGitHubDebRemotes(ctx context.Context) ([]models.Remote, error) {
|
||||
rows, err := db.Pool.Query(ctx, `SELECT `+remoteCols+` FROM remotes WHERE package_type = $1 ORDER BY name`, models.PackageGitHubDeb)
|
||||
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()
|
||||
}
|
||||
|
||||
// ClaimGitHubDebSyncLease atomically claims the per-remote sync lease. It
|
||||
// succeeds only when the remote is due (never synced, or synced longer than
|
||||
// freshness ago) and no live lease is held by another replica. A zero freshness
|
||||
// (prime scans) ignores the recency gate. The returned etag is the stored
|
||||
// releases-list ETag, shared across replicas.
|
||||
func (db *DB) ClaimGitHubDebSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (bool, string, error) {
|
||||
row := db.Pool.QueryRow(ctx, `
|
||||
INSERT INTO github_deb_sync_state AS s (remote_name, sync_lease_owner, sync_lease_expires)
|
||||
VALUES ($1, $2, now() + make_interval(secs => $4))
|
||||
ON CONFLICT (remote_name) DO UPDATE
|
||||
SET sync_lease_owner = $2,
|
||||
sync_lease_expires = now() + make_interval(secs => $4)
|
||||
WHERE (s.last_synced_at IS NULL OR s.last_synced_at < now() - make_interval(secs => $3))
|
||||
AND (s.sync_lease_expires IS NULL OR s.sync_lease_expires < now())
|
||||
RETURNING s.etag
|
||||
`, remoteName, owner, freshness.Seconds(), lease.Seconds())
|
||||
|
||||
var etag string
|
||||
if err := row.Scan(&etag); err != nil {
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return false, "", nil
|
||||
}
|
||||
return false, "", err
|
||||
}
|
||||
return true, etag, nil
|
||||
}
|
||||
|
||||
// ReleaseGitHubDebSyncLease records the completed scan and frees the lease. Only
|
||||
// the owning replica may release; last_synced_at advances so the next poll waits
|
||||
// a full freshness window, and etag is persisted for the next conditional request.
|
||||
func (db *DB) ReleaseGitHubDebSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error {
|
||||
_, err := db.Pool.Exec(ctx, `
|
||||
UPDATE github_deb_sync_state
|
||||
SET last_synced_at = $3, etag = $4, sync_lease_owner = '', sync_lease_expires = NULL
|
||||
WHERE remote_name = $1 AND sync_lease_owner = $2
|
||||
`, remoteName, owner, syncedAt, etag)
|
||||
return err
|
||||
}
|
||||
@@ -0,0 +1,90 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.unkin.net/unkin/artifactapi/pkg/models"
|
||||
)
|
||||
|
||||
func seedGitHubDebRemote(t *testing.T, name string) {
|
||||
t.Helper()
|
||||
if err := testDB.CreateRemote(ctx(), &models.Remote{
|
||||
Name: name, PackageType: models.PackageGitHubDeb, RepoType: models.RepoTypeRemote,
|
||||
BaseURL: "https://api.github.com/repos/acme/tools", ReleasesRemote: "github", MutableTTL: 3600,
|
||||
}); err != nil {
|
||||
t.Fatalf("seed github_deb remote: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestGitHubDebSyncLease exercises the real SQL: exactly one replica may hold the
|
||||
// lease, the recency window blocks a too-soon periodic re-claim, and a prime
|
||||
// (freshness 0) bypasses recency but still respects a live lease.
|
||||
func TestGitHubDebSyncLease(t *testing.T) {
|
||||
requireDB(t)
|
||||
name := "ghdeb-lease-" + time.Now().Format("150405.000000")
|
||||
seedGitHubDebRemote(t, name)
|
||||
|
||||
const lease = 15 * time.Minute
|
||||
freshness := time.Hour
|
||||
|
||||
claimed, etag, err := testDB.ClaimGitHubDebSyncLease(ctx(), name, "replica-1", freshness, lease)
|
||||
if err != nil || !claimed {
|
||||
t.Fatalf("replica-1 first claim: claimed=%v err=%v", claimed, err)
|
||||
}
|
||||
if etag != "" {
|
||||
t.Fatalf("initial etag should be empty, got %q", etag)
|
||||
}
|
||||
|
||||
claimed2, _, err := testDB.ClaimGitHubDebSyncLease(ctx(), name, "replica-2", freshness, lease)
|
||||
if err != nil {
|
||||
t.Fatalf("replica-2 claim err: %v", err)
|
||||
}
|
||||
if claimed2 {
|
||||
t.Fatal("replica-2 claimed while replica-1 holds the lease")
|
||||
}
|
||||
|
||||
if err := testDB.ReleaseGitHubDebSyncLease(ctx(), name, "replica-1", `"etag-1"`, time.Now()); err != nil {
|
||||
t.Fatalf("release: %v", err)
|
||||
}
|
||||
|
||||
claimed3, _, err := testDB.ClaimGitHubDebSyncLease(ctx(), name, "replica-2", freshness, lease)
|
||||
if err != nil {
|
||||
t.Fatalf("replica-2 recency claim err: %v", err)
|
||||
}
|
||||
if claimed3 {
|
||||
t.Fatal("periodic claim succeeded inside the freshness window")
|
||||
}
|
||||
|
||||
claimed4, etag4, err := testDB.ClaimGitHubDebSyncLease(ctx(), name, "replica-2", 0, lease)
|
||||
if err != nil || !claimed4 {
|
||||
t.Fatalf("prime claim: claimed=%v err=%v", claimed4, err)
|
||||
}
|
||||
if etag4 != `"etag-1"` {
|
||||
t.Fatalf("prime claim etag = %q, want persisted \"etag-1\"", etag4)
|
||||
}
|
||||
}
|
||||
|
||||
func TestListGitHubDebRemotes(t *testing.T) {
|
||||
requireDB(t)
|
||||
name := "ghdeb-list-" + time.Now().Format("150405.000000")
|
||||
seedGitHubDebRemote(t, name)
|
||||
seedRemote(t, "generic-"+time.Now().Format("150405.000000"))
|
||||
|
||||
remotes, err := testDB.ListGitHubDebRemotes(ctx())
|
||||
if err != nil {
|
||||
t.Fatalf("list: %v", err)
|
||||
}
|
||||
found := false
|
||||
for _, r := range remotes {
|
||||
if r.PackageType != models.PackageGitHubDeb {
|
||||
t.Fatalf("non-github_deb remote returned: %s (%s)", r.Name, r.PackageType)
|
||||
}
|
||||
if r.Name == name {
|
||||
found = true
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
t.Fatalf("seeded remote %q not returned", name)
|
||||
}
|
||||
}
|
||||
@@ -190,6 +190,14 @@ func (db *DB) migrate() error {
|
||||
sync_lease_expires TIMESTAMPTZ
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS github_deb_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,
|
||||
|
||||
@@ -194,7 +194,14 @@ func extractControl(deb []byte) (string, error) {
|
||||
return "", err
|
||||
}
|
||||
|
||||
tr := tar.NewReader(bytes.NewReader(tarBytes))
|
||||
return readControlParagraph(tarBytes)
|
||||
}
|
||||
|
||||
// readControlParagraph scans a decompressed control.tar and returns the raw
|
||||
// ./control paragraph. Shared by the local upload path (extractControl) and the
|
||||
// github_deb ranged-prefix parser.
|
||||
func readControlParagraph(controlTar []byte) (string, error) {
|
||||
tr := tar.NewReader(bytes.NewReader(controlTar))
|
||||
for {
|
||||
hdr, err := tr.Next()
|
||||
if err == io.EOF {
|
||||
@@ -370,8 +377,12 @@ func generatePackages(metas []provider.DebMetadata) []byte {
|
||||
b.WriteString("\n")
|
||||
fmt.Fprintf(&b, "Filename: %s\n", m.FilePath)
|
||||
fmt.Fprintf(&b, "Size: %d\n", m.Size)
|
||||
fmt.Fprintf(&b, "MD5sum: %s\n", m.MD5)
|
||||
fmt.Fprintf(&b, "SHA256: %s\n", m.SHA256)
|
||||
if m.MD5 != "" {
|
||||
fmt.Fprintf(&b, "MD5sum: %s\n", m.MD5)
|
||||
}
|
||||
if m.SHA256 != "" {
|
||||
fmt.Fprintf(&b, "SHA256: %s\n", m.SHA256)
|
||||
}
|
||||
b.WriteString("\n")
|
||||
}
|
||||
return b.Bytes()
|
||||
|
||||
@@ -0,0 +1,724 @@
|
||||
package deb
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"regexp"
|
||||
"strconv"
|
||||
"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_deb. 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. A .deb is an ar archive whose
|
||||
// control.tar member sits right after the tiny debian-binary member, so a small
|
||||
// front prefix reliably covers it.
|
||||
const (
|
||||
defaultHeaderRangeInitial = 32 << 10 // 32 KiB — covers control.tar of almost every .deb
|
||||
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 .deb assets, derives per-asset control metadata via a ranged prefix fetch
|
||||
// (never downloading whole packages), synthesizes a flat apt repository 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.PackageGitHubDeb }
|
||||
|
||||
func (p *GitHubProvider) Classify(path string) provider.Mutability {
|
||||
switch path {
|
||||
case "Packages", "Packages.gz", "Release", "InRelease", "Release.gpg":
|
||||
return provider.Mutable
|
||||
}
|
||||
return provider.Immutable
|
||||
}
|
||||
|
||||
func (p *GitHubProvider) ContentType(path string) string {
|
||||
switch {
|
||||
case strings.HasSuffix(path, ".deb"):
|
||||
return "application/vnd.debian.binary-package"
|
||||
case strings.HasSuffix(path, ".gz"):
|
||||
return "application/gzip"
|
||||
case path == "Packages" || path == "Release" || path == "InRelease":
|
||||
return "text/plain"
|
||||
}
|
||||
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_deb remote. It refreshes the
|
||||
// derived metadata (bounded by mutable_ttl), serves a synthesized flat apt repo
|
||||
// (Packages/Packages.gz/Release), 404s the signed index variants (the repo is
|
||||
// consumed via [trusted=yes]), and 302-redirects .deb 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)
|
||||
|
||||
// apt appends the flat-repo dist "./" verbatim, so it asks for "./Packages"
|
||||
// etc.; collapse the dot-segment before matching the synthesized index.
|
||||
path := normalizeIndexPath(reqPath)
|
||||
|
||||
switch path {
|
||||
case "Packages", "Packages.gz", "Release":
|
||||
p.serveIndex(w, r, remote, path, store)
|
||||
return true
|
||||
case "InRelease", "Release.gpg":
|
||||
// Unsigned flat repo: apt consumes it with [trusted=yes]. Signal absence
|
||||
// so apt falls back to the plain Release without waiting on a signature.
|
||||
http.Error(w, "not found", http.StatusNotFound)
|
||||
return true
|
||||
}
|
||||
|
||||
if strings.HasSuffix(path, ".deb") {
|
||||
if remote.ReleasesRemote == "" {
|
||||
http.Error(w, "github_deb remote has no releases_remote configured for downloads", http.StatusInternalServerError)
|
||||
return true
|
||||
}
|
||||
loc := strings.TrimRight(proxyBaseURL, "/") + "/api/v1/remote/" + remote.ReleasesRemote + "/" + strings.TrimLeft(path, "/")
|
||||
http.Redirect(w, r, loc, http.StatusFound)
|
||||
return true
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
func (p *GitHubProvider) serveIndex(w http.ResponseWriter, r *http.Request, remote models.Remote, path string, store provider.RemoteMetadataStore) {
|
||||
// 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.DebMetadataReader)
|
||||
if !ok {
|
||||
http.Error(w, "deb metadata not available", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
metas, err := reader.ListDebMetadataEntries(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
|
||||
}
|
||||
|
||||
switch path {
|
||||
case "Packages":
|
||||
w.Header().Set("Content-Type", "text/plain")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
w.Write(generatePackages(metas))
|
||||
case "Packages.gz":
|
||||
w.Header().Set("Content-Type", "application/gzip")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
w.Write(gzipBytes(generatePackages(metas)))
|
||||
case "Release":
|
||||
w.Header().Set("Content-Type", "text/plain")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
w.Write(generateRelease(metas))
|
||||
}
|
||||
}
|
||||
|
||||
// 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.DebMetadataReader)
|
||||
if !ok {
|
||||
return false
|
||||
}
|
||||
rows, err := reader.ListDebMetadataEntries(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_deb: 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) {
|
||||
releases, newEtag, notModified, err := p.fetchReleases(ctx, remote, etag)
|
||||
if err != nil {
|
||||
return etag, false, err
|
||||
}
|
||||
if notModified {
|
||||
return etag, false, nil
|
||||
}
|
||||
|
||||
reader, ok := store.(provider.DebMetadataReader)
|
||||
if !ok {
|
||||
return newEtag, false, errors.New("store does not support deb metadata reads")
|
||||
}
|
||||
existing, err := reader.ListDebMetadataEntries(ctx, remote.Name)
|
||||
if err != nil {
|
||||
return newEtag, false, err
|
||||
}
|
||||
existingByPath := make(map[string]provider.DebMetadata, 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), ".deb") {
|
||||
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
|
||||
}
|
||||
_ = store.DeleteDebMetadata(ctx, remote.Name, fp)
|
||||
}
|
||||
|
||||
meta, err := p.deriveAsset(ctx, remote, asset, fp)
|
||||
if err != nil {
|
||||
slog.Warn("github_deb: derive asset failed", "remote", remote.Name, "asset", asset.Name, "error", err)
|
||||
continue
|
||||
}
|
||||
if err := store.InsertDebMetadata(ctx, meta); err != nil {
|
||||
slog.Error("github_deb: insert metadata failed", "remote", remote.Name, "asset", asset.Name, "error", err)
|
||||
continue
|
||||
}
|
||||
slog.Info("github_deb: derived asset", "remote", remote.Name, "name", meta.Name, "version", meta.Version, "arch", meta.Architecture)
|
||||
}
|
||||
}
|
||||
|
||||
for fp := range existingByPath {
|
||||
if !seen[fp] {
|
||||
_ = store.DeleteDebMetadata(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.DebMetadata, error) {
|
||||
control, err := p.fetchControl(ctx, remote, asset.BrowserDownloadURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
fields := parseControlFields(control)
|
||||
|
||||
meta := &provider.DebMetadata{
|
||||
RepoName: remote.Name,
|
||||
FilePath: fp,
|
||||
Name: fields["Package"],
|
||||
Version: fields["Version"],
|
||||
Architecture: fields["Architecture"],
|
||||
Control: strings.TrimRight(control, "\n"),
|
||||
Size: asset.Size,
|
||||
}
|
||||
if meta.Name == "" {
|
||||
return nil, errors.New("control missing Package field")
|
||||
}
|
||||
|
||||
// The Packages SHA256 must be the sha256 of the whole .deb. Prefer GitHub's
|
||||
// asset digest so we never download the body; only when it is absent (or not
|
||||
// sha256) do we stream the asset once. MD5sum is left unset — apt verifies the
|
||||
// download against SHA256 alone under [trusted=yes].
|
||||
if h, ok := sha256FromDigest(asset.Digest); ok {
|
||||
meta.ContentHash = "sha256:" + h
|
||||
meta.SHA256 = h
|
||||
} else {
|
||||
h, err := p.computeSHA256(ctx, remote, asset.BrowserDownloadURL)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("compute sha256: %w", err)
|
||||
}
|
||||
meta.ContentHash = "sha256:" + h
|
||||
meta.SHA256 = h
|
||||
}
|
||||
|
||||
return meta, nil
|
||||
}
|
||||
|
||||
// fetchControl pulls only the front of the .deb with a ranged GET and extracts
|
||||
// the control paragraph from it. control.tar sits right after the tiny
|
||||
// debian-binary member, so a small prefix suffices; a prefix that truncates the
|
||||
// control member doubles the range and retries.
|
||||
func (p *GitHubProvider) fetchControl(ctx context.Context, remote models.Remote, downloadURL string) (string, error) {
|
||||
n := p.headerInitial
|
||||
for {
|
||||
body, full, err := p.rangeGet(ctx, remote, downloadURL, n)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
control, complete, perr := controlFromPrefix(body)
|
||||
if perr != nil {
|
||||
return "", fmt.Errorf("parse deb control: %w", perr)
|
||||
}
|
||||
if complete {
|
||||
return control, nil
|
||||
}
|
||||
if full || n >= p.headerMax {
|
||||
return "", fmt.Errorf("control.tar not found within %d bytes of %s", n, downloadURL)
|
||||
}
|
||||
n *= 2
|
||||
if n > p.headerMax {
|
||||
n = p.headerMax
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// controlFromPrefix parses the ar members present in a front prefix of a .deb.
|
||||
// It returns the ./control paragraph once control.tar.* is fully covered
|
||||
// (complete=true); a prefix too short to cover it returns complete=false so the
|
||||
// caller can widen the range. Later members (data.tar.*) are ignored.
|
||||
func controlFromPrefix(prefix []byte) (control string, complete bool, err error) {
|
||||
const magic = "!<arch>\n"
|
||||
if len(prefix) < len(magic) {
|
||||
return "", false, nil
|
||||
}
|
||||
if string(prefix[:len(magic)]) != magic {
|
||||
return "", false, errors.New("not an ar archive")
|
||||
}
|
||||
off := len(magic)
|
||||
for {
|
||||
if off+60 > len(prefix) {
|
||||
return "", false, nil
|
||||
}
|
||||
hdr := prefix[off : off+60]
|
||||
off += 60
|
||||
name := strings.TrimSuffix(strings.TrimRight(string(hdr[0:16]), " "), "/")
|
||||
size, err := strconv.ParseInt(strings.TrimSpace(string(hdr[48:58])), 10, 64)
|
||||
if err != nil {
|
||||
return "", false, fmt.Errorf("bad ar size for %q: %w", name, err)
|
||||
}
|
||||
if strings.HasPrefix(name, "control.tar") {
|
||||
if off+int(size) > len(prefix) {
|
||||
return "", false, nil
|
||||
}
|
||||
tarBytes, err := decompress(name, prefix[off:off+int(size)])
|
||||
if err != nil {
|
||||
return "", false, err
|
||||
}
|
||||
c, err := readControlParagraph(tarBytes)
|
||||
if err != nil {
|
||||
return "", false, err
|
||||
}
|
||||
return c, true, nil
|
||||
}
|
||||
if off+int(size) > len(prefix) {
|
||||
return "", false, nil
|
||||
}
|
||||
off += int(size)
|
||||
if size%2 == 1 {
|
||||
off++
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
func (p *GitHubProvider) computeSHA256(ctx context.Context, remote models.Remote, downloadURL string) (string, error) {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, downloadURL, nil)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
hdr, err := p.githubHeaders(ctx, remote, false)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
copyHeaders(req, hdr)
|
||||
|
||||
if err := p.limiterWait(ctx); err != nil {
|
||||
return "", err
|
||||
}
|
||||
resp, err := p.client.Do(req)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return "", fmt.Errorf("GET %s: status %d", downloadURL, resp.StatusCode)
|
||||
}
|
||||
|
||||
h := sha256.New()
|
||||
if _, err := io.Copy(h, resp.Body); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return hex.EncodeToString(h.Sum(nil)), 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
|
||||
// deb_metadata key and the Filename field in the Packages index, so a .deb
|
||||
// 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, "/")
|
||||
}
|
||||
|
||||
func sha256FromDigest(digest string) (string, bool) {
|
||||
if strings.HasPrefix(digest, "sha256:") {
|
||||
return strings.TrimPrefix(digest, "sha256:"), true
|
||||
}
|
||||
return "", false
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
@@ -0,0 +1,462 @@
|
||||
package deb
|
||||
|
||||
import (
|
||||
"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 + DebMetadataReader
|
||||
// keyed by file_path, mirroring the (repo_name, file_path) uniqueness of the
|
||||
// real deb_metadata table.
|
||||
type fakeStore struct {
|
||||
mu sync.Mutex
|
||||
rows map[string]provider.DebMetadata
|
||||
}
|
||||
|
||||
func newFakeStore() *fakeStore { return &fakeStore{rows: map[string]provider.DebMetadata{}} }
|
||||
|
||||
func (f *fakeStore) InsertDebMetadata(_ context.Context, m *provider.DebMetadata) 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) DeleteDebMetadata(_ context.Context, _, filePath string) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
delete(f.rows, filePath)
|
||||
return nil
|
||||
}
|
||||
|
||||
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) ListDebMetadataEntries(ctx context.Context, _ string) ([]provider.DebMetadata, error) {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
out := make([]provider.DebMetadata, 0, len(f.rows))
|
||||
for _, m := range f.rows {
|
||||
out = append(out, m)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// githubFixture serves the releases API and the .deb asset downloads (with Range
|
||||
// support) for a set of packages. digest controls whether the asset carries a
|
||||
// sha256 digest (no-download path) or not (compute path).
|
||||
type githubFixture struct {
|
||||
srv *httptest.Server
|
||||
debBytes 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{
|
||||
debBytes: map[string][]byte{},
|
||||
rangeHit: map[string]int{},
|
||||
fullHit: map[string]int{},
|
||||
}
|
||||
f.debBytes["demo_1.2-3_amd64.deb"] = testsupport.MinimalDeb("demo", "1.2-3", "amd64")
|
||||
|
||||
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.debBytes {
|
||||
a := map[string]any{
|
||||
"name": name,
|
||||
"size": len(f.debBytes[name]),
|
||||
"browser_download_url": f.srv.URL + "/acme/tools/releases/download/v1.2-3/" + name,
|
||||
}
|
||||
if withDigest {
|
||||
sum := sha256.Sum256(f.debBytes[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.debBytes[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-deb",
|
||||
PackageType: models.PackageGitHubDeb,
|
||||
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_amd64.deb"
|
||||
|
||||
func TestGitHubScanDerivesControlFromPrefixAndDigest(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.ListDebMetadataEntries(context.Background(), "acme-deb")
|
||||
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" || m.Architecture != "amd64" {
|
||||
t.Fatalf("bad control fields: %+v", m)
|
||||
}
|
||||
if m.FilePath != demoPath {
|
||||
t.Fatalf("FilePath = %q, want %q", m.FilePath, demoPath)
|
||||
}
|
||||
if int(m.Size) != len(fx.debBytes["demo_1.2-3_amd64.deb"]) {
|
||||
t.Fatalf("Size = %d, want %d", m.Size, len(fx.debBytes["demo_1.2-3_amd64.deb"]))
|
||||
}
|
||||
sum := sha256.Sum256(fx.debBytes["demo_1.2-3_amd64.deb"])
|
||||
if m.SHA256 != hex.EncodeToString(sum[:]) {
|
||||
t.Fatalf("SHA256 = %q, want digest", m.SHA256)
|
||||
}
|
||||
if m.ContentHash != "sha256:"+hex.EncodeToString(sum[:]) {
|
||||
t.Fatalf("ContentHash = %q", m.ContentHash)
|
||||
}
|
||||
if m.MD5 != "" {
|
||||
t.Fatalf("MD5 should be unset for metadata-only derive, got %q", m.MD5)
|
||||
}
|
||||
if fx.fullHit["demo_1.2-3_amd64.deb"] != 0 {
|
||||
t.Fatalf("expected no full download when digest present, got %d", fx.fullHit["demo_1.2-3_amd64.deb"])
|
||||
}
|
||||
if fx.rangeHit["demo_1.2-3_amd64.deb"] == 0 {
|
||||
t.Fatalf("expected ranged control fetch")
|
||||
}
|
||||
if !strings.Contains(m.Control, "Package: demo") {
|
||||
t.Fatalf("raw control not captured: %q", m.Control)
|
||||
}
|
||||
}
|
||||
|
||||
func TestGitHubChecksumComputedWhenDigestAbsent(t *testing.T) {
|
||||
fx := newGitHubFixture(t, false)
|
||||
p := newTestProvider()
|
||||
store := newFakeStore()
|
||||
|
||||
if err := p.scan(context.Background(), fx.remote(), store); err != nil {
|
||||
t.Fatalf("scan: %v", err)
|
||||
}
|
||||
metas, _ := store.ListDebMetadataEntries(context.Background(), "acme-deb")
|
||||
if len(metas) != 1 {
|
||||
t.Fatalf("want 1 row, got %d", len(metas))
|
||||
}
|
||||
sum := sha256.Sum256(fx.debBytes["demo_1.2-3_amd64.deb"])
|
||||
if metas[0].SHA256 != hex.EncodeToString(sum[:]) {
|
||||
t.Fatalf("computed checksum mismatch: %q", metas[0].SHA256)
|
||||
}
|
||||
if fx.fullHit["demo_1.2-3_amd64.deb"] == 0 {
|
||||
t.Fatalf("expected a full download to compute sha256 when digest absent")
|
||||
}
|
||||
}
|
||||
|
||||
func TestGitHubServeRemoteIndexAndRedirect(t *testing.T) {
|
||||
fx := newGitHubFixture(t, true)
|
||||
p := newTestProvider()
|
||||
store := newFakeStore()
|
||||
remote := fx.remote()
|
||||
const proxyBase = "https://artifactapi.example"
|
||||
|
||||
// Release is served and triggers the initial scan.
|
||||
rec := httptest.NewRecorder()
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-deb/Release", nil)
|
||||
if !p.ServeRemote(rec, req, remote, "Release", proxyBase, store) {
|
||||
t.Fatal("ServeRemote did not handle Release")
|
||||
}
|
||||
if rec.Code != 200 || !strings.Contains(rec.Body.String(), "Architectures:") {
|
||||
t.Fatalf("Release bad: code=%d body=%s", rec.Code, rec.Body.String())
|
||||
}
|
||||
if !strings.Contains(rec.Body.String(), "amd64") {
|
||||
t.Fatalf("Release missing arch: %s", rec.Body.String())
|
||||
}
|
||||
|
||||
// Packages carries the package with a Filename that is the github-relative
|
||||
// download path (so it resolves back to this remote and redirects).
|
||||
rec = httptest.NewRecorder()
|
||||
req = httptest.NewRequest(http.MethodGet, "/x", nil)
|
||||
if !p.ServeRemote(rec, req, remote, "Packages", proxyBase, store) {
|
||||
t.Fatal("ServeRemote did not handle Packages")
|
||||
}
|
||||
pkgs := rec.Body.String()
|
||||
if !strings.Contains(pkgs, "Package: demo") {
|
||||
t.Fatalf("Packages missing package: %s", pkgs)
|
||||
}
|
||||
if !strings.Contains(pkgs, "Filename: "+demoPath) {
|
||||
t.Fatalf("Packages missing/incorrect Filename: %s", pkgs)
|
||||
}
|
||||
if !strings.Contains(pkgs, "SHA256: ") {
|
||||
t.Fatalf("Packages missing SHA256: %s", pkgs)
|
||||
}
|
||||
if strings.Contains(pkgs, "MD5sum:") {
|
||||
t.Fatalf("Packages should omit empty MD5sum: %s", pkgs)
|
||||
}
|
||||
|
||||
// Packages.gz decompresses to the same content.
|
||||
rec = httptest.NewRecorder()
|
||||
req = httptest.NewRequest(http.MethodGet, "/x", nil)
|
||||
if !p.ServeRemote(rec, req, remote, "Packages.gz", proxyBase, store) {
|
||||
t.Fatal("ServeRemote did not handle Packages.gz")
|
||||
}
|
||||
gz, err := gzip.NewReader(rec.Body)
|
||||
if err != nil {
|
||||
t.Fatalf("gzip: %v", err)
|
||||
}
|
||||
unz, _ := io.ReadAll(gz)
|
||||
if !strings.Contains(string(unz), "Package: demo") {
|
||||
t.Fatalf("Packages.gz missing package: %s", unz)
|
||||
}
|
||||
|
||||
// InRelease/Release.gpg 404 (unsigned, consumed via [trusted=yes]).
|
||||
for _, sp := range []string{"InRelease", "Release.gpg"} {
|
||||
rec = httptest.NewRecorder()
|
||||
req = httptest.NewRequest(http.MethodGet, "/x", nil)
|
||||
if !p.ServeRemote(rec, req, remote, sp, proxyBase, store) {
|
||||
t.Fatalf("ServeRemote did not handle %s", sp)
|
||||
}
|
||||
if rec.Code != http.StatusNotFound {
|
||||
t.Fatalf("%s want 404, got %d", sp, rec.Code)
|
||||
}
|
||||
}
|
||||
|
||||
// A .deb request redirects to the backend releases_remote.
|
||||
rec = httptest.NewRecorder()
|
||||
req = httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-deb/"+demoPath, nil)
|
||||
if !p.ServeRemote(rec, req, remote, demoPath, proxyBase, store) {
|
||||
t.Fatal("ServeRemote did not handle .deb")
|
||||
}
|
||||
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", got, wantLoc)
|
||||
}
|
||||
}
|
||||
|
||||
// Real apt appends the flat-repo dist "./" verbatim, so the metadata-only remote
|
||||
// receives "./Packages" / "./Release"; ServeRemote must collapse the dot-segment
|
||||
// and synthesize the same index as the un-prefixed request.
|
||||
func TestGitHubServeRemoteAptDotSegment(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-deb/"+path, nil)
|
||||
if !p.ServeRemote(rec, req, remote, path, proxyBase, store) {
|
||||
t.Fatalf("ServeRemote did not handle %q", path)
|
||||
}
|
||||
return rec
|
||||
}
|
||||
|
||||
// Packages is deterministic: byte-identical to the un-prefixed request.
|
||||
plain, dotted := serve("Packages"), serve("./Packages")
|
||||
if plain.Code != 200 || dotted.Code != 200 {
|
||||
t.Fatalf("Packages: plain=%d dotted=%d, want 200/200", plain.Code, dotted.Code)
|
||||
}
|
||||
if !strings.Contains(dotted.Body.String(), "Package: demo") {
|
||||
t.Fatalf("./Packages missing synthesized body: %s", dotted.Body.String())
|
||||
}
|
||||
if !bytes.Equal(plain.Body.Bytes(), dotted.Body.Bytes()) {
|
||||
t.Error("./Packages body differs from Packages body")
|
||||
}
|
||||
|
||||
// Release carries a time.Now() Date: header; compare the rest.
|
||||
rPlain, rDotted := serve("Release"), serve("./Release")
|
||||
if rPlain.Code != 200 || rDotted.Code != 200 {
|
||||
t.Fatalf("Release: plain=%d dotted=%d, want 200/200", rPlain.Code, rDotted.Code)
|
||||
}
|
||||
if stripDate(rPlain.Body.String()) != stripDate(rDotted.Body.String()) {
|
||||
t.Error("./Release body differs from Release body (ignoring Date)")
|
||||
}
|
||||
}
|
||||
|
||||
// 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-deb/Packages", nil).WithContext(ctx)
|
||||
|
||||
if !p.ServeRemote(rec, req, remote, "Packages", "https://x", store) {
|
||||
t.Fatal("ServeRemote did not handle Packages")
|
||||
}
|
||||
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 !strings.Contains(rec.Body.String(), "Package: demo") {
|
||||
t.Fatalf("expected Packages served from cache, got %s", rec.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
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.ListDebMetadataEntries(context.Background(), "acme-deb"); len(rows) != 1 {
|
||||
t.Fatalf("want 1 row after first scan, got %d", len(rows))
|
||||
}
|
||||
|
||||
delete(fx.debBytes, "demo_1.2-3_amd64.deb")
|
||||
if err := p.scan(context.Background(), fx.remote(), store); err != nil {
|
||||
t.Fatalf("rescan: %v", err)
|
||||
}
|
||||
if rows, _ := store.ListDebMetadataEntries(context.Background(), "acme-deb"); len(rows) != 0 {
|
||||
t.Fatalf("want 0 rows after prune, got %d", len(rows))
|
||||
}
|
||||
}
|
||||
|
||||
func TestGitHubAssetPatternFilter(t *testing.T) {
|
||||
fx := newGitHubFixture(t, true)
|
||||
fx.debBytes["other_9_arm64.deb"] = testsupport.MinimalDeb("other", "9", "arm64")
|
||||
p := newTestProvider()
|
||||
store := newFakeStore()
|
||||
remote := fx.remote()
|
||||
remote.Patterns = []string{`^demo_.*_amd64\.deb$`}
|
||||
|
||||
if err := p.scan(context.Background(), remote, store); err != nil {
|
||||
t.Fatalf("scan: %v", err)
|
||||
}
|
||||
rows, _ := store.ListDebMetadataEntries(context.Background(), "acme-deb")
|
||||
if len(rows) != 1 || rows[0].Name != "demo" {
|
||||
t.Fatalf("pattern filter failed, rows=%+v", rows)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,238 @@
|
||||
package deb
|
||||
|
||||
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 deb 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
|
||||
ListGitHubDebRemotes(ctx context.Context) ([]models.Remote, error)
|
||||
ClaimGitHubDebSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (claimed bool, etag string, err error)
|
||||
ReleaseGitHubDebSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error
|
||||
}
|
||||
|
||||
// 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_deb
|
||||
// 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_deb 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_deb 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_deb 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_deb syncer stopped")
|
||||
return
|
||||
case <-ticker.C:
|
||||
s.schedule(ctx)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// schedule enqueues a periodic check for every github_deb 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.ListGitHubDebRemotes(ctx)
|
||||
if err != nil {
|
||||
slog.Error("github_deb 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.ClaimGitHubDebSyncLease(ctx, job.remote.Name, s.owner, freshness, syncLeaseDuration)
|
||||
if err != nil {
|
||||
slog.Error("github_deb 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_deb syncer: scan failed", "remote", job.remote.Name, "error", scanErr)
|
||||
}
|
||||
|
||||
relCtx, relCancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second)
|
||||
defer relCancel()
|
||||
if err := s.store.ReleaseGitHubDebSyncLease(relCtx, job.remote.Name, s.owner, releaseEtag, time.Now()); err != nil {
|
||||
slog.Warn("github_deb syncer: release lease", "remote", job.remote.Name, "error", err)
|
||||
}
|
||||
|
||||
if scanErr == nil && changed {
|
||||
slog.Info("github_deb 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[:])
|
||||
}
|
||||
@@ -0,0 +1,300 @@
|
||||
package deb
|
||||
|
||||
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) ListGitHubDebRemotes(_ context.Context) ([]models.Remote, error) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return append([]models.Remote(nil), f.remotes...), nil
|
||||
}
|
||||
|
||||
func (f *fakeSyncStore) ClaimGitHubDebSyncLease(_ 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) ReleaseGitHubDebSyncLease(_ 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_amd64.deb"]
|
||||
if priorRange == 0 {
|
||||
t.Fatal("first scan should have fetched the asset control")
|
||||
}
|
||||
|
||||
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_amd64.deb"]; got != priorRange {
|
||||
t.Fatalf("304 scan re-fetched asset control: %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_amd64.deb"]
|
||||
|
||||
fx.debBytes["other_9_arm64.deb"] = testsupport.MinimalDeb("other", "9", "arm64")
|
||||
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.ListDebMetadataEntries(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_amd64.deb"]; got != demoRange {
|
||||
t.Fatalf("already-cached asset was re-fetched: %d -> %d", demoRange, got)
|
||||
}
|
||||
if fx.rangeHit["other_9_arm64.deb"] == 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-deb", PackageType: models.PackageGitHubDeb, 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-deb", PackageType: models.PackageGitHubDeb, MutableTTL: 3600}
|
||||
|
||||
s.EnqueuePrime(remote)
|
||||
select {
|
||||
case job := <-s.jobs:
|
||||
if !job.prime || job.remote.Name != "acme-deb" {
|
||||
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.ClaimGitHubDebSyncLease(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.ListDebMetadataEntries(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-deb/Packages", nil)
|
||||
if !p.ServeRemote(rec, req, remote, "Packages", "https://x", store) {
|
||||
t.Fatal("ServeRemote did not handle Packages")
|
||||
}
|
||||
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-deb/Packages", nil)
|
||||
if !p.ServeRemote(rec, req, remote, "Packages", "https://x", store) {
|
||||
t.Fatal("ServeRemote did not handle Packages")
|
||||
}
|
||||
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.ListDebMetadataEntries(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)
|
||||
}
|
||||
}
|
||||
@@ -21,7 +21,7 @@ import (
|
||||
"git.unkin.net/unkin/artifactapi/internal/gc"
|
||||
"git.unkin.net/unkin/artifactapi/internal/githubauth"
|
||||
_ "git.unkin.net/unkin/artifactapi/internal/provider/alpine"
|
||||
_ "git.unkin.net/unkin/artifactapi/internal/provider/deb"
|
||||
"git.unkin.net/unkin/artifactapi/internal/provider/deb"
|
||||
_ "git.unkin.net/unkin/artifactapi/internal/provider/docker"
|
||||
_ "git.unkin.net/unkin/artifactapi/internal/provider/generic"
|
||||
_ "git.unkin.net/unkin/artifactapi/internal/provider/goproxy"
|
||||
@@ -35,6 +35,7 @@ import (
|
||||
"git.unkin.net/unkin/artifactapi/internal/storage"
|
||||
"git.unkin.net/unkin/artifactapi/internal/tfsign"
|
||||
"git.unkin.net/unkin/artifactapi/internal/virtual"
|
||||
"git.unkin.net/unkin/artifactapi/pkg/models"
|
||||
)
|
||||
|
||||
type Server struct {
|
||||
@@ -50,6 +51,7 @@ type Server struct {
|
||||
tfRegistry *tfregistry.Handler
|
||||
gc *gc.Collector
|
||||
syncer *rpm.Syncer
|
||||
debSyncer *deb.Syncer
|
||||
}
|
||||
|
||||
func New(cfg *config.Config, version string) (*Server, error) {
|
||||
@@ -97,6 +99,12 @@ func New(cfg *config.Config, version string) (*Server, error) {
|
||||
Workers: cfg.GitHubSyncWorkers,
|
||||
PollInterval: time.Duration(cfg.GitHubSyncPollInterval) * time.Second,
|
||||
})
|
||||
debSyncer := deb.NewSyncer(db, deb.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
|
||||
@@ -129,6 +137,7 @@ func New(cfg *config.Config, version string) (*Server, error) {
|
||||
tfRegistry: tfRegistry,
|
||||
gc: collector,
|
||||
syncer: syncer,
|
||||
debSyncer: debSyncer,
|
||||
}
|
||||
|
||||
s.router = s.routes()
|
||||
@@ -158,7 +167,10 @@ func (s *Server) routes() chi.Router {
|
||||
r.Mount("/api/v1", proxyHandler.Routes())
|
||||
r.Mount("/v2", proxyHandler.DockerV2Routes())
|
||||
|
||||
remotesHandler := v2.NewRemotesHandler(s.db, s.syncer)
|
||||
remotesHandler := v2.NewRemotesHandler(s.db, map[models.PackageType]v2.Primer{
|
||||
models.PackageGitHubRPM: s.syncer,
|
||||
models.PackageGitHubDeb: s.debSyncer,
|
||||
})
|
||||
virtualsHandler := v2.NewVirtualsHandler(s.db)
|
||||
healthHandler := v2.NewHealthHandler(s.db, s.cache, s.store)
|
||||
statsHandler := v2.NewStatsHandler(s.db)
|
||||
@@ -226,6 +238,7 @@ func (s *Server) newHTTPServer() *http.Server {
|
||||
func (s *Server) Run(ctx context.Context) error {
|
||||
go s.gc.Run(ctx)
|
||||
go s.syncer.Run(ctx)
|
||||
go s.debSyncer.Run(ctx)
|
||||
|
||||
httpServer := s.newHTTPServer()
|
||||
|
||||
@@ -247,6 +260,7 @@ func (s *Server) Run(ctx context.Context) error {
|
||||
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)
|
||||
|
||||
httpServer := s.newHTTPServer()
|
||||
|
||||
|
||||
@@ -17,6 +17,7 @@ const (
|
||||
PackageTerraform PackageType = "terraform"
|
||||
PackageGoProxy PackageType = "goproxy"
|
||||
PackageGitHubRPM PackageType = "github_rpm"
|
||||
PackageGitHubDeb PackageType = "github_deb"
|
||||
)
|
||||
|
||||
var validPackageTypes = map[PackageType]bool{
|
||||
@@ -32,6 +33,7 @@ var validPackageTypes = map[PackageType]bool{
|
||||
PackageTerraform: true,
|
||||
PackageGoProxy: true,
|
||||
PackageGitHubRPM: true,
|
||||
PackageGitHubDeb: true,
|
||||
}
|
||||
|
||||
func (p PackageType) Valid() bool {
|
||||
|
||||
Reference in New Issue
Block a user