| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449 |
- /**
- * 云函数: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'));
- }
|