// Package repository 封装数据库读写。本文件只保留历史 SQLite v8 schema 与迁移, // 供一次性 SQLite→MySQL 工具和历史回归测试使用,不进入生产 Admin 启动路径。 // // 改动本文件前必读 admin/AGENTS.md。三条硬规则: // - 只有本包能写 SQL,handler 和 service 都不许拼 SQL; // - SQL 一律参数化查询(用 ? 占位),禁止字符串拼接; // - 迁移只能往前加,不许在启动时删库重建——data/ 在升级时是保留的, // 里面有人工填了几个月的 PDD 链接和 SKU 映射。 // // 表结构的权威定义在 docs/admin/03-data-model.md,改表要先改文档。 package repository import ( "context" "database/sql" "fmt" "log" "path/filepath" "strings" // 纯 Go 的 SQLite 驱动,注册的驱动名是 "sqlite"(不是 "sqlite3")。 // 不得换成 github.com/mattn/go-sqlite3,那个需要 cgo, // Windows 上要装 gcc,打包 exe 会变麻烦。理由见 admin/AGENTS.md。 _ "modernc.org/sqlite" ) // Execer 让同一个 repository 函数既能直接用 *sql.DB, // 也能在事务里用 *sql.Tx。 // // 需要"多张表要么一起改、要么都不改"时,service 开一个事务, // 把 *sql.Tx 传进来即可,不用为事务再写一套函数。 type Execer interface { Exec(query string, args ...any) (sql.Result, error) Query(query string, args ...any) (*sql.Rows, error) QueryRow(query string, args ...any) *sql.Row } // Open 打开 data/admin.db。 // // **PRAGMA 必须写在 DSN 里,不能用 db.Exec("PRAGMA ...") 设置。** // // 原因是 Go 的 database/sql 是一个**连接池**:db.Exec 只作用于当时 // 拿到的那一条连接,池子后来新开的连接完全没执行过那些 PRAGMA。 // 并发写的时候,没有 busy_timeout 的那些连接会直接报 // "database is locked (SQLITE_BUSY)",而不是等锁释放。 // // 写进 DSN 后,驱动会对**每一条新连接**都应用一遍,这才是对的。 func Open(dataDir string) (*sql.DB, error) { path := filepath.Join(dataDir, "admin.db") // 三条设置的含义见 docs/admin/03-data-model.md §2: // busy_timeout 拿不到锁时最多等 5 秒,而不是立刻报错 // journal_mode WAL 模式,读和写可以同时进行 // foreign_keys 打开外键约束(SQLite 默认是关的) // // _txlock=immediate 是另一个**必须加**的参数,原因见下。 dsn := "file:" + path + "?_pragma=busy_timeout(5000)" + "&_pragma=journal_mode(WAL)" + "&_pragma=foreign_keys(1)" + "&_txlock=immediate" db, err := sql.Open("sqlite", dsn) if err != nil { return nil, fmt.Errorf("打开数据库 %s 失败: %w", path, err) } // SQLite 同一时刻只允许一个写事务。连接数放太开, // 大量连接会互相抢锁、把 busy_timeout 耗光。 // 本项目是单机内部工具,并发量很小,限制在个位数足够。 // // 注意**不要设成 1**:那样在一个事务里再调用需要连接的代码会死锁。 db.SetMaxOpenConns(4) db.SetMaxIdleConns(4) // 关于 _txlock=immediate: // // Go 的 db.Begin() 默认发的是 BEGIN DEFERRED——事务开始时**不拿写锁**, // 等到第一次写才去拿。于是多个事务可以同时开始、各自先读, // 然后同时想升级成写,互相卡死,直接报 SQLITE_BUSY。 // 这种情况 busy_timeout **救不了**:等下去也不可能有结果, // 只能让某个事务整个重来。 // // 加上 _txlock=immediate 后,事务一开始就拿写锁, // 拿不到就按 busy_timeout 排队等——这才是我们要的行为。 // // 实测(6 个并发事务,每个先读后写): // 默认 deferred 失败 5/6 // _txlock=immediate 失败 0/6 // sql.Open 是懒加载的,这里主动连一次,好让配置错误立刻暴露 if err := db.Ping(); err != nil { db.Close() return nil, fmt.Errorf("连接数据库 %s 失败: %w", path, err) } return db, nil } // migrations 按顺序存放每一版的迁移语句。 // // 外层一个元素 = 一个版本;内层是该版本要执行的语句,**一条一执行**。 // 不把多条语句塞进一个字符串,是因为 database/sql 的 Exec 对 // "一次执行多条语句"的支持因驱动而异,拆开最稳妥。 // // 加新版本时**只能往末尾追加**,不许改动已有元素—— // 已经发布出去的库是按旧语句建的,改了会导致新旧库结构不一致。 // // [必须] 这条规则本身也写进了 admin/AGENTS.md:#20 就是因为 v1 被原地改写、 // 而老库的 user_version 已经越过了它,那次改写永远不会在老库上重跑, // 程序拿着一个和代码对不上的库静默启动。v3 不在这个 slice 里 // (见下面 schemaVersion 和 migrateV3 的注释),也是同一个教训的直接结果: // v3 要做的事超出"一串 SQL 顺序执行",硬塞进这个只支持纯 SQL 的结构反而 // 掩盖了它的特殊性。 var migrations = [][]string{ // v1: 初始表结构,对应 docs/admin/03-data-model.md { `CREATE TABLE shopee_products ( goods_id TEXT PRIMARY KEY, title TEXT NOT NULL, shopee_status TEXT, main_sku_code TEXT, -- 下面三个是人工维护的,报表里没有,导入时绝不能覆盖 pdd_goods_url TEXT, pdd_goods_id TEXT, pdd_data TEXT, collect_status TEXT NOT NULL DEFAULT 'no_link' CHECK (collect_status IN ( 'no_link', 'pending', 'collecting', 'collected', 'failed' )), collect_error TEXT, collected_at TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL );`, `CREATE INDEX idx_shopee_products_status ON shopee_products(collect_status);`, `CREATE TABLE shopee_skus ( sku_id TEXT PRIMARY KEY, goods_id TEXT NOT NULL, spec_raw TEXT NOT NULL, color TEXT, size TEXT, advice TEXT, parse_ok INTEGER NOT NULL DEFAULT 0, sku_code TEXT, is_manual INTEGER NOT NULL DEFAULT 0, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, FOREIGN KEY (goods_id) REFERENCES shopee_products(goods_id) ON DELETE CASCADE );`, `CREATE INDEX idx_shopee_skus_goods ON shopee_skus(goods_id);`, `CREATE INDEX idx_shopee_skus_parse ON shopee_skus(parse_ok);`, `CREATE TABLE syb_orders ( syb_id TEXT PRIMARY KEY, order_no TEXT NOT NULL, title TEXT, shopee_goods_id TEXT, shopee_sku_id TEXT, quantity INTEGER NOT NULL CHECK (quantity > 0), price_twd_cent INTEGER CHECK (price_twd_cent IS NULL OR price_twd_cent >= 0), image_url TEXT, syb_data TEXT NOT NULL DEFAULT '{}', created_at TEXT NOT NULL, updated_at TEXT NOT NULL );`, `CREATE INDEX idx_syb_orders_order ON syb_orders(order_no);`, `CREATE INDEX idx_syb_orders_goods ON syb_orders(shopee_goods_id);`, `CREATE INDEX idx_syb_orders_list ON syb_orders(updated_at DESC, syb_id DESC);`, `CREATE TABLE sku_mappings ( shopee_sku_id TEXT PRIMARY KEY, goods_id TEXT NOT NULL, pdd_options TEXT NOT NULL, mapped_at TEXT NOT NULL, mapped_by TEXT, FOREIGN KEY (shopee_sku_id) REFERENCES shopee_skus(sku_id) ON DELETE CASCADE );`, `CREATE INDEX idx_sku_mappings_goods ON sku_mappings(goods_id);`, `CREATE TABLE tasks ( task_id TEXT PRIMARY KEY, task_type TEXT NOT NULL CHECK (task_type IN ('collect', 'purchase')), status TEXT NOT NULL DEFAULT 'pending' CHECK (status IN ('pending', 'assigned', 'claimed', 'succeeded', 'manual_review', 'failed', 'cancelled')), version INTEGER NOT NULL DEFAULT 1 CHECK (version > 0), priority INTEGER NOT NULL DEFAULT 0, assigned_client TEXT, claimed_at TEXT, syb_id TEXT, order_no TEXT, goods_id TEXT, shopee_sku_id TEXT, -- Client 契约要求:pdd_goods_url 必填; -- 采购任务的 quantity 和 max_price_cent 也必填(价格保护) pdd_goods_url TEXT NOT NULL, pdd_goods_id TEXT, pdd_options TEXT, quantity INTEGER CHECK (quantity IS NULL OR quantity > 0), max_price_cent INTEGER CHECK (max_price_cent IS NULL OR max_price_cent > 0), result_data TEXT, error_code TEXT, error_message TEXT, finished_at TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL );`, `CREATE INDEX idx_tasks_claim ON tasks(assigned_client, status, priority DESC, created_at);`, `CREATE INDEX idx_tasks_list ON tasks(updated_at DESC, task_id DESC);`, `CREATE INDEX idx_tasks_order ON tasks(order_no);`, `CREATE TABLE clients ( client_id TEXT PRIMARY KEY, name TEXT, device_address TEXT, platform TEXT, pdd_package TEXT, capabilities TEXT, last_seen_at TEXT NOT NULL, created_at TEXT NOT NULL, updated_at TEXT NOT NULL );`, `CREATE TABLE idempotency_keys ( key TEXT PRIMARY KEY, request_hash TEXT NOT NULL, response_body TEXT NOT NULL, created_at TEXT NOT NULL );`, }, // v2: 领取历史。 // // 为什么需要它:契约要求"只有**从未分配给该客户端**的任务才返回 403" // (docs/admin/04-client-api.md §4.1)。但 tasks.assigned_client 只记 // **当前**归属,任务一旦重派给别人,就查不出原来那台领过—— // 而契约又明确要求"已重派仍要接受原客户端提交的结果"。 // 没有这张表,那条规则根本没法判断。 // // 顺带得到一份审计记录:这个任务被哪几台客户端领过。 { `CREATE TABLE task_claims ( task_id TEXT NOT NULL, client_id TEXT NOT NULL, claimed_at TEXT NOT NULL, PRIMARY KEY (task_id, client_id) );`, `CREATE INDEX idx_task_claims_client ON task_claims(client_id);`, }, } // schemaVersion 是当前代码支持的最新 user_version。 // // 之所以不是 len(migrations),是因为 v3(把 PDD 采集结果从 shopee_products // 拆到独立的 pdd_products、重建 sku_mappings 主键)做的事超出了 // "一串 SQL 顺序执行":它必须在事务外切换 PRAGMA foreign_keys、 // 要在丢弃旧 sku_mappings 前用 Go 数出行数打日志。 // 这些事纯 SQL 表达不了,所以 v3 单独用 migrateV3 函数实现, // 不放进 migrations 这个只支持"一条一条执行 SQL"的结构里。 // // 背景见 #20:v1 曾经被原地改写而不是新增版本,导致已经建过库的机器 // (user_version 已经越过 v1)永远不会重跑改写后的语句,程序拿着一个 // 和代码对不上的库静默启动。 const schemaVersion = 9 // migrationV4 给 PDD 商品增加店铺名。 // // v3 是特殊的 Go 迁移,不能塞进上面的纯 SQL migrations。v4 必须等 v3 // 建好 pdd_products 后再执行,所以单独放在这里。已经发布的 v1/v2 原文 // 保持不动,老库才能可靠地逐版升级。 var migrationV4 = []string{ `ALTER TABLE pdd_products ADD COLUMN shop_name TEXT;`, } // migrationV5 是工单 #46(顺运宝货运单同步)需要的三样东西: // 会话缓存表、同步进度表、给 syb_orders 补一列规格原文。 // // `[必须]` 三条都是新增(新表或 ADD COLUMN),不改动任何已发布的列/表, // 见 admin/AGENTS.md「迁移只追加」。 var migrationV5 = []string{ // 会话缓存。只存 Cookie 就够——08 §3.1 已确认 JWT 从不参与后续请求认证, // 缓存 token 没有意义,缓存 Cookie 才能免登录。 `CREATE TABLE syb_session ( username TEXT PRIMARY KEY, cookies TEXT NOT NULL, -- JSON 数组 expires_at TEXT NOT NULL, -- min(JWT exp, 24h) updated_at TEXT NOT NULL );`, // 上次同步到哪。只允许一行(CHECK id = 1),increment 同步靠它算日期范围。 `CREATE TABLE syb_sync_state ( id INTEGER PRIMARY KEY CHECK (id = 1), last_synced_at TEXT, updated_at TEXT NOT NULL );`, // 规格原文。匹配要用它,放列里才能查、才能在界面显示, // 对应 shopee_skus.spec_raw,两边同一个概念。 `ALTER TABLE syb_orders ADD COLUMN product_spec TEXT;`, } // migrationV6 增加 Admin 网页账号和 Session。只追加新表和索引, // v1-v5 的任何语句都不修改,见工单 #50。 var migrationV6 = []string{ `CREATE TABLE users ( user_id TEXT PRIMARY KEY, username TEXT NOT NULL COLLATE NOCASE UNIQUE, password_hash TEXT NOT NULL, role TEXT NOT NULL CHECK (role IN ('admin', 'purchaser')), status TEXT NOT NULL CHECK (status IN ('active', 'disabled')), last_login_at TEXT, password_changed_at TEXT NOT NULL, created_at TEXT NOT NULL, updated_at TEXT NOT NULL );`, `CREATE TABLE web_sessions ( session_hash TEXT PRIMARY KEY, user_id TEXT NOT NULL, expires_at TEXT NOT NULL, created_at TEXT NOT NULL, last_seen_at TEXT NOT NULL, FOREIGN KEY (user_id) REFERENCES users(user_id) );`, `CREATE INDEX idx_web_sessions_user ON web_sessions(user_id);`, `CREATE INDEX idx_web_sessions_expiry ON web_sessions(expires_at);`, } // migrationV7 记录客户端当前负责人及完整转交历史,见工单 #54。 // client_id 故意不加 clients 外键:客户端记录可删除后由同一稳定编号重新登记, // 归属和审计历史不能因此丢失。 var migrationV7 = []string{ `CREATE TABLE client_user_assignments ( assignment_id TEXT PRIMARY KEY, client_id TEXT NOT NULL, user_id TEXT NOT NULL, started_at TEXT NOT NULL, ended_at TEXT, assigned_by_user_id TEXT NOT NULL, ended_by_user_id TEXT, end_reason TEXT CHECK (end_reason IS NULL OR end_reason IN ('unbind', 'transfer')), FOREIGN KEY (user_id) REFERENCES users(user_id), FOREIGN KEY (assigned_by_user_id) REFERENCES users(user_id), FOREIGN KEY (ended_by_user_id) REFERENCES users(user_id), CHECK ((ended_at IS NULL AND ended_by_user_id IS NULL AND end_reason IS NULL) OR (ended_at IS NOT NULL AND ended_by_user_id IS NOT NULL AND end_reason IS NOT NULL)) );`, `CREATE UNIQUE INDEX idx_client_assignment_current ON client_user_assignments(client_id) WHERE ended_at IS NULL;`, `CREATE INDEX idx_client_assignment_user ON client_user_assignments(user_id, ended_at, client_id);`, `CREATE INDEX idx_client_assignment_history ON client_user_assignments(client_id, started_at DESC);`, } // migrationV8 持久化顺运宝同步记录,见工单 #59。 // 只新增表和索引,不修改 v1-v7 的任何已发布语句。 var migrationV8 = []string{ `CREATE TABLE syb_sync_runs ( run_id TEXT PRIMARY KEY, user_id TEXT NOT NULL, date_from TEXT NOT NULL, date_to TEXT NOT NULL, status TEXT NOT NULL CHECK (status IN ('running', 'succeeded', 'failed', 'interrupted')), stock_count INTEGER NOT NULL DEFAULT 0, detail_count INTEGER NOT NULL DEFAULT 0, created_count INTEGER NOT NULL DEFAULT 0, updated_count INTEGER NOT NULL DEFAULT 0, skipped_count INTEGER NOT NULL DEFAULT 0, error_message TEXT, cursor_advanced INTEGER NOT NULL DEFAULT 0 CHECK (cursor_advanced IN (0, 1)), started_at TEXT NOT NULL, finished_at TEXT, FOREIGN KEY (user_id) REFERENCES users(user_id), CHECK ((status = 'running' AND finished_at IS NULL) OR (status <> 'running' AND finished_at IS NOT NULL)) );`, `CREATE INDEX idx_syb_sync_runs_started ON syb_sync_runs(started_at DESC, run_id DESC);`, } // migrationV9 把顺运宝原始 JSON 里的货运单店铺名提升为可展示、可筛选的独立列。 // 原始 syb_data 永久保留;无效 JSON、缺字段和纯空白值都安全地保持 NULL。 var migrationV9 = []string{ `ALTER TABLE syb_orders ADD COLUMN shop_name TEXT;`, `UPDATE syb_orders SET shop_name = CASE WHEN json_valid(syb_data) THEN NULLIF(TRIM(json_extract(syb_data, '$.stock.shopName')), '') ELSE NULL END WHERE shop_name IS NULL;`, } // Migrate 把数据库升到最新版本。 // 已经是最新的就什么都不做,可以重复调用。 func Migrate(db *sql.DB) error { var current int if err := db.QueryRow("PRAGMA user_version").Scan(¤t); err != nil { return fmt.Errorf("读取 user_version 失败: %w", err) } if current > schemaVersion { return fmt.Errorf( "数据库版本 %d 高于本程序支持的 %d,"+ "说明这个库是更新版本的程序建的,请升级程序而不是降级", current, schemaVersion) } reached := current for v := current; v < len(migrations); v++ { tx, err := db.Begin() if err != nil { return fmt.Errorf("开始迁移 v%d 失败: %w", v+1, err) } for i, stmt := range migrations[v] { if _, err := tx.Exec(stmt); err != nil { tx.Rollback() return fmt.Errorf("执行迁移 v%d 第 %d 条语句失败: %w", v+1, i+1, err) } } // PRAGMA 不支持参数化,这里的值来自循环变量而非外部输入,安全。 if _, err := tx.Exec(fmt.Sprintf("PRAGMA user_version = %d", v+1)); err != nil { tx.Rollback() return fmt.Errorf("更新 user_version 到 %d 失败: %w", v+1, err) } if err := tx.Commit(); err != nil { return fmt.Errorf("提交迁移 v%d 失败: %w", v+1, err) } reached = v + 1 } // v3:见 schemaVersion 的注释,为什么它不在 migrations 里、 // 单独用一个函数处理。全新库也会先走完 v1/v2(拿到旧版 shopee_products / // sku_mappings 结构),再由这一步收敛成最终结构—— // 这样"全新库"和"老库升级"最终跑的是完全相同的 v3 代码, // 不需要分别维护两条路径。 if reached < 3 { if err := migrateV3(db); err != nil { return err } reached = 3 } // v4 是普通的追加列迁移,但必须排在特殊 v3 后面执行。 if reached < 4 { if err := runSQLMigration(db, 4, migrationV4); err != nil { return err } } // v5 同理,是普通的建表 + 追加列迁移,排在 v4 后面执行。 if reached < 5 { if err := runSQLMigration(db, 5, migrationV5); err != nil { return err } reached = 5 } // v6 是纯追加的用户和 Web Session 表。 if reached < 6 { if err := runSQLMigration(db, 6, migrationV6); err != nil { return err } reached = 6 } // v7 是纯追加的客户端归属历史表和索引。 if reached < 7 { if err := runSQLMigration(db, 7, migrationV7); err != nil { return err } reached = 7 } // v8 只新增顺运宝同步记录表和倒序索引。 if reached < 8 { if err := runSQLMigration(db, 8, migrationV8); err != nil { return err } reached = 8 } // v9 增加顺运宝店铺列并从仍保留的原始 JSON 回填。 if reached < 9 { if err := runSQLMigration(db, 9, migrationV9); err != nil { return err } } return nil } // runSQLMigration 在一个事务里执行指定版本的 SQL,并最后更新 user_version。 // 它只接收本文件中写死的版本号和 SQL,不接收外部输入。 func runSQLMigration(db *sql.DB, version int, statements []string) error { tx, err := db.Begin() if err != nil { return fmt.Errorf("开始迁移 v%d 失败: %w", version, err) } defer tx.Rollback() for i, stmt := range statements { if _, err := tx.Exec(stmt); err != nil { return fmt.Errorf("执行迁移 v%d 第 %d 条语句失败: %w", version, i+1, err) } } if _, err := tx.Exec(fmt.Sprintf("PRAGMA user_version = %d", version)); err != nil { return fmt.Errorf("更新 user_version 到 %d 失败: %w", version, err) } if err := tx.Commit(); err != nil { return fmt.Errorf("提交迁移 v%d 失败: %w", version, err) } return nil } // migrateV3 把库收敛成当前结构,对应工单 #20。 // // # 起点不止一种,不能用 user_version 推断结构 // // #16 曾经原地改写过 v1(把 pdd_products 等结构直接塞进 v1,没有新增版本号), // 所以 user_version = 2 现在对应两种不同的**真实**结构: // // 原始 v1 建的库 没有 pdd_products;shopee_products 还带着 // pdd_data / collect_status 等四列;sku_mappings // 还是单列主键。 // #16 改写后的 v1 建的库 结构已经是最终形态,只是 user_version 还停在 2。 // // 只要用 user_version 是不是 2 来判断"要不要迁移"就会踩空——第二种库 // 一旦被当成第一种处理,"建 pdd_products" 这一步会直接撞 "table already // exists"。所以下面每一步都必须先查**实际结构**(sqlite_master / // PRAGMA table_info),再决定做不做;已经是最终形态的库,这个函数应该 // 什么都不改,只把版本号推到 3。 // // # 四件事,各自独立判断要不要做 // // 1. 建 pdd_products 和索引——表已存在就跳过。 // 2. 把 shopee_products 上的 PDD 采集数据搬过去,按 pdd_goods_id 去重—— // shopee_products 已经没有 pdd_data 列就跳过(数据要么搬过了, // 要么这张表从来就没有过)。 // 3. 重建 sku_mappings——已经有 pdd_goods_id 列(说明已经是新结构)就跳过; // 否则旧数据全部丢弃(新主键需要的 pdd_option_key 是 // service.OptionKey() 用 json.Marshal 算出来的,SQL 复现不了,硬凑 // 有静默买错东西的风险,见 admin/AGENTS.md),丢之前先数出行数打日志。 // 4. 重建 shopee_products——还带着旧的四列中任意一个就重建,去掉那四列; // 四列都已经不在就跳过。 // // 重建 shopee_products 时必须按 SQLite 官方的表重建流程处理外键 // (https://www.sqlite.org/lang_altertable.html 的 12 步): // shopee_skus.goods_id 有外键指向 shopee_products(goods_id) ON DELETE CASCADE, // 如果不先关闭外键检查就 DROP/RENAME 这张表,可能把 shopee_skus 的数据带跑。 // 而 PRAGMA foreign_keys 只能在没有打开事务时切换,所以必须先关、 // 再开事务,事务里全部做完再提交、最后恢复。这套开销只在真的要重建 // shopee_products 时才需要——sku_mappings 只是外键的子表(没有别的表 // 指向它),丢它不会牵连别的表,不需要关闭外键检查。 func migrateV3(db *sql.DB) error { const toVersion = 3 ctx := context.Background() needCreatePddProducts, err := tableMissing(db, "pdd_products") if err != nil { return fmt.Errorf("迁移 v3 检查 pdd_products 是否存在失败: %w", err) } shopeeCols, err := tableColumnSet(db, "shopee_products") if err != nil { return fmt.Errorf("迁移 v3 检查 shopee_products 结构失败: %w", err) } needDedup := shopeeCols["pdd_data"] needRebuildShopeeProducts := shopeeCols["pdd_data"] || shopeeCols["collect_status"] || shopeeCols["collect_error"] || shopeeCols["collected_at"] skuCols, err := tableColumnSet(db, "sku_mappings") if err != nil { return fmt.Errorf("迁移 v3 检查 sku_mappings 结构失败: %w", err) } needRebuildSkuMappings := !skuCols["pdd_goods_id"] if !needCreatePddProducts && !needDedup && !needRebuildShopeeProducts && !needRebuildSkuMappings { // 结构已经是最终形态(#16 改写后的 v1 建的库),什么都不用改, // 只需要把版本号推到 3。不打"新增 / 重建"那行日志——那是假话, // 会误导操作员以为数据被动过。 log.Printf("数据库迁移 v3:结构已是最新,仅更新版本号") if _, err := db.Exec(fmt.Sprintf("PRAGMA user_version = %d", toVersion)); err != nil { return fmt.Errorf("更新 user_version 到 %d 失败: %w", toVersion, err) } return nil } // 按**实际做了什么**打日志,不是无条件打同一行——操作员/维护者要能 // 从启动日志一眼看出发生过一次结构迁移,不用等到点坏页面才知道库被动过, // 这正是 #20 要修的"静默"问题;但日志内容必须真实,不能不管做没做都 // 打同一句。 var actions []string if needCreatePddProducts { actions = append(actions, "新增 pdd_products") } if needRebuildShopeeProducts { actions = append(actions, "重建 shopee_products") } if needRebuildSkuMappings { actions = append(actions, "重建 sku_mappings") } log.Printf("数据库迁移 v3:%s", strings.Join(actions, "、")) var conn *sql.Conn if needRebuildShopeeProducts { // 只有要重建 shopee_products 才需要这套连接和 PRAGMA 切换——见函数 // 顶部的注释。单独拿一条连接:PRAGMA foreign_keys=OFF 之后紧接着要在 // **同一条连接**上开事务,database/sql 的连接池不保证 db.Exec 和 // db.Begin 用的是同一条连接。 conn, err = db.Conn(ctx) if err != nil { return fmt.Errorf("迁移 v3 获取专用连接失败: %w", err) } defer conn.Close() if _, err := conn.ExecContext(ctx, "PRAGMA foreign_keys=OFF"); err != nil { return fmt.Errorf("迁移 v3 关闭外键检查失败: %w", err) } // 这条连接用完会回到连接池,后面别的代码还会拿到它继续用。 // 不管上面成功还是失败都要把外键检查恢复成 ON—— // Open() 承诺过整个连接池的外键检查是打开的,这里关了就要负责关回去。 // 这个 defer 注册在 conn.Close() 之后,按 LIFO 顺序会先于 Close 执行。 defer func() { if _, err := conn.ExecContext(context.Background(), "PRAGMA foreign_keys=ON"); err != nil { log.Printf("警告: 迁移 v3 后恢复外键检查失败: %v", err) } }() } var tx *sql.Tx if conn != nil { tx, err = conn.BeginTx(ctx, nil) } else { tx, err = db.BeginTx(ctx, nil) } if err != nil { return fmt.Errorf("开始迁移 v3 事务失败: %w", err) } defer tx.Rollback() // 已提交的事务再 Rollback 是空操作,安全 // ① 建 pdd_products 和索引。内容照搬当前 v1 里的定义,含全部注释—— // 这是最终要收敛到的表结构,不能和 v1 曾经的定义有任何出入。 if needCreatePddProducts { for i, stmt := range migrateV3CreatePddProducts { if _, err := tx.ExecContext(ctx, stmt); err != nil { return fmt.Errorf("迁移 v3 建 pdd_products 第 %d 条语句失败: %w", i+1, err) } } } // ② 把老 shopee_products 上的 PDD 数据搬过去,按 pdd_goods_id 去重。 // 必须在下面重建 shopee_products 之前执行——一旦 shopee_products 被重建, // pdd_data / collect_status / collect_error / collected_at 这几列就没了。 if needDedup { if _, err := tx.ExecContext(ctx, migrateV3DedupIntoPddProducts); err != nil { return fmt.Errorf("迁移 v3 搬运 PDD 数据失败: %w", err) } } // ③ 旧 sku_mappings 全部丢弃:新主键需要的 pdd_option_key 是 // service.OptionKey() 用 json.Marshal 算出来的,SQL 复现不了; // 硬凑一个键有静默买错东西的风险,见 admin/AGENTS.md。 // 丢之前先数出行数打日志,不能悄悄丢——N 为 0 时不打, // 避免每次启动都刷一行没用的日志。 if needRebuildSkuMappings { var discardedMappings int if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM sku_mappings`).Scan(&discardedMappings); err != nil { return fmt.Errorf("迁移 v3 统计旧 sku_mappings 行数失败: %w", err) } if discardedMappings > 0 { log.Printf("迁移 v3:丢弃了 %d 条旧规格映射(缺少 pdd_goods_id / pdd_option_key,无法安全迁移,请重新匹配)", discardedMappings) } for i, stmt := range migrateV3RebuildSkuMappings { if _, err := tx.ExecContext(ctx, stmt); err != nil { return fmt.Errorf("迁移 v3 重建 sku_mappings 第 %d 条语句失败: %w", i+1, err) } } } // ④ 重建 shopee_products:去掉已经搬去 pdd_products 的四个列。 // 按 SQLite 12 步流程:建新表 -> 搬数据 -> 删旧表 -> 改名 -> 建索引。 if needRebuildShopeeProducts { for i, stmt := range migrateV3RebuildShopeeProducts { if _, err := tx.ExecContext(ctx, stmt); err != nil { return fmt.Errorf("迁移 v3 重建 shopee_products 第 %d 条语句失败: %w", i+1, err) } } // 外键检查原本是开着的:重建完必须确认没有把 shopee_skus 的数据带丢 // (比如误伤了它和 shopee_products 之间的外键关系)。 if err := checkForeignKeys(ctx, tx); err != nil { return fmt.Errorf("迁移 v3 外键校验失败: %w", err) } } if _, err := tx.ExecContext(ctx, fmt.Sprintf("PRAGMA user_version = %d", toVersion)); err != nil { return fmt.Errorf("更新 user_version 到 %d 失败: %w", toVersion, err) } if err := tx.Commit(); err != nil { return fmt.Errorf("提交迁移 v3 失败: %w", err) } return nil } // tableMissing 判断某张表是否不存在。 func tableMissing(db *sql.DB, table string) (bool, error) { var name string err := db.QueryRow( `SELECT name FROM sqlite_master WHERE type = 'table' AND name = ?`, table, ).Scan(&name) if err == sql.ErrNoRows { return true, nil } if err != nil { return false, err } return false, nil } // tableColumnSet 返回某张表当前的列名集合,供 migrateV3 判断 // "这张表是不是旧结构" 用——user_version 推断不出真实结构,见 migrateV3 顶部注释。 func tableColumnSet(db *sql.DB, table string) (map[string]bool, error) { // table 只来自本文件里写死的表名常量,不是外部输入,字符串拼接是安全的 // (PRAGMA 本身也不支持参数化,PRAGMA user_version 也是这么处理的)。 rows, err := db.Query(`PRAGMA table_info(` + table + `)`) if err != nil { return nil, err } defer rows.Close() cols := map[string]bool{} for rows.Next() { var cid, notnull, pk int var name, ctype string var dflt sql.NullString if err := rows.Scan(&cid, &name, &ctype, ¬null, &dflt, &pk); err != nil { return nil, err } cols[name] = true } return cols, rows.Err() } // migrateV3CreatePddProducts 建 pdd_products 和索引。 // // [必须] 这段 SQL 文本必须和 `git show 998c06a:admin/repository/db.go` // 里 pdd_products 的定义逐字节一致(含缩进),不能只是"看起来一样"。 // // 原因:#16(998c06a)把这张表直接建进了 v1,那批库上 sqlite_master 存的 // 就是 998c06a 里这段文本原样的字节(连同当时 gofmt 缩进出来的那个前导 // TAB)。这个函数只在表**不存在**时才会执行这条 CREATE(见 migrateV3 里 // needCreatePddProducts 的判断),所以已经建过表的库不会被这段文本重新 // 覆盖——"新建库" 和 "已经带 pdd_products 的老库" 要收敛到完全相同的 // sqlite_master.sql,就必须让新建的这份文本和老库上躺着的那份逐字节相同, // 哪怕只差一个空格/TAB 都会让 TestMigrate_不同起点最终schema一致 失败 // (这个坑已经在 #20 审查阶段被变异测试连同真实文本对比一起抓到过一次)。 var migrateV3CreatePddProducts = []string{ `CREATE TABLE pdd_products ( id INTEGER PRIMARY KEY AUTOINCREMENT, -- 从 PDD 链接里解析出来。它不是主键,所以**必须加 UNIQUE**: -- 少了这条约束,同一个 PDD 商品会被存成好几行, -- 采好几遍,映射还说不清指向哪一行。 goods_id TEXT NOT NULL UNIQUE, url TEXT NOT NULL, -- 操作员填的链接原文 title TEXT, -- 采集回来,人工核对"是不是我要的那个商品" skus_json TEXT, -- schema_version + dimensions + skus -- 注意这里**没有 no_link**:这张表里有这一行,就说明链接已经填了。 -- "未填链接"是蝦皮侧的状态(shopee_products.pdd_goods_id 为空)。 collect_status TEXT NOT NULL DEFAULT 'pending' CHECK (collect_status IN ( 'pending', 'collecting', 'collected', 'failed' )), collect_msg TEXT, -- 失败原因,要能定位问题 artifact_ref TEXT, -- 诊断产物在哪台机器哪个目录 collected_at TEXT, -- 软删除。不硬删是因为 sku_mappings 指向它, -- 硬删会把人工攒了很久的匹配成果一起带走。 deleted_at TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL );`, `CREATE INDEX idx_pdd_products_status ON pdd_products(collect_status);`, } // migrateV3DedupIntoPddProducts 按 pdd_goods_id 去重,把老 // shopee_products 上人工攒的 PDD 数据搬进 pdd_products。 // // - MIN(COALESCE(pdd_goods_url, 空字符串)) / MIN(created_at):多行里任取一个即可, // 用 MIN 只是为了确定性(同一批输入每次跑结果一样,方便排查)。 // - MAX(pdd_data) / MAX(collect_error) / MAX(collected_at):同理, // 只是要"取到某一行的值",用 MAX 是同一个考虑。 // - CASE MAX(collect_status) ...:collect_status 只有全组都是 // 'collected' 时才判定为 collected('collected' 按字符串比较是这几个 // 取值里最小的,只要组里有任何一行不是 collected,MAX 就会取到别的值); // 只要没有 collected 但有 failed 就判定 failed; // 其余(含 collecting / no_link / pending)一律落进 ELSE,判定 pending // —— 这正好同时满足"collecting 映射成 pending"和 // "no_link 映射成 pending,不撞新 CHECK(新表里没有 no_link)"两条要求。 var migrateV3DedupIntoPddProducts = ` INSERT INTO pdd_products (goods_id, url, title, skus_json, collect_status, collect_msg, collected_at, created_at, updated_at) SELECT pdd_goods_id, MIN(COALESCE(pdd_goods_url, '')), NULL, MAX(pdd_data), CASE MAX(collect_status) WHEN 'collected' THEN 'collected' WHEN 'failed' THEN 'failed' ELSE 'pending' END, MAX(collect_error), MAX(collected_at), MIN(created_at), MAX(updated_at) FROM shopee_products WHERE pdd_goods_id IS NOT NULL AND pdd_goods_id <> '' GROUP BY pdd_goods_id;` // migrateV3RebuildShopeeProducts 重建 shopee_products: // 去掉已经搬去 pdd_products 的 pdd_data / collect_status / collect_error / // collected_at 四列,删掉跟着它们的 idx_shopee_products_status, // 换成新结构需要的 idx_shopee_products_pdd。 // // # 步骤顺序为什么是"建 _new → 搬数据 → 删旧表 → 改名"(标准 12 步顺序) // // 这里**必须**按 SQLite 官方 12 步流程的顺序来,不能改成"先把旧表改名 // 让开、再直接用最终表名建新表"(表面上能避开下面说的引号问题, // 实测过、但会坏得更彻底): // // shopee_skus.goods_id 有 `FOREIGN KEY (goods_id) REFERENCES // shopee_products(goods_id)`。SQLite 的 `ALTER TABLE ... RENAME TO` // **默认会连带更新别的表里引用这张表的外键定义**——如果改成先把 // shopee_products RENAME 成 shopee_products_old,shopee_skus 的外键子句 // 会被自动重写成 `REFERENCES "shopee_products_old"(goods_id)`; // 后面一 DROP TABLE shopee_products_old,shopee_skus 就带着一条指向 // 不存在的表的外键,`PRAGMA foreign_key_check` 直接报错, // 而且这个坏结果比"多一对引号"严重得多。 // // 按标准顺序(建 shopee_products_new → 搬数据 → DROP 掉的是旧的 // shopee_products,不是被引用的名字 → RENAME shopee_products_new // 成 shopee_products)不会触发这个重写:shopee_skus 的外键子句 // 全程写的都是 "shopee_products" 这个名字,没有变过,RENAME 结束后 // 这个名字重新存在,直接就能解析上,不需要 SQLite 帮它改写任何东西。 // 实测两种顺序的效果: // // 方案:先建 _new 再 RENAME 成正式名(本文件采用的顺序) // shopee_skus 的外键子句:REFERENCES shopee_products(goods_id) 不变 ✓ // shopee_products 自己的 CREATE TABLE 文本:被 RENAME 加上引号 CREATE TABLE "shopee_products" (...) // // 方案:先把旧表 RENAME 让开,再直接用正式名建新表 // shopee_skus 的外键子句:被自动重写成 REFERENCES "shopee_products_old"(goods_id) ← 引用悬空,PRAGMA foreign_key_check 报错 // // 所以选择"忍受 shopee_products 自己的 CREATE TABLE 文本被套一层引号", // 而不是"外键指向一张已经被删掉的表"。引号这个副作用改在 // assertSchemaEqual 的比较逻辑里通过归一化解决(见 migrate_test.go // 里 normalizeCreateStatement 的注释),不在这里想办法回避—— // 这个坑是 #20 审查阶段三起点收敛测试跑出来才发现的,光读代码看不出来。 // // [必须] 除了表名后缀 _new、以及 RENAME 带来的引号,这段 SQL 里 // shopee_products 的列定义必须和 // `git show 998c06a:admin/repository/db.go` 里的定义逐字节一致 // (含缩进),理由同 migrateV3CreatePddProducts 的注释。 var migrateV3RebuildShopeeProducts = []string{ `CREATE TABLE shopee_products_new ( goods_id TEXT PRIMARY KEY, title TEXT NOT NULL, shopee_status TEXT, main_sku_code TEXT, -- 人工维护的,蝦皮报表里没有这两列,Excel 导入时绝不能覆盖。 -- pdd_goods_id 指向 pdd_products.goods_id,表示"这个蝦皮商品 -- 当前对应哪个 PDD 商品"。PDD 商品下架换新时改这里。 pdd_goods_url TEXT, pdd_goods_id TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL );`, `INSERT INTO shopee_products_new (goods_id, title, shopee_status, main_sku_code, pdd_goods_url, pdd_goods_id, created_at, updated_at) SELECT goods_id, title, shopee_status, main_sku_code, pdd_goods_url, pdd_goods_id, created_at, updated_at FROM shopee_products;`, `DROP TABLE shopee_products;`, `ALTER TABLE shopee_products_new RENAME TO shopee_products;`, `CREATE INDEX idx_shopee_products_pdd ON shopee_products(pdd_goods_id);`, } // migrateV3RebuildSkuMappings 重建 sku_mappings。 // 旧数据全部丢弃(调用方已经数过行数打过日志),不搬任何数据—— // 新主键需要的 pdd_option_key 只有 Go 的 service.OptionKey() 能算, // SQL 里凑不出来,硬凑有静默买错东西的风险。 // // 这里是直接 DROP 旧表、CREATE 新表(同一个最终表名,不经过 RENAME), // 不会有 migrateV3RebuildShopeeProducts 注释里说的引号问题。 // // [必须] 这段 SQL 文本必须和 `git show 998c06a:admin/repository/db.go` // 里 sku_mappings 的定义逐字节一致(含缩进),理由同 // migrateV3CreatePddProducts 的注释。 var migrateV3RebuildSkuMappings = []string{ `DROP TABLE sku_mappings;`, `CREATE TABLE sku_mappings ( shopee_sku_id TEXT NOT NULL, pdd_goods_id TEXT NOT NULL, pdd_option_key TEXT NOT NULL, pdd_options TEXT NOT NULL, goods_id TEXT NOT NULL, mapped_at TEXT NOT NULL, mapped_by TEXT, PRIMARY KEY (shopee_sku_id, pdd_goods_id), FOREIGN KEY (shopee_sku_id) REFERENCES shopee_skus(sku_id) ON DELETE CASCADE );`, `CREATE INDEX idx_sku_mappings_goods ON sku_mappings(goods_id);`, `CREATE INDEX idx_sku_mappings_pdd ON sku_mappings(pdd_goods_id);`, } // checkForeignKeys 跑 PRAGMA foreign_key_check,有任何一行结果 // 就说明外键关系被破坏了(比如子表指向了一个已经不存在的父行)。 func checkForeignKeys(ctx context.Context, tx *sql.Tx) error { rows, err := tx.QueryContext(ctx, "PRAGMA foreign_key_check") if err != nil { return fmt.Errorf("执行外键校验失败: %w", err) } defer rows.Close() if rows.Next() { return fmt.Errorf("存在外键约束冲突,重建表的过程把关联数据带丢了") } return rows.Err() } // requiredTables 是当前代码依赖的全部表。 // Migrate 跑完之后用它做一次自检,见 CheckSchema。 var requiredTables = []string{ "shopee_products", "shopee_skus", "pdd_products", "syb_orders", "sku_mappings", "tasks", "clients", "idempotency_keys", "task_claims", "syb_session", "syb_sync_state", "users", "web_sessions", "client_user_assignments", "syb_sync_runs", } // requiredColumns 只列出不能靠“表存在”发现的关键追加列。 // shop_name 是 v4 新增列;product_spec 是 v5 新增列——两者缺失时 // 查询对应页面会直接失败,因此启动时就应给出明确错误, // 而不是等操作员点到页面才暴露。 var requiredColumns = map[string][]string{ "pdd_products": {"shop_name"}, "syb_orders": {"product_spec"}, "syb_sync_runs": {"user_id", "date_from", "date_to", "status", "cursor_advanced", "started_at", "finished_at"}, } // SQLite 当前测试库比 MySQL v1 多走一版迁移;这里单独检查,避免 MySQL // 在执行 v1 自检时提前要求尚未由 v11 创建的列。 var sqliteOnlyRequiredColumns = map[string][]string{ "syb_orders": {"shop_name"}, } // CheckSchema 在 Migrate 成功后调用,确认代码依赖的表都在。 // // [必须] 缺表就返回错误,调用方要**拒绝启动**,不是打个警告继续跑。 // #20 的教训就是静默启动:程序拿着一个和代码对不上的库正常起来了, // 错误要等操作员点到那个页面才暴露——如果那是个写操作页面, // 暴露出来的就不是报错而是写坏数据。 // // 表名全部检查;关键的追加列也检查,防止 user_version 已更新但迁移未完整 // 落地时,程序拿着缺列的库继续启动。 func CheckSchema(db *sql.DB) error { rows, err := db.Query(`SELECT name FROM sqlite_master WHERE type = 'table'`) if err != nil { return fmt.Errorf("读取数据库表清单失败: %w", err) } defer rows.Close() existing := map[string]bool{} for rows.Next() { var name string if err := rows.Scan(&name); err != nil { return fmt.Errorf("读取数据库表清单失败: %w", err) } existing[name] = true } if err := rows.Err(); err != nil { return fmt.Errorf("读取数据库表清单失败: %w", err) } var missing []string for _, t := range requiredTables { if !existing[t] { missing = append(missing, t) } } if len(missing) > 0 { return fmt.Errorf( "数据库结构与本程序不匹配:缺少表 %s。\n"+ "这通常是数据库比程序旧、而迁移没有覆盖到。\n"+ "请备份 data/admin.db 后删除它让程序重建,或联系维护者。", strings.Join(missing, "、")) } for table, columns := range requiredColumns { existingColumns, err := tableColumnSet(db, table) if err != nil { return fmt.Errorf("检查数据表 %s 的列失败: %w", table, err) } for _, column := range columns { if !existingColumns[column] { return fmt.Errorf( "数据库结构与本程序不匹配:数据表 %s 缺少列 %s。\n"+ "这通常是数据库迁移没有完整执行,请备份 data/admin.db 后联系维护者。", table, column) } } } for table, columns := range sqliteOnlyRequiredColumns { existingColumns, err := tableColumnSet(db, table) if err != nil { return fmt.Errorf("检查数据表 %s 的列失败: %w", table, err) } for _, column := range columns { if !existingColumns[column] { return fmt.Errorf("数据库结构与本程序不匹配:数据表 %s 缺少列 %s", table, column) } } } return nil }