Files
cmsp/internal/erpgo/client.go
T

286 lines
10 KiB
Go

// Package erpgo 消费 ERPGo 的货憨憨查询与 Shopee 商品视频接口。
package erpgo
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"strconv"
"strings"
"time"
"cmsp/internal/config"
"cmsp/internal/store"
)
const pageSize = 200
const maximumPages = 200
// Error 的 JSON 文本跨 Wails 错误通道传递稳定错误码,不透传上游原文。
type Error struct {
Code string `json:"errorCode"`
Message string `json:"message"`
Status int `json:"status,omitempty"`
RequestID string `json:"requestId,omitempty"`
}
func (e *Error) Error() string {
b, _ := json.Marshal(e)
return string(b)
}
func failure(code string, status int, requestID string) *Error {
messages := map[string]string{
"ERPGo_NOT_CONFIGURED": "请在参数设置填写 erpgo 服务地址和 API Key",
"INVALID_ARGUMENT": "查询参数或 erpgo 地址格式不正确,请检查设置",
"API_KEY_INVALID": "erpgo API Key 无效,请检查参数设置",
"SHOP_ACCESS_DENIED": "erpgo 不允许查询此店铺,请检查服务端账号与店铺",
"HHH_AUTH_FAILED": "erpgo 的货憨憨认证失败,请检查服务端账号配置",
"HHH_UPSTREAM_ERROR": "货憨憨查询失败,原数据已保留,请稍后重试",
"HHH_UPSTREAM_TIMEOUT": "货憨憨查询超时,原数据已保留,请稍后重试",
"SERVICE_UNAVAILABLE": "erpgo 服务暂不可用,原数据已保留",
"RATE_LIMITED": "查询过于频繁,请稍后重试",
"INTERNAL_ERROR": "erpgo 内部错误,请联系维护者",
"NETWORK_ERROR": "无法连接 erpgo,请检查服务地址和网络",
"REQUEST_CANCELLED": "查询已取消,原数据已保留",
"INVALID_RESPONSE": "erpgo 返回的数据或分页不符合契约,原数据已保留",
}
message, ok := messages[code]
if !ok {
code, message = "INVALID_RESPONSE", messages["INVALID_RESPONSE"]
}
return &Error{Code: code, Message: message, Status: status, RequestID: requestID}
}
type Client struct {
baseURL string
apiKey string
http *http.Client
}
func NewClient(cfg config.ERPGoConfig, httpClient *http.Client) (*Client, error) {
if err := cfg.Validate(); err != nil {
return nil, failure("INVALID_ARGUMENT", 0, "")
}
if strings.TrimSpace(cfg.BaseURL) == "" || strings.TrimSpace(cfg.APIKey) == "" {
return nil, failure("ERPGo_NOT_CONFIGURED", 0, "")
}
client := http.Client{Timeout: 150 * time.Second}
if httpClient != nil {
client = *httpClient
}
// 防止重定向将 X-API-Key 转发到其他主机;维护者应填写最终服务地址。
client.CheckRedirect = func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }
return &Client{baseURL: strings.TrimRight(strings.TrimSpace(cfg.BaseURL), "/"), apiKey: strings.TrimSpace(cfg.APIKey), http: &client}, nil
}
type envelope struct {
Code int `json:"code"`
ErrorCode string `json:"errorCode"`
RequestID string `json:"requestId"`
Data json.RawMessage `json:"data"`
}
func (c *Client) get(ctx context.Context, endpoint string, query url.Values, dest any) error {
u := c.baseURL + "/api/v1/integrations/huohanhan/" + endpoint
if len(query) > 0 {
u += "?" + query.Encode()
}
req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
if err != nil {
return failure("INVALID_ARGUMENT", 0, "")
}
req.Header.Set("X-API-Key", c.apiKey)
req.Header.Set("Accept", "application/json")
response, err := c.http.Do(req)
if err != nil {
if errors.Is(err, context.Canceled) {
return failure("REQUEST_CANCELLED", 0, "")
}
if errors.Is(err, context.DeadlineExceeded) {
return failure("HHH_UPSTREAM_TIMEOUT", 0, "")
}
return failure("NETWORK_ERROR", 0, "")
}
defer response.Body.Close()
const maxBytes = 16 << 20
raw, err := io.ReadAll(io.LimitReader(response.Body, maxBytes+1))
if err != nil || len(raw) > maxBytes {
return failure("INVALID_RESPONSE", response.StatusCode, "")
}
var result envelope
if json.Unmarshal(raw, &result) != nil {
return failure("INVALID_RESPONSE", response.StatusCode, "")
}
requestID := safeRequestID(result.RequestID, c.apiKey)
if response.StatusCode != http.StatusOK || result.Code != http.StatusOK {
if result.Code != response.StatusCode || result.ErrorCode == "" {
return failure("INVALID_RESPONSE", response.StatusCode, requestID)
}
return failure(result.ErrorCode, response.StatusCode, requestID)
}
if result.ErrorCode != "" || len(result.Data) == 0 || string(result.Data) == "null" || json.Unmarshal(result.Data, dest) != nil {
return failure("INVALID_RESPONSE", response.StatusCode, requestID)
}
return nil
}
func safeRequestID(id, key string) string {
if len(id) > 128 || strings.Contains(id, key) {
return ""
}
for _, r := range id {
if !(r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z' || r >= '0' && r <= '9' || r == '-' || r == '_') {
return ""
}
}
return id
}
type metadata struct {
Source string `json:"source"`
FetchedAt string `json:"fetchedAt"`
}
func (m metadata) valid() bool {
_, err := time.Parse(time.RFC3339, m.FetchedAt)
return m.Source == "huohanhan" && err == nil
}
// Shop 只声明契约中的白名单,避免将上游额外字段透传到本地或界面。
type Shop struct {
ID string `json:"id"`
PlatformShopID string `json:"platformShopId"`
ShopName string `json:"shopName"`
ShopAlias string `json:"shopAlias"`
Region string `json:"region"`
RegionName string `json:"regionName"`
Platform string `json:"platform"`
Status string `json:"status"`
}
func (c *Client) ListShops(ctx context.Context) ([]store.Shop, error) {
var data struct {
metadata
Items []Shop `json:"items"`
}
if err := c.get(ctx, "shops", nil, &data); err != nil {
return nil, err
}
if !data.valid() || data.Items == nil {
return nil, failure("INVALID_RESPONSE", 200, "")
}
shops := make([]store.Shop, 0, len(data.Items))
seen := make(map[string]bool)
for _, s := range data.Items {
if strings.TrimSpace(s.PlatformShopID) == "" || s.Platform != "0" || seen[s.PlatformShopID] {
return nil, failure("INVALID_RESPONSE", 200, "")
}
seen[s.PlatformShopID] = true
shops = append(shops, store.Shop{ID: s.ID, PlatformShopID: s.PlatformShopID, ShopName: s.ShopName, ShopAlias: s.ShopAlias, Region: s.Region, RegionName: s.RegionName, Platform: s.Platform, Status: s.Status})
}
return shops, nil
}
type product struct {
ID string `json:"id"`
ItemID string `json:"itemId"`
PlatformShopID string `json:"platformShopId"`
ShopName string `json:"shopName"`
ItemName string `json:"itemName"`
MainImage string `json:"mainImage"`
Currency string `json:"currency"`
MinSkuPrice float64 `json:"minSkuPrice"`
ItemStatus string `json:"itemStatus"`
CreateTime string `json:"createTime"`
QualityLevel string `json:"qualityLevel"`
VideoDiagnosis string `json:"videoDiagnosis"`
Diagnoses []store.Diagnosis `json:"diagnoses"`
}
type productPage struct {
metadata
PlatformShopID string `json:"platformShopId"`
ItemStatus string `json:"itemStatus"`
Current int `json:"current"`
Size int `json:"size"`
Total int `json:"total"`
Pages int `json:"pages"`
Records []product `json:"records"`
}
// DownloadAllProducts 先完整拉取并校验;任何页失败都不返回可落库的部分结果。
func (c *Client) DownloadAllProducts(ctx context.Context, shopID string, onProgress func(int, int)) ([]store.Product, map[string][]store.Diagnosis, error) {
shopID = strings.TrimSpace(shopID)
if shopID == "" {
return nil, nil, failure("INVALID_ARGUMENT", 0, "")
}
products := make([]store.Product, 0)
diagnoses := make(map[string][]store.Diagnosis)
seen := make(map[string]int)
for current := 1; current <= maximumPages; current++ {
var page productPage
q := url.Values{"platformShopId": {shopID}, "current": {strconv.Itoa(current)}, "size": {strconv.Itoa(pageSize)}}
if err := c.get(ctx, "products", q, &page); err != nil {
return nil, nil, err
}
if !page.valid() || page.PlatformShopID != shopID || page.ItemStatus != "NORMAL" || page.Current != current || page.Size != pageSize || page.Total < 0 || page.Total > maximumPages*pageSize || page.Pages < 0 || page.Pages > maximumPages || page.Records == nil || len(page.Records) > pageSize || page.Pages != (page.Total+pageSize-1)/pageSize {
return nil, nil, failure("INVALID_RESPONSE", 200, "")
}
if len(page.Records) == 0 {
if page.Total == 0 && page.Pages == 0 || current > page.Pages {
return products, diagnoses, nil
}
return nil, nil, failure("INVALID_RESPONSE", 200, "")
}
if current > page.Pages {
return nil, nil, failure("INVALID_RESPONSE", 200, "")
}
for _, p := range page.Records {
if strings.TrimSpace(p.ID) == "" || strings.TrimSpace(p.ItemID) == "" || p.PlatformShopID != shopID || p.ItemStatus != "NORMAL" || p.Diagnoses == nil || (p.VideoDiagnosis != store.VideoDiagnosisMissing && p.VideoDiagnosis != store.VideoDiagnosisOK) {
return nil, nil, failure("INVALID_RESPONSE", 200, "")
}
item := store.Product{ID: p.ID, ItemID: p.ItemID, PlatformShopID: p.PlatformShopID, ShopName: p.ShopName, ItemName: p.ItemName, MainImage: p.MainImage, Currency: p.Currency, MinSkuPrice: p.MinSkuPrice, ItemStatus: p.ItemStatus, CreatedAt: p.CreateTime, QualityLevel: p.QualityLevel, VideoDiagnosis: p.VideoDiagnosis}
if index, ok := seen[p.ID]; ok {
products[index] = item
} else {
seen[p.ID] = len(products)
products = append(products, item)
}
for i := range p.Diagnoses {
p.Diagnoses[i].ProductID = p.ID
}
diagnoses[p.ID] = p.Diagnoses
}
if onProgress != nil {
onProgress(current, page.Pages)
}
if current == page.Pages {
return products, diagnoses, nil
}
}
return nil, nil, failure("INVALID_RESPONSE", 200, "")
}
// BridgeError 为本地数据库等错误提供相同的前端错误协议。
func BridgeError(err error) error {
var apiError *Error
if errors.As(err, &apiError) {
return apiError
}
return &Error{Code: "LOCAL_SYNC_FAILED", Message: "保存同步数据失败,原数据已保留,请查看运行日志"}
}
// LogSummary 不包含上游原文、查询 URL 或凭据。
func LogSummary(err error) string {
var e *Error
if errors.As(err, &e) {
return fmt.Sprintf("%s (HTTP %d, requestId=%s)", e.Code, e.Status, e.RequestID)
}
return "LOCAL_SYNC_FAILED"
}