Files
goauto/server/app/goauto/purchase/reset.go
T

214 lines
8.0 KiB
Go

package purchase
import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"go-admin/app/goauto/device"
"go-admin/app/goauto/models"
"go-admin/app/goauto/purchasecontract"
"go-admin/app/goauto/purchaserule"
"github.com/google/uuid"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
type PurchaseResetRequest struct {
RequestID string `json:"requestId"`
}
type PurchaseResetResponse struct {
TaskID uint64 `json:"taskId"`
TaskNo string `json:"taskNo"`
AttemptNumber int `json:"attemptNumber"`
Status string `json:"status"`
Replayed bool `json:"replayed,omitempty"`
}
// Reset restores one failed purchase task to pending without changing its
// product, target, mapping or price snapshots. It is intentionally separate
// from BatchRetry and AgentRetry, which continue to create replacement tasks.
func (s *Service) Reset(ctx context.Context, taskID uint64, req PurchaseResetRequest) (PurchaseResetResponse, error) {
return s.reset(ctx, taskID, req, nil)
}
func (s *Service) ResetForDevice(ctx context.Context, taskID uint64, req PurchaseResetRequest, token string) (PurchaseResetResponse, error) {
return s.reset(ctx, taskID, req, &token)
}
func (s *Service) reset(ctx context.Context, taskID uint64, req PurchaseResetRequest, token *string) (PurchaseResetResponse, error) {
if taskID == 0 {
return PurchaseResetResponse{}, fail(CodeTaskNotFound, "采购任务不存在")
}
requestID := strings.TrimSpace(req.RequestID)
if _, err := uuid.Parse(requestID); err != nil {
return PurchaseResetResponse{}, fail(CodeInvalidRequest, "requestId 无效")
}
attemptID := purchaseResetAttemptID(requestID, taskID)
response := PurchaseResetResponse{TaskID: taskID, TaskNo: taskNumber(taskID)}
err := s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
var authenticated *models.AgentDevice
if token != nil {
record, err := device.NewService(tx).Authenticate(ctx, *token)
if err != nil {
return err
}
authenticated = &record
}
var task models.PurchaseTask
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&task, taskID).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return fail(CodeTaskNotFound, "采购任务不存在")
}
return internal(err)
}
if authenticated != nil && (task.DeviceID == nil || *task.DeviceID != authenticated.ID || task.CreatedAt.Before(s.Now().AddDate(0, 0, -agentPurchaseHistoryDays))) {
return fail(CodeTaskNotFound, "采购任务不存在")
}
var replay models.PurchaseTaskAttempt
if err := tx.Where("attempt_id = ?", attemptID).First(&replay).Error; err == nil {
if replay.TaskID != task.ID {
return fail(CodeResultConflict, "requestId 已用于其他采购任务")
}
response.AttemptNumber = replay.AttemptNumber
response.Status = models.PurchaseTaskStatusPending
response.Replayed = true
return nil
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
return internal(err)
}
if err := validatePurchaseResetState(tx, task); err != nil {
return err
}
if task.SpecSource != "unresolved" && !hasCompletePurchaseSpec(task) {
return fail(CodeMappingRequired, "任务没有完整的精确商品规格,请创建新采购任务")
}
deviceRecord, err := lockPurchaseResetDevice(tx, task, authenticated)
if err != nil {
return err
}
_, rawRule, rule, err := purchaserule.CurrentRule(ctx, tx, models.PurchaseExecutionModeLive)
if err != nil {
return err
}
required := purchasecontract.RequiredCapabilities(rule)
if err := ensureCapabilities(deviceRecord, required); err != nil {
return err
}
if err := ensureDeviceFree(tx, deviceRecord.ID, task.ID, s.Now()); err != nil {
return err
}
requiredJSON, err := json.Marshal(required)
if err != nil {
return internal(err)
}
var attemptCount int64
if err := tx.Model(&models.PurchaseTaskAttempt{}).Where("task_id = ?", task.ID).Count(&attemptCount).Error; err != nil {
return internal(err)
}
now := s.Now()
if err := task.SetStatus(models.PurchaseTaskStatusPending); err != nil {
return internal(err)
}
task.RuleType = rule.RuleType
task.RuleSchemaVersion = rule.SchemaVersion
task.RequiredCapabilitiesJSON = string(requiredJSON)
task.RuleSnapshot = string(rawRule)
task.LeaseExpiresAt = nil
task.ClaimRequestID = nil
task.ErrorCode = nil
task.ErrorMessage = nil
task.StatusVersion++
task.StatusChangedAt = now
if err := tx.Save(&task).Error; err != nil {
return conflictOrInternal(err)
}
phase := purchaseAttemptPhase(task)
// MySQL normalizes values written to a JSON column. Reload the task before
// hashing so the pending attempt uses the exact representation Start will
// read later, rather than the pre-persistence DefaultLiveRule bytes.
if err := tx.First(&task, task.ID).Error; err != nil {
return internal(err)
}
attempt := models.PurchaseTaskAttempt{
TaskID: task.ID, AttemptID: attemptID, AttemptNumber: int(attemptCount) + 1,
Phase: phase, Status: models.PurchaseAttemptStatusPending, DeviceID: &deviceRecord.ID,
RuleSnapshotHash: purchaseRuleSnapshotHash(task.RuleSnapshot), SpecDecisionSnapshot: task.SpecDecisionSnapshot,
}
if err := tx.Omit("Task").Create(&attempt).Error; err != nil {
return conflictOrInternal(err)
}
response.AttemptNumber = attempt.AttemptNumber
response.Status = task.Status
return nil
})
return response, err
}
func purchaseAttemptPhase(task models.PurchaseTask) string {
if task.SpecSource == "unresolved" || !hasCompletePurchaseSpec(task) {
return models.PurchaseAttemptPhaseSpecProbe
}
return models.PurchaseAttemptPhasePurchase
}
func hasCompletePurchaseSpec(task models.PurchaseTask) bool {
return (strings.TrimSpace(task.TargetColorSnapshot) == "" || strings.TrimSpace(task.MappedColorSnapshot) != "") &&
(strings.TrimSpace(task.TargetSizeSnapshot) == "" || strings.TrimSpace(task.MappedSizeSnapshot) != "")
}
func validatePurchaseResetState(tx *gorm.DB, task models.PurchaseTask) error {
switch task.Status {
case models.PurchaseTaskStatusOrderSubmitStarted, models.PurchaseTaskStatusOrderCreated, models.PurchaseTaskStatusOrderResultUnknown:
return fail(CodeRetryUnsafe, "任务可能已经创建订单,请走“授权重新采购”流程")
}
if task.Status != models.PurchaseTaskStatusFailed {
return fail(CodeRetryNotAllowed, "只有采购失败任务可以重试")
}
if task.TaskType != models.PurchaseTaskTypeSYBOrder || task.ExecutionMode != models.PurchaseExecutionModeLive || task.SYBProductID == nil {
return fail(CodeRetryNotAllowed, "只有正式采购的失败任务可以就地重试")
}
if task.IrreversibleAt != nil || task.OrderSubmitRequestID != nil || task.PDDOrderNo != nil || task.OrderSubmittedAt != nil {
return fail(CodeRetryUnsafe, "任务可能已经创建订单,请走“授权重新采购”流程")
}
var latest models.PurchaseTask
if err := tx.Where("syb_product_id = ?", *task.SYBProductID).Order("id DESC").First(&latest).Error; err != nil {
return internal(err)
}
if latest.ID != task.ID {
return fail(CodeRetryStale, fmt.Sprintf("同一 SYB 商品已有更新任务 %s", taskNumber(latest.ID)))
}
return nil
}
func lockPurchaseResetDevice(tx *gorm.DB, task models.PurchaseTask, authenticated *models.AgentDevice) (models.AgentDevice, error) {
if task.DeviceID == nil {
return models.AgentDevice{}, fail(CodeInvalidRequest, "原任务未分配设备,不能就地重试")
}
var record models.AgentDevice
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&record, *task.DeviceID).Error; err != nil {
return record, internal(err)
}
if authenticated != nil && record.ID != authenticated.ID {
return record, fail(CodeTaskNotFound, "采购任务不存在")
}
if record.Status != models.DeviceStatusOnline || record.TokenRevokedAt != nil {
return record, fail(CodeInvalidRequest, "原设备当前离线,不能就地重试")
}
return record, nil
}
func purchaseResetAttemptID(requestID string, taskID uint64) string {
return uuid.NewSHA1(uuid.NameSpaceOID, []byte(fmt.Sprintf("purchase-reset:%s:%d", requestID, taskID))).String()
}