相关问题优化以及新功能开发
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 7m57s
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 7m57s
迁移 155- 157
This commit is contained in:
185
internal/service/order_package_invalidate/service.go
Normal file
185
internal/service/order_package_invalidate/service.go
Normal file
@@ -0,0 +1,185 @@
|
||||
package order_package_invalidate
|
||||
|
||||
import (
|
||||
"context"
|
||||
"path/filepath"
|
||||
"time"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
|
||||
"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"
|
||||
"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"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/middleware"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/queue"
|
||||
)
|
||||
|
||||
// Service 订单套餐失效任务业务逻辑层
|
||||
type Service struct {
|
||||
taskStore *postgres.OrderPackageInvalidateTaskStore
|
||||
queueClient *queue.Client
|
||||
}
|
||||
|
||||
// New 创建 Service 实例
|
||||
func New(
|
||||
taskStore *postgres.OrderPackageInvalidateTaskStore,
|
||||
queueClient *queue.Client,
|
||||
) *Service {
|
||||
return &Service{
|
||||
taskStore: taskStore,
|
||||
queueClient: queueClient,
|
||||
}
|
||||
}
|
||||
|
||||
// InvalidateTaskPayload Worker 任务载荷
|
||||
type InvalidateTaskPayload struct {
|
||||
TaskID uint `json:"task_id"`
|
||||
}
|
||||
|
||||
// Create 创建订单套餐失效任务
|
||||
func (s *Service) Create(ctx context.Context, req *dto.CreateOrderPackageInvalidateTaskRequest) (*dto.OrderPackageInvalidateTaskResponse, error) {
|
||||
userID := middleware.GetUserIDFromContext(ctx)
|
||||
if userID == 0 {
|
||||
return nil, errors.New(errors.CodeUnauthorized, "未授权访问")
|
||||
}
|
||||
|
||||
task := &model.OrderPackageInvalidateTask{
|
||||
TaskNo: s.taskStore.GenerateTaskNo(),
|
||||
Status: model.ImportTaskStatusPending,
|
||||
FileName: filepath.Base(req.FileKey),
|
||||
StorageKey: req.FileKey,
|
||||
VoucherKeys: model.StringJSONBArray(req.VoucherKeys),
|
||||
Remark: req.Remark,
|
||||
CreatorName: middleware.GetUsernameFromContext(ctx),
|
||||
FailedItems: model.ImportResultItems{},
|
||||
Creator: userID,
|
||||
Updater: userID,
|
||||
}
|
||||
|
||||
if err := s.taskStore.Create(ctx, task); err != nil {
|
||||
return nil, errors.Wrap(errors.CodeInternalError, err, "创建任务失败")
|
||||
}
|
||||
|
||||
payload := InvalidateTaskPayload{TaskID: task.ID}
|
||||
err := s.queueClient.EnqueueTask(
|
||||
ctx,
|
||||
constants.TaskTypeOrderPackageInvalidate,
|
||||
payload,
|
||||
asynq.Queue(constants.QueueForTaskType(constants.TaskTypeOrderPackageInvalidate)),
|
||||
)
|
||||
if err != nil {
|
||||
s.taskStore.UpdateStatus(ctx, task.ID, model.ImportTaskStatusFailed, "任务入队失败: "+err.Error())
|
||||
}
|
||||
|
||||
return s.toResponse(task), nil
|
||||
}
|
||||
|
||||
// List 分页查询任务列表
|
||||
func (s *Service) List(ctx context.Context, req *dto.ListOrderPackageInvalidateTaskRequest) (*dto.OrderPackageInvalidateTaskListResponse, error) {
|
||||
if req.Page <= 0 {
|
||||
req.Page = 1
|
||||
}
|
||||
if req.PageSize <= 0 {
|
||||
req.PageSize = constants.DefaultPageSize
|
||||
}
|
||||
|
||||
opts := &store.QueryOptions{
|
||||
Page: req.Page,
|
||||
PageSize: req.PageSize,
|
||||
}
|
||||
|
||||
filters := map[string]interface{}{}
|
||||
if req.Status != nil {
|
||||
filters["status"] = *req.Status
|
||||
}
|
||||
|
||||
tasks, total, err := s.taskStore.List(ctx, opts, filters)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(errors.CodeInternalError, err, "查询任务列表失败")
|
||||
}
|
||||
|
||||
list := make([]*dto.OrderPackageInvalidateTaskResponse, 0, len(tasks))
|
||||
for _, t := range tasks {
|
||||
list = append(list, s.toResponse(t))
|
||||
}
|
||||
|
||||
return &dto.OrderPackageInvalidateTaskListResponse{
|
||||
Total: total,
|
||||
Page: req.Page,
|
||||
PageSize: req.PageSize,
|
||||
List: list,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// GetByID 查询任务详情(含失败明细)
|
||||
func (s *Service) GetByID(ctx context.Context, id uint) (*dto.OrderPackageInvalidateTaskDetailResponse, error) {
|
||||
task, err := s.taskStore.GetByID(ctx, id)
|
||||
if err != nil {
|
||||
return nil, errors.New(errors.CodeNotFound, "任务不存在")
|
||||
}
|
||||
|
||||
failedItems := make([]dto.InvalidateFailedItem, 0, len(task.FailedItems))
|
||||
for _, item := range task.FailedItems {
|
||||
failedItems = append(failedItems, dto.InvalidateFailedItem{
|
||||
Line: item.Line,
|
||||
OrderNo: item.ICCID, // 复用 ImportResultItem.ICCID 存储 order_no
|
||||
Reason: item.Reason,
|
||||
})
|
||||
}
|
||||
|
||||
return &dto.OrderPackageInvalidateTaskDetailResponse{
|
||||
OrderPackageInvalidateTaskResponse: *s.toResponse(task),
|
||||
FailedItems: failedItems,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// toResponse 转换为响应 DTO
|
||||
func (s *Service) toResponse(task *model.OrderPackageInvalidateTask) *dto.OrderPackageInvalidateTaskResponse {
|
||||
resp := &dto.OrderPackageInvalidateTaskResponse{
|
||||
ID: task.ID,
|
||||
TaskNo: task.TaskNo,
|
||||
Status: task.Status,
|
||||
StatusName: invalidateTaskStatusName(task.Status),
|
||||
TotalCount: task.TotalCount,
|
||||
SuccessCount: task.SuccessCount,
|
||||
FailCount: task.FailCount,
|
||||
FileName: task.FileName,
|
||||
VoucherKeys: []string(task.VoucherKeys),
|
||||
Remark: task.Remark,
|
||||
CreatorName: task.CreatorName,
|
||||
ErrorMessage: task.ErrorMessage,
|
||||
CreatedAt: task.CreatedAt.Format(time.RFC3339),
|
||||
UpdatedAt: task.UpdatedAt.Format(time.RFC3339),
|
||||
}
|
||||
if task.StartedAt != nil {
|
||||
t := task.StartedAt.Format(time.RFC3339)
|
||||
resp.StartedAt = &t
|
||||
}
|
||||
if task.CompletedAt != nil {
|
||||
t := task.CompletedAt.Format(time.RFC3339)
|
||||
resp.CompletedAt = &t
|
||||
}
|
||||
if resp.VoucherKeys == nil {
|
||||
resp.VoucherKeys = []string{}
|
||||
}
|
||||
return resp
|
||||
}
|
||||
|
||||
// invalidateTaskStatusName 状态码转中文
|
||||
func invalidateTaskStatusName(status int) string {
|
||||
switch status {
|
||||
case model.ImportTaskStatusPending:
|
||||
return "待处理"
|
||||
case model.ImportTaskStatusProcessing:
|
||||
return "处理中"
|
||||
case model.ImportTaskStatusCompleted:
|
||||
return "已完成"
|
||||
case model.ImportTaskStatusFailed:
|
||||
return "失败"
|
||||
default:
|
||||
return "未知"
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user