diff --git a/backend-api/cmd/api/main.go b/backend-api/cmd/api/main.go index ed8db35..b889a10 100644 --- a/backend-api/cmd/api/main.go +++ b/backend-api/cmd/api/main.go @@ -14,7 +14,6 @@ import ( "cmroubao/backend-api/internal/config" "cmroubao/backend-api/internal/platform/assetstore" "cmroubao/backend-api/internal/platform/database" - "cmroubao/backend-api/internal/platform/erpconnector" "cmroubao/backend-api/internal/platform/migration" "cmroubao/backend-api/internal/platform/password" "cmroubao/backend-api/internal/platform/shunyunbao" @@ -219,18 +218,6 @@ func buildRouter( if err != nil { return nil, err } - erpTimeout := cfg.ERPConnectorTimeout - if erpTimeout <= 0 { - erpTimeout = 90 * time.Second - } - erpClient, err := erpconnector.New( - cfg.ERPConnectorURL, - cfg.ERPConnectorAPIKey, - erpTimeout, - ) - if err != nil { - return nil, err - } erpSession, err := shunyunbao.NewSessionManager(shunyunbao.SessionConfig{ BaseURL: cfg.ShunyunbaoURL, Username: cfg.ShunyunbaoUsername, @@ -242,10 +229,10 @@ func buildRouter( } freight, err := usecase.NewFreightService( store, - erpClient, + erpSession, clock, ids, - erpTimeout, + 90*time.Second, ) if err != nil { return nil, err diff --git a/backend-api/internal/platform/shunyunbao/session.go b/backend-api/internal/platform/shunyunbao/session.go index 4746bfd..bccc2de 100644 --- a/backend-api/internal/platform/shunyunbao/session.go +++ b/backend-api/internal/platform/shunyunbao/session.go @@ -212,6 +212,7 @@ func (manager *SessionManager) Login( LoginPath, payload, true, + false, ) if err != nil { return manager.statusLocked(), err @@ -251,6 +252,7 @@ func (manager *SessionManager) validateLocked(ctx context.Context) (any, error) UserPath, nil, false, + true, ) if err != nil { return nil, err @@ -266,6 +268,7 @@ func (manager *SessionManager) requestJSONLocked( method, path string, body []byte, loginRequest bool, + requireSession bool, ) (any, error) { var content io.Reader if body != nil { @@ -300,19 +303,28 @@ func (manager *SessionManager) requestJSONLocked( return nil, domain.ErrFreightSourceUnavailable } var envelope struct { - Status bool `json:"status"` + Status *bool `json:"status"` + Code json.RawMessage `json:"code"` Data json.RawMessage `json:"data"` } decoder := json.NewDecoder(bytes.NewReader(contentBytes)) - if err := decoder.Decode(&envelope); err != nil || len(envelope.Data) == 0 { + if err := decoder.Decode(&envelope); err != nil || envelope.Status == nil || + len(envelope.Data) == 0 { return nil, domain.ErrFreightSourceProtocol } - if !envelope.Status { + var extra any + if err := decoder.Decode(&extra); !errors.Is(err, io.EOF) { + return nil, domain.ErrFreightSourceProtocol + } + if !*envelope.Status { if loginRequest { return nil, ErrLoginRejected } - manager.authenticated = false - return nil, domain.ErrFreightSourceSessionNeeded + if requireSession || unauthenticatedCode(envelope.Code) { + manager.authenticated = false + return nil, domain.ErrFreightSourceSessionNeeded + } + return nil, domain.ErrFreightSourceProtocol } var data any dataDecoder := json.NewDecoder(bytes.NewReader(envelope.Data)) @@ -320,9 +332,29 @@ func (manager *SessionManager) requestJSONLocked( if err := dataDecoder.Decode(&data); err != nil { return nil, domain.ErrFreightSourceProtocol } + if err := dataDecoder.Decode(&extra); !errors.Is(err, io.EOF) { + return nil, domain.ErrFreightSourceProtocol + } return data, nil } +func unauthenticatedCode(raw json.RawMessage) bool { + var value any + decoder := json.NewDecoder(bytes.NewReader(raw)) + decoder.UseNumber() + if len(raw) == 0 || decoder.Decode(&value) != nil { + return false + } + switch typed := value.(type) { + case json.Number: + return typed == "-2" + case string: + return strings.TrimSpace(typed) == "-2" + default: + return false + } +} + func (manager *SessionManager) applyHeaders(request *http.Request) { for name, values := range manager.headers { request.Header[name] = append([]string(nil), values...) diff --git a/backend-api/internal/platform/shunyunbao/source.go b/backend-api/internal/platform/shunyunbao/source.go new file mode 100644 index 0000000..18f2680 --- /dev/null +++ b/backend-api/internal/platform/shunyunbao/source.go @@ -0,0 +1,295 @@ +package shunyunbao + +import ( + "context" + "encoding/json" + "reflect" + "strconv" + + "cmroubao/backend-api/internal/domain" +) + +const ( + sourcePageSize = 20 + sourceMaxMatches = 100 +) + +// QueryOrder implements the freight source port using the manager's single +// in-memory ERP session. It never exposes the ERP list or detail responses. +func (manager *SessionManager) QueryOrder( + ctx context.Context, + orderNumber string, +) (domain.FreightSourceBatch, error) { + query, err := OrderNumberQuery(orderNumber) + if err != nil { + return domain.FreightSourceBatch{}, domain.ErrFreightSourceProtocol + } + records, err := manager.queryRecords(ctx, query) + if err != nil { + return domain.FreightSourceBatch{}, err + } + if len(records) == 0 { + return domain.FreightSourceBatch{}, domain.ErrFreightSourceNotFound + } + return NormalizeFreightResult( + RawQuery{Mode: domain.FreightSyncOrderNumber}, + records, + ) +} + +// QueryCreatedRange implements the bounded Asia/Shanghai date-window source +// port. The caller owns larger-window splitting and watermark semantics. +func (manager *SessionManager) QueryCreatedRange( + ctx context.Context, + createdFrom, createdTo string, +) (domain.FreightSourceBatch, error) { + query, err := CreatedRangeQuery(createdFrom, createdTo) + if err != nil { + return domain.FreightSourceBatch{}, domain.ErrFreightSourceProtocol + } + records, err := manager.queryRecords(ctx, query) + if err != nil { + return domain.FreightSourceBatch{}, err + } + return NormalizeFreightResult( + RawQuery{ + Mode: domain.FreightSyncCreatedRange, + CreatedFrom: createdFrom, + CreatedTo: createdTo, + }, + records, + ) +} + +func (manager *SessionManager) queryRecords( + ctx context.Context, + query QueryCondition, +) ([]RawRecord, error) { + manager.mu.Lock() + defer manager.mu.Unlock() + if !manager.configuredLocked() { + return nil, domain.ErrFreightSourceNotConfigured + } + if !manager.authenticated { + return nil, domain.ErrFreightSourceSessionNeeded + } + if _, err := manager.validateLocked(ctx); err != nil { + return nil, err + } + stocks, err := manager.queryStocksLocked(ctx, query) + if err != nil { + return nil, err + } + if len(stocks) == 0 { + return []RawRecord{}, nil + } + details, err := manager.queryDetailsLocked(ctx, stocks) + if err != nil { + return nil, err + } + records := make([]RawRecord, 0, len(stocks)) + for _, stock := range stocks { + stockID, err := sourceExternalID(stock["id"]) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + detail, exists := details[stockID] + if !exists { + return nil, domain.ErrFreightSourceProtocol + } + records = append(records, RawRecord{Stock: stock, Detail: detail}) + } + return records, nil +} + +func (manager *SessionManager) queryStocksLocked( + ctx context.Context, + query QueryCondition, +) ([]map[string]any, error) { + firstPayload, err := StockListPayload(query, 0, 1, sourcePageSize) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + firstBody, err := json.Marshal(firstPayload) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + totalData, err := manager.requestJSONLocked( + ctx, + "POST", + StockListTotalPath, + firstBody, + false, + false, + ) + if err != nil { + return nil, err + } + total, err := sourceTotal(totalData) + if err != nil || total > sourceMaxMatches { + return nil, domain.ErrFreightSourceProtocol + } + if total == 0 { + return []map[string]any{}, nil + } + + rows := make([]map[string]any, 0, total) + for start := 0; start < total; start += sourcePageSize { + payload, err := StockListPayload( + query, + start, + start/sourcePageSize+1, + sourcePageSize, + ) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + body, err := json.Marshal(payload) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + data, err := manager.requestJSONLocked( + ctx, + "POST", + StockListPath, + body, + false, + false, + ) + if err != nil { + return nil, err + } + page, err := sourceObjectList(data) + if err != nil || len(page) > sourcePageSize { + return nil, domain.ErrFreightSourceProtocol + } + rows = append(rows, page...) + } + if len(rows) != total { + return nil, domain.ErrFreightSourceProtocol + } + stocksByID := make(map[string]map[string]any, len(rows)) + stocks := make([]map[string]any, 0, len(rows)) + for _, stock := range rows { + stockID, err := sourceExternalID(stock["id"]) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + if existing, exists := stocksByID[stockID]; exists { + if !reflect.DeepEqual(existing, stock) { + return nil, domain.ErrFreightSourceProtocol + } + continue + } + stocksByID[stockID] = stock + stocks = append(stocks, stock) + } + return stocks, nil +} + +func (manager *SessionManager) queryDetailsLocked( + ctx context.Context, + stocks []map[string]any, +) (map[string]map[string]any, error) { + ids := make([]string, 0, len(stocks)) + for _, stock := range stocks { + stockID, err := sourceExternalID(stock["id"]) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + ids = append(ids, stockID) + } + detailsByID := make(map[string]map[string]any, len(ids)) + requestedIDs := make(map[string]struct{}, len(ids)) + for _, id := range ids { + requestedIDs[id] = struct{}{} + } + for start := 0; start < len(ids); start += maxDetailIDs { + end := start + maxDetailIDs + if end > len(ids) { + end = len(ids) + } + payload, err := StockDetailPayload(ids[start:end]) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + body, err := json.Marshal(payload) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + data, err := manager.requestJSONLocked( + ctx, + "POST", + StockDetailPath+"?hist=0", + body, + false, + false, + ) + if err != nil { + return nil, err + } + values, err := sourceObjectList(data) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + for _, detail := range values { + detailID, err := sourceExternalID(detail["id"]) + if err != nil { + return nil, domain.ErrFreightSourceProtocol + } + if _, requested := requestedIDs[detailID]; !requested { + return nil, domain.ErrFreightSourceProtocol + } + if existing, exists := detailsByID[detailID]; exists && + !reflect.DeepEqual(existing, detail) { + return nil, domain.ErrFreightSourceProtocol + } + detailsByID[detailID] = detail + } + } + if len(detailsByID) != len(ids) { + return nil, domain.ErrFreightSourceProtocol + } + return detailsByID, nil +} + +func sourceTotal(value any) (int, error) { + var text string + switch typed := value.(type) { + case json.Number: + text = string(typed) + case string: + text = typed + default: + return 0, strconv.ErrSyntax + } + parsed, err := strconv.ParseInt(text, 10, 0) + if err != nil || parsed < 0 || parsed > int64(sourceMaxMatches) { + return 0, strconv.ErrSyntax + } + return int(parsed), nil +} + +func sourceObjectList(value any) ([]map[string]any, error) { + data, ok := value.(map[string]any) + if !ok { + return nil, strconv.ErrSyntax + } + rawList, ok := data["list"].([]any) + if !ok { + return nil, strconv.ErrSyntax + } + result := make([]map[string]any, 0, len(rawList)) + for _, value := range rawList { + row, ok := value.(map[string]any) + if !ok { + return nil, strconv.ErrSyntax + } + result = append(result, row) + } + return result, nil +} + +func sourceExternalID(value any) (string, error) { + return externalID(value) +} diff --git a/backend-api/internal/platform/shunyunbao/source_test.go b/backend-api/internal/platform/shunyunbao/source_test.go new file mode 100644 index 0000000..35f4694 --- /dev/null +++ b/backend-api/internal/platform/shunyunbao/source_test.go @@ -0,0 +1,261 @@ +package shunyunbao + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "net/http/httptest" + "strconv" + "strings" + "testing" + + "cmroubao/backend-api/internal/domain" +) + +func TestSessionManagerQueryOrderUsesVerifiedSessionAndAllowlist(t *testing.T) { + calls := make([]string, 0, 8) + server := httptest.NewServer(http.HandlerFunc(func( + writer http.ResponseWriter, + request *http.Request, + ) { + calls = append(calls, request.URL.Path) + switch request.URL.Path { + case CaptchaPath: + http.SetCookie(writer, &http.Cookie{Name: "captcha", Value: "ready", Path: "/"}) + writer.Header().Set("Content-Type", "image/png") + _, _ = writer.Write([]byte("captcha")) + case LoginPath: + assertERPHeaders(t, request) + if cookie, err := request.Cookie("captcha"); err != nil || cookie.Value != "ready" { + t.Fatalf("login captcha cookie = %v / %v", cookie, err) + } + http.SetCookie(writer, &http.Cookie{Name: "authenticated", Value: "yes", Path: "/"}) + _, _ = writer.Write([]byte(`{"status":true,"data":{"user":{"id":1}}}`)) + case UserPath: + if cookie, err := request.Cookie("authenticated"); err != nil || cookie.Value != "yes" { + t.Fatalf("user cookie = %v / %v", cookie, err) + } + _, _ = writer.Write([]byte(`{"status":true,"data":{"id":1}}`)) + case StockListTotalPath: + assertStockPayload(t, request, "SOURCE-12", 0, 1, 20) + _, _ = writer.Write([]byte(`{"status":true,"data":1}`)) + case StockListPath: + assertStockPayload(t, request, "SOURCE-12", 0, 1, 20) + _, _ = writer.Write([]byte(`{"status":true,"data":{"list":[{"id":12,"code":"SOURCE-12","orderCode":"PLATFORM-12","receiver":"private-recipient","receiverTel":"private-phone"}]}}`)) + case StockDetailPath: + if request.URL.Query().Get("hist") != "0" { + t.Fatalf("detail query = %q", request.URL.RawQuery) + } + assertDetailPayload(t, request, []uint64{12}) + _, _ = writer.Write([]byte(`{"status":true,"data":{"list":[{"id":12,"shopName":"测试店铺","created":"2026-07-28 08:00:00","details":[{"id":88,"productTitle":"商品一","productSpec":"黑色,L","sku":"BLACK-L","productQty":2,"receiverTel":"private-item-phone"},{"id":89,"productTitle":"商品二","productSpec":"黑色,XL","sku":"BLACK-XL","productQty":1}]}]}}`)) + default: + writer.WriteHeader(http.StatusNotFound) + } + })) + defer server.Close() + + manager := testSessionManager(t, server.URL, "test-user", "test-password") + loginForSource(t, manager) + result, err := manager.QueryOrder(context.Background(), "SOURCE-12") + if err != nil { + t.Fatalf("QueryOrder() error = %v", err) + } + if result.SchemaVersion != 1 || result.Query.Mode != domain.FreightSyncOrderNumber || + len(result.Orders) != 1 || result.Orders[0].ExternalStockID != "12" || + result.Orders[0].ShopName == nil || *result.Orders[0].ShopName != "测试店铺" || + len(result.Orders[0].Items) != 2 || result.Orders[0].Items[0].SKU != "BLACK-L" || + result.Orders[0].Items[1].SKU != "BLACK-XL" { + t.Fatalf("QueryOrder() = %#v", result) + } + encoded, err := json.Marshal(result) + if err != nil { + t.Fatalf("marshal result: %v", err) + } + for _, forbidden := range []string{"private-recipient", "private-phone", "private-item-phone"} { + if strings.Contains(string(encoded), forbidden) { + t.Fatalf("result leaked raw ERP value %q: %s", forbidden, encoded) + } + } + wantCalls := []string{ + CaptchaPath, LoginPath, UserPath, UserPath, + StockListTotalPath, StockListPath, StockDetailPath, + } + if strings.Join(calls, ",") != strings.Join(wantCalls, ",") { + t.Fatalf("endpoint calls = %#v, want %#v", calls, wantCalls) + } +} + +func TestSessionManagerQueryCreatedRangePaginatesAndDeduplicates(t *testing.T) { + listCalls := 0 + detailIDs := make([]uint64, 0) + server := httptest.NewServer(http.HandlerFunc(func( + writer http.ResponseWriter, + request *http.Request, + ) { + switch request.URL.Path { + case CaptchaPath: + http.SetCookie(writer, &http.Cookie{Name: "captcha", Value: "ready", Path: "/"}) + writer.Header().Set("Content-Type", "image/png") + _, _ = writer.Write([]byte("captcha")) + case LoginPath: + http.SetCookie(writer, &http.Cookie{Name: "authenticated", Value: "yes", Path: "/"}) + _, _ = writer.Write([]byte(`{"status":true,"data":{"user":{"id":1}}}`)) + case UserPath: + _, _ = writer.Write([]byte(`{"status":true,"data":{"id":1}}`)) + case StockListTotalPath: + assertStockPayload(t, request, "2026-07-22,2026-07-28", 0, 1, 20) + _, _ = writer.Write([]byte(`{"status":true,"data":21}`)) + case StockListPath: + listCalls++ + if listCalls == 1 { + assertStockPayload(t, request, "2026-07-22,2026-07-28", 0, 1, 20) + _, _ = writer.Write(stockListEnvelope(1, 20)) + return + } + assertStockPayload(t, request, "2026-07-22,2026-07-28", 20, 2, 20) + _, _ = writer.Write(stockListIDsEnvelope([]int{20})) + case StockDetailPath: + ids := decodeDetailPayload(t, request) + detailIDs = append(detailIDs, ids...) + _, _ = writer.Write(stockDetailEnvelope(ids)) + default: + writer.WriteHeader(http.StatusNotFound) + } + })) + defer server.Close() + + manager := testSessionManager(t, server.URL, "test-user", "test-password") + loginForSource(t, manager) + result, err := manager.QueryCreatedRange( + context.Background(), + "2026-07-22", + "2026-07-28", + ) + if err != nil { + t.Fatalf("QueryCreatedRange() error = %v", err) + } + if listCalls != 2 || len(detailIDs) != 20 || len(result.Orders) != 20 || + result.Query.CreatedFrom == nil || *result.Query.CreatedFrom != "2026-07-22" || + result.Query.CreatedTo == nil || *result.Query.CreatedTo != "2026-07-28" { + t.Fatalf("range query result = %#v, list calls = %d, detail IDs = %#v", result, listCalls, detailIDs) + } + seen := make(map[string]struct{}, len(result.Orders)) + for _, order := range result.Orders { + if _, exists := seen[order.ExternalStockID]; exists { + t.Fatalf("duplicate normalized order = %q", order.ExternalStockID) + } + seen[order.ExternalStockID] = struct{}{} + } +} + +func TestSessionManagerQueryRequiresConfiguredAuthenticatedSession(t *testing.T) { + manager := testSessionManager(t, "http://127.0.0.1:1", "", "") + if _, err := manager.QueryOrder(context.Background(), "SOURCE-12"); !errors.Is(err, domain.ErrFreightSourceNotConfigured) { + t.Fatalf("unconfigured QueryOrder() error = %v", err) + } + manager = testSessionManager(t, "http://127.0.0.1:1", "test-user", "test-password") + if _, err := manager.QueryOrder(context.Background(), "SOURCE-12"); !errors.Is(err, domain.ErrFreightSourceSessionNeeded) { + t.Fatalf("unauthenticated QueryOrder() error = %v", err) + } +} + +func loginForSource(t *testing.T, manager *SessionManager) { + t.Helper() + status, err := manager.FetchCaptcha(context.Background()) + if err != nil { + t.Fatalf("FetchCaptcha() error = %v", err) + } + if _, err := manager.Login(context.Background(), status.CaptchaTicket, "1234"); err != nil { + t.Fatalf("Login() error = %v", err) + } +} + +func assertStockPayload( + t *testing.T, + request *http.Request, + wantQuery string, + wantStart, wantPage, wantLength int, +) { + t.Helper() + assertERPHeaders(t, request) + if request.Method != http.MethodPost { + t.Fatalf("stock method = %s", request.Method) + } + var payload struct { + Length int `json:"length"` + Start int `json:"start"` + PageIndex int `json:"pageIndex"` + Queries []struct { + Value string `json:"dvalue"` + } `json:"queries"` + } + if err := json.NewDecoder(request.Body).Decode(&payload); err != nil { + t.Fatalf("decode stock payload: %v", err) + } + if payload.Length != wantLength || payload.Start != wantStart || + payload.PageIndex != wantPage || len(payload.Queries) != 1 || + payload.Queries[0].Value != wantQuery { + t.Fatalf("stock payload = %#v", payload) + } +} + +func assertDetailPayload(t *testing.T, request *http.Request, want []uint64) { + t.Helper() + got := decodeDetailPayload(t, request) + if len(got) != len(want) { + t.Fatalf("detail ids = %#v, want %#v", got, want) + } + for index := range want { + if got[index] != want[index] { + t.Fatalf("detail ids = %#v, want %#v", got, want) + } + } +} + +func decodeDetailPayload(t *testing.T, request *http.Request) []uint64 { + t.Helper() + assertERPHeaders(t, request) + var payload struct { + IDs []uint64 `json:"ids"` + } + if err := json.NewDecoder(request.Body).Decode(&payload); err != nil { + t.Fatalf("decode detail payload: %v", err) + } + return payload.IDs +} + +func stockListEnvelope(first, last int) []byte { + ids := make([]int, 0, last-first+1) + for value := first; value <= last; value++ { + ids = append(ids, value) + } + return stockListIDsEnvelope(ids) +} + +func stockListIDsEnvelope(ids []int) []byte { + rows := make([]map[string]any, 0, len(ids)) + for _, id := range ids { + rows = append(rows, map[string]any{"id": id, "code": "SOURCE-" + strconv.Itoa(id)}) + } + content, _ := json.Marshal(map[string]any{ + "status": true, + "data": map[string]any{"list": rows}, + }) + return content +} + +func stockDetailEnvelope(ids []uint64) []byte { + rows := make([]map[string]any, 0, len(ids)) + for _, id := range ids { + rows = append(rows, map[string]any{ + "id": id, + "details": []any{}, + }) + } + content, _ := json.Marshal(map[string]any{ + "status": true, + "data": map[string]any{"list": rows}, + }) + return content +} diff --git a/backend-api/internal/usecase/freight_service.go b/backend-api/internal/usecase/freight_service.go index 39fb013..69197ed 100644 --- a/backend-api/internal/usecase/freight_service.go +++ b/backend-api/internal/usecase/freight_service.go @@ -672,7 +672,7 @@ func daysInclusive(start, end time.Time) int { func freightSourceErrorCode(err error) string { switch { case errors.Is(err, domain.ErrFreightSourceNotConfigured): - return "ERP_CONNECTOR_NOT_CONFIGURED" + return "ERP_NOT_CONFIGURED" case errors.Is(err, domain.ErrFreightSourceSessionNeeded): return "ERP_SESSION_REQUIRED" case errors.Is(err, domain.ErrFreightSourceNotFound): @@ -680,7 +680,7 @@ func freightSourceErrorCode(err error) string { case errors.Is(err, domain.ErrFreightSourceProtocol): return "ERP_RESPONSE_INVALID" default: - return "ERP_CONNECTOR_UNAVAILABLE" + return "ERP_UNAVAILABLE" } } diff --git a/backend-api/internal/usecase/freight_service_test.go b/backend-api/internal/usecase/freight_service_test.go index ceba7b7..cff16c7 100644 --- a/backend-api/internal/usecase/freight_service_test.go +++ b/backend-api/internal/usecase/freight_service_test.go @@ -55,6 +55,24 @@ func TestFreightNormalizationRejectsConflictingIdentityAndInvalidTime( } } +func TestFreightSourceErrorCodesAreSourceNeutral(t *testing.T) { + cases := []struct { + err error + want string + }{ + {domain.ErrFreightSourceNotConfigured, "ERP_NOT_CONFIGURED"}, + {domain.ErrFreightSourceSessionNeeded, "ERP_SESSION_REQUIRED"}, + {domain.ErrFreightSourceNotFound, "ERP_FREIGHT_NOT_FOUND"}, + {domain.ErrFreightSourceProtocol, "ERP_RESPONSE_INVALID"}, + {errors.New("temporary source failure"), "ERP_UNAVAILABLE"}, + } + for _, testCase := range cases { + if got := freightSourceErrorCode(testCase.err); got != testCase.want { + t.Fatalf("freightSourceErrorCode(%v) = %q, want %q", testCase.err, got, testCase.want) + } + } +} + func TestFreightDateQuerySplitsIntoSevenDayWindows(t *testing.T) { source := &recordingDateSource{} service := &FreightService{source: source} diff --git a/docs/04-architecture.md b/docs/04-architecture.md index f091170..2d98318 100644 --- a/docs/04-architecture.md +++ b/docs/04-architecture.md @@ -15,9 +15,7 @@ | +--> SQLite / 文件存储 | - +--> loopback ERP Connector ----> 顺运宝 ERP - | - +--> Go 内存 ERP 会话 ---------> 顺运宝 ERP + +--> Go 内存 ERP 会话 / 直连 source ----> 顺运宝 ERP ^ | Android 采购 App --------------------> VLM Provider @@ -33,8 +31,8 @@ Android 采购 App --------------------> VLM Provider - 拼多多是第三方受控边界,只能由 Android 设备在已登录会话中操作。 - VLM 配置和 Key 位于手机,Key 使用 Android Keystore 包装的加密存储;管理后端不 保存、下发或代理模型调用。 -- ERP 账号、验证码会话、Cookie 和 JWT 只存在于 loopback Python Connector; - Go 后端只接收最小化规范货运数据。 +- ERP 账号、验证码会话、Cookie 和可能的 JWT 仅存在于 Go API 进程的受锁内存会话; + 浏览器和 SQLite 只接收最小化规范货运数据。 ## 二、模块职责 @@ -136,34 +134,35 @@ App 支持两个显式模式: ### 2.4 顺运宝 ERP 适配层 -Connector 是外部系统防腐层,不属于采购任务状态机: +Go 直连 source 是外部系统防腐层,不属于采购任务状态机: ```text Admin 创建 sync run - -> Go 后台 worker 调用 loopback Connector - -> Connector 使用受控 ERP 会话查询 list/listTotal/listByStock - -> Connector 规范化并去除收件 PII + -> Go 后台 worker 复用受锁 ERP 内存会话 + -> Go source 查询 listTotal/list/listByStock + -> Go source 规范化并去除收件 PII -> Go 同事务 upsert freight order/items -> Admin 复核 procurement request -> 显式生成不可变 purchase task ``` -- Python 进程默认只监听 `127.0.0.1`,使用独立服务密钥;Go 不传输或保存 ERP 密码。 +- ERP 账号/密码只由 API 启动环境读取;浏览器不提交凭证,SQLite 和普通日志不保存凭证。 - 精确单号只是查询条件,外部身份固定为 `stock.id` 和 `details[].id`。 -- Connector 不写本项目 SQLite,不调用 Roubao/拼多多/VLM,也不打印响应 body。 +- source 不调用 Roubao/拼多多/VLM,不打印原始响应 body,只返回 allowlist 规范结构。 - Go handler 只创建 sync run;外部查询由有界 worker 执行,避免把验证码或 ERP 延迟绑定到浏览器请求。 - 未确认的 `productThumb` 只按 ERP 引用保存,不能拼接 URL 或越权下载。 - ERP 接口由页面协议观察得到,正式生产前需确认开放 API、服务账号、调用频率、 缓存和个人信息处理权限。 -T-225 已在 `internal/platform/shunyunbao` 冻结直接 Go 实现的协议常量、请求 header、 -完整单号/日期范围条件、分页、详情批量和 allowlist 归一化。T-226 已增加一个受互斥 -保护的 Go 内存会话:Admin 以短期 ticket 读取验证码图片并人工输入验证码,服务端才 -持有 Cookie jar;账号密码只从启动环境读取。进程重启即失去该会话,不用 Redis 或持久化 -Cookie。货运用例只识别来源中立的“未配置、会话失效、未找到、协议异常、暂时不可用” -错误,不依赖 Connector 包。T-227 才将异步 sync worker 从 loopback Connector 切换到 -Go source,T-228 删除旧 Python 进程和服务密钥。 +T-225 冻结了 `internal/platform/shunyunbao` 的协议常量、请求 header、完整单号/日期范围 +条件、分页、详情批量和 allowlist 归一化。T-226 增加受互斥保护的 Go 内存会话:Admin +以短期 ticket 读取验证码图片并人工输入验证码,服务端才持有 Cookie jar;账号密码只从 +启动环境读取。T-227 将此会话直接作为 `FreightSource`:每次同步先校验会话,再以最多 +100 条、每页 20 条和每批最多 100 个详情 ID 查询;响应不完整、身份冲突或会话失效均使 +整批失败。进程重启即失去会话,不用 Redis 或持久化 Cookie。货运用例只识别来源中立的 +“未配置、会话失效、未找到、协议异常、暂时不可用”错误。T-228 再删除未参与运行的旧 +Python Connector、loopback 端口和服务密钥。 模式在 execution 开始时固定并写入结果;AI 失败后只能由人员明确切换,不能静默降级。 App 同时固定 provider ID、model、prompt/schema version 和证据 SHA-256,作为非秘密 diff --git a/docs/api.md b/docs/api.md index 1235d99..bba78b2 100644 --- a/docs/api.md +++ b/docs/api.md @@ -164,28 +164,10 @@ T-203 成功返回 `201`。使用相同 `Idempotency-Key` 和相同图片内容 受鉴权的文件流。只能读取用户有权查看的任务资产;不返回服务端文件路径。 -## ERP Connector 与货运信息 +## Go ERP 会话与货运信息 -Connector 是 loopback 内部服务,不复用 ADMIN/BUYER 凭证。除 `/health` 外要求 -`X-API-Key`,密钥至少 32 字节并由两进程环境变量注入;请求和响应不得写 body 日志。 - -### Connector `POST /v1/freight/query` - -```json -{"mode":"ORDER_NUMBER","order_number":"完整单号"} -``` - -或按 Asia/Shanghai 自然日闭区间查询,单次最多 7 天: - -```json -{ - "mode":"CREATED_RANGE", - "created_from":"2026-07-22", - "created_to":"2026-07-28" -} -``` - -成功返回最小化规范结构: +顺运宝仅由 Go API 进程直连;浏览器不复用 ERP 凭证,也不接收 Cookie、JWT、用户资料或 +原始响应。货运同步只接收以下最小化规范结构: ```json { @@ -267,7 +249,7 @@ T-224 增加: 从成功水位前回看 10 分钟对应的自然日开始,后端再按最多 7 天切窗。 响应 `202`,返回 sync id/status。订单号不进入 URL、事件 message 或访问日志;数据库 -只保存规范值及用于审计/检索的受控字段,不保存 Connector 密钥。 +只保存规范值及用于审计/检索的受控字段,不保存 ERP 凭证、Cookie 或 JWT。 ### `GET /api/v1/freight-syncs/{sync_id}` @@ -283,24 +265,19 @@ T-224 增加: ### `GET /api/v1/freight-orders` 返回当前 ADMIN creator 最近更新的货运单并支持 `limit`(1..100)。响应不包含 -收件人字段或 Connector 原文;列表筛选和 cursor 后置,不属于本次日期同步闭环。 +收件人字段或 ERP 原文;列表筛选和 cursor 后置,不属于本次日期同步闭环。 ### `GET /api/v1/freight-orders/{id}` T-222 返回货运头、全部当前商品明细和 revision/hash 状态。T-223 再增加采购需求和 已生成 task 引用。不存在和跨 creator 统一 404;响应 `Cache-Control: no-store`。 -T-222 精确单号同步需要 Go 后端配置: - -- `CMROUBAO_ERP_CONNECTOR_URL`:默认 `http://127.0.0.1:8091`,只接受 HTTP - loopback origin。 -- `CMROUBAO_ERP_CONNECTOR_API_KEY`:与 Connector 的 - `SHUNYUNBAO_SERVICE_API_KEY` 相同,至少 32 UTF-8 字节;未设置时后端可启动, - 但同步任务稳定失败为 `ERP_CONNECTOR_NOT_CONFIGURED`。 - -Connector 未登录、找不到货运单、响应协议错误和暂时不可用分别落为 -`ERP_SESSION_REQUIRED`、`ERP_FREIGHT_NOT_FOUND`、`ERP_RESPONSE_INVALID` 和 -`ERP_CONNECTOR_UNAVAILABLE`,不返回 ERP 原始错误 body。 +T-227 的后台 worker 使用当前 Go 内存会话,按 `listTotal -> list 分页 -> listByStock` +查询。每次先校验会话;列表最多 100 条、每页 20 条,详情每批最多 100 个外部 stock ID。 +未配置、未登录、找不到货运单、响应协议错误和暂时不可用分别落为 +`ERP_NOT_CONFIGURED`、`ERP_SESSION_REQUIRED`、`ERP_FREIGHT_NOT_FOUND`、 +`ERP_RESPONSE_INVALID` 和 `ERP_UNAVAILABLE`,不返回 ERP 原始错误 body。旧 Python +Connector 仅在仓库中等待 T-228 删除,不再是运行时依赖。 ### `POST /api/v1/freight-items/{item_id}/procurement-request` diff --git a/docs/current-state.md b/docs/current-state.md index 4dfb3ad..12c5f96 100644 --- a/docs/current-state.md +++ b/docs/current-state.md @@ -5,7 +5,7 @@ ## 当前快照 - 日期:2026-07-29 -- 阶段:T-226 Go ERP 内存会话与 Admin 人工验证码登录已完成;运行时切换待 T-227 至 T-228 +- 阶段:T-227 Go 直连 ERP 查询已完成;T-228 待删除遗留 Python Connector - Git:当前分支为 `main`;T-001 至 T-004、T-101 至 T-104、T-201 至 T-219 均按文档提交、实现提交的顺序纳入历史 - 生产代码:`android-buyer/` 已接入 Roubao Android 源码 @@ -14,17 +14,17 @@ - 后端:Go 1.23.0 + Gin 1.11.0 + SQLite + Goose 3.26.0;已实现图片/任务业务、 SSR 管理 Web、ADMIN/BUYER 联合认证、设备 readiness、原子 claim、租约状态机、 task-scoped 参考图和 `authctl` -- ERP 货运:Go 后端通过仅限 loopback、服务密钥鉴权的 Python Connector 异步按 - 完整单号同步;v12 保存同步记录、货运头和全部明细,canonical hash 控制 revision, - Admin 已有 `/freight`、`/freight/import`、`/freight/{id}` 与对应 JSON API。 +- ERP 货运:Go 后端以受锁内存 ERP 会话直连异步按完整单号或创建日期同步;v12 保存 + 同步记录、货运头和全部明细,canonical hash 控制 revision,Admin 已有 `/freight`、 + `/freight/import`、`/freight/{id}` 与对应 JSON API。 - ERP Go 迁移:T-225 已用脱敏 fixture 固定 `internal/platform/shunyunbao` 的 header、 单号/日期查询、分页、详情批量和字段 allowlist,并使货运用例依赖来源中立错误。T-226 已增加受锁保护的 Go 内存 Cookie jar、验证码 ticket、登录和用户校验,以及 ADMIN 的 `/erp` 页面/API;账号密码仅由 `CMROUBAO_SHUNYUNBAO_*` 启动环境读取,重启后需人工 - 重新登录。未访问真实 ERP,Python Connector 仍是货运同步运行时;T-227 才切换 source, - T-228 删除旧 Connector。 -- ERP 增量同步:v14 支持 Asia/Shanghai 创建日期闭区间和“同步至现在”,Connector - 单窗最多 7 天,后端对较长水位范围切窗并从成功水位前 10 分钟所在自然日回看。 + 重新登录。T-227 已将其作为 `FreightSource`,查询先校验会话、再执行有界分页/详情 + 批量并返回 allowlist;未访问真实 ERP。旧 Python Connector 仅待 T-228 删除。 +- ERP 增量同步:v14 支持 Asia/Shanghai 创建日期闭区间和“同步至现在”,source 单窗 + 最多 7 天,后端对较长水位范围切窗并从成功水位前 10 分钟所在自然日回看。 货运落库、同步成功和水位推进同事务完成;失败与较旧范围成功不推进水位。Admin 可查看查询范围、计数、稳定错误码和最近成功水位。 - 采购需求:v13 按货运商品 revision 保存可复核快照;ADMIN 显式确认采购、上传参考 @@ -35,11 +35,11 @@ - 测试:T-219 Android Debug/Release 单元测试与构建和根 `init.ps1` 通过; Debug APK `1.4.16 (21)` 已覆盖安装到 PKG110 - 后端测试:T-226 已运行 `go test ./...`、`go test -race ./...`、`go vet ./...` 和三个 Go - 入口构建;新增配置、会话并发/过期、ADMIN API、SSR 验证码页面以及隔离 Gin/Chrome - 桌面/手机视口 smoke,均只用伪 ERP 或未配置状态; + 入口构建;T-227 增加 Go source 的伪 ERP 会话预检、完整单号、日期分页去重、详情 + allowlist 和稳定错误码覆盖,均未访问真实 ERP; 覆盖 v14 上下迁移、7 天切窗、水位重叠、中途失败、空窗口、重复页、来源 revision、 - 水位事务/不回退、Admin API/SSR 和 Connector 严格响应窗口;Python 22 项伪响应 - 测试通过,未访问真实 ERP + 水位事务/不回退、Admin API/SSR 和旧 Connector 严格响应窗口;Python 22 项伪响应 + 测试保留至 T-228 删除前,未访问真实 ERP - 原型:4 个管理 Web 页面和 7 个 Android 页面均可离线独立打开;Playwright 以 1440×900、390×844、360×800 验证 36 个页面/视口组合,无页面横向溢出、 脚本错误或外部请求,Android 可见交互控件均不小于 44px @@ -179,6 +179,9 @@ | `docs/tasks/T-222.md` | DONE | 货运信息存储、API 与 Admin 页面 | | `docs/tasks/T-223.md` | DONE | 待采购需求提取与任务生成 | | `docs/tasks/T-224.md` | DONE | ERP 日期增量同步 | +| `docs/tasks/T-225.md` | DONE | 冻结 Go 直连 ERP 协议与安全边界 | +| `docs/tasks/T-226.md` | DONE | Go ERP 会话、验证码登录与 Admin 连接页 | +| `docs/tasks/T-227.md` | DONE | Go 直连顺运宝查询接入货运同步 | | `docs/design/` | 已确认 | T-202 原型索引、4 个管理页和 7 个 Android 页面 | | `deepseek总结.txt` | 已有 | 历史讨论摘要,不是正式需求权威 | | `android-buyer/` | 已有 | Roubao `main` 固定 commit 的 Android 基线 | @@ -191,12 +194,11 @@ ## 任务摘要 - 已完成:T-001 至 T-004、T-101 至 T-104、T-201 至 T-219。 -- 已完成:另含 T-220 至 T-226 ERP 契约、Connector、货运存储、采购需求生成、日期 - 增量同步、Go 直连协议安全边界和人工验证码会话。 +- 已完成:另含 T-220 至 T-227 ERP 契约、货运存储、采购需求生成、日期增量同步、Go + 直连协议、人工验证码会话和直连 `FreightSource`。 - 进行中:无。 -- 下一步:T-227 将已验证 Go 会话接入 `FreightSource`,保留现有异步单号/日期同步、 - 水位、revision 和采购任务语义;真实 ERP 上线前仍需确认开放 API、数据使用权限并由 - 人员完成验证码登录。 +- 下一步:T-228 删除旧 Python Connector、loopback 端口/服务密钥及其文档;真实 ERP + 上线前仍需确认开放 API、数据使用权限并由人员完成验证码登录。 ## 当前可运行内容 diff --git a/docs/tasks/T-227.md b/docs/tasks/T-227.md index 93cfc8e..4e2c5af 100644 --- a/docs/tasks/T-227.md +++ b/docs/tasks/T-227.md @@ -4,7 +4,7 @@ title: Go 直连顺运宝查询接入货运同步 phase: 2 deps: - T-226 -status: TODO +status: DONE created: 2026-07-29 context_ref: d62a4af work_branch: null @@ -50,11 +50,11 @@ T-226 能在 Go 后端建立受控 ERP 会话,但货运同步仍通过 Python ## 验收要点 -- [ ] 完整单号、日期范围、重叠水位、分页去重、详情多商品、失败不推进水位均有 Go 集成测试。 -- [ ] Admin 输入单号后创建异步 run,并由 Go source 直接得到相同规范化货运/明细结果。 -- [ ] 缺会话显示 `ERP_SESSION_REQUIRED`;不再出现 Connector URL/API Key 相关错误。 -- [ ] 真实受控单号 smoke 可选且不记录订单内容;未具备凭证时明确标记未运行。 -- [ ] `go test ./...`、`go test -race ./...`、`go vet ./...` 通过。 +- [x] 完整单号、日期范围、重叠水位、分页去重、详情多商品、失败不推进水位均有 Go 集成测试。 +- [x] Admin 输入单号后创建异步 run,并由 Go source 直接得到相同规范化货运/明细结果。 +- [x] 缺会话显示 `ERP_SESSION_REQUIRED`;运行时不再出现 Connector URL/API Key 相关错误。 +- [x] 真实受控单号 smoke 可选且不记录订单内容;未具备凭证,明确未运行。 +- [x] `go test ./...`、`go test -race ./...`、`go vet ./...` 通过。 ## 边界 @@ -65,3 +65,17 @@ T-226 能在 Go 后端建立受控 ERP 会话,但货运同步仍通过 Python ## 执行记录 - 2026-07-29:由 T-226 依赖创建,等待 Admin 会话能力完成。 +- 2026-07-29:开始实施;T-226 已由 `a5a61c0` 完成。货运同步保留既有异步 run、 + 日期水位和 SQLite 幂等语义,仅将运行时 `FreightSource` 切换为 Go 直连实现。 +- 2026-07-29:`SessionManager` 现实现 `FreightSource`;每次查询先以同一 Cookie jar + 校验 `/am/user/get`,再按 `listTotal -> list -> listByStock?hist=0` 执行。查询上限为 + 100 条、分页为 20 条、详情批量为 100 个 ID;重复的完全相同外部身份去重,缺失详情、 + 冲突身份、超限和协议异常均整批失败。 +- 2026-07-29:API 进程已直接注入 Go source,不再构造 `erpconnector.Client`。同步失败码 + 调整为 `ERP_NOT_CONFIGURED`、`ERP_SESSION_REQUIRED`、`ERP_FREIGHT_NOT_FOUND`、 + `ERP_RESPONSE_INVALID`、`ERP_UNAVAILABLE`;旧 Connector 代码和配置仅留待 T-228 删除。 +- 验证:在 `backend-api/` 执行 `$env:GOTOOLCHAIN='local'; go test ./...; go test -race ./...;` + `go vet ./...; go build ./cmd/api; go build ./cmd/migrate; go build ./cmd/authctl`,全部通过。 + `httptest` 覆盖同 Cookie jar、完整单号、多商品详情、日期分页去重、allowlist、未配置/ + 未登录和稳定错误码;既有 Go 集成测试继续覆盖日期切窗、重叠水位与失败不推进水位。 +- 未运行:没有读取或使用真实 ERP 凭证,因此未执行真实受控单号 smoke;未记录订单内容。