Files
Toolbox/src/server/box_backup.js
T

1766 lines
64 KiB
JavaScript

const fs = require('fs')
const path = require('path')
const crypto = require('crypto')
const archiver = require('archiver')
const axios = require('axios')
const Database = require('better-sqlite3')
const unzipper = require('unzipper')
const { logJSON } = require('./logger')
require('./weekly')
const weeklyCore = require('./weekly_core')
const aiLib = require('./ai-lib')
const ROOT_DIR = process.cwd()
const CONFIG_PATH = path.join(ROOT_DIR, 'config', 'box_backup.config.json')
const DATA_DIR = path.join(ROOT_DIR, 'data')
const TEMP_DIR = path.join(ROOT_DIR, 'temp', 'box_backup')
const DB_PATH = path.join(DATA_DIR, 'box_backup.db')
const ensureDir = dir => {
if (!fs.existsSync(dir)) fs.mkdirSync(dir, { recursive: true })
}
const streamToBuffer = stream => new Promise((resolve, reject) => {
const chunks = []
stream.on('data', chunk => chunks.push(Buffer.from(chunk)))
stream.on('end', () => resolve(Buffer.concat(chunks)))
stream.on('error', reject)
})
ensureDir(DATA_DIR)
ensureDir(TEMP_DIR)
const db = new Database(DB_PATH)
db.exec(`
CREATE TABLE IF NOT EXISTS backup_runs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_key TEXT NOT NULL UNIQUE,
mode TEXT NOT NULL,
trigger_source TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'running',
stage TEXT NOT NULL DEFAULT '',
started_at TEXT DEFAULT (datetime('now','+8 hours')),
ended_at TEXT,
summary_json TEXT DEFAULT '',
error_text TEXT DEFAULT ''
);
CREATE INDEX IF NOT EXISTS idx_backup_runs_mode_started ON backup_runs(mode, started_at DESC);
CREATE TABLE IF NOT EXISTS backup_run_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL,
level TEXT NOT NULL DEFAULT 'info',
event_type TEXT NOT NULL DEFAULT '',
message TEXT NOT NULL DEFAULT '',
detail_json TEXT DEFAULT '',
created_at TEXT DEFAULT (datetime('now','+8 hours'))
);
CREATE INDEX IF NOT EXISTS idx_backup_run_events_run ON backup_run_events(run_id, id DESC);
CREATE INDEX IF NOT EXISTS idx_backup_run_events_type ON backup_run_events(event_type, created_at DESC);
CREATE TABLE IF NOT EXISTS backup_alerts (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER,
severity TEXT NOT NULL DEFAULT 'warning',
title TEXT NOT NULL DEFAULT '',
message TEXT NOT NULL DEFAULT '',
status TEXT NOT NULL DEFAULT 'open',
created_at TEXT DEFAULT (datetime('now','+8 hours'))
);
CREATE INDEX IF NOT EXISTS idx_backup_alerts_created ON backup_alerts(created_at DESC);
CREATE TABLE IF NOT EXISTS backup_objects (
id INTEGER PRIMARY KEY AUTOINCREMENT,
backup_type TEXT NOT NULL,
relative_path TEXT NOT NULL,
source_path TEXT NOT NULL DEFAULT '',
oss_key TEXT NOT NULL DEFAULT '',
size_bytes INTEGER DEFAULT 0,
mtime_ms INTEGER DEFAULT 0,
sha256 TEXT DEFAULT '',
upload_status TEXT NOT NULL DEFAULT 'pending',
retry_count INTEGER DEFAULT 0,
last_error TEXT DEFAULT '',
last_uploaded_at TEXT,
is_deleted INTEGER DEFAULT 0,
deleted_at TEXT,
meta_json TEXT DEFAULT '',
UNIQUE(backup_type, relative_path)
);
CREATE INDEX IF NOT EXISTS idx_backup_objects_type_status ON backup_objects(backup_type, upload_status, is_deleted);
CREATE TABLE IF NOT EXISTS backup_manual_records (
id INTEGER PRIMARY KEY AUTOINCREMENT,
record_type TEXT NOT NULL,
note TEXT DEFAULT '',
created_at TEXT DEFAULT (datetime('now','+8 hours'))
);
CREATE INDEX IF NOT EXISTS idx_backup_manual_records_type ON backup_manual_records(record_type, created_at DESC);
CREATE TABLE IF NOT EXISTS backup_config_history (
id INTEGER PRIMARY KEY AUTOINCREMENT,
config_json TEXT NOT NULL,
created_at TEXT DEFAULT (datetime('now','+8 hours'))
);
`)
const insertRunStmt = db.prepare(`
INSERT INTO backup_runs (run_key, mode, trigger_source, status, stage, summary_json)
VALUES (@run_key, @mode, @trigger_source, @status, @stage, @summary_json)
`)
const updateRunStmt = db.prepare(`
UPDATE backup_runs
SET status=@status,
stage=@stage,
ended_at=@ended_at,
summary_json=@summary_json,
error_text=@error_text
WHERE id=@id
`)
const insertEventStmt = db.prepare(`
INSERT INTO backup_run_events (run_id, level, event_type, message, detail_json)
VALUES (@run_id, @level, @event_type, @message, @detail_json)
`)
const insertAlertStmt = db.prepare(`
INSERT INTO backup_alerts (run_id, severity, title, message, status)
VALUES (@run_id, @severity, @title, @message, 'open')
`)
const insertManualRecordStmt = db.prepare(`
INSERT INTO backup_manual_records (record_type, note)
VALUES (?, ?)
`)
const insertConfigHistoryStmt = db.prepare(`
INSERT INTO backup_config_history (config_json) VALUES (?)
`)
const getObjectStmt = db.prepare(`
SELECT * FROM backup_objects WHERE backup_type=? AND relative_path=?
`)
const upsertObjectInsertStmt = db.prepare(`
INSERT INTO backup_objects (
backup_type, relative_path, source_path, oss_key, size_bytes, mtime_ms, sha256,
upload_status, retry_count, last_error, last_uploaded_at, is_deleted, deleted_at, meta_json
) VALUES (
@backup_type, @relative_path, @source_path, @oss_key, @size_bytes, @mtime_ms, @sha256,
@upload_status, @retry_count, @last_error, @last_uploaded_at, @is_deleted, @deleted_at, @meta_json
)
`)
const upsertObjectUpdateStmt = db.prepare(`
UPDATE backup_objects
SET source_path=@source_path,
oss_key=@oss_key,
size_bytes=@size_bytes,
mtime_ms=@mtime_ms,
sha256=@sha256,
upload_status=@upload_status,
retry_count=@retry_count,
last_error=@last_error,
last_uploaded_at=@last_uploaded_at,
is_deleted=@is_deleted,
deleted_at=@deleted_at,
meta_json=@meta_json
WHERE id=@id
`)
const DEFAULT_CONFIG = {
version: 1,
scheduler: {
enabled: true,
poll_seconds: 30,
highfreq_cron: '5 1 * * 0',
lowfreq_cron: '35 1 * * 0'
},
paths: {
highfreq: {
code_paths: [
'src',
'scripts',
'config',
'service',
'deploy',
'tools',
'public/common',
'public/assets',
'public/lib',
'public/icons',
'winform/DocHelper',
'winform/auth_token_winform'
],
root_files: [
'package.json',
'package-lock.json',
'webpack.config.js',
'web.config',
'install_service_8977.bat',
'uninstall_service_8977.bat'
],
tool_root: 'public/tools',
tool_code_extensions: ['.html', '.js', '.css', '.json', '.svg', '.ico', '.txt'],
data_config_extensions: ['.json'],
db_dir: 'data',
db_extensions: ['.db'],
exclude_dirs: [
'node_modules',
'.git',
'Toolbox-Weekly-Backup',
'backups',
'temp',
'logs',
'public/images',
'public/download_files',
'public/home-bg',
'data/psc_chunks',
'data/psc_secrets',
'uploads',
'winform/CalendarReminder_winform'
]
},
lowfreq: {
include_paths: [
'uploads/cloud_notes',
'uploads/private_clipboard',
'uploads/plant_home',
'uploads/ai-lib',
'uploads/ccb_private_funds',
'uploads/doc_cloud_keeper',
'uploads/weekly',
'uploads/expense',
'uploads/cloud',
'public/images/korean_references',
'public/images/wave_references',
'public/images/x_fashion',
'public/images/style_check_images',
'public/download_files'
],
exclude_dirs: [
'node_modules',
'.git',
'Toolbox-Weekly-Backup',
'backups',
'temp'
],
keep_deleted_marker: true
}
},
oss: {
enabled: true,
credential_file_name: 'Toolbox_local_creds.env.local',
access_key_id_key: 'BOX_BACKUP_OSS_ACCESS_KEY_ID',
access_key_secret_key: 'BOX_BACKUP_OSS_ACCESS_KEY_SECRET',
endpoint_key: 'BOX_BACKUP_OSS_ENDPOINT',
bucket_key: 'BOX_BACKUP_OSS_BUCKET',
root_prefix_key: 'BOX_BACKUP_OSS_ROOT_PREFIX',
bucket: '',
endpoint: '',
root_prefix: '',
lowfreq_dir_name: 'DPWJ780-36C4-4811-BE92-47A86DF35E10'
},
retry: {
max_attempts: 3,
base_delay_ms: 2000
},
verification: {
head_check: true,
lowfreq_sample_size: 3,
restore_max_download_mb: 2048
},
notifications: {
enabled: false,
channels: []
}
}
const clone = value => JSON.parse(JSON.stringify(value))
const mergeDefaults = (base, extra) => {
if (!extra || typeof extra !== 'object' || Array.isArray(extra)) return clone(base)
const out = Array.isArray(base) ? [] : {}
const keys = new Set([...Object.keys(base || {}), ...Object.keys(extra || {})])
keys.forEach(key => {
const left = base ? base[key] : undefined
const right = extra ? extra[key] : undefined
if (Array.isArray(left)) out[key] = Array.isArray(right) ? right.slice() : left.slice()
else if (left && typeof left === 'object') out[key] = mergeDefaults(left, right)
else out[key] = right === undefined ? left : right
})
return out
}
const ensureConfigFile = () => {
ensureDir(path.dirname(CONFIG_PATH))
if (!fs.existsSync(CONFIG_PATH)) {
fs.writeFileSync(CONFIG_PATH, JSON.stringify(DEFAULT_CONFIG, null, 2))
}
}
const sanitizeStringArray = list => {
const raw = Array.isArray(list) ? list : []
return raw
.map(item => String(item || '').trim().replace(/\\/g, '/'))
.filter(Boolean)
}
const sanitizeConfig = input => {
const merged = mergeDefaults(DEFAULT_CONFIG, input || {})
merged.scheduler.enabled = merged.scheduler.enabled !== false
merged.scheduler.poll_seconds = Math.max(10, parseInt(String(merged.scheduler.poll_seconds || '30'), 10) || 30)
merged.scheduler.highfreq_cron = String(merged.scheduler.highfreq_cron || DEFAULT_CONFIG.scheduler.highfreq_cron).trim()
merged.scheduler.lowfreq_cron = String(merged.scheduler.lowfreq_cron || DEFAULT_CONFIG.scheduler.lowfreq_cron).trim()
merged.paths.highfreq.code_paths = sanitizeStringArray(merged.paths.highfreq.code_paths)
merged.paths.highfreq.root_files = sanitizeStringArray(merged.paths.highfreq.root_files)
merged.paths.highfreq.tool_code_extensions = sanitizeStringArray(merged.paths.highfreq.tool_code_extensions).map(x => x.startsWith('.') ? x.toLowerCase() : `.${x.toLowerCase()}`)
merged.paths.highfreq.data_config_extensions = sanitizeStringArray(merged.paths.highfreq.data_config_extensions).map(x => x.startsWith('.') ? x.toLowerCase() : `.${x.toLowerCase()}`)
merged.paths.highfreq.db_extensions = sanitizeStringArray(merged.paths.highfreq.db_extensions).map(x => x.startsWith('.') ? x.toLowerCase() : `.${x.toLowerCase()}`)
merged.paths.highfreq.exclude_dirs = sanitizeStringArray(merged.paths.highfreq.exclude_dirs)
merged.paths.highfreq.tool_root = String(merged.paths.highfreq.tool_root || 'public/tools').trim().replace(/\\/g, '/')
merged.paths.highfreq.db_dir = String(merged.paths.highfreq.db_dir || 'data').trim().replace(/\\/g, '/')
merged.paths.lowfreq.include_paths = sanitizeStringArray(merged.paths.lowfreq.include_paths)
merged.paths.lowfreq.exclude_dirs = sanitizeStringArray(merged.paths.lowfreq.exclude_dirs)
merged.paths.lowfreq.keep_deleted_marker = merged.paths.lowfreq.keep_deleted_marker !== false
merged.oss.enabled = merged.oss.enabled !== false
;['credential_file_name', 'access_key_id_key', 'access_key_secret_key', 'endpoint_key', 'bucket_key', 'root_prefix_key', 'bucket', 'endpoint', 'root_prefix', 'lowfreq_dir_name'].forEach(key => {
merged.oss[key] = String(merged.oss[key] || '').trim()
})
merged.retry.max_attempts = Math.max(1, parseInt(String(merged.retry.max_attempts || '3'), 10) || 3)
merged.retry.base_delay_ms = Math.max(500, parseInt(String(merged.retry.base_delay_ms || '2000'), 10) || 2000)
merged.verification.head_check = merged.verification.head_check !== false
merged.verification.lowfreq_sample_size = Math.max(1, parseInt(String(merged.verification.lowfreq_sample_size || '3'), 10) || 3)
merged.verification.restore_max_download_mb = Math.max(128, parseInt(String(merged.verification.restore_max_download_mb || '2048'), 10) || 2048)
merged.notifications.enabled = !!merged.notifications.enabled
if (!Array.isArray(merged.notifications.channels)) merged.notifications.channels = []
return merged
}
ensureConfigFile()
let runtimeConfig = sanitizeConfig(JSON.parse(fs.readFileSync(CONFIG_PATH, 'utf-8')))
const beijingParts = date => {
const parts = new Intl.DateTimeFormat('sv-SE', {
timeZone: 'Asia/Shanghai',
hour12: false,
year: 'numeric',
month: '2-digit',
day: '2-digit',
hour: '2-digit',
minute: '2-digit',
second: '2-digit'
}).formatToParts(date || new Date())
const map = {}
parts.forEach(part => {
if (part.type && part.type !== 'literal') map[part.type] = part.value
})
return map
}
const formatBeijingDateTime = date => {
const p = beijingParts(date || new Date())
return `${p.year}-${p.month}-${p.day} ${p.hour}:${p.minute}:${p.second}`
}
const formatCompactStamp = date => {
const p = beijingParts(date || new Date())
return `${p.year}${p.month}${p.day}-${p.hour}${p.minute}${p.second}`
}
const sleep = ms => new Promise(resolve => setTimeout(resolve, ms))
const makeGuid = () => crypto.randomUUID().toUpperCase()
const makeRunKey = mode => `${String(mode || 'run')}-${formatCompactStamp(new Date())}-${crypto.randomBytes(4).toString('hex')}`
const parseCronField = (field, min, max) => {
const value = String(field || '').trim()
if (!value || value === '*') return null
const set = new Set()
value.split(',').forEach(part => {
const token = String(part || '').trim()
if (!token) return
if (/^\d+$/.test(token)) {
const n = parseInt(token, 10)
if (n >= min && n <= max) set.add(n)
return
}
const stepMatch = token.match(/^\*\/(\d+)$/)
if (stepMatch) {
const step = Math.max(1, parseInt(stepMatch[1], 10) || 1)
for (let n = min; n <= max; n += step) set.add(n)
return
}
const rangeMatch = token.match(/^(\d+)-(\d+)$/)
if (rangeMatch) {
const start = parseInt(rangeMatch[1], 10)
const end = parseInt(rangeMatch[2], 10)
for (let n = start; n <= end; n += 1) {
if (n >= min && n <= max) set.add(n)
}
}
})
return set
}
const cronMatches = (expr, date) => {
const parts = String(expr || '').trim().split(/\s+/)
if (parts.length !== 5) return false
const minuteSet = parseCronField(parts[0], 0, 59)
const hourSet = parseCronField(parts[1], 0, 23)
const domSet = parseCronField(parts[2], 1, 31)
const monthSet = parseCronField(parts[3], 1, 12)
const dowSet = parseCronField(parts[4], 0, 6)
const p = beijingParts(date)
const target = new Date(`${p.year}-${p.month}-${p.day}T${p.hour}:${p.minute}:${p.second}+08:00`)
const minute = Number(p.minute)
const hour = Number(p.hour)
const dom = Number(p.day)
const month = Number(p.month)
const dow = target.getUTCDay()
if (minuteSet && !minuteSet.has(minute)) return false
if (hourSet && !hourSet.has(hour)) return false
if (domSet && !domSet.has(dom)) return false
if (monthSet && !monthSet.has(month)) return false
if (dowSet && !dowSet.has(dow)) return false
return true
}
const getNextCronTime = (expr, fromDate = new Date()) => {
const base = new Date(fromDate.getTime())
base.setSeconds(0, 0)
for (let i = 1; i <= 60 * 24 * 21; i += 1) {
const next = new Date(base.getTime() + i * 60 * 1000)
if (cronMatches(expr, next)) return formatBeijingDateTime(next)
}
return ''
}
const readLocalCreds = cfg => {
try {
const os = require('os')
const userProfile = process.env.USERPROFILE || process.env.HOME || os.homedir()
const fileName = ((cfg || {}).oss || {}).credential_file_name || 'Toolbox_local_creds.env.local'
const candidates = [path.join(userProfile, fileName)]
const usersRoot = path.join(path.parse(userProfile).root, 'Users')
try {
if (fs.existsSync(usersRoot)) {
fs.readdirSync(usersRoot).forEach(name => {
const dir = path.join(usersRoot, name)
if (dir !== userProfile && fs.existsSync(dir)) {
candidates.push(path.join(dir, fileName))
}
})
}
} catch {}
for (const candidate of candidates) {
if (!fs.existsSync(candidate)) continue
const text = fs.readFileSync(candidate, 'utf-8')
const out = { __file: candidate }
text.split(/\r?\n/).forEach(line => {
const trimmed = line.trim()
if (!trimmed || trimmed.startsWith('#')) return
const idx = trimmed.indexOf('=')
if (idx <= 0) return
out[trimmed.slice(0, idx).trim()] = trimmed.slice(idx + 1).trim()
})
return out
}
} catch {}
return {}
}
const resolveOssSettings = cfg => {
const local = readLocalCreds(cfg)
const bucket = String(cfg.oss.bucket || local[cfg.oss.bucket_key] || '').trim()
const endpoint = String(cfg.oss.endpoint || local[cfg.oss.endpoint_key] || '').trim()
const accessKeyId = String(local[cfg.oss.access_key_id_key] || '').trim()
const accessKeySecret = String(local[cfg.oss.access_key_secret_key] || '').trim()
const rootPrefix = String(cfg.oss.root_prefix || local[cfg.oss.root_prefix_key] || '').trim().replace(/^\/+|\/+$/g, '')
const lowfreqDirName = String(cfg.oss.lowfreq_dir_name || 'DPWJ780-36C4-4811-BE92-47A86DF35E10').trim().replace(/^\/+|\/+$/g, '')
return {
credential_file: local.__file || '',
bucket,
endpoint,
access_key_id: accessKeyId,
access_key_secret: accessKeySecret,
root_prefix: rootPrefix,
lowfreq_dir_name: lowfreqDirName,
ready: !!(bucket && endpoint && accessKeyId && accessKeySecret && rootPrefix && lowfreqDirName)
}
}
const currentRunState = {
run_id: 0,
run_key: '',
mode: '',
trigger_source: '',
status: 'idle',
stage: '',
message: '',
percent: 0,
started_at: '',
updated_at: '',
stats: {}
}
let currentRunPromise = null
let schedulerTimer = null
let lastSchedulerMinute = ''
const setCurrentRun = patch => {
Object.assign(currentRunState, patch || {})
currentRunState.updated_at = formatBeijingDateTime(new Date())
}
const resetCurrentRun = () => {
setCurrentRun({
run_id: 0,
run_key: '',
mode: '',
trigger_source: '',
status: 'idle',
stage: '',
message: '',
percent: 0,
started_at: '',
stats: {}
})
}
resetCurrentRun()
const recoverInterruptedRuns = () => {
try {
db.prepare(`
UPDATE backup_runs
SET status='failed',
stage='interrupted',
ended_at=?,
error_text=CASE
WHEN error_text IS NULL OR error_text='' THEN '进程中断,任务未正常结束'
ELSE error_text
END
WHERE status='running'
`).run(formatBeijingDateTime(new Date()))
} catch {}
}
recoverInterruptedRuns()
const runSummaryBase = () => ({
highfreq: {
run_guid: '',
archive_path: '',
zip_size_bytes: 0,
file_count: 0,
uploaded: false,
oss_key: '',
meta_oss_key: '',
sha256: '',
verification_ok: false
},
lowfreq: {
scanned: 0,
changed: 0,
uploaded: 0,
skipped: 0,
failed: 0,
deleted_marked: 0
},
verify: {
highfreq_head_ok: false,
sampled: 0,
passed: 0,
failed: 0,
samples: []
},
restore: {
source_run_id: 0,
source_oss_key: '',
downloaded_bytes: 0,
sha256_match: false,
manifest_present: false,
extracted_entries: [],
restore_dir: ''
}
})
const addEvent = (runId, level, eventType, message, detail) => {
const payload = {
run_id: runId,
level: String(level || 'info'),
event_type: String(eventType || ''),
message: String(message || ''),
detail_json: detail ? JSON.stringify(detail) : ''
}
insertEventStmt.run(payload)
try { logJSON(`box_backup.${eventType || 'event'}`, { run_id: runId, level, message, detail: detail || null }, 'box_backup') } catch {}
}
const addAlert = (runId, severity, title, message) => {
insertAlertStmt.run({
run_id: runId || null,
severity: String(severity || 'warning'),
title: String(title || ''),
message: String(message || '')
})
}
const relativeFromRoot = fullPath => path.relative(ROOT_DIR, fullPath).replace(/\\/g, '/')
const sanitizeOssName = value => String(value || '').trim().replace(/[^A-Za-z0-9._-]/g, '_')
const parseObjectMeta = row => {
try { return row && row.meta_json ? JSON.parse(row.meta_json) : {} } catch { return {} }
}
const getEncryptedFileInfo = row => {
const rel = String((row && row.relative_path) || '').replace(/\\/g, '/')
const sourcePath = String((row && row.source_path) || '')
if (rel.startsWith('uploads/weekly/attachments/')) {
return {
encrypted: true,
provider: 'weekly',
label: 'weekly加密附件',
encrypted_path: sourcePath || path.join(ROOT_DIR, rel)
}
}
if (rel.startsWith('uploads/ai-lib/attachments/')) {
return {
encrypted: true,
provider: 'ai-lib',
label: 'ai-lib加密附件',
encrypted_path: sourcePath || path.join(ROOT_DIR, rel)
}
}
return {
encrypted: false,
provider: '',
label: '',
encrypted_path: sourcePath || path.join(ROOT_DIR, rel)
}
}
const getLowfreqStoredName = (existing, file) => {
const meta = parseObjectMeta(existing)
const current = sanitizeOssName(meta.stored_name || '')
if (current) return current
const ext = path.extname(String(file.relativePath || file.fullPath || '')).toLowerCase()
return `${makeGuid()}${ext}`
}
const buildLowfreqOssKey = (oss, storedName) => `${oss.root_prefix}/${oss.lowfreq_dir_name}/lowfreq/${sanitizeOssName(storedName)}`
const isInsideAnyExcluded = (relPath, excludeDirs) => {
const clean = String(relPath || '').replace(/\\/g, '/')
return (excludeDirs || []).some(dir => {
const prefix = String(dir || '').trim().replace(/\\/g, '/').replace(/\/+$/, '')
return prefix && (clean === prefix || clean.startsWith(`${prefix}/`))
})
}
const walkFiles = (rootPath, options = {}) => {
const out = []
const base = path.resolve(ROOT_DIR, rootPath)
if (!fs.existsSync(base)) return out
const excludeDirs = options.excludeDirs || []
const allowedExts = options.allowedExts ? new Set(options.allowedExts.map(x => String(x).toLowerCase())) : null
const stack = [base]
while (stack.length) {
const cur = stack.pop()
let entries = []
try { entries = fs.readdirSync(cur, { withFileTypes: true }) } catch { continue }
for (const entry of entries) {
const full = path.join(cur, entry.name)
const rel = relativeFromRoot(full)
if (entry.isDirectory()) {
if (isInsideAnyExcluded(rel, excludeDirs)) continue
stack.push(full)
continue
}
if (!entry.isFile()) continue
if (isInsideAnyExcluded(rel, excludeDirs)) continue
const ext = path.extname(entry.name || '').toLowerCase()
if (allowedExts && !allowedExts.has(ext)) continue
out.push({ fullPath: full, relativePath: rel })
}
}
return out
}
const collectHighfreqFiles = async (cfg, runId, tempRunDir) => {
const items = []
const seen = new Set()
const addFile = (fullPath, archivePath) => {
const resolved = path.resolve(fullPath)
const archiveName = String(archivePath || relativeFromRoot(fullPath)).replace(/\\/g, '/')
const key = `${resolved}>>${archiveName}`
if (seen.has(key)) return
seen.add(key)
items.push({ fullPath: resolved, archivePath: archiveName })
}
for (const relFile of cfg.paths.highfreq.root_files) {
const full = path.join(ROOT_DIR, relFile)
if (fs.existsSync(full) && fs.statSync(full).isFile()) addFile(full, relFile)
}
for (const relDir of cfg.paths.highfreq.code_paths) {
const files = walkFiles(relDir, { excludeDirs: cfg.paths.highfreq.exclude_dirs })
files.forEach(file => addFile(file.fullPath, file.relativePath))
}
const toolFiles = walkFiles(cfg.paths.highfreq.tool_root, {
excludeDirs: cfg.paths.highfreq.exclude_dirs,
allowedExts: cfg.paths.highfreq.tool_code_extensions
})
toolFiles.forEach(file => addFile(file.fullPath, file.relativePath))
const dataDir = path.join(ROOT_DIR, cfg.paths.highfreq.db_dir)
if (fs.existsSync(dataDir)) {
const topEntries = fs.readdirSync(dataDir, { withFileTypes: true })
for (const entry of topEntries) {
const full = path.join(dataDir, entry.name)
if (!entry.isFile()) continue
const ext = path.extname(entry.name || '').toLowerCase()
if (cfg.paths.highfreq.data_config_extensions.includes(ext)) addFile(full, relativeFromRoot(full))
if (!cfg.paths.highfreq.db_extensions.includes(ext)) continue
const snapshotDir = path.join(tempRunDir, 'db_snapshots')
ensureDir(snapshotDir)
const target = path.join(snapshotDir, entry.name)
try {
const sourceDb = new Database(full, { readonly: true, fileMustExist: true })
try {
await sourceDb.backup(target)
} finally {
sourceDb.close()
}
addFile(target, `data_snapshots/${entry.name}`)
addEvent(runId, 'info', 'db_snapshot', `数据库快照完成: ${entry.name}`)
} catch (error) {
addEvent(runId, 'error', 'db_snapshot_failed', `数据库快照失败: ${entry.name}`, { error: String(error && error.message ? error.message : error) })
throw error
}
}
}
return items.sort((a, b) => a.archivePath.localeCompare(b.archivePath))
}
const sha256File = filePath => new Promise((resolve, reject) => {
const hash = crypto.createHash('sha256')
const stream = fs.createReadStream(filePath)
stream.on('data', chunk => hash.update(chunk))
stream.on('error', reject)
stream.on('end', () => resolve(hash.digest('hex')))
})
const createHighfreqArchive = async (cfg, runId, tempRunDir) => {
const files = await collectHighfreqFiles(cfg, runId, tempRunDir)
const archiveName = `toolbox-highfreq-${formatCompactStamp(new Date())}.zip`
const archivePath = path.join(tempRunDir, archiveName)
const manifest = {
generated_at: formatBeijingDateTime(new Date()),
mode: 'highfreq',
files: []
}
for (const file of files) {
let stat = null
try { stat = fs.statSync(file.fullPath) } catch {}
manifest.files.push({
archive_path: file.archivePath,
source_path: relativeFromRoot(file.fullPath),
size_bytes: stat ? stat.size : 0,
mtime_ms: stat ? stat.mtimeMs : 0
})
}
await new Promise((resolve, reject) => {
const output = fs.createWriteStream(archivePath)
const archive = archiver('zip', { zlib: { level: 9 } })
output.on('close', resolve)
output.on('error', reject)
archive.on('error', reject)
archive.pipe(output)
files.forEach(file => archive.file(file.fullPath, { name: file.archivePath }))
archive.append(JSON.stringify(manifest, null, 2), { name: 'manifest.json' })
archive.finalize()
})
const sha256 = await sha256File(archivePath)
const stat = fs.statSync(archivePath)
return {
archivePath,
archiveName,
fileCount: files.length,
sizeBytes: stat.size,
sha256
}
}
const encodeOssKey = key => String(key || '').split('/').map(part => encodeURIComponent(part)).join('/')
const ossObjectUrl = (bucket, endpoint, key) => `https://${bucket}.${endpoint}/${encodeOssKey(key)}`
const ossObjectResource = (bucket, key) => `/${bucket}/${String(key || '').replace(/^\/+/, '')}`
const signOssRequest = ({ accessKeySecret, method, contentMd5 = '', contentType = '', date, resource, headers = {} }) => {
const canonicalHeaders = Object.keys(headers)
.filter(key => String(key || '').toLowerCase().startsWith('x-oss-'))
.sort((a, b) => a.localeCompare(b))
.map(key => `${String(key).toLowerCase()}:${String(headers[key] || '').trim()}\n`)
.join('')
const stringToSign = [
String(method || 'GET').toUpperCase(),
contentMd5,
contentType,
date,
`${canonicalHeaders}${resource}`
].join('\n')
return crypto.createHmac('sha1', accessKeySecret).update(stringToSign).digest('base64')
}
const headObject = async (oss, key) => {
const date = new Date().toUTCString()
const resource = ossObjectResource(oss.bucket, key)
const authorization = `OSS ${oss.access_key_id}:${signOssRequest({ accessKeySecret: oss.access_key_secret, method: 'HEAD', date, resource })}`
const response = await axios({
method: 'HEAD',
url: ossObjectUrl(oss.bucket, oss.endpoint, key),
headers: { Date: date, Authorization: authorization },
timeout: 60000,
validateStatus: status => status >= 200 && status < 300
})
return response.headers || {}
}
const uploadFileToOss = async ({ oss, localPath, ossKey, contentType }) => {
const stat = fs.statSync(localPath)
const date = new Date().toUTCString()
const resource = ossObjectResource(oss.bucket, ossKey)
const authorization = `OSS ${oss.access_key_id}:${signOssRequest({
accessKeySecret: oss.access_key_secret,
method: 'PUT',
contentType,
date,
resource
})}`
await axios({
method: 'PUT',
url: ossObjectUrl(oss.bucket, oss.endpoint, ossKey),
headers: {
Date: date,
Authorization: authorization,
'Content-Type': contentType,
'Content-Length': stat.size
},
data: fs.createReadStream(localPath),
maxBodyLength: Infinity,
maxContentLength: Infinity,
timeout: 10 * 60 * 1000,
validateStatus: status => status >= 200 && status < 300
})
return stat.size
}
const downloadObjectToFile = async ({ oss, ossKey, targetPath }) => {
ensureDir(path.dirname(targetPath))
const date = new Date().toUTCString()
const resource = ossObjectResource(oss.bucket, ossKey)
const authorization = `OSS ${oss.access_key_id}:${signOssRequest({
accessKeySecret: oss.access_key_secret,
method: 'GET',
date,
resource
})}`
const response = await axios({
method: 'GET',
url: ossObjectUrl(oss.bucket, oss.endpoint, ossKey),
headers: {
Date: date,
Authorization: authorization
},
responseType: 'stream',
timeout: 10 * 60 * 1000,
validateStatus: status => status >= 200 && status < 300
})
await new Promise((resolve, reject) => {
const stream = fs.createWriteStream(targetPath)
response.data.pipe(stream)
response.data.on('error', reject)
stream.on('error', reject)
stream.on('finish', resolve)
})
return fs.statSync(targetPath).size
}
const createOssDownloadStream = async ({ oss, ossKey }) => {
const date = new Date().toUTCString()
const resource = ossObjectResource(oss.bucket, ossKey)
const authorization = `OSS ${oss.access_key_id}:${signOssRequest({
accessKeySecret: oss.access_key_secret,
method: 'GET',
date,
resource
})}`
const response = await axios({
method: 'GET',
url: ossObjectUrl(oss.bucket, oss.endpoint, ossKey),
headers: {
Date: date,
Authorization: authorization
},
responseType: 'stream',
timeout: 10 * 60 * 1000,
validateStatus: status => status >= 200 && status < 300
})
return {
stream: response.data,
headers: response.headers || {}
}
}
const retryWrap = async (cfg, runId, type, message, workFn) => {
const maxAttempts = cfg.retry.max_attempts
const baseDelayMs = cfg.retry.base_delay_ms
let lastError = null
for (let attempt = 1; attempt <= maxAttempts; attempt += 1) {
try {
if (attempt > 1) addEvent(runId, 'warning', 'retry', `${message},第 ${attempt} 次尝试`)
return await workFn(attempt)
} catch (error) {
lastError = error
addEvent(runId, 'error', type, `${message}失败`, { attempt, error: String(error && error.message ? error.message : error) })
if (attempt < maxAttempts) await sleep(baseDelayMs * attempt)
}
}
throw lastError || new Error(`${message}失败`)
}
const upsertObjectRecord = payload => {
const data = {
backup_type: String(payload.backup_type || ''),
relative_path: String(payload.relative_path || ''),
source_path: String(payload.source_path || ''),
oss_key: String(payload.oss_key || ''),
size_bytes: parseInt(String(payload.size_bytes || '0'), 10) || 0,
mtime_ms: Math.floor(Number(payload.mtime_ms || 0)) || 0,
sha256: String(payload.sha256 || ''),
upload_status: String(payload.upload_status || 'pending'),
retry_count: parseInt(String(payload.retry_count || '0'), 10) || 0,
last_error: String(payload.last_error || ''),
last_uploaded_at: payload.last_uploaded_at ? String(payload.last_uploaded_at) : null,
is_deleted: payload.is_deleted ? 1 : 0,
deleted_at: payload.deleted_at ? String(payload.deleted_at) : null,
meta_json: payload.meta_json ? JSON.stringify(payload.meta_json) : ''
}
const existing = getObjectStmt.get(data.backup_type, data.relative_path)
if (!existing) {
upsertObjectInsertStmt.run(data)
return
}
upsertObjectUpdateStmt.run({ id: existing.id, ...data })
}
const buildLowfreqFileList = cfg => {
const files = []
const seen = new Set()
for (const relDir of cfg.paths.lowfreq.include_paths) {
walkFiles(relDir, { excludeDirs: cfg.paths.lowfreq.exclude_dirs }).forEach(file => {
const rel = file.relativePath
if (seen.has(rel)) return
seen.add(rel)
files.push(file)
})
}
return files.sort((a, b) => a.relativePath.localeCompare(b.relativePath))
}
const markDeletedLowfreqObjects = (cfg, runId, liveRelativePaths) => {
const rows = db.prepare(`SELECT * FROM backup_objects WHERE backup_type='lowfreq' AND is_deleted=0`).all()
let count = 0
rows.forEach(row => {
const tracked = String(row.relative_path || '')
if (!tracked) return
const covered = cfg.paths.lowfreq.include_paths.some(prefix => tracked === prefix || tracked.startsWith(`${prefix}/`))
if (!covered) return
if (liveRelativePaths.has(tracked)) return
count += 1
upsertObjectRecord({
backup_type: 'lowfreq',
relative_path: tracked,
source_path: row.source_path || '',
oss_key: row.oss_key || '',
size_bytes: row.size_bytes || 0,
mtime_ms: row.mtime_ms || 0,
sha256: row.sha256 || '',
upload_status: row.upload_status || 'uploaded',
retry_count: row.retry_count || 0,
last_error: row.last_error || '',
last_uploaded_at: row.last_uploaded_at || null,
is_deleted: 1,
deleted_at: formatBeijingDateTime(new Date()),
meta_json: Object.assign({}, parseObjectMeta(row), { deleted_marker: true })
})
addEvent(runId, 'warning', 'deleted_marker', `本地已删除,OSS保留并标记:${tracked}`)
})
return count
}
const uploadLowfreqFile = async (cfg, oss, runId, file) => {
const stat = fs.statSync(file.fullPath)
const existing = getObjectStmt.get('lowfreq', file.relativePath)
const storedName = getLowfreqStoredName(existing, file)
const ossKey = buildLowfreqOssKey(oss, storedName)
const layoutChanged = existing && String(existing.oss_key || '') !== ossKey
const unchanged = existing &&
!layoutChanged &&
Number(existing.is_deleted || 0) === 0 &&
Number(existing.size_bytes || 0) === Number(stat.size || 0) &&
Number(existing.mtime_ms || 0) === Math.floor(Number(stat.mtimeMs || 0)) &&
String(existing.upload_status || '') === 'uploaded'
if (unchanged) {
return { changed: false, skipped: true }
}
const sha256 = await sha256File(file.fullPath)
if (existing &&
!layoutChanged &&
Number(existing.is_deleted || 0) === 0 &&
String(existing.sha256 || '') === sha256 &&
String(existing.upload_status || '') === 'uploaded') {
upsertObjectRecord({
backup_type: 'lowfreq',
relative_path: file.relativePath,
source_path: file.fullPath,
oss_key: existing.oss_key || ossKey,
size_bytes: stat.size,
mtime_ms: stat.mtimeMs,
sha256,
upload_status: 'uploaded',
retry_count: existing.retry_count || 0,
last_error: '',
last_uploaded_at: existing.last_uploaded_at || null,
is_deleted: 0,
deleted_at: null,
meta_json: {
stored_name: storedName,
storage_layout: 'guid-lowfreq-v1',
verified: false
}
})
return { changed: false, skipped: true }
}
try {
await retryWrap(cfg, runId, 'lowfreq_upload_failed', `低频对象上传 ${file.relativePath}`, async () => {
await uploadFileToOss({ oss, localPath: file.fullPath, ossKey, contentType: 'application/octet-stream' })
if (cfg.verification.head_check) {
const headers = await headObject(oss, ossKey)
const remoteSize = parseInt(String(headers['content-length'] || '0'), 10) || 0
if (remoteSize !== stat.size) throw new Error(`HEAD 校验大小不一致: ${remoteSize} != ${stat.size}`)
}
})
upsertObjectRecord({
backup_type: 'lowfreq',
relative_path: file.relativePath,
source_path: file.fullPath,
oss_key: ossKey,
size_bytes: stat.size,
mtime_ms: stat.mtimeMs,
sha256,
upload_status: 'uploaded',
retry_count: 0,
last_error: '',
last_uploaded_at: formatBeijingDateTime(new Date()),
is_deleted: 0,
deleted_at: null,
meta_json: {
stored_name: storedName,
storage_layout: 'guid-lowfreq-v1',
verified: !!cfg.verification.head_check
}
})
return { changed: true, uploaded: true }
} catch (error) {
upsertObjectRecord({
backup_type: 'lowfreq',
relative_path: file.relativePath,
source_path: file.fullPath,
oss_key: ossKey,
size_bytes: stat.size,
mtime_ms: stat.mtimeMs,
sha256,
upload_status: 'failed',
retry_count: cfg.retry.max_attempts,
last_error: String(error && error.message ? error.message : error),
last_uploaded_at: existing && existing.last_uploaded_at ? existing.last_uploaded_at : null,
is_deleted: 0,
deleted_at: null,
meta_json: {
stored_name: storedName,
storage_layout: 'guid-lowfreq-v1',
verified: false
}
})
addAlert(runId, 'warning', '低频增量备份失败', `${file.relativePath} 上传失败:${String(error && error.message ? error.message : error)}`)
return { changed: true, failed: true }
}
}
const runHighfreq = async (cfg, oss, runId, tempRunDir, summary) => {
setCurrentRun({ stage: 'highfreq_prepare', message: '正在准备高频备份', percent: 5 })
addEvent(runId, 'info', 'stage', '开始执行高频备份')
const archive = await createHighfreqArchive(cfg, runId, tempRunDir)
const runGuid = makeGuid()
const highfreqBaseKey = `${oss.root_prefix}/${runGuid}/highfreq`
const ossKey = `${highfreqBaseKey}/archive.zip`
const metaOssKey = `${highfreqBaseKey}/run_meta.json`
const metaPath = path.join(tempRunDir, 'run_meta.json')
summary.highfreq.archive_path = archive.archivePath
summary.highfreq.run_guid = runGuid
summary.highfreq.zip_size_bytes = archive.sizeBytes
summary.highfreq.file_count = archive.fileCount
summary.highfreq.sha256 = archive.sha256
addEvent(runId, 'info', 'archive_created', `高频归档完成,共 ${archive.fileCount} 个文件`)
fs.writeFileSync(metaPath, JSON.stringify({
run_guid: runGuid,
created_at: formatBeijingDateTime(new Date()),
archive_oss_key: ossKey,
archive_sha256: archive.sha256,
zip_size_bytes: archive.sizeBytes,
file_count: archive.fileCount,
source_run_id: runId,
storage_layout: 'guid-highfreq-v1'
}, null, 2))
setCurrentRun({ stage: 'highfreq_upload', message: '正在上传高频 ZIP 到 OSS', percent: 35 })
await retryWrap(cfg, runId, 'highfreq_upload_failed', '高频 ZIP 上传', async () => {
await uploadFileToOss({ oss, localPath: archive.archivePath, ossKey, contentType: 'application/zip' })
await uploadFileToOss({ oss, localPath: metaPath, ossKey: metaOssKey, contentType: 'application/json; charset=utf-8' })
if (cfg.verification.head_check) {
const headers = await headObject(oss, ossKey)
const remoteSize = parseInt(String(headers['content-length'] || '0'), 10) || 0
if (remoteSize !== archive.sizeBytes) throw new Error(`HEAD 校验大小不一致: ${remoteSize} != ${archive.sizeBytes}`)
summary.highfreq.verification_ok = true
}
})
summary.highfreq.uploaded = true
summary.highfreq.oss_key = ossKey
summary.highfreq.meta_oss_key = metaOssKey
addEvent(runId, 'info', 'highfreq_uploaded', `高频备份上传完成:${ossKey}`, { size_bytes: archive.sizeBytes, sha256: archive.sha256, run_guid: runGuid, meta_oss_key: metaOssKey })
}
const runLowfreq = async (cfg, oss, runId, summary) => {
setCurrentRun({ stage: 'lowfreq_scan', message: '正在扫描低频目录', percent: 50 })
addEvent(runId, 'info', 'stage', '开始执行低频增量备份')
const files = buildLowfreqFileList(cfg)
const liveSet = new Set(files.map(file => file.relativePath))
summary.lowfreq.scanned = files.length
let uploaded = 0
let skipped = 0
let changed = 0
let failed = 0
for (let i = 0; i < files.length; i += 1) {
const file = files[i]
const result = await uploadLowfreqFile(cfg, oss, runId, file)
if (result.changed) changed += 1
if (result.uploaded) uploaded += 1
if (result.skipped) skipped += 1
if (result.failed) failed += 1
if ((i + 1) % 10 === 0 || i === files.length - 1) {
setCurrentRun({
stage: 'lowfreq_upload',
message: `正在处理低频文件 ${i + 1}/${files.length}`,
percent: files.length ? Math.min(95, 55 + Math.floor(((i + 1) / files.length) * 40)) : 95,
stats: Object.assign({}, currentRunState.stats, {
lowfreq_total: files.length,
lowfreq_done: i + 1,
lowfreq_uploaded: uploaded,
lowfreq_skipped: skipped,
lowfreq_failed: failed
})
})
}
}
const deletedMarked = cfg.paths.lowfreq.keep_deleted_marker ? markDeletedLowfreqObjects(cfg, runId, liveSet) : 0
summary.lowfreq.changed = changed
summary.lowfreq.uploaded = uploaded
summary.lowfreq.skipped = skipped
summary.lowfreq.failed = failed
summary.lowfreq.deleted_marked = deletedMarked
addEvent(runId, failed > 0 ? 'warning' : 'info', 'lowfreq_done', `低频增量完成:上传 ${uploaded},跳过 ${skipped},失败 ${failed},标记删除 ${deletedMarked}`)
}
const parseRunSummary = row => {
try { return row && row.summary_json ? JSON.parse(row.summary_json) : {} } catch { return {} }
}
const latestRunByExactMode = mode => db.prepare(`
SELECT * FROM backup_runs
WHERE mode=?
AND status <> 'running'
ORDER BY id DESC LIMIT 1
`).get(mode)
const latestSuccessfulHighfreqRun = () => {
const rows = db.prepare(`
SELECT * FROM backup_runs
WHERE mode IN ('high','all')
AND status IN ('success','partial_success')
ORDER BY id DESC
LIMIT 20
`).all()
for (const row of rows) {
const summary = parseRunSummary(row)
if (summary && summary.highfreq && summary.highfreq.uploaded && summary.highfreq.oss_key) {
return { row, summary }
}
}
return null
}
const listLowfreqSampleCandidates = limit => db.prepare(`
SELECT * FROM backup_objects
WHERE backup_type='lowfreq'
AND upload_status='uploaded'
AND is_deleted=0
AND oss_key <> ''
ORDER BY last_uploaded_at DESC, id DESC
LIMIT ?
`).all(Math.max(1, parseInt(String(limit || '3'), 10) || 3))
const runVerify = async (cfg, oss, runId, tempRunDir, summary) => {
setCurrentRun({ stage: 'verify_prepare', message: '正在准备回读抽检', percent: 8 })
addEvent(runId, 'info', 'stage', '开始执行回读抽检')
const latestHigh = latestSuccessfulHighfreqRun()
if (latestHigh && latestHigh.summary && latestHigh.summary.highfreq && latestHigh.summary.highfreq.oss_key) {
setCurrentRun({ stage: 'verify_highfreq_head', message: '正在校验最近一次高频归档', percent: 20 })
const headers = await headObject(oss, latestHigh.summary.highfreq.oss_key)
const remoteSize = parseInt(String(headers['content-length'] || '0'), 10) || 0
const expectedSize = parseInt(String(latestHigh.summary.highfreq.zip_size_bytes || '0'), 10) || 0
if (expectedSize > 0 && remoteSize !== expectedSize) {
throw new Error(`最近高频归档 HEAD 校验失败:${remoteSize} != ${expectedSize}`)
}
summary.verify.highfreq_head_ok = true
addEvent(runId, 'info', 'verify_highfreq_head_ok', '最近一次高频归档 HEAD 校验通过', { oss_key: latestHigh.summary.highfreq.oss_key, remote_size: remoteSize })
}
const samples = listLowfreqSampleCandidates(cfg.verification.lowfreq_sample_size)
summary.verify.sampled = samples.length
for (let i = 0; i < samples.length; i += 1) {
const row = samples[i]
setCurrentRun({
stage: 'verify_lowfreq_sample',
message: `正在抽检低频对象 ${i + 1}/${samples.length}`,
percent: samples.length ? 35 + Math.floor(((i + 1) / samples.length) * 55) : 90
})
const target = path.join(tempRunDir, 'verify_samples', `${i + 1}-${path.basename(String(row.relative_path || 'sample.bin'))}`)
try {
const downloadedBytes = await retryWrap(cfg, runId, 'verify_sample_download_failed', `抽检对象下载 ${row.relative_path}`, async () => {
return await downloadObjectToFile({ oss, ossKey: row.oss_key, targetPath: target })
})
const actualHash = await sha256File(target)
const passed = !row.sha256 || actualHash === row.sha256
summary.verify.samples.push({
relative_path: row.relative_path,
oss_key: row.oss_key,
downloaded_bytes: downloadedBytes,
expected_sha256: row.sha256 || '',
actual_sha256: actualHash,
ok: passed
})
if (!passed) {
summary.verify.failed += 1
addEvent(runId, 'error', 'verify_sample_hash_mismatch', `抽检对象 hash 不一致:${row.relative_path}`)
} else {
summary.verify.passed += 1
addEvent(runId, 'info', 'verify_sample_ok', `抽检通过:${row.relative_path}`)
}
} catch (error) {
summary.verify.failed += 1
summary.verify.samples.push({
relative_path: row.relative_path,
oss_key: row.oss_key,
ok: false,
error: String(error && error.message ? error.message : error)
})
addEvent(runId, 'error', 'verify_sample_failed', `抽检失败:${row.relative_path}`, { error: String(error && error.message ? error.message : error) })
}
}
}
const runRestoreRehearsal = async (cfg, oss, runId, tempRunDir, summary) => {
setCurrentRun({ stage: 'restore_prepare', message: '正在准备恢复演练', percent: 8 })
addEvent(runId, 'info', 'stage', '开始执行恢复演练')
const latestHigh = latestSuccessfulHighfreqRun()
if (!latestHigh) throw new Error('没有可用于恢复演练的高频归档')
const high = latestHigh.summary.highfreq || {}
const expectedBytes = parseInt(String(high.zip_size_bytes || '0'), 10) || 0
const maxBytes = (parseInt(String(cfg.verification.restore_max_download_mb || '2048'), 10) || 2048) * 1024 * 1024
if (expectedBytes > maxBytes) {
throw new Error(`最近高频归档过大,超过恢复演练限制:${expectedBytes} > ${maxBytes}`)
}
summary.restore.source_run_id = latestHigh.row.id
summary.restore.source_oss_key = String(high.oss_key || '')
const downloadPath = path.join(tempRunDir, 'restore_rehearsal', path.basename(String(high.oss_key || 'archive.zip')) || 'archive.zip')
setCurrentRun({ stage: 'restore_download', message: '正在回读高频 ZIP', percent: 35 })
summary.restore.downloaded_bytes = await retryWrap(cfg, runId, 'restore_download_failed', '恢复演练回读 ZIP', async () => {
return await downloadObjectToFile({ oss, ossKey: String(high.oss_key || ''), targetPath: downloadPath })
})
setCurrentRun({ stage: 'restore_verify_hash', message: '正在校验归档 hash', percent: 65 })
const hash = await sha256File(downloadPath)
summary.restore.sha256_match = !high.sha256 || hash === high.sha256
if (!summary.restore.sha256_match) throw new Error('恢复演练失败:下载归档 hash 与原备份记录不一致')
setCurrentRun({ stage: 'restore_open_zip', message: '正在检查归档内容', percent: 82 })
const directory = await unzipper.Open.file(downloadPath)
const files = Array.isArray(directory.files) ? directory.files : []
const manifestEntry = files.find(file => String(file.path || '') === 'manifest.json')
summary.restore.manifest_present = !!manifestEntry
if (!manifestEntry) throw new Error('恢复演练失败:归档缺少 manifest.json')
const restoreDir = path.join(tempRunDir, 'restore_rehearsal', 'extracted')
ensureDir(restoreDir)
summary.restore.restore_dir = restoreDir
const sampleEntries = files.filter(file => file.type !== 'Directory').slice(0, 3)
for (const entry of sampleEntries) {
const outPath = path.join(restoreDir, String(entry.path || '').replace(/\//g, path.sep))
ensureDir(path.dirname(outPath))
await new Promise((resolve, reject) => {
entry.stream()
.pipe(fs.createWriteStream(outPath))
.on('finish', resolve)
.on('error', reject)
})
}
summary.restore.extracted_entries = sampleEntries.map(entry => String(entry.path || ''))
addEvent(runId, 'info', 'restore_rehearsal_ok', '恢复演练完成', {
source_run_id: latestHigh.row.id,
extracted_entries: summary.restore.extracted_entries
})
}
const createRun = (mode, triggerSource) => {
const runKey = makeRunKey(mode)
const info = insertRunStmt.run({
run_key: runKey,
mode,
trigger_source: triggerSource,
status: 'running',
stage: 'init',
summary_json: JSON.stringify(runSummaryBase())
})
return { id: Number(info.lastInsertRowid || 0), runKey }
}
const finishRun = (runId, status, stage, summary, errorText) => {
updateRunStmt.run({
id: runId,
status: String(status || 'done'),
stage: String(stage || ''),
ended_at: formatBeijingDateTime(new Date()),
summary_json: JSON.stringify(summary || runSummaryBase()),
error_text: errorText ? String(errorText) : ''
})
}
const cleanupTempDir = dir => {
try { fs.rmSync(dir, { recursive: true, force: true }) } catch {}
}
const executeRun = async (mode, triggerSource) => {
const run = createRun(mode, triggerSource)
const summary = runSummaryBase()
const tempRunDir = path.join(TEMP_DIR, run.runKey)
ensureDir(tempRunDir)
setCurrentRun({
run_id: run.id,
run_key: run.runKey,
mode,
trigger_source: triggerSource,
status: 'running',
stage: 'init',
message: '任务已启动',
percent: 1,
started_at: formatBeijingDateTime(new Date()),
stats: {}
})
addEvent(run.id, 'info', 'run_start', `开始执行备份:${mode} (${triggerSource})`)
try {
const cfg = runtimeConfig
const oss = resolveOssSettings(cfg)
if (!cfg.oss.enabled) throw new Error('OSS 备份未启用')
if (!oss.ready) throw new Error('OSS 凭证未就绪,请检查本地凭证文件与配置')
if (mode === 'high' || mode === 'all') await runHighfreq(cfg, oss, run.id, tempRunDir, summary)
if (mode === 'low' || mode === 'all') await runLowfreq(cfg, oss, run.id, summary)
if (mode === 'verify') await runVerify(cfg, oss, run.id, tempRunDir, summary)
if (mode === 'restore') await runRestoreRehearsal(cfg, oss, run.id, tempRunDir, summary)
let finalStatus = 'success'
if (mode === 'high' || mode === 'low' || mode === 'all') {
finalStatus = summary.lowfreq.failed > 0 ? 'partial_success' : 'success'
} else if (mode === 'verify') {
finalStatus = summary.verify.failed > 0 ? 'failed' : 'success'
} else if (mode === 'restore') {
finalStatus = (summary.restore.sha256_match && summary.restore.manifest_present) ? 'success' : 'failed'
}
finishRun(run.id, finalStatus, 'done', summary, '')
setCurrentRun({
status: finalStatus,
stage: 'done',
message: finalStatus === 'success'
? (mode === 'verify' ? '回读抽检完成' : (mode === 'restore' ? '恢复演练完成' : '备份完成'))
: (mode === 'all' || mode === 'low' || mode === 'high' ? '备份完成,存在部分失败' : '任务失败'),
percent: 100,
stats: summary
})
addEvent(run.id, finalStatus === 'success' ? 'info' : 'warning', 'run_done', `备份执行完成:${finalStatus}`)
} catch (error) {
const errText = String(error && error.message ? error.message : error)
finishRun(run.id, 'failed', 'failed', summary, errText)
setCurrentRun({
status: 'failed',
stage: 'failed',
message: errText,
percent: currentRunState.percent || 0,
stats: summary
})
addEvent(run.id, 'error', 'run_failed', `备份执行失败:${errText}`)
addAlert(run.id, 'critical', 'BOX容灾备份失败', errText)
} finally {
cleanupTempDir(tempRunDir)
setTimeout(() => {
if (currentRunState.run_id === run.id && currentRunState.status !== 'running') resetCurrentRun()
}, 15000)
}
}
const triggerRun = async (mode, triggerSource = 'manual') => {
const normalized = ['high', 'low', 'all', 'verify', 'restore'].includes(String(mode || '')) ? mode : 'all'
if (currentRunPromise) {
return {
ok: false,
busy: true,
current_run: {
run_id: currentRunState.run_id,
mode: currentRunState.mode,
stage: currentRunState.stage,
status: currentRunState.status
}
}
}
currentRunPromise = executeRun(normalized, triggerSource)
currentRunPromise.finally(() => { currentRunPromise = null })
return {
ok: true,
started: true,
mode: normalized,
run_id: currentRunState.run_id,
run_key: currentRunState.run_key
}
}
const latestRunByMode = mode => db.prepare(`
SELECT * FROM backup_runs
WHERE mode IN (${mode === 'high' ? `'high','all'` : `'low','all'`})
AND status <> 'running'
ORDER BY id DESC LIMIT 1
`).get()
const latestRunStatus = mode => {
const row = latestRunByExactMode(mode)
if (!row) return null
return {
last_run_at: row.ended_at || row.started_at || '',
last_status: row.status,
stage: row.stage || '',
error_text: row.error_text || '',
summary: parseRunSummary(row)
}
}
const latestManualRecord = type => db.prepare(`
SELECT * FROM backup_manual_records WHERE record_type=? ORDER BY id DESC LIMIT 1
`).get(type)
const getOverview = () => {
const cfg = runtimeConfig
const oss = resolveOssSettings(cfg)
const high = latestRunByMode('high')
const low = latestRunByMode('low')
const verify = latestRunStatus('verify')
const restore = latestRunStatus('restore')
const manual = latestManualRecord('baidu_manual')
return {
ok: true,
current_run: clone(currentRunState),
scheduler: {
enabled: !!cfg.scheduler.enabled,
highfreq_cron: cfg.scheduler.highfreq_cron,
lowfreq_cron: cfg.scheduler.lowfreq_cron,
next_highfreq_at: getNextCronTime(cfg.scheduler.highfreq_cron),
next_lowfreq_at: getNextCronTime(cfg.scheduler.lowfreq_cron)
},
highfreq: high ? {
last_run_at: high.ended_at || high.started_at || '',
last_status: high.status,
stage: high.stage || '',
error_text: high.error_text || ''
} : null,
lowfreq: low ? {
last_run_at: low.ended_at || low.started_at || '',
last_status: low.status,
stage: low.stage || '',
error_text: low.error_text || ''
} : null,
verify,
restore,
manual_baidu_backup: manual ? {
last_recorded_at: manual.created_at || '',
note: manual.note || ''
} : null,
alerts_open_count: db.prepare(`SELECT COUNT(1) AS c FROM backup_alerts WHERE status='open'`).get().c || 0,
oss_ready: oss.ready,
credential_file: oss.credential_file || ''
}
}
const listRuns = limit => {
const rows = db.prepare(`
SELECT id, run_key, mode, trigger_source, status, stage, started_at, ended_at, error_text, summary_json
FROM backup_runs
ORDER BY id DESC
LIMIT ?
`).all(Math.max(1, parseInt(String(limit || '50'), 10) || 50))
return rows.map(row => ({
...row,
summary: (() => {
try { return row.summary_json ? JSON.parse(row.summary_json) : {} } catch { return {} }
})()
}))
}
const listEvents = ({ limit = 200, failuresOnly = false } = {}) => {
let sql = `
SELECT e.id, e.run_id, e.level, e.event_type, e.message, e.detail_json, e.created_at,
r.mode, r.trigger_source
FROM backup_run_events e
LEFT JOIN backup_runs r ON r.id = e.run_id
`
const args = []
if (failuresOnly) sql += ` WHERE e.level='error' OR e.event_type IN ('retry','run_failed','highfreq_upload_failed','lowfreq_upload_failed') `
sql += ` ORDER BY e.id DESC LIMIT ? `
args.push(Math.max(1, parseInt(String(limit || '200'), 10) || 200))
return db.prepare(sql).all(...args).map(row => ({
...row,
detail: (() => {
try { return row.detail_json ? JSON.parse(row.detail_json) : null } catch { return null }
})()
}))
}
const listAlerts = limit => db.prepare(`
SELECT id, run_id, severity, title, message, status, created_at
FROM backup_alerts
ORDER BY id DESC
LIMIT ?
`).all(Math.max(1, parseInt(String(limit || '100'), 10) || 100))
const listLowfreqObjects = ({ limit = 200, status = '' } = {}) => {
let sql = `
SELECT id, backup_type, relative_path, source_path, oss_key, size_bytes, mtime_ms, sha256,
upload_status, retry_count, last_error, last_uploaded_at, is_deleted, deleted_at, meta_json
FROM backup_objects
WHERE backup_type='lowfreq'
`
const args = []
const normalizedStatus = String(status || '').trim()
if (normalizedStatus) {
sql += ` AND upload_status=? `
args.push(normalizedStatus)
}
sql += ` ORDER BY id DESC LIMIT ? `
args.push(Math.max(1, parseInt(String(limit || '200'), 10) || 200))
return db.prepare(sql).all(...args).map(row => ({
...row,
meta: parseObjectMeta(row),
encryption: getEncryptedFileInfo(row)
}))
}
const listHighfreqCatalog = limit => {
const rows = db.prepare(`
SELECT id, run_key, mode, trigger_source, status, stage, started_at, ended_at, error_text, summary_json
FROM backup_runs
WHERE mode IN ('high', 'all')
ORDER BY id DESC
LIMIT ?
`).all(Math.max(1, parseInt(String(limit || '50'), 10) || 50))
return rows
.map(row => ({ row, summary: parseRunSummary(row) }))
.filter(item => item.summary && item.summary.highfreq && item.summary.highfreq.oss_key)
.map(item => ({
id: item.row.id,
run_key: item.row.run_key,
mode: item.row.mode,
trigger_source: item.row.trigger_source,
status: item.row.status,
stage: item.row.stage,
started_at: item.row.started_at,
ended_at: item.row.ended_at,
error_text: item.row.error_text || '',
highfreq: item.summary.highfreq || {}
}))
}
const getBackupCatalog = ({ runLimit = 30, objectLimit = 200, status = '' } = {}) => ({
ok: true,
highfreq_runs: listHighfreqCatalog(runLimit),
lowfreq_objects: listLowfreqObjects({ limit: objectLimit, status })
})
const createHighfreqDownload = async runId => {
const row = db.prepare(`SELECT * FROM backup_runs WHERE id=?`).get(parseInt(String(runId || '0'), 10) || 0)
if (!row) throw new Error('not_found')
const summary = parseRunSummary(row)
const high = summary.highfreq || {}
const ossKey = String(high.oss_key || '').trim()
if (!ossKey) throw new Error('highfreq_missing')
const oss = resolveOssSettings(runtimeConfig)
if (!oss.ready) throw new Error('oss_not_ready')
const streamResult = await createOssDownloadStream({ oss, ossKey })
const baseName = sanitizeOssName(`box_backup_highfreq_${high.run_guid || row.id}.zip`) || `box_backup_highfreq_${row.id}.zip`
return {
file_name: baseName,
content_type: 'application/zip',
content_length: streamResult.headers['content-length'] || '',
oss_key: ossKey,
stream: streamResult.stream
}
}
const createLowfreqDownload = async objectId => {
const row = db.prepare(`SELECT * FROM backup_objects WHERE id=? AND backup_type='lowfreq'`).get(parseInt(String(objectId || '0'), 10) || 0)
if (!row) throw new Error('not_found')
if (String(row.upload_status || '') !== 'uploaded') throw new Error('lowfreq_not_ready')
const ossKey = String(row.oss_key || '').trim()
if (!ossKey) throw new Error('lowfreq_missing')
const oss = resolveOssSettings(runtimeConfig)
if (!oss.ready) throw new Error('oss_not_ready')
const streamResult = await createOssDownloadStream({ oss, ossKey })
const meta = parseObjectMeta(row)
const originalName = path.basename(String(row.relative_path || 'backup.bin')) || 'backup.bin'
const fallbackName = sanitizeOssName(meta.stored_name || originalName) || 'backup.bin'
return {
file_name: originalName || fallbackName,
content_type: streamResult.headers['content-type'] || 'application/octet-stream',
content_length: streamResult.headers['content-length'] || '',
oss_key: ossKey,
stream: streamResult.stream
}
}
const createLowfreqDecryptedDownload = async objectId => {
const row = db.prepare(`SELECT * FROM backup_objects WHERE id=? AND backup_type='lowfreq'`).get(parseInt(String(objectId || '0'), 10) || 0)
if (!row) throw new Error('not_found')
const encryption = getEncryptedFileInfo(row)
if (!encryption.encrypted) throw new Error('decrypt_not_supported')
if (String(row.upload_status || '') !== 'uploaded') throw new Error('lowfreq_not_ready')
const ossKey = String(row.oss_key || '').trim()
if (!ossKey) throw new Error('lowfreq_missing')
const oss = resolveOssSettings(runtimeConfig)
if (!oss.ready) throw new Error('oss_not_ready')
const streamResult = await createOssDownloadStream({ oss, ossKey })
const encryptedBuffer = await streamToBuffer(streamResult.stream)
let payload = null
if (encryption.provider === 'weekly') {
payload = weeklyCore.getAttachmentDataByStoredPath(encryption.encrypted_path, encryptedBuffer)
} else if (encryption.provider === 'ai-lib') {
payload = aiLib.getAttachmentDataByStoredPath(encryption.encrypted_path, encryptedBuffer)
}
if (!payload || !payload.data) throw new Error('decrypt_source_not_found')
return {
file_name: String(payload.filename || path.basename(String(row.relative_path || 'file.bin')) || 'file.bin'),
content_type: String(payload.mime || 'application/octet-stream'),
content_length: Buffer.byteLength(payload.data),
data: payload.data,
provider: encryption.provider,
encrypted_label: encryption.label
}
}
const getConfig = () => {
const oss = resolveOssSettings(runtimeConfig)
return {
ok: true,
config: clone(runtimeConfig),
credential_status: {
ready: oss.ready,
credential_file: oss.credential_file || '',
bucket: oss.bucket || '',
endpoint: oss.endpoint || '',
root_prefix: oss.root_prefix || '',
lowfreq_dir_name: oss.lowfreq_dir_name || ''
}
}
}
const saveConfig = nextConfig => {
const sanitized = sanitizeConfig(nextConfig)
const resolvedOss = resolveOssSettings(sanitized)
if (!sanitized.scheduler.highfreq_cron || !sanitized.scheduler.lowfreq_cron) throw new Error('cron 不能为空')
if (!resolvedOss.root_prefix) throw new Error('OSS 根目录不能为空')
if (!resolvedOss.lowfreq_dir_name) throw new Error('低频固定目录不能为空')
fs.writeFileSync(CONFIG_PATH, JSON.stringify(sanitized, null, 2))
runtimeConfig = sanitized
insertConfigHistoryStmt.run(JSON.stringify(sanitized))
restartScheduler()
return getConfig()
}
const recordManualBaiduBackup = note => {
insertManualRecordStmt.run('baidu_manual', String(note || '').trim())
return latestManualRecord('baidu_manual')
}
const getCurrentRun = () => ({
ok: true,
run: clone(currentRunState)
})
const schedulerTick = async () => {
if (!runtimeConfig.scheduler.enabled) return
if (currentRunPromise) return
const now = new Date()
const p = beijingParts(now)
const minuteKey = `${p.year}-${p.month}-${p.day} ${p.hour}:${p.minute}`
if (minuteKey === lastSchedulerMinute) return
lastSchedulerMinute = minuteKey
const dueHigh = cronMatches(runtimeConfig.scheduler.highfreq_cron, now)
const dueLow = cronMatches(runtimeConfig.scheduler.lowfreq_cron, now)
if (!dueHigh && !dueLow) return
const mode = dueHigh && dueLow ? 'all' : (dueHigh ? 'high' : 'low')
await triggerRun(mode, 'cron')
}
function restartScheduler () {
if (schedulerTimer) {
clearInterval(schedulerTimer)
schedulerTimer = null
}
if (String(process.env.PORT || '') !== '8977') return
const pollMs = Math.max(10000, runtimeConfig.scheduler.poll_seconds * 1000)
schedulerTimer = setInterval(() => {
schedulerTick().catch(error => {
try { logJSON('box_backup.scheduler.error', { error: String(error && error.message ? error.message : error) }, 'box_backup') } catch {}
})
}, pollMs)
setTimeout(() => {
schedulerTick().catch(() => {})
}, 3000)
}
restartScheduler()
module.exports = {
getOverview,
getCurrentRun,
getConfig,
saveConfig,
listRuns,
listEvents,
listAlerts,
getBackupCatalog,
createHighfreqDownload,
createLowfreqDownload,
createLowfreqDecryptedDownload,
triggerRun,
recordManualBaiduBackup
}