Files
goauto/server/app/goauto/sybinnercode/apply.go
T

415 lines
17 KiB
Go

package sybinnercode
import (
"context"
"encoding/json"
"errors"
"fmt"
"sort"
"strings"
"time"
"go-admin/app/goauto/models"
"go-admin/app/goauto/sybclient"
"github.com/google/uuid"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
type InnerCodeWriter interface {
MatchReader
DeleteInnerCode(context.Context, int64) error
CreateInnerCodeDetail(context.Context, int64, string) (int64, error)
UpdateDetailCode(context.Context, int64, int64, string) error
}
type ApplyRunner struct {
Factory func(context.Context, *gorm.DB) (InnerCodeWriter, error)
}
func (r ApplyRunner) Start(db *gorm.DB, batchID string) {
go func() {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute)
defer cancel()
writer, err := r.Factory(ctx, db)
if err == nil {
err = RunApplyBatch(ctx, db, writer, batchID)
}
if err != nil {
_ = interruptBatch(db, batchID, err.Error())
}
}()
}
func (s *Service) PreviewApply(ctx context.Context, ids []uint64) (ApplyPreview, error) {
ids = uniqueIDs(ids)
if len(ids) == 0 {
return ApplyPreview{}, invalid("没有选择记录")
}
var records []models.SYBInnerCodeRecord
if err := s.db.WithContext(ctx).Preload("Items").Where("id IN ?", ids).Find(&records).Error; err != nil {
return ApplyPreview{}, internal(err)
}
if len(records) != len(ids) {
return ApplyPreview{}, conflict("部分记录不存在")
}
preview := ApplyPreview{Records: len(records)}
for _, record := range records {
if record.Status != models.SYBInnerCodeReady {
preview.Blocked = append(preview.Blocked, BlockedRecord{ID: record.ID, Status: record.Status})
continue
}
var plan models.SYBInnerCodePlan
if err := s.db.WithContext(ctx).First(&plan, "record_id = ?", record.ID).Error; err != nil {
return ApplyPreview{}, conflict("部分记录缺少有效匹配计划")
}
preview.InboundCodes += len(record.Items)
preview.PlaceholderDetails += plan.PlaceholderCount
preview.ReplaceOldCodes += plan.ReplaceOldCodeCount
}
sort.Slice(preview.Blocked, func(i, j int) bool { return preview.Blocked[i].ID < preview.Blocked[j].ID })
return preview, nil
}
func (s *Service) QueueApply(ctx context.Context, actor uint64, request ApplyRequest) (ApplyResult, error) {
request.RequestID = strings.TrimSpace(request.RequestID)
ids := uniqueIDs(request.IDs)
if actor == 0 || uuid.Validate(request.RequestID) != nil || len(ids) == 0 {
return ApplyResult{}, invalid("requestId、操作人或记录无效")
}
if replay, ok, err := loadMutation[ApplyResult](s.db.WithContext(ctx), request.RequestID, "apply"); err != nil {
return ApplyResult{}, internal(err)
} else if ok {
return replay, nil
}
result := ApplyResult{BatchID: uuid.NewString(), Queued: len(ids)}
err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if replay, ok, replayErr := loadMutation[ApplyResult](tx, request.RequestID, "apply"); 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.SYBInnerCodeReady {
return conflict(fmt.Sprintf("记录 %d 状态为 %s,不能回写", record.ID, record.Status))
}
var count int64
if err := tx.Model(&models.SYBInnerCodePlan{}).Where("record_id = ?", record.ID).Count(&count).Error; err != nil || count != 1 {
return conflict(fmt.Sprintf("记录 %d 缺少唯一有效计划", record.ID))
}
}
batch := models.SYBInnerCodeApplyBatch{ID: result.BatchID, RequestID: request.RequestID, Status: "queued", Requested: len(ids), CreatedBy: actor}
if err := tx.Create(&batch).Error; err != nil {
return err
}
for _, id := range ids {
item := models.SYBInnerCodeApplyItem{BatchID: batch.ID, RecordID: id, Status: models.SYBInnerCodeQueued}
if err := tx.Create(&item).Error; err != nil {
return err
}
}
changed := tx.Model(&models.SYBInnerCodeRecord{}).Where("id IN ? AND status = ?", ids, models.SYBInnerCodeReady).Updates(map[string]any{"status": models.SYBInnerCodeQueued, "apply_batch_id": batch.ID})
if changed.Error != nil {
return changed.Error
}
if changed.RowsAffected != int64(len(ids)) {
return conflict("记录状态并发变化,未提交回写")
}
return saveMutation(tx, request.RequestID, "apply", result)
})
if err != nil {
var target *ServiceError
if errors.As(err, &target) {
return ApplyResult{}, err
}
return ApplyResult{}, internal(err)
}
return result, nil
}
func acquireLease(ctx context.Context, db *gorm.DB, owner string) (bool, error) {
now := time.Now().UTC()
expires := now.Add(45 * time.Second)
seed := models.SYBInnerCodeWorkerLease{ID: 1, OwnerID: "", ExpiresAt: time.Unix(0, 0).UTC()}
if err := db.WithContext(ctx).Clauses(clause.OnConflict{DoNothing: true}).Create(&seed).Error; err != nil {
return false, err
}
result := db.WithContext(ctx).Model(&models.SYBInnerCodeWorkerLease{}).Where("id = 1 AND (expires_at < ? OR owner_id = ?)", now, owner).Updates(map[string]any{"owner_id": owner, "expires_at": expires, "version": gorm.Expr("version + 1")})
return result.RowsAffected == 1, result.Error
}
func releaseLease(db *gorm.DB, owner string) {
db.Model(&models.SYBInnerCodeWorkerLease{}).Where("id = 1 AND owner_id = ?", owner).Updates(map[string]any{"owner_id": "", "expires_at": time.Unix(0, 0).UTC()})
}
func RunApplyBatch(ctx context.Context, db *gorm.DB, writer InnerCodeWriter, batchID string) error {
owner := uuid.NewString()
ok, err := acquireLease(ctx, db, owner)
if err != nil {
return err
}
if !ok {
return conflict("已有其他服务实例正在执行 SYB 回写")
}
defer releaseLease(db, owner)
startedBatch := db.Model(&models.SYBInnerCodeApplyBatch{}).Where("id = ? AND status = ?", batchID, "queued").Updates(map[string]any{"status": "running", "started_at": time.Now().UTC()})
if startedBatch.Error != nil {
return startedBatch.Error
}
if startedBatch.RowsAffected != 1 {
return conflict("回写批次不存在或状态已变化")
}
for {
var item models.SYBInnerCodeApplyItem
err := db.WithContext(ctx).Where("batch_id = ? AND status = ?", batchID, models.SYBInnerCodeQueued).Order("id").First(&item).Error
if errors.Is(err, gorm.ErrRecordNotFound) {
break
}
if err != nil {
return err
}
started := time.Now().UTC()
claim := db.Model(&models.SYBInnerCodeRecord{}).Where("id = ? AND status = ? AND apply_batch_id = ?", item.RecordID, models.SYBInnerCodeQueued, batchID).Updates(map[string]any{"status": models.SYBInnerCodeApplying, "apply_started_at": started})
if claim.Error != nil {
return claim.Error
}
if claim.RowsAffected == 0 {
if err := db.Model(&item).Updates(map[string]any{"status": models.SYBInnerCodeSkipped, "message": "记录状态已变化"}).Error; err != nil {
return err
}
continue
}
if err := db.Model(&item).Update("status", models.SYBInnerCodeApplying).Error; err != nil {
return err
}
status, message := applyOne(ctx, db, writer, item.RecordID)
finished := time.Now().UTC()
if err := db.Transaction(func(tx *gorm.DB) error {
if err := tx.Model(&models.SYBInnerCodeRecord{}).Where("id = ? AND status = ?", item.RecordID, models.SYBInnerCodeApplying).Updates(map[string]any{"status": status, "result_message": message, "applied_at": finished}).Error; err != nil {
return err
}
return tx.Model(&models.SYBInnerCodeApplyItem{}).Where("id = ?", item.ID).Updates(map[string]any{"status": status, "message": message}).Error
}); err != nil {
return err
}
if err := db.Model(&models.SYBInnerCodeApplyBatch{}).Where("id = ?", batchID).Updates(map[string]any{"processed": gorm.Expr("processed + 1")}).Error; err != nil {
return err
}
}
finished := time.Now().UTC()
return db.Model(&models.SYBInnerCodeApplyBatch{}).Where("id = ?", batchID).Updates(map[string]any{"status": "finished", "finished_at": finished}).Error
}
func applyOne(ctx context.Context, db *gorm.DB, writer InnerCodeWriter, recordID uint64) (string, string) {
var record models.SYBInnerCodeRecord
if err := db.Preload("Items", func(q *gorm.DB) *gorm.DB { return q.Order("ordinal") }).First(&record, recordID).Error; err != nil {
return models.SYBInnerCodeNeedsCheck, "读取本地记录失败"
}
var plan models.SYBInnerCodePlan
if err := db.First(&plan, "record_id = ?", recordID).Error; err != nil {
return models.SYBInnerCodeNeedsCheck, "读取匹配计划失败"
}
stock, primary, err := readStock(ctx, writer, plan.StockID, plan.DetailID)
if err != nil {
return models.SYBInnerCodeNeedsCheck, "执行前重读 SYB 失败:" + compact(err.Error(), 800)
}
if primary.ProductSpec != plan.SYBSpec || rawText(primary.Raw["sku"]) != plan.SYBSKU || rawText(primary.Raw["variationSku"]) != plan.SYBVariationSKU {
return models.SYBInnerCodeFailed, "SYB 商品身份已变化,停止回写,请重新匹配"
}
var items []plannedRemoteItem
if err := json.Unmarshal([]byte(plan.RemoteItemsJSON), &items); err != nil || len(items) != len(record.Items) {
return models.SYBInnerCodeNeedsCheck, "匹配计划不完整,禁止回写"
}
current := rawText(primary.Raw["innerExpCode"])
if current != plan.RemoteInnerCode {
return models.SYBInnerCodeNeedsCheck, "原商品入库码在规划后变化,停止回写"
}
for _, target := range items {
if target.DetailID == 0 {
continue
}
remote, found := stockDetail(stock, target.DetailID)
if !found || NormalizeSpecKey(remote.ProductSpec) != NormalizeSpecKey(plan.SYBSpec) || rawText(remote.Raw["purchasePlatform"]) != plan.PurchasePlatform || rawText(remote.Raw["purchaseCode"]) != plan.PurchaseCode || rawText(remote.Raw["innerExpCode"]) != target.RemoteCode {
return models.SYBInnerCodeNeedsCheck, "SYB 商品明细或既有入库码在规划后变化,停止回写"
}
}
if plan.ReplaceOldCodeCount > 0 && current != "" {
if err := checkpoint(db, recordID, "before", "delete_old", items); err != nil {
return models.SYBInnerCodeNeedsCheck, "删除旧码前检查点失败"
}
if err := writer.DeleteInnerCode(ctx, plan.DetailID); err != nil {
return writeFailure(err, "删除旧入库码")
}
if err := checkpoint(db, recordID, "after", "delete_old", items); err != nil {
return models.SYBInnerCodeNeedsCheck, "删除旧码后检查点失败"
}
}
for index := range items {
target := &items[index]
if target.DetailID > 0 && countCode(stock, target.Code) == 1 && detailCode(stock, target.DetailID) == target.Code {
continue
}
if target.DetailID == 0 {
target.Source = "created"
if err := checkpoint(db, recordID, "before", "create_detail", items); err != nil {
return models.SYBInnerCodeNeedsCheck, "创建明细前检查点失败"
}
id, err := writer.CreateInnerCodeDetail(ctx, plan.StockID, fmt.Sprintf("档口入库码-%d-%d", recordID, index+1))
if err != nil {
return writeFailure(err, "创建零价明细")
}
target.DetailID = id
if err := checkpoint(db, recordID, "after", "create_detail", items); err != nil {
return models.SYBInnerCodeNeedsCheck, "创建明细后检查点失败"
}
}
if err := checkpoint(db, recordID, "before", "write_code", items); err != nil {
return models.SYBInnerCodeNeedsCheck, "写码前检查点失败"
}
if err := writer.UpdateDetailCode(ctx, plan.StockID, target.DetailID, target.Code); err != nil {
return writeFailure(err, "写入入库码")
}
verified, _, err := readStock(ctx, writer, plan.StockID, plan.DetailID)
if err != nil {
return models.SYBInnerCodeNeedsCheck, "写入响应后重读失败,禁止自动重试"
}
if countCode(verified, target.Code) != 1 || detailCode(verified, target.DetailID) != target.Code {
return models.SYBInnerCodeNeedsCheck, "写后入库码未唯一出现在预期明细,禁止自动重试"
}
stock = verified
if err := checkpoint(db, recordID, "after", "write_code", items); err != nil {
return models.SYBInnerCodeNeedsCheck, "写入确认后检查点失败"
}
}
return models.SYBInnerCodeUpdated, fmt.Sprintf("回写完成:%d 个入库码均已逐件回读确认", len(items))
}
func writeFailure(err error, action string) (string, string) {
if errors.Is(err, sybclient.ErrWriteResultUnknown) {
return models.SYBInnerCodeNeedsCheck, action + "结果不明确,禁止自动重试"
}
return models.SYBInnerCodeFailed, action + "失败:" + compact(err.Error(), 800)
}
func checkpoint(db *gorm.DB, recordID uint64, phase, action string, state any) error {
raw, err := json.Marshal(state)
if err != nil {
return err
}
return db.Transaction(func(tx *gorm.DB) error {
var max int
tx.Model(&models.SYBInnerCodeCheckpoint{}).Where("record_id = ?", recordID).Select("COALESCE(MAX(sequence),0)").Scan(&max)
return tx.Create(&models.SYBInnerCodeCheckpoint{RecordID: recordID, Sequence: max + 1, Phase: phase, Action: action, StateJSON: string(raw)}).Error
})
}
func readStock(ctx context.Context, reader MatchReader, stockID, detailID int64) (sybclient.StockDetail, sybclient.DetailItem, error) {
stocks, err := reader.DetailListByStock(ctx, []int64{stockID})
if err != nil {
return sybclient.StockDetail{}, sybclient.DetailItem{}, err
}
if len(stocks) != 1 || stocks[0].ID != stockID {
return sybclient.StockDetail{}, sybclient.DetailItem{}, fmt.Errorf("货运单不唯一")
}
var found []sybclient.DetailItem
for _, item := range stocks[0].Details {
if item.ID == detailID {
found = append(found, item)
}
}
if len(found) != 1 {
return sybclient.StockDetail{}, sybclient.DetailItem{}, fmt.Errorf("商品明细不唯一")
}
return stocks[0], found[0], nil
}
func countCode(stock sybclient.StockDetail, code string) int {
count := 0
for _, item := range stock.Details {
if rawText(item.Raw["innerExpCode"]) == code {
count++
}
}
return count
}
func detailCode(stock sybclient.StockDetail, id int64) string {
if item, ok := stockDetail(stock, id); ok {
return rawText(item.Raw["innerExpCode"])
}
return ""
}
func stockDetail(stock sybclient.StockDetail, id int64) (sybclient.DetailItem, bool) {
for _, item := range stock.Details {
if item.ID == id {
return item, true
}
}
return sybclient.DetailItem{}, false
}
func Recheck(ctx context.Context, db *gorm.DB, reader MatchReader, recordID uint64) (string, string, error) {
var record models.SYBInnerCodeRecord
if err := db.Preload("Items").First(&record, recordID).Error; err != nil {
return "", "", err
}
if record.Status != models.SYBInnerCodeNeedsCheck {
return "", "", conflict("只有需复核记录可以执行只读复核")
}
var plan models.SYBInnerCodePlan
if err := db.First(&plan, "record_id = ?", recordID).Error; err != nil {
return "", "", err
}
stock, _, err := readStock(ctx, reader, plan.StockID, plan.DetailID)
if err != nil {
return "", "", err
}
all, none := true, true
for _, item := range record.Items {
count := countCode(stock, item.Code)
if count != 1 {
all = false
}
if count > 0 {
none = false
}
}
status, message := models.SYBInnerCodeNeedsCheck, "远端状态仍不明确,保持需复核"
if all {
status = models.SYBInnerCodeUpdated
message = "只读复核确认全部入库码已正确写入"
} else if none {
status = models.SYBInnerCodeReady
message = "只读复核确认未写入;如需回写必须重新人工确认"
}
err = db.Model(&models.SYBInnerCodeRecord{}).Where("id = ? AND status = ?", recordID, models.SYBInnerCodeNeedsCheck).Updates(map[string]any{"status": status, "result_message": message}).Error
return status, message, err
}
func RecoverInterrupted(db *gorm.DB) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.Model(&models.SYBInnerCodeRecord{}).Where("status = ?", models.SYBInnerCodeQueued).Updates(map[string]any{"status": models.SYBInnerCodeReady, "apply_batch_id": nil, "result_message": "服务重启,未开始的队列已释放,请重新确认"}).Error; err != nil {
return err
}
if err := tx.Model(&models.SYBInnerCodeRecord{}).Where("status = ?", models.SYBInnerCodeApplying).Updates(map[string]any{"status": models.SYBInnerCodeNeedsCheck, "result_message": "服务在远端动作期间重启,只允许只读复核"}).Error; err != nil {
return err
}
return tx.Model(&models.SYBInnerCodeApplyBatch{}).Where("status IN ?", []string{"queued", "running"}).Update("status", "interrupted").Error
})
}
func interruptBatch(db *gorm.DB, batchID, message string) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.Model(&models.SYBInnerCodeRecord{}).Where("apply_batch_id = ? AND status = ?", batchID, models.SYBInnerCodeApplying).Updates(map[string]any{"status": models.SYBInnerCodeNeedsCheck, "result_message": "后台异常中断,只允许只读复核"}).Error; err != nil {
return err
}
if err := tx.Model(&models.SYBInnerCodeRecord{}).Where("apply_batch_id = ? AND status = ?", batchID, models.SYBInnerCodeQueued).Updates(map[string]any{"status": models.SYBInnerCodeReady, "apply_batch_id": nil, "result_message": "后台执行未开始,请重新确认"}).Error; err != nil {
return err
}
return tx.Model(&models.SYBInnerCodeApplyBatch{}).Where("id = ?", batchID).Updates(map[string]any{"status": "interrupted", "finished_at": time.Now().UTC()}).Error
})
}