diff --git a/internal/database/alpine_github_sync.go b/internal/database/alpine_github_sync.go index 3bf0851..1378e16 100644 --- a/internal/database/alpine_github_sync.go +++ b/internal/database/alpine_github_sync.go @@ -2,11 +2,9 @@ package database import ( "context" - "errors" "time" - "github.com/jackc/pgx/v5" - + "git.unkin.net/unkin/artifactapi/internal/provider" "git.unkin.net/unkin/artifactapi/pkg/models" ) @@ -30,42 +28,14 @@ func (db *DB) ListGitHubAlpineRemotes(ctx context.Context) ([]models.Remote, err return remotes, rows.Err() } -// ClaimGitHubAlpineSyncLease atomically claims the per-remote sync lease. It -// succeeds only when the remote is due (never synced, or synced longer than -// freshness ago) and no live lease is held by another replica. A zero freshness -// (prime scans) ignores the recency gate. The returned etag is the stored -// releases-list ETag, shared across replicas. +// ClaimGitHubAlpineSyncLease atomically claims the per-remote github_alpine sync +// lease. See claimSyncLease. func (db *DB) ClaimGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (bool, string, error) { - row := db.Pool.QueryRow(ctx, ` - INSERT INTO github_alpine_sync_state AS s (remote_name, sync_lease_owner, sync_lease_expires) - VALUES ($1, $2, now() + make_interval(secs => $4)) - ON CONFLICT (remote_name) DO UPDATE - SET sync_lease_owner = $2, - sync_lease_expires = now() + make_interval(secs => $4) - WHERE (s.last_synced_at IS NULL OR s.last_synced_at < now() - make_interval(secs => $3)) - AND (s.sync_lease_expires IS NULL OR s.sync_lease_expires < now()) - RETURNING s.etag - `, remoteName, owner, freshness.Seconds(), lease.Seconds()) - - var etag string - if err := row.Scan(&etag); err != nil { - if errors.Is(err, pgx.ErrNoRows) { - return false, "", nil - } - return false, "", err - } - return true, etag, nil + return db.claimSyncLease(ctx, "github_alpine_sync_state", remoteName, owner, freshness, lease) } -// ReleaseGitHubAlpineSyncLease records the completed scan and frees the lease. -// Only the owning replica may release; last_synced_at advances so the next poll -// waits a full freshness window, and etag is persisted for the next conditional -// request. -func (db *DB) ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error { - _, err := db.Pool.Exec(ctx, ` - UPDATE github_alpine_sync_state - SET last_synced_at = $3, etag = $4, sync_lease_owner = '', sync_lease_expires = NULL - WHERE remote_name = $1 AND sync_lease_owner = $2 - `, remoteName, owner, syncedAt, etag) - return err +// ReleaseGitHubAlpineSyncLease records a github_alpine scan outcome and frees the +// lease. See releaseSyncLease. +func (db *DB) ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error { + return db.releaseSyncLease(ctx, "github_alpine_sync_state", remoteName, owner, res) } diff --git a/internal/database/deb_github_sync.go b/internal/database/deb_github_sync.go index 5525e52..3a727fe 100644 --- a/internal/database/deb_github_sync.go +++ b/internal/database/deb_github_sync.go @@ -2,11 +2,9 @@ package database import ( "context" - "errors" "time" - "github.com/jackc/pgx/v5" - + "git.unkin.net/unkin/artifactapi/internal/provider" "git.unkin.net/unkin/artifactapi/pkg/models" ) @@ -30,41 +28,14 @@ func (db *DB) ListGitHubDebRemotes(ctx context.Context) ([]models.Remote, error) 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. +// ClaimGitHubDebSyncLease atomically claims the per-remote github_deb sync +// lease. See claimSyncLease. 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 + return db.claimSyncLease(ctx, "github_deb_sync_state", remoteName, owner, freshness, lease) } -// 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 +// ReleaseGitHubDebSyncLease records a github_deb scan outcome and frees the +// lease. See releaseSyncLease. +func (db *DB) ReleaseGitHubDebSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error { + return db.releaseSyncLease(ctx, "github_deb_sync_state", remoteName, owner, res) } diff --git a/internal/database/deb_github_sync_test.go b/internal/database/deb_github_sync_test.go index 7c73aa1..1049d3a 100644 --- a/internal/database/deb_github_sync_test.go +++ b/internal/database/deb_github_sync_test.go @@ -4,6 +4,7 @@ import ( "testing" "time" + "git.unkin.net/unkin/artifactapi/internal/provider" "git.unkin.net/unkin/artifactapi/pkg/models" ) @@ -44,7 +45,7 @@ func TestGitHubDebSyncLease(t *testing.T) { 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 { + if err := testDB.ReleaseGitHubDebSyncLease(ctx(), name, "replica-1", provider.SyncResult{Etag: `"etag-1"`}); err != nil { t.Fatalf("release: %v", err) } diff --git a/internal/database/github_sync.go b/internal/database/github_sync.go index 75d8590..d3498fc 100644 --- a/internal/database/github_sync.go +++ b/internal/database/github_sync.go @@ -7,6 +7,8 @@ import ( "github.com/jackc/pgx/v5" + "git.unkin.net/unkin/artifactapi/internal/provider" + "git.unkin.net/unkin/artifactapi/pkg/models" ) @@ -30,22 +32,35 @@ func (db *DB) ListGitHubRPMRemotes(ctx context.Context) ([]models.Remote, error) 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. +// ClaimGitHubSyncLease atomically claims the per-remote github_rpm sync lease. +// See claimSyncLease. func (db *DB) ClaimGitHubSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (bool, string, error) { + return db.claimSyncLease(ctx, "github_rpm_sync_state", remoteName, owner, freshness, lease) +} + +// ReleaseGitHubSyncLease records a github_rpm scan outcome and frees the lease. +// See releaseSyncLease. +func (db *DB) ReleaseGitHubSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error { + return db.releaseSyncLease(ctx, "github_rpm_sync_state", remoteName, owner, res) +} + +// claimSyncLease atomically claims a per-remote sync lease in table. It succeeds +// (claimed=true) only when the remote is due and no live lease is held by +// another replica. Due means: a pending retry after a failed scan has reached +// next_retry_at, or, with no retry pending, the remote was never synced or was +// synced longer than freshness ago. A zero freshness (prime scans) ignores the +// recency gate but still honours a pending retry's backoff. The returned etag is +// the stored releases-list ETag, shared across replicas so a conditional request +// can short-circuit an unchanged repo. +func (db *DB) claimSyncLease(ctx context.Context, table, 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) + INSERT INTO `+table+` 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)) + WHERE (CASE WHEN s.next_retry_at IS NOT NULL THEN s.next_retry_at <= now() + ELSE s.last_synced_at IS NULL OR s.last_synced_at < now() - make_interval(secs => $3) END) AND (s.sync_lease_expires IS NULL OR s.sync_lease_expires < now()) RETURNING s.etag `, remoteName, owner, freshness.Seconds(), lease.Seconds()) @@ -60,14 +75,27 @@ func (db *DB) ClaimGitHubSyncLease(ctx context.Context, remoteName, owner string 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 { +// releaseSyncLease records a scan outcome and frees the lease; only the owning +// replica may release. Success advances last_synced_at, persists the etag and +// clears any retry. Failure leaves last_synced_at and etag untouched (the last +// good metadata keeps serving) and schedules next_retry_at with exponential +// backoff from the consecutive-failure count, never before res.RetryAt. +func (db *DB) releaseSyncLease(ctx context.Context, table, remoteName, owner string, res provider.SyncResult) error { + var retryAt *time.Time + if !res.RetryAt.IsZero() { + retryAt = &res.RetryAt + } _, err := db.Pool.Exec(ctx, ` - UPDATE github_rpm_sync_state - SET last_synced_at = $3, etag = $4, sync_lease_owner = '', sync_lease_expires = NULL + UPDATE `+table+` + SET last_synced_at = CASE WHEN $3 THEN last_synced_at ELSE now() END, + etag = CASE WHEN $3 THEN etag ELSE $4 END, + sync_failures = CASE WHEN $3 THEN sync_failures + 1 ELSE 0 END, + next_retry_at = CASE WHEN $3 THEN GREATEST( + now() + make_interval(secs => LEAST($5 * power(2, LEAST(sync_failures, 20)), $6)), + $7::timestamptz) + END, + sync_lease_owner = '', sync_lease_expires = NULL WHERE remote_name = $1 AND sync_lease_owner = $2 - `, remoteName, owner, syncedAt, etag) + `, remoteName, owner, res.Failed, res.Etag, res.Backoff.Seconds(), res.MaxBackoff.Seconds(), retryAt) return err } diff --git a/internal/database/github_sync_test.go b/internal/database/github_sync_test.go index 7c6cc69..c97449b 100644 --- a/internal/database/github_sync_test.go +++ b/internal/database/github_sync_test.go @@ -4,6 +4,7 @@ import ( "testing" "time" + "git.unkin.net/unkin/artifactapi/internal/provider" "git.unkin.net/unkin/artifactapi/pkg/models" ) @@ -47,7 +48,7 @@ func TestGitHubSyncLease(t *testing.T) { } // Replica 1 finishes: record the sync and persist an etag. - if err := testDB.ReleaseGitHubSyncLease(ctx(), name, "replica-1", `"etag-1"`, time.Now()); err != nil { + if err := testDB.ReleaseGitHubSyncLease(ctx(), name, "replica-1", provider.SyncResult{Etag: `"etag-1"`}); err != nil { t.Fatalf("release: %v", err) } @@ -93,3 +94,116 @@ func TestListGitHubRPMRemotes(t *testing.T) { t.Fatalf("seeded remote %q not returned", name) } } + +// A failed scan must not count as a sync: the etag and last_synced_at stay put, +// every claim (prime included) is held off until the backoff elapses, the delay +// doubles per consecutive failure up to the cap, an upstream retry hint pushes +// it later, and a success clears the retry state. +func TestGitHubSyncLeaseFailureBackoff(t *testing.T) { + requireDB(t) + name := "gh-retry-" + time.Now().Format("150405.000000") + seedGitHubRPMRemote(t, name) + + const lease = 15 * time.Minute + freshness := time.Hour + claim := func(f time.Duration) (bool, string) { + t.Helper() + ok, etag, err := testDB.ClaimGitHubSyncLease(ctx(), name, "r1", f, lease) + if err != nil { + t.Fatalf("claim: %v", err) + } + return ok, etag + } + release := func(res provider.SyncResult) { + t.Helper() + if err := testDB.ReleaseGitHubSyncLease(ctx(), name, "r1", res); err != nil { + t.Fatalf("release: %v", err) + } + } + state := func() (failures int, retryIn time.Duration, synced *time.Time, etag string) { + t.Helper() + var next *time.Time + var dbNow time.Time + if err := testDB.Pool.QueryRow(ctx(), `SELECT sync_failures, next_retry_at, last_synced_at, etag, now() FROM github_rpm_sync_state WHERE remote_name = $1`, name). + Scan(&failures, &next, &synced, &etag, &dbNow); err != nil { + t.Fatalf("read state: %v", err) + } + if next != nil { + retryIn = next.Sub(dbNow) + } + return + } + elapse := func() { + t.Helper() + if _, err := testDB.Pool.Exec(ctx(), `UPDATE github_rpm_sync_state SET next_retry_at = now() - interval '1 second' WHERE remote_name = $1`, name); err != nil { + t.Fatalf("elapse backoff: %v", err) + } + } + failed := provider.SyncResult{Failed: true, Etag: `"ignored"`, Backoff: time.Minute, MaxBackoff: 3 * time.Minute} + + // Last good sync. + if ok, _ := claim(freshness); !ok { + t.Fatal("first claim failed") + } + release(provider.SyncResult{Etag: `"good"`}) + _, _, goodSynced, _ := state() + + // First failure: 60s backoff, sync time and etag untouched. + if ok, _ := claim(0); !ok { + t.Fatal("prime claim failed") + } + release(failed) + n, in, synced, etag := state() + if n != 1 || in < 55*time.Second || in > 65*time.Second { + t.Fatalf("after 1st failure: failures=%d retry in %v, want 1 and ~60s", n, in) + } + if etag != `"good"` || synced == nil || !synced.Equal(*goodSynced) { + t.Fatalf("failed scan overwrote sync state: etag=%q synced=%v", etag, synced) + } + if ok, _ := claim(0); ok { + t.Fatal("prime claimed inside the retry backoff") + } + + // Backoff elapsed: claimable despite being inside the freshness window, + // and the second failure doubles the delay. + elapse() + if ok, _ := claim(freshness); !ok { + t.Fatal("periodic claim refused after backoff elapsed") + } + release(failed) + if n, in, _, _ := state(); n != 2 || in < 115*time.Second || in > 125*time.Second { + t.Fatalf("after 2nd failure: failures=%d retry in %v, want 2 and ~120s", n, in) + } + + // Third failure caps at MaxBackoff (3m, not 4m). + elapse() + claim(freshness) + release(failed) + if n, in, _, _ := state(); n != 3 || in < 175*time.Second || in > 185*time.Second { + t.Fatalf("after 3rd failure: failures=%d retry in %v, want 3 and ~180s", n, in) + } + + // An upstream rate-limit reset later than the backoff wins. + elapse() + claim(freshness) + hinted := failed + hinted.RetryAt = time.Now().Add(20 * time.Minute) + release(hinted) + if _, in, _, _ := state(); in < 19*time.Minute || in > 21*time.Minute { + t.Fatalf("retry hint ignored: retry in %v, want ~20m", in) + } + + // Recovery: success persists the new etag and clears the retry, so the + // normal freshness window applies again. + elapse() + if ok, etag := claim(freshness); !ok || etag != `"good"` { + t.Fatalf("recovery claim ok=%v etag=%q, want last good etag", ok, etag) + } + release(provider.SyncResult{Etag: `"new"`}) + if n, in, _, etag := state(); n != 0 || in != 0 || etag != `"new"` { + t.Fatalf("after success: failures=%d retry in %v etag=%q", n, in, etag) + } + if ok, _ := claim(freshness); ok { + t.Fatal("claimed inside freshness window after a successful sync") + } +} diff --git a/internal/provider/alpine/github.go b/internal/provider/alpine/github.go index 13108ca..b60600e 100644 --- a/internal/provider/alpine/github.go +++ b/internal/provider/alpine/github.go @@ -470,7 +470,7 @@ func (p *GitHubProvider) fetchReleases(ctx context.Context, remote models.Remote return nil, "", false, err } if resp.StatusCode != http.StatusOK { - return nil, "", false, fmt.Errorf("github releases API %s: status %d", u, resp.StatusCode) + return nil, "", false, provider.NewUpstreamStatusError("github releases API "+u, resp) } if page == 1 { newEtag = respEtag diff --git a/internal/provider/alpine/syncer.go b/internal/provider/alpine/syncer.go index 8cd9847..d4f3fca 100644 --- a/internal/provider/alpine/syncer.go +++ b/internal/provider/alpine/syncer.go @@ -28,7 +28,7 @@ type SyncStore interface { provider.RemoteMetadataStore ListGitHubAlpineRemotes(ctx context.Context) ([]models.Remote, error) ClaimGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, freshness, lease time.Duration) (claimed bool, etag string, err error) - ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner, etag string, syncedAt time.Time) error + ReleaseGitHubAlpineSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error } // SyncConfig tunes the shared syncer. Zero values fall back to safe defaults. @@ -189,10 +189,11 @@ func (s *Syncer) process(ctx context.Context, job syncJob) { s.mu.Unlock() }() - freshness := time.Duration(job.remote.MutableTTL) * time.Second - if freshness <= 0 { - freshness = defaultSyncFreshness + ttl := time.Duration(job.remote.MutableTTL) * time.Second + if ttl <= 0 { + ttl = defaultSyncFreshness } + freshness := ttl if job.prime { freshness = 0 } @@ -210,16 +211,13 @@ func (s *Syncer) process(ctx context.Context, job syncJob) { defer cancel() newEtag, changed, scanErr := s.prov.scanWithState(scanCtx, job.remote, s.store, etag) - releaseEtag := etag - if scanErr == nil { - releaseEtag = newEtag - } else { + if scanErr != nil { slog.Error("github_alpine syncer: scan failed", "remote", job.remote.Name, "error", scanErr) } relCtx, relCancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second) defer relCancel() - if err := s.store.ReleaseGitHubAlpineSyncLease(relCtx, job.remote.Name, s.owner, releaseEtag, time.Now()); err != nil { + if err := s.store.ReleaseGitHubAlpineSyncLease(relCtx, job.remote.Name, s.owner, provider.NewSyncResult(newEtag, scanErr, ttl)); err != nil { slog.Warn("github_alpine syncer: release lease", "remote", job.remote.Name, "error", err) } diff --git a/internal/provider/alpine/syncer_test.go b/internal/provider/alpine/syncer_test.go index f8c2dc8..0401c52 100644 --- a/internal/provider/alpine/syncer_test.go +++ b/internal/provider/alpine/syncer_test.go @@ -27,6 +27,7 @@ type fakeSyncStore struct { leaseExp map[string]time.Time lastSynced map[string]time.Time etags map[string]string + retryAt map[string]time.Time } func newFakeSyncStore() *fakeSyncStore { @@ -36,6 +37,7 @@ func newFakeSyncStore() *fakeSyncStore { leaseExp: map[string]time.Time{}, lastSynced: map[string]time.Time{}, etags: map[string]string{}, + retryAt: map[string]time.Time{}, } } @@ -52,6 +54,9 @@ func (f *fakeSyncStore) ClaimGitHubAlpineSyncLease(_ context.Context, name, owne ls, hasLS := f.lastSynced[name] exp, hasExp := f.leaseExp[name] freshOK := !hasLS || now.Sub(ls) >= freshness + if ra, pending := f.retryAt[name]; pending { + freshOK = !now.Before(ra) + } leaseOK := !hasExp || exp.Before(now) if freshOK && leaseOK { f.leaseOwner[name] = owner @@ -61,14 +66,19 @@ func (f *fakeSyncStore) ClaimGitHubAlpineSyncLease(_ context.Context, name, owne return false, "", nil } -func (f *fakeSyncStore) ReleaseGitHubAlpineSyncLease(_ context.Context, name, owner, etag string, syncedAt time.Time) error { +func (f *fakeSyncStore) ReleaseGitHubAlpineSyncLease(_ context.Context, name, owner string, res provider.SyncResult) error { f.mu.Lock() defer f.mu.Unlock() if f.leaseOwner[name] != owner { return nil } - f.lastSynced[name] = syncedAt - f.etags[name] = etag + if res.Failed { + f.retryAt[name] = time.Now().Add(res.Backoff) + } else { + f.lastSynced[name] = time.Now() + f.etags[name] = res.Etag + delete(f.retryAt, name) + } delete(f.leaseOwner, name) delete(f.leaseExp, name) return nil diff --git a/internal/provider/deb/github.go b/internal/provider/deb/github.go index 3c8fcad..8731db3 100644 --- a/internal/provider/deb/github.go +++ b/internal/provider/deb/github.go @@ -431,7 +431,7 @@ func (p *GitHubProvider) fetchReleases(ctx context.Context, remote models.Remote return nil, "", false, err } if resp.StatusCode != http.StatusOK { - return nil, "", false, fmt.Errorf("github releases API %s: status %d", u, resp.StatusCode) + return nil, "", false, provider.NewUpstreamStatusError("github releases API "+u, resp) } if page == 1 { newEtag = respEtag diff --git a/internal/provider/deb/syncer.go b/internal/provider/deb/syncer.go index a19f463..b261134 100644 --- a/internal/provider/deb/syncer.go +++ b/internal/provider/deb/syncer.go @@ -28,7 +28,7 @@ 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 + ReleaseGitHubDebSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error } // SyncConfig tunes the shared syncer. Zero values fall back to safe defaults. @@ -189,10 +189,11 @@ func (s *Syncer) process(ctx context.Context, job syncJob) { s.mu.Unlock() }() - freshness := time.Duration(job.remote.MutableTTL) * time.Second - if freshness <= 0 { - freshness = defaultSyncFreshness + ttl := time.Duration(job.remote.MutableTTL) * time.Second + if ttl <= 0 { + ttl = defaultSyncFreshness } + freshness := ttl if job.prime { freshness = 0 } @@ -210,16 +211,13 @@ func (s *Syncer) process(ctx context.Context, job syncJob) { defer cancel() newEtag, changed, scanErr := s.prov.scanWithState(scanCtx, job.remote, s.store, etag) - releaseEtag := etag - if scanErr == nil { - releaseEtag = newEtag - } else { + if scanErr != nil { 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 { + if err := s.store.ReleaseGitHubDebSyncLease(relCtx, job.remote.Name, s.owner, provider.NewSyncResult(newEtag, scanErr, ttl)); err != nil { slog.Warn("github_deb syncer: release lease", "remote", job.remote.Name, "error", err) } diff --git a/internal/provider/deb/syncer_test.go b/internal/provider/deb/syncer_test.go index 12a684d..e3fe662 100644 --- a/internal/provider/deb/syncer_test.go +++ b/internal/provider/deb/syncer_test.go @@ -27,6 +27,7 @@ type fakeSyncStore struct { leaseExp map[string]time.Time lastSynced map[string]time.Time etags map[string]string + retryAt map[string]time.Time } func newFakeSyncStore() *fakeSyncStore { @@ -36,6 +37,7 @@ func newFakeSyncStore() *fakeSyncStore { leaseExp: map[string]time.Time{}, lastSynced: map[string]time.Time{}, etags: map[string]string{}, + retryAt: map[string]time.Time{}, } } @@ -52,6 +54,9 @@ func (f *fakeSyncStore) ClaimGitHubDebSyncLease(_ context.Context, name, owner s ls, hasLS := f.lastSynced[name] exp, hasExp := f.leaseExp[name] freshOK := !hasLS || now.Sub(ls) >= freshness + if ra, pending := f.retryAt[name]; pending { + freshOK = !now.Before(ra) + } leaseOK := !hasExp || exp.Before(now) if freshOK && leaseOK { f.leaseOwner[name] = owner @@ -61,14 +66,19 @@ func (f *fakeSyncStore) ClaimGitHubDebSyncLease(_ context.Context, name, owner s return false, "", nil } -func (f *fakeSyncStore) ReleaseGitHubDebSyncLease(_ context.Context, name, owner, etag string, syncedAt time.Time) error { +func (f *fakeSyncStore) ReleaseGitHubDebSyncLease(_ context.Context, name, owner string, res provider.SyncResult) error { f.mu.Lock() defer f.mu.Unlock() if f.leaseOwner[name] != owner { return nil } - f.lastSynced[name] = syncedAt - f.etags[name] = etag + if res.Failed { + f.retryAt[name] = time.Now().Add(res.Backoff) + } else { + f.lastSynced[name] = time.Now() + f.etags[name] = res.Etag + delete(f.retryAt, name) + } delete(f.leaseOwner, name) delete(f.leaseExp, name) return nil diff --git a/internal/provider/rpm/github.go b/internal/provider/rpm/github.go index e28253f..c41d47d 100644 --- a/internal/provider/rpm/github.go +++ b/internal/provider/rpm/github.go @@ -447,7 +447,7 @@ func (p *GitHubProvider) fetchReleases(ctx context.Context, remote models.Remote return nil, "", false, err } if resp.StatusCode != http.StatusOK { - return nil, "", false, fmt.Errorf("github releases API %s: status %d", u, resp.StatusCode) + return nil, "", false, provider.NewUpstreamStatusError("github releases API "+u, resp) } if page == 1 { newEtag = respEtag diff --git a/internal/provider/rpm/github_test.go b/internal/provider/rpm/github_test.go index 59debd0..33361f5 100644 --- a/internal/provider/rpm/github_test.go +++ b/internal/provider/rpm/github_test.go @@ -78,6 +78,7 @@ type githubFixture struct { notModHit int // releases-list requests answered 304 releaseAuth string // Authorization header seen on the last releases request assetAuth string // Authorization header seen on the last asset request + failStatus int // when set, the releases list answers this status (rate-limit style) mu sync.Mutex } @@ -100,6 +101,14 @@ func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture { f.mu.Lock() f.releasesHit++ f.releaseAuth = r.Header.Get("Authorization") + if f.failStatus != 0 { + status := f.failStatus + f.mu.Unlock() + w.Header().Set("X-RateLimit-Remaining", "0") + w.Header().Set("X-RateLimit-Reset", strconv.FormatInt(time.Now().Add(30*time.Second).Unix(), 10)) + http.Error(w, "API rate limit exceeded", status) + return + } etag := f.etag if etag != "" && r.Header.Get("If-None-Match") == etag { f.notModHit++ diff --git a/internal/provider/rpm/syncer.go b/internal/provider/rpm/syncer.go index f6f5724..ba7cd1f 100644 --- a/internal/provider/rpm/syncer.go +++ b/internal/provider/rpm/syncer.go @@ -35,7 +35,7 @@ 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 + ReleaseGitHubSyncLease(ctx context.Context, remoteName, owner string, res provider.SyncResult) error } // SyncConfig tunes the shared syncer. Zero values fall back to safe defaults. @@ -205,10 +205,11 @@ func (s *Syncer) process(ctx context.Context, job syncJob) { s.mu.Unlock() }() - freshness := time.Duration(job.remote.MutableTTL) * time.Second - if freshness <= 0 { - freshness = defaultSyncFreshness + ttl := time.Duration(job.remote.MutableTTL) * time.Second + if ttl <= 0 { + ttl = defaultSyncFreshness } + freshness := ttl if job.prime { freshness = 0 // prime ignores the recency gate but still respects a live lease } @@ -226,18 +227,15 @@ func (s *Syncer) process(ctx context.Context, job syncJob) { defer cancel() newEtag, changed, scanErr := s.prov.scanWithState(scanCtx, job.remote, s.store, etag) - releaseEtag := etag - if scanErr == nil { - releaseEtag = newEtag - } else { + if scanErr != nil { 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). + // lease and records the outcome (otherwise the lease 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 { + if err := s.store.ReleaseGitHubSyncLease(relCtx, job.remote.Name, s.owner, provider.NewSyncResult(newEtag, scanErr, ttl)); err != nil { slog.Warn("github_rpm syncer: release lease", "remote", job.remote.Name, "error", err) } diff --git a/internal/provider/rpm/syncer_test.go b/internal/provider/rpm/syncer_test.go index 727e0ec..e368bd2 100644 --- a/internal/provider/rpm/syncer_test.go +++ b/internal/provider/rpm/syncer_test.go @@ -27,6 +27,8 @@ type fakeSyncStore struct { leaseExp map[string]time.Time lastSynced map[string]time.Time etags map[string]string + retryAt map[string]time.Time + lastResult provider.SyncResult } func newFakeSyncStore() *fakeSyncStore { @@ -36,6 +38,7 @@ func newFakeSyncStore() *fakeSyncStore { leaseExp: map[string]time.Time{}, lastSynced: map[string]time.Time{}, etags: map[string]string{}, + retryAt: map[string]time.Time{}, } } @@ -52,6 +55,9 @@ func (f *fakeSyncStore) ClaimGitHubSyncLease(_ context.Context, name, owner stri ls, hasLS := f.lastSynced[name] exp, hasExp := f.leaseExp[name] freshOK := !hasLS || now.Sub(ls) >= freshness + if ra, pending := f.retryAt[name]; pending { + freshOK = !now.Before(ra) + } leaseOK := !hasExp || exp.Before(now) if freshOK && leaseOK { f.leaseOwner[name] = owner @@ -61,14 +67,20 @@ func (f *fakeSyncStore) ClaimGitHubSyncLease(_ context.Context, name, owner stri return false, "", nil } -func (f *fakeSyncStore) ReleaseGitHubSyncLease(_ context.Context, name, owner, etag string, syncedAt time.Time) error { +func (f *fakeSyncStore) ReleaseGitHubSyncLease(_ context.Context, name, owner string, res provider.SyncResult) error { f.mu.Lock() defer f.mu.Unlock() if f.leaseOwner[name] != owner { return nil } - f.lastSynced[name] = syncedAt - f.etags[name] = etag + f.lastResult = res + if res.Failed { + f.retryAt[name] = time.Now().Add(res.Backoff) + } else { + f.lastSynced[name] = time.Now() + f.etags[name] = res.Etag + delete(f.retryAt, name) + } delete(f.leaseOwner, name) delete(f.leaseExp, name) return nil @@ -310,3 +322,80 @@ func TestSyncerPrimeBypassesRecencyPeriodicDoesNot(t *testing.T) { t.Fatalf("periodic scan ran inside recency window: %d -> %d releases calls", releasesAfterPrime, fx.releasesHit) } } + +// A rate-limited scan is not recorded as a sync: the last good repodata keeps +// serving, the retry honours the upstream reset hint and backoff, and once the +// upstream recovers the next poll after the backoff re-syncs without waiting +// out mutable_ttl. +func TestSyncerFailedScanRetriesAfterBackoff(t *testing.T) { + fx := newGitHubFixture(t, true) + fx.etag = `"v1"` + store := newFakeSyncStore() + p := newTestProvider() + s := newSyncer(store, p, testSyncConfig()) + remote := fx.remote() + bg := context.Background() + + s.process(bg, syncJob{remote: remote, prime: true}) + if rows, _ := store.ListRPMMetadataEntries(bg, remote.Name); len(rows) != 1 { + t.Fatalf("prime did not derive: %d rows", len(rows)) + } + + // mutable_ttl elapses while GitHub is rate-limiting. + fx.mu.Lock() + fx.failStatus = http.StatusForbidden + fx.etag = `"v2"` + fx.mu.Unlock() + store.mu.Lock() + store.lastSynced[remote.Name] = time.Now().Add(-2 * time.Hour) + store.mu.Unlock() + + s.process(bg, syncJob{remote: remote}) + res := store.lastResult + if !res.Failed { + t.Fatal("403 scan was recorded as a successful sync") + } + if until := time.Until(res.RetryAt); until < 20*time.Second || until > 40*time.Second { + t.Fatalf("retry hint from X-RateLimit-Reset not carried: retry in %v", until) + } + if res.Backoff != time.Minute || res.MaxBackoff >= time.Duration(remote.MutableTTL)*time.Second { + t.Fatalf("backoff %v cap %v, want 1m first retry capped below mutable_ttl", res.Backoff, res.MaxBackoff) + } + if store.etags[remote.Name] != `"v1"` { + t.Fatalf("failed scan replaced the etag: %q", store.etags[remote.Name]) + } + + rec := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodGet, "/api/v1/remote/acme-rpm/repodata/repomd.xml", nil) + p.ServeRemote(rec, req, remote, "repodata/repomd.xml", "https://x", store) + if rec.Code != http.StatusOK { + t.Fatalf("last good repodata not served during failure: %d %s", rec.Code, rec.Body.String()) + } + if rows, _ := store.ListRPMMetadataEntries(bg, remote.Name); len(rows) != 1 { + t.Fatalf("failed scan dropped cached metadata: %d rows", len(rows)) + } + + // Inside the backoff nothing is retried. + hits := fx.releasesHit + s.process(bg, syncJob{remote: remote}) + if fx.releasesHit != hits { + t.Fatal("scan retried inside the backoff window") + } + + // Upstream recovers and the backoff elapses: the next poll re-syncs. + fx.mu.Lock() + fx.failStatus = 0 + fx.rpmBytes["other-9-9.aarch64.rpm"] = testsupport.MinimalRPM("other", "9", "9", "aarch64") + fx.mu.Unlock() + store.mu.Lock() + store.retryAt[remote.Name] = time.Now().Add(-time.Second) + store.mu.Unlock() + + s.process(bg, syncJob{remote: remote}) + if store.lastResult.Failed || store.etags[remote.Name] != `"v2"` { + t.Fatalf("recovery scan not recorded: failed=%v etag=%q", store.lastResult.Failed, store.etags[remote.Name]) + } + if rows, _ := store.ListRPMMetadataEntries(bg, remote.Name); len(rows) != 2 { + t.Fatalf("recovery did not refresh metadata: %d rows", len(rows)) + } +} diff --git a/internal/provider/syncretry.go b/internal/provider/syncretry.go new file mode 100644 index 0000000..dff7070 --- /dev/null +++ b/internal/provider/syncretry.go @@ -0,0 +1,75 @@ +package provider + +import ( + "errors" + "fmt" + "net/http" + "strconv" + "time" +) + +const ( + syncRetryBase = time.Minute + syncRetryMax = 10 * time.Minute +) + +// UpstreamStatusError is a non-success upstream response. RetryAt is the +// upstream's own retry hint (Retry-After, or X-RateLimit-Reset once the quota +// is exhausted); zero when it gave none. +type UpstreamStatusError struct { + URL string + Status int + RetryAt time.Time +} + +func (e *UpstreamStatusError) Error() string { + return fmt.Sprintf("%s: status %d", e.URL, e.Status) +} + +// NewUpstreamStatusError wraps a non-success response, capturing its retry hint. +func NewUpstreamStatusError(url string, resp *http.Response) *UpstreamStatusError { + e := &UpstreamStatusError{URL: url, Status: resp.StatusCode} + if ra := resp.Header.Get("Retry-After"); ra != "" { + if secs, err := strconv.Atoi(ra); err == nil { + e.RetryAt = time.Now().Add(time.Duration(secs) * time.Second) + } else if t, err := http.ParseTime(ra); err == nil { + e.RetryAt = t + } + } else if resp.Header.Get("X-RateLimit-Remaining") == "0" { + if reset, err := strconv.ParseInt(resp.Header.Get("X-RateLimit-Reset"), 10, 64); err == nil { + e.RetryAt = time.Unix(reset, 0) + } + } + return e +} + +// SyncResult is a background scan's outcome, recorded when its sync lease is +// released. A failed scan keeps the prior sync time and ETag and schedules a +// retry after Backoff, doubled per consecutive failure up to MaxBackoff, and +// never earlier than RetryAt. +type SyncResult struct { + Etag string + Failed bool + RetryAt time.Time + Backoff time.Duration + MaxBackoff time.Duration +} + +// NewSyncResult builds the result for a scan against a remote with the given +// mutable_ttl. The retry cap stays well below ttl, and an upstream hint is +// clamped to ttl so a bogus reset can never stall the remote longer than a +// normal sync interval would. +func NewSyncResult(etag string, scanErr error, ttl time.Duration) SyncResult { + if scanErr == nil { + return SyncResult{Etag: etag} + } + res := SyncResult{Failed: true, Backoff: syncRetryBase, MaxBackoff: min(syncRetryMax, max(syncRetryBase, ttl/4))} + var se *UpstreamStatusError + if errors.As(scanErr, &se) && !se.RetryAt.IsZero() { + res.RetryAt = se.RetryAt + if limit := time.Now().Add(ttl); res.RetryAt.After(limit) { + res.RetryAt = limit + } + } + return res +} diff --git a/internal/provider/syncretry_test.go b/internal/provider/syncretry_test.go new file mode 100644 index 0000000..289a1d7 --- /dev/null +++ b/internal/provider/syncretry_test.go @@ -0,0 +1,63 @@ +package provider + +import ( + "errors" + "fmt" + "net/http" + "strconv" + "testing" + "time" +) + +func TestNewUpstreamStatusErrorRetryHint(t *testing.T) { + reset := time.Now().Add(15 * time.Minute).Truncate(time.Second) + cases := []struct { + name string + header map[string]string + want time.Duration // 0 = no hint + }{ + {"retry-after seconds", map[string]string{"Retry-After": "90"}, 90 * time.Second}, + {"retry-after date", map[string]string{"Retry-After": reset.UTC().Format(http.TimeFormat)}, 15 * time.Minute}, + {"quota exhausted", map[string]string{"X-RateLimit-Remaining": "0", "X-RateLimit-Reset": strconv.FormatInt(reset.Unix(), 10)}, 15 * time.Minute}, + {"quota left is not a hint", map[string]string{"X-RateLimit-Remaining": "12", "X-RateLimit-Reset": strconv.FormatInt(reset.Unix(), 10)}, 0}, + {"no headers", nil, 0}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + resp := &http.Response{StatusCode: http.StatusForbidden, Header: http.Header{}} + for k, v := range c.header { + resp.Header.Set(k, v) + } + e := NewUpstreamStatusError("u", resp) + if c.want == 0 { + if !e.RetryAt.IsZero() { + t.Fatalf("unexpected hint %v", e.RetryAt) + } + return + } + if d := time.Until(e.RetryAt) - c.want; d < -2*time.Second || d > 2*time.Second { + t.Fatalf("retry in %v, want ~%v", time.Until(e.RetryAt), c.want) + } + }) + } +} + +func TestNewSyncResult(t *testing.T) { + if r := NewSyncResult(`"e"`, nil, time.Hour); r.Failed || r.Etag != `"e"` { + t.Fatalf("success result = %+v", r) + } + + r := NewSyncResult(`"e"`, errors.New("boom"), time.Hour) + if !r.Failed || r.Backoff != time.Minute || r.MaxBackoff != 10*time.Minute || !r.RetryAt.IsZero() { + t.Fatalf("plain failure = %+v, want 1m backoff capped at 10m, no hint", r) + } + if r := NewSyncResult("", errors.New("boom"), 5*time.Minute); r.MaxBackoff != time.Minute+15*time.Second { + t.Fatalf("short ttl cap = %v, want ttl/4", r.MaxBackoff) + } + + far := &UpstreamStatusError{Status: 403, RetryAt: time.Now().Add(3 * time.Hour)} + r = NewSyncResult("", fmt.Errorf("scan: %w", far), time.Hour) + if until := time.Until(r.RetryAt); until > time.Hour || until < 59*time.Minute { + t.Fatalf("wrapped hint not clamped to ttl: retry in %v", until) + } +} diff --git a/migrations/0002_github_sync_retry.sql b/migrations/0002_github_sync_retry.sql new file mode 100644 index 0000000..2814746 --- /dev/null +++ b/migrations/0002_github_sync_retry.sql @@ -0,0 +1,9 @@ +-- Failed GitHub release scans retry with backoff instead of waiting out a full +-- mutable_ttl: sync_failures counts consecutive failures, next_retry_at gates +-- the next claim while a retry is pending. +ALTER TABLE github_rpm_sync_state ADD COLUMN IF NOT EXISTS sync_failures INT NOT NULL DEFAULT 0; +ALTER TABLE github_rpm_sync_state ADD COLUMN IF NOT EXISTS next_retry_at TIMESTAMPTZ; +ALTER TABLE github_deb_sync_state ADD COLUMN IF NOT EXISTS sync_failures INT NOT NULL DEFAULT 0; +ALTER TABLE github_deb_sync_state ADD COLUMN IF NOT EXISTS next_retry_at TIMESTAMPTZ; +ALTER TABLE github_alpine_sync_state ADD COLUMN IF NOT EXISTS sync_failures INT NOT NULL DEFAULT 0; +ALTER TABLE github_alpine_sync_state ADD COLUMN IF NOT EXISTS next_retry_at TIMESTAMPTZ;