fix(01-02): 修复实名激活任务断链,为 PollingHandler 注入 asynq.Client
- PollingHandler struct 新增 asynqClient *asynq.Client 字段 - NewPollingHandler 构造函数新增 asynqClient 参数(redis 之后,gatewayClient 之前) - triggerFirstRealnameActivation 改用 asynqClient.EnqueueContext,删除 RPush 降级方案 - 删除 _ = task 废弃代码 - registerPollingHandlers 传入 h.asynqClient 参数
This commit is contained in:
@@ -30,6 +30,7 @@ type PollingTaskPayload struct {
|
|||||||
type PollingHandler struct {
|
type PollingHandler struct {
|
||||||
db *gorm.DB
|
db *gorm.DB
|
||||||
redis *redis.Client
|
redis *redis.Client
|
||||||
|
asynqClient *asynq.Client // 用于提交 Asynq 任务(如首次实名激活)
|
||||||
gatewayClient *gateway.Client
|
gatewayClient *gateway.Client
|
||||||
iotCardStore *postgres.IotCardStore
|
iotCardStore *postgres.IotCardStore
|
||||||
concurrencyStore *postgres.PollingConcurrencyConfigStore
|
concurrencyStore *postgres.PollingConcurrencyConfigStore
|
||||||
@@ -44,6 +45,7 @@ type PollingHandler struct {
|
|||||||
func NewPollingHandler(
|
func NewPollingHandler(
|
||||||
db *gorm.DB,
|
db *gorm.DB,
|
||||||
redis *redis.Client,
|
redis *redis.Client,
|
||||||
|
asynqClient *asynq.Client,
|
||||||
gatewayClient *gateway.Client,
|
gatewayClient *gateway.Client,
|
||||||
usageService *packagepkg.UsageService,
|
usageService *packagepkg.UsageService,
|
||||||
logger *zap.Logger,
|
logger *zap.Logger,
|
||||||
@@ -51,6 +53,7 @@ func NewPollingHandler(
|
|||||||
return &PollingHandler{
|
return &PollingHandler{
|
||||||
db: db,
|
db: db,
|
||||||
redis: redis,
|
redis: redis,
|
||||||
|
asynqClient: asynqClient,
|
||||||
gatewayClient: gatewayClient,
|
gatewayClient: gatewayClient,
|
||||||
iotCardStore: postgres.NewIotCardStore(db, redis),
|
iotCardStore: postgres.NewIotCardStore(db, redis),
|
||||||
concurrencyStore: postgres.NewPollingConcurrencyConfigStore(db),
|
concurrencyStore: postgres.NewPollingConcurrencyConfigStore(db),
|
||||||
@@ -1097,11 +1100,8 @@ func (h *PollingHandler) triggerFirstRealnameActivation(ctx context.Context, car
|
|||||||
asynq.Queue(constants.QueueDefault),
|
asynq.Queue(constants.QueueDefault),
|
||||||
)
|
)
|
||||||
|
|
||||||
// 这里需要访问 Asynq Client,暂时使用 Redis 队列
|
if _, err := h.asynqClient.EnqueueContext(ctx, task); err != nil {
|
||||||
// 实际应该通过依赖注入 asynq.Client
|
h.logger.Warn("提交首次实名激活任务失败",
|
||||||
activationKey := constants.RedisPollingManualQueueKey(constants.TaskTypePackageFirstActivation)
|
|
||||||
if err := h.redis.RPush(ctx, activationKey, string(payloadBytes)).Err(); err != nil {
|
|
||||||
h.logger.Warn("提交激活任务失败",
|
|
||||||
zap.Uint("package_usage_id", pkg.ID),
|
zap.Uint("package_usage_id", pkg.ID),
|
||||||
zap.Error(err))
|
zap.Error(err))
|
||||||
continue
|
continue
|
||||||
@@ -1110,8 +1110,5 @@ func (h *PollingHandler) triggerFirstRealnameActivation(ctx context.Context, car
|
|||||||
h.logger.Info("已提交首次实名激活任务",
|
h.logger.Info("已提交首次实名激活任务",
|
||||||
zap.Uint("package_usage_id", pkg.ID),
|
zap.Uint("package_usage_id", pkg.ID),
|
||||||
zap.Uint("card_id", cardID))
|
zap.Uint("card_id", cardID))
|
||||||
|
|
||||||
// 避免未使用变量警告
|
|
||||||
_ = task
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -151,6 +151,7 @@ func (h *Handler) registerPollingHandlers() {
|
|||||||
pollingHandler := task.NewPollingHandler(
|
pollingHandler := task.NewPollingHandler(
|
||||||
h.db,
|
h.db,
|
||||||
h.redis,
|
h.redis,
|
||||||
|
h.asynqClient,
|
||||||
h.gatewayClient,
|
h.gatewayClient,
|
||||||
h.workerResult.Services.UsageService,
|
h.workerResult.Services.UsageService,
|
||||||
h.logger,
|
h.logger,
|
||||||
|
|||||||
Reference in New Issue
Block a user