From b111c2e57fac1c7b484c5a7011e580cf6afff066 Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Fri, 9 Oct 2026 23:28:58 +1100 Subject: [PATCH 1/4] Fail GitHub scans on asset errors and drop sync state with its remote - fail the scan when a release asset fails to derive or insert, so the ETag is kept and the retry re-derives it - cascade sync-state rows on remote delete (migration 0003) --- internal/database/github_sync_test.go | 69 +++++++++++++++++++ internal/provider/alpine/github.go | 7 +- internal/provider/deb/github.go | 7 +- internal/provider/rpm/github.go | 7 +- internal/provider/rpm/github_test.go | 14 +++- internal/provider/rpm/syncer_test.go | 40 +++++++++++ migrations/0003_github_sync_state_cascade.sql | 12 ++++ 7 files changed, 150 insertions(+), 6 deletions(-) create mode 100644 migrations/0003_github_sync_state_cascade.sql diff --git a/internal/database/github_sync_test.go b/internal/database/github_sync_test.go index c97449b..d66dc07 100644 --- a/internal/database/github_sync_test.go +++ b/internal/database/github_sync_test.go @@ -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) + } + }) + } +} diff --git a/internal/provider/alpine/github.go b/internal/provider/alpine/github.go index b60600e..93eb06f 100644 --- a/internal/provider/alpine/github.go +++ b/internal/provider/alpine/github.go @@ -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 { diff --git a/internal/provider/deb/github.go b/internal/provider/deb/github.go index 8731db3..708db2e 100644 --- a/internal/provider/deb/github.go +++ b/internal/provider/deb/github.go @@ -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 { diff --git a/internal/provider/rpm/github.go b/internal/provider/rpm/github.go index c41d47d..ff66d51 100644 --- a/internal/provider/rpm/github.go +++ b/internal/provider/rpm/github.go @@ -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 { diff --git a/internal/provider/rpm/github_test.go b/internal/provider/rpm/github_test.go index 33361f5..fccb666 100644 --- a/internal/provider/rpm/github_test.go +++ b/internal/provider/rpm/github_test.go @@ -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]++ diff --git a/internal/provider/rpm/syncer_test.go b/internal/provider/rpm/syncer_test.go index e368bd2..eb8877b 100644 --- a/internal/provider/rpm/syncer_test.go +++ b/internal/provider/rpm/syncer_test.go @@ -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)) + } +} diff --git a/migrations/0003_github_sync_state_cascade.sql b/migrations/0003_github_sync_state_cascade.sql new file mode 100644 index 0000000..90b9a1a --- /dev/null +++ b/migrations/0003_github_sync_state_cascade.sql @@ -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; -- 2.47.3 From d9f9fd3bca09c00a0a801ed7de1614851585ae91 Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Fri, 9 Oct 2026 23:47:31 +1100 Subject: [PATCH 2/4] 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. --- internal/provider/alpine/github.go | 14 ++-- internal/provider/alpine/github_test.go | 14 +++- internal/provider/alpine/syncer_test.go | 55 ++++++++++++++++ internal/provider/deb/github.go | 14 ++-- internal/provider/deb/github_test.go | 14 +++- internal/provider/deb/syncer_test.go | 55 ++++++++++++++++ internal/provider/rpm/github.go | 10 +-- internal/provider/rpm/syncer_test.go | 85 +++++++++++++++---------- internal/provider/syncretry.go | 4 ++ 9 files changed, 208 insertions(+), 57 deletions(-) diff --git a/internal/provider/alpine/github.go b/internal/provider/alpine/github.go index 93eb06f..9461166 100644 --- a/internal/provider/alpine/github.go +++ b/internal/provider/alpine/github.go @@ -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 { diff --git a/internal/provider/alpine/github_test.go b/internal/provider/alpine/github_test.go index 9f33fb0..53dd777 100644 --- a/internal/provider/alpine/github_test.go +++ b/internal/provider/alpine/github_test.go @@ -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]++ diff --git a/internal/provider/alpine/syncer_test.go b/internal/provider/alpine/syncer_test.go index 0401c52..1a6ce96 100644 --- a/internal/provider/alpine/syncer_test.go +++ b/internal/provider/alpine/syncer_test.go @@ -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)) + } + }) + } +} diff --git a/internal/provider/deb/github.go b/internal/provider/deb/github.go index 708db2e..aa2305a 100644 --- a/internal/provider/deb/github.go +++ b/internal/provider/deb/github.go @@ -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 { diff --git a/internal/provider/deb/github_test.go b/internal/provider/deb/github_test.go index f27fbc7..519a8c8 100644 --- a/internal/provider/deb/github_test.go +++ b/internal/provider/deb/github_test.go @@ -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]++ diff --git a/internal/provider/deb/syncer_test.go b/internal/provider/deb/syncer_test.go index e3fe662..b3919a5 100644 --- a/internal/provider/deb/syncer_test.go +++ b/internal/provider/deb/syncer_test.go @@ -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)) + } + }) + } +} diff --git a/internal/provider/rpm/github.go b/internal/provider/rpm/github.go index ff66d51..e48d4f0 100644 --- a/internal/provider/rpm/github.go +++ b/internal/provider/rpm/github.go @@ -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) } } diff --git a/internal/provider/rpm/syncer_test.go b/internal/provider/rpm/syncer_test.go index eb8877b..a01c785 100644 --- a/internal/provider/rpm/syncer_test.go +++ b/internal/provider/rpm/syncer_test.go @@ -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)) + } + }) } } diff --git a/internal/provider/syncretry.go b/internal/provider/syncretry.go index dff7070..27dd4ba 100644 --- a/internal/provider/syncretry.go +++ b/internal/provider/syncretry.go @@ -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 -- 2.47.3 From 09ab49ef424aaf911954da4783e14300a96d03e1 Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Fri, 9 Oct 2026 23:57:35 +1100 Subject: [PATCH 3/4] Skip missing GitHub assets and treat short reads as transient --- internal/provider/alpine/github.go | 12 ++---- internal/provider/alpine/github_test.go | 11 +++++ internal/provider/deb/github.go | 14 ++----- internal/provider/deb/github_test.go | 11 +++++ internal/provider/rpm/github.go | 14 ++----- internal/provider/rpm/github_test.go | 30 ++++++++++--- internal/provider/rpm/syncer_test.go | 28 ++++++++++--- internal/provider/syncretry.go | 38 +++++++++++++++++ internal/provider/syncretry_test.go | 56 +++++++++++++++++++++++++ 9 files changed, 174 insertions(+), 40 deletions(-) diff --git a/internal/provider/alpine/github.go b/internal/provider/alpine/github.go index 9461166..46aad38 100644 --- a/internal/provider/alpine/github.go +++ b/internal/provider/alpine/github.go @@ -601,7 +601,7 @@ func pkginfoFromPrefix(prefix []byte) (meta *provider.AlpineMetadata, complete b } // rangeGet returns the first n bytes of downloadURL. full is true when the -// response body was shorter than n (i.e. we already have the whole object). +// response holds the whole object. func (p *GitHubProvider) rangeGet(ctx context.Context, remote models.Remote, downloadURL string, n int64) ([]byte, bool, error) { req, err := http.NewRequestWithContext(ctx, http.MethodGet, downloadURL, nil) if err != nil { @@ -623,15 +623,9 @@ func (p *GitHubProvider) rangeGet(ctx context.Context, remote models.Remote, dow } defer resp.Body.Close() if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusPartialContent { - return nil, false, fmt.Errorf("range GET %s: status %d", downloadURL, resp.StatusCode) + return nil, false, fmt.Errorf("range GET: %w", provider.AssetStatusError(downloadURL, resp)) } - - body, err := io.ReadAll(io.LimitReader(resp.Body, n)) - if err != nil { - return nil, false, err - } - full := int64(len(body)) < n - return body, full, nil + return provider.ReadPrefix(resp, n) } // assetPath is the package's location relative to github.com — the path the diff --git a/internal/provider/alpine/github_test.go b/internal/provider/alpine/github_test.go index 53dd777..69f27f7 100644 --- a/internal/provider/alpine/github_test.go +++ b/internal/provider/alpine/github_test.go @@ -8,6 +8,7 @@ import ( "crypto/sha256" "encoding/hex" "encoding/json" + "errors" "fmt" "io" "net/http" @@ -503,3 +504,13 @@ func readAPKIndex(t *testing.T, gzBytes []byte) string { } } } + +func TestGitHubDeriveAssetMissingPkgnameIsInvalid(t *testing.T) { + fx := newGitHubFixture(t, true) + fx.apkBytes["nameless-1-r0.apk"] = testsupport.MinimalApk("", "1-r0", "x86_64") + asset := ghAsset{BrowserDownloadURL: fx.srv.URL + "/acme/tools/releases/download/v1.2.3/nameless-1-r0.apk"} + _, err := newTestProvider().deriveAsset(context.Background(), fx.remote(), asset, "fp") + if !errors.Is(err, provider.ErrInvalidPackage) { + t.Fatalf("want ErrInvalidPackage, got %v", err) + } +} diff --git a/internal/provider/deb/github.go b/internal/provider/deb/github.go index aa2305a..83c85de 100644 --- a/internal/provider/deb/github.go +++ b/internal/provider/deb/github.go @@ -574,7 +574,7 @@ func controlFromPrefix(prefix []byte) (control string, complete bool, err error) } // rangeGet returns the first n bytes of downloadURL. full is true when the -// response body was shorter than n (i.e. we already have the whole object). +// response holds the whole object. func (p *GitHubProvider) rangeGet(ctx context.Context, remote models.Remote, downloadURL string, n int64) ([]byte, bool, error) { req, err := http.NewRequestWithContext(ctx, http.MethodGet, downloadURL, nil) if err != nil { @@ -596,15 +596,9 @@ func (p *GitHubProvider) rangeGet(ctx context.Context, remote models.Remote, dow } defer resp.Body.Close() if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusPartialContent { - return nil, false, fmt.Errorf("range GET %s: status %d", downloadURL, resp.StatusCode) + return nil, false, fmt.Errorf("range GET: %w", provider.AssetStatusError(downloadURL, resp)) } - - body, err := io.ReadAll(io.LimitReader(resp.Body, n)) - if err != nil { - return nil, false, err - } - full := int64(len(body)) < n - return body, full, nil + return provider.ReadPrefix(resp, n) } func (p *GitHubProvider) computeSHA256(ctx context.Context, remote models.Remote, downloadURL string) (string, error) { @@ -627,7 +621,7 @@ func (p *GitHubProvider) computeSHA256(ctx context.Context, remote models.Remote } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { - return "", fmt.Errorf("GET %s: status %d", downloadURL, resp.StatusCode) + return "", fmt.Errorf("GET: %w", provider.AssetStatusError(downloadURL, resp)) } h := sha256.New() diff --git a/internal/provider/deb/github_test.go b/internal/provider/deb/github_test.go index 519a8c8..b03956d 100644 --- a/internal/provider/deb/github_test.go +++ b/internal/provider/deb/github_test.go @@ -7,6 +7,7 @@ import ( "crypto/sha256" "encoding/hex" "encoding/json" + "errors" "fmt" "io" "net/http" @@ -468,3 +469,13 @@ func TestGitHubAssetPatternFilter(t *testing.T) { t.Fatalf("pattern filter failed, rows=%+v", rows) } } + +func TestGitHubDeriveAssetMissingPackageIsInvalid(t *testing.T) { + fx := newGitHubFixture(t, true) + fx.debBytes["nameless_1_amd64.deb"] = testsupport.MinimalDeb("", "1", "amd64") + asset := ghAsset{BrowserDownloadURL: fx.srv.URL + "/acme/tools/releases/download/v1.2-3/nameless_1_amd64.deb", Digest: "sha256:00"} + _, err := newTestProvider().deriveAsset(context.Background(), fx.remote(), asset, "fp") + if !errors.Is(err, provider.ErrInvalidPackage) { + t.Fatalf("want ErrInvalidPackage, got %v", err) + } +} diff --git a/internal/provider/rpm/github.go b/internal/provider/rpm/github.go index e48d4f0..707b54d 100644 --- a/internal/provider/rpm/github.go +++ b/internal/provider/rpm/github.go @@ -581,7 +581,7 @@ func (p *GitHubProvider) fetchHeader(ctx context.Context, remote models.Remote, } // rangeGet returns the first n bytes of downloadURL. full is true when the -// response body was shorter than n (i.e. we already have the whole object). +// response holds the whole object. func (p *GitHubProvider) rangeGet(ctx context.Context, remote models.Remote, downloadURL string, n int64) ([]byte, bool, error) { req, err := http.NewRequestWithContext(ctx, http.MethodGet, downloadURL, nil) if err != nil { @@ -603,15 +603,9 @@ func (p *GitHubProvider) rangeGet(ctx context.Context, remote models.Remote, dow } defer resp.Body.Close() if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusPartialContent { - return nil, false, fmt.Errorf("range GET %s: status %d", downloadURL, resp.StatusCode) + return nil, false, fmt.Errorf("range GET: %w", provider.AssetStatusError(downloadURL, resp)) } - - body, err := io.ReadAll(io.LimitReader(resp.Body, n)) - if err != nil { - return nil, false, err - } - full := int64(len(body)) < n - return body, full, nil + return provider.ReadPrefix(resp, n) } func (p *GitHubProvider) computeSHA256(ctx context.Context, remote models.Remote, downloadURL string) (string, error) { @@ -634,7 +628,7 @@ func (p *GitHubProvider) computeSHA256(ctx context.Context, remote models.Remote } defer resp.Body.Close() if resp.StatusCode != http.StatusOK { - return "", fmt.Errorf("GET %s: status %d", downloadURL, resp.StatusCode) + return "", fmt.Errorf("GET: %w", provider.AssetStatusError(downloadURL, resp)) } h := sha256.New() diff --git a/internal/provider/rpm/github_test.go b/internal/provider/rpm/github_test.go index fccb666..6298c4d 100644 --- a/internal/provider/rpm/github_test.go +++ b/internal/provider/rpm/github_test.go @@ -1,11 +1,13 @@ package rpm import ( + "cmp" "compress/gzip" "context" "crypto/sha256" "encoding/hex" "encoding/json" + "errors" "fmt" "io" "net/http" @@ -24,15 +26,22 @@ import ( // fakeStore is an in-memory provider.RemoteMetadataStore keyed by file_path, // mirroring the (repo_name, file_path) uniqueness of the real table. type fakeStore struct { - mu sync.Mutex - rows map[string]provider.RPMMetadata + mu sync.Mutex + rows map[string]provider.RPMMetadata + insertFail map[string]int // package name -> inserts left to fail } -func newFakeStore() *fakeStore { return &fakeStore{rows: map[string]provider.RPMMetadata{}} } +func newFakeStore() *fakeStore { + return &fakeStore{rows: map[string]provider.RPMMetadata{}, insertFail: map[string]int{}} +} func (f *fakeStore) InsertRPMMetadata(_ context.Context, m *provider.RPMMetadata) error { f.mu.Lock() defer f.mu.Unlock() + if f.insertFail[m.Name] > 0 { + f.insertFail[m.Name]-- + return errors.New("insert failed") + } if _, ok := f.rows[m.FilePath]; ok { return nil // ON CONFLICT DO NOTHING } @@ -79,7 +88,9 @@ 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 + assetFail map[string]int // asset filename -> downloads left to answer failCode + failCode int // status for assetFail downloads; 0 = 500 + shortRead map[string]int // asset filename -> ranged GETs left to close early mu sync.Mutex } @@ -90,6 +101,7 @@ func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture { rangeHit: map[string]int{}, fullHit: map[string]int{}, assetFail: map[string]int{}, + shortRead: map[string]int{}, } f.rpmBytes["demo-1.2-3.x86_64.rpm"] = testsupport.MinimalRPM("demo", "1.2", "3", "x86_64") @@ -149,10 +161,15 @@ func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture { f.mu.Lock() if f.assetFail[name] > 0 { f.assetFail[name]-- + code := cmp.Or(f.failCode, http.StatusInternalServerError) f.mu.Unlock() - http.Error(w, "boom", http.StatusInternalServerError) + http.Error(w, "boom", code) return } + short := rng != "" && f.shortRead[name] > 0 + if short { + f.shortRead[name]-- + } f.assetAuth = r.Header.Get("Authorization") if rng != "" { f.rangeHit[name]++ @@ -173,6 +190,9 @@ func newGitHubFixture(t *testing.T, withDigest bool) *githubFixture { end = len(body) - 1 } w.Header().Set("Content-Range", fmt.Sprintf("bytes 0-%d/%d", end, len(body))) + if short { + end /= 2 + } w.Header().Set("Content-Length", strconv.Itoa(end+1)) w.WriteHeader(http.StatusPartialContent) w.Write(body[:end+1]) diff --git a/internal/provider/rpm/syncer_test.go b/internal/provider/rpm/syncer_test.go index a01c785..be6d9c9 100644 --- a/internal/provider/rpm/syncer_test.go +++ b/internal/provider/rpm/syncer_test.go @@ -407,21 +407,37 @@ func TestSyncerFailedAssetRetry(t *testing.T) { for _, tc := range []struct { name string corrupt bool + failCode int + short bool + insert bool wantRetry bool }{ {name: "transient", wantRetry: true}, + {name: "rate limited", failCode: http.StatusForbidden, wantRetry: true}, + {name: "short read", short: true, wantRetry: true}, + {name: "insert failure", insert: true, wantRetry: true}, {name: "invalid package", corrupt: true}, + {name: "asset not found", failCode: http.StatusNotFound}, + {name: "asset gone", failCode: http.StatusGone}, } { t.Run(tc.name, func(t *testing.T) { + const other = "other-9-9.aarch64.rpm" 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 - } + fx.rpmBytes[other] = testsupport.MinimalRPM("other", "9", "9", "aarch64") store := newFakeSyncStore() + switch { + case tc.corrupt: + fx.rpmBytes[other] = []byte("not a package") + case tc.short: + fx.shortRead[other] = 1 + case tc.insert: + store.insertFail["other"] = 1 + case tc.failCode != 0 && !tc.wantRetry: + fx.failCode, fx.assetFail[other] = tc.failCode, 1<<30 + default: + fx.failCode, fx.assetFail[other] = tc.failCode, 1 + } s := newSyncer(store, newTestProvider(), testSyncConfig()) remote := fx.remote() bg := context.Background() diff --git a/internal/provider/syncretry.go b/internal/provider/syncretry.go index 27dd4ba..cec262e 100644 --- a/internal/provider/syncretry.go +++ b/internal/provider/syncretry.go @@ -3,8 +3,10 @@ package provider import ( "errors" "fmt" + "io" "net/http" "strconv" + "strings" "time" ) @@ -47,6 +49,42 @@ func NewUpstreamStatusError(url string, resp *http.Response) *UpstreamStatusErro // permanent, so a scan skips the asset instead of failing and retrying. var ErrInvalidPackage = errors.New("invalid package") +// AssetStatusError classifies a non-success asset download. A missing or +// unsatisfiable asset is permanent; anything else (403/429 rate limits, 5xx) is +// transient. +func AssetStatusError(url string, resp *http.Response) error { + switch resp.StatusCode { + case http.StatusNotFound, http.StatusGone, http.StatusRequestedRangeNotSatisfiable: + return fmt.Errorf("%w: %s: status %d", ErrInvalidPackage, url, resp.StatusCode) + } + return NewUpstreamStatusError(url, resp) +} + +// ReadPrefix reads up to n bytes of a 200 or 206 asset response. full reports +// that the body is the whole object (per Content-Length or the Content-Range +// total); a body ending before n bytes that is not the whole object is a +// transient short read. +func ReadPrefix(resp *http.Response, n int64) (body []byte, full bool, err error) { + body, err = io.ReadAll(io.LimitReader(resp.Body, n)) + if err != nil { + return nil, false, err + } + total := resp.ContentLength + if resp.StatusCode == http.StatusPartialContent { + total = -1 + if _, t, ok := strings.Cut(resp.Header.Get("Content-Range"), "/"); ok { + if v, perr := strconv.ParseInt(t, 10, 64); perr == nil { + total = v + } + } + } + full = total == int64(len(body)) + if !full && int64(len(body)) < n { + return nil, false, fmt.Errorf("short read: got %d of %d bytes", len(body), n) + } + return body, full, nil +} + // 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 diff --git a/internal/provider/syncretry_test.go b/internal/provider/syncretry_test.go index 289a1d7..4feab63 100644 --- a/internal/provider/syncretry_test.go +++ b/internal/provider/syncretry_test.go @@ -3,8 +3,10 @@ package provider import ( "errors" "fmt" + "io" "net/http" "strconv" + "strings" "testing" "time" ) @@ -61,3 +63,57 @@ func TestNewSyncResult(t *testing.T) { t.Fatalf("wrapped hint not clamped to ttl: retry in %v", until) } } + +func TestAssetStatusErrorClassification(t *testing.T) { + for status, permanent := range map[int]bool{ + http.StatusNotFound: true, http.StatusGone: true, http.StatusRequestedRangeNotSatisfiable: true, + http.StatusForbidden: false, http.StatusTooManyRequests: false, http.StatusInternalServerError: false, http.StatusBadGateway: false, + } { + err := AssetStatusError("u", &http.Response{StatusCode: status, Header: http.Header{}}) + if errors.Is(err, ErrInvalidPackage) != permanent { + t.Errorf("status %d: permanent=%v, want %v", status, !permanent, permanent) + } + } +} + +func TestReadPrefix(t *testing.T) { + cases := []struct { + name string + status int + cl int64 + crange string + body string + n int64 + wantFull bool + wantErr bool + }{ + {name: "200 whole object", status: 200, cl: 5, body: "hello", n: 32, wantFull: true}, + {name: "200 larger than n", status: 200, cl: 100, body: "hello", n: 4}, + {name: "200 unknown length short", status: 200, cl: -1, body: "hello", n: 32, wantErr: true}, + {name: "206 whole object", status: 206, crange: "bytes 0-4/5", body: "hello", n: 32, wantFull: true}, + {name: "206 prefix", status: 206, crange: "bytes 0-3/100", body: "hell", n: 4}, + {name: "206 short read", status: 206, crange: "bytes 0-31/100", body: "hello", n: 32, wantErr: true}, + {name: "206 no content-range short", status: 206, body: "hello", n: 32, wantErr: true}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + resp := &http.Response{StatusCode: tc.status, ContentLength: tc.cl, Header: http.Header{}, Body: io.NopCloser(strings.NewReader(tc.body))} + if tc.crange != "" { + resp.Header.Set("Content-Range", tc.crange) + } + body, full, err := ReadPrefix(resp, tc.n) + if (err != nil) != tc.wantErr { + t.Fatalf("err=%v, wantErr=%v", err, tc.wantErr) + } + if err != nil { + if errors.Is(err, ErrInvalidPackage) { + t.Fatalf("short read must be transient: %v", err) + } + return + } + if full != tc.wantFull || int64(len(body)) > tc.n { + t.Fatalf("full=%v len=%d, want full=%v", full, len(body), tc.wantFull) + } + }) + } +} -- 2.47.3 From fdbf322ca8c67f5bf41745732e3cca1de4374215 Mon Sep 17 00:00:00 2001 From: unkin-agent Date: Sat, 10 Oct 2026 01:10:07 +1100 Subject: [PATCH 4/4] Treat a clean EOF with unknown total as the whole asset --- internal/provider/syncretry.go | 7 +++++-- internal/provider/syncretry_test.go | 31 +++++++++++++++++++++++------ 2 files changed, 30 insertions(+), 8 deletions(-) diff --git a/internal/provider/syncretry.go b/internal/provider/syncretry.go index cec262e..1a2eb47 100644 --- a/internal/provider/syncretry.go +++ b/internal/provider/syncretry.go @@ -62,7 +62,7 @@ func AssetStatusError(url string, resp *http.Response) error { // ReadPrefix reads up to n bytes of a 200 or 206 asset response. full reports // that the body is the whole object (per Content-Length or the Content-Range -// total); a body ending before n bytes that is not the whole object is a +// total, or a clean EOF before n when the total is unknown); a body ending before n bytes that is not the whole object is a // transient short read. func ReadPrefix(resp *http.Response, n int64) (body []byte, full bool, err error) { body, err = io.ReadAll(io.LimitReader(resp.Body, n)) @@ -70,15 +70,18 @@ func ReadPrefix(resp *http.Response, n int64) (body []byte, full bool, err error return nil, false, err } total := resp.ContentLength + unknown := resp.StatusCode == http.StatusOK && total == -1 if resp.StatusCode == http.StatusPartialContent { total = -1 if _, t, ok := strings.Cut(resp.Header.Get("Content-Range"), "/"); ok { if v, perr := strconv.ParseInt(t, 10, 64); perr == nil { total = v } + unknown = t == "*" } } - full = total == int64(len(body)) + // net/http errors on a premature close, so a clean EOF before n with an unknown total is the whole object. + full = total == int64(len(body)) || (unknown && int64(len(body)) < n) if !full && int64(len(body)) < n { return nil, false, fmt.Errorf("short read: got %d of %d bytes", len(body), n) } diff --git a/internal/provider/syncretry_test.go b/internal/provider/syncretry_test.go index 4feab63..906b9e3 100644 --- a/internal/provider/syncretry_test.go +++ b/internal/provider/syncretry_test.go @@ -5,8 +5,8 @@ import ( "fmt" "io" "net/http" + "net/http/httptest" "strconv" - "strings" "testing" "time" ) @@ -88,18 +88,37 @@ func TestReadPrefix(t *testing.T) { wantErr bool }{ {name: "200 whole object", status: 200, cl: 5, body: "hello", n: 32, wantFull: true}, - {name: "200 larger than n", status: 200, cl: 100, body: "hello", n: 4}, - {name: "200 unknown length short", status: 200, cl: -1, body: "hello", n: 32, wantErr: true}, + {name: "200 larger than n", status: 200, cl: 5, body: "hello", n: 4}, + {name: "200 chunked shorter than n", status: 200, cl: -1, body: "hello", n: 32, wantFull: true}, + {name: "200 chunked larger than n", status: 200, cl: -1, body: "hello", n: 4}, {name: "206 whole object", status: 206, crange: "bytes 0-4/5", body: "hello", n: 32, wantFull: true}, {name: "206 prefix", status: 206, crange: "bytes 0-3/100", body: "hell", n: 4}, {name: "206 short read", status: 206, crange: "bytes 0-31/100", body: "hello", n: 32, wantErr: true}, {name: "206 no content-range short", status: 206, body: "hello", n: 32, wantErr: true}, + {name: "206 unknown total shorter than n", status: 206, cl: -1, crange: "bytes 0-4/*", body: "hello", n: 32, wantFull: true}, + {name: "206 unknown total prefix", status: 206, cl: -1, crange: "bytes 0-3/*", body: "hell", n: 4}, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { - resp := &http.Response{StatusCode: tc.status, ContentLength: tc.cl, Header: http.Header{}, Body: io.NopCloser(strings.NewReader(tc.body))} - if tc.crange != "" { - resp.Header.Set("Content-Range", tc.crange) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + if tc.cl > 0 { + w.Header().Set("Content-Length", strconv.FormatInt(tc.cl, 10)) + } + if tc.crange != "" { + w.Header().Set("Content-Range", tc.crange) + } + w.WriteHeader(tc.status) + _, _ = io.WriteString(w, tc.body) + w.(http.Flusher).Flush() + })) + defer srv.Close() + resp, err := http.Get(srv.URL) + if err != nil { + t.Fatal(err) + } + defer resp.Body.Close() + if tc.cl == -1 && resp.ContentLength != -1 { + t.Fatalf("response not chunked: ContentLength=%d", resp.ContentLength) } body, full, err := ReadPrefix(resp, tc.n) if (err != nil) != tc.wantErr { -- 2.47.3