Fail GitHub scans on asset errors and drop sync state with its remote #135

Open
unkin-agent wants to merge 4 commits from benvin/github-asset-fail-etag into master
7 changed files with 150 additions and 6 deletions
Showing only changes of commit b111c2e57f - Show all commits
+69
View File
@@ -207,3 +207,72 @@ func TestGitHubSyncLeaseFailureBackoff(t *testing.T) {
t.Fatal("claimed inside freshness window after a successful sync")
}
}
// Deleting a remote drops its sync state, so a recreated remote of the same name
// starts clean instead of inheriting the old ETag, sync time and retry state.
func TestGitHubSyncStateDeletedWithRemote(t *testing.T) {
requireDB(t)
for _, tc := range []struct {
table string
pt models.PackageType
claim func(name string) (bool, string, error)
release func(name string, res provider.SyncResult) error
}{
{"github_rpm_sync_state", models.PackageGitHubRPM,
func(n string) (bool, string, error) {
return testDB.ClaimGitHubSyncLease(ctx(), n, "r1", time.Hour, time.Minute)
},
func(n string, res provider.SyncResult) error {
return testDB.ReleaseGitHubSyncLease(ctx(), n, "r1", res)
}},
{"github_deb_sync_state", models.PackageGitHubDeb,
func(n string) (bool, string, error) {
return testDB.ClaimGitHubDebSyncLease(ctx(), n, "r1", time.Hour, time.Minute)
},
func(n string, res provider.SyncResult) error {
return testDB.ReleaseGitHubDebSyncLease(ctx(), n, "r1", res)
}},
{"github_alpine_sync_state", models.PackageGitHubAlpine,
func(n string) (bool, string, error) {
return testDB.ClaimGitHubAlpineSyncLease(ctx(), n, "r1", time.Hour, time.Minute)
},
func(n string, res provider.SyncResult) error {
return testDB.ReleaseGitHubAlpineSyncLease(ctx(), n, "r1", res)
}},
} {
t.Run(tc.table, func(t *testing.T) {
name := "gh-recreate-" + time.Now().Format("150405.000000")
create := func() {
t.Helper()
if err := testDB.CreateRemote(ctx(), &models.Remote{
Name: name, PackageType: tc.pt, RepoType: models.RepoTypeRemote,
BaseURL: "https://api.github.com/repos/acme/tools", ReleasesRemote: "github", MutableTTL: 3600,
}); err != nil {
t.Fatalf("create remote: %v", err)
}
}
create()
if ok, _, err := tc.claim(name); err != nil || !ok {
t.Fatalf("claim: ok=%v err=%v", ok, err)
}
if err := tc.release(name, provider.SyncResult{Etag: `"old"`}); err != nil {
t.Fatalf("release: %v", err)
}
if err := testDB.DeleteRemote(ctx(), name); err != nil {
t.Fatalf("delete remote: %v", err)
}
var rows int
if err := testDB.Pool.QueryRow(ctx(), `SELECT count(*) FROM `+tc.table+` WHERE remote_name = $1`, name).Scan(&rows); err != nil || rows != 0 {
t.Fatalf("sync state left after delete: rows=%d err=%v", rows, err)
}
create()
ok, etag, err := tc.claim(name)
if err != nil || !ok || etag != "" {
t.Fatalf("recreated remote claim ok=%v etag=%q err=%v, want a fresh claim with no etag", ok, etag, err)
}
})
}
}
+6 -1
View File
@@ -372,6 +372,9 @@ func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote
return newEtag, false, err
}
// A failed asset fails the scan so the ETag is not advanced and the retry
// re-derives it; assets that did derive are kept and served meanwhile.
var failed []error
seen := map[string]bool{}
for _, rel := range releases {
if rel.Draft {
@@ -400,10 +403,12 @@ func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote
meta, err := p.deriveAsset(ctx, remote, asset, fp)
if err != nil {
slog.Warn("github_alpine: derive asset failed", "remote", remote.Name, "asset", asset.Name, "error", err)
failed = append(failed, fmt.Errorf("%s: %w", asset.Name, err))
continue
}
if err := inserter.InsertAlpineMetadata(ctx, meta); err != nil {
slog.Error("github_alpine: insert metadata failed", "remote", remote.Name, "asset", asset.Name, "error", err)
failed = append(failed, fmt.Errorf("%s: %w", asset.Name, err))
continue
}
slog.Info("github_alpine: derived asset", "remote", remote.Name, "name", meta.Name, "version", meta.Version, "arch", meta.Arch)
@@ -415,7 +420,7 @@ func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote
_ = deleter.DeleteAlpineMetadata(ctx, remote.Name, fp)
}
}
return newEtag, true, nil
return newEtag, true, errors.Join(failed...)
}
type ghRelease struct {
+6 -1
View File
@@ -333,6 +333,9 @@ func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote
return newEtag, false, err
}
// A failed asset fails the scan so the ETag is not advanced and the retry
// re-derives it; assets that did derive are kept and served meanwhile.
var failed []error
seen := map[string]bool{}
for _, rel := range releases {
if rel.Draft {
@@ -361,10 +364,12 @@ func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote
meta, err := p.deriveAsset(ctx, remote, asset, fp)
if err != nil {
slog.Warn("github_deb: derive asset failed", "remote", remote.Name, "asset", asset.Name, "error", err)
failed = append(failed, fmt.Errorf("%s: %w", asset.Name, err))
continue
}
if err := store.InsertDebMetadata(ctx, meta); err != nil {
slog.Error("github_deb: insert metadata failed", "remote", remote.Name, "asset", asset.Name, "error", err)
failed = append(failed, fmt.Errorf("%s: %w", asset.Name, err))
continue
}
slog.Info("github_deb: derived asset", "remote", remote.Name, "name", meta.Name, "version", meta.Version, "arch", meta.Architecture)
@@ -376,7 +381,7 @@ func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote
_ = store.DeleteDebMetadata(ctx, remote.Name, fp)
}
}
return newEtag, true, nil
return newEtag, true, errors.Join(failed...)
}
type ghRelease struct {
+6 -1
View File
@@ -343,6 +343,9 @@ func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote
return newEtag, false, err
}
// A failed asset fails the scan so the ETag is not advanced and the retry
// re-derives it; assets that did derive are kept and served meanwhile.
var failed []error
seen := map[string]bool{}
for _, rel := range releases {
if rel.Draft {
@@ -373,10 +376,12 @@ func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote
meta, err := p.deriveAsset(ctx, remote, asset, fp)
if err != nil {
slog.Warn("github_rpm: derive asset failed", "remote", remote.Name, "asset", asset.Name, "error", err)
failed = append(failed, fmt.Errorf("%s: %w", asset.Name, err))
continue
}
if err := store.InsertRPMMetadata(ctx, meta); err != nil {
slog.Error("github_rpm: insert metadata failed", "remote", remote.Name, "asset", asset.Name, "error", err)
failed = append(failed, fmt.Errorf("%s: %w", asset.Name, err))
continue
}
slog.Info("github_rpm: derived asset", "remote", remote.Name, "name", meta.Name, "version", meta.Version, "arch", meta.Arch)
@@ -388,7 +393,7 @@ func (p *GitHubProvider) scanWithState(ctx context.Context, remote models.Remote
_ = store.DeleteRPMMetadata(ctx, remote.Name, fp)
}
}
return newEtag, true, nil
return newEtag, true, errors.Join(failed...)
}
type ghRelease struct {
+11 -3
View File
@@ -79,15 +79,17 @@ type githubFixture struct {
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)
assetFail map[string]int // asset filename -> downloads left to answer 500
mu sync.Mutex
}
func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture {
t.Helper()
f := &githubFixture{
rpmBytes: map[string][]byte{},
rangeHit: map[string]int{},
fullHit: map[string]int{},
rpmBytes: map[string][]byte{},
rangeHit: map[string]int{},
fullHit: map[string]int{},
assetFail: map[string]int{},
}
f.rpmBytes["demo-1.2-3.x86_64.rpm"] = testsupport.MinimalRPM("demo", "1.2", "3", "x86_64")
@@ -145,6 +147,12 @@ func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture {
}
rng := r.Header.Get("Range")
f.mu.Lock()
if f.assetFail[name] > 0 {
f.assetFail[name]--
f.mu.Unlock()
http.Error(w, "boom", http.StatusInternalServerError)
return
}
f.assetAuth = r.Header.Get("Authorization")
if rng != "" {
f.rangeHit[name]++
+40
View File
@@ -399,3 +399,43 @@ func TestSyncerFailedScanRetriesAfterBackoff(t *testing.T) {
t.Fatalf("recovery did not refresh metadata: %d rows", len(rows))
}
}
// An asset that fails to derive fails the scan: the ETag is not advanced, a
// retry is scheduled, and the retry against the same unchanged release list
// derives the package instead of waiting for upstream to publish again.
func TestSyncerFailedAssetRetriesWithoutAdvancingEtag(t *testing.T) {
fx := newGitHubFixture(t, true)
fx.etag = `"v1"`
fx.rpmBytes["other-9-9.aarch64.rpm"] = testsupport.MinimalRPM("other", "9", "9", "aarch64")
fx.assetFail["other-9-9.aarch64.rpm"] = 1
store := newFakeSyncStore()
p := newTestProvider()
s := newSyncer(store, p, testSyncConfig())
remote := fx.remote()
bg := context.Background()
s.process(bg, syncJob{remote: remote, prime: true})
if !store.lastResult.Failed {
t.Fatal("scan with a failed asset was recorded as a successful sync")
}
if _, saved := store.etags[remote.Name]; saved {
t.Fatalf("failed scan saved etag %q", store.etags[remote.Name])
}
if _, pending := store.retryAt[remote.Name]; !pending {
t.Fatal("no retry scheduled after a failed asset")
}
if rows, _ := store.ListRPMMetadataEntries(bg, remote.Name); len(rows) != 1 {
t.Fatalf("healthy asset not served meanwhile: %d rows", len(rows))
}
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] != `"v1"` {
t.Fatalf("retry 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("retry did not derive the failed asset: %d rows", len(rows))
}
}
@@ -0,0 +1,12 @@
-- Sync state is dropped with its remote, so a recreated remote of the same name
-- starts with no inherited ETag or timing. NOT VALID skips checking rows left by
-- remotes deleted before this migration, which would otherwise fail it.
ALTER TABLE github_rpm_sync_state DROP CONSTRAINT IF EXISTS github_rpm_sync_state_remote_fk;
ALTER TABLE github_rpm_sync_state ADD CONSTRAINT github_rpm_sync_state_remote_fk
FOREIGN KEY (remote_name) REFERENCES remotes(name) ON DELETE CASCADE NOT VALID;
ALTER TABLE github_deb_sync_state DROP CONSTRAINT IF EXISTS github_deb_sync_state_remote_fk;
ALTER TABLE github_deb_sync_state ADD CONSTRAINT github_deb_sync_state_remote_fk
FOREIGN KEY (remote_name) REFERENCES remotes(name) ON DELETE CASCADE NOT VALID;
ALTER TABLE github_alpine_sync_state DROP CONSTRAINT IF EXISTS github_alpine_sync_state_remote_fk;
ALTER TABLE github_alpine_sync_state ADD CONSTRAINT github_alpine_sync_state_remote_fk
FOREIGN KEY (remote_name) REFERENCES remotes(name) ON DELETE CASCADE NOT VALID;