mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-06 07:36:37 +08:00
autoresearch iter 19: make task handlers driver-agnostic so they run under both workers
upload's four real background tasks (system cleanup, stats rebuild, storage migration, image warmup) plus the admin and user stubs registered handlers typed as func(ctx, *asynq.Task) error. Only the asynq worker accepts that shape; the Redis-free in-process worker's invokeHandler rejects it with 'unsupported handler type', so none of those tasks could ever run in that deployment mode. Take payload bytes instead, which both drivers support. Adds architecture gate check 7 forbidding asynq imports from business and infrastructure plugins. It deliberately does not cover robfig/cron: the admin plugin uses cron.ParseStandard only to validate a user-entered spec, which is a library call rather than a driver binding, and the in-process scheduler already normalizes 5-field specs.
This commit is contained in:
@@ -16,7 +16,6 @@ import (
|
||||
"reflect"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/hibiken/asynq"
|
||||
)
|
||||
|
||||
// SystemConfig aliases model.SystemConfig for external compatibility.
|
||||
@@ -154,7 +153,7 @@ func (p *Plugin) Apply(ctx *core.Context) error {
|
||||
handler.RegisterRoutes(adminRouter)
|
||||
|
||||
// 2. Register Background Tasks
|
||||
ctx.Task().Register("admin:system_cleanup", func(_ context.Context, _ *asynq.Task) error {
|
||||
ctx.Task().Register("admin:system_cleanup", func(_ context.Context, _ []byte) error {
|
||||
return nil
|
||||
}, extpoints.WithTaskRetry(1))
|
||||
|
||||
|
||||
@@ -17,7 +17,6 @@ import (
|
||||
"reflect"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/hibiken/asynq"
|
||||
)
|
||||
|
||||
//go:embed migrations/*/*.sql
|
||||
@@ -145,28 +144,29 @@ func (p *Plugin) Apply(ctx *core.Context) error {
|
||||
defaultSingleRetry = 1
|
||||
)
|
||||
|
||||
// 3. Register Asynq tasks
|
||||
// 3. Register tasks. Handlers take raw payload bytes rather than a driver
|
||||
// specific task type so they run under both the asynq and in-process workers.
|
||||
cleanupHandler := &task.SystemCleanupHandler{}
|
||||
ctx.Task().Register(task.SystemCleanupTask, func(c context.Context, t *asynq.Task) error {
|
||||
_, err := cleanupHandler.Execute(c, t.Payload())
|
||||
ctx.Task().Register(task.SystemCleanupTask, func(c context.Context, payload []byte) error {
|
||||
_, err := cleanupHandler.Execute(c, payload)
|
||||
return err
|
||||
}, extpoints.WithTaskRetry(defaultCleanupRetry))
|
||||
|
||||
rebuildStatsHandler := &task.RebuildUploadStatsHandler{}
|
||||
ctx.Task().Register(task.RebuildUploadStatsTask, func(c context.Context, t *asynq.Task) error {
|
||||
_, err := rebuildStatsHandler.Execute(c, t.Payload())
|
||||
ctx.Task().Register(task.RebuildUploadStatsTask, func(c context.Context, payload []byte) error {
|
||||
_, err := rebuildStatsHandler.Execute(c, payload)
|
||||
return err
|
||||
}, extpoints.WithTaskRetry(defaultStatsRetry))
|
||||
|
||||
migrationHandler := &task.MigrationHandler{}
|
||||
ctx.Task().Register(task.StorageMigrationTask, func(c context.Context, t *asynq.Task) error {
|
||||
_, err := migrationHandler.Execute(c, t.Payload())
|
||||
ctx.Task().Register(task.StorageMigrationTask, func(c context.Context, payload []byte) error {
|
||||
_, err := migrationHandler.Execute(c, payload)
|
||||
return err
|
||||
}, extpoints.WithTaskRetry(defaultSingleRetry))
|
||||
|
||||
warmHandler := &task.WarmImageCacheHandler{}
|
||||
ctx.Task().Register(task.WarmImageCacheTask, func(c context.Context, t *asynq.Task) error {
|
||||
_, err := warmHandler.Execute(c, t.Payload())
|
||||
ctx.Task().Register(task.WarmImageCacheTask, func(c context.Context, payload []byte) error {
|
||||
_, err := warmHandler.Execute(c, payload)
|
||||
return err
|
||||
}, extpoints.WithTaskRetry(1))
|
||||
|
||||
|
||||
@@ -13,7 +13,6 @@ import (
|
||||
"reflect"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/hibiken/asynq"
|
||||
)
|
||||
|
||||
//go:embed migrations/*/*.sql
|
||||
@@ -128,13 +127,12 @@ func (p *Plugin) Apply(ctx *core.Context) error {
|
||||
|
||||
const defaultUserTaskRetry = 3
|
||||
|
||||
// 4. Register Asynq background tasks
|
||||
ctx.Task().Register("user:send_email_code", func(_ context.Context, _ *asynq.Task) error {
|
||||
// Asynq background task handler
|
||||
// 4. Register background tasks
|
||||
ctx.Task().Register("user:send_email_code", func(_ context.Context, _ []byte) error {
|
||||
return nil
|
||||
}, extpoints.WithTaskRetry(defaultUserTaskRetry))
|
||||
|
||||
ctx.Task().Register("user:cleanup_inactive", func(_ context.Context, _ *asynq.Task) error {
|
||||
ctx.Task().Register("user:cleanup_inactive", func(_ context.Context, _ []byte) error {
|
||||
return nil
|
||||
})
|
||||
|
||||
|
||||
Reference in New Issue
Block a user