Files
junhong_cmp_fiber/internal/service/package/reset_service.go
break 7029104e5c 让迁移套餐恢复月流量重置调度
缺少 next_reset_at 时,轮询根据已有激活时间或到期时间与套餐天数推算下一重置点;已有值通过查询条件和条件更新双重保护,不会被覆盖。

Constraint: 兼容迁移套餐缺少 activated_at 与 next_reset_at 的历史数据
Rejected: 单次 SQL 人工回填 | 后续迁移数据仍可能再次遗漏
Confidence: high
Scope-risk: narrow
Directive: 保持 next_reset_at 非空记录不可覆盖
Not-tested: 按用户要求未运行测试
2026-08-05 14:33:16 +08:00

338 lines
10 KiB
Go

package packagepkg
import (
"context"
"time"
"github.com/break/junhong_cmp_fiber/internal/model"
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
"github.com/break/junhong_cmp_fiber/pkg/constants"
"github.com/break/junhong_cmp_fiber/pkg/errors"
"github.com/redis/go-redis/v9"
"go.uber.org/zap"
"gorm.io/gorm"
)
type ResetService struct {
db *gorm.DB
redis *redis.Client
packageUsageStore *postgres.PackageUsageStore
logger *zap.Logger
resumeCallback ResumeCallback
}
func NewResetService(
db *gorm.DB,
redis *redis.Client,
packageUsageStore *postgres.PackageUsageStore,
logger *zap.Logger,
) *ResetService {
return &ResetService{
db: db,
redis: redis,
packageUsageStore: packageUsageStore,
logger: logger,
}
}
// SetResumeCallback 设置复机回调,流量重置后对满足条件的卡触发自动复机
func (s *ResetService) SetResumeCallback(callback ResumeCallback) {
s.resumeCallback = callback
}
// ResetDailyUsage 任务 11.2-11.3: 重置日流量
func (s *ResetService) ResetDailyUsage(ctx context.Context) error {
return s.resetDailyUsageWithDB(ctx, s.db)
}
// resetDailyUsageWithDB 内部方法,支持传入 DB/TX
func (s *ResetService) resetDailyUsageWithDB(ctx context.Context, db *gorm.DB) error {
now := time.Now()
var resetPackages []*model.PackageUsage
err := db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var packages []*model.PackageUsage
err := tx.Where("data_reset_cycle = ?", constants.PackageDataResetDaily).
Where("next_reset_at <= ?", now).
Where("status IN ?", []int{constants.PackageUsageStatusActive, constants.PackageUsageStatusDepleted}).
Find(&packages).Error
if err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询待重置套餐失败")
}
if len(packages) == 0 {
s.logger.Info("没有需要重置的日流量套餐")
return nil
}
packageIDs := make([]uint, len(packages))
for i, pkg := range packages {
packageIDs[i] = pkg.ID
}
nextReset := time.Date(now.Year(), now.Month(), now.Day()+1, 0, 0, 0, 0, now.Location())
updates := map[string]interface{}{
"data_usage_mb": 0,
"last_reset_at": now,
"next_reset_at": nextReset,
"status": constants.PackageUsageStatusActive,
}
if err := tx.Model(&model.PackageUsage{}).
Where("id IN ?", packageIDs).
Updates(updates).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "批量重置日流量失败")
}
s.logger.Info("日流量重置完成",
zap.Int("count", len(packages)),
zap.Time("next_reset_at", nextReset))
resetPackages = packages
return nil
})
if err != nil {
return err
}
s.triggerResumeForPackages(ctx, resetPackages)
return nil
}
// ResetMonthlyUsage 任务 11.4-11.5: 重置月流量
func (s *ResetService) ResetMonthlyUsage(ctx context.Context) error {
return s.resetMonthlyUsageWithDB(ctx, s.db)
}
// resetMonthlyUsageWithDB 内部方法,支持传入 DB/TX
func (s *ResetService) resetMonthlyUsageWithDB(ctx context.Context, db *gorm.DB) error {
now := time.Now()
var resetPackages []*model.PackageUsage
err := db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if err := s.backfillMissingMonthlyResetAt(tx, now); err != nil {
return err
}
var packages []*model.PackageUsage
err := tx.Where("data_reset_cycle = ?", constants.PackageDataResetMonthly).
Where("next_reset_at <= ?", now).
Where("status IN ?", []int{constants.PackageUsageStatusActive, constants.PackageUsageStatusDepleted}).
Find(&packages).Error
if err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询待重置套餐失败")
}
if len(packages) == 0 {
s.logger.Info("没有需要重置的月流量套餐")
return nil
}
for _, usage := range packages {
var pkg model.Package
if err := tx.First(&pkg, usage.PackageID).Error; err != nil {
s.logger.Error("查询套餐信息失败",
zap.Uint("usage_id", usage.ID),
zap.Uint("package_id", usage.PackageID),
zap.Error(err))
continue
}
activatedAt := now
if usage.ActivatedAt != nil && !usage.ActivatedAt.IsZero() {
activatedAt = *usage.ActivatedAt
}
nextResetAt := CalculateNextResetTime(constants.PackageDataResetMonthly, pkg.CalendarType, now, activatedAt)
if nextResetAt == nil {
s.logger.Warn("计算下次重置时间失败",
zap.Uint("usage_id", usage.ID))
continue
}
updates := map[string]interface{}{
"data_usage_mb": 0,
"last_reset_at": now,
"next_reset_at": *nextResetAt,
"status": constants.PackageUsageStatusActive,
}
if err := tx.Model(usage).Updates(updates).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "重置月流量失败")
}
s.logger.Info("月流量已重置",
zap.Uint("usage_id", usage.ID),
zap.String("calendar_type", pkg.CalendarType),
zap.Time("next_reset_at", *nextResetAt))
}
resetPackages = packages
return nil
})
if err != nil {
return err
}
s.triggerResumeForPackages(ctx, resetPackages)
return nil
}
type monthlyResetBackfillCandidate struct {
ID uint
ActivatedAt *time.Time
ExpiresAt *time.Time
CalendarType string
DurationMonths int
DurationDays int
}
// backfillMissingMonthlyResetAt 为迁移套餐补齐缺失的下次流量重置时间。
func (s *ResetService) backfillMissingMonthlyResetAt(tx *gorm.DB, now time.Time) error {
var candidates []monthlyResetBackfillCandidate
if err := tx.Table("tb_package_usage AS usage").
Select("usage.id, usage.activated_at, usage.expires_at, pkg.calendar_type, pkg.duration_months, pkg.duration_days").
Joins("JOIN tb_package AS pkg ON pkg.id = usage.package_id AND pkg.deleted_at IS NULL").
Where("usage.data_reset_cycle = ?", constants.PackageDataResetMonthly).
Where("usage.next_reset_at IS NULL").
Where("usage.status IN ?", []int{constants.PackageUsageStatusActive, constants.PackageUsageStatusDepleted}).
Find(&candidates).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询待补齐重置时间的套餐失败")
}
for _, candidate := range candidates {
nextResetAt := calculateMissingMonthlyResetAt(candidate, now)
if nextResetAt == nil {
s.logger.Warn("缺少套餐计时信息,无法补齐月流量重置时间", zap.Uint("usage_id", candidate.ID))
continue
}
result := tx.Model(&model.PackageUsage{}).
Where("id = ? AND next_reset_at IS NULL", candidate.ID).
Update("next_reset_at", *nextResetAt)
if result.Error != nil {
return errors.Wrap(errors.CodeDatabaseError, result.Error, "补齐套餐月流量重置时间失败")
}
if result.RowsAffected > 0 {
s.logger.Info("已补齐套餐月流量重置时间",
zap.Uint("usage_id", candidate.ID),
zap.Time("next_reset_at", *nextResetAt))
}
}
return nil
}
func calculateMissingMonthlyResetAt(candidate monthlyResetBackfillCandidate, now time.Time) *time.Time {
activatedAt := now
if candidate.ActivatedAt != nil && !candidate.ActivatedAt.IsZero() {
activatedAt = *candidate.ActivatedAt
} else if candidate.CalendarType == constants.PackageCalendarTypeByDay {
validityDays := calculateByDayValidityDays(candidate.DurationMonths, candidate.DurationDays)
if candidate.ExpiresAt == nil || candidate.ExpiresAt.IsZero() || validityDays <= 0 {
return nil
}
activatedAt = candidate.ExpiresAt.AddDate(0, 0, -validityDays)
}
return CalculateNextResetTime(constants.PackageDataResetMonthly, candidate.CalendarType, now, activatedAt)
}
// ResetYearlyUsage 任务 11.6-11.7: 重置年流量
func (s *ResetService) ResetYearlyUsage(ctx context.Context) error {
return s.resetYearlyUsageWithDB(ctx, s.db)
}
// resetYearlyUsageWithDB 内部方法,支持传入 DB/TX
func (s *ResetService) resetYearlyUsageWithDB(ctx context.Context, db *gorm.DB) error {
now := time.Now()
var resetPackages []*model.PackageUsage
err := db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var packages []*model.PackageUsage
err := tx.Where("data_reset_cycle = ?", constants.PackageDataResetYearly).
Where("next_reset_at <= ?", now).
Where("status IN ?", []int{constants.PackageUsageStatusActive, constants.PackageUsageStatusDepleted}).
Find(&packages).Error
if err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "查询待重置套餐失败")
}
if len(packages) == 0 {
s.logger.Info("没有需要重置的年流量套餐")
return nil
}
packageIDs := make([]uint, len(packages))
for i, pkg := range packages {
packageIDs[i] = pkg.ID
}
nextReset := time.Date(now.Year()+1, 1, 1, 0, 0, 0, 0, now.Location())
updates := map[string]interface{}{
"data_usage_mb": 0,
"last_reset_at": now,
"next_reset_at": nextReset,
"status": constants.PackageUsageStatusActive,
}
if err := tx.Model(&model.PackageUsage{}).
Where("id IN ?", packageIDs).
Updates(updates).Error; err != nil {
return errors.Wrap(errors.CodeDatabaseError, err, "批量重置年流量失败")
}
s.logger.Info("年流量重置完成",
zap.Int("count", len(packages)),
zap.Time("next_reset_at", nextReset))
resetPackages = packages
return nil
})
if err != nil {
return err
}
s.triggerResumeForPackages(ctx, resetPackages)
return nil
}
// triggerResumeForPackages 对重置后的套餐,遍历其载体并触发自动复机
// 流量重置后,因流量耗尽而停机的卡满足复机条件,此处异步触发避免阻塞重置流程
func (s *ResetService) triggerResumeForPackages(ctx context.Context, packages []*model.PackageUsage) {
if s.resumeCallback == nil || len(packages) == 0 {
return
}
for _, pkg := range packages {
var carrierType string
var carrierID uint
if pkg.DeviceID > 0 {
carrierType = "device"
carrierID = pkg.DeviceID
} else if pkg.IotCardID > 0 {
carrierType = "iot_card"
carrierID = pkg.IotCardID
} else {
continue
}
go func(ct string, cid uint) {
if err := s.resumeCallback.ResumeCardIfStopped(context.Background(), ct, cid); err != nil {
s.logger.Warn("流量重置后自动复机失败",
zap.String("carrier_type", ct),
zap.Uint("carrier_id", cid),
zap.Error(err))
}
}(carrierType, carrierID)
}
}