topfans/backend/services/assetService/main.go
zerosaturation a9ce281d5e feat(asset): emit task:event on mint success (fix daily_mint 悬空)
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 <noreply@anthropic.com>
2026-07-27 15:45:42 +08:00

388 lines
14 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package main
import (
"context"
"flag"
"fmt"
"os"
"os/signal"
"strconv"
"sync"
"syscall"
"time"
"dubbo.apache.org/dubbo-go/v3/client"
_ "dubbo.apache.org/dubbo-go/v3/imports"
"dubbo.apache.org/dubbo-go/v3/protocol"
"dubbo.apache.org/dubbo-go/v3/server"
"github.com/joho/godotenv"
"github.com/topfans/backend/pkg/database"
"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"
pbRanking "github.com/topfans/backend/pkg/proto/ranking"
pbUser "github.com/topfans/backend/pkg/proto/user"
assetClient "github.com/topfans/backend/services/assetService/client"
"github.com/topfans/backend/services/assetService/config"
"github.com/topfans/backend/services/assetService/provider"
"github.com/topfans/backend/services/assetService/repository"
"github.com/topfans/backend/services/assetService/service"
"github.com/topfans/backend/services/assetService/util"
"github.com/topfans/backend/services/assetService/util/ossutil"
"go.uber.org/zap"
)
var (
port = flag.Int("port", getEnvInt("PORT", 20003), "Dubbo service port")
dbHost = flag.String("db-host", getEnv("DB_HOST", "localhost"), "Database host")
dbPort = flag.Int("db-port", getEnvInt("DB_PORT", 5432), "Database port")
dbUser = flag.String("db-user", getEnv("DB_USER", "postgres"), "Database user")
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")
// MQ 配置assetService 是 producer-onlyemit 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)")
)
func getEnv(key, fallback string) string {
if v := os.Getenv(key); v != "" {
return v
}
return fallback
}
func getEnvInt(key string, fallback int) int {
if v := os.Getenv(key); v != "" {
if n, err := strconv.Atoi(v); err == nil {
return n
}
}
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()
flag.Parse()
// 初始化日志(必须在最前面)
env := os.Getenv("ENV")
if env == "" {
env = "development"
}
if err := logger.Init(logger.Config{
ServiceName: "asset-service",
Environment: env,
LogLevel: os.Getenv("LOG_LEVEL"),
}); err != nil {
panic(fmt.Sprintf("Failed to initialize logger: %v", err))
}
defer logger.Sync()
logger.Logger.Info("Starting Asset Service...")
// 初始化数据库
dbConfig := database.Config{
Host: *dbHost,
Port: *dbPort,
User: *dbUser,
Password: *dbPassword,
DBName: *dbName,
SSLMode: "disable",
TimeZone: "Asia/Shanghai",
}
if err := database.Init(dbConfig); err != nil {
logger.Logger.Fatal(fmt.Sprintf("Failed to initialize database: %v", err))
}
logger.Logger.Info("Database initialized successfully")
// 启动健康检查 HTTP 服务器
healthPort := *port + 1000 // e.g., 20003 -> 21003
healthHandler = health.NewHandler("asset-service", healthPort)
healthHandler.Start()
// 自动迁移数据库表
if err := autoMigrate(); err != nil {
logger.Logger.Fatal(fmt.Sprintf("Failed to migrate database: %v", err))
}
// 创建 Repository 层实例
assetRepo := repository.NewAssetRepository(database.GetDB())
mintOrderRepo := repository.NewMintOrderRepository(database.GetDB())
assetLikeRepo := repository.NewAssetLikeRepository(database.GetDB())
rankingRepo := repository.NewRankingRepository(database.GetDB())
materialRepo := repository.NewMaterialRepository(database.GetDB())
relationRepo := repository.NewAssetMaterialRelationRepository(database.GetDB())
mintCostRepo := repository.NewMintCostRepository()
userMintCountRepo := repository.NewUserMintCountRepository()
castloveConfigRepo := repository.NewCastloveConfigRepository(database.GetDB())
shareRepo := repository.NewShareRepo(database.GetDB())
logger.Logger.Info("Repository layer initialized")
// 创建 Dubbo 客户端
cli, err := client.NewClient(
client.WithClientURL(*userServiceURL),
)
if err != nil {
logger.Logger.Fatal(fmt.Sprintf("Failed to create Dubbo client: %v", err))
}
// 获取 User Service RPC 客户端
userServiceClient, err := pbUser.NewUserSocialService(cli)
if err != nil {
logger.Logger.Fatal(fmt.Sprintf("Failed to create User Service RPC client: %v", err))
}
userClient := assetClient.NewUserServiceClient(userServiceClient)
logger.Logger.Info("User Service RPC client initialized")
// 初始化 statisticService SDK事件埋点fire-and-forget
statisticServiceURL := getEnv("STATISTIC_SERVICE_URL", "tri://localhost:20009")
statisticCli, err := client.NewClient(client.WithClientURL(statisticServiceURL))
if err != nil {
logger.Logger.Warn(fmt.Sprintf("statisticService client create failed (events disabled): %v", err))
} else if err := statistic.Init(statisticCli); err != nil {
logger.Logger.Warn(fmt.Sprintf("statistic.Init failed (events disabled): %v", err))
} else {
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()
logger.Logger.Info("AssetLevelProvider initialized")
// 创建 Service 层实例
registryRepo := repository.NewAssetRegistryRepository(database.GetDB())
// 分享服务依赖:OSS QR 上传器 + Redis 缓存 + 落地页 base URL
ossCfg := util.OSSConfig{
Region: os.Getenv("OSS_REGION"),
BucketName: os.Getenv("OSS_BUCKET_NAME"),
RoleArn: os.Getenv("OSS_STS_ROLE_ARN"),
AccessKeyID: os.Getenv("OSS_ACCESS_KEY_ID"),
AccessKeySecret: os.Getenv("OSS_ACCESS_KEY_SECRET"),
}
qrUploader := ossutil.NewOSSQRUploader(ossCfg)
// H5 落地页 base URL:
// - 本地开发: 默认 http://localhost:5173(Vite H5 端口),可在 backend/.env 覆盖
// - 生产环境: 由 deploy/envs/asset.env 注入 https://api.topfans.online
// 注意: 与 frontend/.env.production 的 VITE_LANDING_BASE_URL 保持一致
landingBase := getEnv("LANDING_BASE_URL", "http://localhost:5173")
shareService := service.NewShareService(shareRepo, database.GetRedis(), qrUploader, landingBase)
assetService := service.NewAssetService(assetRepo, mintOrderRepo, assetLikeRepo, userClient, database.GetDB(), registryRepo, shareService)
mintService := service.NewMintService(assetRepo, mintOrderRepo, userClient, database.GetDB(), config.GlobalAssetConfig, registryRepo, mintCostRepo, userMintCountRepo, assetLevelSvc)
assetLikeService := service.NewAssetLikeService(assetRepo, assetLikeRepo, database.GetDB(), assetLevelSvc)
rankingService := service.NewRankingService(rankingRepo, assetLikeRepo, userClient)
materialService := service.NewMaterialService(materialRepo, relationRepo)
castloveConfigService := service.NewCastloveConfigService(castloveConfigRepo)
logger.Logger.Info("Service layer initialized")
// 创建 Provider 层实例
assetProvider := provider.NewAssetProvider(assetService, mintService, assetLikeService, materialService)
rankingProvider := provider.NewRankingProvider(rankingService)
castloveConfigProvider := provider.NewCastloveConfigProvider(castloveConfigService)
logger.Logger.Info("Provider layer initialized")
// 启动赛季重置 Worker每小时检查一次
seasonResetWorker := assetLevelProvider.GetSeasonResetWorker()
go func() {
ticker := time.NewTicker(time.Hour)
defer ticker.Stop()
for {
<-ticker.C
seasonResetWorker.Run()
}
}()
logger.Logger.Info("Season reset worker started")
// ★ mint-task-3 review P1-1: 对账 worker
// - 周期触发 ReconcileStuckMintOrders(扫描陈旧 PROCESSING 订单并补单/标 FAILED)
// - 间隔与陈旧阈值通过 flag/env 可调(默认 10 分钟,生产可放大到 30+ 分钟)
// - interval=0 时关闭(本地/压测)
// - 随服务优雅退出:cancelFunc + WaitGroup
var reconcileWG sync.WaitGroup
var reconcileCancel context.CancelFunc
if *reconcileIntervalSec > 0 {
interval := time.Duration(*reconcileIntervalSec) * time.Second
stale := time.Duration(*reconcileStaleMin) * time.Minute
reconcileCtx, cancel := context.WithCancel(context.Background())
reconcileCancel = cancel
reconcileWG.Add(1)
go func() {
defer reconcileWG.Done()
// 启动后先做一次"热启动对账",把上次服务崩溃遗留的 PROCESSING 订单先处理一波
// (避免冷启动后还要等满 interval 才有第一波对账)
logger.Logger.Info("Mint reconcile worker bootstrapping",
zap.Duration("interval", interval),
zap.Duration("stale_after", stale))
if err := mintService.ReconcileStuckMintOrders(reconcileCtx, stale); err != nil {
logger.Logger.Error("reconcile bootstrap failed", zap.Error(err))
}
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-reconcileCtx.Done():
logger.Logger.Info("Mint reconcile worker exiting")
return
case <-ticker.C:
// 每轮对账本身已有 INFO 日志(scanned/recovered/marked_failed/skipped),
// 顶层不需要再包一行 INFO。
if err := mintService.ReconcileStuckMintOrders(reconcileCtx, stale); err != nil {
logger.Logger.Error("reconcile tick failed", zap.Error(err))
}
}
}
}()
logger.Logger.Info("Mint reconcile worker started",
zap.Duration("interval", interval),
zap.Duration("stale_after", stale))
} else {
logger.Logger.Info("Mint reconcile worker disabled (interval=0)")
}
// 创建 Dubbo 服务器
srv, err := server.NewServer(
server.WithServerProtocol(
protocol.WithPort(*port),
protocol.WithTriple(),
),
)
if err != nil {
logger.Logger.Fatal(fmt.Sprintf("Failed to create Dubbo server: %v", err))
}
// 使用 Triple 协议生成的 RegisterHandler 函数注册服务
if err := pbAsset.RegisterAssetServiceHandler(srv, assetProvider); err != nil {
logger.Logger.Fatal(fmt.Sprintf("Failed to register Asset Service: %v", err))
}
// 注册 Ranking Service
if err := pbRanking.RegisterRankingServiceHandler(srv, rankingProvider); err != nil {
logger.Logger.Fatal(fmt.Sprintf("Failed to register Ranking Service: %v", err))
}
// 注册 Castlove Config Service
if err := pbCastlove.RegisterCastloveConfigServiceHandler(srv, castloveConfigProvider); err != nil {
logger.Logger.Fatal(fmt.Sprintf("Failed to register Castlove Config Service: %v", err))
}
// 启动服务
if err := srv.Serve(); err != nil {
logger.Logger.Fatal(fmt.Sprintf("Failed to start Asset Service: %v", err))
}
logger.Logger.Info(fmt.Sprintf("Asset Service started successfully on port %d", *port))
// 等待退出信号
quit := make(chan os.Signal, 1)
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
<-quit
logger.Logger.Info("Shutting down Asset Service...")
// 关闭健康检查服务器
if healthHandler != nil {
healthHandler.Stop()
}
// 停止 mint 对账 worker(若已启动)
if reconcileCancel != nil {
reconcileCancel()
reconcileWG.Wait()
}
}
// autoMigrate 自动迁移数据库表
func autoMigrate() error {
db := database.GetDB()
if db == nil {
return fmt.Errorf("database is not initialized")
}
// 按顺序迁移资产相关表
tables := []interface{}{
&models.Asset{},
&models.MintOrder{},
&models.AssetLike{},
&models.Material{},
&models.AssetMaterialRelation{},
&models.AssetLevelRecord{},
&models.AssetLevelChangeLog{},
&models.Season{},
&models.SeasonDecayConfig{},
&models.LaserCardTemplate{},
&models.LaserCardInstance{},
&models.LaserCardOperationLog{},
// 铸爱工艺配置(只读,后台直连写库;此处 AutoMigrate 仅做兜底,实际表结构由外部 SQL 维护)
&models.CastloveCategory{},
&models.CastloveCraft{},
}
for _, table := range tables {
if err := db.AutoMigrate(table); err != nil {
return fmt.Errorf("failed to migrate table: %w", err)
}
}
logger.Logger.Info("Database migration completed successfully")
return nil
}