From 0f8fe2c2e2946e5e0b6eb9b4b96aaca29d2d372c Mon Sep 17 00:00:00 2001 From: yangxiangyuan Date: Thu, 6 Aug 2026 22:14:20 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=96=B0=E5=A2=9EOSS=E7=9B=B4=E4=BC=A0?= =?UTF-8?q?=E8=BF=81=E7=A7=BB=E8=83=BD=E5=8A=9B=EF=BC=8C=E5=AE=9E=E7=8E=B0?= =?UTF-8?q?=E5=9B=BE=E7=89=87=E4=BB=8E=E6=9C=AC=E5=9C=B0=E5=88=B0OSS?= =?UTF-8?q?=E7=9A=84=E8=BF=81=E7=A7=BB=E6=B5=81=E7=A8=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. 新增OSS直传工具类函数,实现带凭据校验的PUT/HEAD请求 2. 新增oss_migration_log数据表,用于跟踪迁移状态 3. 重构图片保存逻辑,优先上传OSS失败则降级本地并登记补传任务 4. 新增批量迁移脚本,支持状态查看、断点续传、校验和回写DB 5. 新增OSS直链抽查调试脚本 --- .../debug/debug_oss_spot_check.js | 26 ++ .../runners/run_style_check_oss_migrate.js | 353 ++++++++++++++++++ src/server/style_check.js | 157 ++++++-- 3 files changed, 515 insertions(+), 21 deletions(-) create mode 100644 dev_test_scripts/debug/debug_oss_spot_check.js create mode 100644 dev_test_scripts/runners/run_style_check_oss_migrate.js diff --git a/dev_test_scripts/debug/debug_oss_spot_check.js b/dev_test_scripts/debug/debug_oss_spot_check.js new file mode 100644 index 0000000..8b8d172 --- /dev/null +++ b/dev_test_scripts/debug/debug_oss_spot_check.js @@ -0,0 +1,26 @@ +// 抽查 OSS 直链是否可公共读:每个库取 1 条 oss_url 做 GET,打印状态码与字节数 +'use strict' +const path = require('path') +const https = require('https') +const Database = require('better-sqlite3') + +const db = new Database(path.resolve(__dirname, '..', '..', 'data', 'style_check.db')) +const libs = ['item_images', 'wave', 'korean', 'x_fashion'] + +const get = (url) => new Promise((resolve, reject) => { + https.get(url, res => { + let n = 0 + res.on('data', c => { n += c.length }) + res.on('end', () => resolve({ status: res.statusCode, bytes: n })) + }).on('error', reject) +}) + +;(async () => { + for (const lib of libs) { + const r = db.prepare('SELECT local_path, oss_url FROM oss_migration_log WHERE lib=? LIMIT 1').get(lib) + if (!r) { console.log(lib, '无记录'); continue } + const res = await get(r.oss_url) + console.log(`${lib} -> ${res.status} ${res.bytes} bytes ${r.oss_url}`) + } + db.close() +})() diff --git a/dev_test_scripts/runners/run_style_check_oss_migrate.js b/dev_test_scripts/runners/run_style_check_oss_migrate.js new file mode 100644 index 0000000..4d08069 --- /dev/null +++ b/dev_test_scripts/runners/run_style_check_oss_migrate.js @@ -0,0 +1,353 @@ +// ============================================================ +// 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 }) diff --git a/src/server/style_check.js b/src/server/style_check.js index 06ed3b1..1b259f1 100644 --- a/src/server/style_check.js +++ b/src/server/style_check.js @@ -439,6 +439,124 @@ const localImageRoot = path.join(process.cwd(), 'public', 'images', 'style_check if (!fs.existsSync(localImageRoot)) fs.mkdirSync(localImageRoot, { recursive: true }) const isLocalImageUrl = (url) => String(url || '').startsWith(`${localImageUrlPrefix}/`) + +// ============================================================ +// OSS 直传助手(本 tool 自包含,不引用其他 tool 的代码) +// 凭据复用 EXTERNAL_STORAGE_OSS_*(index.js 启动时已从 ~/Toolbox_local_creds.env.local 加载) +// 桶 android-o1-images,对象公共读,前端直接消费直链 +// key 规范(Rule 5):550E8400-.../style_check/<库GUID>/<相对路径>,四个图片库各用固定 GUID 隔离 +// 上传失败时降级落本地 + 登记 oss_migration_log,由 dev_test_scripts/runners/run_style_check_oss_migrate.js 补传 +// ============================================================ +const httpsOss = require('https') +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' +} +const scOssKey = (lib, relName) => `${SC_OSS_ROOT}/${SC_OSS_LIB_GUID[lib]}/${String(relName || '').replace(/\\/g, '/').replace(/\/+/g, '/')}` +const scOssEncode = k => String(k || '').split('/').map(s => encodeURIComponent(s)).join('/') +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 scOssPut = (key, buffer, 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 = { 'x-oss-object-acl': 'public-read' } + const req = httpsOss.request({ + hostname: `${SC_OSS_BUCKET}.${SC_OSS_ENDPOINT}`, + path: `/${scOssEncode(key)}`, + method: 'PUT', + headers: { + Host: `${SC_OSS_BUCKET}.${SC_OSS_ENDPOINT}`, + Date: date, + 'Content-Type': contentType, + 'Content-Length': buffer.length, + ...extraHeaders, + Authorization: `OSS ${akId}:${scOssSign({ method: 'PUT', contentType, date, key, extraHeaders })}` + } + }, res => { + const chunks = [] + res.on('data', c => chunks.push(c)) + res.on('end', () => { + if (res.statusCode >= 200 && res.statusCode < 300) resolve(true) + else reject(new Error(`oss_put_${res.statusCode}: ${Buffer.concat(chunks).toString('utf8').slice(0, 300)}`)) + }) + }) + req.on('error', reject) + req.setTimeout(120000, () => req.destroy(new Error('oss put timeout'))) + req.end(buffer) +}) + +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); +`) + +// 登记补传(OSS 上传失败降级本地后调用) +const registerOssPending = (lib, localUrl, ossKey, ossUrl, errMsg) => { + try { + db.prepare(`INSERT INTO oss_migration_log (lib, local_path, oss_key, oss_url, status, retry_count, error) + VALUES (?, ?, ?, ?, 'pending', 0, ?) + ON CONFLICT(local_path) DO UPDATE SET oss_key=excluded.oss_key, oss_url=excluded.oss_url, status='pending', error=excluded.error, updated_at=datetime('now')` + ).run(lib, localUrl, ossKey, ossUrl, String(errMsg || '').slice(0, 500)) + } catch (e) { + console.error('[style_check] registerOssPending error:', e.message) + } +} + +// 优先直传 OSS;失败降级本地 + 登记补传,保证不丢图 +const saveImageBufferSmart = async (lib, relName, buffer, contentType) => { + const ossKey = scOssKey(lib, relName) + const ossUrl = scOssUrl(ossKey) + try { + await scOssPut(ossKey, buffer, contentType) + return { url: ossUrl, stored: 'oss' } + } catch (err) { + console.error(`[style_check] OSS 上传失败,降级本地: ${err.message}`) + let localUrl = '' + if (lib === 'item_images') { + const abs = path.join(localImageRoot, relName) + fs.mkdirSync(path.dirname(abs), { recursive: true }) + fs.writeFileSync(abs, buffer) + localUrl = `${localImageUrlPrefix}/${relName}` + } else if (lib === 'wave') { + fs.writeFileSync(path.join(waveImageRoot, relName), buffer) + localUrl = `/images/wave_references/${relName}` + } else if (lib === 'korean') { + fs.writeFileSync(path.join(koreanImageRoot, relName), buffer) + localUrl = `/images/korean_references/${relName}` + } else { + throw err + } + registerOssPending(lib, localUrl, ossKey, ossUrl, err.message) + return { url: localUrl, stored: 'local' } + } +} + const guessExtByContentType = (ct) => { const c = String(ct || '').toLowerCase() if (c.includes('image/jpeg') || c.includes('image/jpg')) return '.jpg' @@ -448,17 +566,6 @@ const guessExtByContentType = (ct) => { return '.jpg' } -const saveImageBufferToLocal = (profileKey, itemId, buffer, ext, orderHint = 1) => { - const dir = path.join(localImageRoot, profileKey, String(itemId)) - if (!fs.existsSync(dir)) fs.mkdirSync(dir, { recursive: true }) - const ts = Date.now() - const rnd = Math.floor(Math.random() * 100000) - const fileName = `${ts}_${orderHint}_${rnd}${ext}` - const abs = path.join(dir, fileName) - fs.writeFileSync(abs, buffer) - return `${localImageUrlPrefix}/${profileKey}/${itemId}/${fileName}` -} - const downloadImageToLocal = async (url, profileKey, itemId, orderHint = 1) => { try { const res = await axios.get(url, { @@ -476,7 +583,11 @@ const downloadImageToLocal = async (url, profileKey, itemId, orderHint = 1) => { const buf = Buffer.from(res.data) if (!buf || buf.length < 2000) throw new Error('image_too_small: ' + buf.length) const ext = guessExtByContentType(ct) - return saveImageBufferToLocal(profileKey, itemId, buf, ext, orderHint) + const ts = Date.now() + const rnd = Math.floor(Math.random() * 100000) + const relName = `${profileKey}/${itemId}/${ts}_${orderHint}_${rnd}${ext}` + const saved = await saveImageBufferSmart('item_images', relName, buf, ct.split(';')[0].trim() || 'image/jpeg') + return saved.url } catch (err) { console.error(`downloadImageToLocal error for ${url}:`, err.message) throw err @@ -565,10 +676,14 @@ const parseImageDataUrl = (dataUrl) => { } } -const saveUploadedLibraryImage = (profileKey, itemId, dataUrl) => { +const saveUploadedLibraryImage = async (profileKey, itemId, dataUrl) => { const parsed = parseImageDataUrl(dataUrl) if (!parsed) throw new Error('invalid_image_data_url') - return saveImageBufferToLocal(profileKey, itemId, parsed.buffer, parsed.ext, 999) + const ts = Date.now() + const rnd = Math.floor(Math.random() * 100000) + const relName = `${profileKey}/${itemId}/${ts}_999_${rnd}${parsed.ext}` + const saved = await saveImageBufferSmart('item_images', relName, parsed.buffer, parsed.mime) + return saved.url } let backgroundFetchTimer = null @@ -1150,10 +1265,10 @@ const startWaveCrawler = async () => { logWave(`AI Approved! Saving...`) const fileName = `wave_${Date.now()}_${Math.floor(Math.random()*1000)}.jpg` - const localPath = `/images/wave_references/${fileName}` - fs.writeFileSync(path.join(waveImageRoot, fileName), buf) + const saved = await saveImageBufferSmart('wave', fileName, buf, 'image/jpeg') + logWave(`Saved to ${saved.stored}: ${saved.url}`) - db.prepare(`INSERT OR IGNORE INTO wave_references (image_url, source_url, evaluation, notes) VALUES (?, ?, ?, ?)`).run(localPath, imgUrl, aiRes.evaluation || '', aiRes.notes || '') + db.prepare(`INSERT OR IGNORE INTO wave_references (image_url, source_url, evaluation, notes) VALUES (?, ?, ?, ?)`).run(saved.url, imgUrl, aiRes.evaluation || '', aiRes.notes || '') } // 如果不是在处理手动队列,才翻页 @@ -1259,10 +1374,10 @@ const startKoreanCrawler = async () => { logKorean(`AI Approved! Saving...`) const fileName = `korean_${Date.now()}_${Math.floor(Math.random()*1000)}.jpg` - const localPath = `/images/korean_references/${fileName}` - fs.writeFileSync(path.join(koreanImageRoot, fileName), buf) + const saved = await saveImageBufferSmart('korean', fileName, buf, 'image/jpeg') + logKorean(`Saved to ${saved.stored}: ${saved.url}`) - db.prepare(`INSERT OR IGNORE INTO korean_references (image_url, source_url, evaluation, notes) VALUES (?, ?, ?, ?)`).run(localPath, imgUrl, aiRes.evaluation || '', aiRes.notes || '') + db.prepare(`INSERT OR IGNORE INTO korean_references (image_url, source_url, evaluation, notes) VALUES (?, ?, ?, ?)`).run(saved.url, imgUrl, aiRes.evaluation || '', aiRes.notes || '') } if (koreanCrawlerState.manualQueue.length === 0) { @@ -1995,7 +2110,7 @@ const bindRoutes = app => { let imageUrl = '' let sourceType = 'manual_url_fetch' if (imageDataUrl) { - imageUrl = saveUploadedLibraryImage(profileKey, itemId, imageDataUrl) + imageUrl = await saveUploadedLibraryImage(profileKey, itemId, imageDataUrl) sourceType = 'manual_upload' } else if (directUrl.startsWith('http://') || directUrl.startsWith('https://')) { imageUrl = await downloadImageToLocal(directUrl, profileKey, itemId, Number(b.sort_order || 999))