/** * 云函数:systemStorageManager * * 通用业务实体、审计和迁移状态的云端入口。 * 数据主存储为 Parse 项目命名空间 Class: * - VideoWorkflowEntity * - VideoWorkflowAudit * - VideoWorkflowMigration * * 所有数据归属都通过 _session.js 从 Parse sessionToken 解析得到。 */ const { requireSession, requireAdmin, assertRequestedUserMatchesSession } = require('./_session'); const { VIDEO_WORKFLOW_CLASSES, createParseClassStore, parseUserFields, parseUserWhere, } = require('./_parseClassStore'); async function handler(request, response) { try { const action = pickParam(request, 'action') || 'list'; if (action === 'adminInspect') { const adminSession = await requireAdmin(request, Psql); const store = createParseClassStore({ sessionToken: adminSession.sessionToken }); return adminInspect(request, response, adminSession, store); } const session = await requireSession(request, Psql); assertRequestedUserMatchesSession(request, session); const store = createParseClassStore({ sessionToken: session.sessionToken }); const userId = session.userId; if (action === 'list') return listEntities(request, response, userId, store); if (action === 'get') return getEntity(request, response, userId, store); if (action === 'upsert') return upsertEntity(request, response, userId, store); if (action === 'patch') return patchEntity(request, response, userId, store); if (action === 'delete') return deleteEntity(request, response, userId, store); if (action === 'purge') return purgeEntity(request, response, userId, store); if (action === 'audit') return writeAudit(request, response, userId, store); if (action === 'auditList') return listAudits(request, response, userId, store); if (action === 'stats') return getStats(request, response, userId, store); if (action === 'migrationGet') return getMigration(request, response, userId, store); if (action === 'migrationSet') return setMigration(request, response, userId, store); return response.json({ code: 400, success: false, error: `未知 action: ${action}` }); } catch (error) { console.error('systemStorageManager failed:', error && error.message); return response.json({ code: error.status || 500, success: false, error: error.message }); } } async function listEntities(request, response, userId, store) { const entityType = requiredText(request, 'entityType', 'type'); const status = pickParam(request, 'status') || 'active'; const limit = safeLimit(pickParam(request, 'limit'), 100, 1000); const where = parseUserWhere(userId, { entityType, ...(status ? { status } : {}), }); const rows = await store.list(VIDEO_WORKFLOW_CLASSES.entity, where, { order: '-updatedAt', limit }); return response.json({ code: 200, success: true, data: rows.map(rowToEntity) }); } async function getEntity(request, response, userId, store) { const entityType = requiredText(request, 'entityType', 'type'); const entityId = requiredText(request, 'entityId', 'id', 'bizId'); const row = await findEntity(store, userId, entityType, entityId); return response.json({ code: 200, success: true, data: row ? rowToEntity(row) : null }); } async function upsertEntity(request, response, userId, store) { const entityType = requiredText(request, 'entityType', 'type'); const entityId = requiredText(request, 'entityId', 'id', 'bizId'); const data = normalizeObject(pickParam(request, 'data', 'entity') || {}); // 行状态只表达存储生命周期,业务状态继续放在 data.status,避免 completed/failed 等业务态污染刷新查询。 const status = pickParam(request, 'status') || 'active'; const schemaVersion = Number(pickParam(request, 'schemaVersion') || data.schemaVersion || 1); const spaceId = pickParam(request, 'spaceId') || data.spaceId || ''; const row = await saveEntity(store, { userId, entityType, entityId, data, status, schemaVersion, spaceId }); return response.json({ code: 200, success: true, data: rowToEntity(row) }); } async function patchEntity(request, response, userId, store) { const entityType = requiredText(request, 'entityType', 'type'); const entityId = requiredText(request, 'entityId', 'id', 'bizId'); const patch = normalizeObject(pickParam(request, 'patch', 'data') || {}); const previousRow = await findEntity(store, userId, entityType, entityId); const previous = previousRow ? normalizeObject(previousRow.data) : {}; const status = pickParam(request, 'status') || previousRow?.status || 'active'; const row = await saveEntity(store, { userId, entityType, entityId, data: { ...previous, ...patch, id: entityId }, status, schemaVersion: Number(patch.schemaVersion || previousRow?.schemaVersion || 1), spaceId: patch.spaceId || previousRow?.spaceId || '', }); return response.json({ code: 200, success: true, data: rowToEntity(row) }); } async function deleteEntity(request, response, userId, store) { const entityType = requiredText(request, 'entityType', 'type'); const entityId = requiredText(request, 'entityId', 'id', 'bizId'); const row = await findEntity(store, userId, entityType, entityId); if (row?.objectId) { await store.update(VIDEO_WORKFLOW_CLASSES.entity, row.objectId, { status: 'deleted' }); } return response.json({ code: 200, success: true, data: { entityType, entityId, deleted: true } }); } async function purgeEntity(request, response, userId, store) { const entityType = requiredText(request, 'entityType', 'type'); const entityId = requiredText(request, 'entityId', 'id', 'bizId'); const reason = String(pickParam(request, 'reason') || '').trim(); if (!reason) { return response.json({ code: 400, success: false, error: '清理必须提供 reason' }); } const row = await findEntity(store, userId, entityType, entityId); const purgedAt = new Date().toISOString(); const data = { purged: true, purgedAt, reason, id: entityId, userId }; let purged = false; if (row?.objectId) { await store.update(VIDEO_WORKFLOW_CLASSES.entity, row.objectId, { status: 'purged', data, }); purged = true; } await createAudit(store, userId, { entityType, entityId, action: 'purge', summary: `清理 ${entityType}:${entityId}`, detail: { reason, purged }, }); return response.json({ code: 200, success: true, data: { entityType, entityId, purged }, }); } async function writeAudit(request, response, userId, store) { const action = requiredText(request, 'auditAction', 'eventAction', 'name'); const entityType = pickParam(request, 'entityType') || ''; const entityId = pickParam(request, 'entityId', 'id') || ''; const summary = String(pickParam(request, 'summary') || '').slice(0, 500); const detail = normalizeObject(pickParam(request, 'detail') || {}); const row = await createAudit(store, userId, { entityType, entityId, action, summary, detail }); return response.json({ code: 200, success: true, data: rowToAudit(row) }); } async function listAudits(request, response, userId, store) { const limit = safeLimit(pickParam(request, 'limit'), 50, 200); const entityType = String(pickParam(request, 'entityType') || '').trim(); const where = parseUserWhere(userId, entityType ? { entityType } : {}); const rows = await store.list(VIDEO_WORKFLOW_CLASSES.audit, where, { order: '-createdAt', limit }); return response.json({ code: 200, success: true, data: rows.map(rowToAudit) }); } async function getStats(request, response, userId, store) { const [entityRows, auditRows, migrationRows] = await Promise.all([ listAll(store, VIDEO_WORKFLOW_CLASSES.entity, parseUserWhere(userId, {}), 5000), listAll(store, VIDEO_WORKFLOW_CLASSES.audit, parseUserWhere(userId, {}), 1000), listAll(store, VIDEO_WORKFLOW_CLASSES.migration, parseUserWhere(userId, {}), 1000), ]); return response.json({ code: 200, success: true, data: buildStats(entityRows, auditRows, migrationRows), }); } async function adminInspect(request, response, adminSession, store) { const targetUserId = String(pickParam(request, 'targetUserId', 'inspectUserId') || '').trim(); const limit = safeLimit(pickParam(request, 'limit'), 100, 500); const userWhere = targetUserId ? { ownerId: targetUserId } : {}; const [entityRows, auditRows, migrationRows] = await Promise.all([ listAll(store, VIDEO_WORKFLOW_CLASSES.entity, parseUserWhereForAdmin(userWhere), limit * 20), listAll(store, VIDEO_WORKFLOW_CLASSES.audit, parseUserWhereForAdmin(userWhere), limit * 20), listAll(store, VIDEO_WORKFLOW_CLASSES.migration, parseUserWhereForAdmin(userWhere), limit * 20), ]); const users = {}; const ensureUser = (userId) => { if (!users[userId]) { users[userId] = { userId, totalEntities: 0, entityTypes: {}, auditCount: 0, latestAuditAt: '', migrations: { total: 0, failedRows: 0, byStatus: {} }, latestUpdatedAt: '', }; } return users[userId]; }; for (const row of entityRows) { const user = ensureUser(row.ownerId || ''); const entityType = row.entityType || ''; const status = row.status || 'active'; if (!entityType || !user.userId) continue; if (!user.entityTypes[entityType]) user.entityTypes[entityType] = { total: 0, byStatus: {} }; user.entityTypes[entityType].total += 1; user.entityTypes[entityType].byStatus[status] = (user.entityTypes[entityType].byStatus[status] || 0) + 1; user.totalEntities += 1; const latest = row.updatedAt || ''; if (latest && String(latest) > String(user.latestUpdatedAt || '')) user.latestUpdatedAt = latest; } for (const row of auditRows) { const user = ensureUser(row.ownerId || ''); if (!user.userId) continue; user.auditCount += 1; const latest = row.createdAt || ''; if (latest && String(latest) > String(user.latestAuditAt || '')) user.latestAuditAt = latest; } for (const row of migrationRows) { const user = ensureUser(row.ownerId || ''); if (!user.userId) continue; const status = row.status || 'unknown'; user.migrations.total += 1; user.migrations.byStatus[status] = (user.migrations.byStatus[status] || 0) + 1; user.migrations.failedRows += Number(row.failedCount || 0); } const rows = Object.values(users) .sort((a, b) => b.totalEntities - a.totalEntities || String(b.latestUpdatedAt || '').localeCompare(String(a.latestUpdatedAt || ''))) .slice(0, limit); return response.json({ code: 200, success: true, data: { adminUserId: adminSession.userId, targetUserId, users: rows, generatedAt: new Date().toISOString(), }, }); } async function getMigration(request, response, userId, store) { const migrationKey = requiredText(request, 'migrationKey'); const sourceKey = requiredText(request, 'sourceKey'); const row = await findMigration(store, userId, migrationKey, sourceKey); return response.json({ code: 200, success: true, data: row ? rowToMigration(row) : null }); } async function setMigration(request, response, userId, store) { const migrationKey = requiredText(request, 'migrationKey'); const sourceKey = requiredText(request, 'sourceKey'); const status = pickParam(request, 'status') || 'completed'; const migratedCount = Number(pickParam(request, 'migratedCount') || 0); const failedCount = Number(pickParam(request, 'failedCount') || 0); const detail = normalizeObject(pickParam(request, 'detail') || {}); const where = parseUserWhere(userId, { migrationKey, sourceKey }); const payload = { ...parseUserFields(userId), migrationKey, sourceKey, status, migratedCount, failedCount, detail, }; const row = await upsertMigrationWithDetailFallback(store, where, payload); return response.json({ code: 200, success: true, data: rowToMigration(row) }); } async function upsertMigrationWithDetailFallback(store, where, payload) { try { return await store.upsertByQuery(VIDEO_WORKFLOW_CLASSES.migration, where, payload); } catch (error) { if (!isParseFieldTypeMismatch(error) || typeof payload.detail === 'string') { throw error; } return store.upsertByQuery(VIDEO_WORKFLOW_CLASSES.migration, where, { ...payload, detail: safeStringify(payload.detail), }); } } async function saveEntity(store, { userId, entityType, entityId, data, status, schemaVersion, spaceId }) { const rawData = normalizeObject(data); const payload = { ...rawData, id: rawData.id || entityId, userId }; const where = parseUserWhere(userId, { entityType, entityId }); return store.upsertByQuery(VIDEO_WORKFLOW_CLASSES.entity, where, { ...parseUserFields(userId), entityType, entityId, status: status || 'active', schemaVersion: schemaVersion || 1, spaceId: spaceId || '', data: payload, }); } async function findEntity(store, userId, entityType, entityId) { return store.findFirst(VIDEO_WORKFLOW_CLASSES.entity, parseUserWhere(userId, { entityType, entityId }), { limit: 1 }); } async function findMigration(store, userId, migrationKey, sourceKey) { return store.findFirst(VIDEO_WORKFLOW_CLASSES.migration, parseUserWhere(userId, { migrationKey, sourceKey }), { limit: 1 }); } async function createAudit(store, userId, { entityType, entityId, action, summary, detail }) { return store.create(VIDEO_WORKFLOW_CLASSES.audit, { ...parseUserFields(userId), entityType: entityType || '', entityId: entityId || '', action, summary: String(summary || '').slice(0, 500), detail: safeStringify(detail || {}), }); } async function listAll(store, className, where, maxRows) { const all = []; const pageSize = 1000; for (let skip = 0; skip < maxRows; skip += pageSize) { const rows = await store.list(className, where, { order: '-updatedAt', limit: Math.min(pageSize, maxRows - skip), skip }); all.push(...rows); if (rows.length < pageSize || all.length >= maxRows) break; } return all.slice(0, maxRows); } function buildStats(entityRows, auditRows, migrationRows) { const entityTypes = {}; let totalEntities = 0; for (const row of entityRows) { const entityType = row.entityType || ''; const status = row.status || 'active'; if (!entityType) continue; if (!entityTypes[entityType]) entityTypes[entityType] = { total: 0, byStatus: {} }; entityTypes[entityType].total += 1; entityTypes[entityType].byStatus[status] = (entityTypes[entityType].byStatus[status] || 0) + 1; totalEntities += 1; } const migrations = { total: 0, byStatus: {} }; for (const row of migrationRows) { const status = row.status || 'unknown'; migrations.total += 1; migrations.byStatus[status] = (migrations.byStatus[status] || 0) + 1; } return { totalEntities, entityTypes, auditCount: auditRows.length, migrations, generatedAt: new Date().toISOString(), }; } function parseUserWhereForAdmin(extra = {}) { return { projectKey: 'video-workflow', ...extra }; } function rowToEntity(row) { const data = normalizeObject(row.data); return { id: row.entityId, entityId: row.entityId, type: row.entityType, entityType: row.entityType, status: row.status || 'active', schemaVersion: row.schemaVersion, spaceId: row.spaceId || '', data, createdAt: row.createdAt, updatedAt: row.updatedAt, }; } function rowToMigration(row) { return { id: row.objectId, migrationKey: row.migrationKey, sourceKey: row.sourceKey, status: row.status, migratedCount: Number(row.migratedCount || 0), failedCount: Number(row.failedCount || 0), detail: normalizeObject(row.detail), createdAt: row.createdAt, updatedAt: row.updatedAt, }; } function rowToAudit(row) { return { id: row.objectId, entityType: row.entityType || '', entityId: row.entityId || '', action: row.action, summary: row.summary || '', detail: normalizeObject(row.detail), createdAt: row.createdAt, }; } function requiredText(request, ...names) { const value = pickParam(request, ...names); const text = String(value || '').trim(); if (!text) { const error = new Error(`缺少 ${names[0]}`); error.status = 400; throw error; } return text; } function pickParam(request, ...names) { const sources = [request.params, request.body, request.query, request]; for (const src of sources) { if (!src || typeof src !== 'object') continue; for (const name of names) { const value = src[name]; if (value !== undefined && value !== null && value !== '') return value; } } return null; } function normalizeObject(value) { if (typeof value === 'string') { try { return JSON.parse(value); } catch { return {}; } } if (!value || typeof value !== 'object' || Array.isArray(value)) return {}; return value; } function safeLimit(value, fallback, max) { const parsed = parseInt(value || fallback, 10); return Number.isFinite(parsed) ? Math.min(Math.max(parsed, 1), max) : fallback; } function safeStringify(value) { const text = JSON.stringify(normalizeObject(value)); return text.length > 8000 ? `${text.slice(0, 7997)}...` : text; } function isParseFieldTypeMismatch(error) { const message = String(error?.message || error?.detail?.error || '').toLowerCase(); return Number(error?.status || 0) === 400 && (message.includes('schema') || message.includes('type') || message.includes('expected')); }