feat: 新增OSS直传迁移能力,实现图片从本地到OSS的迁移流程

1.  新增OSS直传工具类函数,实现带凭据校验的PUT/HEAD请求
2.  新增oss_migration_log数据表,用于跟踪迁移状态
3.  重构图片保存逻辑,优先上传OSS失败则降级本地并登记补传任务
4.  新增批量迁移脚本,支持状态查看、断点续传、校验和回写DB
5.  新增OSS直链抽查调试脚本
This commit is contained in:
yangxiangyuan
2026-08-06 22:14:20 +08:00
parent 876a3b8768
commit 0f8fe2c2e2
3 changed files with 515 additions and 21 deletions
+136 -21
View File
@@ -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))