fix: 代码质量清理 — 统一C端支付枚举、修复静默错误、注入Worker回调、移除废弃同步代码
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 7m6s
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 7m6s
- C端订单支付状态直接使用管理端统一枚举(1/2/3/4),移除冗余映射函数 - 资产查询绑定卡/套餐查询失败时记录Warn日志而非静默忽略 - Worker bootstrap注入停复机回调(流量耗尽停机、套餐激活复机) - 删除废弃的SIM状态同步服务和任务处理器 - 新增开发环境数据清理脚本(full/soft/table三种模式)
This commit is contained in:
@@ -133,7 +133,12 @@ func (s *Service) buildDeviceResolveResponse(ctx context.Context, device *model.
|
||||
slotMap[b.IotCardID] = b.SlotPosition
|
||||
isCurrentMap[b.IotCardID] = b.IsCurrent
|
||||
}
|
||||
cards, _ := s.iotCardStore.GetByIDs(ctx, cardIDs)
|
||||
cards, err := s.iotCardStore.GetByIDs(ctx, cardIDs)
|
||||
if err != nil {
|
||||
logger.GetAppLogger().Warn("查询设备绑定卡信息失败,结果可能不完整",
|
||||
zap.Uints("card_ids", cardIDs),
|
||||
zap.Error(err))
|
||||
}
|
||||
for _, c := range cards {
|
||||
resp.Cards = append(resp.Cards, dto.BoundCardInfo{
|
||||
CardID: c.ID,
|
||||
@@ -285,7 +290,12 @@ func (s *Service) GetRealtimeStatus(ctx context.Context, assetType string, id ui
|
||||
slotMap[b.IotCardID] = b.SlotPosition
|
||||
isCurrentMap[b.IotCardID] = b.IsCurrent
|
||||
}
|
||||
cards, _ := s.iotCardStore.GetByIDs(ctx, cardIDs)
|
||||
cards, err := s.iotCardStore.GetByIDs(ctx, cardIDs)
|
||||
if err != nil {
|
||||
logger.GetAppLogger().Warn("查询设备绑定卡信息失败,结果可能不完整",
|
||||
zap.Uints("card_ids", cardIDs),
|
||||
zap.Error(err))
|
||||
}
|
||||
for _, c := range cards {
|
||||
resp.Cards = append(resp.Cards, dto.BoundCardInfo{
|
||||
CardID: c.ID,
|
||||
@@ -533,7 +543,12 @@ func (s *Service) GetPackages(ctx context.Context, assetType string, id uint, pa
|
||||
for id := range pkgIDSet {
|
||||
pkgIDs = append(pkgIDs, id)
|
||||
}
|
||||
packages, _ := s.packageStore.GetByIDsUnscoped(ctx, pkgIDs)
|
||||
packages, pkgErr := s.packageStore.GetByIDsUnscoped(ctx, pkgIDs)
|
||||
if pkgErr != nil {
|
||||
logger.GetAppLogger().Warn("批量查询套餐信息失败,套餐名称可能缺失",
|
||||
zap.Uints("package_ids", pkgIDs),
|
||||
zap.Error(pkgErr))
|
||||
}
|
||||
pkgMap := make(map[uint]*model.Package, len(packages))
|
||||
for _, p := range packages {
|
||||
pkgMap[p.ID] = p
|
||||
|
||||
@@ -311,7 +311,7 @@ func (s *Service) createPackageOrder(
|
||||
OrderID: order.ID,
|
||||
OrderNo: order.OrderNo,
|
||||
TotalAmount: order.TotalAmount,
|
||||
PaymentStatus: orderStatusToClientStatus(order.PaymentStatus),
|
||||
PaymentStatus: order.PaymentStatus,
|
||||
CreatedAt: formatClientServiceTime(order.CreatedAt),
|
||||
},
|
||||
PayConfig: buildClientPayConfig(appID, payResult.PayConfig),
|
||||
@@ -641,19 +641,6 @@ func buildClientPurchaseBusinessKey(customerID uint, assetInfo *dto.AssetResolve
|
||||
return fmt.Sprintf("%d:%s:%d:%s:%s", customerID, assetInfo.AssetType, assetInfo.AssetID, req.AppType, strings.Join(parts, ","))
|
||||
}
|
||||
|
||||
func orderStatusToClientStatus(status int) int {
|
||||
switch status {
|
||||
case model.PaymentStatusPending:
|
||||
return 0
|
||||
case model.PaymentStatusPaid:
|
||||
return 1
|
||||
case model.PaymentStatusCancelled:
|
||||
return 2
|
||||
default:
|
||||
return status
|
||||
}
|
||||
}
|
||||
|
||||
func rechargeStatusToClientStatus(status int) int {
|
||||
switch status {
|
||||
case 1:
|
||||
|
||||
@@ -1,170 +0,0 @@
|
||||
// Package sync 提供数据同步的业务逻辑服务
|
||||
// 包含批量数据同步、任务调度等功能
|
||||
package sync
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/break/junhong_cmp_fiber/internal/task"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/constants"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/errors"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/queue"
|
||||
"github.com/bytedance/sonic"
|
||||
"github.com/hibiken/asynq"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// Service 同步服务
|
||||
type Service struct {
|
||||
queueClient *queue.Client
|
||||
logger *zap.Logger
|
||||
}
|
||||
|
||||
// NewService 创建同步服务实例
|
||||
func NewService(queueClient *queue.Client, logger *zap.Logger) *Service {
|
||||
return &Service{
|
||||
queueClient: queueClient,
|
||||
logger: logger,
|
||||
}
|
||||
}
|
||||
|
||||
// SyncSIMStatus 同步 SIM 卡状态(异步)
|
||||
func (s *Service) SyncSIMStatus(ctx context.Context, iccids []string, forceSync bool) error {
|
||||
// 构造任务载荷
|
||||
payload := &task.SIMStatusSyncPayload{
|
||||
RequestID: fmt.Sprintf("sim-sync-%d", getCurrentTimestamp()),
|
||||
ICCIDs: iccids,
|
||||
ForceSync: forceSync,
|
||||
}
|
||||
|
||||
payloadBytes, err := sonic.Marshal(payload)
|
||||
if err != nil {
|
||||
s.logger.Error("序列化 SIM 状态同步任务载荷失败",
|
||||
zap.Int("iccid_count", len(iccids)),
|
||||
zap.Bool("force_sync", forceSync),
|
||||
zap.Error(err))
|
||||
return errors.Wrap(errors.CodeInternalError, err, "序列化 SIM 状态同步任务载荷失败")
|
||||
}
|
||||
|
||||
// 提交任务到队列(高优先级)
|
||||
err = s.queueClient.EnqueueTask(
|
||||
ctx,
|
||||
constants.TaskTypeSIMStatusSync,
|
||||
payloadBytes,
|
||||
asynq.Queue(constants.QueueCritical), // SIM 状态同步使用高优先级队列
|
||||
asynq.MaxRetry(constants.DefaultRetryMax),
|
||||
)
|
||||
|
||||
if err != nil {
|
||||
s.logger.Error("提交 SIM 状态同步任务失败",
|
||||
zap.Int("iccid_count", len(iccids)),
|
||||
zap.Error(err))
|
||||
return errors.Wrap(errors.CodeInternalError, err, "提交 SIM 状态同步任务失败")
|
||||
}
|
||||
|
||||
s.logger.Info("SIM 状态同步任务已提交",
|
||||
zap.Int("iccid_count", len(iccids)),
|
||||
zap.Bool("force_sync", forceSync))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// SyncData 通用数据同步(异步)
|
||||
func (s *Service) SyncData(ctx context.Context, syncType string, startDate string, endDate string, batchSize int) error {
|
||||
// 设置默认批量大小
|
||||
if batchSize <= 0 {
|
||||
batchSize = 100 // 默认批量大小
|
||||
}
|
||||
|
||||
// 构造任务载荷
|
||||
payload := &task.DataSyncPayload{
|
||||
RequestID: fmt.Sprintf("data-sync-%s-%d", syncType, getCurrentTimestamp()),
|
||||
SyncType: syncType,
|
||||
StartDate: startDate,
|
||||
EndDate: endDate,
|
||||
BatchSize: batchSize,
|
||||
}
|
||||
|
||||
payloadBytes, err := sonic.Marshal(payload)
|
||||
if err != nil {
|
||||
s.logger.Error("序列化数据同步任务载荷失败",
|
||||
zap.String("sync_type", syncType),
|
||||
zap.String("start_date", startDate),
|
||||
zap.String("end_date", endDate),
|
||||
zap.Error(err))
|
||||
return errors.Wrap(errors.CodeInternalError, err, "序列化数据同步任务载荷失败")
|
||||
}
|
||||
|
||||
// 提交任务到队列(默认优先级)
|
||||
err = s.queueClient.EnqueueTask(
|
||||
ctx,
|
||||
constants.TaskTypeDataSync,
|
||||
payloadBytes,
|
||||
asynq.Queue(constants.QueueDefault),
|
||||
asynq.MaxRetry(constants.DefaultRetryMax),
|
||||
)
|
||||
|
||||
if err != nil {
|
||||
s.logger.Error("提交数据同步任务失败",
|
||||
zap.String("sync_type", syncType),
|
||||
zap.Error(err))
|
||||
return errors.Wrap(errors.CodeInternalError, err, "提交数据同步任务失败")
|
||||
}
|
||||
|
||||
s.logger.Info("数据同步任务已提交",
|
||||
zap.String("sync_type", syncType),
|
||||
zap.String("start_date", startDate),
|
||||
zap.String("end_date", endDate),
|
||||
zap.Int("batch_size", batchSize))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// SyncFlowUsage 同步流量使用数据(异步)
|
||||
func (s *Service) SyncFlowUsage(ctx context.Context, startDate string, endDate string) error {
|
||||
return s.SyncData(ctx, "flow_usage", startDate, endDate, 100)
|
||||
}
|
||||
|
||||
// SyncRealNameInfo 同步实名信息(异步)
|
||||
func (s *Service) SyncRealNameInfo(ctx context.Context, startDate string, endDate string) error {
|
||||
return s.SyncData(ctx, "real_name", startDate, endDate, 50)
|
||||
}
|
||||
|
||||
// SyncBatchSIMStatus 批量同步多个 ICCID 的状态(异步)
|
||||
func (s *Service) SyncBatchSIMStatus(ctx context.Context, iccids []string) error {
|
||||
// 如果 ICCID 列表为空,直接返回
|
||||
if len(iccids) == 0 {
|
||||
s.logger.Warn("批量同步 SIM 状态时 ICCID 列表为空")
|
||||
return nil
|
||||
}
|
||||
|
||||
// 分批处理(每批最多 100 个)
|
||||
batchSize := 100
|
||||
for i := 0; i < len(iccids); i += batchSize {
|
||||
end := i + batchSize
|
||||
if end > len(iccids) {
|
||||
end = len(iccids)
|
||||
}
|
||||
|
||||
batch := iccids[i:end]
|
||||
if err := s.SyncSIMStatus(ctx, batch, false); err != nil {
|
||||
s.logger.Error("批量同步 SIM 状态失败",
|
||||
zap.Int("batch_start", i),
|
||||
zap.Int("batch_end", end),
|
||||
zap.Error(err))
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
s.logger.Info("批量 SIM 状态同步任务已全部提交",
|
||||
zap.Int("total_iccids", len(iccids)),
|
||||
zap.Int("batch_size", batchSize))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// getCurrentTimestamp 获取当前时间戳(毫秒)
|
||||
func getCurrentTimestamp() int64 {
|
||||
return 0 // 实际实现应返回真实时间戳
|
||||
}
|
||||
Reference in New Issue
Block a user