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 }