Files
junhong_cmp_fiber/cmd/worker/main.go
2026-08-18 17:13:20 +08:00

1011 lines
41 KiB
Go

package main
import (
"context"
"fmt"
"os"
"os/signal"
"strconv"
"syscall"
"time"
"github.com/bytedance/sonic"
"github.com/hibiken/asynq"
"github.com/redis/go-redis/v9"
"go.uber.org/zap"
agentrechargeApp "github.com/break/junhong_cmp_fiber/internal/application/agentrecharge"
approvalApp "github.com/break/junhong_cmp_fiber/internal/application/approval"
auditArchiveApp "github.com/break/junhong_cmp_fiber/internal/application/auditarchive"
cardObservationApp "github.com/break/junhong_cmp_fiber/internal/application/cardobservation"
notificationApp "github.com/break/junhong_cmp_fiber/internal/application/notification"
walletApp "github.com/break/junhong_cmp_fiber/internal/application/wallet"
"github.com/break/junhong_cmp_fiber/internal/bootstrap"
"github.com/break/junhong_cmp_fiber/internal/gateway"
approvalInfra "github.com/break/junhong_cmp_fiber/internal/infrastructure/approval"
auditInfra "github.com/break/junhong_cmp_fiber/internal/infrastructure/audit"
cardObservationInfra "github.com/break/junhong_cmp_fiber/internal/infrastructure/cardobservation"
commissionDelivery "github.com/break/junhong_cmp_fiber/internal/infrastructure/commissiondelivery"
"github.com/break/junhong_cmp_fiber/internal/infrastructure/integrationlog"
"github.com/break/junhong_cmp_fiber/internal/infrastructure/messaging/outbox"
notificationInfra "github.com/break/junhong_cmp_fiber/internal/infrastructure/notification"
paymentInfra "github.com/break/junhong_cmp_fiber/internal/infrastructure/payment"
shopInfra "github.com/break/junhong_cmp_fiber/internal/infrastructure/shop"
walletInfra "github.com/break/junhong_cmp_fiber/internal/infrastructure/wallet"
wecomInfra "github.com/break/junhong_cmp_fiber/internal/infrastructure/wecom"
"github.com/break/junhong_cmp_fiber/internal/model"
"github.com/break/junhong_cmp_fiber/internal/polling"
iot_card_svc "github.com/break/junhong_cmp_fiber/internal/service/iot_card"
refundSvc "github.com/break/junhong_cmp_fiber/internal/service/refund"
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
"github.com/break/junhong_cmp_fiber/internal/task"
pkgBootstrap "github.com/break/junhong_cmp_fiber/pkg/bootstrap"
"github.com/break/junhong_cmp_fiber/pkg/config"
"github.com/break/junhong_cmp_fiber/pkg/constants"
"github.com/break/junhong_cmp_fiber/pkg/database"
"github.com/break/junhong_cmp_fiber/pkg/logger"
"github.com/break/junhong_cmp_fiber/pkg/queue"
"github.com/break/junhong_cmp_fiber/pkg/storage"
"github.com/break/junhong_cmp_fiber/pkg/wechat"
"gorm.io/gorm"
)
const (
workerModuleQueueServer = "queue_server"
workerModulePollingInitializer = "polling_initializer"
workerModulePollingScheduler = "polling_scheduler"
workerModuleAsynqScheduler = "asynq_scheduler"
workerModuleOutboxRelay = "outbox_relay"
importRescueLimit = 500 // 启动补偿单次扫描的最大导入任务数
)
// workerModuleStatus 描述当前 Worker 角色下的模块启停计划。
type workerModuleStatus struct {
enabled []string
disabled []string
}
// workerRuntime 保存所有角色共享的 Worker 运行时依赖。
type workerRuntime struct {
redisAddr string
redisClient *redis.Client
db *gorm.DB
storageSvc *storage.Service
gatewayClient *gateway.Client
asynqClient *asynq.Client
workerResult *bootstrap.WorkerBootstrapResult
workerServer *queue.Server
outboxQueueClient *queue.Client
outboxConsumers *outbox.ConsumerRegistry
pollingConfigMgr *polling.PollingConfigManager
pollingQueueMgr *polling.PollingQueueManager
pollingIotCardStore *postgres.IotCardStore
pollingBase *task.PollingBase
lifecycleSvc *polling.PollingLifecycleService
}
func main() {
cfg, err := config.Load()
if err != nil {
panic("加载配置失败: " + err.Error())
}
runWorker(cfg)
}
// runWorker 编排 Worker 进程启动、模块启停和优雅关闭流程。
func runWorker(cfg *config.Config) {
if _, err := pkgBootstrap.EnsureDirectories(cfg, nil); err != nil {
panic("初始化目录失败: " + err.Error())
}
if err := logger.InitLoggers(
cfg.Logging.Level,
cfg.Logging.Development,
logger.LogRotationConfig{
Filename: cfg.Logging.AppLog.Filename,
MaxSize: cfg.Logging.AppLog.MaxSize,
MaxBackups: cfg.Logging.AppLog.MaxBackups,
MaxAge: cfg.Logging.AppLog.MaxAge,
Compress: cfg.Logging.AppLog.Compress,
},
logger.LogRotationConfig{
Filename: cfg.Logging.AccessLog.Filename,
MaxSize: cfg.Logging.AccessLog.MaxSize,
MaxBackups: cfg.Logging.AccessLog.MaxBackups,
MaxAge: cfg.Logging.AccessLog.MaxAge,
Compress: cfg.Logging.AccessLog.Compress,
},
); err != nil {
panic("初始化日志失败: " + err.Error())
}
defer func() {
_ = logger.Sync() // 忽略 sync 错误
}()
appLogger := logger.GetAppLogger()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
moduleStatus := buildWorkerModuleStatus(cfg.Worker.Role)
appLogger.Info("Worker 服务启动中...")
logWorkerRole(appLogger, cfg.Worker.Role, cfg.Worker.InstanceName, moduleStatus)
runtime := initWorkerRuntime(ctx, cfg, appLogger)
defer runtime.close(appLogger)
var pollingInitializer *polling.PollingInitializer
var pollingScheduler *polling.Scheduler
var asynqScheduler *asynq.Scheduler
if runsSingletonModules(cfg.Worker.Role) {
pollingInitializer = startPollingInitializer(ctx, runtime, appLogger)
pollingScheduler = startPollingScheduler(ctx, runtime, pollingInitializer, appLogger)
asynqScheduler = startAsynqScheduler(cfg, runtime.redisAddr, appLogger)
}
taskHandler := createTaskHandler(runtime, appLogger)
taskHandler.RegisterHandlers()
registerWeComApprovalTasks(taskHandler.GetMux(), runtime, cfg, appLogger)
registerAgentRechargeRecoveryTask(taskHandler.GetMux(), runtime, appLogger)
registerAuditArchiveTask(taskHandler.GetMux(), runtime, cfg.Worker.AuditRetentionCleanupEnabled, appLogger)
outboxHandler := outbox.NewHandler(runtime.outboxConsumers)
taskHandler.GetMux().HandleFunc(constants.TaskTypeOutboxDeliver, outboxHandler.Handle)
startOutboxRelay(ctx, runtime, cfg.Worker.InstanceName, appLogger)
rescuePendingImportTasks(ctx, runtime, appLogger)
appLogger.Info("Worker 服务器配置完成",
zap.Int("concurrency", cfg.Queue.Concurrency),
zap.Any("queues", cfg.Queue.Queues))
quit := make(chan os.Signal, 1)
signal.Notify(quit, os.Interrupt, syscall.SIGTERM)
defer signal.Stop(quit)
go func() {
if err := runtime.workerServer.Run(taskHandler.GetMux()); err != nil {
appLogger.Fatal("Worker 服务器运行失败", zap.Error(err))
}
}()
appLogger.Info("Worker 服务器已启动")
<-quit
shutdownWorker(cancel, appLogger, runtime.workerServer, pollingInitializer, pollingScheduler, asynqScheduler)
}
// buildWorkerModuleStatus 根据角色生成启用/禁用模块列表,供启动日志和人工验收使用。
func buildWorkerModuleStatus(role string) workerModuleStatus {
status := workerModuleStatus{
enabled: []string{workerModuleQueueServer, workerModuleOutboxRelay},
}
if runsSingletonModules(role) {
status.enabled = append(
status.enabled,
workerModulePollingInitializer,
workerModulePollingScheduler,
workerModuleAsynqScheduler,
)
return status
}
status.disabled = []string{
workerModulePollingInitializer,
workerModulePollingScheduler,
workerModuleAsynqScheduler,
}
return status
}
// logWorkerRole 输出当前实例角色、实例名以及模块启停计划。
func logWorkerRole(appLogger *zap.Logger, role string, instanceName string, moduleStatus workerModuleStatus) {
appLogger.Info("Worker 角色已确认",
zap.String("role", role),
zap.String("instance_name", instanceName))
appLogger.Info("Worker 模块启停计划",
zap.Strings("enabled_modules", moduleStatus.enabled),
zap.Strings("disabled_modules", moduleStatus.disabled))
}
// runsSingletonModules 判断当前角色是否承担主动调度和初始化职责。
func runsSingletonModules(role string) bool {
return role == constants.WorkerRoleAll || role == constants.WorkerRoleLeader
}
// initWorkerRuntime 初始化所有角色共享的外部依赖和轮询运行时对象。
func initWorkerRuntime(ctx context.Context, cfg *config.Config, appLogger *zap.Logger) *workerRuntime {
redisAddr := cfg.Redis.Address + ":" + strconv.Itoa(cfg.Redis.Port)
redisClient := redis.NewClient(&redis.Options{
Addr: redisAddr,
Password: cfg.Redis.Password,
DB: cfg.Redis.DB,
PoolSize: cfg.Redis.PoolSize,
MinIdleConns: cfg.Redis.MinIdleConns,
DialTimeout: cfg.Redis.DialTimeout,
ReadTimeout: cfg.Redis.ReadTimeout,
WriteTimeout: cfg.Redis.WriteTimeout,
})
if err := redisClient.Ping(ctx).Err(); err != nil {
appLogger.Fatal("连接 Redis 失败", zap.Error(err))
}
appLogger.Info("Redis 已连接", zap.String("address", redisAddr))
db, err := database.InitPostgreSQL(&cfg.Database, appLogger)
if err != nil {
appLogger.Fatal("初始化 PostgreSQL 失败", zap.Error(err))
}
storageSvc := initStorage(cfg, appLogger)
gatewayClient := initGateway(cfg, appLogger)
asynqClient := asynq.NewClient(asynq.RedisClientOpt{
Addr: redisAddr,
Password: cfg.Redis.Password,
DB: cfg.Redis.DB,
})
workerQueueClient := queue.NewClient(redisClient, appLogger)
workerDeps := &bootstrap.WorkerDependencies{
DB: db,
Redis: redisClient,
Logger: appLogger,
AsynqClient: asynqClient,
QueueClient: workerQueueClient,
StorageService: storageSvc,
GatewayClient: gatewayClient,
}
workerResult, err := bootstrap.BootstrapWorker(workerDeps)
if err != nil {
appLogger.Fatal("Worker Bootstrap 失败", zap.Error(err))
}
workerServer := queue.NewServer(redisClient, &cfg.Queue, appLogger)
outboxQueueClient := workerQueueClient
pollingConfigStore := postgres.NewPollingConfigStore(db)
pollingConfigMgr := polling.NewPollingConfigManager(pollingConfigStore, redisClient, appLogger)
if err := pollingConfigMgr.Load(ctx); err != nil {
appLogger.Warn("加载轮询配置失败,使用空配置继续", zap.Error(err))
}
pollingConfigMgr.Start(ctx)
pollingQueueMgr := polling.NewPollingQueueManager(redisClient, constants.PollingShardCount, appLogger)
pollingIotCardStore := postgres.NewIotCardStore(db, redisClient)
pollingBase := task.NewPollingBase(
redisClient,
pollingQueueMgr,
pollingConfigMgr,
pollingIotCardStore,
appLogger,
cfg.Polling.VerboseLog,
)
pollingDeviceSimBindingStore := postgres.NewDeviceSimBindingStore(db, redisClient)
pollingDeviceStore := postgres.NewDeviceStore(db, redisClient)
lifecycleSvc := polling.NewPollingLifecycleService(
pollingQueueMgr,
pollingConfigMgr,
pollingIotCardStore,
pollingDeviceSimBindingStore,
pollingDeviceStore,
appLogger,
)
if stopResumeSvc, ok := workerResult.Services.StopResumeService.(*iot_card_svc.StopResumeService); ok {
stopResumeSvc.SetPollingCallback(lifecycleSvc)
}
runtime := &workerRuntime{
redisAddr: redisAddr,
redisClient: redisClient,
db: db,
storageSvc: storageSvc,
gatewayClient: gatewayClient,
asynqClient: asynqClient,
workerResult: workerResult,
workerServer: workerServer,
outboxQueueClient: outboxQueueClient,
outboxConsumers: outbox.NewConsumerRegistry(),
pollingConfigMgr: pollingConfigMgr,
pollingQueueMgr: pollingQueueMgr,
pollingIotCardStore: pollingIotCardStore,
pollingBase: pollingBase,
lifecycleSvc: lifecycleSvc,
}
registerNotificationOutboxConsumer(runtime, appLogger)
registerWalletOutboxConsumer(runtime, appLogger)
registerCardObservationOutboxConsumer(runtime, appLogger)
registerWeComApprovalOutboxConsumer(runtime, cfg, appLogger)
return runtime
}
// registerWeComApprovalOutboxConsumer 注册企业微信审批提交和标准终态业务消费者。
func registerWeComApprovalOutboxConsumer(runtime *workerRuntime, cfg *config.Config, appLogger *zap.Logger) {
auditWriter, ok := runtime.workerResult.Services.PaymentAudit.(*auditInfra.Writer)
if !ok || auditWriter == nil {
appLogger.Fatal("通用审批统一审计 Writer 未配置")
}
applicationRepository := wecomInfra.NewApplicationRepository(runtime.db)
integrationRepository := integrationlog.NewRepository(runtime.db)
tokenProvider := wecomInfra.NewTokenProvider(
applicationRepository, runtime.redisClient, integrationRepository,
cfg.WeCom.BaseURL, cfg.WeCom.Timeout, appLogger,
)
consumer := wecomInfra.NewApprovalSubmissionConsumer(
wecomInfra.NewApprovalContextRepository(runtime.db, auditWriter),
wecomInfra.NewApprovalSubmissionClient(tokenProvider, integrationRepository, cfg.WeCom.BaseURL, cfg.WeCom.Timeout),
wecomInfra.NewApprovalAttachmentUploader(tokenProvider, integrationRepository, runtime.storageSvc, cfg.WeCom.BaseURL, cfg.WeCom.Timeout),
)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeApprovalSubmissionRequested, consumer); err != nil {
appLogger.Fatal("注册企业微信审批提交 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeApprovalSubmissionRequested), zap.Error(err))
}
owner := cfg.Worker.InstanceName
if owner == "" {
owner = fmt.Sprintf("worker-%d", os.Getpid())
}
walletPosting := walletApp.NewPostingService(
walletInfra.NewCreditEventWriter(outbox.NewRepository(), auditInfra.NewWriter(auditInfra.NewRegistry(), nil)),
nil,
)
stopResumeService, _ := runtime.workerResult.Services.StopResumeService.(*iot_card_svc.StopResumeService)
refundService := refundSvc.New(
runtime.db,
postgres.NewRefundStore(runtime.db),
runtime.workerResult.Stores.Order,
runtime.workerResult.Stores.CommissionRecord,
runtime.workerResult.Stores.AgentWallet,
runtime.workerResult.Stores.AgentWalletTransaction,
stopResumeService,
nil,
runtime.workerResult.Services.ActivationService,
runtime.workerResult.Stores.IotCard,
runtime.workerResult.Stores.Device,
runtime.workerResult.Stores.AssetWallet,
appLogger,
)
refundService.SetAgentWalletRefundService(
walletApp.NewRefundService(walletInfra.NewRefundEventWriter(outbox.NewRepository()), nil),
)
refundService.SetNotificationOutbox(outbox.NewRepository())
refundService.SetLifecycleAudit(auditWriter)
if err := runtime.outboxConsumers.Register(commissionDelivery.EventRefundCommissionDeduct, commissionDelivery.NewRefundConsumer(refundService.ProcessCommissionDeduction, refundService.ProcessAssetPostProcessing)); err != nil {
appLogger.Fatal("注册退款佣金回扣 Outbox 消费者失败", zap.Error(err))
}
if err := runtime.outboxConsumers.Register(commissionDelivery.EventRefundAssetProcess, commissionDelivery.NewRefundConsumer(refundService.ProcessCommissionDeduction, refundService.ProcessAssetPostProcessing)); err != nil {
appLogger.Fatal("注册退款资产后处理 Outbox 消费者失败", zap.Error(err))
}
if err := runtime.outboxConsumers.Register(commissionDelivery.EventCommissionCalculate, commissionDelivery.NewCommissionConsumer(runtime.outboxQueueClient, appLogger)); err != nil {
appLogger.Fatal("注册订单佣金计算 Outbox 消费者失败", zap.Error(err))
}
commissionDelivery.Recover(context.Background(), runtime.db, outbox.NewRepository(), 100, appLogger)
decisionDispatcher := approvalApp.NewDecisionDispatcher(
approvalInfra.NewDecisionDeliveryStore(runtime.db),
map[string]approvalApp.BusinessDecisionHandler{
constants.ApprovalBusinessTypeOfflineRecharge: agentrechargeApp.NewApprovalDecisionHandler(runtime.db, walletPosting, runtime.workerResult.Services.RechargeAudit),
constants.ApprovalBusinessTypeRefund: refundService,
},
owner,
appLogger,
nil,
)
decisionConsumer := approvalInfra.NewTerminalDecisionConsumer(decisionDispatcher)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeApprovalTerminalDecision, decisionConsumer); err != nil {
appLogger.Fatal("注册审批标准终态 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeApprovalTerminalDecision), zap.Error(err))
}
}
// registerWeComApprovalTasks 注册企微权威详情同步和主动恢复任务。
func registerWeComApprovalTasks(mux *asynq.ServeMux, runtime *workerRuntime, cfg *config.Config, appLogger *zap.Logger) {
auditWriter, ok := runtime.workerResult.Services.PaymentAudit.(*auditInfra.Writer)
if !ok || auditWriter == nil {
appLogger.Fatal("通用审批统一审计 Writer 未配置")
}
applicationRepository := wecomInfra.NewApplicationRepository(runtime.db)
integrationRepository := integrationlog.NewRepository(runtime.db)
tokenProvider := wecomInfra.NewTokenProvider(
applicationRepository, runtime.redisClient, integrationRepository,
cfg.WeCom.BaseURL, cfg.WeCom.Timeout, appLogger,
)
decisionSync := approvalApp.NewSyncDecisionService(
runtime.db, approvalInfra.NewRepositoryProvider(),
approvalInfra.NewTerminalEventWriter(outbox.NewRepository()),
approvalInfra.NewDecisionDeliveryStore(runtime.db), nil,
)
decisionSync.SetAuditWriter(auditWriter)
contexts := wecomInfra.NewApprovalContextRepository(runtime.db, auditWriter)
detailHandler := wecomInfra.NewApprovalDetailTaskHandler(
wecomInfra.NewApprovalDetailClient(tokenProvider, integrationRepository, cfg.WeCom.BaseURL, cfg.WeCom.Timeout),
contexts, decisionSync, integrationRepository,
)
recoveryHandler := wecomInfra.NewApprovalRecoveryTaskHandler(
contexts,
wecomInfra.NewApprovalInfoClient(tokenProvider, integrationRepository, cfg.WeCom.BaseURL, cfg.WeCom.Timeout),
runtime.outboxQueueClient,
)
mux.HandleFunc(constants.TaskTypeWeComApprovalSync, detailHandler.Handle)
mux.HandleFunc(constants.TaskTypeWeComApprovalRecovery, recoveryHandler.Handle)
appLogger.Info("注册企业微信审批详情同步任务处理器", zap.String("task_type", constants.TaskTypeWeComApprovalSync))
appLogger.Info("注册企业微信审批主动恢复任务处理器", zap.String("task_type", constants.TaskTypeWeComApprovalRecovery))
}
// registerAgentRechargeRecoveryTask 注册代理在线充值支付恢复任务。
func registerAgentRechargeRecoveryTask(mux *asynq.ServeMux, runtime *workerRuntime, appLogger *zap.Logger) {
integration := integrationlog.NewRepository(runtime.db)
confirm := agentrechargeApp.NewConfirmOnlinePaymentService(
runtime.db,
paymentInfra.NewAgentRechargePaymentEventWriter(outbox.NewRepository()),
runtime.workerResult.Services.PaymentAudit,
)
recovery := agentrechargeApp.NewRecoverOnlinePaymentService(
runtime.db,
paymentInfra.NewWechatWebAdapter(wechat.NewRedisCache(runtime.redisClient), integration, appLogger),
paymentInfra.NewAlipayWapAdapter(integration, appLogger),
paymentInfra.NewFuiouScanAdapter(integration, appLogger),
confirm,
runtime.workerResult.Services.PaymentAudit,
)
handler := paymentInfra.NewAgentRechargeRecoveryTaskHandler(recovery)
mux.HandleFunc(constants.TaskTypeAgentRechargeRecovery, handler.Handle)
appLogger.Info("注册代理在线充值支付恢复任务处理器", zap.String("task_type", constants.TaskTypeAgentRechargeRecovery))
}
// registerCardObservationOutboxConsumer 注册卡观测领域事件消费者。
func registerCardObservationOutboxConsumer(runtime *workerRuntime, appLogger *zap.Logger) {
stopResumeService, _ := runtime.workerResult.Services.StopResumeService.(iot_card_svc.StopResumeServiceInterface)
consumer := cardObservationInfra.NewRealnameChangedConsumer(
runtime.db,
runtime.workerResult.Services.ActivationService,
runtime.workerResult.Stores.DeviceSimBinding,
stopResumeService,
)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeCardRealnameChanged, consumer); err != nil {
appLogger.Fatal("注册卡实名状态变化 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeCardRealnameChanged), zap.Error(err))
}
trafficConsumer := cardObservationInfra.NewTrafficIncrementedConsumer(
runtime.db,
runtime.redisClient,
runtime.workerResult.Services.UsageService,
stopResumeService,
)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeCardTrafficIncremented, trafficConsumer); err != nil {
appLogger.Fatal("注册卡流量正增量 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeCardTrafficIncremented), zap.Error(err))
}
networkConsumer := cardObservationInfra.NewNetworkChangedConsumer(runtime.db, stopResumeService)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeCardNetworkChanged, networkConsumer); err != nil {
appLogger.Fatal("注册卡网络状态变化 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeCardNetworkChanged), zap.Error(err))
}
seriesTrigger := cardObservationApp.NewSeriesTrigger(
cardObservationInfra.NewSeriesCoordinator(runtime.redisClient),
queue.NewCardObservationSeriesScheduler(runtime.outboxQueueClient),
cardObservationInfra.NewSeriesAttemptLogger(integrationlog.NewRepository(runtime.db)),
)
seriesConsumer := cardObservationInfra.NewSeriesRequestedConsumer(seriesTrigger)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeCardSeriesRequested, seriesConsumer); err != nil {
appLogger.Fatal("注册业务卡观测序列 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeCardSeriesRequested), zap.Error(err))
}
}
// registerWalletOutboxConsumer 注册代理主钱包资金事实消费者。
func registerWalletOutboxConsumer(runtime *workerRuntime, appLogger *zap.Logger) {
agentRechargePosting := walletApp.NewPostingService(walletInfra.NewCreditEventWriter(outbox.NewRepository(), auditInfra.NewWriter(auditInfra.NewRegistry(), nil)), nil)
agentRechargeConsumer := paymentInfra.NewAgentRechargePaymentConsumer(runtime.db, agentRechargePosting, runtime.workerResult.Services.RechargeAudit)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeAgentRechargePaymentConfirmed, agentRechargeConsumer); err != nil {
appLogger.Fatal("注册代理在线充值入账 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeAgentRechargePaymentConfirmed), zap.Error(err))
}
debitConsumer := walletInfra.NewDebitEventConsumer(runtime.db)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeAgentMainWalletDebited, debitConsumer); err != nil {
appLogger.Fatal("注册代理主钱包扣款 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeAgentMainWalletDebited), zap.Error(err))
}
reservationConsumer := walletInfra.NewReservationEventConsumer(runtime.db)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeAgentMainWalletReservationChanged, reservationConsumer); err != nil {
appLogger.Fatal("注册代理主钱包预占 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeAgentMainWalletReservationChanged), zap.Error(err))
}
creditConsumer := walletInfra.NewCreditEventConsumer(runtime.db)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeAgentMainWalletCredited, creditConsumer); err != nil {
appLogger.Fatal("注册代理主钱包入账 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeAgentMainWalletCredited), zap.Error(err))
}
refundConsumer := walletInfra.NewRefundEventConsumer(runtime.db)
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeAgentMainWalletRefunded, refundConsumer); err != nil {
appLogger.Fatal("注册代理主钱包退款 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeAgentMainWalletRefunded), zap.Error(err))
}
}
// registerNotificationOutboxConsumer 注册明确、动态后台账号和个人客户通知的稳定 Outbox 消费者。
func registerNotificationOutboxConsumer(runtime *workerRuntime, appLogger *zap.Logger) {
repository := notificationInfra.NewRepository(runtime.db)
registry := notificationInfra.NewRegistry()
shopRecipientResolver := shopInfra.NewRecipientResolver(runtime.db)
dynamicRecipientResolver := notificationInfra.NewDynamicRecipientResolver(runtime.db, shopRecipientResolver)
consumer := notificationApp.NewDeliveryService(repository, registry, dynamicRecipientResolver, appLogger, auditInfra.NewWriter(auditInfra.NewRegistry(), nil))
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeAdminDirectNotification, consumer); err != nil {
appLogger.Fatal("注册站内通知 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeAdminDirectNotification), zap.Error(err))
}
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypePersonalCustomerDirectNotification, consumer); err != nil {
appLogger.Fatal("注册个人客户站内通知 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypePersonalCustomerDirectNotification), zap.Error(err))
}
if err := runtime.outboxConsumers.Register(constants.OutboxEventTypeAdminDynamicNotification, consumer); err != nil {
appLogger.Fatal("注册后台动态接收人站内通知 Outbox 消费者失败",
zap.String("event_type", constants.OutboxEventTypeAdminDynamicNotification), zap.Error(err))
}
}
// close 在 Worker 退出时关闭共享客户端与数据库连接。
func (r *workerRuntime) close(appLogger *zap.Logger) {
if r.outboxQueueClient != nil {
if err := r.outboxQueueClient.Close(); err != nil {
appLogger.Error("关闭 Outbox 队列客户端失败", zap.Error(err))
}
}
if r.asynqClient != nil {
if err := r.asynqClient.Close(); err != nil {
appLogger.Error("关闭 Asynq 客户端失败", zap.Error(err))
}
}
if r.db != nil {
sqlDB, _ := r.db.DB()
if sqlDB != nil {
if err := sqlDB.Close(); err != nil {
appLogger.Error("关闭 PostgreSQL 连接失败", zap.Error(err))
}
}
}
if r.redisClient != nil {
if err := r.redisClient.Close(); err != nil {
appLogger.Error("关闭 Redis 客户端失败", zap.Error(err))
}
}
}
// startOutboxRelay 启动公共 Outbox 的可取消轮询;多个 Worker 依靠 PostgreSQL 租约并发协调。
func startOutboxRelay(ctx context.Context, runtime *workerRuntime, instanceName string, appLogger *zap.Logger) {
owner := instanceName
if owner == "" {
owner = fmt.Sprintf("worker-%d", os.Getpid())
}
relay, err := outbox.NewRelay(
runtime.db,
outbox.NewQueuePublisher(runtime.outboxQueueClient),
appLogger,
outbox.RelayOptions{Owner: owner},
)
if err != nil {
appLogger.Fatal("初始化 Outbox Relay 失败", zap.Error(err))
}
go func() {
ticker := time.NewTicker(constants.OutboxRelayPollInterval)
defer ticker.Stop()
for {
if _, err := relay.ProcessBatch(ctx); err != nil && ctx.Err() == nil {
appLogger.Error("Outbox Relay 批次处理失败",
zap.String("component", "outbox_relay"), zap.String("error_code", "OUTBOX_RELAY_BATCH_FAILED"), zap.Error(err))
}
select {
case <-ctx.Done():
return
case <-ticker.C:
}
}
}()
appLogger.Info("Outbox Relay 已启动", zap.String("lease_owner", owner))
}
// startPollingInitializer 为 leader/all 启动渐进式初始化器,并保留配置变更重启逻辑。
func startPollingInitializer(ctx context.Context, runtime *workerRuntime, appLogger *zap.Logger) *polling.PollingInitializer {
pollingInitializer := polling.NewPollingInitializer(
runtime.pollingIotCardStore,
runtime.redisClient,
runtime.pollingConfigMgr,
runtime.pollingQueueMgr,
appLogger,
)
pollingInitializer.StartBackground(ctx)
runtime.pollingConfigMgr.WatchChanges(ctx, func(hadConfigs, hasConfigs bool) {
if hasConfigs {
appLogger.Info("轮询配置已变更,触发队列重新初始化",
zap.Bool("had_configs", hadConfigs))
pollingInitializer.Restart(ctx)
return
}
appLogger.Info("轮询配置已清空,跳过队列重新初始化")
})
return pollingInitializer
}
// startPollingScheduler 为 leader/all 启动轮询调度器。
func startPollingScheduler(
ctx context.Context,
runtime *workerRuntime,
pollingInitializer *polling.PollingInitializer,
appLogger *zap.Logger,
) *polling.Scheduler {
dataResetHandler := polling.NewDataResetHandler(runtime.workerResult.Services.ResetService, appLogger)
activationHandler := polling.NewPackageActivationHandler(
runtime.db,
runtime.redisClient,
runtime.asynqClient,
runtime.workerResult.Services.ActivationService,
runtime.workerResult.Services.StopResumeService,
appLogger,
)
// 注入套餐失效前流量同步器:复用公共卡观测写入,不在旧 Service 内重复扣减套餐。
trafficSyncer := iot_card_svc.New(
runtime.db,
runtime.pollingIotCardStore,
nil, nil, nil, nil, nil,
runtime.gatewayClient,
appLogger,
)
trafficSyncer.SetRedisClient(runtime.redisClient)
trafficSyncer.SetCardObservationService(runtime.workerResult.Services.CardObservation)
trafficSyncer.SetSpeedTierIntegrationLog(integrationlog.NewRepository(runtime.db))
activationHandler.SetTrafficSyncer(trafficSyncer)
pollingScheduler := polling.NewScheduler(
runtime.redisClient,
runtime.asynqClient,
runtime.pollingQueueMgr,
runtime.pollingConfigMgr,
appLogger,
activationHandler,
dataResetHandler,
)
pollingScheduler.SetInitializer(pollingInitializer)
if err := pollingScheduler.Start(ctx); err != nil {
appLogger.Error("启动轮询调度器失败", zap.Error(err))
return nil
}
return pollingScheduler
}
// startAsynqScheduler 为 leader/all 创建并启动 Asynq Scheduler。
func startAsynqScheduler(cfg *config.Config, redisAddr string, appLogger *zap.Logger) *asynq.Scheduler {
asynqScheduler := asynq.NewScheduler(
asynq.RedisClientOpt{
Addr: redisAddr,
Password: cfg.Redis.Password,
DB: cfg.Redis.DB,
},
&asynq.SchedulerOpts{Location: time.Local},
)
if err := registerAsynqScheduleTasks(asynqScheduler); err != nil {
appLogger.Fatal("注册 Asynq 定时任务失败", zap.Error(err))
}
go func() {
if err := asynqScheduler.Run(); err != nil {
appLogger.Fatal("Asynq Scheduler 启动失败", zap.Error(err))
}
}()
appLogger.Info("Asynq Scheduler 已启动",
zap.Bool("audit_retention_cleanup_enabled", cfg.Worker.AuditRetentionCleanupEnabled))
return asynqScheduler
}
// registerAsynqScheduleTasks 注册 Worker 入口需要的全部定时任务。
func registerAsynqScheduleTasks(asynqScheduler *asynq.Scheduler) error {
if _, err := asynqScheduler.Register("@every 1m", asynq.NewTask(
constants.TaskTypeAgentRechargeRecovery,
nil,
asynq.MaxRetry(3),
asynq.Timeout(10*time.Minute),
asynq.Unique(10*time.Minute),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeAgentRechargeRecovery)),
)); err != nil {
return fmt.Errorf("注册代理在线充值支付恢复定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("@every 1m", asynq.NewTask(
constants.TaskTypeOrderExpire,
nil,
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeOrderExpire)),
)); err != nil {
return fmt.Errorf("注册订单超时定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("@every 1m", asynq.NewTask(
constants.TaskTypeAlertCheck,
nil,
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeAlertCheck)),
)); err != nil {
return fmt.Errorf("注册告警检查定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("@every 2m", asynq.NewTask(
constants.TaskTypeWeComApprovalRecovery,
nil,
asynq.MaxRetry(3),
asynq.Timeout(2*time.Minute),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeWeComApprovalRecovery)),
)); err != nil {
return fmt.Errorf("注册企业微信审批主动恢复定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("0 2 * * *", asynq.NewTask(
constants.TaskTypeDataCleanup,
nil,
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeDataCleanup)),
)); err != nil {
return fmt.Errorf("注册数据清理定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("15 2 * * *", asynq.NewTask(
constants.TaskTypeNotificationCleanup,
nil,
asynq.MaxRetry(3),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeNotificationCleanup)),
)); err != nil {
return fmt.Errorf("注册站内通知保留清理定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("CRON_TZ=Asia/Shanghai 0 3 * * *", asynq.NewTask(
constants.TaskTypePackageExpiryReminder,
nil,
asynq.MaxRetry(3),
asynq.Timeout(10*time.Minute),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypePackageExpiryReminder)),
)); err != nil {
return fmt.Errorf("注册每日套餐临期提醒扫描定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register(
"0 2 * * *",
asynq.NewTask(
constants.TaskTypeDailyTrafficFlush,
nil,
asynq.MaxRetry(3),
asynq.Timeout(5*time.Minute),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeDailyTrafficFlush)),
),
); err != nil {
return fmt.Errorf("注册每日流量落盘定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("CRON_TZ=Asia/Shanghai 0 4 * * *", asynq.NewTask(
constants.TaskTypeAuditDailyArchive,
nil,
asynq.MaxRetry(10),
asynq.Timeout(2*time.Hour),
asynq.Unique(23*time.Hour),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeAuditDailyArchive)),
)); err != nil {
return fmt.Errorf("注册统一审计每日冷归档定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("CRON_TZ=Asia/Shanghai 30 4 * * *", asynq.NewTask(
constants.TaskTypeIntegrationDailyArchive,
nil,
asynq.MaxRetry(10),
asynq.Timeout(2*time.Hour),
asynq.Unique(23*time.Hour),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeIntegrationDailyArchive)),
)); err != nil {
return fmt.Errorf("注册 Integration Log 每日冷归档定时任务失败: %w", err)
}
if _, err := asynqScheduler.Register("CRON_TZ=Asia/Shanghai 0 5 1 * *", asynq.NewTask(
constants.TaskTypeIntegrationMonthlyFinalize,
nil,
asynq.MaxRetry(10),
asynq.Timeout(6*time.Hour),
asynq.Unique(27*24*time.Hour),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeIntegrationMonthlyFinalize)),
)); err != nil {
return fmt.Errorf("注册 Integration Log 月度最终版本复核任务失败: %w", err)
}
if _, err := asynqScheduler.Register("CRON_TZ=Asia/Shanghai 0 6 1 * *", asynq.NewTask(
constants.TaskTypeAuditMonthlyRetention,
nil,
asynq.MaxRetry(10),
asynq.Timeout(12*time.Hour),
asynq.Unique(27*24*time.Hour),
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeAuditMonthlyRetention)),
)); err != nil {
return fmt.Errorf("注册月度日志留存演练或清理任务失败: %w", err)
}
return nil
}
// registerAuditArchiveTask 注册 Audit 与 Integration 冷归档任务处理器。
func registerAuditArchiveTask(mux *asynq.ServeMux, runtime *workerRuntime, cleanupEnabled bool, appLogger *zap.Logger) {
if runtime.storageSvc == nil {
appLogger.Warn("对象存储未配置,审计归档任务将在执行时重试")
mux.HandleFunc(constants.TaskTypeAuditDailyArchive, task.NewAuditDailyArchiveHandler(nil, appLogger).Handle)
integrationHandler := task.NewIntegrationArchiveHandler(nil, appLogger)
mux.HandleFunc(constants.TaskTypeIntegrationDailyArchive, integrationHandler.HandleDaily)
mux.HandleFunc(constants.TaskTypeIntegrationMonthlyFinalize, integrationHandler.HandleMonthlyFinalize)
mux.HandleFunc(constants.TaskTypeAuditMonthlyRetention, task.NewAuditMonthlyRetentionHandler(nil, appLogger, cleanupEnabled).Handle)
return
}
auditWriter, ok := runtime.workerResult.Services.PaymentAudit.(*auditInfra.Writer)
if !ok || auditWriter == nil {
appLogger.Fatal("初始化月度日志留存清理失败:统一审计 Writer 未配置")
}
service, err := auditArchiveApp.NewService(runtime.db, runtime.storageSvc.Provider(), constants.AuditArchiveInstanceID, auditWriter)
if err != nil {
appLogger.Fatal("初始化统一审计归档服务失败", zap.Error(err))
}
mux.HandleFunc(constants.TaskTypeAuditDailyArchive, task.NewAuditDailyArchiveHandler(service, appLogger).Handle)
integrationHandler := task.NewIntegrationArchiveHandler(service, appLogger)
mux.HandleFunc(constants.TaskTypeIntegrationDailyArchive, integrationHandler.HandleDaily)
mux.HandleFunc(constants.TaskTypeIntegrationMonthlyFinalize, integrationHandler.HandleMonthlyFinalize)
mux.HandleFunc(constants.TaskTypeAuditMonthlyRetention, task.NewAuditMonthlyRetentionHandler(service, appLogger, cleanupEnabled).Handle)
appLogger.Info("注册审计归档任务处理器",
zap.String("audit_task_type", constants.TaskTypeAuditDailyArchive),
zap.String("integration_daily_task_type", constants.TaskTypeIntegrationDailyArchive),
zap.String("integration_monthly_task_type", constants.TaskTypeIntegrationMonthlyFinalize),
zap.String("retention_task_type", constants.TaskTypeAuditMonthlyRetention),
zap.Bool("retention_cleanup_enabled", cleanupEnabled))
}
// createTaskHandler 创建并返回包含全部任务处理器的 Asynq Handler。
func createTaskHandler(runtime *workerRuntime, appLogger *zap.Logger) *queue.Handler {
stopResumeSvc, _ := runtime.workerResult.Services.StopResumeService.(iot_card_svc.StopResumeServiceInterface)
return queue.NewHandler(
runtime.db,
runtime.redisClient,
runtime.storageSvc,
runtime.gatewayClient,
runtime.lifecycleSvc,
runtime.workerResult,
runtime.asynqClient,
appLogger,
runtime.pollingBase,
stopResumeSvc,
)
}
// rescuePendingImportTasks 将历史遗留的待处理导入任务重新投递到独立导入队列。
func rescuePendingImportTasks(ctx context.Context, runtime *workerRuntime, appLogger *zap.Logger) {
rescuePendingIotCardImportTasks(ctx, runtime.db, runtime.asynqClient, appLogger)
rescuePendingDeviceImportTasks(ctx, runtime.db, runtime.asynqClient, appLogger)
}
// rescuePendingIotCardImportTasks 补偿仍停留在待处理状态的 IoT 卡导入任务。
func rescuePendingIotCardImportTasks(ctx context.Context, db *gorm.DB, asynqClient *asynq.Client, appLogger *zap.Logger) {
var importTasks []model.IotCardImportTask
if err := db.WithContext(ctx).
Where("status = ?", model.ImportTaskStatusPending).
Limit(importRescueLimit).
Find(&importTasks).Error; err != nil {
appLogger.Warn("扫描待补偿 IoT 卡导入任务失败", zap.Error(err))
return
}
for _, importTask := range importTasks {
payload := task.IotCardImportPayload{TaskID: importTask.ID}
enqueueImportRescueTask(ctx, asynqClient, constants.TaskTypeIotCardImport, payload, importTask.ID, appLogger)
}
}
// rescuePendingDeviceImportTasks 补偿仍停留在待处理状态的设备导入任务。
func rescuePendingDeviceImportTasks(ctx context.Context, db *gorm.DB, asynqClient *asynq.Client, appLogger *zap.Logger) {
var importTasks []model.DeviceImportTask
if err := db.WithContext(ctx).
Where("status = ?", model.ImportTaskStatusPending).
Limit(importRescueLimit).
Find(&importTasks).Error; err != nil {
appLogger.Warn("扫描待补偿设备导入任务失败", zap.Error(err))
return
}
for _, importTask := range importTasks {
payload := task.DeviceImportPayload{TaskID: importTask.ID}
enqueueImportRescueTask(ctx, asynqClient, constants.TaskTypeDeviceImport, payload, importTask.ID, appLogger)
}
}
// enqueueImportRescueTask 将补偿任务提交到任务类型对应的独立队列。
func enqueueImportRescueTask(ctx context.Context, asynqClient *asynq.Client, taskType string, payload any, taskID uint, appLogger *zap.Logger) {
payloadBytes, err := sonic.Marshal(payload)
if err != nil {
appLogger.Warn("序列化导入补偿任务载荷失败",
zap.String("task_type", taskType),
zap.Uint("task_id", taskID),
zap.Error(err))
return
}
queueName := constants.QueueForTaskType(taskType)
taskMessage := asynq.NewTask(
taskType,
payloadBytes,
asynq.Queue(queueName),
asynq.TaskID(fmt.Sprintf("import-rescue:%s:%d", taskType, taskID)),
asynq.Unique(30*time.Minute),
)
if _, err := asynqClient.EnqueueContext(ctx, taskMessage); err != nil {
appLogger.Warn("提交导入补偿任务失败",
zap.String("task_type", taskType),
zap.String("queue", queueName),
zap.Uint("task_id", taskID),
zap.Error(err))
return
}
appLogger.Info("导入补偿任务已提交",
zap.String("task_type", taskType),
zap.String("queue", queueName),
zap.Uint("task_id", taskID))
}
// shutdownWorker 按当前实例实际启动过的模块执行优雅关闭。
func shutdownWorker(
cancel context.CancelFunc,
appLogger *zap.Logger,
workerServer *queue.Server,
pollingInitializer *polling.PollingInitializer,
pollingScheduler *polling.Scheduler,
asynqScheduler *asynq.Scheduler,
) {
appLogger.Info("正在关闭 Worker 服务器...")
if asynqScheduler != nil {
asynqScheduler.Shutdown()
}
if pollingScheduler != nil {
pollingScheduler.Stop()
}
cancel()
if pollingInitializer != nil {
pollingInitializer.Stop()
}
workerServer.Shutdown()
appLogger.Info("Worker 服务器已停止")
}
func initStorage(cfg *config.Config, appLogger *zap.Logger) *storage.Service {
if cfg.Storage.Provider == "" || cfg.Storage.S3.Endpoint == "" {
appLogger.Info("对象存储未配置,跳过初始化")
return nil
}
provider, err := storage.NewS3Provider(&cfg.Storage)
if err != nil {
appLogger.Warn("初始化对象存储失败,功能将不可用", zap.Error(err))
return nil
}
appLogger.Info("对象存储已初始化",
zap.String("provider", cfg.Storage.Provider),
zap.String("bucket", cfg.Storage.S3.Bucket),
)
return storage.NewService(provider, &cfg.Storage)
}
// initGateway 初始化 Gateway 客户端
func initGateway(cfg *config.Config, appLogger *zap.Logger) *gateway.Client {
if cfg.Gateway.BaseURL == "" {
appLogger.Info("Gateway 未配置,跳过初始化(轮询任务将无法查询真实数据)")
return nil
}
client := gateway.NewClient(
cfg.Gateway.BaseURL,
cfg.Gateway.AppID,
cfg.Gateway.AppSecret,
appLogger,
).WithTimeout(time.Duration(cfg.Gateway.Timeout) * time.Second)
appLogger.Info("Gateway 客户端初始化成功",
zap.String("base_url", cfg.Gateway.BaseURL),
zap.String("app_id", cfg.Gateway.AppID))
return client
}