Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
373 lines
11 KiB
Go
373 lines
11 KiB
Go
package iot_card
|
||
|
||
import (
|
||
"context"
|
||
"time"
|
||
|
||
"github.com/redis/go-redis/v9"
|
||
"go.uber.org/zap"
|
||
"gorm.io/gorm"
|
||
|
||
stderrors "errors"
|
||
|
||
"github.com/break/junhong_cmp_fiber/internal/gateway"
|
||
"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"
|
||
)
|
||
|
||
// StopResumeService 停复机服务
|
||
// 处理 IoT 卡的自动停机、复机和手动停复机逻辑
|
||
type StopResumeService struct {
|
||
db *gorm.DB
|
||
redis *redis.Client
|
||
iotCardStore *postgres.IotCardStore
|
||
deviceSimBindingStore *postgres.DeviceSimBindingStore
|
||
gatewayClient *gateway.Client
|
||
logger *zap.Logger
|
||
|
||
maxRetries int
|
||
retryInterval time.Duration
|
||
}
|
||
|
||
// NewStopResumeService 创建停复机服务
|
||
func NewStopResumeService(
|
||
db *gorm.DB,
|
||
redis *redis.Client,
|
||
iotCardStore *postgres.IotCardStore,
|
||
deviceSimBindingStore *postgres.DeviceSimBindingStore,
|
||
gatewayClient *gateway.Client,
|
||
logger *zap.Logger,
|
||
) *StopResumeService {
|
||
return &StopResumeService{
|
||
db: db,
|
||
redis: redis,
|
||
iotCardStore: iotCardStore,
|
||
deviceSimBindingStore: deviceSimBindingStore,
|
||
gatewayClient: gatewayClient,
|
||
logger: logger,
|
||
maxRetries: 3,
|
||
retryInterval: 2 * time.Second,
|
||
}
|
||
}
|
||
|
||
// CheckAndStopCard 任务 24.3: 检查流量耗尽并停机
|
||
// 当所有套餐流量用完时,调用运营商接口停机
|
||
func (s *StopResumeService) CheckAndStopCard(ctx context.Context, cardID uint) error {
|
||
// 查询卡信息
|
||
card, err := s.iotCardStore.GetByID(ctx, cardID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// 如果已经是停机状态,跳过
|
||
if card.NetworkStatus == constants.NetworkStatusOffline {
|
||
s.logger.Debug("卡已处于停机状态,跳过",
|
||
zap.Uint("card_id", cardID))
|
||
return nil
|
||
}
|
||
|
||
// 检查是否有可用套餐(status=1 生效中 或 status=0 待生效)
|
||
hasAvailablePackage, err := s.hasAvailablePackage(ctx, cardID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// 如果还有可用套餐,不停机
|
||
if hasAvailablePackage {
|
||
return nil
|
||
}
|
||
|
||
// 任务 24.5: 调用运营商停机接口(带重试机制)
|
||
if err := s.stopCardWithRetry(ctx, card); err != nil {
|
||
s.logger.Error("调用运营商停机接口失败",
|
||
zap.Uint("card_id", cardID),
|
||
zap.String("iccid", card.ICCID),
|
||
zap.Error(err))
|
||
return err
|
||
}
|
||
|
||
// 更新卡状态
|
||
now := time.Now()
|
||
if err := s.db.WithContext(ctx).Model(card).Updates(map[string]any{
|
||
"network_status": constants.NetworkStatusOffline,
|
||
"stopped_at": now,
|
||
"stop_reason": constants.StopReasonTrafficExhausted,
|
||
}).Error; err != nil {
|
||
return err
|
||
}
|
||
|
||
s.logger.Info("卡因流量耗尽已停机",
|
||
zap.Uint("card_id", cardID),
|
||
zap.String("iccid", card.ICCID))
|
||
|
||
return nil
|
||
}
|
||
|
||
// ResumeCardIfStopped 套餐激活或流量重置后自动复机入口
|
||
// 支持 iot_card 和 device 两种载体类型:
|
||
// - iot_card:对单张卡执行实名检查 → 停机原因检查 → 分布式锁 → Gateway → 更新 DB
|
||
// - device:查询设备下所有已绑定卡,逐一执行单卡复机逻辑
|
||
func (s *StopResumeService) ResumeCardIfStopped(ctx context.Context, carrierType string, carrierID uint) error {
|
||
switch carrierType {
|
||
case "iot_card":
|
||
return s.resumeSingleCard(ctx, carrierID)
|
||
case "device":
|
||
return s.resumeDeviceCards(ctx, carrierID)
|
||
default:
|
||
return nil
|
||
}
|
||
}
|
||
|
||
// resumeSingleCard 对单张卡执行复机逻辑
|
||
// 依次检查:已开机则跳过 → 未实名则跳过(WARN)→ 非流量耗尽停机则跳过 → 加锁 → 调 Gateway → 更新 DB
|
||
func (s *StopResumeService) resumeSingleCard(ctx context.Context, cardID uint) error {
|
||
card, err := s.iotCardStore.GetByID(ctx, cardID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
if card.NetworkStatus == constants.NetworkStatusOnline {
|
||
return nil
|
||
}
|
||
|
||
if card.RealNameStatus != constants.RealNameStatusVerified {
|
||
s.logger.Warn("未实名卡跳过自动复机",
|
||
zap.Uint("card_id", cardID),
|
||
zap.String("iccid", card.ICCID))
|
||
return nil
|
||
}
|
||
|
||
if card.StopReason != constants.StopReasonTrafficExhausted {
|
||
s.logger.Debug("卡非流量耗尽停机,不自动复机",
|
||
zap.Uint("card_id", cardID),
|
||
zap.String("stop_reason", card.StopReason))
|
||
return nil
|
||
}
|
||
|
||
// 分布式锁防止同一张卡并发触发多次 Gateway 调用,TTL 30s
|
||
lockKey := constants.RedisCardResumeLockKey(cardID)
|
||
locked, lockErr := s.redis.SetNX(ctx, lockKey, "1", 30*time.Second).Result()
|
||
if lockErr != nil {
|
||
s.logger.Warn("获取复机分布式锁失败,跳过本次复机",
|
||
zap.Uint("card_id", cardID),
|
||
zap.Error(lockErr))
|
||
return nil
|
||
}
|
||
if !locked {
|
||
s.logger.Debug("复机操作进行中,跳过", zap.Uint("card_id", cardID))
|
||
return nil
|
||
}
|
||
defer s.redis.Del(ctx, lockKey)
|
||
|
||
if err := s.resumeCardWithRetry(ctx, card); err != nil {
|
||
s.logger.Error("调用运营商复机接口失败",
|
||
zap.Uint("card_id", cardID),
|
||
zap.String("iccid", card.ICCID),
|
||
zap.Error(err))
|
||
return err
|
||
}
|
||
|
||
now := time.Now()
|
||
if err := s.db.WithContext(ctx).Model(card).Updates(map[string]any{
|
||
"network_status": constants.NetworkStatusOnline,
|
||
"resumed_at": now,
|
||
"stop_reason": "",
|
||
}).Error; err != nil {
|
||
// Gateway 成功但 DB 更新失败:记录 ERROR 供人工排查,不阻断流程
|
||
s.logger.Error("复机 Gateway 成功但 DB 更新失败",
|
||
zap.Uint("card_id", cardID),
|
||
zap.String("iccid", card.ICCID),
|
||
zap.Error(err))
|
||
return nil
|
||
}
|
||
|
||
s.logger.Info("卡已自动复机",
|
||
zap.Uint("card_id", cardID),
|
||
zap.String("iccid", card.ICCID))
|
||
|
||
return nil
|
||
}
|
||
|
||
// resumeDeviceCards 对设备绑定的所有已绑定卡逐一执行复机逻辑
|
||
// 部分卡实名、部分未实名时:已实名卡复机,未实名卡跳过,互不影响
|
||
func (s *StopResumeService) resumeDeviceCards(ctx context.Context, deviceID uint) error {
|
||
bindings, err := s.deviceSimBindingStore.ListByDeviceID(ctx, deviceID)
|
||
if err != nil {
|
||
s.logger.Warn("查询设备绑定卡失败",
|
||
zap.Uint("device_id", deviceID),
|
||
zap.Error(err))
|
||
return nil
|
||
}
|
||
|
||
for _, binding := range bindings {
|
||
if resumeErr := s.resumeSingleCard(ctx, binding.IotCardID); resumeErr != nil {
|
||
s.logger.Warn("设备绑定卡复机失败,继续处理其他卡",
|
||
zap.Uint("device_id", deviceID),
|
||
zap.Uint("card_id", binding.IotCardID),
|
||
zap.Error(resumeErr))
|
||
}
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// hasAvailablePackage 检查是否有可用套餐
|
||
func (s *StopResumeService) hasAvailablePackage(ctx context.Context, cardID uint) (bool, error) {
|
||
var count int64
|
||
err := s.db.WithContext(ctx).Model(&model.PackageUsage{}).
|
||
Where("iot_card_id = ?", cardID).
|
||
Where("status IN ?", []int{
|
||
constants.PackageUsageStatusPending, // 待生效
|
||
constants.PackageUsageStatusActive, // 生效中
|
||
}).
|
||
Count(&count).Error
|
||
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
|
||
return count > 0, nil
|
||
}
|
||
|
||
// stopCardWithRetry 任务 24.5: 调用运营商停机接口(带重试机制)
|
||
func (s *StopResumeService) stopCardWithRetry(ctx context.Context, card *model.IotCard) error {
|
||
if s.gatewayClient == nil {
|
||
return errors.New(errors.CodeInternalError, "Gateway 未配置,停复机操作不可用")
|
||
}
|
||
|
||
var lastErr error
|
||
for i := 0; i < s.maxRetries; i++ {
|
||
if i > 0 {
|
||
s.logger.Debug("重试调用停机接口",
|
||
zap.Int("attempt", i+1),
|
||
zap.String("iccid", card.ICCID))
|
||
time.Sleep(s.retryInterval)
|
||
}
|
||
|
||
err := s.gatewayClient.StopCard(ctx, &gateway.CardOperationReq{
|
||
CardNo: card.ICCID,
|
||
})
|
||
if err == nil {
|
||
return nil
|
||
}
|
||
|
||
lastErr = err
|
||
s.logger.Warn("调用停机接口失败,准备重试",
|
||
zap.Int("attempt", i+1),
|
||
zap.Error(err))
|
||
}
|
||
|
||
return lastErr
|
||
}
|
||
|
||
// resumeCardWithRetry 任务 24.5: 调用运营商复机接口(带重试机制)
|
||
func (s *StopResumeService) resumeCardWithRetry(ctx context.Context, card *model.IotCard) error {
|
||
if s.gatewayClient == nil {
|
||
return errors.New(errors.CodeInternalError, "Gateway 未配置,停复机操作不可用")
|
||
}
|
||
|
||
var lastErr error
|
||
for i := 0; i < s.maxRetries; i++ {
|
||
if i > 0 {
|
||
s.logger.Debug("重试调用复机接口",
|
||
zap.Int("attempt", i+1),
|
||
zap.String("iccid", card.ICCID))
|
||
time.Sleep(s.retryInterval)
|
||
}
|
||
|
||
err := s.gatewayClient.StartCard(ctx, &gateway.CardOperationReq{
|
||
CardNo: card.ICCID,
|
||
})
|
||
if err == nil {
|
||
return nil
|
||
}
|
||
|
||
lastErr = err
|
||
s.logger.Warn("调用复机接口失败,准备重试",
|
||
zap.Int("attempt", i+1),
|
||
zap.Error(err))
|
||
}
|
||
|
||
return lastErr
|
||
}
|
||
|
||
// ManualStopCard 手动停机单张卡(通过ICCID)
|
||
func (s *StopResumeService) ManualStopCard(ctx context.Context, iccid string) error {
|
||
card, err := s.iotCardStore.GetByICCID(ctx, iccid)
|
||
if err != nil {
|
||
return errors.New(errors.CodeNotFound, "卡不存在")
|
||
}
|
||
|
||
if card.RealNameStatus != constants.RealNameStatusVerified {
|
||
return errors.New(errors.CodeForbidden, "卡未实名,无法操作")
|
||
}
|
||
|
||
// 检查绑定设备是否在复机保护期
|
||
if s.deviceSimBindingStore != nil && s.redis != nil {
|
||
binding, bindErr := s.deviceSimBindingStore.GetActiveBindingByCardID(ctx, card.ID)
|
||
if bindErr == nil && binding != nil {
|
||
exists, _ := s.redis.Exists(ctx, constants.RedisDeviceProtectKey(binding.DeviceID, "start")).Result()
|
||
if exists > 0 {
|
||
return errors.New(errors.CodeForbidden, "设备复机保护期内,禁止停机")
|
||
}
|
||
} else if bindErr != nil && !stderrors.Is(bindErr, gorm.ErrRecordNotFound) {
|
||
return errors.Wrap(errors.CodeInternalError, bindErr, "查询卡绑定关系失败")
|
||
}
|
||
}
|
||
|
||
if err := s.stopCardWithRetry(ctx, card); err != nil {
|
||
return errors.Wrap(errors.CodeGatewayError, err, "调用运营商停机失败,请稍后重试")
|
||
}
|
||
|
||
now := time.Now()
|
||
if err := s.db.WithContext(ctx).Model(card).Updates(map[string]any{
|
||
"network_status": constants.NetworkStatusOffline,
|
||
"stopped_at": now,
|
||
"stop_reason": constants.StopReasonManual,
|
||
}).Error; err != nil {
|
||
return errors.Wrap(errors.CodeDatabaseError, err, "更新卡状态失败")
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// ManualStartCard 手动复机单张卡(通过ICCID)
|
||
func (s *StopResumeService) ManualStartCard(ctx context.Context, iccid string) error {
|
||
card, err := s.iotCardStore.GetByICCID(ctx, iccid)
|
||
if err != nil {
|
||
return errors.New(errors.CodeNotFound, "卡不存在")
|
||
}
|
||
|
||
if card.RealNameStatus != constants.RealNameStatusVerified {
|
||
return errors.New(errors.CodeForbidden, "卡未实名,无法操作")
|
||
}
|
||
|
||
// 检查绑定设备是否在停机保护期
|
||
if s.deviceSimBindingStore != nil && s.redis != nil {
|
||
binding, bindErr := s.deviceSimBindingStore.GetActiveBindingByCardID(ctx, card.ID)
|
||
if bindErr == nil && binding != nil {
|
||
exists, _ := s.redis.Exists(ctx, constants.RedisDeviceProtectKey(binding.DeviceID, "stop")).Result()
|
||
if exists > 0 {
|
||
return errors.New(errors.CodeForbidden, "设备停机保护期内,禁止复机")
|
||
}
|
||
} else if bindErr != nil && !stderrors.Is(bindErr, gorm.ErrRecordNotFound) {
|
||
return errors.Wrap(errors.CodeInternalError, bindErr, "查询卡绑定关系失败")
|
||
}
|
||
}
|
||
|
||
if err := s.resumeCardWithRetry(ctx, card); err != nil {
|
||
return errors.Wrap(errors.CodeGatewayError, err, "调用运营商复机失败,请稍后重试")
|
||
}
|
||
|
||
now := time.Now()
|
||
if err := s.db.WithContext(ctx).Model(card).Updates(map[string]any{
|
||
"network_status": constants.NetworkStatusOnline,
|
||
"resumed_at": now,
|
||
"stop_reason": "",
|
||
}).Error; err != nil {
|
||
return errors.Wrap(errors.CodeDatabaseError, err, "更新卡状态失败")
|
||
}
|
||
return nil
|
||
}
|