From a9ce281d5e920211e88f66df2e24531633eb1f71 Mon Sep 17 00:00:00 2001 From: zerosaturation Date: Mon, 27 Jul 2026 15:45:42 +0800 Subject: [PATCH] =?UTF-8?q?feat(asset):=20emit=20task:event=20on=20mint=20?= =?UTF-8?q?success=20(fix=20daily=5Fmint=20=E6=82=AC=E7=A9=BA)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit assetService 接入 MQ producer:铸造成功时 emit `task:event { event_type: daily_mint }`, 由 taskService consumer 异步调 ProcessTaskEvent 完成判定。 - main.go: 初始化 MQ (Asynq producer-only),沿用 galleryService 模式; 失败仅 Warn 不阻塞主路径 - service/mint_service.go: CreateMintOrder 成功分支 fire-and-forget emit - mq/producer.go: 新增 EnqueueTaskEvent helper,封装 TaskEventPayload marshal + adapter.TaskProducer.Enqueue 配套:spec §4 / Phase F.1 / impl plan Phase D。 Co-Authored-By: Claude Fable 5 --- backend/services/assetService/main.go | 48 +++++++++++++++- backend/services/assetService/mq/producer.go | 56 +++++++++++++++++++ .../assetService/service/mint_service.go | 12 ++++ 3 files changed, 115 insertions(+), 1 deletion(-) create mode 100644 backend/services/assetService/mq/producer.go diff --git a/backend/services/assetService/main.go b/backend/services/assetService/main.go index d367f51..b1f1bde 100644 --- a/backend/services/assetService/main.go +++ b/backend/services/assetService/main.go @@ -21,6 +21,9 @@ import ( "github.com/topfans/backend/pkg/health" "github.com/topfans/backend/pkg/logger" "github.com/topfans/backend/pkg/models" + "github.com/topfans/backend/pkg/mq" + asynqAdapter "github.com/topfans/backend/pkg/mq/asynq" + queueconsts "github.com/topfans/backend/pkg/queue/consts" "github.com/topfans/backend/pkg/statistic" pbAsset "github.com/topfans/backend/pkg/proto/asset" pbCastlove "github.com/topfans/backend/pkg/proto/castlove" @@ -45,7 +48,12 @@ var ( dbPassword = flag.String("db-password", getEnv("DB_PASSWORD", ""), "Database password") dbName = flag.String("db-name", getEnv("DB_NAME", "top-fans"), "Database name") userServiceURL = flag.String("user-service-url", getEnv("USER_SERVICE_URL", "tri://localhost:20000"), "User service URL") - healthHandler *health.Handler + // MQ 配置:assetService 是 producer-only,emit task:event 给 taskService consumer(修复 daily_mint 悬空) + mqRedisAddr = flag.String("mq-redis-addr", getEnv("MQ_REDIS_ADDR", "localhost:6379"), "MQ redis address") + mqRedisDB = flag.Int("mq-redis-db", getEnvInt("MQ_REDIS_DB", 2), "MQ redis db (avoid clashing with app cache db)") + mqRedisPassword = flag.String("mq-redis-password", getEnv("MQ_REDIS_PASSWORD", getEnv("REDIS_PASSWORD", "")), "MQ redis password") + mqConcurrency = flag.Int("mq-concurrency", getEnvInt("MQ_CONCURRENCY", 10), "Asynq worker concurrency (unused for producer-only)") + healthHandler *health.Handler // ★ mint-task-3 review P1-1: 对账 worker 周期(默认 10 分钟;设为 0 = 关闭) reconcileIntervalSec = flag.Int("reconcile-interval-sec", getEnvInt("MINT_RECONCILE_INTERVAL_SEC", 600), "Reconcile stuck mint orders interval (seconds); 0 disables") reconcileStaleMin = flag.Int("reconcile-stale-minutes", getEnvInt("MINT_RECONCILE_STALE_MINUTES", 10), "Stale threshold for stuck PROCESSING orders (minutes)") @@ -67,6 +75,31 @@ func getEnvInt(key string, fallback int) int { return fallback } +// asynqAdapterConfig 提供 Asynq producer 配置。 +// assetService 是 producer-only,只 emit 不消费;Queues 字段在 producer 路径不生效(adapter 校验用)。 +func asynqAdapterConfig() asynqAdapter.Config { + return asynqAdapter.Config{ + RedisAddr: *mqRedisAddr, + RedisDB: *mqRedisDB, + Password: *mqRedisPassword, + Concurrency: *mqConcurrency, + Queues: map[string]int{queueconsts.QueueDefault: 1}, + } +} + +// asynqAdapterStreamsConfig Redis Streams 配置(同 galleryService 模式)。 +// assetService 当前不订阅 stream,但 mq.Init 需要 Streams 字段非 nil,否则 fallback 默认。 +func asynqAdapterStreamsConfig() mq.StreamsConfig { + return mq.StreamsConfig{ + RedisAddr: *mqRedisAddr, + RedisDB: *mqRedisDB, + Password: *mqRedisPassword, + ConsumerGroup: "topfans-service", + MaxLen: 100000, + ReadBlockMS: 1000, + } +} + func main() { // 加载 .env 文件中的环境变量 godotenv.Load() @@ -156,6 +189,19 @@ func main() { logger.Logger.Info("Statistic SDK initialized") } + // 初始化 MQ (Asynq producer-only) — 修复 daily_mint 悬空 + // 失败仅 Warn 日志不阻塞主路径(与 statistic / dubbo 客户端一致) + if err := mq.Init(mq.Config{ + MQDriver: mq.DriverRedis, + Asynq: asynqAdapterConfig(), + Streams: asynqAdapterStreamsConfig(), + }); err != nil { + logger.Logger.Warn(fmt.Sprintf("MQ init failed (task events disabled, RPC endpoints unaffected): %v", err)) + } else { + defer func() { _ = mq.Close() }() + logger.Logger.Info("MQ initialized (producer-only for task:event)") + } + // 创建 Provider 层实例(用于获取 AssetLevelService) assetLevelProvider := provider.NewAssetLevelProvider(database.GetDB()) assetLevelSvc := assetLevelProvider.GetLevelService() diff --git a/backend/services/assetService/mq/producer.go b/backend/services/assetService/mq/producer.go new file mode 100644 index 0000000..3002f0a --- /dev/null +++ b/backend/services/assetService/mq/producer.go @@ -0,0 +1,56 @@ +// Package mq 是 assetService 的消息队列适配入口(producer-only)。 +// +// 业务侧调用示例 (mint_service.CreateMintOrder 成功分支): +// +// assetmq.EnqueueTaskEvent(ctx, userID, starID, tasks.EventDailyMint) +// +// 全部走 pkg/mq/adapter,不直接 import asynq 或 redis-streams 包。 +// assetService 是 producer-only:不注册 handler / 不起 consumer —— +// task:event 消息由 taskService 的 mq/consumer.go 异步消费处理。 +package mq + +import ( + "context" + "fmt" + + "github.com/topfans/backend/pkg/logger" + "github.com/topfans/backend/pkg/mq/adapter" + "github.com/topfans/backend/pkg/mq/tasks" + queueconsts "github.com/topfans/backend/pkg/queue/consts" + "go.uber.org/zap" +) + +// EnqueueTaskEvent 业务事件 emit(spec §4 / Phase F.1)。 +// +// 修复 daily_mint 悬空:铸造成功时 emit `task:event { event_type: daily_mint }`, +// 由 taskService 的 MQ consumer 异步调 DailyTaskService.ProcessTaskEvent 完成判定。 +// 失败仅 Warn 日志不阻塞主路径(fire-and-forget,与现有 EnqueueExhibitionExpire 一致)。 +func EnqueueTaskEvent(ctx context.Context, userID, starID int64, eventType string) error { + payload, err := tasks.MarshalToPayload(tasks.TaskEventPayload{ + UserID: userID, + StarID: starID, + EventType: eventType, + }) + if err != nil { + return fmt.Errorf("asset mq: marshal task event payload: %w", err) + } + task := adapter.Task{ + Queue: queueconsts.QueueDefault, // taskService consumer 注册到 default 队列 + Type: tasks.TypeTaskEvent, + Payload: payload, + MaxRetry: 3, + } + if _, err := adapter.Get().TaskProducer().Enqueue(ctx, task); err != nil { + logger.Logger.Warn("enqueue task:event failed", + zap.Int64("user_id", userID), + zap.Int64("star_id", starID), + zap.String("event_type", eventType), + zap.Error(err)) + return err + } + logger.Logger.Info("task:event enqueued", + zap.Int64("user_id", userID), + zap.Int64("star_id", starID), + zap.String("event_type", eventType)) + return nil +} \ No newline at end of file diff --git a/backend/services/assetService/service/mint_service.go b/backend/services/assetService/service/mint_service.go index 3885fe5..090f7d5 100644 --- a/backend/services/assetService/service/mint_service.go +++ b/backend/services/assetService/service/mint_service.go @@ -25,8 +25,10 @@ import ( "github.com/topfans/backend/pkg/validator" "github.com/topfans/backend/services/assetService/client" "github.com/topfans/backend/services/assetService/config" + assetmq "github.com/topfans/backend/services/assetService/mq" "github.com/topfans/backend/services/assetService/repository" "github.com/topfans/backend/services/assetService/util" + "github.com/topfans/backend/pkg/mq/tasks" "go.uber.org/zap" "gorm.io/gorm" "gorm.io/gorm/clause" @@ -601,6 +603,16 @@ func (s *mintService) CreateMintOrder(req *pb.CreateMintOrderRequest, userID, st }, }) + // 业务事件:触发 daily_mint 每日任务(修复 daily_mint 悬空) + // fire-and-forget,失败仅 Warn 日志不阻塞主路径(spec §6 / Phase F.1) + if err := assetmq.EnqueueTaskEvent(context.Background(), userID, starID, tasks.EventDailyMint); err != nil { + logger.Logger.Warn("emit task:event daily_mint failed", + zap.Int64("user_id", userID), + zap.Int64("star_id", starID), + zap.String("order_id", mintOrder.OrderID), + zap.Error(err)) + } + return response, nil }