350 lines
13 KiB
Go
350 lines
13 KiB
Go
package sybinnercode
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"io"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
|
|
"go-admin/app/goauto/models"
|
|
|
|
"github.com/google/uuid"
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/clause"
|
|
)
|
|
|
|
type MatchEnqueuer interface {
|
|
Start(context.Context, *gorm.DB, string) error
|
|
}
|
|
|
|
type Service struct {
|
|
db *gorm.DB
|
|
clock Clock
|
|
matcher MatchEnqueuer
|
|
}
|
|
|
|
func NewService(db *gorm.DB) *Service { return &Service{db: db, clock: time.Now} }
|
|
func (s *Service) WithMatchEnqueuer(matcher MatchEnqueuer) *Service { s.matcher = matcher; return s }
|
|
|
|
func (s *Service) Import(ctx context.Context, reader io.Reader, filename string, actor uint64, requestID string) (ImportResult, error) {
|
|
requestID = strings.TrimSpace(requestID)
|
|
if uuid.Validate(requestID) != nil || actor == 0 {
|
|
return ImportResult{}, invalid("requestId 或操作人无效")
|
|
}
|
|
if replay, ok, err := loadMutation[ImportResult](s.db.WithContext(ctx), requestID, "import"); err != nil {
|
|
return ImportResult{}, internal(err)
|
|
} else if ok {
|
|
return replay, nil
|
|
}
|
|
records, totalRows, err := parseWorkbook(reader, filename)
|
|
if err != nil {
|
|
return ImportResult{}, err
|
|
}
|
|
result := ImportResult{BusinessDate: records[0].BusinessDate, TotalRows: totalRows, RecordCount: len(records)}
|
|
for _, record := range records {
|
|
result.ItemCount += len(record.Items)
|
|
}
|
|
err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
if replay, ok, replayErr := loadMutation[ImportResult](tx, requestID, "import"); replayErr != nil {
|
|
return replayErr
|
|
} else if ok {
|
|
result = replay
|
|
return nil
|
|
}
|
|
var existing []models.SYBInnerCodeRecord
|
|
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("business_date = ?", result.BusinessDate).Find(&existing).Error; err != nil {
|
|
return err
|
|
}
|
|
if len(existing) > 0 {
|
|
for _, row := range existing {
|
|
if row.Status != models.SYBInnerCodePending {
|
|
return conflict("该营业日期已有匹配或执行证据,不能整批替换;请先核对并删除旧数据")
|
|
}
|
|
}
|
|
ids := make([]uint64, 0, len(existing))
|
|
for _, row := range existing {
|
|
ids = append(ids, row.ID)
|
|
}
|
|
if err := tx.Where("record_id IN ?", ids).Delete(&models.SYBInnerCodeItem{}).Error; err != nil {
|
|
return err
|
|
}
|
|
if err := tx.Where("id IN ?", ids).Delete(&models.SYBInnerCodeRecord{}).Error; err != nil {
|
|
return err
|
|
}
|
|
result.ReplacedCount = len(existing)
|
|
}
|
|
recordIDs := make([]uint64, 0, len(records))
|
|
for _, parsed := range records {
|
|
record := models.SYBInnerCodeRecord{BusinessDate: parsed.BusinessDate, OrderNumber: parsed.OrderNumber, Stall: parsed.Stall, SpecKey: parsed.SpecKey, SpecRaw: parsed.SpecRaw, SourceSKURaw: parsed.SourceSKURaw, ShopName: parsed.ShopName, SourceRow: parsed.SourceRow, PrintSequence: parsed.PrintSequence, Status: models.SYBInnerCodePending, CreatedBy: actor, ImportRequestID: requestID}
|
|
if err := tx.Create(&record).Error; err != nil {
|
|
return err
|
|
}
|
|
recordIDs = append(recordIDs, record.ID)
|
|
items := make([]models.SYBInnerCodeItem, 0, len(parsed.Items))
|
|
for _, item := range parsed.Items {
|
|
items = append(items, models.SYBInnerCodeItem{RecordID: record.ID, BusinessDate: parsed.BusinessDate, Code: item.Code, Ordinal: item.Ordinal, SourceRow: item.SourceRow})
|
|
}
|
|
if err := tx.Create(&items).Error; err != nil {
|
|
return err
|
|
}
|
|
}
|
|
recordIDsJSON, _ := json.Marshal(recordIDs)
|
|
job := models.SYBInnerCodeMatchJob{ID: uuid.NewString(), WorkKey: "import:" + requestID, BusinessDate: result.BusinessDate, RecordIDsJSON: string(recordIDsJSON), Status: "pending", Total: result.RecordCount}
|
|
if err := tx.Create(&job).Error; err != nil {
|
|
return err
|
|
}
|
|
result.MatchJobID = job.ID
|
|
return saveMutation(tx, requestID, "import", result)
|
|
})
|
|
if err != nil {
|
|
var serviceErr *ServiceError
|
|
if errors.As(err, &serviceErr) {
|
|
return ImportResult{}, err
|
|
}
|
|
return ImportResult{}, internal(err)
|
|
}
|
|
if s.matcher != nil {
|
|
if enqueueErr := s.matcher.Start(ctx, s.db, result.MatchJobID); enqueueErr != nil {
|
|
return ImportResult{}, internal(enqueueErr)
|
|
}
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func (s *Service) List(ctx context.Context, request ListRequest) (ListResult, error) {
|
|
if request.Page < 1 {
|
|
request.Page = 1
|
|
}
|
|
if request.PageSize == 0 {
|
|
request.PageSize = 100
|
|
}
|
|
if request.PageSize != 20 && request.PageSize != 50 && request.PageSize != 100 && request.PageSize != 200 {
|
|
return ListResult{}, invalid("pageSize 只允许 20、50、100、200")
|
|
}
|
|
query := s.db.WithContext(ctx).Model(&models.SYBInnerCodeRecord{})
|
|
if request.DateFrom != "" {
|
|
query = query.Where("business_date >= ?", request.DateFrom)
|
|
}
|
|
if request.DateTo != "" {
|
|
query = query.Where("business_date <= ?", request.DateTo)
|
|
}
|
|
if request.Status != "" {
|
|
query = query.Where("status = ?", request.Status)
|
|
}
|
|
if keyword := strings.TrimSpace(request.Keyword); keyword != "" {
|
|
like := "%" + keyword + "%"
|
|
query = query.Where("order_number LIKE ? OR EXISTS (SELECT 1 FROM syb_inner_code_item i WHERE i.record_id = syb_inner_code_record.id AND i.code LIKE ?)", like, like)
|
|
}
|
|
var total int64
|
|
if err := query.Count(&total).Error; err != nil {
|
|
return ListResult{}, internal(err)
|
|
}
|
|
var items []models.SYBInnerCodeRecord
|
|
if err := query.Preload("Items", func(db *gorm.DB) *gorm.DB { return db.Order("ordinal ASC") }).Order("business_date DESC, source_row ASC, id ASC").Limit(request.PageSize).Offset((request.Page - 1) * request.PageSize).Find(&items).Error; err != nil {
|
|
return ListResult{}, internal(err)
|
|
}
|
|
if len(items) > 0 {
|
|
ids := make([]uint64, 0, len(items))
|
|
byID := map[uint64]*models.SYBInnerCodeRecord{}
|
|
for i := range items {
|
|
ids = append(ids, items[i].ID)
|
|
byID[items[i].ID] = &items[i]
|
|
}
|
|
var plans []models.SYBInnerCodePlan
|
|
if err := s.db.WithContext(ctx).Where("record_id IN ?", ids).Find(&plans).Error; err != nil {
|
|
return ListResult{}, internal(err)
|
|
}
|
|
for i := range plans {
|
|
byID[plans[i].RecordID].Plan = &plans[i]
|
|
}
|
|
}
|
|
return ListResult{Items: items, Total: total, Page: request.Page, PageSize: request.PageSize}, nil
|
|
}
|
|
|
|
func (s *Service) Detail(ctx context.Context, id uint64) (models.SYBInnerCodeRecord, error) {
|
|
var record models.SYBInnerCodeRecord
|
|
if id == 0 {
|
|
return record, invalid("记录 ID 无效")
|
|
}
|
|
if err := s.db.WithContext(ctx).Preload("Items", func(db *gorm.DB) *gorm.DB { return db.Order("ordinal ASC") }).Preload("Plan").First(&record, id).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return record, &ServiceError{Code: CodeNotFound, Message: "档口入库码记录不存在"}
|
|
}
|
|
return record, internal(err)
|
|
}
|
|
return record, nil
|
|
}
|
|
|
|
func (s *Service) QueueRematch(ctx context.Context, request RematchRequest) (RematchResult, error) {
|
|
request.RequestID = strings.TrimSpace(request.RequestID)
|
|
ids := uniqueIDs(request.IDs)
|
|
if uuid.Validate(request.RequestID) != nil || len(ids) == 0 {
|
|
return RematchResult{}, invalid("requestId 或记录无效")
|
|
}
|
|
if replay, ok, err := loadMutation[RematchResult](s.db.WithContext(ctx), request.RequestID, "rematch"); err != nil {
|
|
return RematchResult{}, internal(err)
|
|
} else if ok {
|
|
return replay, nil
|
|
}
|
|
result := RematchResult{MatchJobID: uuid.NewString(), Queued: len(ids)}
|
|
var date string
|
|
err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
var records []models.SYBInnerCodeRecord
|
|
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id IN ?", ids).Find(&records).Error; err != nil {
|
|
return err
|
|
}
|
|
if len(records) != len(ids) {
|
|
return conflict("部分记录不存在")
|
|
}
|
|
for _, record := range records {
|
|
if record.Status != models.SYBInnerCodeFailed && record.Status != models.SYBInnerCodeSkipped {
|
|
return conflict("只有读取失败或匹配受限记录可以重新匹配")
|
|
}
|
|
if date == "" {
|
|
date = record.BusinessDate
|
|
} else if date != record.BusinessDate {
|
|
return conflict("重新匹配记录必须属于同一营业日期")
|
|
}
|
|
}
|
|
if err := tx.Where("record_id IN ?", ids).Delete(&models.SYBInnerCodePlan{}).Error; err != nil {
|
|
return err
|
|
}
|
|
if err := tx.Model(&models.SYBInnerCodeRecord{}).Where("id IN ?", ids).Updates(map[string]any{"status": models.SYBInnerCodePending, "result_message": "等待重新匹配"}).Error; err != nil {
|
|
return err
|
|
}
|
|
recordIDsJSON, _ := json.Marshal(ids)
|
|
job := models.SYBInnerCodeMatchJob{ID: result.MatchJobID, WorkKey: "rematch:" + request.RequestID, BusinessDate: date, RecordIDsJSON: string(recordIDsJSON), Status: "pending", Total: len(ids)}
|
|
if err := tx.Create(&job).Error; err != nil {
|
|
return err
|
|
}
|
|
return saveMutation(tx, request.RequestID, "rematch", result)
|
|
})
|
|
if err != nil {
|
|
var target *ServiceError
|
|
if errors.As(err, &target) {
|
|
return RematchResult{}, err
|
|
}
|
|
return RematchResult{}, internal(err)
|
|
}
|
|
if s.matcher != nil {
|
|
if err := s.matcher.Start(ctx, s.db, result.MatchJobID); err != nil {
|
|
return RematchResult{}, internal(err)
|
|
}
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func (s *Service) Delete(ctx context.Context, actor uint64, request DeleteRequest) (DeleteResult, error) {
|
|
request.RequestID = strings.TrimSpace(request.RequestID)
|
|
if actor == 0 || uuid.Validate(request.RequestID) != nil {
|
|
return DeleteResult{}, invalid("requestId 或操作人无效")
|
|
}
|
|
ids := uniqueIDs(request.IDs)
|
|
if len(ids) == 0 || len(ids) > 200 {
|
|
return DeleteResult{}, invalid("一次必须选择 1 至 200 条记录")
|
|
}
|
|
if replay, ok, err := loadMutation[DeleteResult](s.db.WithContext(ctx), request.RequestID, "delete"); err != nil {
|
|
return DeleteResult{}, internal(err)
|
|
} else if ok {
|
|
return replay, nil
|
|
}
|
|
result := DeleteResult{}
|
|
err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
if replay, ok, replayErr := loadMutation[DeleteResult](tx, request.RequestID, "delete"); replayErr != nil {
|
|
return replayErr
|
|
} else if ok {
|
|
result = replay
|
|
return nil
|
|
}
|
|
var records []models.SYBInnerCodeRecord
|
|
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id IN ?", ids).Find(&records).Error; err != nil {
|
|
return err
|
|
}
|
|
if len(records) != len(ids) {
|
|
return conflict("部分记录不存在,未删除任何数据")
|
|
}
|
|
for _, record := range records {
|
|
if record.Status == models.SYBInnerCodeQueued || record.Status == models.SYBInnerCodeApplying || record.Status == models.SYBInnerCodeNeedsCheck {
|
|
result.Blocked = append(result.Blocked, BlockedRecord{ID: record.ID, Status: record.Status})
|
|
}
|
|
}
|
|
if len(result.Blocked) > 0 {
|
|
sort.Slice(result.Blocked, func(i, j int) bool { return result.Blocked[i].ID < result.Blocked[j].ID })
|
|
return &ServiceError{Code: CodeConflict, Message: "选中记录包含排队中、回写中或需复核状态,未删除任何数据", Details: map[string]any{"blocked": result.Blocked}}
|
|
}
|
|
// The state gate above is the dynamic restriction for active writeback
|
|
// evidence. Terminal evidence belongs to imported data and is physically
|
|
// removed in the same transaction as the selected records.
|
|
if err := tx.Where("record_id IN ?", ids).Delete(&models.SYBInnerCodeCheckpoint{}).Error; err != nil {
|
|
return err
|
|
}
|
|
if err := tx.Where("record_id IN ?", ids).Delete(&models.SYBInnerCodeApplyItem{}).Error; err != nil {
|
|
return err
|
|
}
|
|
if err := tx.Where("record_id IN ?", ids).Delete(&models.SYBInnerCodePlan{}).Error; err != nil {
|
|
return err
|
|
}
|
|
if err := tx.Where("record_id IN ?", ids).Delete(&models.SYBInnerCodeItem{}).Error; err != nil {
|
|
return err
|
|
}
|
|
deleted := tx.Where("id IN ?", ids).Delete(&models.SYBInnerCodeRecord{})
|
|
if deleted.Error != nil {
|
|
return deleted.Error
|
|
}
|
|
if deleted.RowsAffected != int64(len(ids)) {
|
|
return conflict("记录状态发生变化,未完成删除")
|
|
}
|
|
result.Deleted = int(deleted.RowsAffected)
|
|
return saveMutation(tx, request.RequestID, "delete", result)
|
|
})
|
|
if err != nil {
|
|
var serviceErr *ServiceError
|
|
if errors.As(err, &serviceErr) {
|
|
return result, err
|
|
}
|
|
return result, internal(err)
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func uniqueIDs(ids []uint64) []uint64 {
|
|
seen := map[uint64]bool{}
|
|
result := make([]uint64, 0, len(ids))
|
|
for _, id := range ids {
|
|
if id > 0 && !seen[id] {
|
|
seen[id] = true
|
|
result = append(result, id)
|
|
}
|
|
}
|
|
sort.Slice(result, func(i, j int) bool { return result[i] < result[j] })
|
|
return result
|
|
}
|
|
func saveMutation(db *gorm.DB, requestID, action string, result any) error {
|
|
payload, err := json.Marshal(result)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return db.Create(&models.SYBInnerCodeMutation{RequestID: requestID, Action: action, ResultJSON: string(payload)}).Error
|
|
}
|
|
func loadMutation[T any](db *gorm.DB, requestID, action string) (T, bool, error) {
|
|
var zero T
|
|
var row models.SYBInnerCodeMutation
|
|
err := db.Where("request_id = ?", requestID).First(&row).Error
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return zero, false, nil
|
|
}
|
|
if err != nil {
|
|
return zero, false, err
|
|
}
|
|
if row.Action != action {
|
|
return zero, false, conflict("requestId 已用于其他操作")
|
|
}
|
|
if err := json.Unmarshal([]byte(row.ResultJSON), &zero); err != nil {
|
|
return zero, false, err
|
|
}
|
|
return zero, true, nil
|
|
}
|