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 + # ============================================================================== # 总结与判定 # ==============================================================================