From 1e8a7476e0050b5223a8a04e3843e81d209229c9 Mon Sep 17 00:00:00 2001 From: zerosaturation Date: Mon, 27 Jul 2026 15:42:10 +0800 Subject: [PATCH] feat(task): register task:event MQ consumer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 --- backend/services/taskService/main.go | 2 +- backend/services/taskService/mq/consumer.go | 39 ++++++++++++++++++--- 2 files changed, 36 insertions(+), 5 deletions(-) diff --git a/backend/services/taskService/main.go b/backend/services/taskService/main.go index 3774abc..1ccba2c 100644 --- a/backend/services/taskService/main.go +++ b/backend/services/taskService/main.go @@ -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 handlers(async tasks disabled): %v", err)) } mqCtx, mqCancel := context.WithCancel(context.Background()) diff --git a/backend/services/taskService/mq/consumer.go b/backend/services/taskService/mq/consumer.go index ac4385f..d01484d 100644 --- a/backend/services/taskService/mq/consumer.go +++ b/backend/services/taskService/mq/consumer.go @@ -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 {