Retry failed GitHub release scans with backoff (#133)
ci/woodpecker/tag/docker Pipeline was successful
ci/woodpecker/tag/docker Pipeline was successful
Releasing a GitHub sync lease always advanced `last_synced_at`, even after a failed scan (e.g. a rate-limit 403). A remote that failed once waited a full `mutable_ttl` before retrying, and a cold remote kept returning 503 until then. - record scan outcomes in one shared lease helper for github_rpm/deb/alpine - keep `last_synced_at` and the ETag on failure; retry from 60s with exponential backoff, capped at min(10m, ttl/4) - honour `Retry-After` / `X-RateLimit-Reset`, clamped to `mutable_ttl` - add `sync_failures` / `next_retry_at` columns (migration 0002) Reviewed-on: #133 Co-authored-by: unkin-agent <unkin-agent@unkin.net> Co-committed-by: unkin-agent <unkin-agent@unkin.net>
This commit was merged in pull request #133.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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++
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user