// ============================================================ // run_style_check_oss_migrate.js - 穿搭镜检图片本地 -> OSS 一次性迁移 runner // 用法: // node dev_test_scripts/runners/run_style_check_oss_migrate.js status // node dev_test_scripts/runners/run_style_check_oss_migrate.js upload (并发6,断点续传,可反复执行) // node dev_test_scripts/runners/run_style_check_oss_migrate.js verify [N] (默认抽样200张HEAD校验+计数比对) // node dev_test_scripts/runners/run_style_check_oss_migrate.js rewrite (全绿后回写DB+data.json,自动备份) // 凭据:自加载 ~/Toolbox_local_creds.env.local(EXTERNAL_STORAGE_OSS_*) // 桶/android-o1-images 与 key 规范与 src/server/style_check.js 中 SC_OSS_* 常量严格一致 // ============================================================ 'use strict' const fs = require('fs') const path = require('path') const os = require('os') const crypto = require('crypto') const https = require('https') const Database = require('better-sqlite3') const ROOT = path.resolve(__dirname, '..', '..') const DB_PATH = path.join(ROOT, 'data', 'style_check.db') const IMG_ROOT = path.join(ROOT, 'public', 'images') const X_FASHION_DATA_ROOT = path.join(ROOT, 'public', 'data', 'x_fashion') // ---------- 凭据加载(与 index.js 同款逻辑) ---------- const loadCreds = () => { const userProfile = process.env.USERPROFILE || process.env.HOME || os.homedir() const p = path.join(userProfile, 'Toolbox_local_creds.env.local') if (!fs.existsSync(p)) throw new Error('找不到凭据文件: ' + p) fs.readFileSync(p, 'utf-8').split(/\r?\n/).forEach(line => { const t = line.trim() if (!t || t.startsWith('#')) return const i = t.indexOf('=') if (i <= 0) return const k = t.slice(0, i).trim() const v = t.slice(i + 1).trim() if (k && !Object.prototype.hasOwnProperty.call(process.env, k)) process.env[k] = v }) } loadCreds() // ---------- OSS 常量(必须与 style_check.js 保持一致) ---------- const SC_OSS_BUCKET = 'android-o1-images' const SC_OSS_ENDPOINT = 'oss-accelerate.aliyuncs.com' const SC_OSS_ROOT = '550E8400-E29B-41D4-A716-446655440000/style_check' const SC_OSS_LIB_GUID = { item_images: '9F2C6A41-3B8D-4E5A-9D17-2C4E8A6B0F31', wave: 'A7D3E912-6C4F-4B8E-8A25-5F1B9C3D7E60', korean: 'C4B8F263-9A1E-4D7C-B396-8E2A5D0F4C17', x_fashion: 'E1A5C874-2F6B-4A9D-9C48-7B3E6F1A8D52' } // ---------- 四个图片库:本地目录 / URL 前缀 ---------- const LIBS = [ { lib: 'item_images', dir: path.join(IMG_ROOT, 'style_check_images'), prefix: '/images/style_check_images', recursive: true }, { lib: 'wave', dir: path.join(IMG_ROOT, 'wave_references'), prefix: '/images/wave_references', recursive: false }, { lib: 'korean', dir: path.join(IMG_ROOT, 'korean_references'), prefix: '/images/korean_references', recursive: false }, { lib: 'x_fashion', dir: path.join(IMG_ROOT, 'x_fashion'), prefix: '/images/x_fashion', recursive: true } ] const IMG_EXT = new Set(['.jpg', '.jpeg', '.png', '.webp', '.gif']) const contentTypeByExt = (ext) => { const e = String(ext || '').toLowerCase() if (e === '.png') return 'image/png' if (e === '.webp') return 'image/webp' if (e === '.gif') return 'image/gif' return 'image/jpeg' } // ---------- 扫描本地文件 ---------- const walk = (dir) => { const out = [] if (!fs.existsSync(dir)) return out for (const name of fs.readdirSync(dir)) { const abs = path.join(dir, name) const st = fs.statSync(abs) if (st.isDirectory()) out.push(...walk(abs)) else if (IMG_EXT.has(path.extname(name).toLowerCase())) out.push(abs) } return out } const scanAll = () => { const files = [] for (const L of LIBS) { for (const abs of walk(L.dir)) { const rel = path.relative(L.dir, abs).replace(/\\/g, '/') files.push({ lib: L.lib, abs, rel, localUrl: `${L.prefix}/${rel}` }) } } return files } // ---------- OSS 请求 ---------- const scOssEncode = k => String(k || '').split('/').map(s => encodeURIComponent(s)).join('/') const scOssKey = (lib, rel) => `${SC_OSS_ROOT}/${SC_OSS_LIB_GUID[lib]}/${rel}` const scOssUrl = key => `https://${SC_OSS_BUCKET}.${SC_OSS_ENDPOINT}/${scOssEncode(key)}` const scOssSign = ({ method, contentType = '', date, key, extraHeaders = {} }) => { const canonical = Object.keys(extraHeaders) .filter(k => k.toLowerCase().startsWith('x-oss-')) .sort((a, b) => a.localeCompare(b)) .map(k => `${k.toLowerCase()}:${String(extraHeaders[k] || '').trim()}`) .join('\n') const stringToSign = `${method}\n\n${contentType}\n${date}\n${canonical ? canonical + '\n' : ''}/${SC_OSS_BUCKET}/${key}` return crypto.createHmac('sha1', process.env.EXTERNAL_STORAGE_OSS_ACCESS_KEY_SECRET || '').update(stringToSign).digest('base64') } const ossRequest = (method, key, { body = null, contentType = '' } = {}) => new Promise((resolve, reject) => { const akId = process.env.EXTERNAL_STORAGE_OSS_ACCESS_KEY_ID || '' if (!akId) { reject(new Error('missing EXTERNAL_STORAGE_OSS credentials')); return } const date = new Date().toUTCString() const extraHeaders = method === 'PUT' ? { 'x-oss-object-acl': 'public-read' } : {} const headers = { Host: `${SC_OSS_BUCKET}.${SC_OSS_ENDPOINT}`, Date: date, ...extraHeaders, Authorization: `OSS ${akId}:${scOssSign({ method, contentType, date, key, extraHeaders })}` } if (body) { headers['Content-Type'] = contentType headers['Content-Length'] = body.length } const req = https.request({ hostname: `${SC_OSS_BUCKET}.${SC_OSS_ENDPOINT}`, path: `/${scOssEncode(key)}`, method, headers }, res => { const chunks = [] res.on('data', c => chunks.push(c)) res.on('end', () => resolve({ status: res.statusCode, body: Buffer.concat(chunks) })) }) req.on('error', reject) req.setTimeout(300000, () => req.destroy(new Error('oss request timeout'))) req.end(body || undefined) }) // ---------- 并发池 ---------- const runPool = async (items, worker, concurrency) => { let idx = 0 const workers = Array.from({ length: Math.min(concurrency, items.length) }, async () => { while (idx < items.length) { const my = idx++ await worker(items[my], my) } }) await Promise.all(workers) } const log = (...a) => console.log(`[${new Date().toLocaleTimeString('zh-CN', { hour12: false })}]`, ...a) // ---------- DB ---------- const openDb = () => new Database(DB_PATH) const ensureLogTable = (db) => db.exec(` CREATE TABLE IF NOT EXISTS oss_migration_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, lib TEXT NOT NULL, local_path TEXT NOT NULL, oss_key TEXT NOT NULL, oss_url TEXT NOT NULL, status TEXT NOT NULL DEFAULT 'pending', retry_count INTEGER NOT NULL DEFAULT 0, error TEXT DEFAULT '', created_at TEXT DEFAULT (datetime('now')), updated_at TEXT DEFAULT (datetime('now')), UNIQUE(local_path) ); CREATE INDEX IF NOT EXISTS idx_oss_migration_status ON oss_migration_log(status); `) // ============================================================ // status // ============================================================ const cmdStatus = () => { const db = openDb() ensureLogTable(db) const rows = db.prepare(`SELECT lib, status, COUNT(*) AS n FROM oss_migration_log GROUP BY lib, status`).all() const local = scanAll() const localCount = {} local.forEach(f => { localCount[f.lib] = (localCount[f.lib] || 0) + 1 }) log('本地文件数:', JSON.stringify(localCount)) if (rows.length === 0) log('日志表为空,尚未上传') rows.forEach(r => log(` ${r.lib} / ${r.status}: ${r.n}`)) const dbLocal = db.prepare(` SELECT (SELECT COUNT(*) FROM apparel_item_images WHERE image_url LIKE '/images/%') AS item_local, (SELECT COUNT(*) FROM wave_references WHERE image_url LIKE '/images/%') AS wave_local, (SELECT COUNT(*) FROM korean_references WHERE image_url LIKE '/images/%') AS korean_local `).get() log('DB 仍指向本地的行数:', JSON.stringify(dbLocal)) db.close() } // ============================================================ // upload(断点续传) // ============================================================ const cmdUpload = async () => { const db = openDb() ensureLogTable(db) const done = new Set(db.prepare(`SELECT local_path FROM oss_migration_log WHERE status IN ('uploaded','rewritten')`).all().map(r => r.local_path)) const files = scanAll().filter(f => !done.has(f.localUrl)) log(`本地共 ${done.size + files.length} 个文件,已完成 ${done.size},待上传 ${files.length}`) if (files.length === 0) { log('无需上传'); db.close(); return } const upsert = db.prepare(`INSERT INTO oss_migration_log (lib, local_path, oss_key, oss_url, status, retry_count, error) VALUES (?, ?, ?, ?, ?, ?, ?) ON CONFLICT(local_path) DO UPDATE SET oss_key=excluded.oss_key, oss_url=excluded.oss_url, status=excluded.status, retry_count=excluded.retry_count, error=excluded.error, updated_at=datetime('now')`) const getRow = db.prepare(`SELECT retry_count FROM oss_migration_log WHERE local_path=?`) let ok = 0, fail = 0 await runPool(files, async (f) => { const key = scOssKey(f.lib, f.rel) const url = scOssUrl(key) try { const buf = fs.readFileSync(f.abs) const resp = await ossRequest('PUT', key, { body: buf, contentType: contentTypeByExt(path.extname(f.abs)) }) if (resp.status < 200 || resp.status >= 300) throw new Error(`oss_put_${resp.status}`) upsert.run(f.lib, f.localUrl, key, url, 'uploaded', 0, '') ok++ } catch (e) { const prev = getRow.get(f.localUrl) upsert.run(f.lib, f.localUrl, key, url, 'pending', (prev ? prev.retry_count : 0) + 1, String(e.message || e).slice(0, 500)) fail++ if (fail <= 20) log(`上传失败 ${f.localUrl}: ${e.message}`) } const total = ok + fail if (total % 100 === 0) log(`进度 ${total}/${files.length}(成功 ${ok} 失败 ${fail})`) }, 6) log(`上传结束:成功 ${ok},失败 ${fail},累计完成 ${done.size + ok}/${done.size + files.length}`) db.close() } // ============================================================ // verify [N] // ============================================================ const cmdVerify = async (sampleN) => { const db = openDb() ensureLogTable(db) const local = scanAll() const localCount = {} local.forEach(f => { localCount[f.lib] = (localCount[f.lib] || 0) + 1 }) const logCount = {} db.prepare(`SELECT lib, status, COUNT(*) AS n FROM oss_migration_log GROUP BY lib, status`).all() .forEach(r => { logCount[`${r.lib}/${r.status}`] = r.n }) log('本地文件数:', JSON.stringify(localCount)) log('日志状态数:', JSON.stringify(logCount)) const bad = db.prepare(`SELECT lib, local_path, status, error FROM oss_migration_log WHERE status NOT IN ('uploaded','rewritten')`).all() if (bad.length > 0) { log(`存在未成功行 ${bad.length} 条(rewrite 被阻止):`) bad.slice(0, 10).forEach(b => log(` ${b.status} ${b.local_path} ${b.error}`)) } else { log('无 pending/failed 行') } // 计数比对 let countOk = true for (const L of LIBS) { const uploaded = db.prepare(`SELECT COUNT(*) AS n FROM oss_migration_log WHERE lib=? AND status IN ('uploaded','rewritten')`).get(L.lib).n const lc = localCount[L.lib] || 0 if (uploaded < lc) { countOk = false; log(`计数不符 ${L.lib}: 本地 ${lc} > 已上传 ${uploaded}`) } else log(`计数OK ${L.lib}: 本地 ${lc} / 已上传 ${uploaded}`) } // 抽样 HEAD 校验 const rows = db.prepare(`SELECT local_path, oss_key FROM oss_migration_log WHERE status IN ('uploaded','rewritten')`).all() const n = Math.min(sampleN, rows.length) const sample = [] const step = rows.length > n ? Math.floor(rows.length / n) : 1 for (let i = 0; i < rows.length && sample.length < n; i += step) sample.push(rows[i]) let headOk = 0, headFail = 0 await runPool(sample, async (r) => { try { const resp = await ossRequest('HEAD', r.oss_key) if (resp.status === 200) headOk++ else { headFail++; log(`HEAD ${resp.status}: ${r.local_path}`) } } catch (e) { headFail++; log(`HEAD error: ${r.local_path} ${e.message}`) } }, 6) log(`抽样HEAD:${headOk} 成功 / ${headFail} 失败(共 ${sample.length})`) const allGreen = bad.length === 0 && countOk && headFail === 0 log(allGreen ? 'VERIFY ALL GREEN,可执行 rewrite' : 'VERIFY 未全绿,禁止 rewrite') db.close() process.exitCode = allGreen ? 0 : 1 } // ============================================================ // rewrite(回写 DB + data.json,先备份) // ============================================================ const cmdRewrite = () => { const db = openDb() ensureLogTable(db) const bad = db.prepare(`SELECT COUNT(*) AS n FROM oss_migration_log WHERE status NOT IN ('uploaded','rewritten')`).get().n if (bad > 0) { log(`存在 ${bad} 条未成功行,禁止 rewrite。先跑 upload/verify`); db.close(); process.exitCode = 1; return } const ts = new Date().toISOString().replace(/[:.]/g, '-').slice(0, 19) // 备份 DB const dbBak = `${DB_PATH}.pre_oss_${ts}.bak` fs.copyFileSync(DB_PATH, dbBak) log('DB 已备份:', dbBak) // 备份 data.json const jsonFiles = [] if (fs.existsSync(X_FASHION_DATA_ROOT)) { for (const acc of fs.readdirSync(X_FASHION_DATA_ROOT)) { const p = path.join(X_FASHION_DATA_ROOT, acc, 'data.json') if (fs.existsSync(p)) jsonFiles.push(p) } } jsonFiles.forEach(p => fs.copyFileSync(p, `${p}.pre_oss_${ts}.bak`)) log(`data.json 已备份 ${jsonFiles.length} 个`) const rows = db.prepare(`SELECT local_path, oss_url FROM oss_migration_log WHERE status='uploaded'`).all() log(`待回写 ${rows.length} 行`) const upItem = db.prepare(`UPDATE apparel_item_images SET image_url=? WHERE image_url=?`) const upWave = db.prepare(`UPDATE wave_references SET image_url=? WHERE image_url=?`) const upKorean = db.prepare(`UPDATE korean_references SET image_url=? WHERE image_url=?`) const markDone = db.prepare(`UPDATE oss_migration_log SET status='rewritten', updated_at=datetime('now') WHERE local_path=?`) // x_fashion data.json 内存替换 const jsonContents = jsonFiles.map(p => ({ p, text: fs.readFileSync(p, 'utf-8') })) let dbHits = 0, jsonHits = 0 const tx = db.transaction(() => { for (const r of rows) { dbHits += upItem.run(r.oss_url, r.local_path).changes + upWave.run(r.oss_url, r.local_path).changes + upKorean.run(r.oss_url, r.local_path).changes for (const jc of jsonContents) { if (jc.text.includes(r.local_path)) { jc.text = jc.text.split(r.local_path).join(r.oss_url) jsonHits++ } } markDone.run(r.local_path) } }) tx() jsonContents.forEach(jc => fs.writeFileSync(jc.p, jc.text)) log(`回写完成:DB 命中 ${dbHits} 行,data.json 命中 ${jsonHits} 处`) const remain = db.prepare(`SELECT COUNT(*) AS n FROM apparel_item_images WHERE image_url LIKE '/images/%'`).get().n + db.prepare(`SELECT COUNT(*) AS n FROM wave_references WHERE image_url LIKE '/images/%'`).get().n + db.prepare(`SELECT COUNT(*) AS n FROM korean_references WHERE image_url LIKE '/images/%'`).get().n log(`DB 中仍指向 /images/ 本地的行数: ${remain}`) db.close() } // ---------- 入口 ---------- const [,, mode, argN] = process.argv const main = async () => { if (mode === 'status') cmdStatus() else if (mode === 'upload') await cmdUpload() else if (mode === 'verify') await cmdVerify(parseInt(argN || '200', 10) || 200) else if (mode === 'rewrite') cmdRewrite() else { console.log('用法: node run_style_check_oss_migrate.js ') process.exitCode = 1 } } main().catch(e => { console.error('runner error:', e); process.exitCode = 1 })