490 lines
16 KiB
Go
490 lines
16 KiB
Go
package sybimport
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"go-admin/app/goauto/aimatching"
|
|
"go-admin/app/goauto/models"
|
|
"go-admin/app/goauto/shopeeproduct"
|
|
|
|
"github.com/google/uuid"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
const (
|
|
SpecAIParseInvokeTarget = "GoAutoSYBSpecAIParse"
|
|
defaultAIParseBatchLimit = 20
|
|
aiParseLeaseDuration = 30 * time.Minute
|
|
aiParseRetryDelay = time.Hour
|
|
maxAIParseAttempts = 3
|
|
)
|
|
|
|
var (
|
|
errAIParseWorkNotClaimed = errors.New("syb spec ai parse work not claimed")
|
|
errAIParseInputChanged = errors.New("syb spec ai parse input changed")
|
|
)
|
|
|
|
type aiParseInput struct {
|
|
ProductSpec string
|
|
Colors []string
|
|
Sizes []string
|
|
Fingerprint string
|
|
}
|
|
|
|
func StartSpecAIParseRun(ctx context.Context, db *gorm.DB, requestID string, batchLimit int) (models.SYBSpecAIParseRun, bool, error) {
|
|
if _, err := uuid.Parse(strings.TrimSpace(requestID)); err != nil {
|
|
return models.SYBSpecAIParseRun{}, false, fmt.Errorf("requestId 必须是 UUID")
|
|
}
|
|
if batchLimit <= 0 {
|
|
batchLimit = defaultAIParseBatchLimit
|
|
}
|
|
if batchLimit > 100 {
|
|
return models.SYBSpecAIParseRun{}, false, fmt.Errorf("batchLimit 不能超过 100")
|
|
}
|
|
now := time.Now().UTC()
|
|
lease := now.Add(aiParseLeaseDuration)
|
|
one := uint8(1)
|
|
owner := uuid.NewString()
|
|
var run models.SYBSpecAIParseRun
|
|
created := false
|
|
err := db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
if err := tx.Model(&models.SYBSpecAIParseRun{}).
|
|
Where("status = ? AND active_slot = ? AND lease_expires_at < ?", "running", 1, now).
|
|
Updates(map[string]any{"status": "failed", "active_slot": nil, "lease_owner": "", "lease_expires_at": nil, "error_summary": "上次运行租约过期,已安全释放", "finished_at": now}).Error; err != nil {
|
|
return err
|
|
}
|
|
if err := tx.Where("request_id = ?", requestID).First(&run).Error; err == nil {
|
|
return nil
|
|
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return err
|
|
}
|
|
if err := tx.Where("status = ? AND active_slot = ?", "running", 1).First(&run).Error; err == nil {
|
|
return nil
|
|
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return err
|
|
}
|
|
run = models.SYBSpecAIParseRun{
|
|
RequestID: requestID, Trigger: "scheduled", Status: "running", ActiveSlot: &one,
|
|
LeaseOwner: owner, LeaseExpiresAt: &lease, BatchLimit: batchLimit, StartedAt: now,
|
|
}
|
|
if err := tx.Create(&run).Error; err != nil {
|
|
return err
|
|
}
|
|
created = true
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
if findErr := db.WithContext(ctx).Where("status = ? AND active_slot = ?", "running", 1).First(&run).Error; findErr == nil {
|
|
return run, false, nil
|
|
}
|
|
return models.SYBSpecAIParseRun{}, false, err
|
|
}
|
|
return run, created, nil
|
|
}
|
|
|
|
func ProcessSpecAIParseRun(ctx context.Context, db *gorm.DB, runID uint64) error {
|
|
var run models.SYBSpecAIParseRun
|
|
if err := db.WithContext(ctx).First(&run, runID).Error; err != nil {
|
|
return err
|
|
}
|
|
if run.Status != "running" || run.ActiveSlot == nil || *run.ActiveSlot != 1 {
|
|
return nil
|
|
}
|
|
limit := run.BatchLimit
|
|
if limit <= 0 || limit > 100 {
|
|
limit = defaultAIParseBatchLimit
|
|
}
|
|
queryLimit := limit * 25
|
|
if queryLimit < 100 {
|
|
queryLimit = 100
|
|
}
|
|
if queryLimit > 1000 {
|
|
queryLimit = 1000
|
|
}
|
|
var candidates []models.SYBProduct
|
|
if err := db.WithContext(ctx).
|
|
Where("manually_confirmed = ?", false).
|
|
Where("parse_status IN ?", []string{models.SYBParseStatusUncertain, models.SYBParseStatusFailed}).
|
|
Order("updated_at ASC, id ASC").Limit(queryLimit).Find(&candidates).Error; err != nil {
|
|
finishSpecAIParseRun(db, run, "failed", 0, 0, 0, 0, 0, 1, "扫描异常规格失败")
|
|
return err
|
|
}
|
|
|
|
eligible, processed, confirmed, unmatched, failed := 0, 0, 0, 0, 0
|
|
firstError := ""
|
|
for _, candidate := range candidates {
|
|
if processed >= limit {
|
|
break
|
|
}
|
|
if candidate.AIConfirmed {
|
|
current, currentErr := aiConfirmationTargetsCurrent(ctx, db, candidate)
|
|
if currentErr != nil {
|
|
failed++
|
|
continue
|
|
}
|
|
if current {
|
|
continue
|
|
}
|
|
if err := db.WithContext(ctx).Model(&models.SYBProduct{}).
|
|
Where("id = ? AND manually_confirmed = ?", candidate.ID, false).
|
|
Updates(map[string]any{"ai_confirmed": false, "ai_confidence": nil, "ai_reason": "", "ai_confirmed_at": nil, "ai_input_fingerprint": ""}).Error; err != nil {
|
|
failed++
|
|
continue
|
|
}
|
|
candidate.AIConfirmed = false
|
|
}
|
|
outcome, err := Reparse(ctx, db, candidate.ID, false)
|
|
if err != nil {
|
|
failed++
|
|
if firstError == "" {
|
|
firstError = "确定性重新解析失败"
|
|
}
|
|
continue
|
|
}
|
|
if outcome.NewStatus == models.SYBParseStatusSuccess {
|
|
processed++
|
|
confirmed++
|
|
continue
|
|
}
|
|
if err := db.WithContext(ctx).First(&candidate, candidate.ID).Error; err != nil {
|
|
failed++
|
|
continue
|
|
}
|
|
input, ok, err := buildAIParseInput(ctx, db, candidate)
|
|
if err != nil {
|
|
failed++
|
|
if firstError == "" {
|
|
firstError = "读取 AI 解析上下文失败"
|
|
}
|
|
continue
|
|
}
|
|
if !ok {
|
|
continue
|
|
}
|
|
eligible++
|
|
work, claimed, err := claimSpecAIParseWork(ctx, db, run, candidate.ID, input.Fingerprint)
|
|
if err != nil {
|
|
failed++
|
|
continue
|
|
}
|
|
if !claimed {
|
|
continue
|
|
}
|
|
processed++
|
|
renewSpecAIParseRun(db, run)
|
|
matcher := aimatching.NewService(db)
|
|
result, matchErr := matcher.ResolveSYBSpec(ctx, aimatching.SYBSpecParseRequest{
|
|
ProductSpec: input.ProductSpec, Colors: input.Colors, Sizes: input.Sizes,
|
|
})
|
|
if matchErr == nil {
|
|
settings, settingsErr := matcher.Settings(ctx)
|
|
if settingsErr != nil {
|
|
matchErr = settingsErr
|
|
} else if result.Confidence == nil || *result.Confidence < settings.AutoConfirmMinConfidence || strings.TrimSpace(result.Reason) == "" {
|
|
unmatched++
|
|
completeSpecAIParseWork(db, work, input.Fingerprint, false, nil)
|
|
continue
|
|
} else {
|
|
matchErr = applyAIParseResult(ctx, db, candidate.ID, input.Fingerprint, result)
|
|
}
|
|
}
|
|
if matchErr != nil {
|
|
if isNoAIParseMatch(matchErr) || errors.Is(matchErr, errAIParseInputChanged) {
|
|
unmatched++
|
|
completeSpecAIParseWork(db, work, input.Fingerprint, false, nil)
|
|
continue
|
|
}
|
|
failed++
|
|
if firstError == "" {
|
|
firstError = safeAIParseError(matchErr)
|
|
}
|
|
completeSpecAIParseWork(db, work, input.Fingerprint, false, matchErr)
|
|
continue
|
|
}
|
|
confirmed++
|
|
completeSpecAIParseWork(db, work, input.Fingerprint, true, nil)
|
|
}
|
|
status := "completed"
|
|
if failed > 0 {
|
|
status = "completed_partial"
|
|
}
|
|
return finishSpecAIParseRun(db, run, status, len(candidates), eligible, processed, confirmed, unmatched, failed, firstError)
|
|
}
|
|
|
|
func buildAIParseInput(ctx context.Context, db *gorm.DB, record models.SYBProduct) (aiParseInput, bool, error) {
|
|
if record.ManuallyConfirmed || (record.ParseStatus != models.SYBParseStatusUncertain && record.ParseStatus != models.SYBParseStatusFailed) || record.ShopeeProductID == nil {
|
|
return aiParseInput{}, false, nil
|
|
}
|
|
var raw rawDetailSpec
|
|
if err := json.Unmarshal([]byte(record.RawJSON), &raw); err != nil {
|
|
return aiParseInput{}, false, err
|
|
}
|
|
raw.ProductSpec = strings.TrimSpace(raw.ProductSpec)
|
|
if raw.ProductSpec == "" {
|
|
return aiParseInput{}, false, nil
|
|
}
|
|
var product models.ShopeeProduct
|
|
if err := db.WithContext(ctx).First(&product, *record.ShopeeProductID).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return aiParseInput{}, false, nil
|
|
}
|
|
return aiParseInput{}, false, err
|
|
}
|
|
specs, err := shopeeproduct.Unmarshal(product.SpecsJSON)
|
|
if err != nil {
|
|
return aiParseInput{}, false, err
|
|
}
|
|
colors, sizes, ambiguous := closedShopeeCandidates(specs)
|
|
if ambiguous || (len(colors) == 0 && len(sizes) == 0) {
|
|
return aiParseInput{}, false, nil
|
|
}
|
|
var setting models.AIMatchingSetting
|
|
if err := db.WithContext(ctx).First(&setting, 1).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return aiParseInput{}, false, nil
|
|
}
|
|
return aiParseInput{}, false, err
|
|
}
|
|
if !setting.Enabled || strings.TrimSpace(setting.APIKey) == "" {
|
|
return aiParseInput{}, false, nil
|
|
}
|
|
fingerprintPayload := struct {
|
|
ProductSpec string
|
|
ShopeeProductID uint64
|
|
ShopeeSpecsJSON string
|
|
SettingUpdatedAt string
|
|
}{raw.ProductSpec, product.ID, product.SpecsJSON, setting.UpdatedAt.UTC().Format(time.RFC3339Nano)}
|
|
encoded, _ := json.Marshal(fingerprintPayload)
|
|
hash := sha256.Sum256(encoded)
|
|
return aiParseInput{ProductSpec: raw.ProductSpec, Colors: colors, Sizes: sizes, Fingerprint: hex.EncodeToString(hash[:])}, true, nil
|
|
}
|
|
|
|
func aiConfirmationTargetsCurrent(ctx context.Context, db *gorm.DB, record models.SYBProduct) (bool, error) {
|
|
if !record.AIConfirmed || record.ShopeeProductID == nil {
|
|
return false, nil
|
|
}
|
|
var product models.ShopeeProduct
|
|
if err := db.WithContext(ctx).First(&product, *record.ShopeeProductID).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return false, nil
|
|
}
|
|
return false, err
|
|
}
|
|
specs, err := shopeeproduct.Unmarshal(product.SpecsJSON)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
colors, sizes, ambiguous := closedShopeeCandidates(specs)
|
|
if ambiguous {
|
|
return false, nil
|
|
}
|
|
return closedCandidateContains(record.TargetColor, colors) && closedCandidateContains(record.TargetSize, sizes) &&
|
|
(strings.TrimSpace(record.TargetColor) != "" || strings.TrimSpace(record.TargetSize) != ""), nil
|
|
}
|
|
|
|
func closedCandidateContains(value string, candidates []string) bool {
|
|
value = strings.TrimSpace(value)
|
|
if len(candidates) == 0 {
|
|
return value == ""
|
|
}
|
|
for _, candidate := range candidates {
|
|
if candidate == value {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
func closedShopeeCandidates(specs []shopeeproduct.SpecDimension) (colors, sizes []string, ambiguous bool) {
|
|
roleDimensions := map[string]int{}
|
|
for _, dimension := range specs {
|
|
if dimension.Role != shopeeproduct.RoleColor && dimension.Role != shopeeproduct.RoleSize {
|
|
continue
|
|
}
|
|
values := make([]string, 0, len(dimension.Values))
|
|
seen := map[string]bool{}
|
|
for _, value := range dimension.Values {
|
|
name := strings.TrimSpace(value.Name)
|
|
if name != "" && !seen[name] {
|
|
seen[name] = true
|
|
values = append(values, name)
|
|
}
|
|
}
|
|
if len(values) == 0 {
|
|
continue
|
|
}
|
|
roleDimensions[dimension.Role]++
|
|
if roleDimensions[dimension.Role] > 1 {
|
|
return nil, nil, true
|
|
}
|
|
if dimension.Role == shopeeproduct.RoleColor {
|
|
colors = values
|
|
} else {
|
|
sizes = values
|
|
}
|
|
}
|
|
return colors, sizes, false
|
|
}
|
|
|
|
func applyAIParseResult(ctx context.Context, db *gorm.DB, id uint64, fingerprint string, result aimatching.SYBSpecParseResult) error {
|
|
now := time.Now().UTC()
|
|
return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
var record models.SYBProduct
|
|
if err := tx.First(&record, id).Error; err != nil {
|
|
return err
|
|
}
|
|
input, ok, err := buildAIParseInput(ctx, tx, record)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !ok || input.Fingerprint != fingerprint {
|
|
return errAIParseInputChanged
|
|
}
|
|
if result.Confidence == nil || strings.TrimSpace(result.Reason) == "" {
|
|
return errAIParseInputChanged
|
|
}
|
|
write := tx.Model(&models.SYBProduct{}).Where("id = ? AND manually_confirmed = ? AND ai_confirmed = ?", id, false, false).Updates(map[string]any{
|
|
"target_color": result.Color, "target_size": result.Size,
|
|
"ai_confirmed": true, "ai_confidence": *result.Confidence, "ai_reason": truncateAIParseText(result.Reason),
|
|
"ai_confirmed_at": now, "ai_input_fingerprint": fingerprint,
|
|
})
|
|
if write.Error != nil {
|
|
return write.Error
|
|
}
|
|
if write.RowsAffected != 1 {
|
|
return errAIParseInputChanged
|
|
}
|
|
return mergeParsedSpec(tx, *record.ShopeeProductID, ParseResult{Color: result.Color, Size: result.Size, Status: models.SYBParseStatusSuccess})
|
|
})
|
|
}
|
|
|
|
func claimSpecAIParseWork(ctx context.Context, db *gorm.DB, run models.SYBSpecAIParseRun, productID uint64, fingerprint string) (models.SYBSpecAIParseWorkItem, bool, error) {
|
|
now := time.Now().UTC()
|
|
lease := now.Add(aiParseLeaseDuration)
|
|
var work models.SYBSpecAIParseWorkItem
|
|
err := db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
err := tx.Where("syb_product_id = ?", productID).First(&work).Error
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
work = models.SYBSpecAIParseWorkItem{SYBProductID: productID, RunID: &run.ID, InputFingerprint: fingerprint, Status: "running", AttemptCount: 1, LeaseOwner: run.LeaseOwner, LeaseExpiresAt: &lease}
|
|
return tx.Create(&work).Error
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if work.InputFingerprint == fingerprint {
|
|
if work.Status == "completed" || work.Status == "unmatched" || work.AttemptCount >= maxAIParseAttempts || (work.NextAttemptAt != nil && work.NextAttemptAt.After(now)) || (work.Status == "running" && work.LeaseExpiresAt != nil && work.LeaseExpiresAt.After(now)) {
|
|
return errAIParseWorkNotClaimed
|
|
}
|
|
} else {
|
|
work.AttemptCount = 0
|
|
}
|
|
updates := map[string]any{"run_id": run.ID, "input_fingerprint": fingerprint, "status": "running", "attempt_count": work.AttemptCount + 1, "next_attempt_at": nil, "lease_owner": run.LeaseOwner, "lease_expires_at": lease, "last_error_code": "", "last_error": ""}
|
|
if err := tx.Model(&models.SYBSpecAIParseWorkItem{}).Where("id = ?", work.ID).Updates(updates).Error; err != nil {
|
|
return err
|
|
}
|
|
return tx.First(&work, work.ID).Error
|
|
})
|
|
if errors.Is(err, errAIParseWorkNotClaimed) {
|
|
return work, false, nil
|
|
}
|
|
return work, err == nil, err
|
|
}
|
|
|
|
func completeSpecAIParseWork(db *gorm.DB, work models.SYBSpecAIParseWorkItem, fingerprint string, confirmed bool, parseErr error) {
|
|
now := time.Now().UTC()
|
|
updates := map[string]any{"input_fingerprint": fingerprint, "lease_owner": "", "lease_expires_at": nil}
|
|
if parseErr == nil {
|
|
if confirmed {
|
|
updates["status"] = "completed"
|
|
} else {
|
|
updates["status"] = "unmatched"
|
|
}
|
|
updates["next_attempt_at"], updates["last_error_code"], updates["last_error"] = nil, "", ""
|
|
} else {
|
|
code := aiParseErrorCode(parseErr)
|
|
updates["status"], updates["last_error_code"], updates["last_error"] = "failed", code, safeAIParseError(parseErr)
|
|
if code == aimatching.CodeProviderUnavailable && work.AttemptCount < maxAIParseAttempts {
|
|
next := now.Add(aiParseRetryDelay)
|
|
updates["next_attempt_at"] = next
|
|
} else {
|
|
updates["next_attempt_at"] = nil
|
|
}
|
|
}
|
|
_ = db.Model(&models.SYBSpecAIParseWorkItem{}).Where("id = ?", work.ID).Updates(updates).Error
|
|
}
|
|
|
|
func renewSpecAIParseRun(db *gorm.DB, run models.SYBSpecAIParseRun) {
|
|
lease := time.Now().UTC().Add(aiParseLeaseDuration)
|
|
_ = db.Model(&models.SYBSpecAIParseRun{}).Where("id = ? AND status = ? AND lease_owner = ?", run.ID, "running", run.LeaseOwner).Update("lease_expires_at", lease).Error
|
|
}
|
|
|
|
func finishSpecAIParseRun(db *gorm.DB, run models.SYBSpecAIParseRun, status string, scanned, eligible, processed, confirmed, unmatched, failed int, summary string) error {
|
|
now := time.Now().UTC()
|
|
return db.Model(&models.SYBSpecAIParseRun{}).Where("id = ? AND status = ? AND lease_owner = ?", run.ID, "running", run.LeaseOwner).Updates(map[string]any{
|
|
"status": status, "active_slot": nil, "lease_owner": "", "lease_expires_at": nil,
|
|
"scanned_count": scanned, "eligible_count": eligible, "processed_count": processed,
|
|
"confirmed_count": confirmed, "unmatched_count": unmatched, "failed_count": failed,
|
|
"error_summary": truncateAIParseText(summary), "finished_at": now,
|
|
}).Error
|
|
}
|
|
|
|
func isNoAIParseMatch(err error) bool { return aiParseErrorCode(err) == aimatching.CodeNoMatch }
|
|
|
|
func aiParseErrorCode(err error) string {
|
|
var target *aimatching.Error
|
|
if errors.As(err, &target) {
|
|
return target.Code
|
|
}
|
|
return CodeInternal
|
|
}
|
|
|
|
func safeAIParseError(err error) string {
|
|
var target *aimatching.Error
|
|
if errors.As(err, &target) {
|
|
return truncateAIParseText(target.Message)
|
|
}
|
|
return "服务端处理失败"
|
|
}
|
|
|
|
func truncateAIParseText(value string) string {
|
|
runes := []rune(strings.TrimSpace(value))
|
|
if len(runes) > 500 {
|
|
runes = runes[:500]
|
|
}
|
|
return string(runes)
|
|
}
|
|
|
|
type scheduledSpecAIParseArgs struct {
|
|
BatchLimit int `json:"batchLimit"`
|
|
}
|
|
|
|
type ScheduledSpecAIParseJob struct{}
|
|
|
|
func (ScheduledSpecAIParseJob) Exec(_ interface{}) error {
|
|
return errors.New("SYB 规格 AI 解析定时任务缺少数据库连接")
|
|
}
|
|
|
|
func (ScheduledSpecAIParseJob) ExecWithDB(db *gorm.DB, arg interface{}) error {
|
|
args := scheduledSpecAIParseArgs{BatchLimit: defaultAIParseBatchLimit}
|
|
if raw, ok := arg.(string); ok && strings.TrimSpace(raw) != "" {
|
|
if err := json.Unmarshal([]byte(raw), &args); err != nil {
|
|
return fmt.Errorf("SYB 规格 AI 解析参数不是合法 JSON: %w", err)
|
|
}
|
|
}
|
|
if args.BatchLimit < 1 || args.BatchLimit > 100 {
|
|
return errors.New("SYB 规格 AI 解析 batchLimit 必须在 1 到 100 之间")
|
|
}
|
|
run, created, err := StartSpecAIParseRun(context.Background(), db, uuid.NewString(), args.BatchLimit)
|
|
if err != nil || !created {
|
|
return err
|
|
}
|
|
return ProcessSpecAIParseRun(context.Background(), db, run.ID)
|
|
}
|