From 166adc6fa84fc51457169426a1590a47a2b6f2b0 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Sun, 23 Aug 2026 23:59:12 +0800 Subject: [PATCH] feat(#70): import cmautobuy product catalog --- docs/02-architecture-and-code-map.md | 13 +- docs/03-business-rules-and-glossary.md | 15 +- docs/04-local-development-and-verification.md | 27 +- docs/09-delivery-issues.md | 11 +- server/app/goauto/cmautobuyimport/config.go | 107 ++++++++ server/app/goauto/cmautobuyimport/convert.go | 163 ++++++++++++ .../goauto/cmautobuyimport/convert_test.go | 95 +++++++ server/app/goauto/cmautobuyimport/importer.go | 246 ++++++++++++++++++ .../cmautobuyimport/mysql_integration_test.go | 85 ++++++ server/app/goauto/cmautobuyimport/source.go | 110 ++++++++ server/app/goauto/cmautobuyimport/types.go | 88 +++++++ server/cmd/import-cmautobuy-products/main.go | 93 +++++++ 12 files changed, 1045 insertions(+), 8 deletions(-) create mode 100644 server/app/goauto/cmautobuyimport/config.go create mode 100644 server/app/goauto/cmautobuyimport/convert.go create mode 100644 server/app/goauto/cmautobuyimport/convert_test.go create mode 100644 server/app/goauto/cmautobuyimport/importer.go create mode 100644 server/app/goauto/cmautobuyimport/mysql_integration_test.go create mode 100644 server/app/goauto/cmautobuyimport/source.go create mode 100644 server/app/goauto/cmautobuyimport/types.go create mode 100644 server/cmd/import-cmautobuy-products/main.go diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md index 36c224d..3c4d05f 100644 --- a/docs/02-architecture-and-code-map.md +++ b/docs/02-architecture-and-code-map.md @@ -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 # 架构与代码地图 @@ -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、唯一键和关联完整性。 diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md index 92913be..9c06daa 100644 --- a/docs/03-business-rules-and-glossary.md +++ b/docs/03-business-rules-and-glossary.md @@ -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 # 业务规则与术语 @@ -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、唯一键和关联完整性。 diff --git a/docs/04-local-development-and-verification.md b/docs/04-local-development-and-verification.md index 20d9a20..1dd8170 100644 --- a/docs/04-local-development-and-verification.md +++ b/docs/04-local-development-and-verification.md @@ -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 # 本地开发与验证 @@ -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 +``` diff --git a/docs/09-delivery-issues.md b/docs/09-delivery-issues.md index a10336e..012ea06 100644 --- a/docs/09-delivery-issues.md +++ b/docs/09-delivery-issues.md @@ -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 # 当前 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 等待用户查看精确影响后再次确认 | diff --git a/server/app/goauto/cmautobuyimport/config.go b/server/app/goauto/cmautobuyimport/config.go new file mode 100644 index 0000000..b36d3e4 --- /dev/null +++ b/server/app/goauto/cmautobuyimport/config.go @@ -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) + } +} diff --git a/server/app/goauto/cmautobuyimport/convert.go b/server/app/goauto/cmautobuyimport/convert.go new file mode 100644 index 0000000..0c1c618 --- /dev/null +++ b/server/app/goauto/cmautobuyimport/convert.go @@ -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 +} diff --git a/server/app/goauto/cmautobuyimport/convert_test.go b/server/app/goauto/cmautobuyimport/convert_test.go new file mode 100644 index 0000000..2657938 --- /dev/null +++ b/server/app/goauto/cmautobuyimport/convert_test.go @@ -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") + } +} diff --git a/server/app/goauto/cmautobuyimport/importer.go b/server/app/goauto/cmautobuyimport/importer.go new file mode 100644 index 0000000..f112837 --- /dev/null +++ b/server/app/goauto/cmautobuyimport/importer.go @@ -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 +} diff --git a/server/app/goauto/cmautobuyimport/mysql_integration_test.go b/server/app/goauto/cmautobuyimport/mysql_integration_test.go new file mode 100644 index 0000000..4dfe61e --- /dev/null +++ b/server/app/goauto/cmautobuyimport/mysql_integration_test.go @@ -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) + } +} diff --git a/server/app/goauto/cmautobuyimport/source.go b/server/app/goauto/cmautobuyimport/source.go new file mode 100644 index 0000000..68330d3 --- /dev/null +++ b/server/app/goauto/cmautobuyimport/source.go @@ -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 +} diff --git a/server/app/goauto/cmautobuyimport/types.go b/server/app/goauto/cmautobuyimport/types.go new file mode 100644 index 0000000..46d1453 --- /dev/null +++ b/server/app/goauto/cmautobuyimport/types.go @@ -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 +} diff --git a/server/cmd/import-cmautobuy-products/main.go b/server/cmd/import-cmautobuy-products/main.go new file mode 100644 index 0000000..755eb5f --- /dev/null +++ b/server/cmd/import-cmautobuy-products/main.go @@ -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) +}