Skip to content

Commit bf958fa

Browse files
Gustedmfenniak
authored andcommitted
fix: make package cleanup work again (#12446)
- Regression of forgejo/forgejo!11776 (and forgejo/forgejo!11881) - Scope of the transaction is moved to a per-package cleanup rule basis. This is also a enhancement for scaling (already deployed on Codeberg for a while). - Package cleanup is now run with `RetryTx`, because rebuilding repository files runs `RetryTx` and it could indicate to retry the whole transaction. - Previously it would error and say running `RetryTx` in a transaction was not possible, this is now possible. Nested `RetryTx` is always allowed, matching of which errors to retry is still the responsible of the inner `RetryTx`. Co-authored-by: Mathieu Fenniak <mathieu@fenniak.net> Reviewed-on: https://codeberg.org/forgejo/forgejo/pulls/12446 Reviewed-by: Mathieu Fenniak <mfenniak@noreply.codeberg.org>
1 parent 69cf1f3 commit bf958fa

4 files changed

Lines changed: 171 additions & 75 deletions

File tree

models/db/context.go

Lines changed: 38 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -510,31 +510,57 @@ type RetryConfig struct {
510510
AttemptCount int
511511
}
512512

513+
var ErrNestedRetryTxFailure = errors.New("(nested)")
514+
515+
type nestedRetryTxState int
516+
517+
var nestedRetryTx nestedRetryTxState
518+
513519
// Execute the given function in a transaction. RetryConfig will retry the function on an error, if it matches the
514520
// ErrorIs parameter, up to the total of AttemptCount number of tries. RetryTx cannot be invoked when already within a
515521
// transaction and will return an error immediately.
522+
//
523+
// ErrNestedRetryTxFailure is an error type that will occur when RetryTx is nested within each other, and indicates that
524+
// an inner RetryTx encountered an error that matched its error list.
516525
func RetryTx(ctx context.Context, config RetryConfig, f func(ctx context.Context) error) error {
517-
if InTransaction(ctx) {
526+
matchError := func(err error) bool {
527+
for _, possibleError := range config.ErrorIs {
528+
if errors.Is(err, possibleError) {
529+
return true
530+
}
531+
}
532+
return false
533+
}
534+
535+
// Accept `ErrNestedRetryTxFailure` as error to retry on, means that a nested
536+
// RetryTx indicated to retry the whole transaction.
537+
config.ErrorIs = append(config.ErrorIs, ErrNestedRetryTxFailure)
538+
539+
withinRetryTx, present := ctx.Value(nestedRetryTx).(bool)
540+
if present && withinRetryTx {
541+
// If a caller already started `RetryTx`, then we assume we don't have to actually perform retries here -- we
542+
// can attempt the requested function once, and if an error is returned that matches the configured error list,
543+
// we'll return that error + ErrNestedRetryTxFailure wrapping.
544+
err := f(ctx)
545+
if err == nil {
546+
return nil
547+
} else if matchError(err) {
548+
return fmt.Errorf("nested RetryTx; internal Tx failed with error that won't be retried: %w %w", err, ErrNestedRetryTxFailure)
549+
}
550+
return err
551+
} else if InTransaction(ctx) {
518552
return errors.New("unsupported operation: attempted to use RetryTx while already within a transaction")
519553
} else if config.AttemptCount == 0 {
520554
return errors.New("unsupported operation: attempted to use RetryTx with 0 attempts")
521555
}
522556

557+
innerCtx := context.WithValue(ctx, nestedRetryTx, true)
523558
var lastError error
524559
for range config.AttemptCount {
525-
err := WithTx(ctx, f)
560+
err := WithTx(innerCtx, f)
526561
if err == nil {
527562
return nil
528-
}
529-
530-
foundMatch := false
531-
for _, possibleError := range config.ErrorIs {
532-
if errors.Is(err, possibleError) {
533-
foundMatch = true
534-
break
535-
}
536-
}
537-
if !foundMatch {
563+
} else if !matchError(err) {
538564
return err
539565
}
540566

models/db/context_test.go

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -275,4 +275,48 @@ func TestRetryTx(t *testing.T) {
275275
require.NoError(t, err)
276276
assert.Equal(t, 2, attemptCount)
277277
})
278+
279+
t.Run("nested", func(t *testing.T) {
280+
attemptCount := 0
281+
testError := errors.New("hello")
282+
err := db.RetryTx(t.Context(), db.RetryConfig{
283+
AttemptCount: 2,
284+
}, func(ctx context.Context) error {
285+
attemptCount++
286+
return db.RetryTx(ctx, db.RetryConfig{
287+
AttemptCount: 2,
288+
ErrorIs: []error{testError},
289+
}, func(ctx context.Context) error {
290+
if attemptCount == 2 {
291+
return nil
292+
}
293+
return testError
294+
})
295+
})
296+
297+
require.NoError(t, err)
298+
assert.Equal(t, 2, attemptCount)
299+
})
300+
301+
t.Run("inner RetryTx decides on error", func(t *testing.T) {
302+
attemptCount := 0
303+
testError := errors.New("hello")
304+
err := db.RetryTx(t.Context(), db.RetryConfig{
305+
AttemptCount: 2,
306+
ErrorIs: []error{},
307+
}, func(ctx context.Context) error {
308+
attemptCount++
309+
return db.RetryTx(ctx, db.RetryConfig{
310+
AttemptCount: 2,
311+
}, func(ctx context.Context) error {
312+
if attemptCount == 2 {
313+
return nil
314+
}
315+
return testError
316+
})
317+
})
318+
319+
require.ErrorIs(t, err, testError)
320+
assert.Equal(t, 1, attemptCount)
321+
})
278322
}

services/packages/cleanup/cleanup.go

Lines changed: 53 additions & 63 deletions
Original file line numberDiff line numberDiff line change
@@ -33,78 +33,68 @@ func CleanupTask(ctx context.Context, olderThan time.Duration) error {
3333
return CleanupExpiredData(ctx, olderThan)
3434
}
3535

36-
func ExecuteCleanupRules(outerCtx context.Context) error {
37-
ctx, committer, err := db.TxContext(outerCtx)
38-
if err != nil {
39-
return err
40-
}
41-
defer committer.Close()
42-
43-
err = packages_model.IterateEnabledCleanupRules(ctx, func(ctx context.Context, pcr *packages_model.PackageCleanupRule) error {
44-
select {
45-
case <-outerCtx.Done():
46-
return db.ErrCancelledf("While processing package cleanup rules")
47-
default:
48-
}
49-
50-
versionsToRemove, err := GetCleanupTargets(ctx, pcr, true)
51-
if err != nil {
52-
return fmt.Errorf("CleanupRule [%d]: GetCleanupTargets failed: %w", pcr.ID, err)
53-
}
54-
55-
anyVersionDeleted := false
56-
packageWithVersionDeleted := make(map[int64]bool) // set of Package.ID's where at least one package version was removed
57-
for _, ct := range versionsToRemove {
58-
if err := packages_service.DeletePackageVersionAndReferences(ctx, ct.PackageVersion); err != nil {
59-
return fmt.Errorf("CleanupRule [%d]: DeletePackageVersionAndReferences failed: %w", pcr.ID, err)
36+
func ExecuteCleanupRules(ctx context.Context) error {
37+
return packages_model.IterateEnabledCleanupRules(ctx, func(ctx context.Context, pcr *packages_model.PackageCleanupRule) error {
38+
// We have no errors to retry on, because we have no evidence we need any in
39+
// this area. What we do retry on is when a nested `db.RetryTx` indicates to
40+
// retry the whole transaction.
41+
return db.RetryTx(ctx, db.RetryConfig{
42+
AttemptCount: 3,
43+
}, func(ctx context.Context) error {
44+
versionsToRemove, err := GetCleanupTargets(ctx, pcr, true)
45+
if err != nil {
46+
return fmt.Errorf("CleanupRule [%d]: GetCleanupTargets failed: %w", pcr.ID, err)
6047
}
61-
packageWithVersionDeleted[ct.Package.ID] = true
62-
anyVersionDeleted = true
63-
}
6448

65-
if pcr.Type == packages_model.TypeCargo {
66-
for packageID := range packageWithVersionDeleted {
67-
owner, err := user_model.GetUserByID(ctx, pcr.OwnerID)
68-
if err != nil {
69-
return fmt.Errorf("GetUserByID failed: %w", err)
70-
}
71-
if err := cargo_service.UpdatePackageIndexIfExists(ctx, owner, owner, packageID); err != nil {
72-
return fmt.Errorf("CleanupRule [%d]: cargo.UpdatePackageIndexIfExists failed: %w", pcr.ID, err)
49+
anyVersionDeleted := false
50+
packageWithVersionDeleted := make(map[int64]bool) // set of Package.ID's where at least one package version was removed
51+
for _, ct := range versionsToRemove {
52+
if err := packages_service.DeletePackageVersionAndReferences(ctx, ct.PackageVersion); err != nil {
53+
return fmt.Errorf("CleanupRule [%d]: DeletePackageVersionAndReferences failed: %w", pcr.ID, err)
7354
}
55+
packageWithVersionDeleted[ct.Package.ID] = true
56+
anyVersionDeleted = true
7457
}
75-
}
7658

77-
if anyVersionDeleted {
78-
switch pcr.Type {
79-
case packages_model.TypeDebian:
80-
if err := debian_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
81-
return fmt.Errorf("CleanupRule [%d]: debian.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
82-
}
83-
case packages_model.TypeAlpine:
84-
if err := alpine_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
85-
return fmt.Errorf("CleanupRule [%d]: alpine.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
86-
}
87-
case packages_model.TypeRpm:
88-
if err := rpm_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
89-
return fmt.Errorf("CleanupRule [%d]: rpm.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
59+
if pcr.Type == packages_model.TypeCargo {
60+
for packageID := range packageWithVersionDeleted {
61+
owner, err := user_model.GetUserByID(ctx, pcr.OwnerID)
62+
if err != nil {
63+
return fmt.Errorf("GetUserByID failed: %w", err)
64+
}
65+
if err := cargo_service.UpdatePackageIndexIfExists(ctx, owner, owner, packageID); err != nil {
66+
return fmt.Errorf("CleanupRule [%d]: cargo.UpdatePackageIndexIfExists failed: %w", pcr.ID, err)
67+
}
9068
}
91-
case packages_model.TypeArch:
92-
if err := arch_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
93-
return fmt.Errorf("CleanupRule [%d]: arch.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
94-
}
95-
case packages_model.TypeAlt:
96-
if err := alt_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
97-
return fmt.Errorf("CleanupRule [%d]: alt.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
69+
}
70+
71+
if anyVersionDeleted {
72+
switch pcr.Type {
73+
case packages_model.TypeDebian:
74+
if err := debian_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
75+
return fmt.Errorf("CleanupRule [%d]: debian.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
76+
}
77+
case packages_model.TypeAlpine:
78+
if err := alpine_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
79+
return fmt.Errorf("CleanupRule [%d]: alpine.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
80+
}
81+
case packages_model.TypeRpm:
82+
if err := rpm_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
83+
return fmt.Errorf("CleanupRule [%d]: rpm.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
84+
}
85+
case packages_model.TypeArch:
86+
if err := arch_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
87+
return fmt.Errorf("CleanupRule [%d]: arch.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
88+
}
89+
case packages_model.TypeAlt:
90+
if err := alt_service.BuildAllRepositoryFiles(ctx, pcr.OwnerID); err != nil {
91+
return fmt.Errorf("CleanupRule [%d]: alt.BuildAllRepositoryFiles failed: %w", pcr.ID, err)
92+
}
9893
}
9994
}
100-
}
101-
return nil
95+
return nil
96+
})
10297
})
103-
if err != nil {
104-
return err
105-
}
106-
107-
return committer.Commit()
10898
}
10999

110100
type CleanupTarget struct {

tests/integration/api_packages_test.go

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -527,6 +527,42 @@ func TestPackageCleanup(t *testing.T) {
527527

528528
duration, _ := time.ParseDuration("-1h")
529529

530+
t.Run("Debian", func(t *testing.T) {
531+
defer tests.PrintCurrentTest(t)()
532+
// Debian does a repository rebuild.
533+
534+
distribution := "forgejo"
535+
component := "main"
536+
architecture := "amd64"
537+
packageName := "runner"
538+
packageDescription := "Forgejo Runner"
539+
540+
rootURL := fmt.Sprintf("/api/packages/%s/debian", user.Name)
541+
uploadURL := fmt.Sprintf("%s/pool/%s/%s/upload", rootURL, distribution, component)
542+
543+
req := NewRequestWithBody(t, "PUT", uploadURL,
544+
createDebianArchive(packageName, "1.0.0", architecture, packageDescription)).
545+
AddBasicAuth(user.Name)
546+
MakeRequest(t, req, http.StatusCreated)
547+
548+
resp := MakeRequest(t, NewRequestf(t, "GET", "%s/dists/%s/%s/binary-%s/Packages", rootURL, distribution, component, architecture), http.StatusOK)
549+
assert.Contains(t, resp.Body.String(), "pool/forgejo/main/runner_1.0.0_amd64.deb")
550+
551+
pcr, err := packages_model.InsertCleanupRule(t.Context(), &packages_model.PackageCleanupRule{
552+
Enabled: true,
553+
RemovePattern: `.+`,
554+
OwnerID: user.ID,
555+
Type: packages_model.TypeDebian,
556+
})
557+
require.NoError(t, err)
558+
559+
require.NoError(t, packages_cleanup_service.CleanupTask(t.Context(), duration))
560+
561+
MakeRequest(t, NewRequestf(t, "GET", "%s/dists/%s/%s/binary-%s/Packages", rootURL, distribution, component, architecture), http.StatusNotFound)
562+
563+
require.NoError(t, packages_model.DeleteCleanupRuleByID(t.Context(), pcr.ID))
564+
})
565+
530566
t.Run("Common", func(t *testing.T) {
531567
defer tests.PrintCurrentTest(t)()
532568

0 commit comments

Comments
 (0)