feat(task): register task:event MQ consumer

- consumer.go: +newHandleTaskEvent(dailySvc) adapter.TaskHandler
  - unmarshal TaskEventPayload -> delegate to ProcessTaskEvent
  - returns non-nil error on failure -> Asynq MaxRetry=3
- consumer.go: RegisterHandlers signature +dailySvc parameter
- consumer.go: register TaskEvent with QueueDefault, MaxRetry=3
- main.go: pass dailySvc to taskmq.RegisterHandlers

Now task:event MQ messages produced by:
- assetService/mq/producer.go (Phase F.1) for daily_mint
- frontend reportEvent RPC -> gateway (Phase F.2 callers) for daily_login/browse_asset/place_asset
are consumed asynchronously and routed to ProcessTaskEvent.

spec: docs/superpowers/specs/2026-07-21-daily-task-config-driven-design.md §3 F3 + §4.2
plan: docs/superpowers/plans/2026-07-21-daily-task-config-driven-impl.md Phase E

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
zerosaturation 2026-07-27 15:42:10 +08:00
parent ecdd96da5b
commit 1e8a7476e0
2 changed files with 36 additions and 5 deletions

View File

@ -192,7 +192,7 @@ func main() {
} else {
defer func() { _ = mq.Close() }()
// 7.0.1 注册 MQ handler 并启动 consumer
if err := taskmq.RegisterHandlers(revenueSvc); err != nil {
if err := taskmq.RegisterHandlers(revenueSvc, dailySvc); err != nil {
logger.Logger.Warn(fmt.Sprintf("Failed to register task MQ handlersasync tasks disabled: %v", err))
}
mqCtx, mqCancel := context.WithCancel(context.Background())

View File

@ -2,7 +2,7 @@
//
// 业务侧用法 (main.go 启动时):
//
// if err := mq.RegisterHandlers(&revenueSvc); err != nil { ... }
// if err := mq.RegisterHandlers(revenueSvc, dailySvc); err != nil { ... }
// go mq.StartConsumers(ctx)
//
// 业务侧通过 pkg/mq/adapter 接口对接 Asynq,不直接 import asynq。
@ -22,8 +22,8 @@ import (
// RegisterHandlers 注册 taskService 涉及的所有 MQ handler。
//
// 必须传入已经初始化好的 RevenueService 实例,因为 handler 内部会调用。
func RegisterHandlers(revenueSvc service.RevenueService) error {
// 必须传入已经初始化好的 RevenueService + DailyTaskService 实例,因为 handler 内部会调用。
func RegisterHandlers(revenueSvc service.RevenueService, dailySvc service.DailyTaskService) error {
tc := adapter.Get().TaskConsumer()
if err := tc.RegisterTask(tasks.TypeRevenueExhibition, newHandleRevenueExhibition(revenueSvc), adapter.TaskRegisterOptions{
MaxRetry: 3,
@ -37,8 +37,14 @@ func RegisterHandlers(revenueSvc service.RevenueService) error {
}); err != nil {
return fmt.Errorf("register revenue:like-bet: %w", err)
}
if err := tc.RegisterTask(tasks.TypeTaskEvent, newHandleTaskEvent(dailySvc), adapter.TaskRegisterOptions{
MaxRetry: 3,
Queue: queueconsts.QueueDefault,
}); err != nil {
return fmt.Errorf("register task:event: %w", err)
}
logger.Logger.Info("task mq handlers registered",
zap.String("types", tasks.TypeRevenueExhibition+","+tasks.TypeRevenueLikeBet))
zap.String("types", tasks.TypeRevenueExhibition+","+tasks.TypeRevenueLikeBet+","+tasks.TypeTaskEvent))
return nil
}
@ -79,6 +85,31 @@ func newHandleRevenueExhibition(revenueSvc service.RevenueService) adapter.TaskH
}
}
// newHandleTaskEvent 包装 DailyTaskService.ProcessTaskEvent。
//
// 异步消费 task:event 消息spec §4.3 引擎调用链之主路径):
// 来自 assetService 铸造成功 emit / 前端 reportEvent → gateway 转发 MQ → 此 handler
// 调用 ProcessTaskEvent 后忽略返回值MQ consumer 不回包给前端)。
// 失败返回 non-nil error 触发 Asynq MaxRetry=3 重试。
func newHandleTaskEvent(dailySvc service.DailyTaskService) adapter.TaskHandler {
return func(ctx context.Context, t *adapter.Task) error {
var p tasks.TaskEventPayload
if err := tasks.UnmarshalPayload(t.Payload, &p); err != nil {
logger.Logger.Error("handle task:event: unmarshal failed", zap.Error(err))
return fmt.Errorf("unmarshal: %w", err)
}
if _, err := dailySvc.ProcessTaskEvent(ctx, p.UserID, p.StarID, p.EventType); err != nil {
logger.Logger.Error("handle task:event: ProcessTaskEvent failed",
zap.Int64("user_id", p.UserID),
zap.Int64("star_id", p.StarID),
zap.String("event_type", p.EventType),
zap.Error(err))
return err
}
return nil
}
}
// newHandleRevenueLikeBet 包装 RevenueService.RecordLikeBetRevenue。
func newHandleRevenueLikeBet(revenueSvc service.RevenueService) adapter.TaskHandler {
return func(ctx context.Context, t *adapter.Task) error {