diff --git a/server/app/goauto/sybinnercode/match.go b/server/app/goauto/sybinnercode/match.go index 544b6f5..9665fba 100644 --- a/server/app/goauto/sybinnercode/match.go +++ b/server/app/goauto/sybinnercode/match.go @@ -17,6 +17,7 @@ import ( "github.com/google/uuid" "gorm.io/gorm" + "gorm.io/gorm/clause" ) type MatchReader interface { @@ -45,12 +46,19 @@ func (m Matcher) runBackground(db *gorm.DB, jobID string) { err = RunMatchJob(ctx, db, reader, jobID) } if err != nil { - now := time.Now().UTC() - _ = db.Model(&models.SYBInnerCodeMatchJob{}).Where("id = ? AND status IN ?", jobID, []string{"pending", "running"}).Updates(map[string]any{"status": "failed", "error_message": compact(err.Error(), 1000), "finished_at": now}).Error + _ = failMatchJob(db, jobID, "匹配任务中断,请重新匹配") } } -func RunMatchJob(ctx context.Context, db *gorm.DB, reader MatchReader, jobID string) error { +func RunMatchJob(ctx context.Context, db *gorm.DB, reader MatchReader, jobID string) (runErr error) { + claimedJob := false + defer func() { + if runErr != nil && claimedJob { + if err := failMatchJob(db, jobID, "匹配任务中断,请重新匹配"); err != nil { + runErr = errors.Join(runErr, err) + } + } + }() now := time.Now().UTC() claimed := db.WithContext(ctx).Model(&models.SYBInnerCodeMatchJob{}).Where("id = ? AND status = ?", jobID, "pending").Updates(map[string]any{"status": "running", "started_at": now}) if claimed.Error != nil { @@ -66,6 +74,7 @@ func RunMatchJob(ctx context.Context, db *gorm.DB, reader MatchReader, jobID str } return conflict("匹配任务状态不允许执行") } + claimedJob = true var job models.SYBInnerCodeMatchJob if err := db.WithContext(ctx).First(&job, "id = ?", jobID).Error; err != nil { return err @@ -75,7 +84,7 @@ func RunMatchJob(ctx context.Context, db *gorm.DB, reader MatchReader, jobID str return fmt.Errorf("匹配任务记录范围无效") } var records []models.SYBInnerCodeRecord - if err := db.WithContext(ctx).Preload("Items", func(q *gorm.DB) *gorm.DB { return q.Order("ordinal ASC") }).Where("id IN ? AND business_date = ? AND status IN ?", recordIDs, job.BusinessDate, []string{models.SYBInnerCodePending, models.SYBInnerCodeFailed, models.SYBInnerCodeSkipped}).Order("source_row,id").Find(&records).Error; err != nil { + if err := db.WithContext(ctx).Preload("Items", func(q *gorm.DB) *gorm.DB { return q.Order("ordinal ASC") }).Where("id IN ? AND business_date = ? AND status IN ?", recordIDs, job.BusinessDate, []string{models.SYBInnerCodePending, models.SYBInnerCodeMatching, models.SYBInnerCodeFailed, models.SYBInnerCodeSkipped}).Order("source_row,id").Find(&records).Error; err != nil { return err } used, err := loadReservedDetails(ctx, db, records) @@ -84,7 +93,10 @@ func RunMatchJob(ctx context.Context, db *gorm.DB, reader MatchReader, jobID str } ready, failed := 0, 0 for _, record := range records { - plan, status, message, planErr := planRecord(ctx, reader, record, used) + if err := ctx.Err(); err != nil { + return err + } + plan, status, message, planErr := safePlanRecord(ctx, reader, record, used) if planErr != nil { status = models.SYBInnerCodeFailed message = "读取 SYB 失败:" + compact(planErr.Error(), 900) @@ -111,12 +123,62 @@ func RunMatchJob(ctx context.Context, db *gorm.DB, reader MatchReader, jobID str } else { failed++ } - db.Model(&models.SYBInnerCodeMatchJob{}).Where("id = ?", jobID).Updates(map[string]any{"processed": gorm.Expr("processed + 1"), "ready": ready, "failed": failed}) + if err := db.WithContext(ctx).Model(&models.SYBInnerCodeMatchJob{}).Where("id = ?", jobID).Updates(map[string]any{"processed": gorm.Expr("processed + 1"), "ready": ready, "failed": failed}).Error; err != nil { + return err + } } finished := time.Now().UTC() return db.WithContext(ctx).Model(&models.SYBInnerCodeMatchJob{}).Where("id = ? AND status = ?", jobID, "running").Updates(map[string]any{"status": "succeeded", "finished_at": finished, "ready": ready, "failed": failed}).Error } +// Only the read/plan step is isolated: no remote write is retried here. +func safePlanRecord(ctx context.Context, reader MatchReader, record models.SYBInnerCodeRecord, used map[int64]bool) (plan *models.SYBInnerCodePlan, status, message string, err error) { + defer func() { + if recover() != nil { + plan, status, message = nil, models.SYBInnerCodeFailed, "匹配处理异常,请重新匹配" + err = nil + } + }() + return planRecord(ctx, reader, record, used) +} + +// Cleanup must not inherit an expired job or HTTP request context. +func failMatchJob(db *gorm.DB, jobID, message string) error { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + var job models.SYBInnerCodeMatchJob + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&job, "id = ?", jobID).Error; err != nil { + return err + } + if job.Status != "pending" && job.Status != "running" { + return nil + } + var ids []uint64 + if err := json.Unmarshal([]byte(job.RecordIDsJSON), &ids); err != nil { + return err + } + if err := tx.Model(&models.SYBInnerCodeRecord{}).Where("id IN ? AND business_date = ? AND status IN ?", ids, job.BusinessDate, []string{models.SYBInnerCodePending, models.SYBInnerCodeMatching}).Updates(map[string]any{"status": models.SYBInnerCodeFailed, "result_message": message}).Error; err != nil { + return err + } + return tx.Model(&job).Updates(map[string]any{"status": "failed", "error_message": message, "finished_at": time.Now().UTC()}).Error + }) +} + +// Called at startup before new jobs can be submitted; never resumes remote writes. +func RecoverInterruptedMatches(db *gorm.DB) error { + var jobs []models.SYBInnerCodeMatchJob + if err := db.Where("status IN ?", []string{"pending", "running"}).Find(&jobs).Error; err != nil { + return err + } + for _, job := range jobs { + if err := failMatchJob(db, job.ID, "服务重启,匹配任务中断,请重新匹配"); err != nil { + return err + } + } + return nil +} + func loadReservedDetails(ctx context.Context, db *gorm.DB, selected []models.SYBInnerCodeRecord) (map[int64]bool, error) { selectedIDs := map[uint64]bool{} for _, r := range selected { @@ -186,6 +248,9 @@ func planRecord(ctx context.Context, reader MatchReader, record models.SYBInnerC if reason != "" { return nil, models.SYBInnerCodeSkipped, reason, nil } + if len(matches) == 0 { + return nil, models.SYBInnerCodeSkipped, "相同规格候选的原始 SKU 或档口货号未匹配", nil + } count := len(record.Items) if count == 0 { return nil, models.SYBInnerCodeSkipped, "记录没有入库码", nil @@ -213,6 +278,9 @@ func planRecord(ctx context.Context, reader MatchReader, record models.SYBInnerC } else { return nil, models.SYBInnerCodeSkipped, fmt.Sprintf("SYB 商品数量与入库码数量不一致(%d/%d)", matches[0].ProductQty, count), nil } + if len(chosen) == 0 { + return nil, models.SYBInnerCodeSkipped, "未形成唯一的商品分配,不能自动选择", nil + } primary := chosen[0] items := make([]plannedRemoteItem, 0, count) placeholder := 0 diff --git a/server/app/goauto/sybinnercode/match_recovery_test.go b/server/app/goauto/sybinnercode/match_recovery_test.go new file mode 100644 index 0000000..d7c2ae9 --- /dev/null +++ b/server/app/goauto/sybinnercode/match_recovery_test.go @@ -0,0 +1,181 @@ +package sybinnercode + +import ( + "context" + "errors" + "testing" + + "github.com/google/uuid" + "go-admin/app/goauto/models" + "go-admin/app/goauto/sybclient" +) + +type interruptingReader struct { + *fakeMatchReader + panicOrder string + cancel context.CancelFunc +} + +func (r interruptingReader) ListByOrderNumber(ctx context.Context, order string) ([]sybclient.StockRow, error) { + if order == r.panicOrder { + panic("must not escape or be persisted") + } + if r.cancel != nil { + r.cancel() + return nil, ctx.Err() + } + return r.fakeMatchReader.ListByOrderNumber(ctx, order) +} + +func TestEmptyEvidenceDoesNotPanicAndNextRecordContinues(t *testing.T) { + db := testDB(t) + job := createMatchJob(t, db, []models.SYBInnerCodeRecord{matchRecord("BAD", "NO-SKU", "NO#9", "IC-1"), matchRecord("GOOD", "SKU-A", "A#1", "IC-2")}) + r := &fakeMatchReader{rows: map[string][]sybclient.StockRow{"BAD": {{ID: 10}}, "GOOD": {{ID: 11}}}, stocks: map[int64]sybclient.StockDetail{10: {ID: 10, Details: []sybclient.DetailItem{detail(20, "黑色,L", 1, "SKU-A", "A#1", "")}}, 11: {ID: 11, Details: []sybclient.DetailItem{detail(21, "黑色,L", 1, "SKU-A", "A#1", "")}}}} + if err := RunMatchJob(context.Background(), db, r, job); err != nil { + t.Fatal(err) + } + var rows []models.SYBInnerCodeRecord + db.Order("id").Find(&rows) + if rows[0].Status != models.SYBInnerCodeSkipped || rows[1].Status != models.SYBInnerCodeReady { + t.Fatal("empty evidence must be limited, next row ready") + } + var plans int64 + db.Model(&models.SYBInnerCodePlan{}).Count(&plans) + if plans != 1 { + t.Fatalf("unexpected plans=%d", plans) + } + var j models.SYBInnerCodeMatchJob + db.First(&j, "id = ?", job) + if j.Processed != 2 || j.Ready != 1 || j.Failed != 1 || j.Status != "succeeded" { + t.Fatalf("counts=%+v", j) + } +} + +func TestSingleRecordPanicIsIsolated(t *testing.T) { + db := testDB(t) + job := createMatchJob(t, db, []models.SYBInnerCodeRecord{matchRecord("PANIC", "S", "A", "IC-1"), matchRecord("NEXT", "S", "A", "IC-2")}) + r := interruptingReader{fakeMatchReader: &fakeMatchReader{rows: map[string][]sybclient.StockRow{}}, panicOrder: "PANIC"} + if err := RunMatchJob(context.Background(), db, r, job); err != nil { + t.Fatal(err) + } + var rows []models.SYBInnerCodeRecord + db.Order("id").Find(&rows) + if rows[0].Status != models.SYBInnerCodeFailed || rows[0].ResultMessage != "匹配处理异常,请重新匹配" || rows[1].Status != models.SYBInnerCodeFailed { + t.Fatal("panic was not safely persisted or next row not processed") + } + var j models.SYBInnerCodeMatchJob + db.First(&j, "id = ?", job) + if j.Processed != 2 { + t.Fatal("remaining row was not processed") + } +} + +func TestCancelledMatchPersistsFailureWithIndependentContext(t *testing.T) { + db := testDB(t) + job := createMatchJob(t, db, []models.SYBInnerCodeRecord{matchRecord("A", "S", "A", "IC-1"), matchRecord("B", "S", "B", "IC-2")}) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + err := RunMatchJob(ctx, db, interruptingReader{fakeMatchReader: &fakeMatchReader{}, cancel: cancel}, job) + if !errors.Is(err, context.Canceled) { + t.Fatalf("err=%v", err) + } + var rows []models.SYBInnerCodeRecord + db.Find(&rows) + for _, r := range rows { + if r.Status != models.SYBInnerCodeFailed { + t.Fatal("cancel left pending row") + } + } + var j models.SYBInnerCodeMatchJob + db.First(&j, "id = ?", job) + if j.Status != "failed" || j.FinishedAt == nil { + t.Fatal("cancel left running job") + } +} + +func TestRestartRecoveryPreservesSuccessfulRecordsAndScope(t *testing.T) { + db := testDB(t) + job := createMatchJob(t, db, []models.SYBInnerCodeRecord{matchRecord("A", "S", "A", "IC-1"), matchRecord("B", "S", "B", "IC-2")}) + db.Model(&models.SYBInnerCodeMatchJob{}).Where("id = ?", job).Update("status", "running") + db.Model(&models.SYBInnerCodeRecord{}).Where("order_number = ?", "A").Update("status", models.SYBInnerCodeUpdated) + outside := matchRecord("OUTSIDE", "S", "C", "IC-3") + db.Create(&outside) + if err := RecoverInterruptedMatches(db); err != nil { + t.Fatal(err) + } + if err := RecoverInterruptedMatches(db); err != nil { + t.Fatal(err) + } + var rows []models.SYBInnerCodeRecord + db.Order("id").Find(&rows) + if rows[0].Status != models.SYBInnerCodeUpdated || rows[1].Status != models.SYBInnerCodeFailed || rows[2].Status != models.SYBInnerCodePending { + t.Fatal("recovery altered success or outside scope") + } + var j models.SYBInnerCodeMatchJob + db.First(&j, "id = ?", job) + if j.Status != "failed" || j.FinishedAt == nil { + t.Fatal("recovery left running job") + } +} + +func TestBatchRematchClaimsPendingFailedSkippedAndIsIdempotent(t *testing.T) { + db := testDB(t) + statuses := []string{models.SYBInnerCodePending, models.SYBInnerCodeFailed, models.SYBInnerCodeSkipped} + ids := []uint64{} + for i, status := range statuses { + r := matchRecord(string(rune('A'+i)), "S", "A", string(rune('X'+i))) + r.Status = status + if err := db.Create(&r).Error; err != nil { + t.Fatal(err) + } + ids = append(ids, r.ID) + } + s := NewService(db) + req := RematchRequest{RequestID: uuid.NewString(), IDs: ids} + result, err := s.QueueRematch(context.Background(), req) + if err != nil || result.Queued != 3 { + t.Fatalf("result=%+v err=%v", result, err) + } + replay, err := s.QueueRematch(context.Background(), req) + if err != nil || replay.MatchJobID != result.MatchJobID { + t.Fatal("idempotent replay failed") + } + if _, err := s.QueueRematch(context.Background(), RematchRequest{RequestID: uuid.NewString(), IDs: ids}); err == nil { + t.Fatal("overlapping job accepted") + } + var rows []models.SYBInnerCodeRecord + db.Find(&rows) + for _, r := range rows { + if r.Status != models.SYBInnerCodeMatching { + t.Fatal("not claimed") + } + } + if _, err := s.Delete(context.Background(), 1, DeleteRequest{RequestID: uuid.NewString(), IDs: ids}); err == nil { + t.Fatal("matching records can be deleted during execution") + } + if err := RunMatchJob(context.Background(), db, &fakeMatchReader{rows: map[string][]sybclient.StockRow{}}, result.MatchJobID); err != nil { + t.Fatal(err) + } + var job models.SYBInnerCodeMatchJob + db.First(&job, "id = ?", result.MatchJobID) + if job.Processed != 3 || job.Status != "succeeded" { + t.Fatal("claimed matching records were not executed") + } +} + +func TestRematchRejectsPendingOwnedByImportAndProtectedStates(t *testing.T) { + db := testDB(t) + createMatchJob(t, db, []models.SYBInnerCodeRecord{matchRecord("A", "S", "A", "X")}) + var r models.SYBInnerCodeRecord + db.First(&r) + s := NewService(db) + if _, err := s.QueueRematch(context.Background(), RematchRequest{RequestID: uuid.NewString(), IDs: []uint64{r.ID}}); err == nil { + t.Fatal("active import overlapped") + } + for _, status := range []string{models.SYBInnerCodeReady, models.SYBInnerCodeUpdated, models.SYBInnerCodeAlreadyFilled, models.SYBInnerCodeQueued, models.SYBInnerCodeApplying, models.SYBInnerCodeNeedsCheck} { + db.Model(&r).Update("status", status) + if _, err := s.QueueRematch(context.Background(), RematchRequest{RequestID: uuid.NewString(), IDs: []uint64{r.ID}}); err == nil { + t.Fatalf("protected status accepted: %s", status) + } + } +} diff --git a/server/app/goauto/sybinnercode/service.go b/server/app/goauto/sybinnercode/service.go index dc7d78d..98c97ab 100644 --- a/server/app/goauto/sybinnercode/service.go +++ b/server/app/goauto/sybinnercode/service.go @@ -200,8 +200,8 @@ func (s *Service) QueueRematch(ctx context.Context, request RematchRequest) (Rem return conflict("部分记录不存在") } for _, record := range records { - if record.Status != models.SYBInnerCodeFailed && record.Status != models.SYBInnerCodeSkipped { - return conflict("只有读取失败或匹配受限记录可以重新匹配") + if record.Status != models.SYBInnerCodePending && record.Status != models.SYBInnerCodeFailed && record.Status != models.SYBInnerCodeSkipped { + return conflict("只有待匹配、读取失败或匹配受限记录可以重新匹配") } if date == "" { date = record.BusinessDate @@ -209,10 +209,29 @@ func (s *Service) QueueRematch(ctx context.Context, request RematchRequest) (Rem return conflict("重新匹配记录必须属于同一营业日期") } } + var activeJobs []models.SYBInnerCodeMatchJob + if err := tx.Where("business_date = ? AND status IN ?", date, []string{"pending", "running"}).Find(&activeJobs).Error; err != nil { + return err + } + selected := make(map[uint64]bool, len(ids)) + for _, id := range ids { + selected[id] = true + } + for _, job := range activeJobs { + var jobIDs []uint64 + if err := json.Unmarshal([]byte(job.RecordIDsJSON), &jobIDs); err != nil { + return err + } + for _, id := range jobIDs { + if selected[id] { + 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 { + if err := tx.Model(&models.SYBInnerCodeRecord{}).Where("id IN ?", ids).Updates(map[string]any{"status": models.SYBInnerCodeMatching, "result_message": "等待重新匹配"}).Error; err != nil { return err } recordIDsJSON, _ := json.Marshal(ids) @@ -267,13 +286,13 @@ func (s *Service) Delete(ctx context.Context, actor uint64, request DeleteReques return conflict("部分记录不存在,未删除任何数据") } for _, record := range records { - if record.Status == models.SYBInnerCodeQueued || record.Status == models.SYBInnerCodeApplying || record.Status == models.SYBInnerCodeNeedsCheck { + if record.Status == models.SYBInnerCodeMatching || 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}} + 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 diff --git a/server/cmd/api/server.go b/server/cmd/api/server.go index e7c93b7..e435dae 100644 --- a/server/cmd/api/server.go +++ b/server/cmd/api/server.go @@ -117,6 +117,9 @@ func run() error { if err := goautosybinnercode.RecoverInterrupted(db); err != nil { return fmt.Errorf("recover interrupted SYB inner-code writes: %w", err) } + if err := goautosybinnercode.RecoverInterruptedMatches(db); err != nil { + return fmt.Errorf("recover interrupted SYB inner-code matches: %w", err) + } goautoreplacement.RecoverMatching(db) goautopurchase.RecoverPurchaseMatching(db) goautopurchase.RecoverOrderWritebacks(db) diff --git a/web/src/views/goauto/syb-inner-codes/index.vue b/web/src/views/goauto/syb-inner-codes/index.vue index a400603..8717c18 100644 --- a/web/src/views/goauto/syb-inner-codes/index.vue +++ b/web/src/views/goauto/syb-inner-codes/index.vue @@ -5,14 +5,15 @@
导入 Excel 后自动匹配 SYB 商品;确认后逐件回写,结果不明确时只读复核。