diff --git a/server/app/goauto/access/modules.go b/server/app/goauto/access/modules.go index 1adf63a..cf7793d 100644 --- a/server/app/goauto/access/modules.go +++ b/server/app/goauto/access/modules.go @@ -136,6 +136,8 @@ func moduleKeyForAPI(path string) string { return ModulePDDProducts case strings.HasPrefix(path, "/api/admin/v1/shopee-products"): return ModuleShopeeProducts + case strings.HasPrefix(path, "/api/admin/v1/shopee-spec-auto-match"): + return ModuleShopeeProducts case strings.HasPrefix(path, "/api/admin/v1/syb-products/sync-runs"): return ModuleSYBSyncRuns case strings.HasPrefix(path, "/api/admin/v1/syb-products"): diff --git a/server/app/goauto/access/purchaser.go b/server/app/goauto/access/purchaser.go index 958e12a..264a291 100644 --- a/server/app/goauto/access/purchaser.go +++ b/server/app/goauto/access/purchaser.go @@ -48,6 +48,8 @@ var AdminAPIs = []APIPermission{ {"AI 建议颜色映射", "/api/admin/v1/shopee-products/:productId/specs/mapping/suggest-colors", "POST", true}, {"AI 建议尺码映射", "/api/admin/v1/shopee-products/:productId/specs/mapping/suggest-sizes", "POST", true}, {"一键匹配并确认颜色尺码", "/api/admin/v1/shopee-products/:productId/specs/mapping/auto-match", "POST", true}, + {"手动执行虾皮规格自动匹配", "/api/admin/v1/shopee-spec-auto-match/runs", "POST", false}, + {"查看最近虾皮规格自动匹配", "/api/admin/v1/shopee-spec-auto-match/runs/latest", "GET", false}, {"查看 SYB 商品", "/api/admin/v1/syb-products", "GET", true}, {"查看 SYB 商品详情", "/api/admin/v1/syb-products/:productId", "GET", true}, diff --git a/server/app/goauto/access/purchaser_test.go b/server/app/goauto/access/purchaser_test.go index 7cd208f..89ce9aa 100644 --- a/server/app/goauto/access/purchaser_test.go +++ b/server/app/goauto/access/purchaser_test.go @@ -15,10 +15,12 @@ func TestPurchaserPermissionMatrixHasNoDuplicates(t *testing.T) { func TestPurchaserExcludesAdministratorOperations(t *testing.T) { denied := map[string]bool{ - "POST /api/admin/v1/devices/:deviceId/disable": true, - "POST /api/admin/v1/syb-products/import": true, - "POST /api/admin/v1/collection-rules": true, - "PUT /api/admin/v1/ai-matching-settings": true, + "POST /api/admin/v1/devices/:deviceId/disable": true, + "POST /api/admin/v1/syb-products/import": true, + "POST /api/admin/v1/collection-rules": true, + "PUT /api/admin/v1/ai-matching-settings": true, + "POST /api/admin/v1/shopee-spec-auto-match/runs": true, + "GET /api/admin/v1/shopee-spec-auto-match/runs/latest": true, } for _, permission := range PurchaserAPIs() { if denied[permission.Method+" "+permission.Path] { diff --git a/server/app/goauto/aimatching/service.go b/server/app/goauto/aimatching/service.go index 4255494..dccc430 100644 --- a/server/app/goauto/aimatching/service.go +++ b/server/app/goauto/aimatching/service.go @@ -25,12 +25,28 @@ const ( defaultAutoConfirmMinConfidence = 0.9 ) +// MaxProviderTimeout is also the total budget used by composite synchronous +// AI operations. This keeps their HTTP response inside the Admin and API +// transport windows even when an operation needs more than one provider call. +const MaxProviderTimeout = 600 * time.Second + // Service owns the internal AI Provider configuration. The API key exception // is deliberately narrow: it is plain text only in the dedicated settings // table and is returned only by the administrator settings handler. type Service struct { - DB *gorm.DB - HTTPClient *http.Client + DB *gorm.DB + HTTPClient *http.Client + ProviderFailureLogger func(ProviderFailureDiagnostic) +} + +// ProviderFailureDiagnostic deliberately contains no URL, model, prompt, +// candidates, response body or credential. It is safe for operational logs. +type ProviderFailureDiagnostic struct { + CallID string + Operation string + Kind string + StatusCode int + Duration time.Duration } func NewService(db *gorm.DB) *Service { @@ -179,6 +195,58 @@ func (s *Service) Resolve(ctx context.Context, request MatchRequest) (MatchResul return result, nil } +// ResolveSYBSpec parses one SYB productSpec into the linked Shopee product's +// exact color/size labels. Unlike Resolve, this operation does not map to PDD: +// every non-empty answer must be an exact member of the supplied Shopee set. +func (s *Service) ResolveSYBSpec(ctx context.Context, request SYBSpecParseRequest) (SYBSpecParseResult, error) { + request.ProductSpec = strings.TrimSpace(request.ProductSpec) + request.Colors = usableCandidates(request.Colors) + request.Sizes = usableCandidates(request.Sizes) + if request.ProductSpec == "" || (len(request.Colors) == 0 && len(request.Sizes) == 0) { + return SYBSpecParseResult{}, fail(CodeNoMatch, "SYB 采购规格缺少可判断的原文或蝦皮候选") + } + setting, apiKey, err := s.activeSetting(ctx) + if err != nil { + return SYBSpecParseResult{}, err + } + payload := openAIChatRequest{Model: setting.Model, Temperature: 0, Messages: []openAIMessage{ + {Role: "system", Content: "你只负责把一条 SYB 商品规格原文解析成给定蝦皮候选中的原始颜色和尺码。不得猜测、不得改写候选、不得返回候选外文本。只返回 JSON:{\"color\":\"颜色候选原文或空\",\"size\":\"尺码候选原文或空\",\"reason\":\"简短原因\",\"confidence\":0到1}。提供了某角色候选时必须唯一可靠地选择一个,否则对应字段留空。"}, + {Role: "user", Content: sybSpecParsePrompt(request)}, + }} + body, err := json.Marshal(payload) + if err != nil { + return SYBSpecParseResult{}, &Error{Code: CodeProviderUnavailable, Message: "SYB 规格 AI 解析请求生成失败", Cause: err} + } + ctx, cancel := context.WithTimeout(ctx, time.Duration(setting.TimeoutSeconds)*time.Second) + defer cancel() + httpRequest, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint(setting.BaseURL, "chat/completions"), bytes.NewReader(body)) + if err != nil { + return SYBSpecParseResult{}, fail(CodeInvalidSetting, "AI 服务地址无效") + } + httpRequest.Header.Set("Authorization", "Bearer "+apiKey) + httpRequest.Header.Set("Content-Type", "application/json") + response, err := s.httpClient().Do(httpRequest) + if err != nil { + return SYBSpecParseResult{}, &Error{Code: CodeProviderUnavailable, Message: "SYB 规格 AI 解析服务暂时不可用", Cause: err} + } + defer response.Body.Close() + responseBody, readErr := io.ReadAll(io.LimitReader(response.Body, 1<<20)) + if readErr != nil || response.StatusCode < http.StatusOK || response.StatusCode >= http.StatusMultipleChoices { + return SYBSpecParseResult{}, fail(CodeProviderUnavailable, "SYB 规格 AI 解析服务暂时不可用") + } + choice, err := parseProviderChoice(responseBody) + if err != nil || !validClosedChoice(choice.Color, request.Colors) || !validClosedChoice(choice.Size, request.Sizes) { + return SYBSpecParseResult{}, fail(CodeNoMatch, "AI 未能在蝦皮候选中唯一解析采购规格") + } + if choice.Confidence == nil || *choice.Confidence < 0 || *choice.Confidence > 1 || strings.TrimSpace(choice.Reason) == "" { + return SYBSpecParseResult{}, fail(CodeNoMatch, "AI 解析结果缺少有效置信度或理由") + } + return SYBSpecParseResult{ + Color: choice.Color, Size: choice.Size, Provider: ProviderOpenAICompatible, + Model: setting.Model, Reason: safeReason(choice.Reason), Confidence: choice.Confidence, + }, nil +} + func (s *Service) activeSetting(ctx context.Context) (models.AIMatchingSetting, string, error) { setting, err := s.setting(ctx) if errors.Is(err, gorm.ErrRecordNotFound) || !setting.Enabled { @@ -294,6 +362,22 @@ func validChoice(target, selected string, candidates []string) bool { return false } +func validClosedChoice(selected string, candidates []string) bool { + selected = strings.TrimSpace(selected) + if len(candidates) == 0 { + return selected == "" + } + if selected == "" { + return false + } + for _, candidate := range candidates { + if candidate == selected { + return true + } + } + return false +} + func safeReason(reason string) string { reason = strings.TrimSpace(reason) if reason == "" { @@ -316,6 +400,16 @@ func matchPrompt(request MatchRequest) string { return string(raw) } +func sybSpecParsePrompt(request SYBSpecParseRequest) string { + payload := struct { + ProductSpec string `json:"productSpec"` + Colors []string `json:"shopeeColorCandidates,omitempty"` + Sizes []string `json:"shopeeSizeCandidates,omitempty"` + }{request.ProductSpec, request.Colors, request.Sizes} + raw, _ := json.Marshal(payload) + return string(raw) +} + type openAIMessage struct { Role string `json:"role"` Content string `json:"content"` diff --git a/server/app/goauto/aimatching/service_test.go b/server/app/goauto/aimatching/service_test.go index 0238445..1c99746 100644 --- a/server/app/goauto/aimatching/service_test.go +++ b/server/app/goauto/aimatching/service_test.go @@ -2,6 +2,7 @@ package aimatching import ( "context" + "encoding/json" "errors" "io" "net/http" @@ -135,3 +136,61 @@ func TestAutoConfirmThresholdDefaultsPersistsAndValidates(t *testing.T) { t.Fatalf("invalid threshold err=%v", err) } } + +func TestResolveSYBSpecUsesOnlyProductSpecAndClosedShopeeCandidates(t *testing.T) { + service := matcherTestService(t) + if _, err := service.SaveSettings(context.Background(), SaveSettingsRequest{Enabled: true, BaseURL: "https://provider.example/v1", Model: "test-model", APIKey: "test-secret", TimeoutSeconds: 8}, 7); err != nil { + t.Fatal(err) + } + var sent map[string]any + service.HTTPClient = &http.Client{Transport: roundTripper(func(request *http.Request) (*http.Response, error) { + raw, err := io.ReadAll(request.Body) + if err != nil { + t.Fatal(err) + } + if err := json.Unmarshal(raw, &sent); err != nil { + t.Fatal(err) + } + body := `{"choices":[{"message":{"content":"{\"color\":\"黑色\",\"size\":\"XL\",\"reason\":\"原文对应唯一候选\",\"confidence\":0.95}"}}]}` + return &http.Response{StatusCode: http.StatusOK, Header: make(http.Header), Body: io.NopCloser(strings.NewReader(body)), Request: request}, nil + })} + result, err := service.ResolveSYBSpec(context.Background(), SYBSpecParseRequest{ProductSpec: "黑色 XL【备注】", Colors: []string{"黑色", "白色"}, Sizes: []string{"L", "XL"}}) + if err != nil || result.Color != "黑色" || result.Size != "XL" || result.Confidence == nil || *result.Confidence != 0.95 { + t.Fatalf("result=%+v err=%v", result, err) + } + encoded, _ := json.Marshal(sent) + for _, forbidden := range []string{"orderCode", "address", "rawJson", "price", "test-secret"} { + if strings.Contains(string(encoded), forbidden) { + t.Fatalf("provider payload leaked forbidden field %q: %s", forbidden, encoded) + } + } + for _, required := range []string{"productSpec", "shopeeColorCandidates", "shopeeSizeCandidates"} { + if !strings.Contains(string(encoded), required) { + t.Fatalf("provider payload missing %q: %s", required, encoded) + } + } +} + +func TestResolveSYBSpecRejectsCandidateOutsideClosedSetAndMissingConfidence(t *testing.T) { + service := matcherTestService(t) + if _, err := service.SaveSettings(context.Background(), SaveSettingsRequest{Enabled: true, BaseURL: "https://provider.example/v1", Model: "test-model", APIKey: "test-secret", TimeoutSeconds: 8}, 7); err != nil { + t.Fatal(err) + } + responses := []string{ + `{"choices":[{"message":{"content":"{\"color\":\"灰色\",\"size\":\"XL\",\"reason\":\"猜测\",\"confidence\":0.99}"}}]}`, + `{"choices":[{"message":{"content":"{\"color\":\"黑色\",\"size\":\"XL\",\"reason\":\"候选\"}"}}]}`, + } + service.HTTPClient = &http.Client{Transport: roundTripper(func(request *http.Request) (*http.Response, error) { + body := responses[0] + responses = responses[1:] + return &http.Response{StatusCode: http.StatusOK, Header: make(http.Header), Body: io.NopCloser(strings.NewReader(body)), Request: request}, nil + })} + request := SYBSpecParseRequest{ProductSpec: "黑 XL", Colors: []string{"黑色"}, Sizes: []string{"XL"}} + for i := 0; i < 2; i++ { + _, err := service.ResolveSYBSpec(context.Background(), request) + var target *Error + if !errors.As(err, &target) || target.Code != CodeNoMatch { + t.Fatalf("attempt %d err=%v", i, err) + } + } +} diff --git a/server/app/goauto/aimatching/suggest.go b/server/app/goauto/aimatching/suggest.go index d132b1c..908ea72 100644 --- a/server/app/goauto/aimatching/suggest.go +++ b/server/app/goauto/aimatching/suggest.go @@ -10,6 +10,9 @@ import ( "net/http" "strings" "time" + + log "github.com/go-admin-team/go-admin-core/logger" + "github.com/google/uuid" ) // Suggestion limits are enforced defensively here too, even though callers @@ -95,18 +98,26 @@ func (s *Service) SuggestBatch(ctx context.Context, request SuggestRequest) (Sug } httpRequest.Header.Set("Authorization", "Bearer "+apiKey) httpRequest.Header.Set("Content-Type", "application/json") + callID, startedAt := uuid.NewString(), time.Now() response, err := s.httpClient().Do(httpRequest) if err != nil { + s.logProviderFailure(ProviderFailureDiagnostic{CallID: callID, Operation: "suggest_batch", Kind: providerNetworkErrorKind(err), Duration: time.Since(startedAt)}) return SuggestResult{}, &Error{Code: CodeProviderUnavailable, Message: "AI 建议服务暂时不可用", Cause: err} } defer response.Body.Close() limited := io.LimitReader(response.Body, 1<<20) responseBody, readErr := io.ReadAll(limited) - if readErr != nil || response.StatusCode < http.StatusOK || response.StatusCode >= http.StatusMultipleChoices { + if readErr != nil { + s.logProviderFailure(ProviderFailureDiagnostic{CallID: callID, Operation: "suggest_batch", Kind: "read_error", StatusCode: response.StatusCode, Duration: time.Since(startedAt)}) + return SuggestResult{}, fail(CodeProviderUnavailable, "AI 建议服务暂时不可用") + } + if response.StatusCode < http.StatusOK || response.StatusCode >= http.StatusMultipleChoices { + s.logProviderFailure(ProviderFailureDiagnostic{CallID: callID, Operation: "suggest_batch", Kind: "http_status", StatusCode: response.StatusCode, Duration: time.Since(startedAt)}) return SuggestResult{}, fail(CodeProviderUnavailable, "AI 建议服务暂时不可用") } raw, err := parseSuggestChoices(responseBody) if err != nil { + s.logProviderFailure(ProviderFailureDiagnostic{CallID: callID, Operation: "suggest_batch", Kind: "invalid_response", StatusCode: response.StatusCode, Duration: time.Since(startedAt)}) return SuggestResult{}, fail(CodeProviderUnavailable, "AI 建议响应无效") } @@ -145,6 +156,22 @@ func (s *Service) SuggestBatch(ctx context.Context, request SuggestRequest) (Sug return SuggestResult{Decisions: decisions, Provider: ProviderOpenAICompatible, Model: setting.Model}, nil } +func (s *Service) logProviderFailure(diagnostic ProviderFailureDiagnostic) { + if s.ProviderFailureLogger != nil { + s.ProviderFailureLogger(diagnostic) + return + } + log.Warnf("AI provider call failed: call_id=%s operation=%s kind=%s status=%d duration_ms=%d", + diagnostic.CallID, diagnostic.Operation, diagnostic.Kind, diagnostic.StatusCode, diagnostic.Duration.Milliseconds()) +} + +func providerNetworkErrorKind(err error) string { + if errors.Is(err, context.DeadlineExceeded) { + return "timeout" + } + return "network_error" +} + func suggestSystemPrompt(dimension string) string { noun := "颜色或尺码" switch dimension { diff --git a/server/app/goauto/aimatching/suggest_test.go b/server/app/goauto/aimatching/suggest_test.go index e2db94b..6eb3d0d 100644 --- a/server/app/goauto/aimatching/suggest_test.go +++ b/server/app/goauto/aimatching/suggest_test.go @@ -6,7 +6,9 @@ import ( "fmt" "net/http" "net/http/httptest" + "strings" "testing" + "time" "go-admin/app/goauto/models" @@ -134,3 +136,93 @@ func TestSuggestBatchRequiresConfiguredProvider(t *testing.T) { t.Fatalf("expected CodeNotConfigured, got %v", target.Code) } } + +func TestSuggestBatchLogsSafeDiagnosticForProvider502(t *testing.T) { + const sensitiveBody = "api-key-and-provider-body-must-not-be-logged" + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + http.Error(w, sensitiveBody, http.StatusBadGateway) + })) + defer server.Close() + + db := openSuggestTestDB(t) + seedEnabledSetting(t, db, server.URL) + var diagnostic ProviderFailureDiagnostic + service := NewService(db) + service.ProviderFailureLogger = func(value ProviderFailureDiagnostic) { diagnostic = value } + + _, err := service.SuggestBatch(context.Background(), SuggestRequest{ + Sources: []SuggestSource{{ID: "s1", Label: "sensitive-source"}}, + Candidates: []SuggestCandidate{{ID: "c1", Label: "sensitive-candidate"}}, + }) + if err == nil { + t.Fatal("expected provider failure") + } + if diagnostic.Operation != "suggest_batch" || diagnostic.Kind != "http_status" || diagnostic.StatusCode != http.StatusBadGateway || diagnostic.CallID == "" { + t.Fatalf("unexpected diagnostic: %+v", diagnostic) + } + printed := fmt.Sprintf("%+v", diagnostic) + for _, secret := range []string{sensitiveBody, "sensitive-source", "sensitive-candidate", "test-key", server.URL} { + if strings.Contains(printed, secret) { + t.Fatalf("diagnostic leaked %q: %s", secret, printed) + } + } +} + +func TestSuggestBatchClassifiesProviderTimeout(t *testing.T) { + db := openSuggestTestDB(t) + seedEnabledSetting(t, db, "http://provider.invalid") + var diagnostic ProviderFailureDiagnostic + service := NewService(db) + service.HTTPClient = &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) { + return nil, context.DeadlineExceeded + })} + service.ProviderFailureLogger = func(value ProviderFailureDiagnostic) { diagnostic = value } + + _, err := service.SuggestBatch(context.Background(), SuggestRequest{ + Sources: []SuggestSource{{ID: "s1", Label: "黑色"}}, Candidates: []SuggestCandidate{{ID: "c1", Label: "黑色"}}, + }) + if err == nil || diagnostic.Kind != "timeout" || diagnostic.StatusCode != 0 { + t.Fatalf("timeout was not safely classified: diagnostic=%+v err=%v", diagnostic, err) + } +} + +func TestSuggestBatchClassifiesInvalidResponseWithoutLoggingBody(t *testing.T) { + const sensitiveBody = "not-json-with-sensitive-provider-details" + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + _, _ = w.Write([]byte(sensitiveBody)) + })) + defer server.Close() + db := openSuggestTestDB(t) + seedEnabledSetting(t, db, server.URL) + var diagnostic ProviderFailureDiagnostic + service := NewService(db) + service.ProviderFailureLogger = func(value ProviderFailureDiagnostic) { diagnostic = value } + + _, err := service.SuggestBatch(context.Background(), SuggestRequest{ + Sources: []SuggestSource{{ID: "s1", Label: "黑色"}}, Candidates: []SuggestCandidate{{ID: "c1", Label: "黑色"}}, + }) + if err == nil || diagnostic.Kind != "invalid_response" || strings.Contains(fmt.Sprintf("%+v", diagnostic), sensitiveBody) { + t.Fatalf("invalid response diagnostic is unsafe or missing: diagnostic=%+v err=%v", diagnostic, err) + } +} + +func TestSuggestBatchAllowsProviderResponseAfterTwoSeconds(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + time.Sleep(2100 * time.Millisecond) + chatCompletionResponder(`{"suggestions":[{"sourceId":"s1","candidateId":"c1","confidence":0.95,"reason":"match"}]}`)(w, r) + })) + defer server.Close() + + db := openSuggestTestDB(t) + seedEnabledSetting(t, db, server.URL) + result, err := NewService(db).SuggestBatch(context.Background(), SuggestRequest{ + Sources: []SuggestSource{{ID: "s1", Label: "黑色"}}, Candidates: []SuggestCandidate{{ID: "c1", Label: "黑色"}}, + }) + if err != nil || result.Decisions["s1"].CandidateID != "c1" { + t.Fatalf("delayed provider response failed: result=%+v err=%v", result, err) + } +} + +type roundTripFunc func(*http.Request) (*http.Response, error) + +func (fn roundTripFunc) RoundTrip(request *http.Request) (*http.Response, error) { return fn(request) } diff --git a/server/app/goauto/aimatching/types.go b/server/app/goauto/aimatching/types.go index 327533f..0614ca5 100644 --- a/server/app/goauto/aimatching/types.go +++ b/server/app/goauto/aimatching/types.go @@ -29,6 +29,24 @@ type MatchRequest struct { Sizes []string } +// SYBSpecParseRequest contains the only source text and closed Shopee +// candidate sets that may leave GoAuto for an AI-assisted SYB parse. It must +// never contain the shipment/order, account, address, price or full raw JSON. +type SYBSpecParseRequest struct { + ProductSpec string + Colors []string + Sizes []string +} + +type SYBSpecParseResult struct { + Color string + Size string + Provider string + Model string + Reason string + Confidence *float64 +} + type CandidateSnapshot struct { Colors []string `json:"colors,omitempty"` Sizes []string `json:"sizes,omitempty"` diff --git a/server/app/goauto/migrations/migrate.go b/server/app/goauto/migrations/migrate.go index c0c0d67..8289c27 100644 --- a/server/app/goauto/migrations/migrate.go +++ b/server/app/goauto/migrations/migrate.go @@ -34,7 +34,11 @@ func MigratedModels() []any { &models.PDDProduct{}, &models.AIMatchingSetting{}, &models.ShopeeProduct{}, + &models.ShopeeSpecAutoMatchRun{}, + &models.ShopeeSpecAutoMatchWorkItem{}, &models.SYBProduct{}, + &models.SYBSpecAIParseRun{}, + &models.SYBSpecAIParseWorkItem{}, &models.SYBSession{}, &models.SYBShop{}, &models.SYBSyncRun{}, diff --git a/server/app/goauto/models/schema.go b/server/app/goauto/models/schema.go index 851b0d6..9652cee 100644 --- a/server/app/goauto/models/schema.go +++ b/server/app/goauto/models/schema.go @@ -536,6 +536,15 @@ type SYBProduct struct { // recording the parser's own last output for audit even after a manual // correction; it is not overwritten by the correction itself. ManuallyConfirmed bool `json:"manuallyConfirmed" gorm:"not null;default:false"` + // AIConfirmed is independent of ParseStatus: ParseStatus remains the + // deterministic parser's audit result, while these fields record a closed- + // candidate, high-confidence AI decision. Human correction always clears + // and supersedes this decision. + AIConfirmed bool `json:"aiConfirmed" gorm:"not null;default:false;index"` + AIConfidence *float64 `json:"aiConfidence,omitempty"` + AIReason string `json:"aiReason,omitempty" gorm:"size:500;not null;default:''"` + AIConfirmedAt *time.Time `json:"aiConfirmedAt,omitempty"` + AIInputFingerprint string `json:"-" gorm:"size:64;not null;default:'';index"` // RawJSON is the untouched `details[]` element as SYB returned it. It is // what reparse (#41: "适用于解析规则更新后批量重跑,只读取已保存的原始 diff --git a/server/app/goauto/models/shopee_spec_auto_match.go b/server/app/goauto/models/shopee_spec_auto_match.go new file mode 100644 index 0000000..215207e --- /dev/null +++ b/server/app/goauto/models/shopee_spec_auto_match.go @@ -0,0 +1,56 @@ +package models + +import "time" + +// ShopeeSpecAutoMatchRun is one scheduled or administrator-triggered batch. +// ActiveSlot is 1 only while running; its nullable unique index is the +// database-level cross-process mutex shared by both trigger paths. +type ShopeeSpecAutoMatchRun struct { + ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"` + RequestID string `json:"requestId" gorm:"size:36;not null;uniqueIndex:ux_shopee_spec_auto_match_run_request"` + Trigger string `json:"trigger" gorm:"size:16;not null;index"` + Status string `json:"status" gorm:"size:24;not null;index"` + ActiveSlot *uint8 `json:"-" gorm:"uniqueIndex:ux_shopee_spec_auto_match_run_active"` + LeaseOwner string `json:"-" gorm:"size:64;not null;default:''"` + LeaseExpiresAt *time.Time `json:"-" gorm:"index"` + RequestedBy *uint64 `json:"requestedBy,omitempty"` + BatchLimit int `json:"batchLimit" gorm:"not null;default:20"` + ScannedCount int `json:"scannedCount" gorm:"not null;default:0"` + EligibleCount int `json:"eligibleCount" gorm:"not null;default:0"` + ProcessedCount int `json:"processedCount" gorm:"not null;default:0"` + ConfirmedCount int `json:"confirmedCount" gorm:"not null;default:0"` + UnmatchedCount int `json:"unmatchedCount" gorm:"not null;default:0"` + FailedCount int `json:"failedCount" gorm:"not null;default:0"` + ErrorSummary string `json:"errorSummary,omitempty" gorm:"size:500;not null;default:''"` + StartedAt time.Time `json:"startedAt" gorm:"not null"` + FinishedAt *time.Time `json:"finishedAt,omitempty"` + CreatedAt time.Time `json:"createdAt"` + UpdatedAt time.Time `json:"updatedAt"` +} + +func (ShopeeSpecAutoMatchRun) TableName() string { return "shopee_spec_auto_match_run" } + +// ShopeeSpecAutoMatchWorkItem remembers the last input fingerprint and retry +// state for each product, preventing unchanged low-confidence inputs from +// repeatedly spending AI calls. +type ShopeeSpecAutoMatchWorkItem struct { + ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"` + ShopeeProductID uint64 `json:"shopeeProductId" gorm:"not null;uniqueIndex:ux_shopee_spec_auto_match_work_product"` + RunID *uint64 `json:"runId,omitempty" gorm:"index"` + InputFingerprint string `json:"inputFingerprint" gorm:"size:128;not null;default:'';index"` + Status string `json:"status" gorm:"size:24;not null;index"` + AttemptCount int `json:"attemptCount" gorm:"not null;default:0"` + NextAttemptAt *time.Time `json:"nextAttemptAt,omitempty" gorm:"index"` + LeaseOwner string `json:"-" gorm:"size:64;not null;default:''"` + LeaseExpiresAt *time.Time `json:"-" gorm:"index"` + ConfirmedCount int `json:"confirmedCount" gorm:"not null;default:0"` + UnmatchedCount int `json:"unmatchedCount" gorm:"not null;default:0"` + LastErrorCode string `json:"lastErrorCode,omitempty" gorm:"size:64;not null;default:''"` + LastError string `json:"lastError,omitempty" gorm:"size:500;not null;default:''"` + CreatedAt time.Time `json:"createdAt"` + UpdatedAt time.Time `json:"updatedAt"` +} + +func (ShopeeSpecAutoMatchWorkItem) TableName() string { + return "shopee_spec_auto_match_work_item" +} diff --git a/server/app/goauto/models/syb_spec_ai_parse.go b/server/app/goauto/models/syb_spec_ai_parse.go new file mode 100644 index 0000000..37e6343 --- /dev/null +++ b/server/app/goauto/models/syb_spec_ai_parse.go @@ -0,0 +1,48 @@ +package models + +import "time" + +// SYBSpecAIParseRun is one globally serialized scheduled batch. +type SYBSpecAIParseRun struct { + ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"` + RequestID string `json:"requestId" gorm:"size:36;not null;uniqueIndex:ux_syb_spec_ai_parse_run_request"` + Trigger string `json:"trigger" gorm:"size:16;not null;index"` + Status string `json:"status" gorm:"size:24;not null;index"` + ActiveSlot *uint8 `json:"-" gorm:"uniqueIndex:ux_syb_spec_ai_parse_run_active"` + LeaseOwner string `json:"-" gorm:"size:64;not null;default:''"` + LeaseExpiresAt *time.Time `json:"-" gorm:"index"` + BatchLimit int `json:"batchLimit" gorm:"not null;default:20"` + ScannedCount int `json:"scannedCount" gorm:"not null;default:0"` + EligibleCount int `json:"eligibleCount" gorm:"not null;default:0"` + ProcessedCount int `json:"processedCount" gorm:"not null;default:0"` + ConfirmedCount int `json:"confirmedCount" gorm:"not null;default:0"` + UnmatchedCount int `json:"unmatchedCount" gorm:"not null;default:0"` + FailedCount int `json:"failedCount" gorm:"not null;default:0"` + ErrorSummary string `json:"errorSummary,omitempty" gorm:"size:500;not null;default:''"` + StartedAt time.Time `json:"startedAt" gorm:"not null"` + FinishedAt *time.Time `json:"finishedAt,omitempty"` + CreatedAt time.Time `json:"createdAt"` + UpdatedAt time.Time `json:"updatedAt"` +} + +func (SYBSpecAIParseRun) TableName() string { return "syb_spec_ai_parse_run" } + +// SYBSpecAIParseWorkItem prevents unchanged ambiguous input from repeatedly +// spending provider calls and owns the per-row recovery lease. +type SYBSpecAIParseWorkItem struct { + ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"` + SYBProductID uint64 `json:"sybProductId" gorm:"not null;uniqueIndex:ux_syb_spec_ai_parse_work_product"` + RunID *uint64 `json:"runId,omitempty" gorm:"index"` + InputFingerprint string `json:"inputFingerprint" gorm:"size:64;not null;default:'';index"` + Status string `json:"status" gorm:"size:24;not null;index"` + AttemptCount int `json:"attemptCount" gorm:"not null;default:0"` + NextAttemptAt *time.Time `json:"nextAttemptAt,omitempty" gorm:"index"` + LeaseOwner string `json:"-" gorm:"size:64;not null;default:''"` + LeaseExpiresAt *time.Time `json:"-" gorm:"index"` + LastErrorCode string `json:"lastErrorCode,omitempty" gorm:"size:64;not null;default:''"` + LastError string `json:"lastError,omitempty" gorm:"size:500;not null;default:''"` + CreatedAt time.Time `json:"createdAt"` + UpdatedAt time.Time `json:"updatedAt"` +} + +func (SYBSpecAIParseWorkItem) TableName() string { return "syb_spec_ai_parse_work_item" } diff --git a/server/app/goauto/purchase/ai_match_eligibility.go b/server/app/goauto/purchase/ai_match_eligibility.go index 094a664..add170d 100644 --- a/server/app/goauto/purchase/ai_match_eligibility.go +++ b/server/app/goauto/purchase/ai_match_eligibility.go @@ -32,10 +32,44 @@ type skuCombinationRow struct { } func sybSpecsTrusted(syb models.SYBProduct) bool { - if syb.ParseStatus == models.SYBParseStatusFailed { + if strings.TrimSpace(syb.TargetColor) == "" && strings.TrimSpace(syb.TargetSize) == "" { return false } - return strings.TrimSpace(syb.TargetColor) != "" || strings.TrimSpace(syb.TargetSize) != "" + return syb.ParseStatus != models.SYBParseStatusFailed || syb.ManuallyConfirmed || syb.AIConfirmed +} + +func aiConfirmedSpecsCurrent(syb models.SYBProduct, shopee models.ShopeeProduct) bool { + if !syb.AIConfirmed { + return true + } + specs, err := shopeeproduct.Unmarshal(shopee.SpecsJSON) + if err != nil { + return false + } + wanted := map[string]string{shopeeproduct.RoleColor: strings.TrimSpace(syb.TargetColor), shopeeproduct.RoleSize: strings.TrimSpace(syb.TargetSize)} + foundAny := false + for role, target := range wanted { + if target == "" { + continue + } + foundAny = true + found := false + for _, dimension := range specs { + if dimension.Role != role { + continue + } + for _, value := range dimension.Values { + if value.Name == target { + found = true + break + } + } + } + if !found { + return false + } + } + return foundAny } func (s *Service) loadLatestSKUCombinations(ctx context.Context, pddIDs []uint64, dataset *batchPreviewDataset) error { diff --git a/server/app/goauto/purchase/batch.go b/server/app/goauto/purchase/batch.go index d371b94..683a072 100644 --- a/server/app/goauto/purchase/batch.go +++ b/server/app/goauto/purchase/batch.go @@ -388,6 +388,10 @@ func (s *Service) previewFromDataset(id uint64, dataset batchPreviewDataset, gua item.ReasonCode, item.Reason, item.NextAction = "SHOPEE_NOT_FOUND", "关联的蝦皮商品不存在,请先处理商品档案", "open_shopee" return item } + if !aiConfirmedSpecsCurrent(syb, shopee) { + item.ReasonCode, item.Reason, item.NextAction = "SYB_PARSE_FAILED", "AI 解析依据已变化,请等待重新解析或人工修正", "reparse" + return item + } if shopee.PDDProductID == nil { item.ReasonCode, item.Reason, item.NextAction = "PDD_NOT_LINKED", "尚未关联 PDD 商品,请先关联", processActionOpenPDDLink return item diff --git a/server/app/goauto/purchase/batch_test.go b/server/app/goauto/purchase/batch_test.go index cbe46f0..069cf83 100644 --- a/server/app/goauto/purchase/batch_test.go +++ b/server/app/goauto/purchase/batch_test.go @@ -54,6 +54,9 @@ func TestSybSpecsTrustedOnlyBlocksFailedOrEmptyExtraction(t *testing.T) { {"uncertain with color", models.SYBProduct{ParseStatus: models.SYBParseStatusUncertain, TargetColor: "套装"}, true}, {"uncertain with size", models.SYBProduct{ParseStatus: models.SYBParseStatusUncertain, TargetSize: "均码"}, true}, {"failed with values", models.SYBProduct{ParseStatus: models.SYBParseStatusFailed, TargetColor: "黑色"}, false}, + {"failed AI confirmed", models.SYBProduct{ParseStatus: models.SYBParseStatusFailed, TargetColor: "黑色", AIConfirmed: true}, true}, + {"failed manually confirmed", models.SYBProduct{ParseStatus: models.SYBParseStatusFailed, TargetColor: "黑色", ManuallyConfirmed: true}, true}, + {"AI confirmed without values", models.SYBProduct{ParseStatus: models.SYBParseStatusFailed, AIConfirmed: true}, false}, {"uncertain without values", models.SYBProduct{ParseStatus: models.SYBParseStatusUncertain}, false}, } for _, tt := range tests { @@ -65,6 +68,28 @@ func TestSybSpecsTrustedOnlyBlocksFailedOrEmptyExtraction(t *testing.T) { } } +func TestSYBSpecsTrustedAcceptsAIWithoutCallingItManual(t *testing.T) { + row := models.SYBProduct{ParseStatus: models.SYBParseStatusUncertain, TargetColor: "黑色", AIConfirmed: true} + if !sybSpecsTrusted(row) { + t.Fatal("a valid AI-confirmed parse must pass the purchase parse gate") + } + if row.ManuallyConfirmed { + t.Fatal("AI confirmation must not be represented as human confirmation") + } +} + +func TestAIConfirmedSpecsBecomeUntrustedWhenShopeeCandidateDisappears(t *testing.T) { + row := models.SYBProduct{TargetColor: "黑色", TargetSize: "XL", ParseStatus: models.SYBParseStatusUncertain, AIConfirmed: true} + product := models.ShopeeProduct{SpecsJSON: `[{"name":"颜色","role":"color","values":[{"name":"白色","source":"import"}]},{"name":"尺码","role":"size","values":[{"name":"XL","source":"import"}]}]`} + if aiConfirmedSpecsCurrent(row, product) { + t.Fatal("removed Shopee candidate must invalidate the AI parse gate") + } + product.SpecsJSON = `[{"name":"颜色","role":"color","values":[{"name":"黑色","source":"import"}]},{"name":"尺码","role":"size","values":[{"name":"XL","source":"import"}]}]` + if !aiConfirmedSpecsCurrent(row, product) { + t.Fatal("unchanged exact Shopee candidates should keep AI parse valid") + } +} + func TestBatchPreviewExposesIndependentCollectionEligibility(t *testing.T) { db := testDB(t) fixture := seed(t, db, liveCaps(), true) diff --git a/server/app/goauto/purchase/process_stage.go b/server/app/goauto/purchase/process_stage.go index f2b7f64..0fec4d4 100644 --- a/server/app/goauto/purchase/process_stage.go +++ b/server/app/goauto/purchase/process_stage.go @@ -136,7 +136,7 @@ func processStageFromDataset(id uint64, dataset batchPreviewDataset, preview Bat if !ok { return stage(ProcessStageManualAction, "SYB 商品不存在或已删除", "refresh") } - if syb.ParseStatus == models.SYBParseStatusFailed { + if syb.ParseStatus == models.SYBParseStatusFailed && !syb.ManuallyConfirmed && !syb.AIConfirmed { return stage(ProcessStageManualAction, "解析失败,请先处理", "reparse") } if syb.ShopeeProductID == nil { diff --git a/server/app/goauto/shopeeproduct/auto_match.go b/server/app/goauto/shopeeproduct/auto_match.go index 91998fe..bbc42e7 100644 --- a/server/app/goauto/shopeeproduct/auto_match.go +++ b/server/app/goauto/shopeeproduct/auto_match.go @@ -48,6 +48,9 @@ type autoMatchRef struct { // before the transaction; the transaction rechecks the complete spec context // and PDD candidate set so a stale decision can never be written. func (service *Service) AutoMatchMappings(ctx context.Context, id uint64, request AutoMatchRequest) (AutoMatchResponse, error) { + ctx, cancel := context.WithTimeout(ctx, aimatching.MaxProviderTimeout) + defer cancel() + requestID := strings.TrimSpace(request.RequestID) if _, err := uuid.Parse(requestID); err != nil { return AutoMatchResponse{}, invalidRequest("requestId 必须是 UUID") diff --git a/server/app/goauto/shopeeproduct/auto_match_batch.go b/server/app/goauto/shopeeproduct/auto_match_batch.go new file mode 100644 index 0000000..aaac6ba --- /dev/null +++ b/server/app/goauto/shopeeproduct/auto_match_batch.go @@ -0,0 +1,360 @@ +package shopeeproduct + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "strings" + "time" + + "go-admin/app/goauto/models" + + "github.com/google/uuid" + "gorm.io/gorm" +) + +const ( + SpecAutoMatchInvokeTarget = "GoAutoShopeeSpecAutoMatch" + defaultAutoMatchBatchLimit = 20 + autoMatchLeaseDuration = 30 * time.Minute + autoMatchRetryDelay = time.Hour + maxAutoMatchAttempts = 3 +) + +type AutoMatchRunView struct { + models.ShopeeSpecAutoMatchRun + AlreadyRunning bool `json:"alreadyRunning,omitempty"` + Replayed bool `json:"replayed,omitempty"` +} + +// StartAutoMatchRun acquires the single database-backed activity slot. A +// repeated requestId is idempotent; a concurrent trigger receives the current +// run instead of starting a second batch. +func (service *Service) StartAutoMatchRun(ctx context.Context, trigger, requestID string, requestedBy *uint64, batchLimit int) (AutoMatchRunView, bool, error) { + if _, err := uuid.Parse(strings.TrimSpace(requestID)); err != nil { + return AutoMatchRunView{}, false, invalidRequest("requestId 必须是 UUID") + } + if trigger != "manual" && trigger != "scheduled" { + return AutoMatchRunView{}, false, invalidRequest("trigger 无效") + } + if batchLimit <= 0 { + batchLimit = defaultAutoMatchBatchLimit + } + if batchLimit > 100 { + return AutoMatchRunView{}, false, invalidRequest("batchLimit 不能超过 100") + } + now := time.Now().UTC() + lease := now.Add(autoMatchLeaseDuration) + owner := uuid.NewString() + one := uint8(1) + var result models.ShopeeSpecAutoMatchRun + created := false + err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if err := tx.Model(&models.ShopeeSpecAutoMatchRun{}). + 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(&result).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(&result).Error; err == nil { + result.ActiveSlot = &one + return nil + } else if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + result = models.ShopeeSpecAutoMatchRun{RequestID: requestID, Trigger: trigger, Status: "running", ActiveSlot: &one, LeaseOwner: owner, LeaseExpiresAt: &lease, RequestedBy: requestedBy, BatchLimit: batchLimit, StartedAt: now} + if err := tx.Create(&result).Error; err != nil { + return err + } + created = true + return nil + }) + if err != nil { + // A unique-slot race means another instance won after our read. Return + // its run as the stable, non-error result. + if findErr := service.DB.WithContext(ctx).Where("status = ? AND active_slot = ?", "running", 1).First(&result).Error; findErr == nil { + return AutoMatchRunView{ShopeeSpecAutoMatchRun: result, AlreadyRunning: true}, false, nil + } + return AutoMatchRunView{}, false, internalError(err) + } + view := AutoMatchRunView{ShopeeSpecAutoMatchRun: result} + if !created { + view.AlreadyRunning = result.RequestID != requestID + view.Replayed = result.RequestID == requestID + } + return view, created, nil +} + +func (service *Service) LatestAutoMatchRun(ctx context.Context) (*AutoMatchRunView, error) { + var run models.ShopeeSpecAutoMatchRun + err := service.DB.WithContext(ctx).Order("id DESC").First(&run).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + if err != nil { + return nil, internalError(err) + } + return &AutoMatchRunView{ShopeeSpecAutoMatchRun: run}, nil +} + +// ProcessAutoMatchRun performs a bounded batch. It is safe to call from an +// HTTP-launched goroutine or the scheduler because only the run owning the +// active slot may update and finish itself. +func (service *Service) ProcessAutoMatchRun(ctx context.Context, runID uint64) error { + var run models.ShopeeSpecAutoMatchRun + if err := service.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 = defaultAutoMatchBatchLimit + } + var candidates []models.ShopeeProduct + queryLimit := limit * 25 + if queryLimit < 100 { + queryLimit = 100 + } + if queryLimit > 1000 { + queryLimit = 1000 + } + if err := service.DB.WithContext(ctx). + Joins("JOIN pdd_product ON pdd_product.id = shopee_product.pdd_product_id AND pdd_product.status = ?", "active"). + Where("shopee_product.pdd_product_id IS NOT NULL"). + Order("shopee_product.updated_at ASC, shopee_product.id ASC").Limit(queryLimit).Find(&candidates).Error; err != nil { + service.finishAutoMatchRun(run, "failed", 0, 0, 0, 0, 0, 1, "扫描符合条件的商品失败") + return err + } + + eligible, processed, confirmed, unmatched, failed := 0, 0, 0, 0, 0 + firstError := "" + for _, product := range candidates { + if processed >= limit { + break + } + fingerprint, ok, err := service.autoMatchEligibility(ctx, product) + if err != nil { + failed++ + if firstError == "" { + firstError = safeBatchError(err) + } + continue + } + if !ok { + continue + } + eligible++ + work, claimed, err := service.claimAutoMatchWork(ctx, run, product.ID, fingerprint) + if err != nil { + failed++ + if firstError == "" { + firstError = safeBatchError(err) + } + continue + } + if !claimed { + continue + } + processed++ + service.renewAutoMatchRun(run) + response, matchErr := service.AutoMatchMappings(ctx, product.ID, AutoMatchRequest{RequestID: uuid.NewString(), SpecContextVersion: fingerprint[:64]}) + // fingerprint begins with the 64-character context version. + postFingerprint := fingerprint + if next, _, nextErr := service.autoMatchEligibility(ctx, product); nextErr == nil && next != "" { + postFingerprint = next + } + if matchErr != nil { + failed++ + if firstError == "" { + firstError = safeBatchError(matchErr) + } + service.completeAutoMatchWork(work, postFingerprint, 0, 0, matchErr) + continue + } + confirmed += response.ConfirmedCount + unmatched += response.UnmatchedCount + service.completeAutoMatchWork(work, postFingerprint, response.ConfirmedCount, response.UnmatchedCount, nil) + } + status := "completed" + if failed > 0 { + status = "completed_partial" + } + return service.finishAutoMatchRun(run, status, len(candidates), eligible, processed, confirmed, unmatched, failed, firstError) +} + +func (service *Service) autoMatchEligibility(ctx context.Context, product models.ShopeeProduct) (string, bool, error) { + if product.PDDProductID == nil { + return "", false, nil + } + var pdd models.PDDProduct + if err := service.DB.WithContext(ctx).First(&pdd, *product.PDDProductID).Error; err != nil { + return "", false, err + } + if pdd.Status != "active" { + return "", false, nil + } + shopeeSpecs, err := Unmarshal(product.SpecsJSON) + if err != nil { + return "", false, err + } + shared, needsMatch := false, false + for _, role := range []string{RoleColor, RoleSize} { + pddValues, err := selectablePDDValues(pdd.SpecsJSON, role) + if err != nil { + return "", false, err + } + if len(pddValues) == 0 { + continue + } + for _, dimension := range shopeeSpecs { + if dimension.Role != role || len(dimension.Values) == 0 { + continue + } + shared = true + for _, value := range dimension.Values { + if value.Mapping == nil || value.Mapping.Status != MappingStatusConfirmed || !pddValues[value.Mapping.PDDValue] { + needsMatch = true + } + } + } + } + if !shared || !needsMatch { + return "", false, nil + } + contextVersion := computeSpecContextVersion(product.PDDProductID, product.SpecsJSON, pdd.SpecsJSON) + var setting struct{ UpdatedAt time.Time } + _ = service.DB.WithContext(ctx).Table((models.AIMatchingSetting{}).TableName()).Select("updated_at").Where("id = ?", 1).Scan(&setting).Error + h := sha256.Sum256([]byte(contextVersion + "\x00" + setting.UpdatedAt.UTC().Format(time.RFC3339Nano))) + // Keeping the context version as a prefix lets ProcessAutoMatchRun pass the + // exact version to #194 without re-reading a potentially drifting input. + return contextVersion + hex.EncodeToString(h[:]), true, nil +} + +func (service *Service) claimAutoMatchWork(ctx context.Context, run models.ShopeeSpecAutoMatchRun, productID uint64, fingerprint string) (models.ShopeeSpecAutoMatchWorkItem, bool, error) { + now := time.Now().UTC() + lease := now.Add(autoMatchLeaseDuration) + var work models.ShopeeSpecAutoMatchWorkItem + err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + err := tx.Where("shopee_product_id = ?", productID).First(&work).Error + if errors.Is(err, gorm.ErrRecordNotFound) { + work = models.ShopeeSpecAutoMatchWorkItem{ShopeeProductID: 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 >= maxAutoMatchAttempts || (work.NextAttemptAt != nil && work.NextAttemptAt.After(now)) || (work.Status == "running" && work.LeaseExpiresAt != nil && work.LeaseExpiresAt.After(now)) { + return errWorkNotClaimed + } + } 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.ShopeeSpecAutoMatchWorkItem{}).Where("id = ?", work.ID).Updates(updates).Error; err != nil { + return err + } + return tx.First(&work, work.ID).Error + }) + if errors.Is(err, errWorkNotClaimed) { + return work, false, nil + } + return work, err == nil, err +} + +var errWorkNotClaimed = errors.New("auto match work not claimed") + +func (service *Service) completeAutoMatchWork(work models.ShopeeSpecAutoMatchWorkItem, fingerprint string, confirmed, unmatched int, matchErr error) { + now := time.Now().UTC() + updates := map[string]any{"input_fingerprint": fingerprint, "lease_owner": "", "lease_expires_at": nil, "confirmed_count": confirmed, "unmatched_count": unmatched} + if matchErr == nil { + if unmatched > 0 { + updates["status"] = "unmatched" + } else { + updates["status"] = "completed" + } + updates["next_attempt_at"], updates["last_error_code"], updates["last_error"] = nil, "", "" + } else { + code := batchErrorCode(matchErr) + updates["status"], updates["last_error_code"], updates["last_error"] = "failed", code, safeBatchError(matchErr) + if code == CodeAIUnavailable && work.AttemptCount < maxAutoMatchAttempts { + next := now.Add(autoMatchRetryDelay) + updates["next_attempt_at"] = next + } else { + updates["next_attempt_at"] = nil + } + } + _ = service.DB.Model(&models.ShopeeSpecAutoMatchWorkItem{}).Where("id = ?", work.ID).Updates(updates).Error +} + +func (service *Service) renewAutoMatchRun(run models.ShopeeSpecAutoMatchRun) { + lease := time.Now().UTC().Add(autoMatchLeaseDuration) + _ = service.DB.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ? AND status = ? AND lease_owner = ?", run.ID, "running", run.LeaseOwner).Update("lease_expires_at", lease).Error +} + +func (service *Service) finishAutoMatchRun(run models.ShopeeSpecAutoMatchRun, status string, scanned, eligible, processed, confirmed, unmatched, failed int, summary string) error { + now := time.Now().UTC() + 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": truncateBatchText(summary), "finished_at": now} + return service.DB.Model(&models.ShopeeSpecAutoMatchRun{}).Where("id = ? AND status = ? AND lease_owner = ?", run.ID, "running", run.LeaseOwner).Updates(updates).Error +} + +func batchErrorCode(err error) string { + var serviceErr *ServiceError + if errors.As(err, &serviceErr) { + return serviceErr.Code + } + return CodeInternal +} + +func safeBatchError(err error) string { + var serviceErr *ServiceError + if errors.As(err, &serviceErr) { + return truncateBatchText(serviceErr.Message) + } + return "服务端处理失败" +} + +func truncateBatchText(value string) string { + runes := []rune(strings.TrimSpace(value)) + if len(runes) > 500 { + runes = runes[:500] + } + return string(runes) +} + +type scheduledAutoMatchArgs struct { + BatchLimit int `json:"batchLimit"` +} + +type ScheduledAutoMatchJob struct{} + +func (ScheduledAutoMatchJob) Exec(_ interface{}) error { + return errors.New("规格自动匹配定时任务缺少数据库连接") +} + +func (ScheduledAutoMatchJob) ExecWithDB(db *gorm.DB, arg interface{}) error { + args := scheduledAutoMatchArgs{BatchLimit: defaultAutoMatchBatchLimit} + if raw, ok := arg.(string); ok && strings.TrimSpace(raw) != "" { + if err := json.Unmarshal([]byte(raw), &args); err != nil { + return fmt.Errorf("规格自动匹配参数不是合法 JSON: %w", err) + } + } + if args.BatchLimit < 1 || args.BatchLimit > 100 { + return errors.New("规格自动匹配 batchLimit 必须在 1 到 100 之间") + } + service := NewService(db) + run, created, err := service.StartAutoMatchRun(context.Background(), "scheduled", uuid.NewString(), nil, args.BatchLimit) + if err != nil || !created { + return err + } + return service.ProcessAutoMatchRun(context.Background(), run.ID) +} diff --git a/server/app/goauto/shopeeproduct/auto_match_batch_test.go b/server/app/goauto/shopeeproduct/auto_match_batch_test.go new file mode 100644 index 0000000..cdf9cff --- /dev/null +++ b/server/app/goauto/shopeeproduct/auto_match_batch_test.go @@ -0,0 +1,86 @@ +package shopeeproduct + +import ( + "context" + "testing" + + "go-admin/app/goauto/models" + + "github.com/google/uuid" +) + +func TestAutoMatchRunIsIdempotentAndGloballySerialized(t *testing.T) { + db := openTestDB(t) + service := NewService(db) + requestID := uuid.NewString() + first, created, err := service.StartAutoMatchRun(context.Background(), "manual", requestID, nil, 20) + if err != nil || !created { + t.Fatalf("first=%+v created=%v err=%v", first, created, err) + } + replay, created, err := service.StartAutoMatchRun(context.Background(), "manual", requestID, nil, 20) + if err != nil || created || !replay.Replayed || replay.ID != first.ID { + t.Fatalf("replay=%+v created=%v err=%v", replay, created, err) + } + concurrent, created, err := service.StartAutoMatchRun(context.Background(), "scheduled", uuid.NewString(), nil, 20) + if err != nil || created || !concurrent.AlreadyRunning || concurrent.ID != first.ID { + t.Fatalf("concurrent=%+v created=%v err=%v", concurrent, created, err) + } +} + +func TestProcessAutoMatchRunConfirmsExactSizeAndFinishes(t *testing.T) { + db := openTestDB(t) + pdd := seedPDDProduct(t, db, "active") + service := NewService(db) + createdProduct, err := service.Create(context.Background(), CreateRequest{ + RequestID: uuid.NewString(), ShopeeItemID: "SP-BATCH-EXACT", PDDProductID: &pdd.ID, + Specs: []SpecDimension{{Name: "尺码", Role: RoleSize, Values: []SpecValue{{Name: " xl ", Source: ValueSourceImport}}}}, + }) + if err != nil { + t.Fatal(err) + } + run, started, err := service.StartAutoMatchRun(context.Background(), "manual", uuid.NewString(), nil, 20) + if err != nil || !started { + t.Fatalf("run=%+v started=%v err=%v", run, started, err) + } + if err := service.ProcessAutoMatchRun(context.Background(), run.ID); err != nil { + t.Fatal(err) + } + latest, err := service.LatestAutoMatchRun(context.Background()) + if err != nil { + t.Fatal(err) + } + if latest == nil || latest.Status != "completed" || latest.ProcessedCount != 1 || latest.ConfirmedCount != 1 || latest.ActiveSlot != nil { + t.Fatalf("latest=%+v", latest) + } + detail, err := service.Detail(context.Background(), createdProduct.Product.ID) + if err != nil { + t.Fatal(err) + } + mapping := detail.Product.Specs[0].Values[0].Mapping + if mapping == nil || mapping.Status != MappingStatusConfirmed || mapping.PDDValue != "XL" { + t.Fatalf("mapping=%+v", mapping) + } +} + +func TestUnchangedUnmatchedWorkIsNotClaimedAgain(t *testing.T) { + db := openTestDB(t) + service := NewService(db) + one := uint8(1) + run := models.ShopeeSpecAutoMatchRun{RequestID: uuid.NewString(), Trigger: "manual", Status: "running", ActiveSlot: &one, LeaseOwner: uuid.NewString(), BatchLimit: 20} + if err := db.Create(&run).Error; err != nil { + t.Fatal(err) + } + work, claimed, err := service.claimAutoMatchWork(context.Background(), run, 99, "fingerprint") + if err != nil || !claimed { + t.Fatalf("work=%+v claimed=%v err=%v", work, claimed, err) + } + service.completeAutoMatchWork(work, "fingerprint", 0, 1, nil) + _, claimed, err = service.claimAutoMatchWork(context.Background(), run, 99, "fingerprint") + if err != nil || claimed { + t.Fatalf("unchanged unmatched claimed=%v err=%v", claimed, err) + } + _, claimed, err = service.claimAutoMatchWork(context.Background(), run, 99, "changed") + if err != nil || !claimed { + t.Fatalf("changed input claimed=%v err=%v", claimed, err) + } +} diff --git a/server/app/goauto/shopeeproduct/auto_match_test.go b/server/app/goauto/shopeeproduct/auto_match_test.go index 46bfaa1..6ea3d0a 100644 --- a/server/app/goauto/shopeeproduct/auto_match_test.go +++ b/server/app/goauto/shopeeproduct/auto_match_test.go @@ -5,11 +5,13 @@ import ( "encoding/json" "net/http" "net/http/httptest" + "strings" "sync/atomic" "testing" "go-admin/app/goauto/models" + "github.com/gin-gonic/gin" "github.com/google/uuid" ) @@ -116,6 +118,20 @@ func TestAutoMatchMappingsProviderFailureDoesNotWrite(t *testing.T) { } } +func TestWriteErrorMapsAIUnavailableToStructured503(t *testing.T) { + gin.SetMode(gin.TestMode) + recorder := httptest.NewRecorder() + context, _ := gin.CreateTestContext(recorder) + writeError(context, aiUnavailable("AI 匹配服务暂时不可用,请稍后重试")) + + if recorder.Code != http.StatusServiceUnavailable { + t.Fatalf("status = %d, body = %s", recorder.Code, recorder.Body.String()) + } + if !strings.Contains(recorder.Body.String(), `"code":"AI_MATCHING_UNAVAILABLE"`) { + t.Fatalf("response is not structured: %s", recorder.Body.String()) + } +} + func TestAutoMatchMappingsRejectsContextDriftBeforeAtomicWrite(t *testing.T) { db := openTestDB(t) pdd := seedPDDProduct(t, db, "active") diff --git a/server/app/goauto/shopeeproduct/handler.go b/server/app/goauto/shopeeproduct/handler.go index 82e1e19..5d751a5 100644 --- a/server/app/goauto/shopeeproduct/handler.go +++ b/server/app/goauto/shopeeproduct/handler.go @@ -1,6 +1,7 @@ package shopeeproduct import ( + "context" "encoding/json" "errors" "io" @@ -300,6 +301,45 @@ func (handler Handler) AutoMatchMappings(c *gin.Context) { c.JSON(http.StatusOK, gin.H{"code": 200, "data": response}) } +func (handler Handler) StartAutoMatchRun(c *gin.Context) { + var request struct { + RequestID string `json:"requestId"` + } + if err := decodeJSON(c, &request); err != nil { + writeError(c, invalidRequest("请求 JSON 无效")) + return + } + service, ok := handler.service(c) + if !ok { + return + } + operator := currentUserID(c) + run, created, err := service.StartAutoMatchRun(c.Request.Context(), "manual", request.RequestID, &operator, defaultAutoMatchBatchLimit) + if err != nil { + writeError(c, err) + return + } + if created { + go func(runID uint64, db *gorm.DB) { + _ = NewService(db).ProcessAutoMatchRun(context.Background(), runID) + }(run.ID, service.DB) + } + c.JSON(http.StatusAccepted, gin.H{"code": 200, "data": gin.H{"run": run}}) +} + +func (handler Handler) LatestAutoMatchRun(c *gin.Context) { + service, ok := handler.service(c) + if !ok { + return + } + run, err := service.LatestAutoMatchRun(c.Request.Context()) + if err != nil { + writeError(c, err) + return + } + c.JSON(http.StatusOK, gin.H{"code": 200, "data": gin.H{"run": run}}) +} + func (handler Handler) BatchDelete(c *gin.Context) { var request BatchDeleteRequest if err := decodeJSON(c, &request); err != nil { diff --git a/server/app/goauto/shopeeproduct/router.go b/server/app/goauto/shopeeproduct/router.go index ec5982b..5ec039a 100644 --- a/server/app/goauto/shopeeproduct/router.go +++ b/server/app/goauto/shopeeproduct/router.go @@ -10,6 +10,9 @@ import ( func InitRouter(engine *gin.Engine, auth *jwt.GinJWTMiddleware) { handler := Handler{} admin := engine.Group("/api/admin/v1/shopee-products").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()) + adminOnlyRuns := engine.Group("/api/admin/v1/shopee-spec-auto-match/runs").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(middleware.RequireRoleKey("admin")) + adminOnlyRuns.POST("", handler.StartAutoMatchRun) + adminOnlyRuns.GET("/latest", handler.LatestAutoMatchRun) admin.GET("", handler.List) admin.POST("", handler.Create) admin.POST("/batch-delete", handler.BatchDelete) diff --git a/server/app/goauto/sybimport/ai_parse_batch.go b/server/app/goauto/sybimport/ai_parse_batch.go new file mode 100644 index 0000000..be815a9 --- /dev/null +++ b/server/app/goauto/sybimport/ai_parse_batch.go @@ -0,0 +1,489 @@ +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) +} diff --git a/server/app/goauto/sybimport/ai_parse_batch_test.go b/server/app/goauto/sybimport/ai_parse_batch_test.go new file mode 100644 index 0000000..cdbdc8f --- /dev/null +++ b/server/app/goauto/sybimport/ai_parse_batch_test.go @@ -0,0 +1,204 @@ +package sybimport_test + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + "go-admin/app/goauto/models" + "go-admin/app/goauto/shopeeproduct" + "go-admin/app/goauto/sybimport" + + "github.com/google/uuid" + "gorm.io/gorm" +) + +func seedAIParseCandidate(t *testing.T, db *gorm.DB, serverURL string, detailID uint64, productSpec string) models.SYBProduct { + t.Helper() + detail := realDetailA() + detail.ID = detailID + detail.ProductSpec = productSpec + detail.Raw = []byte(fmt.Sprintf(`{"id":%d,"productSpec":%q}`, detail.ID, productSpec)) + applied, err := sybimport.ApplyDetail(context.Background(), db, realOrder(), detail) + if err != nil { + t.Fatal(err) + } + specs, err := shopeeproduct.Marshal([]shopeeproduct.SpecDimension{ + {Name: "颜色", Role: shopeeproduct.RoleColor, Values: []shopeeproduct.SpecValue{{Name: "黑色", Source: shopeeproduct.ValueSourceImport}, {Name: "白色", Source: shopeeproduct.ValueSourceImport}}}, + {Name: "尺码", Role: shopeeproduct.RoleSize, Values: []shopeeproduct.SpecValue{{Name: "L", Source: shopeeproduct.ValueSourceImport}, {Name: "XL", Source: shopeeproduct.ValueSourceImport}}}, + }) + if err != nil { + t.Fatal(err) + } + if err := db.Model(&models.ShopeeProduct{}).Where("id = ?", *applied.SYBProduct.ShopeeProductID).Update("specs_json", specs).Error; err != nil { + t.Fatal(err) + } + setting := models.AIMatchingSetting{ID: 1, Enabled: true, Provider: "openai_compatible", BaseURL: serverURL, Model: "test-model", APIKey: "test-key", TimeoutSeconds: 5, AutoConfirmMinConfidence: 0.9} + if err := db.Save(&setting).Error; err != nil { + t.Fatal(err) + } + return applied.SYBProduct +} + +func TestScheduledAIParseConfirmsClosedCandidatesAndDoesNotRepeat(t *testing.T) { + var calls atomic.Int32 + provider := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + calls.Add(1) + response.Header().Set("Content-Type", "application/json") + _, _ = response.Write([]byte(`{"choices":[{"message":{"content":"{\"color\":\"黑色\",\"size\":\"XL\",\"reason\":\"原文对应唯一候选\",\"confidence\":0.95}"}}]}`)) + })) + defer provider.Close() + db := openTestDB(t) + record := seedAIParseCandidate(t, db, provider.URL, 19801, "黑色 XL") + rawBefore := record.RawJSON + run, created, err := sybimport.StartSpecAIParseRun(context.Background(), db, uuid.NewString(), 20) + if err != nil || !created { + t.Fatalf("start run: created=%v err=%v", created, err) + } + if err := sybimport.ProcessSpecAIParseRun(context.Background(), db, run.ID); err != nil { + t.Fatal(err) + } + if err := db.First(&record, record.ID).Error; err != nil { + t.Fatal(err) + } + if !record.AIConfirmed || record.ManuallyConfirmed || record.ParseStatus != models.SYBParseStatusUncertain || record.TargetColor != "黑色" || record.TargetSize != "XL" || record.AIConfidence == nil || *record.AIConfidence != 0.95 || record.AIReason == "" || record.AIInputFingerprint == "" { + t.Fatalf("unexpected confirmed record: %+v", record) + } + if record.RawJSON != rawBefore { + t.Fatal("AI confirmation must not rewrite RawJSON") + } + var finished models.SYBSpecAIParseRun + if err := db.First(&finished, run.ID).Error; err != nil { + t.Fatal(err) + } + if finished.Status != "completed" || finished.ConfirmedCount != 1 || finished.ProcessedCount != 1 { + t.Fatalf("unexpected run: %+v", finished) + } + second, created, err := sybimport.StartSpecAIParseRun(context.Background(), db, uuid.NewString(), 20) + if err != nil || !created { + t.Fatalf("second start: created=%v err=%v", created, err) + } + if err := sybimport.ProcessSpecAIParseRun(context.Background(), db, second.ID); err != nil { + t.Fatal(err) + } + if calls.Load() != 1 { + t.Fatalf("unchanged confirmed input called provider %d times", calls.Load()) + } +} + +func TestScheduledAIParseLeavesLowConfidenceUnmatchedForSameFingerprint(t *testing.T) { + var calls atomic.Int32 + provider := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + calls.Add(1) + _, _ = response.Write([]byte(`{"choices":[{"message":{"content":"{\"color\":\"黑色\",\"size\":\"XL\",\"reason\":\"仍有歧义\",\"confidence\":0.4}"}}]}`)) + })) + defer provider.Close() + db := openTestDB(t) + record := seedAIParseCandidate(t, db, provider.URL, 19802, "黑色 XL") + for i := 0; i < 2; i++ { + run, created, err := sybimport.StartSpecAIParseRun(context.Background(), db, uuid.NewString(), 20) + if err != nil || !created { + t.Fatalf("start %d: created=%v err=%v", i, created, err) + } + if err := sybimport.ProcessSpecAIParseRun(context.Background(), db, run.ID); err != nil { + t.Fatal(err) + } + } + if err := db.First(&record, record.ID).Error; err != nil { + t.Fatal(err) + } + if record.AIConfirmed || calls.Load() != 1 { + t.Fatalf("low confidence must remain unconfirmed and not repeat: confirmed=%v calls=%d", record.AIConfirmed, calls.Load()) + } + var work models.SYBSpecAIParseWorkItem + if err := db.Where("syb_product_id = ?", record.ID).First(&work).Error; err != nil { + t.Fatal(err) + } + if work.Status != "unmatched" { + t.Fatalf("work status=%s", work.Status) + } +} + +func TestScheduledAIParseSkipsEmptySourceAndManualConfirmation(t *testing.T) { + var calls atomic.Int32 + provider := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + calls.Add(1) + response.WriteHeader(http.StatusInternalServerError) + })) + defer provider.Close() + db := openTestDB(t) + empty := seedAIParseCandidate(t, db, provider.URL, 19803, "") + manual := seedAIParseCandidate(t, db, provider.URL, 19804, "黑色 XL") + if _, err := sybimport.ManualCorrect(context.Background(), db, manual.ID, "黑色", "XL"); err != nil { + t.Fatal(err) + } + run, created, err := sybimport.StartSpecAIParseRun(context.Background(), db, uuid.NewString(), 20) + if err != nil || !created { + t.Fatalf("start: created=%v err=%v", created, err) + } + if err := sybimport.ProcessSpecAIParseRun(context.Background(), db, run.ID); err != nil { + t.Fatal(err) + } + if calls.Load() != 0 { + t.Fatalf("ineligible rows called provider %d times", calls.Load()) + } + var workCount int64 + if err := db.Model(&models.SYBSpecAIParseWorkItem{}).Where("syb_product_id IN ?", []uint64{empty.ID, manual.ID}).Count(&workCount).Error; err != nil { + t.Fatal(err) + } + if workCount != 0 { + t.Fatalf("ineligible rows created %d work items", workCount) + } +} + +func TestSpecAIParseRunHasSingleGlobalActiveSlot(t *testing.T) { + db := openTestDB(t) + first, created, err := sybimport.StartSpecAIParseRun(context.Background(), db, uuid.NewString(), 20) + if err != nil || !created { + t.Fatalf("first: created=%v err=%v", created, err) + } + second, created, err := sybimport.StartSpecAIParseRun(context.Background(), db, uuid.NewString(), 20) + if err != nil || created || second.ID != first.ID { + t.Fatalf("second must reuse active run: first=%d second=%d created=%v err=%v", first.ID, second.ID, created, err) + } +} + +func TestScheduledAIParseRetriesProviderFailureAtMostThreeTimes(t *testing.T) { + var calls atomic.Int32 + provider := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + calls.Add(1) + response.WriteHeader(http.StatusBadGateway) + })) + defer provider.Close() + db := openTestDB(t) + record := seedAIParseCandidate(t, db, provider.URL, 19805, "黑色 XL") + for attempt := 1; attempt <= 4; attempt++ { + run, created, err := sybimport.StartSpecAIParseRun(context.Background(), db, uuid.NewString(), 20) + if err != nil || !created { + t.Fatalf("start %d: created=%v err=%v", attempt, created, err) + } + if err := sybimport.ProcessSpecAIParseRun(context.Background(), db, run.ID); err != nil { + t.Fatal(err) + } + if attempt < 3 { + past := time.Now().UTC().Add(-time.Minute) + if err := db.Model(&models.SYBSpecAIParseWorkItem{}).Where("syb_product_id = ?", record.ID).Update("next_attempt_at", past).Error; err != nil { + t.Fatal(err) + } + } + } + if calls.Load() != 3 { + t.Fatalf("provider calls=%d, want 3", calls.Load()) + } + var work models.SYBSpecAIParseWorkItem + if err := db.Where("syb_product_id = ?", record.ID).First(&work).Error; err != nil { + t.Fatal(err) + } + if work.AttemptCount != 3 || work.NextAttemptAt != nil { + t.Fatalf("retry state=%+v", work) + } +} diff --git a/server/app/goauto/sybimport/apply.go b/server/app/goauto/sybimport/apply.go index ed27d0a..8e18a2b 100644 --- a/server/app/goauto/sybimport/apply.go +++ b/server/app/goauto/sybimport/apply.go @@ -130,13 +130,37 @@ func ApplyDetail(ctx context.Context, db *gorm.DB, order OrderInput, detail Deta result.Outcome = OutcomeCreated case err == nil: record.ID = existing.ID - if err := tx.Model(&models.SYBProduct{}).Where("id = ?", existing.ID).Updates(map[string]any{ + // Human-confirmed target values are authoritative and survive every + // source re-import. ParseStatus/ParseNote below still record what the + // current deterministic parser observed for audit. + if existing.ManuallyConfirmed { + record.TargetColor, record.TargetSize = existing.TargetColor, existing.TargetSize + record.ManuallyConfirmed = true + } + updates := map[string]any{ "stock_id": record.StockID, "shop_name": record.ShopName, "shopee_item_id": record.ShopeeItemID, "shopee_product_id": record.ShopeeProductID, "product_title": record.ProductTitle, "target_color": record.TargetColor, "target_size": record.TargetSize, "quantity": record.Quantity, "unit_price_cent": record.UnitPriceCent, "image_url": record.ImageURL, "parse_status": record.ParseStatus, "parse_note": record.ParseNote, "raw_json": record.RawJSON, - }).Error; err != nil { + } + // An identical re-import keeps a valid AI decision. Changed source, + // link, or a newly deterministic parse invalidates it atomically. + preserveAI := existing.AIConfirmed && !existing.ManuallyConfirmed && parsed.Status != models.SYBParseStatusSuccess && + existing.RawJSON == record.RawJSON && sameOptionalID(existing.ShopeeProductID, record.ShopeeProductID) + if preserveAI { + record.TargetColor, record.TargetSize = existing.TargetColor, existing.TargetSize + updates["target_color"], updates["target_size"] = existing.TargetColor, existing.TargetSize + record.AIConfirmed, record.AIConfidence, record.AIReason = true, existing.AIConfidence, existing.AIReason + record.AIConfirmedAt, record.AIInputFingerprint = existing.AIConfirmedAt, existing.AIInputFingerprint + } else { + updates["ai_confirmed"], updates["ai_confidence"], updates["ai_reason"] = false, nil, "" + updates["ai_confirmed_at"], updates["ai_input_fingerprint"] = nil, "" + } + if existing.ManuallyConfirmed { + updates["target_color"], updates["target_size"] = existing.TargetColor, existing.TargetSize + } + if err := tx.Model(&models.SYBProduct{}).Where("id = ?", existing.ID).Updates(updates).Error; err != nil { return err } result.Outcome = OutcomeUpdated @@ -161,6 +185,13 @@ func ApplyDetail(ctx context.Context, db *gorm.DB, order OrderInput, detail Deta return result, nil } +func sameOptionalID(left, right *uint64) bool { + if left == nil || right == nil { + return left == nil && right == nil + } + return *left == *right +} + // findOrCreateShopeeProduct implements #40's revival rule: a live match wins, // a soft-deleted match is revived (keeping its prior mapping), and only when // neither exists does the import create a minimal archive. On an existing diff --git a/server/app/goauto/sybimport/reparse.go b/server/app/goauto/sybimport/reparse.go index 4ade256..195056d 100644 --- a/server/app/goauto/sybimport/reparse.go +++ b/server/app/goauto/sybimport/reparse.go @@ -28,7 +28,7 @@ type rawDetailSpec struct { // result page's per-line feedback. type ReparseOutcome struct { SYBProductID uint64 `json:"sybProductId"` - Outcome string `json:"outcome"` // reparsed | skipped_manual | unchanged + Outcome string `json:"outcome"` // reparsed | skipped_manual | skipped_ai | unchanged OldStatus string `json:"oldStatus"` NewStatus string `json:"newStatus"` } @@ -36,6 +36,7 @@ type ReparseOutcome struct { const ( ReparseOutcomeReparsed = "reparsed" ReparseOutcomeSkippedManual = "skipped_manual" + ReparseOutcomeSkippedAI = "skipped_ai" ReparseOutcomeUnchanged = "unchanged" ) @@ -65,6 +66,11 @@ func Reparse(ctx context.Context, db *gorm.DB, sybProductID uint64, force bool) outcome.NewStatus = record.ParseStatus return nil } + if record.AIConfirmed && !force { + outcome.Outcome = ReparseOutcomeSkippedAI + outcome.NewStatus = record.ParseStatus + return nil + } var raw rawDetailSpec if err := json.Unmarshal([]byte(record.RawJSON), &raw); err != nil { @@ -76,8 +82,11 @@ func Reparse(ctx context.Context, db *gorm.DB, sybProductID uint64, force bool) if parsed.Color == record.TargetColor && parsed.Size == record.TargetSize && parsed.Status == record.ParseStatus { outcome.Outcome = ReparseOutcomeUnchanged if force { - record.ManuallyConfirmed = false - if err := tx.Model(&models.SYBProduct{}).Where("id = ?", record.ID).Update("manually_confirmed", false).Error; err != nil { + record.ManuallyConfirmed, record.AIConfirmed = false, false + if err := tx.Model(&models.SYBProduct{}).Where("id = ?", record.ID).Updates(map[string]any{ + "manually_confirmed": false, "ai_confirmed": false, "ai_confidence": nil, + "ai_reason": "", "ai_confirmed_at": nil, "ai_input_fingerprint": "", + }).Error; err != nil { return err } } @@ -87,6 +96,8 @@ func Reparse(ctx context.Context, db *gorm.DB, sybProductID uint64, force bool) updates := map[string]any{ "target_color": parsed.Color, "target_size": parsed.Size, "parse_status": parsed.Status, "parse_note": parsed.Note, "manually_confirmed": false, + "ai_confirmed": false, "ai_confidence": nil, "ai_reason": "", + "ai_confirmed_at": nil, "ai_input_fingerprint": "", } if err := tx.Model(&models.SYBProduct{}).Where("id = ?", record.ID).Updates(updates).Error; err != nil { return err @@ -140,8 +151,12 @@ func ManualCorrect(ctx context.Context, db *gorm.DB, sybProductID uint64, color, return err } record.TargetColor, record.TargetSize, record.ManuallyConfirmed = color, size, true + record.AIConfirmed, record.AIConfidence, record.AIReason = false, nil, "" + record.AIConfirmedAt, record.AIInputFingerprint = nil, "" if err := tx.Model(&models.SYBProduct{}).Where("id = ?", sybProductID).Updates(map[string]any{ "target_color": color, "target_size": size, "manually_confirmed": true, + "ai_confirmed": false, "ai_confidence": nil, "ai_reason": "", + "ai_confirmed_at": nil, "ai_input_fingerprint": "", }).Error; err != nil { return err } diff --git a/server/app/goauto/sybimport/reparse_test.go b/server/app/goauto/sybimport/reparse_test.go index 8f9c712..ebd5b97 100644 --- a/server/app/goauto/sybimport/reparse_test.go +++ b/server/app/goauto/sybimport/reparse_test.go @@ -60,6 +60,35 @@ func TestReparseSkipsManuallyConfirmedRowByDefault(t *testing.T) { } } +func TestReparseSkipsAIConfirmationUnlessForced(t *testing.T) { + db := openTestDB(t) + applied, err := sybimport.ApplyDetail(context.Background(), db, realOrder(), realDetailB()) + if err != nil { + t.Fatal(err) + } + confidence := 0.95 + if err := db.Model(&models.SYBProduct{}).Where("id = ?", applied.SYBProduct.ID).Updates(map[string]any{ + "target_color": "AI颜色", "target_size": "AI尺码", "ai_confirmed": true, + "ai_confidence": confidence, "ai_reason": "AI 结果", "ai_input_fingerprint": strings.Repeat("c", 64), + }).Error; err != nil { + t.Fatal(err) + } + outcome, err := sybimport.Reparse(context.Background(), db, applied.SYBProduct.ID, false) + if err != nil || outcome.Outcome != sybimport.ReparseOutcomeSkippedAI { + t.Fatalf("unforced outcome=%+v err=%v", outcome, err) + } + if _, err := sybimport.Reparse(context.Background(), db, applied.SYBProduct.ID, true); err != nil { + t.Fatal(err) + } + var record models.SYBProduct + if err := db.First(&record, applied.SYBProduct.ID).Error; err != nil { + t.Fatal(err) + } + if record.AIConfirmed || record.AIConfidence != nil || record.AIReason != "" || record.AIInputFingerprint != "" { + t.Fatalf("forced reparse retained AI state: %+v", record) + } +} + // 可勾选强制覆盖. func TestReparseWithForceOverridesManualCorrection(t *testing.T) { db := openTestDB(t) @@ -151,6 +180,83 @@ func TestManualCorrectMergesIntoArchiveLikeASuccessfulParse(t *testing.T) { } } +func TestManualCorrectSupersedesAIConfirmation(t *testing.T) { + db := openTestDB(t) + applied, err := sybimport.ApplyDetail(context.Background(), db, realOrder(), realDetailB()) + if err != nil { + t.Fatal(err) + } + confidence := 0.96 + if err := db.Model(&models.SYBProduct{}).Where("id = ?", applied.SYBProduct.ID).Updates(map[string]any{ + "ai_confirmed": true, "ai_confidence": confidence, "ai_reason": "旧 AI 结果", "ai_input_fingerprint": strings.Repeat("a", 64), + }).Error; err != nil { + t.Fatal(err) + } + corrected, err := sybimport.ManualCorrect(context.Background(), db, applied.SYBProduct.ID, "人工颜色", "人工尺码") + if err != nil { + t.Fatal(err) + } + if !corrected.ManuallyConfirmed || corrected.AIConfirmed || corrected.AIConfidence != nil || corrected.AIReason != "" || corrected.AIInputFingerprint != "" { + t.Fatalf("human correction did not supersede AI state: %+v", corrected) + } +} + +func TestReimportPreservesIdenticalAIInputAndInvalidatesChangedSource(t *testing.T) { + db := openTestDB(t) + order, detail := realOrder(), realDetailB() + first, err := sybimport.ApplyDetail(context.Background(), db, order, detail) + if err != nil { + t.Fatal(err) + } + confidence := 0.95 + if err := db.Model(&models.SYBProduct{}).Where("id = ?", first.SYBProduct.ID).Updates(map[string]any{ + "target_color": "AI颜色", "target_size": "AI尺码", "ai_confirmed": true, + "ai_confidence": confidence, "ai_reason": "已确认", "ai_input_fingerprint": strings.Repeat("b", 64), + }).Error; err != nil { + t.Fatal(err) + } + same, err := sybimport.ApplyDetail(context.Background(), db, order, detail) + if err != nil { + t.Fatal(err) + } + if !same.SYBProduct.AIConfirmed { + t.Fatal("identical re-import must preserve AI confirmation") + } + if same.SYBProduct.TargetColor != "AI颜色" || same.SYBProduct.TargetSize != "AI尺码" { + t.Fatal("identical re-import must preserve AI-confirmed target values") + } + detail.ProductSpec += " 新备注" + detail.Raw = []byte(`{"productSpec":"changed"}`) + changed, err := sybimport.ApplyDetail(context.Background(), db, order, detail) + if err != nil { + t.Fatal(err) + } + if changed.SYBProduct.AIConfirmed || changed.SYBProduct.AIConfidence != nil || changed.SYBProduct.AIReason != "" { + t.Fatalf("changed source retained stale AI state: %+v", changed.SYBProduct) + } +} + +func TestReimportNeverOverwritesManualTargetValues(t *testing.T) { + db := openTestDB(t) + order, detail := realOrder(), realDetailB() + first, err := sybimport.ApplyDetail(context.Background(), db, order, detail) + if err != nil { + t.Fatal(err) + } + if _, err := sybimport.ManualCorrect(context.Background(), db, first.SYBProduct.ID, "人工颜色", "人工尺码"); err != nil { + t.Fatal(err) + } + detail.ProductSpec = "来源新颜色,来源新尺码" + detail.Raw = []byte(`{"productSpec":"来源新颜色,来源新尺码"}`) + updated, err := sybimport.ApplyDetail(context.Background(), db, order, detail) + if err != nil { + t.Fatal(err) + } + if !updated.SYBProduct.ManuallyConfirmed || updated.SYBProduct.TargetColor != "人工颜色" || updated.SYBProduct.TargetSize != "人工尺码" { + t.Fatalf("re-import overwrote human decision: %+v", updated.SYBProduct) + } +} + // Regression test: ReparseOutcome originally had no json tags at all, so Go's // default marshaling produced PascalCase keys ("SYBProductID", "OldStatus") // instead of the camelCase the rest of this API and the admin frontend use. diff --git a/server/app/jobs/apis/execution_log.go b/server/app/jobs/apis/execution_log.go new file mode 100644 index 0000000..86cf9cd --- /dev/null +++ b/server/app/jobs/apis/execution_log.go @@ -0,0 +1,87 @@ +package apis + +import ( + "errors" + "net/http" + "strconv" + "strings" + "time" + + "github.com/gin-gonic/gin" + + jobservice "go-admin/app/jobs/service" +) + +func (e SysJob) ListExecutionLogs(c *gin.Context) { + jobID, err := strconv.Atoi(c.Param("id")) + if err != nil || jobID < 1 { + c.JSON(http.StatusUnprocessableEntity, gin.H{"code": 400, "msg": "jobId 无效"}) + return + } + page, err := positiveQueryInt(c.Query("pageIndex"), 1) + if err != nil { + c.JSON(http.StatusUnprocessableEntity, gin.H{"code": 400, "msg": "pageIndex 必须是正整数"}) + return + } + pageSize, err := positiveQueryInt(c.Query("pageSize"), 20) + if err != nil || pageSize > 100 { + c.JSON(http.StatusUnprocessableEntity, gin.H{"code": 400, "msg": "pageSize 必须是 1 到 100 的整数"}) + return + } + startedFrom, err := optionalTime(c.Query("startedFrom")) + if err != nil { + c.JSON(http.StatusUnprocessableEntity, gin.H{"code": 400, "msg": "startedFrom 必须是 RFC3339 时间"}) + return + } + startedTo, err := optionalTime(c.Query("startedTo")) + if err != nil { + c.JSON(http.StatusUnprocessableEntity, gin.H{"code": 400, "msg": "startedTo 必须是 RFC3339 时间"}) + return + } + + e.MakeContext(c) + db, err := e.GetOrm() + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "msg": "服务端处理失败"}) + return + } + response, err := jobservice.NewExecutionLogService(db).List(c.Request.Context(), jobID, jobservice.ExecutionLogListRequest{ + Page: page, PageSize: pageSize, Status: strings.TrimSpace(c.Query("status")), + StartedFrom: startedFrom, StartedTo: startedTo, + }) + if err != nil { + switch { + case errors.Is(err, jobservice.ErrExecutionLogInvalidRequest): + c.JSON(http.StatusUnprocessableEntity, gin.H{"code": 400, "msg": err.Error()}) + case errors.Is(err, jobservice.ErrExecutionLogJobNotFound): + c.JSON(http.StatusNotFound, gin.H{"code": 404, "msg": "定时任务不存在"}) + default: + e.GetLogger().Errorf("list scheduled job execution logs failed job_id=%d: %v", jobID, err) + c.JSON(http.StatusInternalServerError, gin.H{"code": 500, "msg": "服务端处理失败"}) + } + return + } + c.JSON(http.StatusOK, gin.H{"code": 200, "data": response}) +} + +func positiveQueryInt(value string, fallback int) (int, error) { + if strings.TrimSpace(value) == "" { + return fallback, nil + } + parsed, err := strconv.Atoi(value) + if err != nil || parsed < 1 { + return 0, errors.New("invalid positive integer") + } + return parsed, nil +} + +func optionalTime(value string) (*time.Time, error) { + if strings.TrimSpace(value) == "" { + return nil, nil + } + parsed, err := time.Parse(time.RFC3339, value) + if err != nil { + return nil, err + } + return &parsed, nil +} diff --git a/server/app/jobs/examples.go b/server/app/jobs/examples.go index 5937819..8036b15 100644 --- a/server/app/jobs/examples.go +++ b/server/app/jobs/examples.go @@ -4,6 +4,7 @@ import ( "fmt" "time" + "go-admin/app/goauto/shopeeproduct" "go-admin/app/goauto/sybimport" ) @@ -12,8 +13,10 @@ import ( // 字典 key 可以配置到 自动任务 调用目标 中; func InitJob() { jobList = map[string]JobExec{ - "ExamplesOne": ExamplesOne{}, - sybimport.HourlySyncInvokeTarget: sybimport.HourlySyncJob{}, + "ExamplesOne": ExamplesOne{}, + sybimport.HourlySyncInvokeTarget: sybimport.HourlySyncJob{}, + sybimport.SpecAIParseInvokeTarget: sybimport.ScheduledSpecAIParseJob{}, + shopeeproduct.SpecAutoMatchInvokeTarget: shopeeproduct.ScheduledAutoMatchJob{}, // ... } } diff --git a/server/app/jobs/execution_log.go b/server/app/jobs/execution_log.go new file mode 100644 index 0000000..6f305c6 --- /dev/null +++ b/server/app/jobs/execution_log.go @@ -0,0 +1,132 @@ +package jobs + +import ( + "context" + "errors" + "fmt" + "time" + + log "github.com/go-admin-team/go-admin-core/logger" + "github.com/google/uuid" + "gorm.io/gorm" + + "go-admin/app/jobs/models" +) + +const ( + executionErrorTargetMissing = "JOB_TARGET_NOT_FOUND" + executionErrorExecFailed = "JOB_EXECUTION_FAILED" + executionErrorHTTPFailed = "JOB_HTTP_FAILED" + executionErrorPanicked = "JOB_EXECUTION_PANICKED" + executionErrorInterrupted = "JOB_INTERRUPTED" +) + +type executionFailure struct { + code string + message string + cause error +} + +func (failure *executionFailure) Error() string { + if failure.cause != nil { + return failure.cause.Error() + } + return failure.message +} + +func newExecutionFailure(code, message string, cause error) error { + return &executionFailure{code: code, message: message, cause: cause} +} + +func runWithExecutionLog(db *gorm.DB, core JobCore, execute func() error) (executionErr error) { + startedAt := time.Now().UTC() + record := &models.SysJobExecutionLog{ + ExecutionID: uuid.NewString(), JobID: core.JobId, + JobNameSnapshot: core.Name, InvokeTargetSnapshot: core.InvokeTarget, + TriggerType: models.JobTriggerScheduled, Status: models.JobExecutionRunning, + StartedAt: startedAt, + } + created := false + if db != nil { + if err := db.WithContext(context.Background()).Create(record).Error; err != nil { + log.Errorf("[Job] execution log create failed job_id=%d: %v", core.JobId, err) + } else { + created = true + } + } + + defer func() { + if recovered := recover(); recovered != nil { + executionErr = newExecutionFailure(executionErrorPanicked, "任务执行异常中断", nil) + if created { + finishExecutionLog(db, core.JobId, record, startedAt, executionErr) + } + return + } + if created { + finishExecutionLog(db, core.JobId, record, startedAt, executionErr) + } + }() + executionErr = execute() + return executionErr +} + +func finishExecutionLog(db *gorm.DB, jobID int, record *models.SysJobExecutionLog, startedAt time.Time, executionErr error) { + finishedAt := time.Now().UTC() + updates := map[string]any{ + "status": models.JobExecutionSucceeded, "finished_at": finishedAt, + "duration_ms": finishedAt.Sub(startedAt).Milliseconds(), "error_code": "", "error_message": "", + } + if executionErr != nil { + code, message := publicExecutionFailure(executionErr) + updates["status"] = models.JobExecutionFailed + updates["error_code"] = code + updates["error_message"] = message + } + if err := db.WithContext(context.Background()).Model(&models.SysJobExecutionLog{}). + Where("id = ? AND status = ?", record.ID, models.JobExecutionRunning).Updates(updates).Error; err != nil { + log.Errorf("[Job] execution log finish failed job_id=%d execution_id=%s: %v", jobID, record.ExecutionID, err) + } +} + +func publicExecutionFailure(err error) (string, string) { + var failure *executionFailure + if errors.As(err, &failure) { + return failure.code, failure.message + } + return executionErrorExecFailed, "任务执行失败,请查看受控服务日志" +} + +// RecoverInterruptedExecutionLogs closes invocations left running by the +// previous process. Production currently runs one scheduler per database. +func RecoverInterruptedExecutionLogs(db *gorm.DB) error { + if db == nil { + return nil + } + now := time.Now().UTC() + var records []models.SysJobExecutionLog + if err := db.WithContext(context.Background()).Where("status = ?", models.JobExecutionRunning).Find(&records).Error; err != nil { + return err + } + return db.WithContext(context.Background()).Transaction(func(tx *gorm.DB) error { + for _, record := range records { + duration := now.Sub(record.StartedAt).Milliseconds() + if duration < 0 { + duration = 0 + } + if err := tx.Model(&models.SysJobExecutionLog{}). + Where("id = ? AND status = ?", record.ID, models.JobExecutionRunning).Updates(map[string]any{ + "status": models.JobExecutionInterrupted, "finished_at": now, + "duration_ms": duration, "error_code": executionErrorInterrupted, + "error_message": "服务重启前任务未完成", + }).Error; err != nil { + return err + } + } + return nil + }) +} + +func missingExecutionTarget(target string) error { + return newExecutionFailure(executionErrorTargetMissing, "任务调用目标未注册", fmt.Errorf("job target %q is not registered", target)) +} diff --git a/server/app/jobs/execution_log_test.go b/server/app/jobs/execution_log_test.go new file mode 100644 index 0000000..51bf61d --- /dev/null +++ b/server/app/jobs/execution_log_test.go @@ -0,0 +1,113 @@ +package jobs + +import ( + "errors" + "testing" + "time" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" + + "go-admin/app/jobs/models" +) + +func jobExecutionTestDB(t *testing.T) *gorm.DB { + t.Helper() + db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&models.SysJobExecutionLog{}); err != nil { + t.Fatal(err) + } + return db +} + +func TestRunWithExecutionLogRecordsSuccessAndSanitizedFailure(t *testing.T) { + db := jobExecutionTestDB(t) + core := JobCore{JobId: 7, Name: "测试任务", InvokeTarget: "TestTarget"} + if err := runWithExecutionLog(db, core, func() error { return nil }); err != nil { + t.Fatal(err) + } + rawSecret := "token=secret-value productSpec=private" + if err := runWithExecutionLog(db, core, func() error { return errors.New(rawSecret) }); err == nil { + t.Fatal("failed execution must return its error") + } + var records []models.SysJobExecutionLog + if err := db.Order("id asc").Find(&records).Error; err != nil { + t.Fatal(err) + } + if len(records) != 2 || records[0].Status != models.JobExecutionSucceeded || records[1].Status != models.JobExecutionFailed { + t.Fatalf("unexpected records: %+v", records) + } + if records[1].ErrorCode != executionErrorExecFailed || records[1].ErrorMessage == rawSecret || records[1].ErrorMessage == "" { + t.Fatalf("failure was not safely summarized: %+v", records[1]) + } + if records[0].FinishedAt == nil || records[1].FinishedAt == nil || records[0].ExecutionID == records[1].ExecutionID { + t.Fatal("execution lifecycle or unique IDs were not recorded") + } +} + +func TestRecoverInterruptedExecutionLogs(t *testing.T) { + db := jobExecutionTestDB(t) + started := time.Now().UTC().Add(-2 * time.Second) + record := models.SysJobExecutionLog{ + ExecutionID: "00000000-0000-4000-8000-000000000010", JobID: 9, + JobNameSnapshot: "中断任务", InvokeTargetSnapshot: "Interrupted", + TriggerType: models.JobTriggerScheduled, Status: models.JobExecutionRunning, StartedAt: started, + } + if err := db.Create(&record).Error; err != nil { + t.Fatal(err) + } + if err := RecoverInterruptedExecutionLogs(db); err != nil { + t.Fatal(err) + } + if err := db.First(&record, record.ID).Error; err != nil { + t.Fatal(err) + } + if record.Status != models.JobExecutionInterrupted || record.FinishedAt == nil || record.ErrorCode != executionErrorInterrupted || record.DurationMS < 1000 { + t.Fatalf("record was not safely interrupted: %+v", record) + } +} + +func TestMissingExecutionTargetHasPublicSafeMessage(t *testing.T) { + code, message := publicExecutionFailure(missingExecutionTarget("SecretTarget")) + if code != executionErrorTargetMissing || message != "任务调用目标未注册" { + t.Fatalf("unexpected target failure %s %s", code, message) + } +} + +func TestExecJobRecordsMissingTargetAsFailure(t *testing.T) { + db := jobExecutionTestDB(t) + previousJobList := jobList + jobList = map[string]JobExec{} + t.Cleanup(func() { jobList = previousJobList }) + + job := &ExecJob{JobCore: JobCore{JobId: 10, Name: "未注册任务", InvokeTarget: "MissingTarget"}, DB: db} + job.Run() + + var record models.SysJobExecutionLog + if err := db.First(&record).Error; err != nil { + t.Fatal(err) + } + if record.Status != models.JobExecutionFailed || record.ErrorCode != executionErrorTargetMissing || record.ErrorMessage != "任务调用目标未注册" { + t.Fatalf("missing target was not safely recorded: %+v", record) + } +} + +func TestRunWithExecutionLogSafelyRecordsPanic(t *testing.T) { + db := jobExecutionTestDB(t) + err := runWithExecutionLog(db, JobCore{JobId: 11, Name: "异常任务", InvokeTarget: "Panic"}, func() error { + panic("provider-secret") + }) + if code, message := publicExecutionFailure(err); code != executionErrorPanicked || message != "任务执行异常中断" { + t.Fatalf("panic was not returned as a safe failure: %s %s", code, message) + } + var record models.SysJobExecutionLog + if err := db.First(&record).Error; err != nil { + t.Fatal(err) + } + if record.Status != models.JobExecutionFailed || record.ErrorCode != executionErrorPanicked || record.ErrorMessage != "任务执行异常中断" { + t.Fatalf("panic was not safely persisted: %+v", record) + } +} diff --git a/server/app/jobs/jobbase.go b/server/app/jobs/jobbase.go index 3b3a613..f2020b8 100644 --- a/server/app/jobs/jobbase.go +++ b/server/app/jobs/jobbase.go @@ -33,6 +33,7 @@ type JobCore struct { // HttpJob 任务类型 http type HttpJob struct { JobCore + DB *gorm.DB } type ExecJob struct { @@ -42,14 +43,16 @@ type ExecJob struct { func (e *ExecJob) Run() { startTime := time.Now() - var obj = jobList[e.InvokeTarget] - if obj == nil { - log.Warn("[Job] ExecJob Run job nil") - return - } - err := CallExecWithDB(obj.(JobExec), e.DB, e.Args) + err := runWithExecutionLog(e.DB, e.JobCore, func() error { + obj := jobList[e.InvokeTarget] + if obj == nil { + return missingExecutionTarget(e.InvokeTarget) + } + return CallExecWithDB(obj, e.DB, e.Args) + }) if err != nil { - log.Errorf("[Job] JobCore %s failed: %v", e.Name, err) + code, _ := publicExecutionFailure(err) + log.Errorf("[Job] JobCore %s failed error_code=%s", e.Name, code) return } // 结束时间 @@ -68,22 +71,26 @@ func (e *ExecJob) Run() { func (h *HttpJob) Run() { startTime := time.Now() - var count = 0 - var err error - var str string - /* 循环 */ -LOOP: - if count < retryCount { - /* 跳过迭代 */ - str, err = pkg.Get(h.InvokeTarget) - if err != nil { - // 如果失败暂停一段时间重试 - log.Warnf("[Job] mission failed! %v", err) - log.Warnf("[Job] Retry after the task fails %d seconds! %s \n", (count+1)*5, str) - time.Sleep(time.Duration(count+1) * 5 * time.Second) - count = count + 1 - goto LOOP + err := runWithExecutionLog(h.DB, h.JobCore, func() error { + var lastErr error + for count := 0; count < retryCount; count++ { + _, requestErr := pkg.Get(h.InvokeTarget) + if requestErr == nil { + return nil + } + lastErr = requestErr + log.Warnf("[Job] HTTP mission failed attempt=%d", count+1) + if count+1 < retryCount { + log.Warnf("[Job] Retry after the task fails %d seconds!\n", (count+1)*5) + time.Sleep(time.Duration(count+1) * 5 * time.Second) + } } + return newExecutionFailure(executionErrorHTTPFailed, "HTTP 任务请求失败", lastErr) + }) + if err != nil { + code, _ := publicExecutionFailure(err) + log.Errorf("[Job] JobCore %s failed error_code=%s", h.Name, code) + return } // 结束时间 endTime := time.Now() @@ -102,6 +109,9 @@ func Setup(dbs map[string]*gorm.DB) { fmt.Println(time.Now().Format(timeFormat), " [INFO] JobCore Starting...") for k, db := range dbs { + if err := RecoverInterruptedExecutionLogs(db); err != nil { + log.Errorf("[Job] recover interrupted execution logs failed: %v", err) + } sdk.Runtime.SetCrontab(k, cronjob.NewWithSeconds()) setup(k, db) } @@ -127,6 +137,7 @@ func setup(key string, db *gorm.DB) { for i := 0; i < len(jobList); i++ { if jobList[i].JobType == 1 { j := &HttpJob{} + j.DB = db j.InvokeTarget = jobList[i].InvokeTarget j.CronExpression = jobList[i].CronExpression j.JobId = jobList[i].JobId diff --git a/server/app/jobs/models/sys_job_execution_log.go b/server/app/jobs/models/sys_job_execution_log.go new file mode 100644 index 0000000..2b79d01 --- /dev/null +++ b/server/app/jobs/models/sys_job_execution_log.go @@ -0,0 +1,32 @@ +package models + +import "time" + +const ( + JobExecutionRunning = "running" + JobExecutionSucceeded = "succeeded" + JobExecutionFailed = "failed" + JobExecutionInterrupted = "interrupted" + JobTriggerScheduled = "scheduled" +) + +// SysJobExecutionLog stores one scheduler invocation. Job arguments and raw +// provider responses are deliberately excluded from this audit record. +type SysJobExecutionLog struct { + ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"` + ExecutionID string `json:"executionId" gorm:"size:36;not null;uniqueIndex:ux_sys_job_execution_id"` + JobID int `json:"jobId" gorm:"not null;index:idx_sys_job_execution_job_started,priority:1"` + JobNameSnapshot string `json:"jobName" gorm:"size:255;not null"` + InvokeTargetSnapshot string `json:"invokeTarget" gorm:"size:255;not null"` + TriggerType string `json:"triggerType" gorm:"size:16;not null"` + Status string `json:"status" gorm:"size:16;not null;index:idx_sys_job_execution_status_started,priority:1"` + StartedAt time.Time `json:"startedAt" gorm:"not null;index:idx_sys_job_execution_job_started,priority:2;index:idx_sys_job_execution_status_started,priority:2"` + FinishedAt *time.Time `json:"finishedAt,omitempty"` + DurationMS int64 `json:"durationMs" gorm:"not null;default:0"` + ErrorCode string `json:"errorCode,omitempty" gorm:"size:64;not null;default:''"` + ErrorMessage string `json:"errorMessage,omitempty" gorm:"size:500;not null;default:''"` + CreatedAt time.Time `json:"createdAt"` + UpdatedAt time.Time `json:"updatedAt"` +} + +func (SysJobExecutionLog) TableName() string { return "sys_job_execution_log" } diff --git a/server/app/jobs/router/sys_job.go b/server/app/jobs/router/sys_job.go index 89723d4..6bddaab 100644 --- a/server/app/jobs/router/sys_job.go +++ b/server/app/jobs/router/sys_job.go @@ -24,6 +24,8 @@ func registerSysJobRouter(v1 *gin.RouterGroup, authMiddleware *jwt.GinJWTMiddlew list := make([]models2.SysJob, 0) return &list })) + jobAPI := apis.SysJob{} + r.GET("/:id/execution-logs", actions.PermissionAction(), jobAPI.ListExecutionLogs) r.GET("/:id", actions.PermissionAction(), actions.ViewAction(new(dto2.SysJobById), func() interface{} { return &dto2.SysJobItem{} })) diff --git a/server/app/jobs/service/execution_log.go b/server/app/jobs/service/execution_log.go new file mode 100644 index 0000000..69d76d0 --- /dev/null +++ b/server/app/jobs/service/execution_log.go @@ -0,0 +1,107 @@ +package service + +import ( + "context" + "errors" + "fmt" + "time" + + "gorm.io/gorm" + + "go-admin/app/jobs/models" +) + +var ( + ErrExecutionLogInvalidRequest = errors.New("invalid execution log request") + ErrExecutionLogJobNotFound = errors.New("scheduled job not found") +) + +type ExecutionLogListRequest struct { + Page int + PageSize int + Status string + StartedFrom *time.Time + StartedTo *time.Time +} + +type ExecutionLogJob struct { + JobID int `json:"jobId"` + JobName string `json:"jobName"` + InvokeTarget string `json:"invokeTarget"` + Deleted bool `json:"deleted"` +} + +type ExecutionLogListResponse struct { + Job ExecutionLogJob `json:"job"` + Items []models.SysJobExecutionLog `json:"items"` + Total int64 `json:"total"` + Page int `json:"page"` + PageSize int `json:"pageSize"` +} + +type ExecutionLogService struct{ db *gorm.DB } + +func NewExecutionLogService(db *gorm.DB) *ExecutionLogService { + return &ExecutionLogService{db: db} +} + +func (service *ExecutionLogService) List(ctx context.Context, jobID int, request ExecutionLogListRequest) (ExecutionLogListResponse, error) { + if service.db == nil || jobID < 1 { + return ExecutionLogListResponse{}, fmt.Errorf("%w: jobId 无效", ErrExecutionLogInvalidRequest) + } + if request.Page < 1 { + request.Page = 1 + } + if request.PageSize < 1 { + request.PageSize = 20 + } + if request.PageSize > 100 { + return ExecutionLogListResponse{}, fmt.Errorf("%w: pageSize 必须是 1 到 100 的整数", ErrExecutionLogInvalidRequest) + } + if request.Status != "" && !validExecutionStatus(request.Status) { + return ExecutionLogListResponse{}, fmt.Errorf("%w: status 无效", ErrExecutionLogInvalidRequest) + } + if request.StartedFrom != nil && request.StartedTo != nil && request.StartedFrom.After(*request.StartedTo) { + return ExecutionLogListResponse{}, fmt.Errorf("%w: 开始时间范围无效", ErrExecutionLogInvalidRequest) + } + + var job models.SysJob + if err := service.db.WithContext(ctx).Unscoped().First(&job, jobID).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return ExecutionLogListResponse{}, ErrExecutionLogJobNotFound + } + return ExecutionLogListResponse{}, err + } + + query := service.db.WithContext(ctx).Model(&models.SysJobExecutionLog{}).Where("job_id = ?", jobID) + if request.Status != "" { + query = query.Where("status = ?", request.Status) + } + if request.StartedFrom != nil { + query = query.Where("started_at >= ?", request.StartedFrom.UTC()) + } + if request.StartedTo != nil { + query = query.Where("started_at <= ?", request.StartedTo.UTC()) + } + var total int64 + if err := query.Count(&total).Error; err != nil { + return ExecutionLogListResponse{}, err + } + items := make([]models.SysJobExecutionLog, 0, request.PageSize) + if err := query.Order("started_at DESC, id DESC").Offset((request.Page - 1) * request.PageSize).Limit(request.PageSize).Find(&items).Error; err != nil { + return ExecutionLogListResponse{}, err + } + return ExecutionLogListResponse{ + Job: ExecutionLogJob{JobID: job.JobId, JobName: job.JobName, InvokeTarget: job.InvokeTarget, Deleted: job.DeletedAt.Valid}, + Items: items, Total: total, Page: request.Page, PageSize: request.PageSize, + }, nil +} + +func validExecutionStatus(status string) bool { + switch status { + case models.JobExecutionRunning, models.JobExecutionSucceeded, models.JobExecutionFailed, models.JobExecutionInterrupted: + return true + default: + return false + } +} diff --git a/server/app/jobs/service/execution_log_test.go b/server/app/jobs/service/execution_log_test.go new file mode 100644 index 0000000..ba832a5 --- /dev/null +++ b/server/app/jobs/service/execution_log_test.go @@ -0,0 +1,87 @@ +package service + +import ( + "context" + "errors" + "testing" + "time" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" + + "go-admin/app/jobs/models" +) + +func executionLogTestDB(t *testing.T) *gorm.DB { + t.Helper() + db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&models.SysJob{}, &models.SysJobExecutionLog{}); err != nil { + t.Fatal(err) + } + return db +} + +func TestExecutionLogListFiltersPaginatesAndKeepsDeletedJob(t *testing.T) { + db := executionLogTestDB(t) + job := models.SysJob{JobName: "测试任务", InvokeTarget: "TestTarget"} + if err := db.Create(&job).Error; err != nil { + t.Fatal(err) + } + base := time.Date(2026, 9, 2, 10, 0, 0, 0, time.UTC) + records := []models.SysJobExecutionLog{ + {ExecutionID: "00000000-0000-4000-8000-000000000001", JobID: job.JobId, JobNameSnapshot: job.JobName, InvokeTargetSnapshot: job.InvokeTarget, TriggerType: models.JobTriggerScheduled, Status: models.JobExecutionSucceeded, StartedAt: base}, + {ExecutionID: "00000000-0000-4000-8000-000000000002", JobID: job.JobId, JobNameSnapshot: job.JobName, InvokeTargetSnapshot: job.InvokeTarget, TriggerType: models.JobTriggerScheduled, Status: models.JobExecutionFailed, StartedAt: base.Add(time.Hour)}, + {ExecutionID: "00000000-0000-4000-8000-000000000003", JobID: job.JobId, JobNameSnapshot: job.JobName, InvokeTargetSnapshot: job.InvokeTarget, TriggerType: models.JobTriggerScheduled, Status: models.JobExecutionFailed, StartedAt: base.Add(2 * time.Hour)}, + } + if err := db.Create(&records).Error; err != nil { + t.Fatal(err) + } + from, to := base.Add(30*time.Minute), base.Add(3*time.Hour) + result, err := NewExecutionLogService(db).List(context.Background(), job.JobId, ExecutionLogListRequest{ + Page: 1, PageSize: 1, Status: models.JobExecutionFailed, StartedFrom: &from, StartedTo: &to, + }) + if err != nil { + t.Fatal(err) + } + if result.Total != 2 || len(result.Items) != 1 || result.Items[0].ExecutionID != records[2].ExecutionID { + t.Fatalf("unexpected filtered page: %+v", result) + } + if err := db.Delete(&job).Error; err != nil { + t.Fatal(err) + } + deleted, err := NewExecutionLogService(db).List(context.Background(), job.JobId, ExecutionLogListRequest{Page: 1, PageSize: 20}) + if err != nil { + t.Fatal(err) + } + if !deleted.Job.Deleted || deleted.Job.JobName != job.JobName || deleted.Total != 3 { + t.Fatalf("deleted job history unavailable: %+v", deleted) + } +} + +func TestExecutionLogListValidationAndNotFound(t *testing.T) { + db := executionLogTestDB(t) + service := NewExecutionLogService(db) + if _, err := service.List(context.Background(), 0, ExecutionLogListRequest{}); !errors.Is(err, ErrExecutionLogInvalidRequest) { + t.Fatalf("invalid job id error = %v", err) + } + if _, err := service.List(context.Background(), 1, ExecutionLogListRequest{PageSize: 101}); !errors.Is(err, ErrExecutionLogInvalidRequest) { + t.Fatalf("invalid page size error = %v", err) + } + job := models.SysJob{JobName: "测试任务", InvokeTarget: "TestTarget"} + if err := db.Create(&job).Error; err != nil { + t.Fatal(err) + } + if _, err := service.List(context.Background(), job.JobId, ExecutionLogListRequest{Status: "unknown"}); !errors.Is(err, ErrExecutionLogInvalidRequest) { + t.Fatalf("invalid status error = %v", err) + } + from, to := time.Now(), time.Now().Add(-time.Hour) + if _, err := service.List(context.Background(), job.JobId, ExecutionLogListRequest{StartedFrom: &from, StartedTo: &to}); !errors.Is(err, ErrExecutionLogInvalidRequest) { + t.Fatalf("invalid range error = %v", err) + } + if _, err := service.List(context.Background(), 999, ExecutionLogListRequest{}); !errors.Is(err, ErrExecutionLogJobNotFound) { + t.Fatalf("not found error = %v", err) + } +} diff --git a/server/app/jobs/service/sys_job.go b/server/app/jobs/service/sys_job.go index 356ea00..03d72fd 100644 --- a/server/app/jobs/service/sys_job.go +++ b/server/app/jobs/service/sys_job.go @@ -62,6 +62,7 @@ func (e *SysJob) StartJob(c *dto.GeneralGetDto) error { if data.JobType == 1 { var j = &jobs.HttpJob{} + j.DB = e.Orm.WithContext(context.Background()) j.InvokeTarget = data.InvokeTarget j.CronExpression = data.CronExpression j.JobId = data.JobId diff --git a/server/cmd/api/server.go b/server/cmd/api/server.go index e6c1613..3266e0c 100644 --- a/server/cmd/api/server.go +++ b/server/cmd/api/server.go @@ -54,6 +54,12 @@ var ( } ) +// The synchronous Admin AI endpoints allow a provider timeout of up to 600 +// seconds and the browser waits 610 seconds. Keep the HTTP server alive a +// little longer so it can return the domain response instead of truncating +// the connection and surfacing a proxy-level 502. +const minimumAPIWriteTimeout = 620 * time.Second + var AppRouters = make([]func(), 0) func init() { @@ -125,11 +131,15 @@ func run() error { ) } + writeTimeout, err := validatedAPIWriteTimeout(config.ApplicationConfig.WriterTimeout) + if err != nil { + return err + } srv := &http.Server{ Addr: fmt.Sprintf("%s:%d", config.ApplicationConfig.Host, config.ApplicationConfig.Port), Handler: sdk.Runtime.GetEngine(), ReadTimeout: time.Duration(config.ApplicationConfig.ReadTimeout) * time.Second, - WriteTimeout: time.Duration(config.ApplicationConfig.WriterTimeout) * time.Second, + WriteTimeout: writeTimeout, } go func() { @@ -193,6 +203,14 @@ func run() error { return nil } +func validatedAPIWriteTimeout(seconds int) (time.Duration, error) { + timeout := time.Duration(seconds) * time.Second + if timeout < minimumAPIWriteTimeout { + return 0, fmt.Errorf("application writetimeout must be at least %s for synchronous AI requests", minimumAPIWriteTimeout) + } + return timeout, nil +} + type policyLoader interface { LoadPolicy() error } diff --git a/server/cmd/api/server_timeout_test.go b/server/cmd/api/server_timeout_test.go new file mode 100644 index 0000000..5034b6e --- /dev/null +++ b/server/cmd/api/server_timeout_test.go @@ -0,0 +1,19 @@ +package api + +import ( + "testing" + "time" +) + +func TestValidatedAPIWriteTimeoutProtectsSynchronousAIRequests(t *testing.T) { + if _, err := validatedAPIWriteTimeout(2); err == nil { + t.Fatal("two-second write timeout must be rejected") + } + got, err := validatedAPIWriteTimeout(620) + if err != nil { + t.Fatalf("620-second write timeout should be accepted: %v", err) + } + if got != 620*time.Second { + t.Fatalf("write timeout = %s, want 620s", got) + } +} diff --git a/server/cmd/migrate/migration/version-local/1788290000000_shopee_spec_auto_match.go b/server/cmd/migrate/migration/version-local/1788290000000_shopee_spec_auto_match.go new file mode 100644 index 0000000..4c50b1a --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1788290000000_shopee_spec_auto_match.go @@ -0,0 +1,47 @@ +package version_local + +import ( + "errors" + "runtime" + + goautomigrations "go-admin/app/goauto/migrations" + "go-admin/app/goauto/shopeeproduct" + jobsmodels "go-admin/app/jobs/models" + "go-admin/cmd/migrate/migration" + common "go-admin/common/models" + + "gorm.io/gorm" +) + +func init() { + _, fileName, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateShopeeSpecAutoMatch) +} + +func migrateShopeeSpecAutoMatch(db *gorm.DB, version string) error { + return db.Transaction(func(tx *gorm.DB) error { + if err := goautomigrations.Migrate(tx); err != nil { + return err + } + if err := ensureShopeeSpecAutoMatchJob(tx); err != nil { + return err + } + return tx.Create(&common.Migration{Version: version}).Error + }) +} + +func ensureShopeeSpecAutoMatchJob(db *gorm.DB) error { + var existing jobsmodels.SysJob + err := db.Where("invoke_target = ?", shopeeproduct.SpecAutoMatchInvokeTarget).First(&existing).Error + if err == nil { + return nil + } + if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + return db.Create(&jobsmodels.SysJob{ + JobName: "蝦皮规格自动匹配", JobGroup: "GoAuto", JobType: 2, + CronExpression: "0 15 * * * *", InvokeTarget: shopeeproduct.SpecAutoMatchInvokeTarget, + Args: `{"batchLimit":20}`, MisfirePolicy: 1, Concurrent: 1, Status: 1, + }).Error +} diff --git a/server/cmd/migrate/migration/version-local/1788290000000_shopee_spec_auto_match_test.go b/server/cmd/migrate/migration/version-local/1788290000000_shopee_spec_auto_match_test.go new file mode 100644 index 0000000..7a2b28e --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1788290000000_shopee_spec_auto_match_test.go @@ -0,0 +1,50 @@ +package version_local + +import ( + "testing" + + "go-admin/app/goauto/shopeeproduct" + jobsmodels "go-admin/app/jobs/models" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +func TestEnsureShopeeSpecAutoMatchJobIsDisabledIdempotentAndPreservesChanges(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:shopee-spec-auto-match-job?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&jobsmodels.SysJob{}); err != nil { + t.Fatal(err) + } + if err := ensureShopeeSpecAutoMatchJob(db); err != nil { + t.Fatal(err) + } + var job jobsmodels.SysJob + if err := db.Where("invoke_target = ?", shopeeproduct.SpecAutoMatchInvokeTarget).First(&job).Error; err != nil { + t.Fatal(err) + } + if job.Status != 1 || job.CronExpression != "0 15 * * * *" || job.Args != `{"batchLimit":20}` { + t.Fatalf("unexpected seed: %+v", job) + } + if err := db.Model(&job).Updates(map[string]any{"status": 2, "cron_expression": "0 30 * * * *", "args": `{"batchLimit":5}`}).Error; err != nil { + t.Fatal(err) + } + if err := ensureShopeeSpecAutoMatchJob(db); err != nil { + t.Fatal(err) + } + var count int64 + if err := db.Model(&jobsmodels.SysJob{}).Where("invoke_target = ?", shopeeproduct.SpecAutoMatchInvokeTarget).Count(&count).Error; err != nil { + t.Fatal(err) + } + if count != 1 { + t.Fatalf("count=%d", count) + } + if err := db.Where("invoke_target = ?", shopeeproduct.SpecAutoMatchInvokeTarget).First(&job).Error; err != nil { + t.Fatal(err) + } + if job.Status != 2 || job.CronExpression != "0 30 * * * *" || job.Args != `{"batchLimit":5}` { + t.Fatalf("admin changes overwritten: %+v", job) + } +} diff --git a/server/cmd/migrate/migration/version-local/1788354164329_syb_spec_ai_parse.go b/server/cmd/migrate/migration/version-local/1788354164329_syb_spec_ai_parse.go new file mode 100644 index 0000000..40936df --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1788354164329_syb_spec_ai_parse.go @@ -0,0 +1,47 @@ +package version_local + +import ( + "errors" + "runtime" + + goautomigrations "go-admin/app/goauto/migrations" + "go-admin/app/goauto/sybimport" + jobsmodels "go-admin/app/jobs/models" + "go-admin/cmd/migrate/migration" + common "go-admin/common/models" + + "gorm.io/gorm" +) + +func init() { + _, fileName, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateSYBSpecAIParse) +} + +func migrateSYBSpecAIParse(db *gorm.DB, version string) error { + return db.Transaction(func(tx *gorm.DB) error { + if err := goautomigrations.Migrate(tx); err != nil { + return err + } + if err := ensureSYBSpecAIParseJob(tx); err != nil { + return err + } + return tx.Create(&common.Migration{Version: version}).Error + }) +} + +func ensureSYBSpecAIParseJob(db *gorm.DB) error { + var existing jobsmodels.SysJob + err := db.Where("invoke_target = ?", sybimport.SpecAIParseInvokeTarget).First(&existing).Error + if err == nil { + return nil + } + if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + return db.Create(&jobsmodels.SysJob{ + JobName: "SYB 异常规格 AI 解析", JobGroup: "GoAuto", JobType: 2, + CronExpression: "0 5 * * * *", InvokeTarget: sybimport.SpecAIParseInvokeTarget, + Args: `{"batchLimit":20}`, MisfirePolicy: 1, Concurrent: 1, Status: 1, + }).Error +} diff --git a/server/cmd/migrate/migration/version-local/1788354164329_syb_spec_ai_parse_test.go b/server/cmd/migrate/migration/version-local/1788354164329_syb_spec_ai_parse_test.go new file mode 100644 index 0000000..f77c2f1 --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1788354164329_syb_spec_ai_parse_test.go @@ -0,0 +1,50 @@ +package version_local + +import ( + "testing" + + "go-admin/app/goauto/sybimport" + jobsmodels "go-admin/app/jobs/models" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +func TestEnsureSYBSpecAIParseJobIsDisabledIdempotentAndPreservesChanges(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&jobsmodels.SysJob{}); err != nil { + t.Fatal(err) + } + if err := ensureSYBSpecAIParseJob(db); err != nil { + t.Fatal(err) + } + var job jobsmodels.SysJob + if err := db.Where("invoke_target = ?", sybimport.SpecAIParseInvokeTarget).First(&job).Error; err != nil { + t.Fatal(err) + } + if job.Status != 1 || job.CronExpression != "0 5 * * * *" || job.Args != `{"batchLimit":20}` { + t.Fatalf("unexpected initial job: %+v", job) + } + if err := db.Model(&job).Updates(map[string]any{"status": 2, "cron_expression": "0 7 * * * *", "args": `{"batchLimit":7}`}).Error; err != nil { + t.Fatal(err) + } + if err := ensureSYBSpecAIParseJob(db); err != nil { + t.Fatal(err) + } + var count int64 + if err := db.Model(&jobsmodels.SysJob{}).Where("invoke_target = ?", sybimport.SpecAIParseInvokeTarget).Count(&count).Error; err != nil { + t.Fatal(err) + } + if count != 1 { + t.Fatalf("job count=%d, want 1", count) + } + if err := db.First(&job, job.JobId).Error; err != nil { + t.Fatal(err) + } + if job.Status != 2 || job.CronExpression != "0 7 * * * *" || job.Args != `{"batchLimit":7}` { + t.Fatalf("existing administrator changes were overwritten: %+v", job) + } +} diff --git a/server/cmd/migrate/migration/version-local/1788357000000_sys_job_execution_log.go b/server/cmd/migrate/migration/version-local/1788357000000_sys_job_execution_log.go new file mode 100644 index 0000000..81b93a9 --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1788357000000_sys_job_execution_log.go @@ -0,0 +1,68 @@ +package version_local + +import ( + "fmt" + "runtime" + + jobsmodels "go-admin/app/jobs/models" + "go-admin/cmd/migrate/migration" + migrationmodels "go-admin/cmd/migrate/migration/models" + common "go-admin/common/models" + + "gorm.io/gorm" +) + +const jobExecutionLogsAPIPath = "/api/v1/sysjob/:id/execution-logs" + +func init() { + _, fileName, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateSysJobExecutionLog) +} + +func migrateSysJobExecutionLog(db *gorm.DB, version string) error { + return db.Transaction(func(tx *gorm.DB) error { + if err := ensureSysJobExecutionLog(tx); err != nil { + return err + } + return tx.Create(&common.Migration{Version: version}).Error + }) +} + +func ensureSysJobExecutionLog(db *gorm.DB) error { + if err := db.AutoMigrate(&jobsmodels.SysJobExecutionLog{}); err != nil { + return err + } + + var logMenu migrationmodels.SysMenu + if err := db.Where("menu_name = ?", "JobLog").First(&logMenu).Error; err != nil { + return fmt.Errorf("find JobLog menu: %w", err) + } + api := migrationmodels.SysApi{} + if err := db.Where(migrationmodels.SysApi{Path: jobExecutionLogsAPIPath, Action: "GET"}). + Attrs(migrationmodels.SysApi{Title: "查询定时任务执行日志", Type: "BUS"}).FirstOrCreate(&api).Error; err != nil { + return err + } + if err := db.Model(&logMenu).Association("SysApi").Append(&api); err != nil { + return err + } + + var roleIDs []int + if err := db.Table("sys_role_menu").Where("menu_id = ?", logMenu.MenuId).Distinct().Pluck("role_id", &roleIDs).Error; err != nil { + return err + } + if len(roleIDs) == 0 { + return nil + } + var roles []migrationmodels.SysRole + if err := db.Where("role_id IN ?", roleIDs).Find(&roles).Error; err != nil { + return err + } + for _, role := range roles { + rule := purchaserCasbinRule{Ptype: "p", V0: role.RoleKey, V1: jobExecutionLogsAPIPath, V2: "GET"} + if err := db.Where("ptype = ? AND v0 = ? AND v1 = ? AND v2 = ?", rule.Ptype, rule.V0, rule.V1, rule.V2). + FirstOrCreate(&rule).Error; err != nil { + return err + } + } + return nil +} diff --git a/server/cmd/migrate/migration/version-local/1788357000000_sys_job_execution_log_test.go b/server/cmd/migrate/migration/version-local/1788357000000_sys_job_execution_log_test.go new file mode 100644 index 0000000..b6b363e --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1788357000000_sys_job_execution_log_test.go @@ -0,0 +1,78 @@ +package version_local + +import ( + "testing" + + jobsmodels "go-admin/app/jobs/models" + migrationmodels "go-admin/cmd/migrate/migration/models" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +func TestEnsureSysJobExecutionLogIsIdempotentAndGrantsOnlyBoundRoles(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate( + &migrationmodels.SysMenu{}, &migrationmodels.SysApi{}, &migrationmodels.SysRole{}, &purchaserCasbinRule{}, + ); err != nil { + t.Fatal(err) + } + menu := migrationmodels.SysMenu{MenuName: "JobLog", Title: "日志", Path: "/schedule/log", Component: "/schedule/log"} + if err := db.Create(&menu).Error; err != nil { + t.Fatal(err) + } + bound := migrationmodels.SysRole{RoleName: "日志管理员", RoleKey: "job_logger"} + unbound := migrationmodels.SysRole{RoleName: "无日志权限", RoleKey: "no_job_logs"} + if err := db.Create(&bound).Error; err != nil { + t.Fatal(err) + } + if err := db.Create(&unbound).Error; err != nil { + t.Fatal(err) + } + if err := db.Model(&bound).Association("SysMenu").Append(&menu); err != nil { + t.Fatal(err) + } + custom := purchaserCasbinRule{Ptype: "p", V0: bound.RoleKey, V1: "/custom", V2: "GET"} + if err := db.Create(&custom).Error; err != nil { + t.Fatal(err) + } + + if err := ensureSysJobExecutionLog(db); err != nil { + t.Fatal(err) + } + if err := ensureSysJobExecutionLog(db); err != nil { + t.Fatalf("migration helper must be repeatable: %v", err) + } + if !db.Migrator().HasTable(&jobsmodels.SysJobExecutionLog{}) { + t.Fatal("execution log table was not created") + } + var apiCount, menuAPI, boundPolicy, unboundPolicy, customPolicy int64 + db.Model(&migrationmodels.SysApi{}).Where("path = ? AND action = ?", jobExecutionLogsAPIPath, "GET").Count(&apiCount) + var api migrationmodels.SysApi + if err := db.Where("path = ? AND action = ?", jobExecutionLogsAPIPath, "GET").First(&api).Error; err != nil { + t.Fatal(err) + } + db.Table("sys_menu_api_rule").Where("sys_menu_menu_id = ? AND sys_api_id = ?", menu.MenuId, api.Id).Count(&menuAPI) + db.Model(&purchaserCasbinRule{}).Where("v0 = ? AND v1 = ? AND v2 = ?", bound.RoleKey, jobExecutionLogsAPIPath, "GET").Count(&boundPolicy) + db.Model(&purchaserCasbinRule{}).Where("v0 = ? AND v1 = ? AND v2 = ?", unbound.RoleKey, jobExecutionLogsAPIPath, "GET").Count(&unboundPolicy) + db.Model(&purchaserCasbinRule{}).Where("v0 = ? AND v1 = ?", bound.RoleKey, "/custom").Count(&customPolicy) + if apiCount != 1 || menuAPI != 1 || boundPolicy != 1 || unboundPolicy != 0 || customPolicy != 1 { + t.Fatalf("unexpected permission state api=%d menu=%d bound=%d unbound=%d custom=%d", apiCount, menuAPI, boundPolicy, unboundPolicy, customPolicy) + } +} + +func TestEnsureSysJobExecutionLogRequiresExistingMenu(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&migrationmodels.SysMenu{}, &migrationmodels.SysApi{}, &migrationmodels.SysRole{}, &purchaserCasbinRule{}); err != nil { + t.Fatal(err) + } + if err := ensureSysJobExecutionLog(db); err == nil { + t.Fatal("missing JobLog menu must fail closed") + } +} diff --git a/server/config/settings.full.yml b/server/config/settings.full.yml index bcca7e2..3722d2a 100644 --- a/server/config/settings.full.yml +++ b/server/config/settings.full.yml @@ -9,7 +9,7 @@ settings: # 端口号 port: 8000 # 服务端口号 readtimeout: 1 - writertimeout: 2 + writertimeout: 620 # 数据权限功能开关 enabledp: false ssl: diff --git a/server/config/settings.yml b/server/config/settings.yml index 0533615..9c0a32e 100644 --- a/server/config/settings.yml +++ b/server/config/settings.yml @@ -9,7 +9,7 @@ settings: # 端口号 port: 8000 # 服务端口号 readtimeout: 1 - writertimeout: 2 + writertimeout: 620 # 数据权限功能开关 enabledp: false logger: @@ -84,4 +84,4 @@ settings: # blockingTimeout: 5 # reclaimInterval: 1 locker: - redis: \ No newline at end of file + redis: diff --git a/web/src/api/goauto/shopee-products.js b/web/src/api/goauto/shopee-products.js index 77b0983..c6956e2 100644 --- a/web/src/api/goauto/shopee-products.js +++ b/web/src/api/goauto/shopee-products.js @@ -70,6 +70,14 @@ export function autoMatchShopeeSpecMappings(productId, data) { return request({ url: `/api/admin/v1/shopee-products/${productId}/specs/mapping/auto-match`, method: 'post', data, timeout: aiSuggestTimeoutMs }) } +export function startShopeeSpecAutoMatchRun(data) { + return request({ url: '/api/admin/v1/shopee-spec-auto-match/runs', method: 'post', data }) +} + +export function getLatestShopeeSpecAutoMatchRun() { + return request({ url: '/api/admin/v1/shopee-spec-auto-match/runs/latest', method: 'get' }) +} + export function batchDeleteShopeeProducts(data) { return request({ url: '/api/admin/v1/shopee-products/batch-delete', method: 'post', data }) } diff --git a/web/src/api/job/sys-job.js b/web/src/api/job/sys-job.js index 83e5b5c..f30f8e4 100644 --- a/web/src/api/job/sys-job.js +++ b/web/src/api/job/sys-job.js @@ -17,6 +17,15 @@ export function getSysJob(jobId) { }) } +// 查询指定定时任务的持久化执行历史 +export function listJobExecutionLogs(jobId, query) { + return request({ + url: '/api/v1/sysjob/' + jobId + '/execution-logs', + method: 'get', + params: query + }) +} + // 新增SysJob export function addSysJob(data) { return request({ @@ -59,4 +68,3 @@ export function startJob(jobId) { method: 'get' }) } - diff --git a/web/src/views/goauto/shopee-products/index.vue b/web/src/views/goauto/shopee-products/index.vue index e30bdfc..2b2ac0f 100644 --- a/web/src/views/goauto/shopee-products/index.vue +++ b/web/src/views/goauto/shopee-products/index.vue @@ -1,7 +1,7 @@