fix: 套餐激活和流量重置后注入并调用复机回调
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
@@ -13,11 +13,10 @@ import (
|
|||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
)
|
)
|
||||||
|
|
||||||
// ResumeCallback 任务 24.7: 复机回调接口
|
// ResumeCallback 复机回调接口
|
||||||
// 用于在套餐激活后触发自动复机
|
// 用于在套餐激活或流量重置后触发自动复机
|
||||||
type ResumeCallback interface {
|
type ResumeCallback interface {
|
||||||
// ResumeCardIfStopped 购买套餐后自动复机
|
ResumeCardIfStopped(ctx context.Context, carrierType string, carrierID uint) error
|
||||||
ResumeCardIfStopped(ctx context.Context, cardID uint) error
|
|
||||||
}
|
}
|
||||||
|
|
||||||
type ActivationService struct {
|
type ActivationService struct {
|
||||||
@@ -149,16 +148,15 @@ func (s *ActivationService) ActivateByRealname(ctx context.Context, carrierType
|
|||||||
zap.Time("activated_at", activatedAt),
|
zap.Time("activated_at", activatedAt),
|
||||||
zap.Time("expires_at", expiresAt))
|
zap.Time("expires_at", expiresAt))
|
||||||
|
|
||||||
// 任务 24.7: 在套餐激活后触发自动复机
|
if s.resumeCallback != nil {
|
||||||
if s.resumeCallback != nil && carrierType == "iot_card" {
|
go func(ct string, cid uint) {
|
||||||
go func(cardID uint) {
|
if err := s.resumeCallback.ResumeCardIfStopped(context.Background(), ct, cid); err != nil {
|
||||||
resumeCtx := context.Background()
|
s.logger.Error("实名激活后自动复机失败",
|
||||||
if err := s.resumeCallback.ResumeCardIfStopped(resumeCtx, cardID); err != nil {
|
zap.String("carrier_type", ct),
|
||||||
s.logger.Error("自动复机失败",
|
zap.Uint("carrier_id", cid),
|
||||||
zap.Uint("card_id", cardID),
|
|
||||||
zap.Error(err))
|
zap.Error(err))
|
||||||
}
|
}
|
||||||
}(carrierID)
|
}(carrierType, carrierID)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -310,16 +308,15 @@ func (s *ActivationService) activateNextMainPackage(ctx context.Context, tx *gor
|
|||||||
zap.Time("activated_at", activatedAt),
|
zap.Time("activated_at", activatedAt),
|
||||||
zap.Time("expires_at", expiresAt))
|
zap.Time("expires_at", expiresAt))
|
||||||
|
|
||||||
// 任务 24.7: 在套餐激活后触发自动复机
|
if s.resumeCallback != nil {
|
||||||
if s.resumeCallback != nil && carrierType == "iot_card" {
|
go func(ct string, cid uint) {
|
||||||
go func(cardID uint) {
|
if err := s.resumeCallback.ResumeCardIfStopped(context.Background(), ct, cid); err != nil {
|
||||||
resumeCtx := context.Background()
|
|
||||||
if err := s.resumeCallback.ResumeCardIfStopped(resumeCtx, cardID); err != nil {
|
|
||||||
s.logger.Error("排队激活后自动复机失败",
|
s.logger.Error("排队激活后自动复机失败",
|
||||||
zap.Uint("card_id", cardID),
|
zap.String("carrier_type", ct),
|
||||||
|
zap.Uint("carrier_id", cid),
|
||||||
zap.Error(err))
|
zap.Error(err))
|
||||||
}
|
}
|
||||||
}(carrierID)
|
}(carrierType, carrierID)
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ type ResetService struct {
|
|||||||
redis *redis.Client
|
redis *redis.Client
|
||||||
packageUsageStore *postgres.PackageUsageStore
|
packageUsageStore *postgres.PackageUsageStore
|
||||||
logger *zap.Logger
|
logger *zap.Logger
|
||||||
|
resumeCallback ResumeCallback
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewResetService(
|
func NewResetService(
|
||||||
@@ -34,6 +35,11 @@ func NewResetService(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SetResumeCallback 设置复机回调,流量重置后对满足条件的卡触发自动复机
|
||||||
|
func (s *ResetService) SetResumeCallback(callback ResumeCallback) {
|
||||||
|
s.resumeCallback = callback
|
||||||
|
}
|
||||||
|
|
||||||
// ResetDailyUsage 任务 11.2-11.3: 重置日流量
|
// ResetDailyUsage 任务 11.2-11.3: 重置日流量
|
||||||
func (s *ResetService) ResetDailyUsage(ctx context.Context) error {
|
func (s *ResetService) ResetDailyUsage(ctx context.Context) error {
|
||||||
return s.resetDailyUsageWithDB(ctx, s.db)
|
return s.resetDailyUsageWithDB(ctx, s.db)
|
||||||
@@ -43,8 +49,8 @@ func (s *ResetService) ResetDailyUsage(ctx context.Context) error {
|
|||||||
func (s *ResetService) resetDailyUsageWithDB(ctx context.Context, db *gorm.DB) error {
|
func (s *ResetService) resetDailyUsageWithDB(ctx context.Context, db *gorm.DB) error {
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
|
|
||||||
return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
var resetPackages []*model.PackageUsage
|
||||||
// 查询需要重置的套餐
|
err := db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||||
var packages []*model.PackageUsage
|
var packages []*model.PackageUsage
|
||||||
err := tx.Where("data_reset_cycle = ?", constants.PackageDataResetDaily).
|
err := tx.Where("data_reset_cycle = ?", constants.PackageDataResetDaily).
|
||||||
Where("next_reset_at <= ?", now).
|
Where("next_reset_at <= ?", now).
|
||||||
@@ -60,21 +66,18 @@ func (s *ResetService) resetDailyUsageWithDB(ctx context.Context, db *gorm.DB) e
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// 批量重置
|
|
||||||
packageIDs := make([]uint, len(packages))
|
packageIDs := make([]uint, len(packages))
|
||||||
for i, pkg := range packages {
|
for i, pkg := range packages {
|
||||||
packageIDs[i] = pkg.ID
|
packageIDs[i] = pkg.ID
|
||||||
}
|
}
|
||||||
|
|
||||||
// 计算下次重置时间(明天 00:00:00)
|
|
||||||
nextReset := time.Date(now.Year(), now.Month(), now.Day()+1, 0, 0, 0, 0, now.Location())
|
nextReset := time.Date(now.Year(), now.Month(), now.Day()+1, 0, 0, 0, 0, now.Location())
|
||||||
|
|
||||||
// 批量更新
|
|
||||||
updates := map[string]interface{}{
|
updates := map[string]interface{}{
|
||||||
"data_usage_mb": 0,
|
"data_usage_mb": 0,
|
||||||
"last_reset_at": now,
|
"last_reset_at": now,
|
||||||
"next_reset_at": nextReset,
|
"next_reset_at": nextReset,
|
||||||
"status": constants.PackageUsageStatusActive, // 重置后恢复为生效中
|
"status": constants.PackageUsageStatusActive,
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := tx.Model(&model.PackageUsage{}).
|
if err := tx.Model(&model.PackageUsage{}).
|
||||||
@@ -87,8 +90,16 @@ func (s *ResetService) resetDailyUsageWithDB(ctx context.Context, db *gorm.DB) e
|
|||||||
zap.Int("count", len(packages)),
|
zap.Int("count", len(packages)),
|
||||||
zap.Time("next_reset_at", nextReset))
|
zap.Time("next_reset_at", nextReset))
|
||||||
|
|
||||||
|
resetPackages = packages
|
||||||
return nil
|
return nil
|
||||||
})
|
})
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
s.triggerResumeForPackages(ctx, resetPackages)
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// ResetMonthlyUsage 任务 11.4-11.5: 重置月流量
|
// ResetMonthlyUsage 任务 11.4-11.5: 重置月流量
|
||||||
@@ -100,8 +111,8 @@ func (s *ResetService) ResetMonthlyUsage(ctx context.Context) error {
|
|||||||
func (s *ResetService) resetMonthlyUsageWithDB(ctx context.Context, db *gorm.DB) error {
|
func (s *ResetService) resetMonthlyUsageWithDB(ctx context.Context, db *gorm.DB) error {
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
|
|
||||||
return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
var resetPackages []*model.PackageUsage
|
||||||
// 查询需要重置的套餐
|
err := db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||||
var packages []*model.PackageUsage
|
var packages []*model.PackageUsage
|
||||||
err := tx.Where("data_reset_cycle = ?", constants.PackageDataResetMonthly).
|
err := tx.Where("data_reset_cycle = ?", constants.PackageDataResetMonthly).
|
||||||
Where("next_reset_at <= ?", now).
|
Where("next_reset_at <= ?", now).
|
||||||
@@ -117,9 +128,7 @@ func (s *ResetService) resetMonthlyUsageWithDB(ctx context.Context, db *gorm.DB)
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// 按套餐分组处理(根据套餐周期类型计算下次重置时间)
|
|
||||||
for _, usage := range packages {
|
for _, usage := range packages {
|
||||||
// 查询套餐信息,获取 calendar_type
|
|
||||||
var pkg model.Package
|
var pkg model.Package
|
||||||
if err := tx.First(&pkg, usage.PackageID).Error; err != nil {
|
if err := tx.First(&pkg, usage.PackageID).Error; err != nil {
|
||||||
s.logger.Error("查询套餐信息失败",
|
s.logger.Error("查询套餐信息失败",
|
||||||
@@ -129,12 +138,9 @@ func (s *ResetService) resetMonthlyUsageWithDB(ctx context.Context, db *gorm.DB)
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
// 计算下次重置时间(基于套餐周期类型)
|
|
||||||
// 自然月套餐:每月1号重置
|
|
||||||
// 按天套餐:每30天重置
|
|
||||||
activatedAt := usage.ActivatedAt
|
activatedAt := usage.ActivatedAt
|
||||||
if activatedAt.IsZero() {
|
if activatedAt.IsZero() {
|
||||||
activatedAt = now // 兜底处理
|
activatedAt = now
|
||||||
}
|
}
|
||||||
nextResetAt := CalculateNextResetTime(constants.PackageDataResetMonthly, pkg.CalendarType, now, activatedAt)
|
nextResetAt := CalculateNextResetTime(constants.PackageDataResetMonthly, pkg.CalendarType, now, activatedAt)
|
||||||
if nextResetAt == nil {
|
if nextResetAt == nil {
|
||||||
@@ -143,12 +149,11 @@ func (s *ResetService) resetMonthlyUsageWithDB(ctx context.Context, db *gorm.DB)
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
// 更新套餐
|
|
||||||
updates := map[string]interface{}{
|
updates := map[string]interface{}{
|
||||||
"data_usage_mb": 0,
|
"data_usage_mb": 0,
|
||||||
"last_reset_at": now,
|
"last_reset_at": now,
|
||||||
"next_reset_at": *nextResetAt,
|
"next_reset_at": *nextResetAt,
|
||||||
"status": constants.PackageUsageStatusActive, // 重置后恢复为生效中
|
"status": constants.PackageUsageStatusActive,
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := tx.Model(usage).Updates(updates).Error; err != nil {
|
if err := tx.Model(usage).Updates(updates).Error; err != nil {
|
||||||
@@ -161,8 +166,16 @@ func (s *ResetService) resetMonthlyUsageWithDB(ctx context.Context, db *gorm.DB)
|
|||||||
zap.Time("next_reset_at", *nextResetAt))
|
zap.Time("next_reset_at", *nextResetAt))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
resetPackages = packages
|
||||||
return nil
|
return nil
|
||||||
})
|
})
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
s.triggerResumeForPackages(ctx, resetPackages)
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// ResetYearlyUsage 任务 11.6-11.7: 重置年流量
|
// ResetYearlyUsage 任务 11.6-11.7: 重置年流量
|
||||||
@@ -174,8 +187,8 @@ func (s *ResetService) ResetYearlyUsage(ctx context.Context) error {
|
|||||||
func (s *ResetService) resetYearlyUsageWithDB(ctx context.Context, db *gorm.DB) error {
|
func (s *ResetService) resetYearlyUsageWithDB(ctx context.Context, db *gorm.DB) error {
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
|
|
||||||
return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
var resetPackages []*model.PackageUsage
|
||||||
// 查询需要重置的套餐
|
err := db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||||
var packages []*model.PackageUsage
|
var packages []*model.PackageUsage
|
||||||
err := tx.Where("data_reset_cycle = ?", constants.PackageDataResetYearly).
|
err := tx.Where("data_reset_cycle = ?", constants.PackageDataResetYearly).
|
||||||
Where("next_reset_at <= ?", now).
|
Where("next_reset_at <= ?", now).
|
||||||
@@ -191,21 +204,18 @@ func (s *ResetService) resetYearlyUsageWithDB(ctx context.Context, db *gorm.DB)
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// 批量重置
|
|
||||||
packageIDs := make([]uint, len(packages))
|
packageIDs := make([]uint, len(packages))
|
||||||
for i, pkg := range packages {
|
for i, pkg := range packages {
|
||||||
packageIDs[i] = pkg.ID
|
packageIDs[i] = pkg.ID
|
||||||
}
|
}
|
||||||
|
|
||||||
// 计算下次重置时间(明年 1月1日 00:00:00)
|
|
||||||
nextReset := time.Date(now.Year()+1, 1, 1, 0, 0, 0, 0, now.Location())
|
nextReset := time.Date(now.Year()+1, 1, 1, 0, 0, 0, 0, now.Location())
|
||||||
|
|
||||||
// 批量更新
|
|
||||||
updates := map[string]interface{}{
|
updates := map[string]interface{}{
|
||||||
"data_usage_mb": 0,
|
"data_usage_mb": 0,
|
||||||
"last_reset_at": now,
|
"last_reset_at": now,
|
||||||
"next_reset_at": nextReset,
|
"next_reset_at": nextReset,
|
||||||
"status": constants.PackageUsageStatusActive, // 重置后恢复为生效中
|
"status": constants.PackageUsageStatusActive,
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := tx.Model(&model.PackageUsage{}).
|
if err := tx.Model(&model.PackageUsage{}).
|
||||||
@@ -218,6 +228,46 @@ func (s *ResetService) resetYearlyUsageWithDB(ctx context.Context, db *gorm.DB)
|
|||||||
zap.Int("count", len(packages)),
|
zap.Int("count", len(packages)),
|
||||||
zap.Time("next_reset_at", nextReset))
|
zap.Time("next_reset_at", nextReset))
|
||||||
|
|
||||||
|
resetPackages = packages
|
||||||
return nil
|
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)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user