Loading registry/handlers/app.go +27 −4 Original line number Diff line number Diff line Loading @@ -250,7 +250,7 @@ func NewApp(ctx context.Context, config *configuration.Configuration) (_ *App, e } } if err := startUploadPurger(app, app.driver, log, purgeConfig); err != nil { if err := app.startUploadPurger(app.driver, log, purgeConfig); err != nil { return nil, err } Loading Loading @@ -2084,7 +2084,7 @@ func badPurgeUploadConfig(reason string) error { // startUploadPurger schedules a goroutine which will periodically // check upload directories for old files and delete them func startUploadPurger(ctx context.Context, storageDriver storagedriver.StorageDriver, log dcontext.Logger, config map[any]any) error { func (app *App) startUploadPurger(storageDriver storagedriver.StorageDriver, log dcontext.Logger, config map[any]any) error { if v, ok := (config["enabled"]).(bool); ok && !v { return nil } Loading Loading @@ -2131,19 +2131,42 @@ func startUploadPurger(ctx context.Context, storageDriver storagedriver.StorageD return badPurgeUploadConfig("cannot parse dryrun") } ctx, cancel := context.WithCancel(app) done := make(chan struct{}) go func() { defer close(done) // nolint: gosec // used only for jitter calculation jitter := time.Duration(rand.Int()%60) * time.Minute log.Infof("Starting upload purge in %s", jitter) time.Sleep(jitter) select { case <-time.After(jitter): case <-ctx.Done(): return } for { storage.PurgeUploads(ctx, storageDriver, time.Now().Add(-purgeAgeDuration), !dryRunBool) log.Infof("Starting upload purge in %s", intervalDuration) time.Sleep(intervalDuration) select { case <-time.After(intervalDuration): case <-ctx.Done(): return } } }() app.registerShutdownFunc( func(_ *App, errCh chan error, l dlog.Logger) { l.Info("stopping upload purger") cancel() <-done l.Info("upload purger has been shut down") errCh <- nil }, time.Duration(0), ) return nil } Loading registry/handlers/app_test.go +48 −0 Original line number Diff line number Diff line Loading @@ -24,6 +24,7 @@ import ( "github.com/docker/distribution/configuration" dcontext "github.com/docker/distribution/context" dlog "github.com/docker/distribution/log" "github.com/docker/distribution/reference" "github.com/docker/distribution/registry/api/errcode" v1 "github.com/docker/distribution/registry/api/gitlab/v1" Loading Loading @@ -1600,6 +1601,53 @@ func (s *ApplyStorageMiddlewareTestSuite) TestMultipleShutdownFuncs() { require.True(s.T(), shutdown2Called, "Second shutdown function should have been called") } func Test_startUploadPurger_Shutdown(t *testing.T) { ctx := dtestutil.NewContextWithLogger(t) app := &App{ Context: ctx, shutdownConfigs: make([]shutdownConfig, 0), } config := map[any]any{ "enabled": true, "age": "168h", "interval": "24h", "dryrun": false, } err := app.startUploadPurger(inmemory.New(), dcontext.GetLogger(ctx), config) require.NoError(t, err) require.Len(t, app.shutdownConfigs, 1, "expected one shutdown func to be registered") errCh := make(chan error, 1) l := dlog.GetLogger(dlog.WithContext(ctx)) app.shutdownConfigs[0].sFunc(app, errCh, l) select { case err := <-errCh: require.NoError(t, err) case <-time.After(5 * time.Second): t.Fatal("shutdown func did not complete within timeout") } } func Test_startUploadPurger_Disabled(t *testing.T) { ctx := dtestutil.NewContextWithLogger(t) app := &App{ Context: ctx, shutdownConfigs: make([]shutdownConfig, 0), } config := map[any]any{ "enabled": false, } err := app.startUploadPurger(inmemory.New(), dcontext.GetLogger(ctx), config) require.NoError(t, err) require.Empty(t, app.shutdownConfigs, "no shutdown func should be registered when purger is disabled") } func TestApplyStorageMiddlewareTestSuite(t *testing.T) { suite.Run(t, new(ApplyStorageMiddlewareTestSuite)) } Loading
registry/handlers/app.go +27 −4 Original line number Diff line number Diff line Loading @@ -250,7 +250,7 @@ func NewApp(ctx context.Context, config *configuration.Configuration) (_ *App, e } } if err := startUploadPurger(app, app.driver, log, purgeConfig); err != nil { if err := app.startUploadPurger(app.driver, log, purgeConfig); err != nil { return nil, err } Loading Loading @@ -2084,7 +2084,7 @@ func badPurgeUploadConfig(reason string) error { // startUploadPurger schedules a goroutine which will periodically // check upload directories for old files and delete them func startUploadPurger(ctx context.Context, storageDriver storagedriver.StorageDriver, log dcontext.Logger, config map[any]any) error { func (app *App) startUploadPurger(storageDriver storagedriver.StorageDriver, log dcontext.Logger, config map[any]any) error { if v, ok := (config["enabled"]).(bool); ok && !v { return nil } Loading Loading @@ -2131,19 +2131,42 @@ func startUploadPurger(ctx context.Context, storageDriver storagedriver.StorageD return badPurgeUploadConfig("cannot parse dryrun") } ctx, cancel := context.WithCancel(app) done := make(chan struct{}) go func() { defer close(done) // nolint: gosec // used only for jitter calculation jitter := time.Duration(rand.Int()%60) * time.Minute log.Infof("Starting upload purge in %s", jitter) time.Sleep(jitter) select { case <-time.After(jitter): case <-ctx.Done(): return } for { storage.PurgeUploads(ctx, storageDriver, time.Now().Add(-purgeAgeDuration), !dryRunBool) log.Infof("Starting upload purge in %s", intervalDuration) time.Sleep(intervalDuration) select { case <-time.After(intervalDuration): case <-ctx.Done(): return } } }() app.registerShutdownFunc( func(_ *App, errCh chan error, l dlog.Logger) { l.Info("stopping upload purger") cancel() <-done l.Info("upload purger has been shut down") errCh <- nil }, time.Duration(0), ) return nil } Loading
registry/handlers/app_test.go +48 −0 Original line number Diff line number Diff line Loading @@ -24,6 +24,7 @@ import ( "github.com/docker/distribution/configuration" dcontext "github.com/docker/distribution/context" dlog "github.com/docker/distribution/log" "github.com/docker/distribution/reference" "github.com/docker/distribution/registry/api/errcode" v1 "github.com/docker/distribution/registry/api/gitlab/v1" Loading Loading @@ -1600,6 +1601,53 @@ func (s *ApplyStorageMiddlewareTestSuite) TestMultipleShutdownFuncs() { require.True(s.T(), shutdown2Called, "Second shutdown function should have been called") } func Test_startUploadPurger_Shutdown(t *testing.T) { ctx := dtestutil.NewContextWithLogger(t) app := &App{ Context: ctx, shutdownConfigs: make([]shutdownConfig, 0), } config := map[any]any{ "enabled": true, "age": "168h", "interval": "24h", "dryrun": false, } err := app.startUploadPurger(inmemory.New(), dcontext.GetLogger(ctx), config) require.NoError(t, err) require.Len(t, app.shutdownConfigs, 1, "expected one shutdown func to be registered") errCh := make(chan error, 1) l := dlog.GetLogger(dlog.WithContext(ctx)) app.shutdownConfigs[0].sFunc(app, errCh, l) select { case err := <-errCh: require.NoError(t, err) case <-time.After(5 * time.Second): t.Fatal("shutdown func did not complete within timeout") } } func Test_startUploadPurger_Disabled(t *testing.T) { ctx := dtestutil.NewContextWithLogger(t) app := &App{ Context: ctx, shutdownConfigs: make([]shutdownConfig, 0), } config := map[any]any{ "enabled": false, } err := app.startUploadPurger(inmemory.New(), dcontext.GetLogger(ctx), config) require.NoError(t, err) require.Empty(t, app.shutdownConfigs, "no shutdown func should be registered when purger is disabled") } func TestApplyStorageMiddlewareTestSuite(t *testing.T) { suite.Run(t, new(ApplyStorageMiddlewareTestSuite)) }