diff --git a/README.md b/README.md index 2d19ddb..c850030 100644 --- a/README.md +++ b/README.md @@ -26,6 +26,10 @@ 后台设备列表 `GET /api/admin/devices` 支持按 `virtual_no`、`imei`、设备名称、归属状态、激活状态、店铺、套餐系列等条件筛选。 +### 套餐接续可靠性 + +`main` 分支使用 Asynq 在旧主套餐事务提交后投递下一套餐,并通过公平孤儿扫描补偿漏投;部署与回滚说明见 [Main 分支套餐接续饥饿热修](docs/hotfix-main-package-activation-recovery/功能总结.md)。 + ### 三种客户类型 | 客户类型 | 业务特点 | 典型场景 | 钱包归属 | diff --git a/docs/hotfix-main-package-activation-recovery/功能总结.md b/docs/hotfix-main-package-activation-recovery/功能总结.md new file mode 100644 index 0000000..016d4d8 --- /dev/null +++ b/docs/hotfix-main-package-activation-recovery/功能总结.md @@ -0,0 +1,44 @@ +# Main 分支套餐接续饥饿热修 + +## 修复边界 + +本热修仅适用于 `main` 的纯 Asynq 套餐接续链路,不引入 Outbox、数据库迁移、新任务类型或新依赖。 + +- 孤儿扫描先在 PostgreSQL 中按卡或设备选择队首套餐,并排除仍有 `status IN (1,2)` 占位主套餐的载体,最后取 100 个真实孤儿。 +- 旧主套餐及加油包状态事务提交后,再投递现有 `package:queue:activation` 任务。 +- Redis 激活锁冲突返回套餐激活冲突错误,由现有 `MaxRetry(3)` 重试。 +- 只有套餐实际从待生效推进为生效中时,Handler 才记录“套餐激活成功”。 + +## 部署观察 + +部署 Worker 后至少观察两个套餐轮询周期: + +1. `孤儿套餐扫描完成` 的 `orphan_count` 应能覆盖真实无占位套餐,不再固定被同一批占位载体挡住。 +2. `已提交套餐激活任务` 后应出现实际激活、明确跳过或可重试错误,不再出现未改状态却打印成功。 +3. Redis 锁冲突应进入 Asynq 重试,不应确认任务成功。 + +可使用以下只读 SQL 检查仍未恢复的真实孤儿数量: + +```sql +SELECT COUNT(*) AS orphan_pending_count +FROM tb_package_usage AS pending +WHERE pending.status = 0 + 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 (1, 2) + 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) + ) + ); +``` + +## 回滚 + +本次无数据库迁移。回滚热修提交并重新部署 Worker 即可;已经正确激活的套餐属于有效业务事实,不执行反向 SQL。 diff --git a/internal/polling/package_activation_handler.go b/internal/polling/package_activation_handler.go index a4642ad..fe84e6d 100644 --- a/internal/polling/package_activation_handler.go +++ b/internal/polling/package_activation_handler.go @@ -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 } diff --git a/internal/service/package/activation_service.go b/internal/service/package/activation_service.go index e6a1a90..01d41a8 100644 --- a/internal/service/package/activation_service.go +++ b/internal/service/package/activation_service.go @@ -228,14 +228,14 @@ func (s *ActivationService) ActivateQueuedPackage(ctx context.Context, carrierTy // ActivateSpecificPackage 任务 4: 激活指定的套餐使用记录 // 根据 PackageUsageID 精准激活目标套餐,而非重跑"查找过期包"流程 -func (s *ActivationService) ActivateSpecificPackage(ctx context.Context, packageUsageID uint) error { +func (s *ActivationService) ActivateSpecificPackage(ctx context.Context, packageUsageID uint) (bool, error) { // 加载 PackageUsage var usage model.PackageUsage if err := s.db.WithContext(ctx).First(&usage, packageUsageID).Error; err != nil { if err == gorm.ErrRecordNotFound { - return errors.New(errors.CodeNotFound, "套餐使用记录不存在") + return false, errors.New(errors.CodeNotFound, "套餐使用记录不存在") } - return errors.Wrap(errors.CodeDatabaseError, err, "查询套餐使用记录失败") + return false, errors.Wrap(errors.CodeDatabaseError, err, "查询套餐使用记录失败") } // 幂等检查:非待生效状态直接返回 @@ -243,7 +243,7 @@ func (s *ActivationService) ActivateSpecificPackage(ctx context.Context, package s.logger.Info("套餐无需激活(非待生效状态)", zap.Uint("usage_id", usage.ID), zap.Int("status", usage.Status)) - return nil + return false, nil } // 确定载体类型和 ID @@ -256,7 +256,7 @@ func (s *ActivationService) ActivateSpecificPackage(ctx context.Context, package carrierType = "device" carrierID = usage.DeviceID } else { - return errors.New(errors.CodeInvalidParam, "套餐使用记录缺少载体信息") + return false, errors.New(errors.CodeInvalidParam, "套餐使用记录缺少载体信息") } // 获取分布式锁 @@ -264,24 +264,25 @@ func (s *ActivationService) ActivateSpecificPackage(ctx context.Context, package lockValue := time.Now().String() locked, err := s.redis.SetNX(ctx, lockKey, lockValue, 30*time.Second).Result() if err != nil { - return errors.Wrap(errors.CodeRedisError, err, "获取分布式锁失败") + return false, errors.Wrap(errors.CodeRedisError, err, "获取分布式锁失败") } if !locked { s.logger.Warn("套餐激活正在进行中,跳过", zap.String("carrier_type", carrierType), zap.Uint("carrier_id", carrierID)) - return nil + return false, errors.New(errors.CodePackageActivationConflict) } defer s.redis.Del(ctx, lockKey) // 加载关联 Package var pkg model.Package if err := s.db.WithContext(ctx).First(&pkg, usage.PackageID).Error; err != nil { - return errors.Wrap(errors.CodeDatabaseError, err, "查询套餐信息失败") + return false, errors.Wrap(errors.CodeDatabaseError, err, "查询套餐信息失败") } // 事务内激活 - return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + activated := false + err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { // 再次幂等检查(防止锁释放期间被其他操作激活) var currentUsage model.PackageUsage if err := tx.First(¤tUsage, packageUsageID).Error; err != nil { @@ -319,6 +320,7 @@ func (s *ActivationService) ActivateSpecificPackage(ctx context.Context, package if err := s.activatePendingUsage(ctx, tx, ¤tUsage, &pkg, carrierType, carrierID, time.Now(), "指定套餐已激活"); err != nil { return err } + activated = true // 异步复机 if s.resumeCallback != nil { @@ -334,6 +336,11 @@ func (s *ActivationService) ActivateSpecificPackage(ctx context.Context, package return nil }) + if err != nil { + return false, err + } + + return activated, nil } // ActivateNextPendingMainPackage 按购买顺序激活载体的下一个待生效主套餐。 diff --git a/openspec/changes/fix-main-package-activation-starvation/.openspec.yaml b/openspec/changes/fix-main-package-activation-starvation/.openspec.yaml new file mode 100644 index 0000000..e08b5f8 --- /dev/null +++ b/openspec/changes/fix-main-package-activation-starvation/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-08-03 diff --git a/openspec/changes/fix-main-package-activation-starvation/design.md b/openspec/changes/fix-main-package-activation-starvation/design.md new file mode 100644 index 0000000..021b9f0 --- /dev/null +++ b/openspec/changes/fix-main-package-activation-starvation/design.md @@ -0,0 +1,106 @@ +## Context + +线上只读数据得到旧孤儿扫描窗口 `100/100/0`:固定取出的 100 条待生效记录全部仍有 `status IN (1,2)` 占位主套餐,而真正无占位套餐的记录位于窗口之外。当前实现先 `LIMIT 100`,再逐条查询占位状态,因此同一批无效候选会永久挡住真实孤儿。 + +过期路径还在数据库事务内调用 `enqueueActivationTask`。Asynq 消费者可能早于事务提交读取旧主套餐,随后 `ActivateSpecificPackage` 判断“已有生效主套餐”并返回 `nil`;Handler 继续记录“套餐激活成功”,任务不再重试。Redis 锁冲突也返回 `nil`,存在另一条假成功路径。 + +本热修保持当前线上 `Polling Handler → GORM transaction → Asynq → Package Activation Service` 架构,不引入新的可靠投递设施。 + +## Goals / Non-Goals + +**Goals:** + +- 每轮 100 个恢复名额只用于真实孤儿载体,同一载体只选择队首套餐。 +- 旧套餐过期事实提交后才投递 Asynq,消除事务可见性竞态。 +- Redis 锁冲突返回错误,由现有 Asynq 重试。 +- 成功日志只对应实际激活或明确的已完成幂等事实。 + +**Non-Goals:** + +- 不新增 Outbox、迁移、索引、依赖、队列或任务类型。 +- 不重构套餐购买、实名激活、流量扣减、退款和停复机。 +- 不修改 API、DTO、路由或前端。 +- 不批量修复历史数据。 +- 按用户要求,不新增、修改或运行自动化测试。 + +## Decisions + +### 决策 1:数据库先筛选真实孤儿队首,再执行 LIMIT + +使用 GORM `Raw` 执行 PostgreSQL CTE/窗口查询: + +1. 从有效 `status=0` 主套餐按卡/设备载体分组。 +2. 每组按 `priority ASC, created_at ASC, id ASC` 选 `ROW_NUMBER()=1`。 +3. 使用相关 `NOT EXISTS` 排除同载体 `status IN (1,2)` 主套餐。 +4. 稳定排序后 `LIMIT 100`。 + +删除现有逐条 `Count` 和 Go map 分组。查询仍返回完整 `PackageUsage`,沿用现有投递循环。 + +**拒绝:扩大 LIMIT。** 只会推迟复现并增加 N+1。 + +**拒绝:分页遍历所有 pending。** 需要游标状态,复杂度高于一次正确查询。 + +### 决策 2:过期事务提交后再投递现有 Asynq + +`processExpiredPackage` 事务只更新旧主套餐和关联加油包。提交成功后,使用普通数据库句柄调用现有 `activateNextPackage` 查询队首并入队。 + +提交后入队失败时返回错误并记录上下文;同一轮末尾及后续轮询的真实孤儿扫描会再次发现该载体,提供持久状态驱动的补偿。 + +**拒绝:任务增加固定延迟。** 固定延迟不能证明事务已提交。 + +**拒绝:引入 Outbox。** 当前线上没有该基础设施,热修不扩张架构。 + +### 决策 3:锁冲突必须触发 Asynq 重试 + +`ActivateSpecificPackage` 未取得 `RedisPackageActivationLockKey` 时返回现有 `CodePackageActivationConflict`。Handler 原样返回错误,由任务已有 `MaxRetry(3)` 处理。 + +Redis Key 继续使用 `pkg/constants/redis.go` 的生成函数,不新增硬编码 Key。 + +### 决策 4:显式返回本次是否激活 + +`ActivateSpecificPackage` 返回 `(bool, error)`: + +- `true,nil`:本次把待生效套餐推进为生效中; +- `false,nil`:记录已非待生效、存在占位套餐或条件暂不满足; +- `false,error`:数据库、Redis 或锁冲突,应由任务重试或记录失败。 + +Handler 仅在 `true,nil` 时记录“套餐激活成功”。Handler 已在调用前识别 `status=1` 的重复任务并记录幂等跳过。 + +### 决策 5:依赖注入和事务边界保持不变 + +`PackageActivationHandler` 继续通过结构体字段持有 `*gorm.DB`、`*redis.Client`、`*asynq.Client`、`*ActivationService` 和 Zap Logger;不新增单实现接口或工厂。套餐激活仍由 Service 自己开启 GORM 事务,Handler 不直接更新新套餐状态。 + +### 决策 6:公共能力与验证 + +- Audit Event:N/A,系统自动生命周期推进。 +- Domain Ledger:N/A,`tb_package_usage` 是权威事实。 +- Integration Log:N/A,无新增外部调用。 +- Outbox:N/A,保持当前 Asynq + 周期自愈。 + +按用户要求不写或运行自动化测试。验证使用 `gofmt`、`git diff --check`、`go build ./...`、只读 SQL、查询计划和日志检查。 + +## Risks / Trade-offs + +- **[风险] CTE 扫描大量 pending** → 用 `EXPLAIN (ANALYZE, BUFFERS)` 验证;无证据不新增索引。 +- **[风险] 提交后、入队前进程退出** → 下一轮真实孤儿扫描恢复,最长增加一个轮询周期。 +- **[风险] 多实例重复入队** → 载体 Redis 锁和套餐状态幂等保证只实际激活一次。 +- **[风险] 方法签名变化遗漏调用点** → 使用 `rg` 检查全部调用方并以全量构建证明编译契约。 +- **[权衡] Asynq 最终失败后仍依赖轮询重新入队** → 这是当前线上架构的既有补偿边界,本热修不扩建基础设施。 +- **[权衡] 不新增自动化测试** → 遵循用户边界,以构建、SQL 和日志证据替代。 + +## Migration Plan + +1. 实施真实孤儿查询、提交后入队、锁冲突重试和准确日志。 +2. 执行格式化、静态检查、`go build ./...` 和只读 SQL语义检查。 +3. 形成独立中文 Lore 热修提交。 +4. 部署 Worker,观察至少两个轮询周期内真实孤儿收敛和任务日志。 + +### 回滚 + +- 无数据库迁移,revert 热修提交并重新部署 Worker。 +- 已正确激活的套餐保持业务事实,不执行反向 SQL。 +- 回滚后新增孤儿继续使用带状态保护的单卡 SQL逐条恢复。 + +## Open Questions + +无。 diff --git a/openspec/changes/fix-main-package-activation-starvation/proposal.md b/openspec/changes/fix-main-package-activation-starvation/proposal.md new file mode 100644 index 0000000..b821ae9 --- /dev/null +++ b/openspec/changes/fix-main-package-activation-starvation/proposal.md @@ -0,0 +1,36 @@ +## Why + +功能 ID:`hotfix-main-package-activation-recovery` + +线上已第二次出现原主套餐成功过期、队首待生效套餐仍长期停留在 `status=0` 的故障。只读诊断确认当前孤儿扫描固定取出的 100 条记录全部仍有占位套餐,真实孤儿永远无法进入恢复窗口;同时,过期事务提交前投递 Asynq 会让消费者读到旧套餐仍为生效中并把跳过误判为成功。 + +## What Changes + +- PostgreSQL 在 `LIMIT 100` 前完成每个载体队首选择和 `status IN (1,2)` 占位排除,删除逐条检查的 N+1 查询。 +- 旧主套餐过期和加油包失效事务提交成功后,才查询队首套餐并投递现有 `package:queue:activation` Asynq 任务。 +- Redis 激活锁冲突返回现有套餐激活冲突错误,让 Asynq 按既有策略重试,不再确认假成功。 +- Asynq Handler 仅在实际推进套餐或确认已完成幂等事实时记录成功;占位阻塞、等待实名和锁冲突记录明确原因。 +- 不新增 Outbox、数据库迁移、依赖或新任务类型,保持线上现有纯 Asynq 架构。 +- 按用户明确要求,不新增、修改或运行自动化测试;使用全量构建、只读 SQL、查询计划和日志核验。 + +## Capabilities + +### New Capabilities + +无。 + +### Modified Capabilities + +- `package-queue-activation`:补充纯 Asynq 过期接续的提交后投递、真实孤儿公平扫描、锁冲突重试和准确成功语义。 + +## Impact + +- **适用范围**:仅当前 `main` 线上代码。 +- **架构通道**:旧套餐轮询复杂写用例,沿用 `Polling Handler → GORM transaction → Asynq → Package Activation Service`;不迁移未触碰的套餐模块。 +- **代码**:`internal/polling/package_activation_handler.go`、`internal/service/package/activation_service.go`。 +- **数据库**:无迁移;仅调整 `tb_package_usage` 查询顺序与过滤。 +- **API/前端**:无改动。 +- **依赖**:继续使用 GORM、PostgreSQL、Redis、Asynq 和 Zap,不新增依赖。 +- **性能**:孤儿扫描由最多 101 次查询收敛为一次候选查询和有限任务投递;使用查询计划确认数据库耗时。 +- **审计与可靠性**:Audit Event、Domain Ledger、Integration Log、Outbox 均为 N/A;`tb_package_usage` 是权威状态,现有 Asynq + 周期孤儿扫描负责最终恢复。 +- **验证**:不写自动化测试;执行 `gofmt`、静态检查、`go build ./...`、只读 SQL及日志核验。 diff --git a/openspec/changes/fix-main-package-activation-starvation/specs/package-queue-activation/spec.md b/openspec/changes/fix-main-package-activation-starvation/specs/package-queue-activation/spec.md new file mode 100644 index 0000000..3f77be6 --- /dev/null +++ b/openspec/changes/fix-main-package-activation-starvation/specs/package-queue-activation/spec.md @@ -0,0 +1,77 @@ +## MODIFIED Requirements + +### Requirement: 当前主套餐过期后自动激活下一个 + +系统 SHALL 在生效中或已用完主套餐到期时,先提交旧主套餐及关联加油包的状态事务,再投递同一载体队首待生效套餐的现有 Asynq 激活任务;系统 MUST NOT 在旧套餐事务提交前投递任务。 + +#### Scenario: 事务提交后投递队首套餐 + +- **WHEN** 轮询处理一个到期的 `status=1` 或 `status=2` 主套餐,且存在队首待生效套餐 +- **THEN** 系统先提交旧主套餐 `status=3` 和关联加油包 `status=4` 的事务 +- **AND** 提交成功后查询稳定队首并投递 `package:queue:activation` +- **AND** 消费者读取时旧主套餐不再处于 `status IN (1,2)` + +#### Scenario: 过期事务失败 + +- **WHEN** 更新旧主套餐或关联加油包失败 +- **THEN** 事务回滚 +- **AND** 系统不得投递下一套餐激活任务 +- **AND** 下一轮过期扫描仍可重新处理 + +#### Scenario: 提交后入队失败 + +- **WHEN** 旧套餐事务已经提交,但 Asynq 入队失败或进程退出 +- **THEN** 队首套餐保持 `status=0` +- **AND** 周期孤儿扫描重新发现并补投该套餐 + +#### Scenario: 无待生效套餐 + +- **WHEN** 旧套餐过期提交后不存在有效队首待生效套餐 +- **THEN** 系统不投递激活任务 +- **AND** 沿用既有无套餐停机检查 + +## ADDED Requirements + +### Requirement: 孤儿待生效套餐必须公平恢复 + +系统 SHALL 在数据库中先选择每个卡或设备载体唯一的队首待生效主套餐,并排除仍有 `status IN (1,2)` 占位主套餐的载体,最后才执行单轮 100 个真实孤儿上限。队首顺序 MUST 为 `priority ASC, created_at ASC, id ASC`。 + +#### Scenario: 固定窗口全部为非孤儿 + +- **WHEN** 排序靠前的 100 条待生效记录均有占位主套餐,窗口之后存在真实孤儿 +- **THEN** 数据库先排除前 100 条非孤儿 +- **AND** 窗口后的真实孤儿进入本轮候选并被投递 + +#### Scenario: 同一载体有多条 pending + +- **WHEN** 一个真实孤儿载体存在多条待生效主套餐 +- **THEN** 本轮只选择 priority 最小、created_at 最早、id 最小的一条 +- **AND** 该载体只占一个恢复名额 + +#### Scenario: 生效中或已用完套餐占位 + +- **WHEN** pending 所属载体存在 `status=1` 或 `status=2` 主套餐 +- **THEN** 该载体不得进入孤儿候选 + +### Requirement: Asynq 激活结果必须准确且可重试 + +系统 SHALL 仅在本次实际把套餐推进为 `status=1` 时记录新激活成功。Redis 激活锁冲突 MUST 返回错误,使 Asynq 按既有 `MaxRetry(3)` 重试,不得确认未执行任务成功。 + +#### Scenario: 锁冲突触发重试 + +- **WHEN** 消费者未取得载体级套餐激活锁 +- **THEN** Service 返回套餐激活冲突错误 +- **AND** Handler 将错误返回 Asynq +- **AND** 本次不记录激活成功 + +#### Scenario: 本次实际激活 + +- **WHEN** 套餐为待生效、载体无占位主套餐且满足激活条件 +- **THEN** Service 返回 `activated=true` +- **AND** Handler 记录包含套餐使用记录和触发类型的成功日志 + +#### Scenario: 条件暂不满足 + +- **WHEN** Service 复检发现占位套餐或等待实名条件 +- **THEN** Service 返回 `activated=false`且不修改套餐 +- **AND** Handler 记录未激活原因,不记录成功 diff --git a/openspec/changes/fix-main-package-activation-starvation/tasks.md b/openspec/changes/fix-main-package-activation-starvation/tasks.md new file mode 100644 index 0000000..5484557 --- /dev/null +++ b/openspec/changes/fix-main-package-activation-starvation/tasks.md @@ -0,0 +1,20 @@ +## 1. Main 分支隔离与基线 + +- [x] 1.1 确认变更仅面向 `main` 的纯 Asynq 套餐接续链路,并保持 Outbox 与新任务基础设施为非目标。 +- [x] 1.2 检索套餐激活方法的全部调用点及现有错误码、Redis 键和任务重试配置,锁定最小修改边界。 + +## 2. 过期接续与孤儿恢复 + +- [x] 2.1 将孤儿恢复改为数据库先按载体选择队首并排除 `status IN (1,2)` 占位套餐,最后限制 100 个真实孤儿,删除逐条占位查询。 +- [x] 2.2 将旧主套餐过期后的下一套餐投递移到事务提交后,保留现有停机检查和周期孤儿补偿。 + +## 3. 激活结果与重试语义 + +- [x] 3.1 让指定套餐激活显式返回是否实际激活,并在 Redis 激活锁冲突时返回现有套餐激活冲突错误。 +- [x] 3.2 调整 Asynq Handler 日志,仅在实际激活时记录成功,未激活时记录明确的跳过信息。 + +## 4. 文档、验证与提交 + +- [x] 4.1 更新功能总结和 README,说明 `main` 纯 Asynq 热修边界、部署观察项及回滚方式。 +- [x] 4.2 执行 `gofmt`、`git diff --check`、`go build ./...` 和 OpenSpec 严格校验;按用户要求不新增、修改或运行自动化测试。 +- [x] 4.3 仅暂存本热修代码、文档和独立 OpenSpec,创建符合 Lore 协议的中文提交。