feat(#70): import cmautobuy product catalog
This commit is contained in:
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Architecture-and-Code-Map
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Architecture-and-Code-Map.-
|
||||
wiki_revision: 5fc5fcee07160d6aecdba455c262016907848af7
|
||||
synchronized_at: 2026-08-23T15:20:30Z
|
||||
wiki_revision: 6863d83cfbd8dd6577475374f09a799a6f358346
|
||||
synchronized_at: 2026-08-23T15:56:21Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 架构与代码地图
|
||||
@@ -139,3 +139,12 @@ Android Portal/Agent
|
||||
| 三端统一验证 | `scripts/verify.ps1` |
|
||||
|
||||
以上入口均已落地;新增 API、表或执行动作时必须同步更新本代码地图与共享契约。
|
||||
|
||||
|
||||
## cmautobuy 商品离线导入
|
||||
|
||||
`server/cmd/import-cmautobuy-products/` 是 #70 的独立单向导入入口,转换逻辑位于 `server/app/goauto/cmautobuyimport/`。它不注册 HTTP 路由、不参与服务启动,也不是数据库 schema migration。
|
||||
|
||||
导入器从 cmautobuy MySQL 只读一致性快照读取有效 `pdd_products`、`shopee_products` 和 `shopee_skus`,再写入 GoAuto 的 `pdd_product`、`shopee_product`。旧 PDD 组合 SKU 聚合为 GoAuto 通用维度;同一颜色出现多个不同价格时不猜测价格。旧组合 `spec_mappings`、SYB、任务、订单和其他域数据不进入导入范围。
|
||||
|
||||
默认模式是 dry-run;显式 `--apply` 才会在目标 MySQL 单事务覆盖同业务键商品。来源与目标配置都来自未跟踪 YAML,连接串和密码不输出。来源蝦皮的字符串 `pdd_goods_id` 必须通过目标 PDD `goods_id` 换算为数字 `pdd_product_id`,提交前重新校验 JSON、唯一键和关联完整性。
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Business-Rules-and-Glossary
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Business-Rules-and-Glossary.-
|
||||
wiki_revision: 632e0c2d197884971925fe275e7a7495ec2ba31f
|
||||
synchronized_at: 2026-08-23T15:20:34Z
|
||||
wiki_revision: dcdf1dab84108752889e1c75bfd6bbc2a207ba01
|
||||
synchronized_at: 2026-08-23T15:56:25Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 业务规则与术语
|
||||
@@ -180,3 +180,14 @@ synchronized_at: 2026-08-23T15:20:34Z
|
||||
| 两趟执行 | 慢路径的实现方式:第一趟采集规格后结束并释放设备,服务端离线匹配,第二趟重新派发下单,复用 `task_attempt` 机制 |
|
||||
| 采购探测 | 采购时对 PDD 商品页的实时规格采集,同时作为一次商品信息更新写回档案,按 `completed` / `completed_partial` 语义处理 |
|
||||
| spec_source | 采购任务的规格来源标记:`manual_mapping` / `exact_match` / `ai_match`,用于事后批量追溯 |
|
||||
|
||||
|
||||
## cmautobuy 商品导入
|
||||
|
||||
- 这是人工触发的单向离线导入,不是持续双写或双库同步;不随服务启动执行。
|
||||
- 来源只读取未软删除的 PDD、蝦皮商品及蝦皮 SKU;不导入 SYB、任务、订单、物流、设备、用户、规则和旧组合规格映射。
|
||||
- 目标已有相同 `goods_id` 或存活 `shopee_item_id` 时,按用户确认由来源商品档案覆盖;PDD 人工 `disabled` 状态保留,不能被导入自动启用。
|
||||
- PDD 销量和评价、蝦皮参考售价在来源没有对应值时覆盖为空;蝦皮币种读取 GoAuto 系统设置,缺省 TWD。
|
||||
- 同一颜色只有唯一有效价格时才保存颜色价格;出现多个价格时价格留空、PDD 商品保持 `pending` 并报告,禁止取最低价、最高价或平均价。
|
||||
- 蝦皮只聚合来源已解析的颜色和尺码,未解析内容不猜测。来源 PDD 关联不存在、已删除或无法转换时,蝦皮记录作为冲突跳过,不能静默关联其他商品。
|
||||
- 默认 dry-run 不写数据库;`--apply` 前必须查看精确影响并再次人工确认。apply 在目标单事务中执行,提交前校验 JSON、唯一键和关联完整性。
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Local-Development-and-Verification
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Local-Development-and-Verification.-
|
||||
wiki_revision: 26b3a4aa0d95089e689020647938e5c54a4f733f
|
||||
synchronized_at: 2026-08-21T09:19:50Z
|
||||
wiki_revision: a5518705c35e3a2d71c239de7882ff501ba495ff
|
||||
synchronized_at: 2026-08-23T15:56:30Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 本地开发与验证
|
||||
@@ -114,3 +114,26 @@ Android 骨架使用 Kotlin 1.9.22、AGP 8.2.0、Java 17 和 SDK 34。
|
||||
- 浏览器打开 PDD、两层确认、商品详情识别和规格遍历。
|
||||
- 网络断开、登录失效、验证码、风控和规则删除后的快照执行。
|
||||
- 一台设备串行任务和 20 台设备连接稳定性。
|
||||
|
||||
|
||||
## cmautobuy 商品导入
|
||||
|
||||
从 `server/` 运行 #70 独立命令。两个配置文件都必须是未被 Git 跟踪的本地文件;密码不接受命令行参数。
|
||||
|
||||
```powershell
|
||||
# 默认只读预检,不写两个数据库
|
||||
go run ./cmd/import-cmautobuy-products --source-config "D:\chengma\cmautobuy\admin\config.yaml" --target-config "D:\OPC\goauto\config.yaml"
|
||||
|
||||
# 仅在检查 dry-run 精确数量、备份目标商品表并再次人工确认后运行
|
||||
go run ./cmd/import-cmautobuy-products --source-config "D:\chengma\cmautobuy\admin\config.yaml" --target-config "D:\OPC\goauto\config.yaml" --apply
|
||||
```
|
||||
|
||||
也可分别用 `GOAUTO_CMAUTOBUY_CONFIG` 和 `GOAUTO_CONFIG` 指定路径。来源配置支持 `disabled` 或 `verify_ca`;相对 CA 路径按来源配置文件目录解析。导入报告不得包含 DSN、密码或原始业务响应。
|
||||
|
||||
MySQL 8.4 写入路径的隔离测试只允许连接名称以 `_test` 结尾的数据库:
|
||||
|
||||
```powershell
|
||||
$env:GOAUTO_IMPORT_MYSQL_TEST_DSN="<仅由安全环境注入的 _test DSN>"
|
||||
go test ./app/goauto/cmautobuyimport -run TestRunMySQL84DryRunAndApply -count=1
|
||||
Remove-Item Env:GOAUTO_IMPORT_MYSQL_TEST_DSN
|
||||
```
|
||||
|
||||
@@ -2,8 +2,8 @@
|
||||
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
|
||||
wiki_page: Delivery-Issues
|
||||
wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Delivery-Issues.-
|
||||
wiki_revision: 1e4f063a6c557e11c81a37ea626d75020c23b3cc
|
||||
synchronized_at: 2026-08-23T15:20:51Z
|
||||
wiki_revision: ea99bb3501bfd66f6876e195c4b799b4a798965b
|
||||
synchronized_at: 2026-08-23T15:56:52Z
|
||||
<!-- gitea-wiki-mirror:end -->
|
||||
|
||||
# 当前 MVP 交付工单索引
|
||||
@@ -90,3 +90,10 @@ synchronized_at: 2026-08-23T15:20:51Z
|
||||
6. T17:一加/ColorOS 真机验收。
|
||||
|
||||
任何实现工作必须先把对应活动 Task 置为进行中,并在工单中记录验证证据。
|
||||
|
||||
|
||||
## 独立商品数据导入
|
||||
|
||||
| 顺序 | 工单 | 交付项 | 主要依赖 / 门禁 |
|
||||
|---|---|---|---|
|
||||
| T57 | [#70](https://git.ilapage.cn/OPC/goauto/issues/70) | 从 cmautobuy 线上库导入 PDD 与蝦皮商品 | #31、#40/#41;代码和 dry-run 已实施,正式 apply 等待用户查看精确影响后再次确认 |
|
||||
|
||||
@@ -0,0 +1,107 @@
|
||||
package cmautobuyimport
|
||||
|
||||
import (
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/go-sql-driver/mysql"
|
||||
"gopkg.in/yaml.v2"
|
||||
)
|
||||
|
||||
type databaseFile struct {
|
||||
Database map[string]any `yaml:"database"`
|
||||
}
|
||||
|
||||
// DSNFromConfig reads the same ignored YAML database block used by both
|
||||
// projects. It returns a DSN but callers must never print it.
|
||||
func DSNFromConfig(path string, source bool) (string, error) {
|
||||
raw, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("读取数据库配置失败: %w", err)
|
||||
}
|
||||
var file databaseFile
|
||||
if err := yaml.Unmarshal(raw, &file); err != nil {
|
||||
return "", fmt.Errorf("解析数据库配置失败: %w", err)
|
||||
}
|
||||
value := func(key string) string { return configScalar(file.Database[key]) }
|
||||
cfg := mysql.NewConfig()
|
||||
cfg.Net = "tcp"
|
||||
cfg.Addr = value("host") + ":" + value("port")
|
||||
cfg.User = value("user")
|
||||
cfg.Passwd = configScalar(file.Database["password"])
|
||||
cfg.DBName = value("name")
|
||||
cfg.ParseTime = true
|
||||
cfg.Loc = time.UTC
|
||||
cfg.Params = map[string]string{"charset": "utf8mb4", "time_zone": "'+00:00'"}
|
||||
cfg.Timeout = 5 * time.Second
|
||||
if cfg.Addr == ":" || cfg.User == "" || cfg.DBName == "" {
|
||||
return "", fmt.Errorf("数据库配置缺少 host、port、user 或 name")
|
||||
}
|
||||
if source {
|
||||
mode := strings.TrimSpace(value("tls_mode"))
|
||||
switch mode {
|
||||
case "", "disabled":
|
||||
case "verify_ca":
|
||||
caPath := strings.TrimSpace(value("tls_ca"))
|
||||
if caPath != "" && !filepath.IsAbs(caPath) {
|
||||
caPath = filepath.Join(filepath.Dir(path), caPath)
|
||||
}
|
||||
pem, readErr := os.ReadFile(caPath)
|
||||
if readErr != nil {
|
||||
return "", fmt.Errorf("读取来源数据库 CA 失败: %w", readErr)
|
||||
}
|
||||
pool := x509.NewCertPool()
|
||||
if !pool.AppendCertsFromPEM(pem) {
|
||||
return "", fmt.Errorf("来源数据库 CA 不是有效 PEM")
|
||||
}
|
||||
name := "goauto-cmautobuy-verify-ca"
|
||||
// cmautobuy's MySQL 8 auto-generated server certificate has no
|
||||
// SAN. Match its existing verify_ca behavior: skip only the built-in
|
||||
// hostname check, then explicitly verify the complete chain against
|
||||
// the configured server-specific CA. TLS is still mandatory and can
|
||||
// never fall back to plaintext.
|
||||
tlsConfig := &tls.Config{MinVersion: tls.VersionTLS12, InsecureSkipVerify: true}
|
||||
tlsConfig.VerifyConnection = func(state tls.ConnectionState) error {
|
||||
if len(state.PeerCertificates) == 0 {
|
||||
return fmt.Errorf("来源 MySQL 没有提供 TLS 证书")
|
||||
}
|
||||
intermediates := x509.NewCertPool()
|
||||
for _, certificate := range state.PeerCertificates[1:] {
|
||||
intermediates.AddCert(certificate)
|
||||
}
|
||||
_, err := state.PeerCertificates[0].Verify(x509.VerifyOptions{Roots: pool, Intermediates: intermediates, KeyUsages: []x509.ExtKeyUsage{x509.ExtKeyUsageServerAuth}})
|
||||
return err
|
||||
}
|
||||
if err := mysql.RegisterTLSConfig(name, tlsConfig); err != nil {
|
||||
return "", fmt.Errorf("注册来源数据库 TLS 配置失败: %w", err)
|
||||
}
|
||||
cfg.TLSConfig = name
|
||||
default:
|
||||
return "", fmt.Errorf("来源数据库 tls_mode 只支持 disabled 或 verify_ca")
|
||||
}
|
||||
}
|
||||
return cfg.FormatDSN(), nil
|
||||
}
|
||||
|
||||
func configScalar(value any) string {
|
||||
switch v := value.(type) {
|
||||
case nil:
|
||||
return ""
|
||||
case string:
|
||||
return v
|
||||
case int:
|
||||
return strconv.Itoa(v)
|
||||
case int64:
|
||||
return strconv.FormatInt(v, 10)
|
||||
case float64:
|
||||
return strconv.FormatFloat(v, 'f', -1, 64)
|
||||
default:
|
||||
return fmt.Sprint(v)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,163 @@
|
||||
package cmautobuyimport
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"sort"
|
||||
"strings"
|
||||
|
||||
"go-admin/app/goauto/product"
|
||||
"go-admin/app/goauto/shopeeproduct"
|
||||
)
|
||||
|
||||
type legacyPDDData struct {
|
||||
Dimensions []struct {
|
||||
Key string `json:"key"`
|
||||
Name string `json:"name"`
|
||||
} `json:"dimensions"`
|
||||
SKUs []struct {
|
||||
Options map[string]string `json:"options"`
|
||||
PriceCent *int64 `json:"price_cent"`
|
||||
Available bool `json:"available"`
|
||||
} `json:"skus"`
|
||||
}
|
||||
|
||||
func ConvertPDD(source SourcePDDProduct) (PDDRecord, error) {
|
||||
canonical, parsedID, err := product.NormalizeURL(source.URL)
|
||||
if err != nil || parsedID != strings.TrimSpace(source.GoodsID) {
|
||||
return PDDRecord{}, fmt.Errorf("URL 与 goods_id 不一致或无效")
|
||||
}
|
||||
if len([]rune(strings.TrimSpace(source.Title))) > 500 || len([]rune(strings.TrimSpace(source.ShopName))) > 255 {
|
||||
return PDDRecord{}, fmt.Errorf("标题或店铺名称超过目标字段长度")
|
||||
}
|
||||
|
||||
specs := []product.SpecDimension{}
|
||||
priceConflict := false
|
||||
if strings.TrimSpace(source.SKUsJSON) != "" {
|
||||
var legacy legacyPDDData
|
||||
if err := json.Unmarshal([]byte(source.SKUsJSON), &legacy); err != nil {
|
||||
return PDDRecord{}, fmt.Errorf("skus_json 不是合法 JSON: %w", err)
|
||||
}
|
||||
for _, dimension := range legacy.Dimensions {
|
||||
key := strings.TrimSpace(dimension.Key)
|
||||
name := strings.TrimSpace(dimension.Name)
|
||||
if name == "" {
|
||||
name = key
|
||||
}
|
||||
role := roleFor(key, name)
|
||||
values := make([]product.SpecValue, 0)
|
||||
seen := map[string]bool{}
|
||||
for _, sku := range legacy.SKUs {
|
||||
value := strings.TrimSpace(sku.Options[key])
|
||||
if value == "" || seen[value] {
|
||||
continue
|
||||
}
|
||||
seen[value] = true
|
||||
selectable := false
|
||||
prices := map[int64]bool{}
|
||||
for _, candidate := range legacy.SKUs {
|
||||
if strings.TrimSpace(candidate.Options[key]) != value {
|
||||
continue
|
||||
}
|
||||
selectable = selectable || candidate.Available
|
||||
if role == "color" && candidate.Available && candidate.PriceCent != nil && *candidate.PriceCent >= 0 {
|
||||
prices[*candidate.PriceCent] = true
|
||||
}
|
||||
}
|
||||
var price *int64
|
||||
if role == "color" && len(prices) == 1 {
|
||||
for amount := range prices {
|
||||
v := amount
|
||||
price = &v
|
||||
}
|
||||
} else if role == "color" && len(prices) > 1 {
|
||||
priceConflict = true
|
||||
}
|
||||
values = append(values, product.SpecValue{Name: value, Selectable: selectable, PriceCent: price})
|
||||
}
|
||||
if len(values) > 0 {
|
||||
specs = append(specs, product.SpecDimension{Name: name, Role: role, Values: values})
|
||||
}
|
||||
}
|
||||
}
|
||||
if err := product.ValidateSpecs(specs); err != nil {
|
||||
return PDDRecord{}, fmt.Errorf("规格无法转换: %w", err)
|
||||
}
|
||||
raw, _ := json.Marshal(specs)
|
||||
status := "pending"
|
||||
if source.CollectStatus == "collected" && len(specs) > 0 && !priceConflict {
|
||||
status = "active"
|
||||
}
|
||||
return PDDRecord{GoodsID: parsedID, URL: canonical, Title: strings.TrimSpace(source.Title),
|
||||
ShopName: strings.TrimSpace(source.ShopName), Status: status, SpecsJSON: string(raw), PriceConflict: priceConflict}, nil
|
||||
}
|
||||
|
||||
func ConvertShopee(source SourceShopeeProduct, currency string) (ShopeeRecord, error) {
|
||||
itemID := strings.TrimSpace(source.GoodsID)
|
||||
if itemID == "" || len(itemID) > 64 {
|
||||
return ShopeeRecord{}, fmt.Errorf("蝦皮商品 ID 为空或超过 64 个字符")
|
||||
}
|
||||
if len([]rune(strings.TrimSpace(source.Title))) > 500 || len([]rune(strings.TrimSpace(source.ShopName))) > 255 {
|
||||
return ShopeeRecord{}, fmt.Errorf("标题或店铺名称超过目标字段长度")
|
||||
}
|
||||
colors, sizes := map[string]string{}, map[string]string{}
|
||||
skipped := 0
|
||||
for _, sku := range source.SKUs {
|
||||
if !sku.ParseOK {
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
origin := shopeeproduct.ValueSourceImport
|
||||
if sku.IsManual {
|
||||
origin = shopeeproduct.ValueSourceManual
|
||||
}
|
||||
if value := strings.TrimSpace(sku.Color); value != "" {
|
||||
if colors[value] != shopeeproduct.ValueSourceManual {
|
||||
colors[value] = origin
|
||||
}
|
||||
}
|
||||
if value := strings.TrimSpace(sku.Size); value != "" {
|
||||
if sizes[value] != shopeeproduct.ValueSourceManual {
|
||||
sizes[value] = origin
|
||||
}
|
||||
}
|
||||
}
|
||||
specs := []shopeeproduct.SpecDimension{}
|
||||
if values := shopeeValues(colors); len(values) > 0 {
|
||||
specs = append(specs, shopeeproduct.SpecDimension{Name: "颜色", Role: shopeeproduct.RoleColor, Values: values})
|
||||
}
|
||||
if values := shopeeValues(sizes); len(values) > 0 {
|
||||
specs = append(specs, shopeeproduct.SpecDimension{Name: "尺码", Role: shopeeproduct.RoleSize, Values: values})
|
||||
}
|
||||
if err := shopeeproduct.Validate(specs); err != nil {
|
||||
return ShopeeRecord{}, fmt.Errorf("规格无法转换: %w", err)
|
||||
}
|
||||
raw, _ := shopeeproduct.Marshal(specs)
|
||||
return ShopeeRecord{ItemID: itemID, Title: strings.TrimSpace(source.Title), ShopName: strings.TrimSpace(source.ShopName),
|
||||
ImageURL: strings.TrimSpace(source.ImageURL), Currency: currency, SpecsJSON: raw,
|
||||
PDDGoodsID: strings.TrimSpace(source.PDDGoodsID), SkippedSKU: skipped}, nil
|
||||
}
|
||||
|
||||
func roleFor(key, name string) string {
|
||||
joined := strings.ToLower(key + " " + name)
|
||||
if strings.Contains(joined, "color") || strings.Contains(joined, "颜色") || strings.Contains(joined, "顏色") {
|
||||
return "color"
|
||||
}
|
||||
if strings.Contains(joined, "size") || strings.Contains(joined, "尺码") || strings.Contains(joined, "尺碼") {
|
||||
return "size"
|
||||
}
|
||||
return "other"
|
||||
}
|
||||
|
||||
func shopeeValues(source map[string]string) []shopeeproduct.SpecValue {
|
||||
names := make([]string, 0, len(source))
|
||||
for name := range source {
|
||||
names = append(names, name)
|
||||
}
|
||||
sort.Strings(names)
|
||||
values := make([]shopeeproduct.SpecValue, 0, len(names))
|
||||
for _, name := range names {
|
||||
values = append(values, shopeeproduct.SpecValue{Name: name, Source: source[name]})
|
||||
}
|
||||
return values
|
||||
}
|
||||
@@ -0,0 +1,95 @@
|
||||
package cmautobuyimport
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
|
||||
"go-admin/app/goauto/product"
|
||||
"go-admin/app/goauto/shopeeproduct"
|
||||
)
|
||||
|
||||
func TestConvertPDDAggregatesLegacySKUCombinations(t *testing.T) {
|
||||
source := SourcePDDProduct{GoodsID: "719834019024", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=719834019024&utm_source=old",
|
||||
Title: "商品", ShopName: "店铺", CollectStatus: "collected", SKUsJSON: `{
|
||||
"dimensions":[{"key":"color","name":"颜色分类"},{"key":"size","name":"尺码"}],
|
||||
"skus":[
|
||||
{"options":{"color":"黑色","size":"M"},"price_cent":2000,"available":true},
|
||||
{"options":{"color":"黑色","size":"L"},"price_cent":2000,"available":true},
|
||||
{"options":{"color":"白色","size":"M"},"price_cent":2200,"available":false}
|
||||
]}`}
|
||||
record, err := ConvertPDD(source)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if record.Status != "active" || record.PriceConflict || record.URL != "https://mobile.yangkeduo.com/goods.html?goods_id=719834019024" {
|
||||
t.Fatalf("unexpected record: %+v", record)
|
||||
}
|
||||
var specs []product.SpecDimension
|
||||
if err := json.Unmarshal([]byte(record.SpecsJSON), &specs); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(specs) != 2 || specs[0].Role != "color" || len(specs[0].Values) != 2 || specs[0].Values[0].PriceCent == nil || *specs[0].Values[0].PriceCent != 2000 {
|
||||
t.Fatalf("unexpected specs: %+v", specs)
|
||||
}
|
||||
if specs[0].Values[1].Selectable || specs[0].Values[1].PriceCent != nil {
|
||||
t.Fatalf("unavailable color must not get a purchasing price: %+v", specs[0].Values[1])
|
||||
}
|
||||
}
|
||||
|
||||
func TestConvertPDDDoesNotGuessConflictingColorPrice(t *testing.T) {
|
||||
record, err := ConvertPDD(SourcePDDProduct{GoodsID: "719834019024", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=719834019024",
|
||||
CollectStatus: "collected", SKUsJSON: `{"dimensions":[{"key":"color","name":"颜色"},{"key":"size","name":"尺码"}],"skus":[
|
||||
{"options":{"color":"黑色","size":"M"},"price_cent":2000,"available":true},
|
||||
{"options":{"color":"黑色","size":"L"},"price_cent":2100,"available":true}]}`})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !record.PriceConflict || record.Status != "pending" {
|
||||
t.Fatalf("conflicting prices must keep product pending: %+v", record)
|
||||
}
|
||||
var specs []product.SpecDimension
|
||||
_ = json.Unmarshal([]byte(record.SpecsJSON), &specs)
|
||||
if specs[0].Values[0].PriceCent != nil {
|
||||
t.Fatal("conflicting price must be nil")
|
||||
}
|
||||
}
|
||||
|
||||
func TestConvertPDDRejectsGoodsIDMismatch(t *testing.T) {
|
||||
_, err := ConvertPDD(SourcePDDProduct{GoodsID: "719834019024", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=972800403573"})
|
||||
if err == nil {
|
||||
t.Fatal("mismatched goods_id was accepted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestConvertShopeeAggregatesSpecsAndKeepsManualOrigin(t *testing.T) {
|
||||
record, err := ConvertShopee(SourceShopeeProduct{GoodsID: "26154802794", Title: "蝦皮商品", PDDGoodsID: "719834019024", SKUs: []SourceShopeeSKU{
|
||||
{Color: "黑色", Size: "XL", ParseOK: true},
|
||||
{Color: "黑色", Size: "2XL", ParseOK: true, IsManual: true},
|
||||
{Color: "白色", Size: "L", ParseOK: false},
|
||||
}}, "TWD")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if record.SkippedSKU != 1 || record.PDDGoodsID != "719834019024" {
|
||||
t.Fatalf("unexpected record: %+v", record)
|
||||
}
|
||||
var specs []shopeeproduct.SpecDimension
|
||||
if err := json.Unmarshal([]byte(record.SpecsJSON), &specs); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(specs) != 2 || len(specs[0].Values) != 1 || specs[0].Values[0].Name != "黑色" {
|
||||
t.Fatalf("unexpected specs: %+v", specs)
|
||||
}
|
||||
if specs[1].Values[1].Name != "XL" && specs[1].Values[0].Name != "2XL" {
|
||||
t.Fatalf("sizes were not aggregated: %+v", specs[1].Values)
|
||||
}
|
||||
manualFound := false
|
||||
for _, value := range specs[1].Values {
|
||||
if value.Name == "2XL" && value.Source == shopeeproduct.ValueSourceManual {
|
||||
manualFound = true
|
||||
}
|
||||
}
|
||||
if !manualFound {
|
||||
t.Fatal("manual source was lost")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,246 @@
|
||||
package cmautobuyimport
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
var targetColumns = map[string][]string{
|
||||
"pdd_product": {"id", "goods_id", "url", "title", "shop_name", "sales_count", "review_count", "status", "specs_json", "created_at", "updated_at"},
|
||||
"shopee_product": {"id", "shopee_item_id", "title", "shop_name", "pdd_product_id", "image_url", "sale_price_cent", "currency", "specs_json", "deleted_at", "deleted_flag", "created_at", "updated_at"},
|
||||
}
|
||||
|
||||
type plannedPDD struct {
|
||||
record PDDRecord
|
||||
existing bool
|
||||
disabled bool
|
||||
}
|
||||
|
||||
type plannedShopee struct {
|
||||
record ShopeeRecord
|
||||
existing bool
|
||||
}
|
||||
|
||||
func ResolveTargetCurrency(ctx context.Context, target *sql.DB) (string, error) {
|
||||
var exists int
|
||||
if err := target.QueryRowContext(ctx, `SELECT COUNT(*) FROM information_schema.tables
|
||||
WHERE table_schema=DATABASE() AND table_name='sys_config'`).Scan(&exists); err != nil {
|
||||
return "", fmt.Errorf("检查目标币种配置失败: %w", err)
|
||||
}
|
||||
if exists == 0 {
|
||||
return "TWD", nil
|
||||
}
|
||||
var value sql.NullString
|
||||
err := target.QueryRowContext(ctx, `SELECT config_value FROM sys_config WHERE config_key='shopee_default_currency' LIMIT 1`).Scan(&value)
|
||||
if err == sql.ErrNoRows || strings.TrimSpace(value.String) == "" {
|
||||
return "TWD", nil
|
||||
}
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("读取目标币种配置失败: %w", err)
|
||||
}
|
||||
currency := strings.ToUpper(strings.TrimSpace(value.String))
|
||||
if len(currency) != 3 {
|
||||
return "", fmt.Errorf("目标蝦皮默认币种配置无效")
|
||||
}
|
||||
return currency, nil
|
||||
}
|
||||
|
||||
// Run plans the complete import first. Dry-run returns the same counts and
|
||||
// conflicts as apply but performs no write. Apply writes both product domains
|
||||
// in one target transaction and rolls back on any error.
|
||||
func Run(ctx context.Context, target *sql.DB, dataset Dataset, options Options) (Report, error) {
|
||||
if err := checkColumns(ctx, target, targetColumns); err != nil {
|
||||
return Report{}, fmt.Errorf("目标数据库结构不兼容: %w", err)
|
||||
}
|
||||
now := options.Now
|
||||
if now == nil {
|
||||
now = time.Now
|
||||
}
|
||||
currency := strings.ToUpper(strings.TrimSpace(options.Currency))
|
||||
if currency == "" {
|
||||
currency = "TWD"
|
||||
}
|
||||
if len(currency) != 3 {
|
||||
return Report{}, fmt.Errorf("目标币种必须是 3 位 ISO 4217 代码")
|
||||
}
|
||||
report := Report{Mode: "dry-run", GeneratedAt: now().UTC(), Issues: []Issue{}}
|
||||
if options.Apply {
|
||||
report.Mode = "apply"
|
||||
}
|
||||
report.PDD.Source, report.Shopee.Source = len(dataset.PDD), len(dataset.Shopee)
|
||||
|
||||
pddPlan := make([]plannedPDD, 0, len(dataset.PDD))
|
||||
validPDD := map[string]bool{}
|
||||
for _, source := range dataset.PDD {
|
||||
record, err := ConvertPDD(source)
|
||||
if err != nil {
|
||||
report.PDD.Conflicts++
|
||||
report.Issues = append(report.Issues, Issue{Entity: "pdd", ID: source.GoodsID, Code: "PDD_INVALID", Reason: err.Error()})
|
||||
continue
|
||||
}
|
||||
var count int
|
||||
var status sql.NullString
|
||||
err = target.QueryRowContext(ctx, `SELECT COUNT(*),MAX(status) FROM pdd_product WHERE goods_id=?`, record.GoodsID).Scan(&count, &status)
|
||||
if err != nil {
|
||||
return Report{}, fmt.Errorf("检查目标 PDD 商品失败: %w", err)
|
||||
}
|
||||
plan := plannedPDD{record: record, existing: count > 0, disabled: status.String == "disabled"}
|
||||
if plan.existing {
|
||||
report.PDD.Overwrote++
|
||||
} else {
|
||||
report.PDD.Created++
|
||||
}
|
||||
if record.PriceConflict {
|
||||
report.PriceConflicts++
|
||||
report.Issues = append(report.Issues, Issue{Entity: "pdd", ID: record.GoodsID, Code: "PDD_COLOR_PRICE_CONFLICT", Reason: "同一颜色存在多个价格,价格已留空且商品保持待采集"})
|
||||
}
|
||||
validPDD[record.GoodsID] = true
|
||||
pddPlan = append(pddPlan, plan)
|
||||
}
|
||||
|
||||
shopeePlan := make([]plannedShopee, 0, len(dataset.Shopee))
|
||||
for _, source := range dataset.Shopee {
|
||||
record, err := ConvertShopee(source, currency)
|
||||
if err != nil {
|
||||
report.Shopee.Conflicts++
|
||||
report.Issues = append(report.Issues, Issue{Entity: "shopee", ID: source.GoodsID, Code: "SHOPEE_INVALID", Reason: err.Error()})
|
||||
continue
|
||||
}
|
||||
if record.PDDGoodsID != "" && !validPDD[record.PDDGoodsID] {
|
||||
report.Shopee.Conflicts++
|
||||
report.Issues = append(report.Issues, Issue{Entity: "shopee", ID: record.ItemID, Code: "PDD_ASSOCIATION_INVALID", Reason: "来源关联的 PDD 商品不存在、已删除或无法转换"})
|
||||
continue
|
||||
}
|
||||
var live, deleted int
|
||||
if err := target.QueryRowContext(ctx, `SELECT
|
||||
COALESCE(SUM(CASE WHEN deleted_at IS NULL THEN 1 ELSE 0 END),0),
|
||||
COALESCE(SUM(CASE WHEN deleted_at IS NOT NULL THEN 1 ELSE 0 END),0)
|
||||
FROM shopee_product WHERE shopee_item_id=?`, record.ItemID).Scan(&live, &deleted); err != nil {
|
||||
return Report{}, fmt.Errorf("检查目标蝦皮商品失败: %w", err)
|
||||
}
|
||||
if live == 0 && deleted > 0 {
|
||||
report.Shopee.Conflicts++
|
||||
report.Issues = append(report.Issues, Issue{Entity: "shopee", ID: record.ItemID, Code: "SHOPEE_TARGET_DELETED", Reason: "目标仅有软删除记录,未自动恢复"})
|
||||
continue
|
||||
}
|
||||
plan := plannedShopee{record: record, existing: live > 0}
|
||||
if plan.existing {
|
||||
report.Shopee.Overwrote++
|
||||
} else {
|
||||
report.Shopee.Created++
|
||||
}
|
||||
report.UnparsedShopeeSKU += record.SkippedSKU
|
||||
if record.SkippedSKU > 0 {
|
||||
report.Issues = append(report.Issues, Issue{Entity: "shopee", ID: record.ItemID, Code: "SHOPEE_SKU_SKIPPED", Reason: fmt.Sprintf("%d 条未解析规格未导入", record.SkippedSKU)})
|
||||
}
|
||||
if record.PDDGoodsID != "" {
|
||||
report.Associations++
|
||||
}
|
||||
shopeePlan = append(shopeePlan, plan)
|
||||
}
|
||||
if !options.Apply {
|
||||
return report, nil
|
||||
}
|
||||
if err := applyPlan(ctx, target, pddPlan, shopeePlan, report.GeneratedAt); err != nil {
|
||||
return Report{}, err
|
||||
}
|
||||
return report, nil
|
||||
}
|
||||
|
||||
func applyPlan(ctx context.Context, target *sql.DB, pddPlan []plannedPDD, shopeePlan []plannedShopee, now time.Time) error {
|
||||
tx, err := target.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelReadCommitted})
|
||||
if err != nil {
|
||||
return fmt.Errorf("开始目标事务失败: %w", err)
|
||||
}
|
||||
defer tx.Rollback()
|
||||
for _, item := range pddPlan {
|
||||
status := item.record.Status
|
||||
if item.disabled {
|
||||
status = "disabled"
|
||||
}
|
||||
_, err := tx.ExecContext(ctx, `INSERT INTO pdd_product
|
||||
(goods_id,url,title,shop_name,sales_count,review_count,status,specs_json,created_at,updated_at)
|
||||
VALUES (?,?,?,?,NULL,NULL,?,?,?,?)
|
||||
ON DUPLICATE KEY UPDATE url=VALUES(url),title=VALUES(title),shop_name=VALUES(shop_name),
|
||||
sales_count=NULL,review_count=NULL,status=IF(status='disabled','disabled',VALUES(status)),
|
||||
specs_json=VALUES(specs_json),updated_at=VALUES(updated_at)`, item.record.GoodsID, item.record.URL,
|
||||
item.record.Title, item.record.ShopName, status, item.record.SpecsJSON, now, now)
|
||||
if err != nil {
|
||||
return fmt.Errorf("写入 PDD 商品 %s 失败: %w", item.record.GoodsID, err)
|
||||
}
|
||||
}
|
||||
pddIDs := map[string]uint64{}
|
||||
rows, err := tx.QueryContext(ctx, `SELECT id,goods_id FROM pdd_product`)
|
||||
if err != nil {
|
||||
return fmt.Errorf("读取目标 PDD ID 失败: %w", err)
|
||||
}
|
||||
for rows.Next() {
|
||||
var id uint64
|
||||
var goodsID string
|
||||
if err := rows.Scan(&id, &goodsID); err != nil {
|
||||
rows.Close()
|
||||
return err
|
||||
}
|
||||
pddIDs[goodsID] = id
|
||||
}
|
||||
rows.Close()
|
||||
for _, item := range shopeePlan {
|
||||
var pddID any
|
||||
if item.record.PDDGoodsID != "" {
|
||||
id, ok := pddIDs[item.record.PDDGoodsID]
|
||||
if !ok {
|
||||
return fmt.Errorf("蝦皮商品 %s 的 PDD 关联在目标事务中不存在", item.record.ItemID)
|
||||
}
|
||||
pddID = id
|
||||
}
|
||||
if item.existing {
|
||||
result, err := tx.ExecContext(ctx, `UPDATE shopee_product SET title=?,shop_name=?,pdd_product_id=?,image_url=?,
|
||||
sale_price_cent=NULL,currency=?,specs_json=?,updated_at=? WHERE shopee_item_id=? AND deleted_at IS NULL`,
|
||||
item.record.Title, item.record.ShopName, pddID, item.record.ImageURL, item.record.Currency,
|
||||
item.record.SpecsJSON, now, item.record.ItemID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("覆盖蝦皮商品 %s 失败: %w", item.record.ItemID, err)
|
||||
}
|
||||
if affected, _ := result.RowsAffected(); affected != 1 {
|
||||
return fmt.Errorf("覆盖蝦皮商品 %s 时目标记录发生变化", item.record.ItemID)
|
||||
}
|
||||
} else {
|
||||
_, err := tx.ExecContext(ctx, `INSERT INTO shopee_product
|
||||
(shopee_item_id,title,shop_name,pdd_product_id,image_url,sale_price_cent,currency,specs_json,
|
||||
deleted_flag,created_at,updated_at) VALUES (?,?,?,?,?,NULL,?,?,0,?,?)`, item.record.ItemID,
|
||||
item.record.Title, item.record.ShopName, pddID, item.record.ImageURL, item.record.Currency,
|
||||
item.record.SpecsJSON, now, now)
|
||||
if err != nil {
|
||||
return fmt.Errorf("新增蝦皮商品 %s 失败: %w", item.record.ItemID, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
if err := validateWrittenPlan(ctx, tx, pddPlan, shopeePlan); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return fmt.Errorf("提交目标事务失败: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateWrittenPlan(ctx context.Context, tx *sql.Tx, pddPlan []plannedPDD, shopeePlan []plannedShopee) error {
|
||||
for _, item := range pddPlan {
|
||||
var count int
|
||||
if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM pdd_product WHERE goods_id=? AND JSON_VALID(specs_json)`, item.record.GoodsID).Scan(&count); err != nil || count != 1 {
|
||||
return fmt.Errorf("PDD 商品 %s 提交前校验失败", item.record.GoodsID)
|
||||
}
|
||||
}
|
||||
for _, item := range shopeePlan {
|
||||
var count int
|
||||
if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM shopee_product sp LEFT JOIN pdd_product pp ON pp.id=sp.pdd_product_id
|
||||
WHERE sp.shopee_item_id=? AND sp.deleted_at IS NULL AND JSON_VALID(sp.specs_json)
|
||||
AND (sp.pdd_product_id IS NULL OR pp.id IS NOT NULL)`, item.record.ItemID).Scan(&count); err != nil || count != 1 {
|
||||
return fmt.Errorf("蝦皮商品 %s 提交前校验失败", item.record.ItemID)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,85 @@
|
||||
package cmautobuyimport
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestRunMySQL84DryRunAndApply(t *testing.T) {
|
||||
dsn := os.Getenv("GOAUTO_IMPORT_MYSQL_TEST_DSN")
|
||||
if dsn == "" {
|
||||
t.Skip("GOAUTO_IMPORT_MYSQL_TEST_DSN not set")
|
||||
}
|
||||
db, err := sql.Open("mysql", dsn)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
var databaseName string
|
||||
if err := db.QueryRow(`SELECT DATABASE()`).Scan(&databaseName); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !strings.HasSuffix(databaseName, "_test") {
|
||||
t.Fatalf("refusing destructive fixture setup outside *_test database: %s", databaseName)
|
||||
}
|
||||
for _, statement := range []string{
|
||||
`DROP TABLE IF EXISTS shopee_product`, `DROP TABLE IF EXISTS pdd_product`, `DROP TABLE IF EXISTS sys_config`,
|
||||
`CREATE TABLE pdd_product (
|
||||
id BIGINT UNSIGNED PRIMARY KEY AUTO_INCREMENT, goods_id VARCHAR(32) NOT NULL UNIQUE, url TEXT NOT NULL,
|
||||
title VARCHAR(500) NOT NULL DEFAULT '', shop_name VARCHAR(255) NOT NULL DEFAULT '', sales_count BIGINT NULL,
|
||||
review_count BIGINT NULL, status VARCHAR(16) NOT NULL, specs_json JSON NOT NULL,
|
||||
created_at DATETIME(3) NOT NULL, updated_at DATETIME(3) NOT NULL)`,
|
||||
`CREATE TABLE shopee_product (
|
||||
id BIGINT UNSIGNED PRIMARY KEY AUTO_INCREMENT, shopee_item_id VARCHAR(64) NOT NULL,
|
||||
title VARCHAR(500) NOT NULL DEFAULT '', shop_name VARCHAR(255) NOT NULL DEFAULT '', pdd_product_id BIGINT UNSIGNED NULL,
|
||||
image_url TEXT NOT NULL, sale_price_cent BIGINT NULL, currency VARCHAR(3) NOT NULL, specs_json JSON NOT NULL,
|
||||
last_create_request_id VARCHAR(36) NULL, last_update_request_id VARCHAR(36) NULL, deleted_at DATETIME(3) NULL,
|
||||
deleted_flag BIGINT UNSIGNED NOT NULL DEFAULT 0, deleted_by BIGINT UNSIGNED NULL,
|
||||
created_at DATETIME(3) NOT NULL, updated_at DATETIME(3) NOT NULL,
|
||||
UNIQUE KEY ux_shopee_product_item_id(shopee_item_id,deleted_flag))`,
|
||||
`CREATE TABLE sys_config (config_key VARCHAR(128) PRIMARY KEY, config_value TEXT)`,
|
||||
`INSERT INTO sys_config(config_key,config_value) VALUES ('shopee_default_currency','TWD')`,
|
||||
`INSERT INTO pdd_product(goods_id,url,title,shop_name,sales_count,review_count,status,specs_json,created_at,updated_at)
|
||||
VALUES ('719834019024','https://mobile.yangkeduo.com/goods.html?goods_id=719834019024','旧标题','旧店铺',9,8,'active','[]',NOW(3),NOW(3))`,
|
||||
`INSERT INTO shopee_product(shopee_item_id,title,shop_name,image_url,sale_price_cent,currency,specs_json,deleted_flag,created_at,updated_at)
|
||||
VALUES ('26154802794','旧蝦皮','旧店铺','',123,'TWD','[]',0,NOW(3),NOW(3))`,
|
||||
} {
|
||||
if _, err := db.Exec(statement); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
dataset := Dataset{
|
||||
PDD: []SourcePDDProduct{{GoodsID: "719834019024", URL: "https://mobile.yangkeduo.com/goods.html?goods_id=719834019024", Title: "新标题", ShopName: "新店铺", CollectStatus: "collected", SKUsJSON: `{"dimensions":[{"key":"color","name":"颜色"}],"skus":[{"options":{"color":"黑色"},"price_cent":2000,"available":true}]}`}},
|
||||
Shopee: []SourceShopeeProduct{{GoodsID: "26154802794", Title: "新蝦皮", ShopName: "新店铺", PDDGoodsID: "719834019024", SKUs: []SourceShopeeSKU{{Color: "黑色", ParseOK: true}}}},
|
||||
}
|
||||
ctx := context.Background()
|
||||
dry, err := Run(ctx, db, dataset, Options{Currency: "TWD", Now: func() time.Time { return time.Unix(1, 0) }})
|
||||
if err != nil || dry.PDD.Overwrote != 1 || dry.Shopee.Overwrote != 1 {
|
||||
t.Fatalf("dry-run=%+v err=%v", dry, err)
|
||||
}
|
||||
var title string
|
||||
if err := db.QueryRow(`SELECT title FROM pdd_product WHERE goods_id='719834019024'`).Scan(&title); err != nil || title != "旧标题" {
|
||||
t.Fatalf("dry-run wrote target: title=%q err=%v", title, err)
|
||||
}
|
||||
if _, err := Run(ctx, db, dataset, Options{Apply: true, Currency: "TWD", Now: func() time.Time { return time.Unix(2, 0) }}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var sales, reviews, price sql.NullInt64
|
||||
var linked uint64
|
||||
if err := db.QueryRow(`SELECT title,sales_count,review_count FROM pdd_product WHERE goods_id='719834019024'`).Scan(&title, &sales, &reviews); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if title != "新标题" || sales.Valid || reviews.Valid {
|
||||
t.Fatalf("PDD overwrite failed: title=%q sales=%v reviews=%v", title, sales, reviews)
|
||||
}
|
||||
if err := db.QueryRow(`SELECT sp.title,sp.sale_price_cent,sp.pdd_product_id FROM shopee_product sp WHERE shopee_item_id='26154802794' AND deleted_at IS NULL`).Scan(&title, &price, &linked); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if title != "新蝦皮" || price.Valid || linked == 0 {
|
||||
t.Fatalf("Shopee overwrite failed: title=%q price=%v linked=%d", title, price, linked)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
package cmautobuyimport
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"strings"
|
||||
)
|
||||
|
||||
var sourceColumns = map[string][]string{
|
||||
"pdd_products": {"goods_id", "url", "title", "shop_name", "skus_json", "collect_status", "deleted_at"},
|
||||
"shopee_products": {"goods_id", "title", "image_url", "shopee_shop_name", "pdd_goods_id", "deleted_at"},
|
||||
"shopee_skus": {"goods_id", "color", "size", "parse_ok", "is_manual"},
|
||||
}
|
||||
|
||||
func LoadSource(ctx context.Context, db *sql.DB) (Dataset, error) {
|
||||
if err := checkColumns(ctx, db, sourceColumns); err != nil {
|
||||
return Dataset{}, fmt.Errorf("来源数据库结构不兼容: %w", err)
|
||||
}
|
||||
tx, err := db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelRepeatableRead, ReadOnly: true})
|
||||
if err != nil {
|
||||
return Dataset{}, fmt.Errorf("开始来源只读事务失败: %w", err)
|
||||
}
|
||||
defer tx.Rollback()
|
||||
dataset := Dataset{}
|
||||
rows, err := tx.QueryContext(ctx, `SELECT goods_id,url,COALESCE(title,''),COALESCE(shop_name,''),
|
||||
COALESCE(skus_json,''),collect_status FROM pdd_products WHERE deleted_at IS NULL ORDER BY goods_id`)
|
||||
if err != nil {
|
||||
return Dataset{}, fmt.Errorf("读取来源 PDD 商品失败: %w", err)
|
||||
}
|
||||
for rows.Next() {
|
||||
var item SourcePDDProduct
|
||||
if err := rows.Scan(&item.GoodsID, &item.URL, &item.Title, &item.ShopName, &item.SKUsJSON, &item.CollectStatus); err != nil {
|
||||
rows.Close()
|
||||
return Dataset{}, fmt.Errorf("读取来源 PDD 商品行失败: %w", err)
|
||||
}
|
||||
dataset.PDD = append(dataset.PDD, item)
|
||||
}
|
||||
if err := rows.Close(); err != nil {
|
||||
return Dataset{}, err
|
||||
}
|
||||
|
||||
products := map[string]int{}
|
||||
rows, err = tx.QueryContext(ctx, `SELECT goods_id,title,COALESCE(shopee_shop_name,''),COALESCE(image_url,''),
|
||||
COALESCE(pdd_goods_id,'') FROM shopee_products WHERE deleted_at IS NULL ORDER BY goods_id`)
|
||||
if err != nil {
|
||||
return Dataset{}, fmt.Errorf("读取来源蝦皮商品失败: %w", err)
|
||||
}
|
||||
for rows.Next() {
|
||||
var item SourceShopeeProduct
|
||||
if err := rows.Scan(&item.GoodsID, &item.Title, &item.ShopName, &item.ImageURL, &item.PDDGoodsID); err != nil {
|
||||
rows.Close()
|
||||
return Dataset{}, fmt.Errorf("读取来源蝦皮商品行失败: %w", err)
|
||||
}
|
||||
dataset.Shopee = append(dataset.Shopee, item)
|
||||
products[item.GoodsID] = len(dataset.Shopee) - 1
|
||||
}
|
||||
if err := rows.Close(); err != nil {
|
||||
return Dataset{}, err
|
||||
}
|
||||
rows, err = tx.QueryContext(ctx, `SELECT goods_id,COALESCE(color,''),COALESCE(size,''),parse_ok,is_manual
|
||||
FROM shopee_skus ORDER BY goods_id,sku_id`)
|
||||
if err != nil {
|
||||
return Dataset{}, fmt.Errorf("读取来源蝦皮规格失败: %w", err)
|
||||
}
|
||||
for rows.Next() {
|
||||
var goodsID string
|
||||
var sku SourceShopeeSKU
|
||||
if err := rows.Scan(&goodsID, &sku.Color, &sku.Size, &sku.ParseOK, &sku.IsManual); err != nil {
|
||||
rows.Close()
|
||||
return Dataset{}, fmt.Errorf("读取来源蝦皮规格行失败: %w", err)
|
||||
}
|
||||
if index, ok := products[goodsID]; ok {
|
||||
dataset.Shopee[index].SKUs = append(dataset.Shopee[index].SKUs, sku)
|
||||
}
|
||||
}
|
||||
if err := rows.Close(); err != nil {
|
||||
return Dataset{}, err
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return Dataset{}, fmt.Errorf("结束来源只读事务失败: %w", err)
|
||||
}
|
||||
return dataset, nil
|
||||
}
|
||||
|
||||
func checkColumns(ctx context.Context, db *sql.DB, required map[string][]string) error {
|
||||
for table, columns := range required {
|
||||
rows, err := db.QueryContext(ctx, `SELECT column_name FROM information_schema.columns
|
||||
WHERE table_schema=DATABASE() AND table_name=?`, table)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
found := map[string]bool{}
|
||||
for rows.Next() {
|
||||
var column string
|
||||
if err := rows.Scan(&column); err != nil {
|
||||
rows.Close()
|
||||
return err
|
||||
}
|
||||
found[strings.ToLower(column)] = true
|
||||
}
|
||||
rows.Close()
|
||||
for _, column := range columns {
|
||||
if !found[column] {
|
||||
return fmt.Errorf("%s.%s 缺失", table, column)
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
// Package cmautobuyimport implements the #70 offline, one-way product import.
|
||||
// It deliberately has no HTTP route and is never called during server startup.
|
||||
package cmautobuyimport
|
||||
|
||||
import "time"
|
||||
|
||||
type SourcePDDProduct struct {
|
||||
GoodsID string
|
||||
URL string
|
||||
Title string
|
||||
ShopName string
|
||||
SKUsJSON string
|
||||
CollectStatus string
|
||||
}
|
||||
|
||||
type SourceShopeeProduct struct {
|
||||
GoodsID string
|
||||
Title string
|
||||
ShopName string
|
||||
ImageURL string
|
||||
PDDGoodsID string
|
||||
SKUs []SourceShopeeSKU
|
||||
}
|
||||
|
||||
type SourceShopeeSKU struct {
|
||||
Color string
|
||||
Size string
|
||||
ParseOK bool
|
||||
IsManual bool
|
||||
}
|
||||
|
||||
type Dataset struct {
|
||||
PDD []SourcePDDProduct
|
||||
Shopee []SourceShopeeProduct
|
||||
}
|
||||
|
||||
type PDDRecord struct {
|
||||
GoodsID string
|
||||
URL string
|
||||
Title string
|
||||
ShopName string
|
||||
Status string
|
||||
SpecsJSON string
|
||||
PriceConflict bool
|
||||
}
|
||||
|
||||
type ShopeeRecord struct {
|
||||
ItemID string
|
||||
Title string
|
||||
ShopName string
|
||||
ImageURL string
|
||||
Currency string
|
||||
SpecsJSON string
|
||||
PDDGoodsID string
|
||||
SkippedSKU int
|
||||
}
|
||||
|
||||
type Issue struct {
|
||||
Entity string `json:"entity"`
|
||||
ID string `json:"id"`
|
||||
Code string `json:"code"`
|
||||
Reason string `json:"reason"`
|
||||
}
|
||||
|
||||
type EntityCounts struct {
|
||||
Source int `json:"source"`
|
||||
Created int `json:"created"`
|
||||
Overwrote int `json:"overwrote"`
|
||||
Skipped int `json:"skipped"`
|
||||
Conflicts int `json:"conflicts"`
|
||||
}
|
||||
|
||||
type Report struct {
|
||||
Mode string `json:"mode"`
|
||||
GeneratedAt time.Time `json:"generatedAt"`
|
||||
PDD EntityCounts `json:"pdd"`
|
||||
Shopee EntityCounts `json:"shopee"`
|
||||
Associations int `json:"associations"`
|
||||
PriceConflicts int `json:"priceConflicts"`
|
||||
UnparsedShopeeSKU int `json:"unparsedShopeeSku"`
|
||||
Issues []Issue `json:"issues"`
|
||||
}
|
||||
|
||||
type Options struct {
|
||||
Apply bool
|
||||
Currency string
|
||||
Now func() time.Time
|
||||
}
|
||||
@@ -0,0 +1,93 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"flag"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"go-admin/app/goauto/cmautobuyimport"
|
||||
)
|
||||
|
||||
func main() {
|
||||
var sourceConfig, targetConfig string
|
||||
var apply bool
|
||||
flag.StringVar(&sourceConfig, "source-config", os.Getenv("GOAUTO_CMAUTOBUY_CONFIG"), "cmautobuy 未跟踪 config.yaml 路径")
|
||||
flag.StringVar(&targetConfig, "target-config", os.Getenv("GOAUTO_CONFIG"), "GoAuto 未跟踪 config.yaml 路径")
|
||||
flag.BoolVar(&apply, "apply", false, "确认写入目标库;缺省只执行 dry-run")
|
||||
flag.Parse()
|
||||
if sourceConfig == "" || targetConfig == "" {
|
||||
fatal("必须通过 --source-config/--target-config 或对应环境变量提供两个未跟踪配置文件")
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute)
|
||||
defer cancel()
|
||||
sourceDSN, err := cmautobuyimport.DSNFromConfig(sourceConfig, true)
|
||||
if err != nil {
|
||||
fatal(err.Error())
|
||||
}
|
||||
targetDSN, err := cmautobuyimport.DSNFromConfig(targetConfig, false)
|
||||
if err != nil {
|
||||
fatal(err.Error())
|
||||
}
|
||||
source, err := sql.Open("mysql", sourceDSN)
|
||||
if err != nil {
|
||||
fatal("打开来源数据库失败")
|
||||
}
|
||||
defer source.Close()
|
||||
target, err := sql.Open("mysql", targetDSN)
|
||||
if err != nil {
|
||||
fatal("打开目标数据库失败")
|
||||
}
|
||||
defer target.Close()
|
||||
if err := source.PingContext(ctx); err != nil {
|
||||
fatal("连接来源数据库失败:" + safeConnectionError(err))
|
||||
}
|
||||
if err := target.PingContext(ctx); err != nil {
|
||||
fatal("连接目标数据库失败:" + safeConnectionError(err))
|
||||
}
|
||||
dataset, err := cmautobuyimport.LoadSource(ctx, source)
|
||||
if err != nil {
|
||||
fatal(err.Error())
|
||||
}
|
||||
currency, err := cmautobuyimport.ResolveTargetCurrency(ctx, target)
|
||||
if err != nil {
|
||||
fatal(err.Error())
|
||||
}
|
||||
report, err := cmautobuyimport.Run(ctx, target, dataset, cmautobuyimport.Options{Apply: apply, Currency: currency})
|
||||
if err != nil {
|
||||
fatal(err.Error())
|
||||
}
|
||||
raw, err := json.MarshalIndent(report, "", " ")
|
||||
if err != nil {
|
||||
fatal("生成导入报告失败")
|
||||
}
|
||||
fmt.Println(string(raw))
|
||||
if !apply {
|
||||
fmt.Println("dry-run 完成:没有写入数据库。查看报告并再次确认后才能使用 --apply。")
|
||||
}
|
||||
}
|
||||
|
||||
func safeConnectionError(err error) string {
|
||||
message := strings.ToLower(err.Error())
|
||||
switch {
|
||||
case strings.Contains(message, "access denied") || strings.Contains(message, "authentication"):
|
||||
return "认证失败(凭据和 DSN 未输出)"
|
||||
case strings.Contains(message, "certificate") || strings.Contains(message, "tls") || strings.Contains(message, "x509"):
|
||||
return "TLS 证书校验失败(凭据和 DSN 未输出)"
|
||||
case strings.Contains(message, "refused"):
|
||||
return "服务器拒绝连接(凭据和 DSN 未输出)"
|
||||
case strings.Contains(message, "timeout") || strings.Contains(message, "deadline"):
|
||||
return "连接超时(凭据和 DSN 未输出)"
|
||||
default:
|
||||
return "网络或数据库连接异常(凭据和 DSN 未输出)"
|
||||
}
|
||||
}
|
||||
|
||||
func fatal(message string) {
|
||||
fmt.Fprintln(os.Stderr, "导入失败:"+message)
|
||||
os.Exit(1)
|
||||
}
|
||||
Reference in New Issue
Block a user