/** * 云函数共享的 Parse Class 存储工具。 * * 本项目的结构化业务主数据写入带 VideoWorkflow 命名空间的 Parse Class。 * 这里封装 REST 调用、项目命名空间字段和 owner Pointer,避免业务云函数 * 继续直接依赖 PSQL。 */ const PARSE_CLASS_PROJECT_KEY = 'video-workflow'; const VIDEO_WORKFLOW_CLASSES = { entity: 'VideoWorkflowEntity', audit: 'VideoWorkflowAudit', migration: 'VideoWorkflowMigration', fileAsset: 'VideoWorkflowFileAsset', }; function createParseClassStore(options = {}) { const parseApiHost = normalizeParseApiHost(options.parseApiHost || parseStoreEnv('PARSE_API_HOST', parseStoreEnv('PARSE_BASE_URL', 'https://server.fmode.cn'))); const appId = options.appId || parseStoreEnv('PARSE_APP_ID', parseStoreEnv('PARSE_APPLICATION_ID', 'ncloudmaster')); const sessionToken = options.sessionToken || ''; return { findFirst: (className, where, opts = {}) => parseFindFirst({ parseApiHost, appId, sessionToken, className, where, opts }), list: (className, where, opts = {}) => parseList({ parseApiHost, appId, sessionToken, className, where, opts }), create: (className, data) => parseCreate({ parseApiHost, appId, sessionToken, className, data }), update: (className, objectId, patch) => parseUpdate({ parseApiHost, appId, sessionToken, className, objectId, patch }), deleteRecord: (className, objectId) => parseDelete({ parseApiHost, appId, sessionToken, className, objectId }), upsertByQuery: (className, where, data, opts = {}) => parseUpsertByQuery({ parseApiHost, appId, sessionToken, className, where, data, opts }), userFields: (userId) => parseUserFields(userId), projectKey: PARSE_CLASS_PROJECT_KEY, }; } function parseUserFields(userId) { return { projectKey: PARSE_CLASS_PROJECT_KEY, owner: { __type: 'Pointer', className: '_User', objectId: String(userId || '') }, ownerId: String(userId || ''), }; } function parseUserWhere(userId, extra = {}) { return { projectKey: PARSE_CLASS_PROJECT_KEY, owner: { __type: 'Pointer', className: '_User', objectId: String(userId || '') }, ...extra, }; } async function parseFindFirst(ctx) { const rows = await parseList({ ...ctx, opts: { ...ctx.opts, limit: 1 } }); return rows[0] || null; } async function parseList({ parseApiHost, appId, sessionToken, className, where, opts = {} }) { const query = new URLSearchParams(); query.set('where', JSON.stringify(where || {})); query.set('limit', String(safeParseLimit(opts.limit, 100, 1000))); if (opts.skip !== undefined) query.set('skip', String(Math.max(0, parseInt(opts.skip || 0, 10) || 0))); if (opts.order) query.set('order', String(opts.order)); if (opts.include) query.set('include', String(opts.include)); if (opts.keys) query.set('keys', String(opts.keys)); const data = await parseRequest({ parseApiHost, appId, sessionToken, method: 'GET', path: `/classes/${encodeURIComponent(className)}?${query.toString()}`, }); return Array.isArray(data.results) ? data.results : []; } async function parseCreate({ parseApiHost, appId, sessionToken, className, data }) { const created = await parseRequest({ parseApiHost, appId, sessionToken, method: 'POST', path: `/classes/${encodeURIComponent(className)}`, body: sanitizeParseWrite(data), }); return { ...sanitizeParseWrite(data), ...created }; } async function parseUpdate({ parseApiHost, appId, sessionToken, className, objectId, patch }) { await parseRequest({ parseApiHost, appId, sessionToken, method: 'PUT', path: `/classes/${encodeURIComponent(className)}/${encodeURIComponent(objectId)}`, body: sanitizeParseWrite(patch), }); return parseFindFirst({ parseApiHost, appId, sessionToken, className, where: { objectId }, opts: { limit: 1 }, }); } async function parseDelete({ parseApiHost, appId, sessionToken, className, objectId }) { return parseRequest({ parseApiHost, appId, sessionToken, method: 'DELETE', path: `/classes/${encodeURIComponent(className)}/${encodeURIComponent(objectId)}`, }); } async function parseUpsertByQuery({ parseApiHost, appId, sessionToken, className, where, data, opts = {} }) { const findExisting = () => parseFindFirst({ parseApiHost, appId, sessionToken, className, where, opts: { order: opts.order || '-updatedAt', limit: 1 }, }); const existing = await findExisting(); if (existing && existing.objectId) { return parseUpdate({ parseApiHost, appId, sessionToken, className, objectId: existing.objectId, patch: data, }); } try { return await parseCreate({ parseApiHost, appId, sessionToken, className, data: { ...where, ...data } }); } catch (error) { if (!isRetryableParseError(error)) throw error; await parseStoreSleep(500); const created = await findExisting(); if (created && created.objectId) { return parseUpdate({ parseApiHost, appId, sessionToken, className, objectId: created.objectId, patch: data, }); } return parseCreate({ parseApiHost, appId, sessionToken, className, data: { ...where, ...data } }); } } async function parseRequest({ parseApiHost, appId, sessionToken, method, path, body }) { const headers = { 'Accept': 'application/json', 'X-Parse-Application-Id': appId, }; if (sessionToken) headers['X-Parse-Session-Token'] = sessionToken; if (body !== undefined) headers['Content-Type'] = 'application/json'; const maxAttempts = isSafeParseRetryMethod(method) ? 3 : 1; let lastError = null; for (let attempt = 1; attempt <= maxAttempts; attempt += 1) { try { const resp = await fetch(`${parseApiHost}/parse${path}`, { method, headers, body: body !== undefined ? JSON.stringify(body) : undefined, }); const data = await resp.json().catch(() => ({})); if (!resp.ok || data.error) { const error = new Error(data.error || data.message || `Parse ${method} ${path} failed`); error.status = resp.status || 500; error.detail = data; if (attempt < maxAttempts && isRetryableParseError(error)) { await parseStoreSleep(350 * attempt); continue; } throw error; } return data; } catch (error) { lastError = error; if (attempt < maxAttempts && isRetryableParseError(error)) { await parseStoreSleep(350 * attempt); continue; } throw error; } } throw lastError || new Error(`Parse ${method} ${path} failed`); } function sanitizeParseWrite(value) { const source = normalizeParseObject(value); const next = {}; for (const [key, item] of Object.entries(source)) { if (['objectId', 'createdAt', 'updatedAt'].includes(key)) continue; if (item === undefined) continue; next[key] = item; } return next; } function normalizeParseObject(value) { if (!value || typeof value !== 'object' || Array.isArray(value)) return {}; return value; } function normalizeParseApiHost(value) { return String(value || 'https://server.fmode.cn').replace(/\/+$/, '').replace(/\/parse$/i, ''); } function safeParseLimit(value, fallback, max) { const parsed = parseInt(value || fallback, 10); return Number.isFinite(parsed) ? Math.min(Math.max(parsed, 1), max) : fallback; } function parseStoreEnv(name, fallback = '') { if (typeof process !== 'undefined' && process.env && process.env[name] !== undefined) { return process.env[name]; } return fallback; } function isSafeParseRetryMethod(method) { return ['GET', 'PUT', 'DELETE'].includes(String(method || '').toUpperCase()); } function isRetryableParseError(error) { const status = Number(error?.status || 0); const message = `${error?.message || ''} ${error?.detail?.error || ''} ${error?.detail?.message || ''}`; return status >= 500 || /fetch failed|Failed to fetch|NetworkError|Load failed|ECONNRESET|ETIMEDOUT|EAI_AGAIN/i.test(message); } function parseStoreSleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); } module.exports = { PARSE_CLASS_PROJECT_KEY, VIDEO_WORKFLOW_CLASSES, createParseClassStore, parseUserFields, parseUserWhere, };