From e24c35f53489f4a4319ae2a706ad1914aeab2499 Mon Sep 17 00:00:00 2001 From: Ben Vincent Date: Mon, 10 Aug 2026 21:31:24 +1000 Subject: [PATCH] feat: background syncer for github_rpm remotes (#108) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Why Builds on #107 (merged), which derives `github_rpm` RPM metadata lazily on the client request path, single-flighted per replica. Two problems remain: the derive still happens per replica, so across a multi-replica deployment the same releases are scanned and re-derived N times, multiplying GitHub queries; and a cold cache blocks the first request on a full derive. GitHub's rate limits are low (~60/hr unauthenticated, ~5000/hr authenticated), so this needs a single coordinated syncer with a shared rate limit and conditional requests. ## How - Add a single per-process background syncer (started at boot, cleanly stopped on shutdown) that owns a deduped/coalescing work queue, a worker pool, and one global token-bucket rate limiter (`golang.org/x/time/rate`) bound onto the github provider so every GitHub call (releases list + each ranged asset GET) acquires a token first. - Re-check each `github_rpm` remote for new/changed releases on its existing `mutable_ttl` cadence; derive only new/changed assets incrementally and prune assets that disappear upstream. Repodata is served from primed DB rows. - Prime metadata in the background on remote creation; the create call returns immediately. - Send the stored releases-list `ETag` as `If-None-Match`; a `304` derives nothing and is not counted against GitHub's rate limit, so an unchanged repo is nearly free. - Coordinate replicas through a `github_rpm_sync_state` row (`last_synced_at`, `etag`, `sync_lease_owner`, `sync_lease_expires`): a periodic scan runs only for the replica that atomically claims the lease, bounding total GitHub load to ~once per `mutable_ttl` regardless of replica count; the ETag is shared through the same row. - Keep the request path fast: serve current cache, enqueue a prime on an empty cache, and return a bounded wait then a retryable `503` rather than blocking on a cold derive. - Add `GITHUB_SYNC_RATE` / `GITHUB_SYNC_BURST` / `GITHUB_SYNC_WORKERS` / `GITHUB_SYNC_POLL_INTERVAL` config with conservative defaults (1 req/s, burst 5, 3 workers, 60s tick) and document the syncer in the README. ## Tests - Unit (httptest, Range/ETag-aware fixture): `304` releases response derives nothing; incremental derive fetches only the newly added asset; the shared limiter caps request rate; work-queue enqueues coalesce to one job; prime enqueues a job; a held lease stops a second replica from scanning; cold-start serves `503` while warm cache serves `200`. - DB integration (testcontainers postgres): the real lease SQL — one holder at a time, recency gate blocks a too-soon periodic re-claim, prime (freshness 0) bypasses recency but respects a live lease. - Docker e2e re-run: `dnf install dotvault` works; prime-on-create derives in the background at ~1 req/s (global limiter); `dnf makecache` served fast from the priming cache (no cold block); clean shutdown mid-scan, no panics. ## Notes - Reuses `mutable_ttl` as the check interval (no new per-remote field), per brief. Reviewed-on: https://git.unkin.net/unkin/artifactapi/pulls/108 Co-authored-by: Ben Vincent Co-committed-by: Ben Vincent --- README.md | 37 ++- go.mod | 1 + go.sum | 2 + internal/api/v2/errorpaths_test.go | 2 +- internal/api/v2/remotes.go | 20 +- internal/config/env.go | 34 +++ internal/database/github_sync.go | 73 ++++++ internal/database/github_sync_test.go | 95 ++++++++ internal/database/postgres.go | 8 + internal/provider/rpm/github.go | 163 ++++++++++++-- internal/provider/rpm/github_test.go | 26 ++- internal/provider/rpm/syncer.go | 256 +++++++++++++++++++++ internal/provider/rpm/syncer_test.go | 312 ++++++++++++++++++++++++++ internal/server/server.go | 14 +- 14 files changed, 1015 insertions(+), 28 deletions(-) create mode 100644 internal/database/github_sync.go create mode 100644 internal/database/github_sync_test.go create mode 100644 internal/provider/rpm/syncer.go create mode 100644 internal/provider/rpm/syncer_test.go diff --git a/README.md b/README.md index 86aed17..0f5edde 100644 --- a/README.md +++ b/README.md @@ -46,12 +46,43 @@ metadata comes from a **ranged GET of just the RPM header** (the header sits at the front of the file, so the whole package is never downloaded); the sha256 checksum comes from the GitHub asset `digest` when present, else a one-time lazy stream. Derived metadata is cached (keyed by asset) so repodata generation -is cheap on repeat, and refreshed no more often than `mutable_ttl`. +is served from primed DB rows, never a cold on-demand derive. Each package's `` points back at the remote, which **302-redirects** the download to the `releases_remote` — an existing generic `github.com` remote that streams the actual bytes. `dnf` follows the redirect transparently. +#### Background syncer + +A single process-wide **background syncer** keeps every `github_rpm` remote's +derived metadata current off the client request path: + +- **Prime on create.** Creating a `github_rpm` remote enqueues a background prime + scan, so its metadata is derived right away without blocking the create call. + The first `dnf` request is served from cache. If a request arrives before the + prime lands, it returns a retryable `503` (with `Retry-After`) rather than + serving an empty repo or blocking on a multi-minute derive. +- **Periodic re-check, driven by `mutable_ttl`.** Each remote is re-checked for + new or changed releases no more often than its `mutable_ttl`. New/changed + assets are derived incrementally; assets already cached are never re-fetched, + and assets that disappear upstream are pruned. +- **ETag / 304 conditional requests.** The releases-list `ETag` is stored per + remote and sent as `If-None-Match`; a `304 Not Modified` means nothing changed + and the syncer derives nothing. GitHub does not count `304` conditional + responses against the rate limit, so an unchanged repo is nearly free — this is + the main lever keeping GitHub traffic low. +- **Global rate limit.** Every GitHub call (releases list + each ranged asset + header GET) passes through a single token-bucket limiter **shared across all + remotes**, so GitHub is never hammered. Configure a token (`password`) on the + remote for the higher authenticated rate limit (~5000/hr vs ~60/hr + unauthenticated). +- **Multi-replica coordination.** State is shared through the database. Before a + periodic scan a replica must atomically claim a per-remote lease + (`github_rpm_sync_state`: `last_synced_at`, `etag`, `sync_lease_owner`, + `sync_lease_expires`); only the winner scans. This bounds total GitHub load to + ~once per `mutable_ttl` regardless of replica count, and the shared `etag` + lets any replica issue the conditional request. + ```hcl # Backend that serves the actual .rpm bytes from github.com. resource "artifactapi_remote_generic" "github" { @@ -242,6 +273,10 @@ S3 client supports MinIO, Ceph RGW, and AWS S3 (via minio-go). | `MINIO_BUCKET` | `artifacts` | S3 bucket | | `MINIO_SECURE` | `false` | Use HTTPS for S3 | | `MINIO_REGION` | | S3 region (AWS) | +| `GITHUB_SYNC_RATE` | `1` | `github_rpm` syncer global GitHub request rate (req/s), shared across all remotes. `1`/s = 3600/hr, under an authenticated token's ~5000/hr; unauthenticated (~60/hr) relies on ETag/304 | +| `GITHUB_SYNC_BURST` | `5` | Token-bucket burst for the shared limiter | +| `GITHUB_SYNC_WORKERS` | `3` | Concurrent `github_rpm` scan workers | +| `GITHUB_SYNC_POLL_INTERVAL` | `60` | Base scheduler tick in seconds; per-remote cadence is its `mutable_ttl`, enforced by the DB lease | ## Development diff --git a/go.mod b/go.mod index daba034..f20fac6 100644 --- a/go.mod +++ b/go.mod @@ -15,6 +15,7 @@ require ( github.com/testcontainers/testcontainers-go/modules/postgres v0.42.0 github.com/testcontainers/testcontainers-go/modules/redis v0.42.0 golang.org/x/crypto v0.51.0 + golang.org/x/time v0.15.0 gopkg.in/yaml.v3 v3.0.1 ) diff --git a/go.sum b/go.sum index d909271..d14c2bf 100644 --- a/go.sum +++ b/go.sum @@ -234,6 +234,8 @@ golang.org/x/term v0.43.0 h1:S4RLU2sB31O/NCl+zFN9Aru9A/Cq2aqKpTZJ6B+DwT4= golang.org/x/term v0.43.0/go.mod h1:lrhlHNdQJHO+1qVYiHfFKVuVioJIheAc3fBSMFYEIsk= golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= +golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= +golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= diff --git a/internal/api/v2/errorpaths_test.go b/internal/api/v2/errorpaths_test.go index abedb70..aecf463 100644 --- a/internal/api/v2/errorpaths_test.go +++ b/internal/api/v2/errorpaths_test.go @@ -57,7 +57,7 @@ func do(t *testing.T, h http.Handler, method, path, body string) int { } func TestRemotesErrorPaths(t *testing.T) { - h := NewRemotesHandler(closedDB(t)).Routes() + h := NewRemotesHandler(closedDB(t), nil).Routes() if c := do(t, h, "GET", "/", ""); c != 500 { t.Errorf("list with dead db = %d, want 500", c) } diff --git a/internal/api/v2/remotes.go b/internal/api/v2/remotes.go index 6ebd8ba..da2e735 100644 --- a/internal/api/v2/remotes.go +++ b/internal/api/v2/remotes.go @@ -11,12 +11,19 @@ import ( "git.unkin.net/unkin/artifactapi/pkg/models" ) -type RemotesHandler struct { - db *database.DB +// Primer enqueues a background metadata prime for a newly created remote so the +// create call never blocks on a derive. *rpm.Syncer satisfies it. +type Primer interface { + EnqueuePrime(remote models.Remote) } -func NewRemotesHandler(db *database.DB) *RemotesHandler { - return &RemotesHandler{db: db} +type RemotesHandler struct { + db *database.DB + primer Primer +} + +func NewRemotesHandler(db *database.DB, primer Primer) *RemotesHandler { + return &RemotesHandler{db: db, primer: primer} } func (h *RemotesHandler) Routes() chi.Router { @@ -77,6 +84,11 @@ 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) + } writeJSON(w, http.StatusCreated, remote) } diff --git a/internal/config/env.go b/internal/config/env.go index c2ec5bd..1addad3 100644 --- a/internal/config/env.go +++ b/internal/config/env.go @@ -32,6 +32,18 @@ type Config struct { TFSigningKeyPath string TFSigningKeyPassphrase string TFProviderProtocols string + + // github_rpm background syncer. The syncer keeps derived RPM metadata for + // every github_rpm remote fresh off the client request path, sharing a + // single global token-bucket limiter across all remotes so GitHub is never + // hammered. Defaults are conservative: 1 req/s (3600/hr) sits well under an + // authenticated token's 5000/hr. Unauthenticated remotes (60/hr) lean on + // ETag/304 — an unchanged repo costs nothing — so keep those repos small or + // configure a token. + GitHubSyncRatePerSec float64 + GitHubSyncBurst int + GitHubSyncWorkers int + GitHubSyncPollInterval int } func (c *Config) DatabaseDSN() string { @@ -49,6 +61,23 @@ func Load() (*Config, error) { s3Secure, _ := strconv.ParseBool(getenv("MINIO_SECURE", "false")) + syncRate, err := strconv.ParseFloat(getenv("GITHUB_SYNC_RATE", "1"), 64) + if err != nil { + return nil, fmt.Errorf("invalid GITHUB_SYNC_RATE: %w", err) + } + syncBurst, err := strconv.Atoi(getenv("GITHUB_SYNC_BURST", "5")) + if err != nil { + return nil, fmt.Errorf("invalid GITHUB_SYNC_BURST: %w", err) + } + syncWorkers, err := strconv.Atoi(getenv("GITHUB_SYNC_WORKERS", "3")) + if err != nil { + return nil, fmt.Errorf("invalid GITHUB_SYNC_WORKERS: %w", err) + } + syncPoll, err := strconv.Atoi(getenv("GITHUB_SYNC_POLL_INTERVAL", "60")) + if err != nil { + return nil, fmt.Errorf("invalid GITHUB_SYNC_POLL_INTERVAL: %w", err) + } + cfg := &Config{ ListenAddr: getenv("LISTEN_ADDR", ":8000"), @@ -71,6 +100,11 @@ func Load() (*Config, error) { TFSigningKeyPath: getenv("TF_SIGNING_KEY_PATH", ""), TFSigningKeyPassphrase: getenv("TF_SIGNING_KEY_PASSPHRASE", ""), TFProviderProtocols: getenv("TF_PROVIDER_PROTOCOLS", "5.0,6.0"), + + GitHubSyncRatePerSec: syncRate, + GitHubSyncBurst: syncBurst, + GitHubSyncWorkers: syncWorkers, + GitHubSyncPollInterval: syncPoll, } return cfg, nil diff --git a/internal/database/github_sync.go b/internal/database/github_sync.go new file mode 100644 index 0000000..75d8590 --- /dev/null +++ b/internal/database/github_sync.go @@ -0,0 +1,73 @@ +package database + +import ( + "context" + "errors" + "time" + + "github.com/jackc/pgx/v5" + + "git.unkin.net/unkin/artifactapi/pkg/models" +) + +// ListGitHubRPMRemotes returns every github_rpm remote so the syncer can sweep +// them on each poll tick. +func (db *DB) ListGitHubRPMRemotes(ctx context.Context) ([]models.Remote, error) { + rows, err := db.Pool.Query(ctx, `SELECT `+remoteCols+` FROM remotes WHERE package_type = $1 ORDER BY name`, models.PackageGitHubRPM) + 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() +} + +// ClaimGitHubSyncLease atomically claims the per-remote sync lease. It succeeds +// (claimed=true) only when the remote is due — never synced, or synced longer +// than freshness ago — and no live lease is held by another replica. This bounds +// total GitHub load to roughly one scan per freshness window regardless of how +// many replicas poll. The returned etag is the stored releases-list ETag, shared +// across replicas so a conditional request can short-circuit an unchanged repo. +// A zero freshness (used for prime scans) ignores the recency gate and claims +// whenever no live lease is held. +func (db *DB) ClaimGitHubSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (bool, string, error) { + row := db.Pool.QueryRow(ctx, ` + INSERT INTO github_rpm_sync_state AS s (remote_name, sync_lease_owner, sync_lease_expires) + 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 +} + +// ReleaseGitHubSyncLease records the completed scan and frees the lease. Only the +// owning replica may release; last_synced_at advances so the next poll waits a +// full freshness window, and etag is persisted for the next conditional request. +func (db *DB) ReleaseGitHubSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error { + _, err := db.Pool.Exec(ctx, ` + UPDATE github_rpm_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 +} diff --git a/internal/database/github_sync_test.go b/internal/database/github_sync_test.go new file mode 100644 index 0000000..7c6cc69 --- /dev/null +++ b/internal/database/github_sync_test.go @@ -0,0 +1,95 @@ +package database + +import ( + "testing" + "time" + + "git.unkin.net/unkin/artifactapi/pkg/models" +) + +func seedGitHubRPMRemote(t *testing.T, name string) { + t.Helper() + if err := testDB.CreateRemote(ctx(), &models.Remote{ + Name: name, PackageType: models.PackageGitHubRPM, RepoType: models.RepoTypeRemote, + BaseURL: "https://api.github.com/repos/acme/tools", ReleasesRemote: "github", MutableTTL: 3600, + }); err != nil { + t.Fatalf("seed github_rpm remote: %v", err) + } +} + +// TestGitHubSyncLease 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 TestGitHubSyncLease(t *testing.T) { + requireDB(t) + name := "gh-lease-" + time.Now().Format("150405.000000") + seedGitHubRPMRemote(t, name) + + const lease = 15 * time.Minute + freshness := time.Hour + + // First claim on a never-synced remote wins; etag starts empty. + claimed, etag, err := testDB.ClaimGitHubSyncLease(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) + } + + // A second replica cannot claim while the lease is held. + claimed2, _, err := testDB.ClaimGitHubSyncLease(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") + } + + // Replica 1 finishes: record the sync and persist an etag. + if err := testDB.ReleaseGitHubSyncLease(ctx(), name, "replica-1", `"etag-1"`, time.Now()); err != nil { + t.Fatalf("release: %v", err) + } + + // A periodic re-claim inside the freshness window is blocked by recency. + claimed3, _, err := testDB.ClaimGitHubSyncLease(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") + } + + // A prime (freshness 0) bypasses recency and reads the persisted etag. + claimed4, etag4, err := testDB.ClaimGitHubSyncLease(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 TestListGitHubRPMRemotes(t *testing.T) { + requireDB(t) + name := "gh-list-" + time.Now().Format("150405.000000") + seedGitHubRPMRemote(t, name) + seedRemote(t, "generic-"+time.Now().Format("150405.000000")) + + remotes, err := testDB.ListGitHubRPMRemotes(ctx()) + if err != nil { + t.Fatalf("list: %v", err) + } + found := false + for _, r := range remotes { + if r.PackageType != models.PackageGitHubRPM { + t.Fatalf("non-github_rpm remote returned: %s (%s)", r.Name, r.PackageType) + } + if r.Name == name { + found = true + } + } + if !found { + t.Fatalf("seeded remote %q not returned", name) + } +} diff --git a/internal/database/postgres.go b/internal/database/postgres.go index baf2303..cbfb38c 100644 --- a/internal/database/postgres.go +++ b/internal/database/postgres.go @@ -164,6 +164,14 @@ func (db *DB) migrate() error { ALTER TABLE rpm_metadata ADD COLUMN IF NOT EXISTS conflicts JSONB DEFAULT '[]'; ALTER TABLE rpm_metadata ADD COLUMN IF NOT EXISTS obsoletes JSONB DEFAULT '[]'; + CREATE TABLE IF NOT EXISTS github_rpm_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, diff --git a/internal/provider/rpm/github.go b/internal/provider/rpm/github.go index 3788072..03343dd 100644 --- a/internal/provider/rpm/github.go +++ b/internal/provider/rpm/github.go @@ -18,13 +18,19 @@ import ( "time" rpmlib "github.com/cavaliergopher/rpm" + "golang.org/x/time/rate" "git.unkin.net/unkin/artifactapi/internal/provider" "git.unkin.net/unkin/artifactapi/pkg/models" ) +// gitHubProvider is the process-wide singleton. 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(newGitHubProvider()) + provider.Register(gitHubProvider) } // Tuning knobs for the no-precache header fetch. Fields (not consts) so tests @@ -39,6 +45,12 @@ const ( defaultScanTimeout = 10 * time.Minute // defaultServeTimeout bounds a repodata DB read served on a detached context. defaultServeTimeout = 30 * time.Second + + // defaultColdWait bounds how long a repodata request blocks waiting for a + // just-enqueued prime to populate an empty cache before returning a + // retryable 503. Kept short so a client never hangs on a rate-limited derive + // of a large repo; small repos usually prime within this window. + defaultColdWait = 8 * time.Second ) // GitHubProvider is a metadata-only remote: it scans a GitHub repo's releases @@ -54,6 +66,15 @@ type GitHubProvider struct { pageCap int scanTimeout time.Duration serveTimeout time.Duration + coldWait time.Duration + + // limiter, when set by the Syncer, gates every GitHub HTTP call (releases + // list + each ranged asset fetch) through a single process-wide token bucket. + // nil means unlimited (direct provider use / unit tests). + limiter *rate.Limiter + // syncer, when set, routes freshness refresh and cold-start priming through + // the shared background work queue instead of an inline per-replica scan. + syncer *Syncer mu sync.Mutex scanning map[string]bool @@ -68,11 +89,21 @@ func newGitHubProvider() *GitHubProvider { pageCap: defaultReleasePageCap, scanTimeout: defaultScanTimeout, serveTimeout: defaultServeTimeout, + coldWait: defaultColdWait, scanning: map[string]bool{}, lastScan: map[string]time.Time{}, } } +// limiterWait blocks until the shared rate limiter grants a token, or returns +// the context error if it is canceled first. A nil limiter is a no-op. +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.PackageGitHubRPM } // Classify/ContentType/UpstreamURL/RewriteResponse/AuthHeaders satisfy the @@ -114,7 +145,7 @@ func (p *GitHubProvider) AuthHeaders(_ context.Context, remote models.Remote) (h // 302-redirects .rpm downloads to the backend releases_remote. Returns false // only for paths it does not own, letting the normal proxy path take over. func (p *GitHubProvider) ServeRemote(w http.ResponseWriter, r *http.Request, remote models.Remote, path, proxyBaseURL string, store provider.RemoteMetadataStore) bool { - p.refresh(remote, store) + p.onRequest(remote, store) if strings.HasPrefix(path, "repodata/") { // Serve repodata on a context detached from the inbound request: a @@ -124,6 +155,15 @@ func (p *GitHubProvider) ServeRemote(w http.ResponseWriter, r *http.Request, rem defer cancel() sr := r.WithContext(sctx) + // Cold start: with the syncer wired, an empty cache means the prime has + // not landed yet. Enqueue it and wait briefly rather than serving empty + // repodata; if it still has not primed, return a retryable 503. + 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 true + } + tail := strings.TrimPrefix(path, "repodata/") lp := &Provider{} switch { @@ -154,6 +194,56 @@ func (p *GitHubProvider) ServeRemote(w http.ResponseWriter, r *http.Request, rem return false } +// onRequest keeps a remote's derived metadata fresh off the request path. With +// the background syncer wired it enqueues a deduped, rate-limited, lease-gated +// refresh and returns immediately; the request always serves the current cache. +// Without a syncer (direct provider use / unit tests) it falls back to the +// legacy inline single-flight scan. +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 metadata +// row. On an empty cache it enqueues a prime and polls briefly for it to land, +// so the very first client after a remote is created gets real repodata instead +// of an empty index or a blocking multi-minute derive. Returns false if the +// cache is still empty after the bounded wait. +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 { + rows, err := store.ListRPMMetadataEntries(ctx, name) + if err != nil { + // Treat a failed read as "not empty" so a transient DB error becomes a + // normal serve attempt (which reports its own error) rather than a 503. + return false + } + return len(rows) == 0 +} + // refresh brings the derived metadata up to date without coupling the scan to // the inbound request. When the cache is stale it single-flights a scan: if the // cache already holds rows the scan runs in the background and the caller serves @@ -212,15 +302,31 @@ func (p *GitHubProvider) runScan(remote models.Remote, store provider.RemoteMeta p.mu.Unlock() } +// scan runs a full unconditional derive. Retained for the legacy inline refresh +// path and existing tests; the syncer uses scanWithState to pass and receive the +// releases-list ETag. func (p *GitHubProvider) scan(ctx context.Context, remote models.Remote, store provider.RemoteMetadataStore) error { - releases, err := p.fetchReleases(ctx, remote) + _, _, 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, so it returns +// (etag, changed=false) without a single asset fetch. 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 err + return etag, false, err + } + if notModified { + return etag, false, nil } existing, err := store.ListRPMMetadataEntries(ctx, remote.Name) if err != nil { - return err + return newEtag, false, err } existingByPath := make(map[string]provider.RPMMetadata, len(existing)) for _, m := range existing { @@ -229,7 +335,7 @@ func (p *GitHubProvider) scan(ctx context.Context, remote models.Remote, store p allow, err := compilePatterns(remote.Patterns) if err != nil { - return err + return newEtag, false, err } seen := map[string]bool{} @@ -277,7 +383,7 @@ func (p *GitHubProvider) scan(ctx context.Context, remote models.Remote, store p _ = store.DeleteRPMMetadata(ctx, remote.Name, fp) } } - return nil + return newEtag, true, nil } type ghRelease struct { @@ -293,32 +399,53 @@ type ghAsset struct { Digest string `json:"digest"` } -func (p *GitHubProvider) fetchReleases(ctx context.Context, remote models.Remote) ([]ghRelease, error) { +// fetchReleases lists a repo's releases. It sends the prior ETag as +// If-None-Match on page 1 (the newest releases, where a new one first appears): +// a 304 there means the repo is unchanged, so it returns notModified without +// paging further — GitHub does not count 304 conditional responses against the +// rate limit, making an unchanged repo nearly free. On a 200 it captures the +// page-1 ETag and pages through the rest normally. 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" - var all []ghRelease 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, err + return nil, "", false, err } copyHeaders(req, githubHeaders(remote, true)) + 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, err + 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, err + return nil, "", false, err } if resp.StatusCode != http.StatusOK { - return nil, fmt.Errorf("github releases API %s: status %d", u, resp.StatusCode) + 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, fmt.Errorf("decode releases: %w", err) + return nil, "", false, fmt.Errorf("decode releases: %w", err) } if len(releases) == 0 { break @@ -328,7 +455,7 @@ func (p *GitHubProvider) fetchReleases(ctx context.Context, remote models.Remote break } } - return all, nil + return all, newEtag, false, nil } func (p *GitHubProvider) deriveAsset(ctx context.Context, remote models.Remote, asset ghAsset, fp string) (*provider.RPMMetadata, error) { @@ -447,6 +574,9 @@ func (p *GitHubProvider) rangeGet(ctx context.Context, remote models.Remote, dow copyHeaders(req, githubHeaders(remote, false)) 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 @@ -471,6 +601,9 @@ func (p *GitHubProvider) computeSHA256(ctx context.Context, remote models.Remote } copyHeaders(req, githubHeaders(remote, false)) + if err := p.limiterWait(ctx); err != nil { + return "", err + } resp, err := p.client.Do(req) if err != nil { return "", err diff --git a/internal/provider/rpm/github_test.go b/internal/provider/rpm/github_test.go index 3a30278..f2eba4e 100644 --- a/internal/provider/rpm/github_test.go +++ b/internal/provider/rpm/github_test.go @@ -66,11 +66,14 @@ func (f *fakeStore) ListRPMMetadataEntries(ctx context.Context, _ string) ([]pro // 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 - rpmBytes map[string][]byte // asset filename -> bytes - rangeHit map[string]int // asset filename -> number of ranged GETs - fullHit map[string]int // asset filename -> number of full GETs - mu sync.Mutex + srv *httptest.Server + rpmBytes map[string][]byte // asset filename -> bytes + rangeHit map[string]int // asset filename -> number of ranged GETs + fullHit map[string]int // asset filename -> number of full GETs + etag string // when set, served as ETag; matching If-None-Match yields 304 + releasesHit int // total releases-list requests (200 + 304) + notModHit int // releases-list requests answered 304 + mu sync.Mutex } func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture { @@ -89,6 +92,19 @@ func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture { w.Write([]byte("[]")) return } + f.mu.Lock() + f.releasesHit++ + 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.rpmBytes { a := map[string]any{ diff --git a/internal/provider/rpm/syncer.go b/internal/provider/rpm/syncer.go new file mode 100644 index 0000000..f6f5724 --- /dev/null +++ b/internal/provider/rpm/syncer.go @@ -0,0 +1,256 @@ +package rpm + +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 is how long a claimed sync lease is held before it is + // considered abandoned. It comfortably exceeds a scan's own timeout so a live + // scan never loses its lease, while a crashed replica's lease still expires. + syncLeaseDuration = 15 * time.Minute + // defaultSyncFreshness is the periodic re-check interval used when a remote's + // mutable_ttl is unset. + defaultSyncFreshness = 5 * time.Minute + // jobQueueDepth bounds the pending work queue; enqueues past it are dropped + // (a later poll re-enqueues), never blocking the caller. + jobQueueDepth = 256 +) + +// SyncStore is the persistence surface the 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 + ListGitHubRPMRemotes(ctx context.Context) ([]models.Remote, error) + ClaimGitHubSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (claimed bool, etag string, err error) + ReleaseGitHubSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error +} + +// SyncConfig tunes the shared syncer. Zero values fall back to safe defaults. +type SyncConfig struct { + RatePerSec float64 // global GitHub request rate (req/s) + Burst int // token-bucket burst + Workers int // concurrent scan workers + PollInterval time.Duration // base scheduler tick; per-remote cadence is mutable_ttl +} + +type syncJob struct { + remote models.Remote + prime bool +} + +// Syncer is the single per-process background worker that keeps every +// github_rpm 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 provider so every GitHub call it makes +// passes through the same bucket. 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 // remotes queued or in-flight, for dedup/coalescing +} + +// NewSyncer builds the syncer bound to the process-wide github 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{}, + } + // Bind the shared limiter and back-reference so the request path routes + // through this syncer and every derive HTTP call is rate limited. + 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_rpm 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) // sweep at boot so existing remotes are checked immediately + for { + select { + case <-ctx.Done(): + wg.Wait() + slog.Info("github_rpm syncer stopped") + return + case <-ticker.C: + s.schedule(ctx) + } + } +} + +// schedule enqueues a periodic check for every github_rpm remote. The DB lease +// (claimed in the worker) enforces the per-remote mutable_ttl cadence and cross +// replica coordination, so enqueuing every tick is cheap: a not-yet-due remote +// simply fails to claim and is skipped. +func (s *Syncer) schedule(ctx context.Context) { + remotes, err := s.store.ListGitHubRPMRemotes(ctx) + if err != nil { + slog.Error("github_rpm 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 so its metadata is derived without blocking the create call. +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: a full queue drops the +// job (a later poll re-enqueues it) after clearing the dedup slot. +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. The +// lease bounds total GitHub load to one scan per freshness window across all +// replicas; 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 // prime ignores the recency gate but still respects a live lease + } + + claimed, etag, err := s.store.ClaimGitHubSyncLease(ctx, job.remote.Name, s.owner, freshness, syncLeaseDuration) + if err != nil { + slog.Error("github_rpm 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_rpm syncer: scan failed", "remote", job.remote.Name, "error", scanErr) + } + + // Release on a detached context so a clean shutdown mid-scan still frees the + // lease and advances last_synced_at (otherwise it simply expires). + relCtx, relCancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second) + defer relCancel() + if err := s.store.ReleaseGitHubSyncLease(relCtx, job.remote.Name, s.owner, releaseEtag, time.Now()); err != nil { + slog.Warn("github_rpm syncer: release lease", "remote", job.remote.Name, "error", err) + } + + if scanErr == nil && changed { + slog.Info("github_rpm 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[:]) +} diff --git a/internal/provider/rpm/syncer_test.go b/internal/provider/rpm/syncer_test.go new file mode 100644 index 0000000..727e0ec --- /dev/null +++ b/internal/provider/rpm/syncer_test.go @@ -0,0 +1,312 @@ +package rpm + +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) ListGitHubRPMRemotes(_ context.Context) ([]models.Remote, error) { + f.mu.Lock() + defer f.mu.Unlock() + return append([]models.Remote(nil), f.remotes...), nil +} + +func (f *fakeSyncStore) ClaimGitHubSyncLease(_ 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) ReleaseGitHubSyncLease(_ 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 header GETs 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.x86_64.rpm"] + if priorRange == 0 { + t.Fatal("first scan should have fetched the asset header") + } + + 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.x86_64.rpm"]; got != priorRange { + t.Fatalf("304 scan re-fetched asset header: %d -> %d", priorRange, got) + } +} + +// (b) On a real change, only the newly added asset is derived; assets already +// cached are never re-fetched. +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.x86_64.rpm"] + + // Add a new asset and bump the ETag so the conditional request returns 200. + fx.rpmBytes["other-9-9.aarch64.rpm"] = testsupport.MinimalRPM("other", "9", "9", "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.ListRPMMetadataEntries(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.x86_64.rpm"]; got != demoRange { + t.Fatalf("already-cached asset was re-fetched: %d -> %d", demoRange, got) + } + if fx.rangeHit["other-9-9.aarch64.rpm"] == 0 { + t.Fatal("newly added asset was not derived") + } +} + +// (c) The shared limiter caps the request rate: three gated releases calls at +// one token per 120ms cannot complete faster than ~2 gaps. +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-rpm", PackageType: models.PackageGitHubRPM, 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-rpm", PackageType: models.PackageGitHubRPM, MutableTTL: 3600} + + s.EnqueuePrime(remote) + select { + case job := <-s.jobs: + if !job.prime || job.remote.Name != "acme-rpm" { + 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: with the lease owned +// by another replica, process claims nothing and makes zero GitHub calls. +func TestSyncerLeasePreventsSecondReplica(t *testing.T) { + fx := newGitHubFixture(t, true) + fx.etag = `"v1"` + store := newFakeSyncStore() + p := newTestProvider() + s := newSyncer(store, p, testSyncConfig()) + remote := fx.remote() + + // Replica 1 holds the lease. + claimed, _, err := store.ClaimGitHubSyncLease(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) + } + + // Replica 2 (this syncer) tries to process the same remote; it must skip. + 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.ListRPMMetadataEntries(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, a repodata request enqueues a +// prime and, when it has not landed within the bounded cold wait, returns a +// retryable 503 rather than serving empty repodata (and without regressing the +// detached-context serve). +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-rpm/repodata/repomd.xml", nil) + if !p.ServeRemote(rec, req, remote, "repodata/repomd.xml", "https://x", store) { + t.Fatal("ServeRemote did not handle repomd.xml") + } + 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") + } + // The prime was enqueued. + 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 repodata 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-rpm/repodata/repomd.xml", nil) + if !p.ServeRemote(rec, req, remote, "repodata/repomd.xml", "https://x", store) { + t.Fatal("ServeRemote did not handle repomd.xml") + } + 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, deriving metadata, +// while 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 + + // Prime derives despite no prior sync. + s.process(context.Background(), syncJob{remote: remote, prime: true}) + if rows, _ := store.ListRPMMetadataEntries(context.Background(), remote.Name); len(rows) != 1 { + t.Fatalf("prime did not derive: %d rows", len(rows)) + } + releasesAfterPrime := fx.releasesHit + + // A periodic job immediately after is gated by mutable_ttl recency: no new + // releases call. + 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) + } +} diff --git a/internal/server/server.go b/internal/server/server.go index 91114ec..4037d5a 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -27,7 +27,7 @@ import ( _ "git.unkin.net/unkin/artifactapi/internal/provider/npm" _ "git.unkin.net/unkin/artifactapi/internal/provider/puppet" _ "git.unkin.net/unkin/artifactapi/internal/provider/pypi" - _ "git.unkin.net/unkin/artifactapi/internal/provider/rpm" + "git.unkin.net/unkin/artifactapi/internal/provider/rpm" _ "git.unkin.net/unkin/artifactapi/internal/provider/terraform" "git.unkin.net/unkin/artifactapi/internal/proxy" "git.unkin.net/unkin/artifactapi/internal/storage" @@ -47,6 +47,7 @@ type Server struct { localHandler *v2.LocalHandler tfRegistry *tfregistry.Handler gc *gc.Collector + syncer *rpm.Syncer } func New(cfg *config.Config, version string) (*Server, error) { @@ -69,6 +70,12 @@ func New(cfg *config.Config, version string) (*Server, error) { localHandler := v2.NewLocalHandler(db, s3) virtEngine := virtual.NewEngine(db, engine) collector := gc.New(db, s3, 1*time.Hour) + syncer := rpm.NewSyncer(db, rpm.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 @@ -100,6 +107,7 @@ func New(cfg *config.Config, version string) (*Server, error) { localHandler: localHandler, tfRegistry: tfRegistry, gc: collector, + syncer: syncer, } s.router = s.routes() @@ -129,7 +137,7 @@ func (s *Server) routes() chi.Router { r.Mount("/api/v1", proxyHandler.Routes()) r.Mount("/v2", proxyHandler.DockerV2Routes()) - remotesHandler := v2.NewRemotesHandler(s.db) + remotesHandler := v2.NewRemotesHandler(s.db, s.syncer) virtualsHandler := v2.NewVirtualsHandler(s.db) healthHandler := v2.NewHealthHandler(s.db, s.cache, s.store) statsHandler := v2.NewStatsHandler(s.db) @@ -196,6 +204,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) httpServer := s.newHTTPServer() @@ -216,6 +225,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) httpServer := s.newHTTPServer()