避免套餐过期后排队权益永久失联
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m6s

线上保持现有纯 Asynq 架构,以公平孤儿扫描和提交后投递消除永久饥饿及事务可见性竞态。

Constraint: 线上保持现有纯 Asynq 架构,不引入 Outbox、迁移或新任务基础设施。

Rejected: 事务内投递或扩大扫描 LIMIT | 无法消除竞态和永久饥饿。

Confidence: high

Scope-risk: narrow

Directive: 后续分支整合时按目标分支的套餐接续架构独立处理,不混用本热修实现。

Tested: go build ./...(退出码 0);git diff --check;openspec validate fix-main-package-activation-starvation --strict。

Not-tested: 按用户要求未新增、修改或运行自动化测试;线上 SQL、查询计划和日志待部署后核验。
This commit is contained in:
2026-08-03 09:58:05 +08:00
parent 1efb665619
commit a0de08d789
9 changed files with 369 additions and 103 deletions

View File

@@ -17,6 +17,9 @@ import (
"github.com/break/junhong_cmp_fiber/pkg/errors"
)
// orphanPackageScanLimit 单轮最多恢复的真实孤儿载体数量。
const orphanPackageScanLimit = 100
// TrafficSyncer 套餐失效前流量同步接口
// 在主套餐标记为已过期之前,从 Gateway 拉取最新流量写入 DB
type TrafficSyncer interface {
@@ -156,17 +159,48 @@ func (h *PackageActivationHandler) HandlePackageActivationCheck(ctx context.Cont
}
// findAndActivateOrphanPackages 任务 6.1-6.2: 查找孤儿载体并触发激活
// 孤儿定义:存在 status=0 的主套餐,但不存在 status=1 的主套餐。
// 孤儿定义:存在待生效主套餐,但不存在生效中或已用完的占位主套餐。
func (h *PackageActivationHandler) findAndActivateOrphanPackages(ctx context.Context) (int, error) {
// 查询孤儿待生效主套餐(无生效主套餐但有待生效主套餐)
var orphanUsages []*model.PackageUsage
err := h.db.WithContext(ctx).
Where("status = ?", constants.PackageUsageStatusPending).
Where("master_usage_id IS NULL").
Where("deleted_at IS NULL").
Order("priority ASC, created_at ASC").
Limit(100).
Find(&orphanUsages).Error
err := h.db.WithContext(ctx).Raw(`
WITH pending_queue AS (
SELECT pending.id,
pending.priority,
pending.created_at,
ROW_NUMBER() OVER (
PARTITION BY
CASE WHEN COALESCE(pending.iot_card_id, 0) > 0 THEN 'iot_card' ELSE 'device' END,
CASE WHEN COALESCE(pending.iot_card_id, 0) > 0 THEN pending.iot_card_id ELSE pending.device_id END
ORDER BY pending.priority ASC, pending.created_at ASC, pending.id ASC
) AS queue_position
FROM tb_package_usage AS pending
WHERE pending.status = ?
AND pending.master_usage_id IS NULL
AND pending.deleted_at IS NULL
AND (COALESCE(pending.iot_card_id, 0) > 0 OR COALESCE(pending.device_id, 0) > 0)
AND NOT EXISTS (
SELECT 1
FROM tb_package_usage AS occupied
WHERE occupied.status IN (?, ?)
AND occupied.master_usage_id IS NULL
AND occupied.deleted_at IS NULL
AND (
(COALESCE(pending.iot_card_id, 0) > 0 AND occupied.iot_card_id = pending.iot_card_id)
OR (COALESCE(pending.iot_card_id, 0) = 0 AND pending.device_id > 0 AND occupied.device_id = pending.device_id)
)
)
)
SELECT usage.*
FROM pending_queue AS candidate
JOIN tb_package_usage AS usage ON usage.id = candidate.id
WHERE candidate.queue_position = 1
ORDER BY candidate.priority ASC, candidate.created_at ASC, candidate.id ASC
LIMIT ?`,
constants.PackageUsageStatusPending,
constants.PackageUsageStatusActive,
constants.PackageUsageStatusDepleted,
orphanPackageScanLimit,
).Scan(&orphanUsages).Error
if err != nil {
return 0, errors.Wrap(errors.CodeDatabaseError, err, "查询孤儿套餐失败")
}
@@ -175,75 +209,14 @@ func (h *PackageActivationHandler) findAndActivateOrphanPackages(ctx context.Con
return 0, nil
}
// 按载体分组,去重
type carrierKey struct {
carrierType string
carrierID uint
}
carrierMap := make(map[carrierKey]*model.PackageUsage) // 保留 priority 最低的套餐
for _, usage := range orphanUsages {
key := carrierKey{}
if usage.IotCardID > 0 {
key.carrierType = "iot_card"
key.carrierID = usage.IotCardID
} else if usage.DeviceID > 0 {
key.carrierType = "device"
key.carrierID = usage.DeviceID
} else {
continue
}
// 检查该载体是否已有占位主套餐(生效中或已用完均视为占位,不允许激活待生效套餐)
var activeCount int64
var countErr error
occupiedStatuses := []int{constants.PackageUsageStatusActive, constants.PackageUsageStatusDepleted}
if key.carrierType == "iot_card" {
countErr = h.db.WithContext(ctx).
Model(&model.PackageUsage{}).
Where("status IN ?", occupiedStatuses).
Where("master_usage_id IS NULL").
Where("iot_card_id = ?", key.carrierID).
Count(&activeCount).Error
} else {
countErr = h.db.WithContext(ctx).
Model(&model.PackageUsage{}).
Where("status IN ?", occupiedStatuses).
Where("master_usage_id IS NULL").
Where("device_id = ?", key.carrierID).
Count(&activeCount).Error
}
if countErr != nil {
h.logger.Warn("检查载体生效套餐失败",
zap.String("carrier_type", key.carrierType),
zap.Uint("carrier_id", key.carrierID),
zap.Error(countErr))
continue
}
// 已有生效或已用完(占位)套餐,跳过
if activeCount > 0 {
continue
}
// 保留购买顺序最靠前的套餐,队首未满足实名条件时不跳过。
if existing, ok := carrierMap[key]; ok {
if usage.Priority < existing.Priority || (usage.Priority == existing.Priority && usage.CreatedAt.Before(existing.CreatedAt)) {
carrierMap[key] = usage
}
} else {
carrierMap[key] = usage
}
}
// 为每个孤儿载体提交激活任务
count := 0
for key, usage := range carrierMap {
if err := h.enqueueActivationTask(ctx, usage.ID, key.carrierType, key.carrierID, "orphan_recovery"); err != nil {
for _, usage := range orphanUsages {
carrierType, carrierID := h.getCarrierInfo(usage)
if err := h.enqueueActivationTask(ctx, usage.ID, carrierType, carrierID, "orphan_recovery"); err != nil {
h.logger.Warn("提交孤儿套餐激活任务失败",
zap.Uint("package_usage_id", usage.ID),
zap.String("carrier_type", key.carrierType),
zap.Uint("carrier_id", key.carrierID),
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID),
zap.Error(err))
continue
}
@@ -269,7 +242,7 @@ func (h *PackageActivationHandler) findExpiredMainPackages(ctx context.Context)
}
// processExpiredPackage 处理单个过期套餐
// 流程:先同步最新流量 → 事务内标记过期/失效/激活下一个事务提交后触发停机
// 流程:先同步最新流量 → 事务内标记过期失效加油包 → 提交后投递下一套餐并触发停机检查
func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pkg *model.PackageUsage) error {
carrierType, carrierID := h.getCarrierInfo(pkg)
@@ -294,20 +267,7 @@ func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pk
// 任务 19.4: 加油包级联失效
if err := h.invalidateAddons(ctx, tx, pkg.ID); err != nil {
h.logger.Warn("加油包级联失效失败",
zap.Uint("master_usage_id", pkg.ID),
zap.Error(err))
}
// 任务 19.5: 查询并激活下一个待生效主套餐
if carrierType != "" && carrierID > 0 {
if err := h.activateNextPackage(ctx, tx, carrierType, carrierID); err != nil {
h.logger.Warn("激活下一个待生效套餐失败",
zap.String("carrier_type", carrierType),
zap.Uint("carrier_id", carrierID),
zap.Error(err))
}
return err
}
return nil
@@ -317,9 +277,11 @@ func (h *PackageActivationHandler) processExpiredPackage(ctx context.Context, pk
return err
}
// 事务提交后再触发异步停机,确保 CheckAndStopCard 读到最新的套餐状态
// 事务提交后再投递,确保消费者只能读取到旧套餐已经过期的状态
if carrierType != "" && carrierID > 0 {
activationErr := h.activateNextPackage(ctx, h.db.WithContext(ctx), carrierType, carrierID)
h.triggerStopAfterExpiry(ctx, carrierType, carrierID)
return activationErr
}
return nil
@@ -368,7 +330,7 @@ func (h *PackageActivationHandler) activateNextPackage(ctx context.Context, tx *
var nextPkg model.PackageUsage
query := tx.Where("status = ?", constants.PackageUsageStatusPending).
Where("master_usage_id IS NULL"). // 主套餐
Order("priority ASC").
Order("priority ASC, created_at ASC, id ASC").
Limit(1)
if carrierType == "iot_card" {
@@ -496,12 +458,19 @@ func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Cont
// 调用 ActivationService 执行激活
if h.activationService != nil {
if err := h.activationService.ActivateSpecificPackage(ctx, payload.PackageUsageID); err != nil {
activated, err := h.activationService.ActivateSpecificPackage(ctx, payload.PackageUsageID)
if err != nil {
h.logger.Error("套餐激活失败",
zap.Uint("package_usage_id", payload.PackageUsageID),
zap.Error(err))
return err
}
if !activated {
h.logger.Info("套餐本次未激活,等待后续检查",
zap.Uint("package_usage_id", payload.PackageUsageID),
zap.String("activation_type", payload.ActivationType))
return nil
}
} else {
// ActivationService 未注入,无法安全激活(缺少 expires_at 计算)
h.logger.Error("激活服务未注入,无法执行套餐激活",
@@ -510,7 +479,8 @@ func (h *PackageActivationHandler) HandlePackageQueueActivation(ctx context.Cont
}
h.logger.Info("套餐激活成功",
zap.Uint("package_usage_id", payload.PackageUsageID))
zap.Uint("package_usage_id", payload.PackageUsageID),
zap.String("activation_type", payload.ActivationType))
return nil
}