修复套餐接续停机竞态与轮询兜底
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 11m13s
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 11m13s
This commit is contained in:
@@ -253,7 +253,8 @@ func (h *PackageActivationHandler) findExpiredMainPackages(ctx context.Context)
|
||||
}
|
||||
|
||||
// processExpiredPackage 处理单个过期套餐
|
||||
// 流程:先同步最新流量 → 事务内标记过期和失效加油包 → 提交后投递下一套餐并触发停机检查
|
||||
// 流程:先同步最新流量 → 事务内标记过期和失效加油包 → 提交后投递下一套餐;
|
||||
// 仅明确没有后续套餐时才触发停机检查,避免与异步激活任务竞态。
|
||||
func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pkg *model.PackageUsage) error {
|
||||
carrierType, carrierID := h.getCarrierInfo(pkg)
|
||||
|
||||
@@ -311,9 +312,15 @@ func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pk
|
||||
|
||||
// 事务提交后再投递,确保消费者只能读取到旧套餐已经过期的状态。
|
||||
if carrierType != "" && carrierID > 0 {
|
||||
activationErr := h.activateNextPackage(ctx, h.db.WithContext(ctx), carrierType, carrierID)
|
||||
activationEnqueued, err := h.activateNextPackage(ctx, h.db.WithContext(ctx), carrierType, carrierID)
|
||||
if err != nil {
|
||||
// 激活结果未知时不能将其当作无套餐;孤儿扫描和套餐轮询会继续兜底。
|
||||
return err
|
||||
}
|
||||
if activationEnqueued {
|
||||
return nil
|
||||
}
|
||||
h.triggerStopAfterExpiry(ctx, carrierType, carrierID)
|
||||
return activationErr
|
||||
}
|
||||
|
||||
return nil
|
||||
@@ -368,7 +375,8 @@ func (h *PackageActivationHandler) getCarrierInfo(pkg *model.PackageUsage) (stri
|
||||
}
|
||||
|
||||
// activateNextPackage 提交下一个待生效主套餐的异步激活任务。
|
||||
func (h *PackageActivationHandler) activateNextPackage(ctx context.Context, tx *gorm.DB, carrierType string, carrierID uint) error {
|
||||
// 返回 true 表示已找到并成功提交后续套餐;false 表示不存在后续套餐。
|
||||
func (h *PackageActivationHandler) activateNextPackage(ctx context.Context, tx *gorm.DB, carrierType string, carrierID uint) (bool, error) {
|
||||
var nextPkg model.PackageUsage
|
||||
query := tx.Where("status = ?", constants.PackageUsageStatusPending).
|
||||
Where("master_usage_id IS NULL").
|
||||
@@ -383,16 +391,19 @@ func (h *PackageActivationHandler) activateNextPackage(ctx context.Context, tx *
|
||||
|
||||
if err := query.First(&nextPkg).Error; err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return nil
|
||||
return false, nil
|
||||
}
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
|
||||
return h.enqueueActivationTask(ctx, nextPkg.ID, carrierType, carrierID, "queue")
|
||||
if err := h.enqueueActivationTask(ctx, nextPkg.ID, carrierType, carrierID, "queue"); err != nil {
|
||||
return false, err
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// triggerStopAfterExpiry 套餐过期后异步触发停机检查
|
||||
// 仅在确认无后续生效套餐时有效;CheckAndStopCard 内部有幂等保护,重复调用安全
|
||||
// triggerStopAfterExpiry 在明确无后续套餐时异步触发停机检查。
|
||||
// 后续套餐存在时,停机重评估由激活任务在完成后顺序执行。
|
||||
func (h *PackageActivationHandler) triggerStopAfterExpiry(ctx context.Context, carrierType string, carrierID uint) {
|
||||
if h.stopResumeCallback == nil {
|
||||
return
|
||||
@@ -431,6 +442,42 @@ func (h *PackageActivationHandler) triggerStopAfterExpiry(ctx context.Context, c
|
||||
}
|
||||
}
|
||||
|
||||
// reconcileCarrierAfterActivation 在套餐激活任务完成后按最新事实重评估停复机。
|
||||
// 它必须在激活调用返回后执行,不能与激活任务并发读取过期权益快照。
|
||||
func (h *PackageActivationHandler) reconcileCarrierAfterActivation(ctx context.Context, packageUsageID uint, carrierType string, carrierID uint) error {
|
||||
if h.stopResumeCallback == nil {
|
||||
h.logger.Warn("套餐激活后停复机回调未注入,跳过重评估",
|
||||
zap.Uint("package_usage_id", packageUsageID),
|
||||
zap.String("carrier_type", carrierType),
|
||||
zap.Uint("carrier_id", carrierID))
|
||||
return nil
|
||||
}
|
||||
if carrierID == 0 {
|
||||
return errors.New(errors.CodeInvalidParam, "套餐使用记录缺少有效载体")
|
||||
}
|
||||
|
||||
if carrierType == constants.AssetTypeIotCard {
|
||||
if err := h.stopResumeCallback.CheckAndStopCard(ctx, carrierID); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if carrierType != constants.AssetTypeDevice {
|
||||
return errors.New(errors.CodeInvalidParam, "套餐使用记录载体类型无效")
|
||||
}
|
||||
|
||||
bindings, err := h.deviceSimBinding.ListByDeviceID(ctx, carrierID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, binding := range bindings {
|
||||
if err := h.stopResumeCallback.CheckAndStopCard(ctx, binding.IotCardID); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// enqueueActivationTask 提交套餐激活任务到 Asynq
|
||||
func (h *PackageActivationHandler) enqueueActivationTask(ctx context.Context, packageUsageID uint, carrierType string, carrierID uint, activationType string) error {
|
||||
linkage := auditcontext.From(ctx)
|
||||
@@ -499,14 +546,7 @@ func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Cont
|
||||
return err
|
||||
}
|
||||
|
||||
// 幂等性检查:如果已经是生效状态,跳过
|
||||
if pkg.Status == constants.PackageUsageStatusActive {
|
||||
h.logger.Info("套餐已激活,跳过",
|
||||
zap.Uint("package_usage_id", payload.PackageUsageID))
|
||||
return nil
|
||||
}
|
||||
|
||||
// 调用 ActivationService 执行激活
|
||||
// 调用 ActivationService 执行激活。即使套餐已由其他任务激活,仍须重新评估停复机。
|
||||
if h.activationService != nil {
|
||||
if err := h.activationService.ActivateSpecificPackage(ctx, payload.PackageUsageID); err != nil {
|
||||
h.logger.Error("套餐激活失败",
|
||||
@@ -521,8 +561,20 @@ func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Cont
|
||||
return errors.New(errors.CodeInternalError, "激活服务未注入,无法执行套餐激活")
|
||||
}
|
||||
|
||||
h.logger.Info("套餐激活成功",
|
||||
zap.Uint("package_usage_id", payload.PackageUsageID))
|
||||
carrierType, carrierID := h.getCarrierInfo(&pkg)
|
||||
if err := h.reconcileCarrierAfterActivation(ctx, payload.PackageUsageID, carrierType, carrierID); err != nil {
|
||||
h.logger.Error("套餐激活后停复机重评估失败",
|
||||
zap.Uint("package_usage_id", payload.PackageUsageID),
|
||||
zap.String("carrier_type", carrierType),
|
||||
zap.Uint("carrier_id", carrierID),
|
||||
zap.Error(err))
|
||||
return err
|
||||
}
|
||||
|
||||
h.logger.Info("套餐激活及停复机重评估完成",
|
||||
zap.Uint("package_usage_id", payload.PackageUsageID),
|
||||
zap.String("carrier_type", carrierType),
|
||||
zap.Uint("carrier_id", carrierID))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -98,6 +98,20 @@ func (m *PollingQueueManager) Requeue(ctx context.Context, cardID uint, taskType
|
||||
}).Err()
|
||||
}
|
||||
|
||||
// EnsureQueued 仅在任务当前不在分片队列中时补入任务,不覆盖已有任务的执行时间。
|
||||
func (m *PollingQueueManager) EnsureQueued(ctx context.Context, cardID uint, taskType string, nextCheckAt time.Time) (bool, error) {
|
||||
shardID := int(cardID) % m.shardCount
|
||||
key := constants.RedisPollingShardQueueKey(shardID, taskType)
|
||||
added, err := m.redis.ZAddArgs(ctx, key, redis.ZAddArgs{
|
||||
NX: true,
|
||||
Members: []redis.Z{{
|
||||
Score: float64(nextCheckAt.Unix()),
|
||||
Member: fmt.Sprintf("%d", cardID),
|
||||
}},
|
||||
}).Result()
|
||||
return added > 0, err
|
||||
}
|
||||
|
||||
// RemoveFromAllQueues 从所有分片的所有5个队列(realname/carddata/package/protect/card_status)移除指定卡
|
||||
// 修复 Bug3:旧实现漏掉 protect 队列
|
||||
func (m *PollingQueueManager) RemoveFromAllQueues(ctx context.Context, cardID uint) error {
|
||||
|
||||
@@ -334,10 +334,10 @@ func (s *ActivationService) ActivateSpecificPackage(ctx context.Context, package
|
||||
return errors.Wrap(errors.CodeRedisError, err, "获取分布式锁失败")
|
||||
}
|
||||
if !locked {
|
||||
s.logger.Warn("套餐激活正在进行中,跳过",
|
||||
s.logger.Warn("套餐激活正在进行中,等待任务重试",
|
||||
zap.String("carrier_type", carrierType),
|
||||
zap.Uint("carrier_id", carrierID))
|
||||
return nil
|
||||
return errors.New(errors.CodePackageActivationConflict)
|
||||
}
|
||||
defer s.redis.Del(ctx, lockKey)
|
||||
|
||||
|
||||
@@ -132,6 +132,29 @@ func (b *PollingBase) releaseCardTrafficSyncLock(_ context.Context, cardID uint,
|
||||
}
|
||||
}
|
||||
|
||||
// ensureMissingTask 根据数据库最新卡状态,仅在对应分片队列缺失时补入任务。
|
||||
func (b *PollingBase) ensureMissingTask(ctx context.Context, cardID uint, taskType string) error {
|
||||
card, err := b.iotCardStore.GetByID(ctx, cardID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
info, ok := b.configMgr.MergedTaskIntervals(card)[taskType]
|
||||
if !ok || info.Interval <= 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
added, err := b.queueMgr.EnsureQueued(ctx, cardID, taskType, time.Now())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if added {
|
||||
b.logger.Info("卡状态轮询补齐缺失套餐任务",
|
||||
zap.Uint("card_id", cardID), zap.String("task_type", taskType))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// requeueCardAt 使用独立短超时上下文执行真正的 ZADD 重入队。
|
||||
func (b *PollingBase) requeueCardAt(cardID uint, taskType string, nextCheckAt time.Time) error {
|
||||
ctx, cancel := pollingFallbackContext()
|
||||
|
||||
@@ -106,6 +106,12 @@ func (h *PollingCardStatusHandler) Handle(ctx context.Context, task *asynq.Task)
|
||||
h.base.logger.Info("独立卡命中风险状态,已关闭轮询", zap.Uint("card_id", cardID), zap.String("gateway_extend", decision.GatewayExtend))
|
||||
return nil
|
||||
}
|
||||
|
||||
// 卡状态任务仍正常但套餐任务丢失时,仅补入缺失项;不改写已有套餐任务的执行时间。
|
||||
h.base.invalidateCardCache(ctx, cardID)
|
||||
if err := h.base.ensureMissingTask(ctx, cardID, constants.TaskTypePollingPackage); err != nil {
|
||||
h.base.logger.Warn("卡状态轮询补齐套餐任务失败", zap.Uint("card_id", cardID), zap.Error(err))
|
||||
}
|
||||
return h.base.requeueCard(ctx, cardID, constants.TaskTypePollingCardStatus)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user