From 84eaf3f555468d1c6b55cbe41b42554a521c67e6 Mon Sep 17 00:00:00 2001 From: ryan Date: Sat, 29 Aug 2026 09:23:01 +0800 Subject: [PATCH] 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. --- backend/plugins/domain/admin/plugin.go | 3 +-- backend/plugins/domain/upload/plugin.go | 20 ++++++++++---------- backend/plugins/domain/user/plugin.go | 8 +++----- scripts/check_cordis_architecture.sh | 19 +++++++++++++++++++ 4 files changed, 33 insertions(+), 17 deletions(-) diff --git a/backend/plugins/domain/admin/plugin.go b/backend/plugins/domain/admin/plugin.go index 6797579a..25b15110 100644 --- a/backend/plugins/domain/admin/plugin.go +++ b/backend/plugins/domain/admin/plugin.go @@ -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)) diff --git a/backend/plugins/domain/upload/plugin.go b/backend/plugins/domain/upload/plugin.go index bc0c224b..a4c1f39c 100644 --- a/backend/plugins/domain/upload/plugin.go +++ b/backend/plugins/domain/upload/plugin.go @@ -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)) diff --git a/backend/plugins/domain/user/plugin.go b/backend/plugins/domain/user/plugin.go index 42dde7f6..ffb5d2ea 100644 --- a/backend/plugins/domain/user/plugin.go +++ b/backend/plugins/domain/user/plugin.go @@ -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 }) diff --git a/scripts/check_cordis_architecture.sh b/scripts/check_cordis_architecture.sh index 3c040d25..0874a908 100755 --- a/scripts/check_cordis_architecture.sh +++ b/scripts/check_cordis_architecture.sh @@ -197,6 +197,25 @@ else log_pass "并发调用统一使用 util.Go 具备 panic 恢复能力" fi +# ============================================================================== +# 7. 插件驱动无关性 (Driver-Agnostic Plugins) +# ============================================================================== +log_check "7. 检查业务/基础设施插件的驱动无关性 (禁止绑定具体 Worker/调度器运行时类型)..." + +# 业务插件只能通过 ctx.Task() / ctx.Schedule() 扩展点声明工作,不得 import Worker +# 驱动运行时类型:绑定 *asynq.Task 之类的签名只被 asynq 驱动满足,换用 in-process +# Worker 后 invokeHandler 会以 "unsupported handler type" 直接拒绝,任务永不执行。 +RUNTIME_BOUND=$(rg -n '"github.com/hibiken/asynq"' \ + "${BACKEND_DIR}/plugins/domain/" "${BACKEND_DIR}/plugins/infra/" \ + --glob '*.go' -g '!*_test.go' 2>/dev/null || true) + +if [ -n "${RUNTIME_BOUND}" ]; then + log_fail "业务/基础设施插件严禁依赖具体 Worker/调度器运行时(必须面向 ctx.Task()/ctx.Schedule() 扩展点编程,否则更换驱动后任务无法执行):" + echo "${RUNTIME_BOUND}" >&2 +else + log_pass "业务/基础设施插件保持驱动无关,可自由切换 Worker 驱动" +fi + # ============================================================================== # 总结与判定 # ==============================================================================