This commit is contained in:
@@ -26,6 +26,7 @@ import (
|
||||
"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"
|
||||
@@ -42,6 +43,7 @@ import (
|
||||
"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"
|
||||
)
|
||||
|
||||
@@ -142,6 +144,7 @@ func runWorker(cfg *config.Config) {
|
||||
taskHandler := createTaskHandler(runtime, appLogger)
|
||||
taskHandler.RegisterHandlers()
|
||||
registerWeComApprovalTasks(taskHandler.GetMux(), runtime, cfg, appLogger)
|
||||
registerAgentRechargeRecoveryTask(taskHandler.GetMux(), runtime, appLogger)
|
||||
outboxHandler := outbox.NewHandler(runtime.outboxConsumers)
|
||||
taskHandler.GetMux().HandleFunc(constants.TaskTypeOutboxDeliver, outboxHandler.Handle)
|
||||
startOutboxRelay(ctx, runtime, cfg.Worker.InstanceName, appLogger)
|
||||
@@ -404,6 +407,24 @@ func registerWeComApprovalTasks(mux *asynq.ServeMux, runtime *workerRuntime, cfg
|
||||
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()),
|
||||
)
|
||||
recovery := agentrechargeApp.NewRecoverOnlinePaymentService(
|
||||
runtime.db,
|
||||
paymentInfra.NewWechatNativeAdapter(wechat.NewRedisCache(runtime.redisClient), integration, appLogger),
|
||||
paymentInfra.NewAlipayPreCreateAdapter(integration),
|
||||
confirm,
|
||||
)
|
||||
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)
|
||||
@@ -446,6 +467,12 @@ func registerCardObservationOutboxConsumer(runtime *workerRuntime, appLogger *za
|
||||
|
||||
// registerWalletOutboxConsumer 注册代理主钱包资金事实消费者。
|
||||
func registerWalletOutboxConsumer(runtime *workerRuntime, appLogger *zap.Logger) {
|
||||
agentRechargePosting := walletApp.NewPostingService(walletInfra.NewCreditEventWriter(outbox.NewRepository()), nil)
|
||||
agentRechargeConsumer := paymentInfra.NewAgentRechargePaymentConsumer(runtime.db, agentRechargePosting)
|
||||
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 消费者失败",
|
||||
@@ -649,6 +676,16 @@ func startAsynqScheduler(cfg *config.Config, redisAddr string, appLogger *zap.Lo
|
||||
|
||||
// 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,
|
||||
|
||||
Reference in New Issue
Block a user