Skip invalid GitHub packages instead of failing the scan
Wrap package parse errors in provider.ErrInvalidPackage so only transient asset and insert errors fail the scan; cover rpm, deb and alpine.
This commit is contained in:
@@ -372,8 +372,8 @@ 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.
|
||||
// A transiently failed asset fails the scan so the ETag is not advanced and
|
||||
// the retry re-derives it; an invalid package is skipped as no retry fixes it.
|
||||
var failed []error
|
||||
seen := map[string]bool{}
|
||||
for _, rel := range releases {
|
||||
@@ -403,7 +403,9 @@ 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))
|
||||
if !errors.Is(err, provider.ErrInvalidPackage) {
|
||||
failed = append(failed, fmt.Errorf("%s: %w", asset.Name, err))
|
||||
}
|
||||
continue
|
||||
}
|
||||
if err := inserter.InsertAlpineMetadata(ctx, meta); err != nil {
|
||||
@@ -501,7 +503,7 @@ func (p *GitHubProvider) deriveAsset(ctx context.Context, remote models.Remote,
|
||||
return nil, err
|
||||
}
|
||||
if meta.Name == "" || meta.Arch == "" {
|
||||
return nil, errors.New(".PKGINFO missing pkgname/arch")
|
||||
return nil, fmt.Errorf("%w: .PKGINFO missing pkgname/arch", provider.ErrInvalidPackage)
|
||||
}
|
||||
|
||||
meta.RepoName = remote.Name
|
||||
@@ -531,13 +533,13 @@ func (p *GitHubProvider) fetchPkginfo(ctx context.Context, remote models.Remote,
|
||||
}
|
||||
meta, complete, perr := pkginfoFromPrefix(body)
|
||||
if perr != nil {
|
||||
return nil, fmt.Errorf("parse apk .PKGINFO: %w", perr)
|
||||
return nil, fmt.Errorf("%w: parse apk .PKGINFO: %w", provider.ErrInvalidPackage, perr)
|
||||
}
|
||||
if complete {
|
||||
return meta, nil
|
||||
}
|
||||
if full || n >= p.headerMax {
|
||||
return nil, fmt.Errorf(".PKGINFO not found within %d bytes of %s", n, downloadURL)
|
||||
return nil, fmt.Errorf("%w: .PKGINFO not found within %d bytes of %s", provider.ErrInvalidPackage, n, downloadURL)
|
||||
}
|
||||
n *= 2
|
||||
if n > p.headerMax {
|
||||
|
||||
@@ -88,15 +88,17 @@ type githubFixture struct {
|
||||
notModHit int
|
||||
releaseAuth string
|
||||
assetAuth string
|
||||
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{
|
||||
apkBytes: map[string][]byte{},
|
||||
rangeHit: map[string]int{},
|
||||
fullHit: map[string]int{},
|
||||
apkBytes: map[string][]byte{},
|
||||
rangeHit: map[string]int{},
|
||||
fullHit: map[string]int{},
|
||||
assetFail: map[string]int{},
|
||||
}
|
||||
f.apkBytes["demo-1.2.3-r0.apk"] = testsupport.MinimalApk("demo", "1.2.3-r0", "x86_64")
|
||||
|
||||
@@ -146,6 +148,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]++
|
||||
|
||||
@@ -308,3 +308,58 @@ func TestSyncerPrimeBypassesRecencyPeriodicDoesNot(t *testing.T) {
|
||||
t.Fatalf("periodic scan ran inside recency window: %d -> %d releases calls", releasesAfterPrime, fx.releasesHit)
|
||||
}
|
||||
}
|
||||
|
||||
// A transiently failed asset fails the scan: the ETag is not advanced, a retry
|
||||
// is scheduled, and the retry against the unchanged release list derives it. An
|
||||
// invalid package is skipped instead, since no retry can fix it.
|
||||
func TestSyncerFailedAssetRetry(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
corrupt bool
|
||||
wantRetry bool
|
||||
}{
|
||||
{name: "transient", wantRetry: true},
|
||||
{name: "invalid package", corrupt: true},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
fx := newGitHubFixture(t, true)
|
||||
fx.etag = `"v1"`
|
||||
if tc.corrupt {
|
||||
fx.apkBytes["other-9-r0.apk"] = []byte("not a package")
|
||||
} else {
|
||||
fx.apkBytes["other-9-r0.apk"] = testsupport.MinimalApk("other", "9-r0", "aarch64")
|
||||
fx.assetFail["other-9-r0.apk"] = 1
|
||||
}
|
||||
store := newFakeSyncStore()
|
||||
s := newSyncer(store, newTestProvider(), testSyncConfig())
|
||||
remote := fx.remote()
|
||||
bg := context.Background()
|
||||
|
||||
s.process(bg, syncJob{remote: remote, prime: true})
|
||||
_, pending := store.retryAt[remote.Name]
|
||||
if pending != tc.wantRetry {
|
||||
t.Fatalf("retry pending=%v, want %v", pending, tc.wantRetry)
|
||||
}
|
||||
if _, saved := store.etags[remote.Name]; saved == tc.wantRetry {
|
||||
t.Fatalf("etag saved=%v after first scan", saved)
|
||||
}
|
||||
if rows, _ := store.ListAlpineMetadataEntries(bg, remote.Name); len(rows) != 1 {
|
||||
t.Fatalf("healthy asset not served: %d rows", len(rows))
|
||||
}
|
||||
if !tc.wantRetry {
|
||||
return
|
||||
}
|
||||
|
||||
store.mu.Lock()
|
||||
store.retryAt[remote.Name] = time.Now().Add(-time.Second)
|
||||
store.mu.Unlock()
|
||||
s.process(bg, syncJob{remote: remote})
|
||||
if _, pending := store.retryAt[remote.Name]; pending || store.etags[remote.Name] != `"v1"` {
|
||||
t.Fatalf("retry not recorded: pending=%v etag=%q", pending, store.etags[remote.Name])
|
||||
}
|
||||
if rows, _ := store.ListAlpineMetadataEntries(bg, remote.Name); len(rows) != 2 {
|
||||
t.Fatalf("retry did not derive the failed asset: %d rows", len(rows))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -333,8 +333,8 @@ 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.
|
||||
// A transiently failed asset fails the scan so the ETag is not advanced and
|
||||
// the retry re-derives it; an invalid package is skipped as no retry fixes it.
|
||||
var failed []error
|
||||
seen := map[string]bool{}
|
||||
for _, rel := range releases {
|
||||
@@ -364,7 +364,9 @@ 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))
|
||||
if !errors.Is(err, provider.ErrInvalidPackage) {
|
||||
failed = append(failed, fmt.Errorf("%s: %w", asset.Name, err))
|
||||
}
|
||||
continue
|
||||
}
|
||||
if err := store.InsertDebMetadata(ctx, meta); err != nil {
|
||||
@@ -473,7 +475,7 @@ func (p *GitHubProvider) deriveAsset(ctx context.Context, remote models.Remote,
|
||||
Size: asset.Size,
|
||||
}
|
||||
if meta.Name == "" {
|
||||
return nil, errors.New("control missing Package field")
|
||||
return nil, fmt.Errorf("%w: control missing Package field", provider.ErrInvalidPackage)
|
||||
}
|
||||
|
||||
// The Packages SHA256 must be the sha256 of the whole .deb. Prefer GitHub's
|
||||
@@ -508,13 +510,13 @@ func (p *GitHubProvider) fetchControl(ctx context.Context, remote models.Remote,
|
||||
}
|
||||
control, complete, perr := controlFromPrefix(body)
|
||||
if perr != nil {
|
||||
return "", fmt.Errorf("parse deb control: %w", perr)
|
||||
return "", fmt.Errorf("%w: parse deb control: %w", provider.ErrInvalidPackage, perr)
|
||||
}
|
||||
if complete {
|
||||
return control, nil
|
||||
}
|
||||
if full || n >= p.headerMax {
|
||||
return "", fmt.Errorf("control.tar not found within %d bytes of %s", n, downloadURL)
|
||||
return "", fmt.Errorf("%w: control.tar not found within %d bytes of %s", provider.ErrInvalidPackage, n, downloadURL)
|
||||
}
|
||||
n *= 2
|
||||
if n > p.headerMax {
|
||||
|
||||
@@ -81,15 +81,17 @@ type githubFixture struct {
|
||||
notModHit int
|
||||
releaseAuth string
|
||||
assetAuth string
|
||||
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{
|
||||
debBytes: map[string][]byte{},
|
||||
rangeHit: map[string]int{},
|
||||
fullHit: map[string]int{},
|
||||
debBytes: map[string][]byte{},
|
||||
rangeHit: map[string]int{},
|
||||
fullHit: map[string]int{},
|
||||
assetFail: map[string]int{},
|
||||
}
|
||||
f.debBytes["demo_1.2-3_amd64.deb"] = testsupport.MinimalDeb("demo", "1.2-3", "amd64")
|
||||
|
||||
@@ -139,6 +141,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]++
|
||||
|
||||
@@ -308,3 +308,58 @@ func TestSyncerPrimeBypassesRecencyPeriodicDoesNot(t *testing.T) {
|
||||
t.Fatalf("periodic scan ran inside recency window: %d -> %d releases calls", releasesAfterPrime, fx.releasesHit)
|
||||
}
|
||||
}
|
||||
|
||||
// A transiently failed asset fails the scan: the ETag is not advanced, a retry
|
||||
// is scheduled, and the retry against the unchanged release list derives it. An
|
||||
// invalid package is skipped instead, since no retry can fix it.
|
||||
func TestSyncerFailedAssetRetry(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
corrupt bool
|
||||
wantRetry bool
|
||||
}{
|
||||
{name: "transient", wantRetry: true},
|
||||
{name: "invalid package", corrupt: true},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
fx := newGitHubFixture(t, true)
|
||||
fx.etag = `"v1"`
|
||||
if tc.corrupt {
|
||||
fx.debBytes["other_9-9_arm64.deb"] = []byte("not a package")
|
||||
} else {
|
||||
fx.debBytes["other_9-9_arm64.deb"] = testsupport.MinimalDeb("other", "9-9", "arm64")
|
||||
fx.assetFail["other_9-9_arm64.deb"] = 1
|
||||
}
|
||||
store := newFakeSyncStore()
|
||||
s := newSyncer(store, newTestProvider(), testSyncConfig())
|
||||
remote := fx.remote()
|
||||
bg := context.Background()
|
||||
|
||||
s.process(bg, syncJob{remote: remote, prime: true})
|
||||
_, pending := store.retryAt[remote.Name]
|
||||
if pending != tc.wantRetry {
|
||||
t.Fatalf("retry pending=%v, want %v", pending, tc.wantRetry)
|
||||
}
|
||||
if _, saved := store.etags[remote.Name]; saved == tc.wantRetry {
|
||||
t.Fatalf("etag saved=%v after first scan", saved)
|
||||
}
|
||||
if rows, _ := store.ListDebMetadataEntries(bg, remote.Name); len(rows) != 1 {
|
||||
t.Fatalf("healthy asset not served: %d rows", len(rows))
|
||||
}
|
||||
if !tc.wantRetry {
|
||||
return
|
||||
}
|
||||
|
||||
store.mu.Lock()
|
||||
store.retryAt[remote.Name] = time.Now().Add(-time.Second)
|
||||
store.mu.Unlock()
|
||||
s.process(bg, syncJob{remote: remote})
|
||||
if _, pending := store.retryAt[remote.Name]; pending || store.etags[remote.Name] != `"v1"` {
|
||||
t.Fatalf("retry not recorded: pending=%v etag=%q", pending, store.etags[remote.Name])
|
||||
}
|
||||
if rows, _ := store.ListDebMetadataEntries(bg, remote.Name); len(rows) != 2 {
|
||||
t.Fatalf("retry did not derive the failed asset: %d rows", len(rows))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -343,8 +343,8 @@ 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.
|
||||
// A transiently failed asset fails the scan so the ETag is not advanced and
|
||||
// the retry re-derives it; an invalid package is skipped as no retry fixes it.
|
||||
var failed []error
|
||||
seen := map[string]bool{}
|
||||
for _, rel := range releases {
|
||||
@@ -376,7 +376,9 @@ 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))
|
||||
if !errors.Is(err, provider.ErrInvalidPackage) {
|
||||
failed = append(failed, fmt.Errorf("%s: %w", asset.Name, err))
|
||||
}
|
||||
continue
|
||||
}
|
||||
if err := store.InsertRPMMetadata(ctx, meta); err != nil {
|
||||
@@ -574,7 +576,7 @@ func (p *GitHubProvider) fetchHeader(ctx context.Context, remote models.Remote,
|
||||
}
|
||||
continue
|
||||
}
|
||||
return nil, fmt.Errorf("parse rpm header: %w", perr)
|
||||
return nil, fmt.Errorf("%w: parse rpm header: %w", provider.ErrInvalidPackage, perr)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -400,42 +400,57 @@ func TestSyncerFailedScanRetriesAfterBackoff(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// 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()
|
||||
// A transiently failed asset fails the scan: the ETag is not advanced, a retry
|
||||
// is scheduled, and the retry against the unchanged release list derives it. An
|
||||
// invalid package is skipped instead, since no retry can fix it.
|
||||
func TestSyncerFailedAssetRetry(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
corrupt bool
|
||||
wantRetry bool
|
||||
}{
|
||||
{name: "transient", wantRetry: true},
|
||||
{name: "invalid package", corrupt: true},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
fx := newGitHubFixture(t, true)
|
||||
fx.etag = `"v1"`
|
||||
if tc.corrupt {
|
||||
fx.rpmBytes["other-9-9.aarch64.rpm"] = []byte("not a package")
|
||||
} else {
|
||||
fx.rpmBytes["other-9-9.aarch64.rpm"] = testsupport.MinimalRPM("other", "9", "9", "aarch64")
|
||||
fx.assetFail["other-9-9.aarch64.rpm"] = 1
|
||||
}
|
||||
store := newFakeSyncStore()
|
||||
s := newSyncer(store, newTestProvider(), 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))
|
||||
}
|
||||
s.process(bg, syncJob{remote: remote, prime: true})
|
||||
_, pending := store.retryAt[remote.Name]
|
||||
if pending != tc.wantRetry {
|
||||
t.Fatalf("retry pending=%v, want %v", pending, tc.wantRetry)
|
||||
}
|
||||
if _, saved := store.etags[remote.Name]; saved == tc.wantRetry {
|
||||
t.Fatalf("etag saved=%v after first scan", saved)
|
||||
}
|
||||
if rows, _ := store.ListRPMMetadataEntries(bg, remote.Name); len(rows) != 1 {
|
||||
t.Fatalf("healthy asset not served: %d rows", len(rows))
|
||||
}
|
||||
if !tc.wantRetry {
|
||||
return
|
||||
}
|
||||
|
||||
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))
|
||||
store.mu.Lock()
|
||||
store.retryAt[remote.Name] = time.Now().Add(-time.Second)
|
||||
store.mu.Unlock()
|
||||
s.process(bg, syncJob{remote: remote})
|
||||
if _, pending := store.retryAt[remote.Name]; pending || store.etags[remote.Name] != `"v1"` {
|
||||
t.Fatalf("retry not recorded: pending=%v etag=%q", pending, 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))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -43,6 +43,10 @@ func NewUpstreamStatusError(url string, resp *http.Response) *UpstreamStatusErro
|
||||
return e
|
||||
}
|
||||
|
||||
// ErrInvalidPackage marks an asset whose bytes are not a valid package. It is
|
||||
// permanent, so a scan skips the asset instead of failing and retrying.
|
||||
var ErrInvalidPackage = errors.New("invalid package")
|
||||
|
||||
// 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
|
||||
|
||||
Reference in New Issue
Block a user