309 lines
13 KiB
Go
309 lines
13 KiB
Go
package task
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"encoding/csv"
|
|
stderrors "errors"
|
|
"io"
|
|
"os"
|
|
"strconv"
|
|
"strings"
|
|
"unicode/utf8"
|
|
|
|
"github.com/bytedance/sonic"
|
|
"github.com/hibiken/asynq"
|
|
"go.uber.org/zap"
|
|
"gorm.io/gorm"
|
|
|
|
"github.com/break/junhong_cmp_fiber/internal/infrastructure/audit"
|
|
"github.com/break/junhong_cmp_fiber/internal/model"
|
|
"github.com/break/junhong_cmp_fiber/internal/model/dto"
|
|
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
|
|
"github.com/break/junhong_cmp_fiber/pkg/asynctask"
|
|
"github.com/break/junhong_cmp_fiber/pkg/auditcontext"
|
|
"github.com/break/junhong_cmp_fiber/pkg/constants"
|
|
apperrors "github.com/break/junhong_cmp_fiber/pkg/errors"
|
|
"github.com/break/junhong_cmp_fiber/pkg/middleware"
|
|
"github.com/break/junhong_cmp_fiber/pkg/storage"
|
|
)
|
|
|
|
// AssetPackageBatchOrderCreator 定义 Worker 复用现有后台订单规则的最小接口。
|
|
type AssetPackageBatchOrderCreator interface {
|
|
CreateAdminOrder(ctx context.Context, req *dto.CreateAdminOrderRequest, buyerType string, buyerID uint) (*dto.OrderResponse, error)
|
|
}
|
|
|
|
// AssetPackageBatchOrderPayload 资产套餐批量订购任务载荷。
|
|
type AssetPackageBatchOrderPayload struct {
|
|
TaskID uint `json:"task_id"`
|
|
}
|
|
|
|
// AssetPackageBatchOrderHandler 资产套餐批量订购任务处理器。
|
|
type AssetPackageBatchOrderHandler struct {
|
|
taskStore *postgres.AssetPackageBatchOrderTaskStore
|
|
shopStore *postgres.ShopStore
|
|
orderCreator AssetPackageBatchOrderCreator
|
|
storageService *storage.Service
|
|
logger *zap.Logger
|
|
auditWriter *audit.Writer
|
|
}
|
|
|
|
// NewAssetPackageBatchOrderHandler 创建资产套餐批量订购任务处理器。
|
|
func NewAssetPackageBatchOrderHandler(taskStore *postgres.AssetPackageBatchOrderTaskStore, shopStore *postgres.ShopStore, orderCreator AssetPackageBatchOrderCreator, storageService *storage.Service, logger *zap.Logger, auditWriters ...*audit.Writer) *AssetPackageBatchOrderHandler {
|
|
handler := &AssetPackageBatchOrderHandler{taskStore: taskStore, shopStore: shopStore, orderCreator: orderCreator, storageService: storageService, logger: logger}
|
|
if len(auditWriters) > 0 {
|
|
handler.auditWriter = auditWriters[0]
|
|
}
|
|
return handler
|
|
}
|
|
|
|
// Handle 处理资产套餐批量订购任务。
|
|
func (h *AssetPackageBatchOrderHandler) Handle(ctx context.Context, taskMessage *asynq.Task) error {
|
|
var payload AssetPackageBatchOrderPayload
|
|
if err := sonic.Unmarshal(taskMessage.Payload(), &payload); err != nil {
|
|
h.logger.Error("解析资产套餐批量订购任务载荷失败", zap.Error(err))
|
|
return asynq.SkipRetry
|
|
}
|
|
taskRecord, err := h.taskStore.GetByID(ctx, payload.TaskID)
|
|
if err != nil {
|
|
h.logger.Error("查询资产套餐批量订购任务失败", zap.Uint("task_id", payload.TaskID), zap.Error(err))
|
|
return asynq.SkipRetry
|
|
}
|
|
rootEventID := audit.TaskEventID(constants.AuditResourceAssetPackageBatchOrderTask, taskRecord.ID, "completed")
|
|
ctx = auditcontext.With(ctx, auditcontext.Context{
|
|
ActorKind: constants.AuditActorSystemTask, ActorID: constants.TaskTypeAssetPackageBatchOrder,
|
|
ActorName: "资产套餐批量订购任务", Source: constants.AuditSourceWorker,
|
|
CorrelationID: taskRecord.TaskNo, ParentEventID: rootEventID,
|
|
})
|
|
claimed, err := h.taskStore.Claim(ctx, taskRecord.ID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !claimed {
|
|
h.logger.Info("资产套餐批量订购任务已被领取或终结,跳过重复消费", zap.Uint("task_id", taskRecord.ID))
|
|
return nil
|
|
}
|
|
rows, err := h.downloadAndParse(ctx, taskRecord.StorageKey)
|
|
if err != nil {
|
|
h.logger.Error("下载或解析批量订购CSV失败", zap.Uint("task_id", taskRecord.ID), zap.Error(err))
|
|
if finishErr := h.finishBatchOrderTask(ctx, taskRecord, nil, 0, 1, asynctask.StatusFailed, err.Error()); finishErr != nil {
|
|
h.resetBatchOrderTaskForRetry(ctx, taskRecord.ID)
|
|
return finishErr
|
|
}
|
|
return asynq.SkipRetry
|
|
}
|
|
items, successCount, failCount := h.processRows(ctx, taskRecord, rows)
|
|
if err := h.finishBatchOrderTask(ctx, taskRecord, items, successCount, failCount, asynctask.StatusCompleted, ""); err != nil {
|
|
h.resetBatchOrderTaskForRetry(ctx, taskRecord.ID)
|
|
return err
|
|
}
|
|
h.logger.Info("资产套餐批量订购任务完成", zap.Uint("task_id", taskRecord.ID), zap.Int("success", successCount), zap.Int("fail", failCount))
|
|
return nil
|
|
}
|
|
|
|
func (h *AssetPackageBatchOrderHandler) resetBatchOrderTaskForRetry(ctx context.Context, taskID uint) {
|
|
_ = h.taskStore.DB().WithContext(ctx).Model(&model.AssetPackageBatchOrderTask{}).
|
|
Where("id = ? AND status = ?", taskID, asynctask.StatusProcessing).
|
|
Updates(map[string]any{"status": asynctask.StatusPending, "started_at": nil}).Error
|
|
}
|
|
|
|
func (h *AssetPackageBatchOrderHandler) finishBatchOrderTask(ctx context.Context, taskRecord *model.AssetPackageBatchOrderTask, items model.AssetPackageBatchOrderResultItems, successCount, failCount, status int, errorMessage string) error {
|
|
if h.auditWriter == nil {
|
|
return apperrors.New(apperrors.CodeInvalidStatus, "资产套餐批量订购统一审计接缝未配置")
|
|
}
|
|
return h.taskStore.DB().WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
txStore := h.taskStore.WithTx(tx)
|
|
if status == asynctask.StatusFailed {
|
|
if err := txStore.MarkFailed(ctx, taskRecord.ID, errorMessage); err != nil {
|
|
return err
|
|
}
|
|
} else if err := txStore.Complete(ctx, taskRecord.ID, items, successCount, failCount); err != nil {
|
|
return err
|
|
}
|
|
rootID := audit.TaskEventID(constants.AuditResourceAssetPackageBatchOrderTask, taskRecord.ID, "completed")
|
|
result := batchAuditResult(successCount, failCount)
|
|
return h.auditWriter.WriteTask(ctx, tx, audit.TaskInput{
|
|
EventID: rootID, ActionCode: constants.AuditActionAssetPackageBatchOrderTaskCompleted,
|
|
Summary: "完成资产套餐批量订购任务", TaskID: taskRecord.ID, TaskNo: taskRecord.TaskNo,
|
|
Result: result, CorrelationID: taskRecord.TaskNo,
|
|
ParentEventID: audit.TaskEventID(constants.AuditResourceAssetPackageBatchOrderTask, taskRecord.ID, "created"),
|
|
BatchTotal: len(items), SuccessCount: successCount, FailCount: failCount,
|
|
IdentitySnapshot: map[string]any{
|
|
"id": taskRecord.ID, "task_no": taskRecord.TaskNo, "file_name": taskRecord.FileName,
|
|
"package_id": taskRecord.PackageID, "package_code": taskRecord.PackageCode,
|
|
"package_name": taskRecord.PackageName, "payment_method": taskRecord.PaymentMethod,
|
|
},
|
|
BeforeData: map[string]any{"status": asynctask.StatusProcessing},
|
|
AfterData: map[string]any{
|
|
"status": status, "total_count": len(items), "success_count": successCount, "fail_count": failCount,
|
|
},
|
|
Metadata: map[string]any{"task_success_count": successCount, "task_fail_count": failCount},
|
|
})
|
|
})
|
|
}
|
|
|
|
type assetPackageBatchOrderRow struct {
|
|
Line int
|
|
Identifier string
|
|
}
|
|
|
|
func (h *AssetPackageBatchOrderHandler) downloadAndParse(ctx context.Context, key string) ([]assetPackageBatchOrderRow, error) {
|
|
if h.storageService == nil {
|
|
return nil, assetPackageBatchOrderError("对象存储服务未配置")
|
|
}
|
|
localPath, cleanup, err := h.storageService.DownloadToTemp(ctx, key)
|
|
if err != nil {
|
|
return nil, assetPackageBatchOrderError("下载批量订购CSV失败")
|
|
}
|
|
defer cleanup()
|
|
info, err := os.Stat(localPath)
|
|
if err != nil || info.Size() > constants.AssetPackageBatchOrderMaxFileSize {
|
|
return nil, assetPackageBatchOrderError("批量订购CSV不存在或超过10MB")
|
|
}
|
|
data, err := os.ReadFile(localPath)
|
|
if err != nil {
|
|
return nil, assetPackageBatchOrderError("读取批量订购CSV失败")
|
|
}
|
|
return parseAssetPackageBatchOrderCSV(data)
|
|
}
|
|
|
|
func parseAssetPackageBatchOrderCSV(data []byte) ([]assetPackageBatchOrderRow, error) {
|
|
data = bytes.TrimPrefix(data, []byte{0xEF, 0xBB, 0xBF})
|
|
if !utf8.Valid(data) {
|
|
return nil, assetPackageBatchOrderError("批量订购CSV必须使用UTF-8编码")
|
|
}
|
|
reader := csv.NewReader(bytes.NewReader(data))
|
|
reader.TrimLeadingSpace = true
|
|
rows := make([]assetPackageBatchOrderRow, 0)
|
|
line := 0
|
|
for {
|
|
record, err := reader.Read()
|
|
if err == io.EOF {
|
|
break
|
|
}
|
|
line++
|
|
if err != nil {
|
|
return nil, assetPackageBatchOrderError("批量订购CSV格式错误")
|
|
}
|
|
if len(record) != 1 {
|
|
return nil, assetPackageBatchOrderError("批量订购CSV必须只有一列资产标识")
|
|
}
|
|
identifier := strings.TrimSpace(record[0])
|
|
if line == 1 && isAssetIdentifierHeader(identifier) {
|
|
continue
|
|
}
|
|
rows = append(rows, assetPackageBatchOrderRow{Line: line, Identifier: identifier})
|
|
if len(rows) > constants.AssetPackageBatchOrderMaxRows {
|
|
return nil, assetPackageBatchOrderError("批量订购CSV最多包含1000行资产")
|
|
}
|
|
}
|
|
if len(rows) == 0 {
|
|
return nil, assetPackageBatchOrderError("批量订购CSV没有有效数据行")
|
|
}
|
|
return rows, nil
|
|
}
|
|
|
|
func (h *AssetPackageBatchOrderHandler) processRows(ctx context.Context, taskRecord *model.AssetPackageBatchOrderTask, rows []assetPackageBatchOrderRow) (model.AssetPackageBatchOrderResultItems, int, int) {
|
|
subordinateShopIDs := h.resolveSubordinateShopIDs(ctx, taskRecord)
|
|
workerCtx := middleware.SetUserContext(ctx, &middleware.UserContextInfo{
|
|
UserID: taskRecord.Creator, UserType: taskRecord.CreatorUserType,
|
|
Username: taskRecord.CreatorName, ShopID: taskRecord.CreatorShopID,
|
|
SubordinateShopIDs: subordinateShopIDs,
|
|
})
|
|
workerCtx = auditcontext.With(workerCtx, auditcontext.Context{
|
|
ActorKind: constants.AuditActorSystemTask, ActorID: constants.TaskTypeAssetPackageBatchOrder,
|
|
ActorName: "资产套餐批量订购任务", Source: constants.AuditSourceWorker,
|
|
CorrelationID: taskRecord.TaskNo,
|
|
ParentEventID: audit.TaskEventID(constants.AuditResourceAssetPackageBatchOrderTask, taskRecord.ID, "completed"),
|
|
})
|
|
buyerType, buyerID := "", uint(0)
|
|
if taskRecord.CreatorUserType == constants.UserTypeAgent {
|
|
buyerType, buyerID = model.BuyerTypeAgent, taskRecord.CreatorShopID
|
|
}
|
|
items := make(model.AssetPackageBatchOrderResultItems, 0, len(rows))
|
|
seenInputs, seenAssets := make(map[string]int), make(map[string]int)
|
|
successCount := 0
|
|
for _, row := range rows {
|
|
item := h.processOne(workerCtx, taskRecord, row, buyerType, buyerID, seenInputs, seenAssets)
|
|
items = append(items, item)
|
|
if item.Status == constants.AssetPackageBatchOrderItemStatusSuccess {
|
|
successCount++
|
|
}
|
|
}
|
|
return items, successCount, len(items) - successCount
|
|
}
|
|
|
|
func (h *AssetPackageBatchOrderHandler) resolveSubordinateShopIDs(ctx context.Context, taskRecord *model.AssetPackageBatchOrderTask) []uint {
|
|
if taskRecord.CreatorUserType != constants.UserTypeAgent {
|
|
return nil
|
|
}
|
|
fallback := []uint{taskRecord.CreatorShopID}
|
|
if h.shopStore == nil {
|
|
h.logger.Warn("批量订购任务未配置店铺存储,代理数据权限降级为仅自己店铺", zap.Uint("task_id", taskRecord.ID), zap.Uint("shop_id", taskRecord.CreatorShopID))
|
|
return fallback
|
|
}
|
|
shopIDs, err := h.shopStore.GetSubordinateShopIDs(ctx, taskRecord.CreatorShopID)
|
|
if err != nil || len(shopIDs) == 0 {
|
|
h.logger.Warn("查询批量订购任务代理数据权限失败,降级为仅自己店铺", zap.Uint("task_id", taskRecord.ID), zap.Uint("shop_id", taskRecord.CreatorShopID), zap.Error(err))
|
|
return fallback
|
|
}
|
|
return shopIDs
|
|
}
|
|
|
|
func (h *AssetPackageBatchOrderHandler) processOne(ctx context.Context, taskRecord *model.AssetPackageBatchOrderTask, row assetPackageBatchOrderRow, buyerType string, buyerID uint, seenInputs, seenAssets map[string]int) model.AssetPackageBatchOrderResultItem {
|
|
item := model.AssetPackageBatchOrderResultItem{Line: row.Line, AssetIdentifier: row.Identifier}
|
|
if row.Identifier == "" {
|
|
return failedBatchOrderItem(item, "资产标识不能为空")
|
|
}
|
|
normalized := strings.ToLower(row.Identifier)
|
|
if firstLine, exists := seenInputs[normalized]; exists {
|
|
return failedBatchOrderItem(item, "资产标识与第"+strconv.Itoa(firstLine)+"行重复")
|
|
}
|
|
seenInputs[normalized] = row.Line
|
|
order, err := h.orderCreator.CreateAdminOrder(ctx, &dto.CreateAdminOrderRequest{
|
|
Identifier: row.Identifier, PackageIDs: []uint{taskRecord.PackageID},
|
|
PaymentMethod: taskRecord.PaymentMethod, PaymentVoucherKey: []string(taskRecord.VoucherKeys),
|
|
}, buyerType, buyerID)
|
|
if err != nil {
|
|
return failedBatchOrderItem(item, publicBatchOrderError(err))
|
|
}
|
|
assetKey := order.AssetType + ":"
|
|
if order.IotCardID != nil {
|
|
assetKey += strconv.FormatUint(uint64(*order.IotCardID), 10)
|
|
} else if order.DeviceID != nil {
|
|
assetKey += strconv.FormatUint(uint64(*order.DeviceID), 10)
|
|
}
|
|
if firstLine, exists := seenAssets[assetKey]; assetKey != ":" && exists {
|
|
return failedBatchOrderItem(item, "资产与第"+strconv.Itoa(firstLine)+"行解析为同一资产")
|
|
}
|
|
seenAssets[assetKey] = row.Line
|
|
item.Status, item.OrderID, item.OrderNo, item.Amount = constants.AssetPackageBatchOrderItemStatusSuccess, order.ID, order.OrderNo, order.TotalAmount
|
|
return item
|
|
}
|
|
|
|
func failedBatchOrderItem(item model.AssetPackageBatchOrderResultItem, reason string) model.AssetPackageBatchOrderResultItem {
|
|
item.Status, item.Reason = constants.AssetPackageBatchOrderItemStatusFailed, reason
|
|
return item
|
|
}
|
|
|
|
func publicBatchOrderError(err error) string {
|
|
var appErr *apperrors.AppError
|
|
if stderrors.As(err, &appErr) && appErr.Message != "" {
|
|
return appErr.Message
|
|
}
|
|
return "创建订单失败"
|
|
}
|
|
|
|
func isAssetIdentifierHeader(value string) bool {
|
|
switch strings.ToLower(strings.TrimSpace(value)) {
|
|
case "资产标识", "identifier", "iccid", "iccid/虚拟号":
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
type assetPackageBatchOrderError string
|
|
|
|
func (e assetPackageBatchOrderError) Error() string { return string(e) }
|